Annotation of OSKit-Mach/oskit/ds_asyncio.c, revision 1.1.1.1

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: }

unix.superglobalmegacorp.com

This archive runs on limited infrastructure. Preserving old code on modern bandwidth. Automated agents are requested to crawl responsibly.