#include "u.h" #include "../port/lib.h" #include "mem.h" #include "dat.h" #include "fns.h" #include "../port/error.h" enum { Minpow = 7, Maxpow = 16, }; struct Pool { Lock; Block* list; int had; int have; int want; int goal; }; Pool pool[Maxpow]; /* * IO queues */ typedef struct Queue Queue; struct Queue { Lock; Block* bfirst; /* buffer */ Block* blast; int len; /* bytes in queue */ int limit; /* max bytes in queue */ int state; int eof; /* number of eofs read by user */ void (*kick)(void*); /* restart output */ void* arg; /* argument to kick */ QLock rlock; /* mutex for reading processes */ Rendez rr; /* process waiting to read */ QLock wlock; /* mutex for writing processes */ Rendez wr; /* process waiting to write */ uchar* syncbuf; /* synchronous IO buffer */ int synclen; /* syncbuf length */ }; enum { /* Block.flag */ Bfilled=1, /* block filled */ /* Queue.state */ Qstarve= (1<<0), /* consumer starved */ Qmsg= (1<<1), /* message stream */ Qclosed= (1<<2), Qflow= (1<<3), }; void checkb(Block *b, char *msg) { if(b->base > b->lim) panic("checkb 0 %s %lux %lux", msg, b->base, b->lim); if(b->rp < b->base) panic("checkb 1 %s %lux %lux", msg, b->base, b->rp); if(b->wp < b->base) panic("checkb 2 %s %lux %lux", msg, b->base, b->wp); if(b->rp > b->lim) panic("checkb 3 %s %lux %lux", msg, b->rp, b->lim); if(b->wp > b->lim) panic("checkb 4 %s %lux %lux", msg, b->wp, b->lim); } void poison(Block *b) { b->next = (void*)0xdeadbabe; b->rp = (void*)0xdeadbabe; b->wp = (void*)0xdeadbabe; b->lim = (void*)0xdeadbabe; b->base = (void*)0xdeadbabe; } /* * Manage interrupt level memory allocation. */ static void iallocmgr(void) { int pow; attention = 0; spllo(); for(pow = Minpow; pow <= Maxpow; pow++) { p = &pool[pow]; /* Low pass filter */ delta = 3 * (p->had - p->have); delta = (delta/2) + p->want; if(delta < 0) { lock(p); p->have -= delta; bp = p->list; while(delta--) p->list = p->list->next; unlock(p); spllo(); while(bp) { next = bp->next; free(bp); bp = next; } splhi(); } else { spllo(); n = 0; s = sizeof(Block)+(1<base = (uchar*)addr; b->rp = b->base; b->wp = b->base; b->lim = ((uchar*)b)+size; b->size = pow; n++; } spllo(); lock(p); unlock(p); splhi(); } p->had = p->have; p->want = 0; } } void qinit(void) { } void ixsummary(void) { int pow; Pool *p; print("size have/goal\n"); for(pow = Minpow; pow <= Maxpow; pow++){ cl = &pool[pow]; print("%d %d/%d\n", 1<have, p->goal); } print("\n"); } /* * interrupt time allocation */ Block* iallocb(int size) { int pow; Block *bp; for(pow = Minpow; pow < Maxpow; pow++) if(size >= (1<list; if(p == 0) { p->want++; unlock(p); if(attention == 0) { attention++; newcallback(iallocmgr); } return 0; } p->have--; p->list = bp->next; unlock(p); bp->wp = b->base; bp->rp = b->base; bp->list = 0; bp->next = 0; return bp; } void ifreeb(Block *bp) { Pool *p; p = &pool[bp->size]; lock(p); bp->next = p->list; p->list = bp; p->have++; unlock(p); } /* * allocate queues and blocks (round data base address to 64 bit boundary) */ Block* allocb(int size) { Block *b; ulong addr; size += sizeof(Block) + 7; b = malloc(size); if(b == 0) exhausted("Blocks"); addr = (ulong)b; addr = (addr + sizeof(Block) + 7) & ~7; b->base = (uchar*)addr; b->rp = b->base; b->wp = b->base; b->lim = ((uchar*)b) + size; return b; } /* * Interrupt level copy out of a queue, return # bytes copied. If drop is * set, any bytes left in a block afer a consume are discarded. */ int qconsume(Queue *q, void *vp, int len) { Block *b; int n, dowakeup; uchar *p = vp; /* sync with qwrite */ lock(q); b = q->bfirst; if(b == 0){ q->state |= Qstarve; unlock(q); return -1; } checkb(b, "qconsume 1"); n = BLEN(b); if(n < len) len = n; memmove(p, b->rp, len); if((q->state & Qmsg) || len == n) q->bfirst = b->next; b->rp += len; q->len -= len; /* if writer flow controlled, restart */ if((q->state & Qflow) && q->len < q->limit/2){ q->state &= ~Qflow; dowakeup = 1; } else dowakeup = 0; unlock(q); if(dowakeup) wakeup(&q->wr); checkb(b, "qconsume 2"); /* discard the block if we're done with it */ if((q->state & Qmsg) || len == n) { poison(b); ifree(b); } return len; } int qpass(Queue *q, Block *b) { int s, i, len, dowakeup; s = splhi(); len = BLEN(b); /* sync with qread */ dowakeup = 0; lock(q); if(q->syncbuf){ /* synchronous communications, just copy into buffer */ if(len < q->synclen) q->synclen = len; i = q->synclen; memmove(q->syncbuf, b->rp, i); q->syncbuf = 0; /* tell reader buffer is full */ len -= i; if(len <= 0 || (q->state & Qmsg)){ unlock(q); wakeup(&q->rr); free(b); splx(s); return i; } /* queue anything that's left */ dowakeup = 1; b->rp += i; } /* no waiting receivers, room in buffer? */ if(q->len >= q->limit){ unlock(q); splx(s); return -1; } /* save in buffer */ if(q->bfirst) q->blast->next = b; else q->bfirst = b; q->blast = b; q->len += len; checkb(b, "qpass"); if(q->state & Qstarve){ q->state &= ~Qstarve; dowakeup = 1; } unlock(q); if(dowakeup){ if(q->kick) (*q->kick)(q->arg); wakeup(&q->rr); } splx(s); return len; } int qproduce(Queue *q, void *vp, int len) { Block *b; int i, dowakeup; uchar *p = vp; /* sync with qread */ dowakeup = 0; lock(q); if(q->syncbuf){ /* synchronous communications, just copy into buffer */ if(len < q->synclen) q->synclen = len; i = q->synclen; memmove(q->syncbuf, p, i); q->syncbuf = 0; /* tell reader buffer is full */ len -= i; if(len <= 0 || (q->state & Qmsg)){ unlock(q); wakeup(&q->rr); return i; } /* queue anything that's left */ dowakeup = 1; p += i; } /* no waiting receivers, room in buffer? */ if(q->len >= q->limit){ unlock(q); return -1; } /* save in buffer */ b = iallocb(len); if(b == 0){ unlock(q); return -2; } memmove(b->wp, p, len); b->wp += len; if(q->bfirst) q->blast->next = b; else q->bfirst = b; q->blast = b; q->len += len; checkb(b, "qproduce"); if(q->state & Qstarve){ q->state &= ~Qstarve; dowakeup = 1; } unlock(q); if(dowakeup){ if(q->kick) (*q->kick)(q->arg); wakeup(&q->rr); } return len; } /* * called by non-interrupt code */ Queue* qopen(int limit, int msg, void (*kick)(void*), void *arg) { Queue *q; q = malloc(sizeof(Queue)); if(q == 0) return 0; memset(q, 0, sizeof(Queue)); q->limit = limit; q->kick = kick; q->arg = arg; q->state = msg ? Qmsg : 0; q->state |= Qstarve; q->eof = 0; return q; } static int filled(void *a) { Queue *q = a; return q->syncbuf == 0; } static int notempty(void *a) { Queue *q = a; return q->bfirst != 0; } /* * read a queue. if no data is queued, post a Block * and wait on its Rendez. */ long qread(Queue *q, void *vp, int len) { Block *b; int x, n, dowakeup; uchar *p = vp; qlock(&q->rlock); if(waserror()){ /* can't let go if the buffer is in use */ if(q->syncbuf){ qlock(&q->wlock); x = splhi(); lock(q); q->syncbuf = 0; unlock(q); splx(x); qunlock(&q->wlock); } qunlock(&q->rlock); nexterror(); } /* wait for data */ for(;;){ /* sync with qwrite/qproduce */ x = splhi(); lock(q); b = q->bfirst; if(b) break; if(q->state & Qclosed){ unlock(q); splx(x); poperror(); qunlock(&q->rlock); if(++q->eof > 3) error(Ehungup); return 0; } if(globalmem(vp)){ /* just let the writer fill the buffer directly */ q->synclen = len; q->syncbuf = vp; unlock(q); splx(x); sleep(&q->rr, filled, q); len = q->synclen; poperror(); qunlock(&q->rlock); return len; } else { q->state |= Qstarve; unlock(q); splx(x); sleep(&q->rr, notempty, q); } } checkb(b, "qread 1"); /* remove a buffered block */ q->bfirst = b->next; n = BLEN(b); q->len -= n; /* if writer flow controlled, restart */ if((q->state & Qflow) && q->len < q->limit/2){ q->state &= ~Qflow; dowakeup = 1; } else dowakeup = 0; unlock(q); splx(x); /* do this outside of the lock(q)! */ if(n > len) n = len; memmove(p, b->rp, n); b->rp += n; checkb(b, "qread 2"); /* free it or put what's left on the queue */ if(b->rp >= b->wp || (q->state&Qmsg)) { poison(b); free(b); } else { x = splhi(); lock(q); b->next = q->bfirst; q->bfirst = b; q->len += BLEN(b); unlock(q); splx(x); } /* wakeup flow controlled writers (with a bit of histeresis) */ if(dowakeup) wakeup(&q->wr); poperror(); qunlock(&q->rlock); return n; } static int qnotfull(void *a) { Queue *q = a; return q->len < q->limit || (q->state & Qclosed); } /* * write to a queue. if no reader blocks are posted * queue the data. * * all copies should be outside of spl since they can fault. */ long qwrite(Queue *q, void *vp, int len, int nowait) { int n, sofar, x, dowakeup; Block *b; uchar *p = vp; dowakeup = 0; if(waserror()){ qunlock(&q->wlock); nexterror(); }; qlock(&q->wlock); sofar = 0; if(q->syncbuf){ if(len < q->synclen) sofar = len; else sofar = q->synclen; memmove(q->syncbuf, p, sofar); q->synclen = sofar; q->syncbuf = 0; wakeup(&q->rr); if(len == sofar || (q->state & Qmsg)){ qunlock(&q->wlock); poperror(); return len; } } do { n = len-sofar; if(n > 128*1024) n = 128*1024; b = allocb(n); memmove(b->wp, p+sofar, n); b->wp += n; /* flow control */ while(!qnotfull(q)){ if(nowait){ free(b); qunlock(&q->wlock); poperror(); return len; } q->state |= Qflow; sleep(&q->wr, qnotfull, q); } x = splhi(); lock(q); if(q->state & Qclosed){ unlock(q); splx(x); error(Ehungup); } checkb(b, "qwrite"); if(q->syncbuf){ /* we guessed wrong and did an extra copy */ if(n > q->synclen) n = q->synclen; memmove(q->syncbuf, b->rp, n); q->synclen = n; q->syncbuf = 0; dowakeup = 1; free(b); } else { /* we guessed right, queue it */ if(q->bfirst) q->blast->next = b; else q->bfirst = b; q->blast = b; q->len += n; if(q->state & Qstarve){ q->state &= ~Qstarve; dowakeup = 1; } } unlock(q); splx(x); if(dowakeup){ if(q->kick) (*q->kick)(q->arg); wakeup(&q->rr); } sofar += n; } while(sofar < len && (q->state & Qmsg) == 0); qunlock(&q->wlock); poperror(); return len; } /* * Mark a queue as closed. No further IO is permitted. * All blocks are released. */ void qclose(Queue *q) { int x; Block *b, *bfirst; /* mark it */ x = splhi(); lock(q); q->state |= Qclosed; bfirst = q->bfirst; q->bfirst = 0; q->len = 0; unlock(q); splx(x); /* free queued blocks */ while(bfirst){ b = bfirst->next; poison(bfirst); free(bfirst); bfirst = b; } /* wake up readers/writers */ wakeup(&q->rr); wakeup(&q->wr); } /* * Mark a queue as closed. Wakeup any readers. Don't remove queued * blocks. */ void qhangup(Queue *q) { int x; /* mark it */ x = splhi(); lock(q); q->state |= Qclosed; unlock(q); splx(x); /* wake up readers/writers */ wakeup(&q->rr); wakeup(&q->wr); } /* * mark a queue as no longer hung up */ void qreopen(Queue *q) { q->state &= ~Qclosed; q->state |= Qstarve; q->eof = 0; } /* * return bytes queued */ int qlen(Queue *q) { return q->len; } /* * return true if we can read without blocking */ int qcanread(Queue *q) { return q->bfirst!=0; }