|
|
1.1 root 1: /* In keeping with Mach's old chario behavior, we just ignore RECNUM. */
2:
3: #include <stddef.h>
4: #include <string.h>
5:
6: #include <machine/spl.h>
7: #include <mach/mig_errors.h>
8:
9: #include "device_reply.h"
10: #include "device_error_reply.h"
11:
12: #include "ds_oskit.h"
13: #include "ds_request.h"
14:
15:
16: static struct oskit_listener_ops listener_ops; /* forward decl */
17:
18:
19: static void
20: queue_request (device_t dev, struct pending_request *req,
21: oskit_s32_t rw, queue_t queue)
22: {
23: spl_t s;
24:
25: s = splio ();
26: simple_lock (&device_ready_queue_lock); /* locks all request queues! */
27:
28: queue_enter (queue, req, struct pending_request *, chain);
29:
30: simple_unlock (&device_ready_queue_lock);
31: splx (s);
32:
33: /* The driver's asyncio interface is responsible for being interrupt-safe. */
34: if ((dev->com.stream.listening & rw) == 0)
35: {
36: oskit_s32_t mask;
37:
38: if (dev->com.stream.listening != 0)
39: /* There is an old listener installed for just the other direction,
40: but now we are interested in both directions. */
41: oskit_asyncio_remove_listener (dev->com.stream.aio,
42: &dev->com.stream.listener);
43:
44: dev->com.stream.listener.ops = &listener_ops;
45: dev->com.stream.listening |= rw;
46: mask = oskit_asyncio_add_listener (dev->com.stream.aio,
47: &dev->com.stream.listener,
48: dev->com.stream.listening);
49: if (OSKIT_FAILED (mask))
50: panic ("asyncio_add_listener: %x", mask);
51: if (mask & rw)
52: ds_device_ready (dev);
53: }
54: }
55: #define requeue_read_request(dev, req) \
56: queue_request(dev, req, OSKIT_ASYNCIO_READABLE, &dev->com.stream.read_queue)
57: #define requeue_write_request(dev, req) \
58: queue_request(dev, req, OSKIT_ASYNCIO_WRITABLE, &dev->com.stream.write_queue)
59:
60:
61: static io_return_t
62: new_request (device_t dev, oskit_s32_t rw, queue_t queue,
63: void (*completer) (device_t, struct pending_request *),
64: ipc_port_t reply_port, mach_msg_type_name_t reply_port_type,
65: oskit_size_t count, oskit_size_t offset,
66: union device_data u)
67: {
68: struct pending_request *req;
69:
70: req = request_allocate ();
71: if (!req)
72: return D_NO_MEMORY;
73:
74: req->completer = completer;
75: req->reply_port = reply_port;
76: req->reply_port_type = reply_port_type;
77: req->count = count;
78: req->offset = offset;
79: req->data = u;
80:
81: queue_request (dev, req, rw, queue);
82:
83: /* We add a reference to the device for each live request.
84: The reference lives as long as the request structure, even
85: if it is momentarily removed and later requeued. */
86: device_reference (dev);
87:
88: return 0;
89: }
90:
91: static void
92: request_done (device_t dev, struct pending_request *req)
93: {
94: request_free (req);
95: device_deallocate (dev);
96: }
97:
98: static io_return_t
99: new_request_inband_buf (device_t dev, oskit_s32_t rw, queue_t queue,
100: void (*completer_inband) (device_t,
101: struct pending_request *),
102: void (*completer_small) (device_t,
103: struct pending_request *),
104: ipc_port_t reply_port,
105: mach_msg_type_name_t reply_port_type,
106: oskit_size_t count, oskit_size_t offset,
107: const char *data, oskit_size_t copycnt)
108: {
109: union device_data u;
110: io_return_t err;
111:
112: if (count <= IO_SMALL_MAX)
113: {
114: unsigned int i;
115: for (i = 0; i < copycnt; ++i)
116: u.small[i] = data[i];
117: err = new_request (dev, rw, queue, completer_small,
118: reply_port, reply_port_type,
119: count, offset, u);
120: }
121: else
122: {
123: u.inbando = zalloc (io_inband_zone);
124: if (u.inbando == 0)
125: return D_NO_MEMORY;
126: assert (count <= IO_INBAND_MAX);
127: memcpy (u.inband, data, copycnt);
128: err = new_request (dev, rw, queue, completer_inband,
129: reply_port, reply_port_type,
130: count, offset, u);
131: if (err)
132: zfree (io_inband_zone, u.inbando);
133: }
134:
135: return err;
136: }
137:
138:
139: /* Ascertain if this request needs a reply message. Returns nonzero iff
140: there is a live reply port. If there is no reply port or the reply port
141: has died, returns zero after cleaning up the reply port. */
142: static int
143: need_reply (struct pending_request *req)
144: {
145: if (!IP_VALID (req->reply_port))
146: return 0;
147:
148: ip_lock (req->reply_port);
149: if (ip_active (req->reply_port))
150: {
151: ip_unlock (req->reply_port);
152: return 1;
153: }
154:
155: ip_release (req->reply_port);
156: ip_check_unlock (req->reply_port);
157: return 0;
158: }
159:
160:
161: static int
162: ds_asyncio_complete_read_inband_1 (device_t dev, struct pending_request *req,
163: char *data)
164: {
165: inline void error (oskit_error_t rc)
166: {
167: if (req->offset == 0)
168: ds_device_read_inband_error_reply (req->reply_port,
169: req->reply_port_type,
170: oskit_to_mach_error (rc));
171: else
172: ds_device_read_reply_inband (req->reply_port,
173: req->reply_port_type,
174: 0, data, req->offset);
175: }
176:
177: if (need_reply (req))
178: {
179: oskit_s32_t n = oskit_asyncio_readable (dev->com.stream.aio);
180: if (OSKIT_FAILED (n))
181: error (n);
182: else
183: {
184: oskit_error_t rc;
185: oskit_u32_t nread;
186:
187: if (n > req->count - req->offset)
188: n = req->count - req->offset;
189:
190: rc = oskit_stream_read (dev->com.stream.io,
191: &data[req->offset], n, &nread);
192: if (rc == OSKIT_EWOULDBLOCK)
193: {
194: requeue_read_request (dev, req);
195: return 0;
196: }
197:
198: if (rc)
199: error (rc);
200: else
201: {
202: req->offset += nread;
203: #if 0 /* XXX */
204: if (req->offset < req->count)
205: {
206: requeue_read_request (dev, req);
207: return 0;
208: }
209: #endif
210:
211: ds_device_read_reply_inband (req->reply_port,
212: req->reply_port_type,
213: 0, data, req->offset);
214: }
215: }
216: }
217:
218: request_done (dev, req);
219: return 1;
220: }
221:
222: static void
223: ds_asyncio_complete_read_inband_small (device_t dev,
224: struct pending_request *req)
225: {
226: ds_asyncio_complete_read_inband_1 (dev, req, req->data.small);
227: }
228:
229: static void
230: ds_asyncio_complete_read_inband (device_t dev, struct pending_request *req)
231: {
232: vm_offset_t ptr = req->data.inbando;
233: if (ds_asyncio_complete_read_inband_1 (dev, req, req->data.inband))
234: zfree (io_inband_zone, ptr);
235: }
236:
237: io_return_t
238: ds_asyncio_read_inband (device_t dev, ipc_port_t reply_port,
239: mach_msg_type_name_t reply_port_type, dev_mode_t mode,
240: recnum_t recnum, int count, char *data,
241: unsigned *bytes_read)
242: {
243: oskit_error_t rc;
244: io_return_t err;
245: oskit_s32_t n;
246:
247: n = oskit_asyncio_readable (dev->com.stream.aio);
248: if (OSKIT_FAILED (n))
249: return oskit_to_mach_error (n);
250: else if (n > 0)
251: {
252: if (n > count)
253: n = count;
254: rc = oskit_stream_read (dev->com.stream.io, data, n, bytes_read);
255: if (OSKIT_FAILED (rc) && (rc != OSKIT_EWOULDBLOCK /*|| (mode & D_NOWAIT)*/))
256: return oskit_to_mach_error (rc);
257: if(rc==0)
258: if (mode & D_NOWAIT) /* return just what we got */
259: return D_SUCCESS;
260: }
261: else
262: {
263: #if 0
264: if (mode & D_NOWAIT) /* Old Mach chario doesn't do this. */
265: return D_WOULD_BLOCK;
266: #endif
267:
268: *bytes_read = 0;
269: }
270:
271: err = new_request_inband_buf (dev, OSKIT_ASYNCIO_READABLE,
272: &dev->com.stream.read_queue,
273: ds_asyncio_complete_read_inband,
274: ds_asyncio_complete_read_inband_small,
275: reply_port, reply_port_type,
276: count, *bytes_read, data, *bytes_read);
277:
278: return err ?: MIG_NO_REPLY;
279: }
280:
281:
282:
283: static int
284: ds_asyncio_complete_write_inband_1 (device_t dev,
285: struct pending_request *req, char *data)
286: {
287: oskit_u32_t wrote;
288: oskit_error_t rc;
289:
290: inline void error (oskit_error_t rc)
291: {
292: if (req->offset == 0)
293: ds_device_write_inband_error_reply (req->reply_port,
294: req->reply_port_type,
295: oskit_to_mach_error (rc));
296: else
297: ds_device_write_reply_inband (req->reply_port,
298: req->reply_port_type,
299: 0, req->offset);
300: }
301:
302: rc = oskit_stream_write (dev->com.stream.io, &data[req->offset],
303: req->count - req->offset, &wrote);
304: if (rc == OSKIT_EWOULDBLOCK)
305: {
306: requeue_write_request (dev, req);
307: return 0;
308: }
309:
310: if (rc)
311: error (rc);
312: else
313: {
314: req->offset += wrote;
315: if (req->offset < req->count)
316: {
317: requeue_write_request (dev, req);
318: return 0;
319: }
320:
321: ds_device_write_reply_inband (req->reply_port,
322: req->reply_port_type,
323: 0, req->offset);
324: }
325:
326: request_done (dev, req);
327: return 1;
328: }
329:
330: static void
331: ds_asyncio_complete_write_inband_small (device_t dev,
332: struct pending_request *req)
333: {
334: ds_asyncio_complete_write_inband_1 (dev, req, req->data.small);
335: }
336:
337: static void
338: ds_asyncio_complete_write_inband (device_t dev, struct pending_request *req)
339: {
340: vm_offset_t ptr = req->data.inbando;
341: if (ds_asyncio_complete_write_inband_1 (dev, req, req->data.inband))
342: zfree (io_inband_zone, ptr);
343: }
344:
345: io_return_t
346: ds_asyncio_write_inband (device_t dev, ipc_port_t reply_port,
347: mach_msg_type_name_t reply_port_type, dev_mode_t mode,
348: recnum_t recnum,
349: io_buf_ptr_t data, unsigned int count,
350: int *bytes_written)
351: {
352: oskit_error_t rc;
353: io_return_t err;
354: oskit_u32_t wrote;
355:
356: rc = oskit_stream_write (dev->com.stream.io, data, count, &wrote);
357: if (rc && (rc != OSKIT_EWOULDBLOCK || (mode & D_NOWAIT)))
358: return oskit_to_mach_error (rc);
359: if (rc == 0)
360: {
361: if (wrote == count || (mode & D_NOWAIT))
362: {
363: *bytes_written = wrote;
364: return D_SUCCESS;
365: }
366:
367: data += wrote;
368: count -= wrote;
369: }
370:
371: err = new_request_inband_buf (dev, OSKIT_ASYNCIO_WRITABLE,
372: &dev->com.stream.write_queue,
373: ds_asyncio_complete_write_inband,
374: ds_asyncio_complete_write_inband_small,
375: reply_port, reply_port_type,
376: count, 0, data, count);
377:
378: return err ?: MIG_NO_REPLY;
379: }
380:
381: /* Kludge just for kmsg. */
382: void
383: ds_asyncio_close (device_t dev)
384: {
385: if ((void *) dev->com_device == kmsg_stream && (dev->mode & D_READ))
386: --kmsg_readers;
387: }
388:
389: const struct device_ops asyncio_device_ops =
390: {
391: write_inband: ds_asyncio_write_inband,
392: read_inband: ds_asyncio_read_inband,
393: close: ds_asyncio_close
394: };
395:
396:
397: static device_t
398: listener_device (oskit_listener_t *listener)
399: {
400: return (device_t) ((char *) listener
401: - offsetof (struct device, com.stream.listener));
402: }
403:
404: static OSKIT_COMDECL
405: listener_query(oskit_listener_t *io, const oskit_iid_t *iid,
406: void **out_ihandle)
407: {
408: device_t dev = listener_device (io);
409:
410: if (memcmp (iid, &oskit_iunknown_iid, sizeof(*iid)) == 0 ||
411: memcmp (iid, &oskit_listener_iid, sizeof(*iid)) == 0)
412: {
413: *out_ihandle = &dev->com.stream.listener;
414: device_reference (dev);
415: return 0;
416: }
417:
418: *out_ihandle = NULL;
419: return OSKIT_E_NOINTERFACE;
420: }
421:
422: static OSKIT_COMDECL_U
423: listener_addref(oskit_listener_t *io)
424: {
425: device_t dev = listener_device (io);
426: device_reference (dev);
427: return dev->ref_count;
428: }
429:
430: static OSKIT_COMDECL_U
431: listener_release(oskit_listener_t *io)
432: {
433: device_t dev = listener_device (io);
434: oskit_u32_t n = dev->ref_count - 1;
435: device_deallocate (dev);
436: return n;
437: }
438:
439: static OSKIT_COMDECL
440: listener_notify (oskit_listener_t *io, oskit_iunknown_t *obj)
441: {
442: device_t dev = listener_device (io);
443:
444: ds_device_ready (dev);
445:
446: return 0;
447: }
448:
449: static struct oskit_listener_ops listener_ops =
450: { listener_query, listener_addref, listener_release, listener_notify };
451:
452:
453:
454: static struct pending_request *
455: dequeue_request (device_t dev, queue_t queue)
456: {
457: struct pending_request *req;
458:
459: spl_t s = splio ();
460: simple_lock (&device_ready_queue_lock); /* locks all request queues! */
461:
462: if (queue_empty (queue))
463: req = 0;
464: else
465: queue_remove_first (queue, req, struct pending_request *, chain);
466:
467: simple_unlock (&device_ready_queue_lock);
468: splx (s);
469:
470: return req;
471: }
472:
473:
474: /* This gets called (at spl0) by the io_done_thread when one of our
475: devices was put on the device_ready_queue by a listener firing. */
476:
477: void
478: ds_asyncio_ready (device_t dev)
479: {
480: oskit_s32_t avail;
481: struct pending_request *req;
482:
483: do
484: {
485: avail = oskit_asyncio_poll (dev->com.stream.aio);
486: if (OSKIT_FAILED (avail))
487: panic ("oskit_asyncio_poll: %x", avail);
488: poll_results:
489: if (avail & OSKIT_ASYNCIO_READABLE)
490: {
491: /* Pluck off a read request and call its completer.
492: This will dispatch the request and either free
493: or requeue the structure. */
494: req = dequeue_request (dev, &dev->com.stream.read_queue);
495: if (req)
496: (*req->completer) (dev, req);
497: else
498: avail &= ~OSKIT_ASYNCIO_READABLE; /* ignore availability */
499: }
500: if (avail & OSKIT_ASYNCIO_WRITABLE)
501: {
502: req = dequeue_request (dev, &dev->com.stream.write_queue);
503: if (req)
504: (*req->completer) (dev, req);
505: else
506: avail &= ~OSKIT_ASYNCIO_WRITABLE;
507: }
508: } while (avail & (OSKIT_ASYNCIO_READABLE|OSKIT_ASYNCIO_WRITABLE));
509:
510: avail = dev->com.stream.listening;
511: if (queue_empty (&dev->com.stream.read_queue))
512: avail &= ~OSKIT_ASYNCIO_READABLE;
513: if (queue_empty (&dev->com.stream.write_queue))
514: avail &= ~OSKIT_ASYNCIO_WRITABLE;
515: if (avail != dev->com.stream.listening)
516: {
517: /* We don't need the same listener any more. */
518: oskit_asyncio_remove_listener (dev->com.stream.aio,
519: &dev->com.stream.listener);
520:
521: dev->com.stream.listening = avail;
522: if (avail != 0)
523: {
524: avail = oskit_asyncio_add_listener (dev->com.stream.aio,
525: &dev->com.stream.listener,
526: dev->com.stream.listening);
527: if (OSKIT_FAILED (avail))
528: panic ("asyncio_add_listener: %x", avail);
529: if (avail)
530: goto poll_results;
531: }
532: }
533: }
This archive runs on limited infrastructure. Preserving old code on modern bandwidth. Automated agents are requested to crawl responsibly.