Annotation of coherent/f/etc/conf/streams/src/strhead.c, revision 1.1.1.1

1.1       root        1: #define        _DDI_DKI        1
                      2: #define        _SYSV4          1
                      3: 
                      4: /*
                      5:  * This file contains the stream head put and service routines.
                      6:  */
                      7: 
                      8: #include <common/ccompat.h>
                      9: #include <kernel/strmlib.h>
                     10: #include <sys/debug.h>
                     11: #include <sys/types.h>
                     12: #include <sys/stream.h>
                     13: #include <sys/stropts.h>
                     14: #include <sys/poll.h>
                     15: #include <sys/errno.h>
                     16: #include <sys/signal.h>
                     17: #include <stddef.h>
                     18: 
                     19: 
                     20: #define        QUEUE_POLL(q)           ((shead_t *) (q)->q_ptr)->sh_pollhead
                     21: 
                     22: 
                     23: /*
                     24:  * This function processes a "set options" request for a stream.
                     25:  */
                     26: 
                     27: #if    __USE_PROTO__
                     28: __LOCAL__ void (STREAM_SETOPT) (queue_t * q, struct stroptions * so)
                     29: #else
                     30: __LOCAL__ void
                     31: STREAM_SETOPT __ARGS ((q, so))
                     32: queue_t              * q;
                     33: struct stroptions
                     34:              * so;
                     35: #endif
                     36: {
                     37:        shead_t       * sheadp;
                     38:        ulong_t         flags;
                     39:        pl_t            prev_pl;
                     40: 
                     41:        ASSERT (q != NULL);
                     42:        ASSERT (so != NULL);
                     43: 
                     44:        sheadp = (shead_t *) q->q_ptr;
                     45:        flags = so->so_flags;
                     46: 
                     47:        prev_pl = SHEAD_LOCK (sheadp);
                     48: 
                     49:        if ((flags & SO_READOPT) != 0) {
                     50:                /*
                     51:                 * Set the read option flags. If multiple of the flag values
                     52:                 * are specified, collapse to one.
                     53:                 */
                     54: 
                     55:                SHEAD_SRDOPT (sheadp, so->so_readopt);
                     56:        }
                     57: 
                     58:        if ((flags & SO_WROFF) != 0)
                     59:                sheadp->sh_wroff = so->so_wroff;
                     60: 
                     61:        if ((flags & SO_MREADOFF) != 0)
                     62:                sheadp->sh_flags &= ~ SH_READMSG;
                     63:        else if ((flags & SO_MREADON) != 0)
                     64:                sheadp->sh_flags |= SH_READMSG;
                     65: 
                     66:        if ((flags & SO_NDELOFF) != 0)
                     67:                sheadp->sh_flags &= ~ SH_NDELAY;
                     68:        else if ((flags & SO_NDELON) != 0)
                     69:                sheadp->sh_flags |= SH_NDELAY;
                     70: 
                     71:        if ((flags & SO_ISNTTY) != 0)
                     72:                sheadp->sh_flags &= ~ SH_TTY;
                     73:        else if ((flags & SO_ISTTY) != 0)
                     74:                sheadp->sh_flags |= SH_TTY;
                     75: 
                     76:        if ((flags & SO_TONSTOP) != 0)
                     77:                sheadp->sh_flags &= ~ SH_TOSTOP;
                     78:        else if ((flags & SO_TOSTOP) != 0)
                     79:                sheadp->sh_flags |= SH_TOSTOP;
                     80: 
                     81:        SHEAD_UNLOCK (sheadp, prev_pl);
                     82: 
                     83: 
                     84:        prev_pl = QFREEZE_TRACE (q, "STREAM_SETOPT");
                     85: 
                     86:        if ((flags & SO_MINPSZ) != 0)
                     87:                q->q_minpsz = so->so_minpsz;
                     88: 
                     89:        if ((flags & SO_MAXPSZ) != 0)
                     90:                q->q_maxpsz = so->so_maxpsz;
                     91: 
                     92:        QUNFREEZE_TRACE (q, prev_pl);
                     93: 
                     94:        QBAND_SETOPT (q, so);
                     95: }
                     96: 
                     97: 
                     98: /*
                     99:  * The read side put procedure, where data enters the stream head. Process
                    100:  * the priority messages, queue the normal ones for processes to grab.
                    101:  *
                    102:  * We use the stream head service procedure purely for doing the less time-
                    103:  * critical stuff such as waking up processes waiting for data and dealing
                    104:  * with polling.
                    105:  */
                    106: 
                    107: #if    __USE_PROTO__
                    108: __LOCAL__ void headrput (queue_t * q, mblk_t * mp)
                    109: #else
                    110: __LOCAL__ void
                    111: headrput (q, mp)
                    112: queue_t              * q;
                    113: mblk_t       * mp;
                    114: #endif
                    115: {
                    116:        shead_t       * sheadp = (shead_t *) q->q_ptr;
                    117:        int             tmp;
                    118:         pl_t           prev_pl;
                    119: 
                    120:        switch (mp->b_datap->db_type) {
                    121: 
                    122:        case M_PCPROTO:
                    123:                /*
                    124:                 * Basically like normal data/protocol messages, but we
                    125:                 * discard later arrivals if there is already such a message
                    126:                 * queued at the stream head.
                    127:                 */
                    128: 
                    129:                prev_pl = QFREEZE_TRACE (q, "headrput");
                    130: 
                    131:                tmp = q->q_first->b_datap->db_type == M_PCPROTO;
                    132: 
                    133:                QUNFREEZE_TRACE (q, prev_pl);
                    134: 
                    135:                if (tmp)
                    136:                        break;
                    137: 
                    138:                /* FALL INTO */
                    139:        case M_PROTO:
                    140:        case M_DATA:
                    141:                putq (q, mp);
                    142:                mp = NULL;
                    143:                break;
                    144: 
                    145:        case M_PCSIG:
                    146:                /*
                    147:                 * Priority signals get generated straight away.
                    148:                 */
                    149: 
                    150:                SHEAD_SIGNAL (sheadp, * mp->b_rptr);
                    151:                break;
                    152: 
                    153:        case M_SIG:
                    154:                /*
                    155:                 * Non-priority signals are generated in-band when the signal
                    156:                 * message is dequeued by a reading process.
                    157:                 */
                    158: 
                    159:                putq (q, mp);
                    160:                mp = NULL;
                    161:                break;
                    162: 
                    163: 
                    164:        case M_IOCACK:
                    165:        case M_IOCNAK:
                    166:                /*
                    167:                 * Send this message to the specific waiting process based
                    168:                 * on the ID field.
                    169:                 */
                    170: 
                    171:                prev_pl = SHEAD_LOCK (sheadp);
                    172: 
                    173:                if (sheadp->sh_ioc_seq ==
                    174:                            ((struct iocblk *) mp->b_rptr)->ioc_id &&
                    175:                    sheadp->sh_ioc_msg == NULL) {
                    176:                        /*
                    177:                         * Park the message and let the waiting process know
                    178:                         * about it.
                    179:                         */
                    180: 
                    181:                        sheadp->sh_ioc_msg = mp;
                    182:                        mp = NULL;
                    183:                }
                    184: 
                    185:                SHEAD_UNLOCK (sheadp, prev_pl);
                    186: 
                    187:                SHEAD_WAKE (sheadp, SH_IOCTL_WAIT);
                    188:                break;
                    189: 
                    190:        case M_IOCTL:
                    191:                /*
                    192:                 * Thanks to the operation of STREAMS pipes, we can see
                    193:                 * messages sent from the other end.
                    194:                 */
                    195: 
                    196:                mp->b_datap->db_type = M_IOCNAK;
                    197:                qreply (q, mp);
                    198:                mp = NULL;
                    199:                break;
                    200: 
                    201:        case M_FLUSH:
                    202:                /*
                    203:                 * Check to see what action we should perform;
                    204:                 */
                    205: 
                    206:                tmp = * mp->b_rptr;
                    207: 
                    208:                if ((tmp & FLUSHR) != 0) {
                    209: 
                    210:                        if ((tmp & FLUSHBAND) != 0)
                    211:                                flushband (q, mp->b_rptr [1], FLUSHALL);
                    212:                        else
                    213:                                flushq (q, FLUSHALL);
                    214: 
                    215:                        tmp = * mp->b_rptr &= ~ FLUSHR;
                    216:                }
                    217: 
                    218:                if ((tmp & FLUSHW) != 0) {
                    219: 
                    220:                        qreply (q, mp);
                    221:                        mp = NULL;
                    222:                }
                    223:                break;
                    224: 
                    225:        case M_ERROR:
                    226:                prev_pl = SHEAD_LOCK (sheadp);
                    227: 
                    228:                if (mp->b_wptr > mp->b_rptr + 1) {
                    229:                        /*
                    230:                         * Two-byte form of the M_ERROR message.
                    231:                         */
                    232: 
                    233:                        if (mp->b_rptr [0] != NOERROR)
                    234:                                sheadp->sh_rerrcode = mp->b_rptr [0];
                    235: 
                    236:                        if (mp->b_rptr [1] != NOERROR)
                    237:                                sheadp->sh_rerrcode = mp->b_rptr [1];
                    238:                } else {
                    239:                        /*
                    240:                         * One-byte form of M_ERROR; NOERROR or 0 are not
                    241:                         * valid values for the error code.
                    242:                         */
                    243: 
                    244:                        tmp = (uchar_t) mp->b_rptr [0];
                    245: 
                    246:                        if (tmp == 0 || tmp == NOERROR)
                    247:                                tmp = ENXIO;
                    248: 
                    249:                        sheadp->sh_rerrcode = tmp;
                    250:                        sheadp->sh_werrcode = tmp;
                    251:                }
                    252: 
                    253: 
                    254:                /*
                    255:                 * Since we already have the stream head locked, and because
                    256:                 * we want to wake up everyone, we call SV_BROADCAST ()
                    257:                 * directly rather than going through SHEAD_WAKE ().
                    258:                 */
                    259: 
                    260:                sheadp->sh_lock_mask &= ~ SH_WAIT_MASK;
                    261:                SV_BROADCAST (sheadp->sh_wait_sv, 0);
                    262: 
                    263:                SHEAD_UNLOCK (sheadp, prev_pl);
                    264: 
                    265:                tmp = (sheadp->sh_rerrcode != 0 ? FLUSHR : 0) |
                    266:                        (sheadp->sh_werrcode != 0 ? FLUSHW : 0);
                    267: 
                    268:                if (tmp != 0) {
                    269: 
                    270:                        pollwakeup (sheadp->sh_pollhead, POLLERR);
                    271:                        putctl1 (W (sheadp->sh_head), M_FLUSH, tmp);
                    272:                }
                    273:                break;
                    274: 
                    275:        case M_HANGUP:
                    276:                prev_pl = SHEAD_LOCK (sheadp);
                    277: 
                    278:                sheadp->sh_flags |= SH_HANGUP;
                    279: 
                    280:                /*
                    281:                 * Since we already have the stream head locked, and because
                    282:                 * we want to wake up everyone, we call SV_BROADCAST ()
                    283:                 * directly rather than going through SHEAD_WAKE ().
                    284:                 */
                    285: 
                    286:                sheadp->sh_lock_mask &= ~ SH_WAIT_MASK;
                    287:                SV_BROADCAST (sheadp->sh_wait_sv, 0);
                    288: 
                    289:                SHEAD_UNLOCK (sheadp, prev_pl);
                    290: 
                    291:                if (sheadp->sh_controller != NULL) {
                    292:                        /*
                    293:                         * If we have acquired a controlling process, send a
                    294:                         * SIGHUP to the process (and not the foreground
                    295:                         * process group).
                    296:                         */
                    297: 
                    298:                        proc_signal (sheadp->sh_controller, SIGHUP);
                    299:                }
                    300: 
                    301:                break;
                    302: 
                    303:        case M_SETOPTS:
                    304:                STREAM_SETOPT (q, (struct stroptions *) mp->b_rptr);
                    305:                break;
                    306:        }
                    307: 
                    308:        if (mp != NULL)
                    309:                freemsg (mp);
                    310: }
                    311: 
                    312: 
                    313: /*
                    314:  * The stream head read side service procedure exists to defer notification of
                    315:  * data arrivals at the stream head.
                    316:  */
                    317: 
                    318: #if    __USE_PROTO__
                    319: __LOCAL__ void headrsrv (queue_t * q)
                    320: #else
                    321: __LOCAL__ void
                    322: headrsrv (q)
                    323: queue_t              * q;
                    324: #endif
                    325: {
                    326:        mblk_t        * msg;
                    327:        pl_t            prev_pl;
                    328: 
                    329:        prev_pl = QFREEZE_TRACE (q, "headrsrv");
                    330: 
                    331:        if ((msg = q->q_first) != NULL && datamsg (msg->b_datap->db_type)) {
                    332: 
                    333:                if (! pcmsg (msg->b_datap->db_type)) {
                    334: 
                    335:                        pollwakeup (QUEUE_POLL (q), POLLIN);
                    336: 
                    337:                        pollwakeup (QUEUE_POLL (q),
                    338:                                    msg->b_band == 0 ? POLLRDNORM :
                    339:                                                       POLLRDBAND);
                    340:                } else
                    341:                        pollwakeup (QUEUE_POLL (q), POLLPRI);
                    342: 
                    343:        }
                    344: 
                    345:        QUNFREEZE_TRACE (q, prev_pl);
                    346: 
                    347:        if (msg != NULL)
                    348:                SHEAD_WAKE ((shead_t *) q->q_ptr, SH_READ_WAIT);
                    349: }
                    350: 
                    351: 
                    352: /*
                    353:  * The stream head write put procedure is not normally used; stream head
                    354:  * actions call putq () directly. However, M_FLUSH processing comes through
                    355:  * here.
                    356:  */
                    357: 
                    358: #if    __USE_PROTO__
                    359: __LOCAL__ void headwput (queue_t * q, mblk_t * mp)
                    360: #else
                    361: __LOCAL__ void
                    362: headwput (q, mp)
                    363: queue_t              * q;
                    364: mblk_t       * mp;
                    365: #endif
                    366: {
                    367:        if (mp->b_datap->db_type == M_FLUSH &&
                    368:            (* mp->b_rptr & FLUSHW) != 0) {
                    369: 
                    370:                if ((* mp->b_rptr & FLUSHBAND) != 0)
                    371:                        flushband (q, mp->b_rptr [1], FLUSHALL);
                    372:                else
                    373:                        flushq (q, FLUSHALL);
                    374: 
                    375:                pollwakeup (QUEUE_POLL (q), POLLOUT);
                    376:                SHEAD_WAKE ((shead_t *) q->q_ptr, SH_WRITE_WAIT);
                    377:        }
                    378: 
                    379:        putq (q, mp);
                    380: }
                    381: 
                    382: 
                    383: /*
                    384:  * The stream head write service procedure just forwards messages on down the
                    385:  * line, checking for (band) flow control.
                    386:  */
                    387: 
                    388: #if    __USE_PROTO__
                    389: __LOCAL__ void headwsrv (queue_t * q)
                    390: #else
                    391: __LOCAL__ void
                    392: headwsrv (q)
                    393: queue_t              * q;
                    394: #endif
                    395: {
                    396:        mblk_t        * mp;
                    397: 
                    398:        pollwakeup (QUEUE_POLL (q), POLLWRBAND);
                    399: 
                    400:        while ((mp = getq (q)) != NULL) {
                    401: 
                    402:                if (pcmsg (mp->b_datap->db_type) ||
                    403:                    (mp->b_band > 0 && bcanputnext (q, mp->b_band)) ||
                    404:                    canputnext (q))
                    405:                        putnext (q, mp);
                    406:                else {
                    407:                        putbq (q, mp);
                    408:                        return ;
                    409:                }
                    410:        }
                    411: 
                    412:        /*
                    413:         * Flow control has been relieved on this queue; let the processes
                    414:         * know about it!
                    415:         */
                    416: 
                    417:        pollwakeup (QUEUE_POLL (q), POLLOUT);
                    418:        SHEAD_WAKE ((shead_t *) q->q_ptr, SH_WRITE_WAIT);
                    419: }
                    420: 
                    421: 
                    422: /*
                    423:  * since the stream head is a queue pair like any other, we need a
                    424:  * streamtab/qinit for it.
                    425:  */
                    426: 
                    427: __LOCAL__ struct module_info headrinfo = {
                    428:        0, "standard head", 0, 256, 512, 0
                    429: };
                    430: 
                    431: __LOCAL__ struct qinit defstreamread = {
                    432:        headrput, headrsrv, NULL, NULL, NULL,
                    433:        & headrinfo, NULL
                    434: }, defstreamwrite = {
                    435:        headwput, headwsrv, NULL, NULL, NULL,
                    436:        & headrinfo, NULL
                    437: };
                    438: 
                    439: struct streamtab headinfo = {
                    440:        & defstreamread, & defstreamwrite, NULL, NULL
                    441: };
                    442: 

unix.superglobalmegacorp.com

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