Annotation of OSKit-Mach/oskit/ds_asyncio.c, revision 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.