From 767f56e7eb2486e985649ee6c83c5bf2619578c8 Mon Sep 17 00:00:00 2001 From: David du Colombier <0intro@gmail.com> Date: Thu, 27 May 1993 00:00:00 +0000 Subject: [PATCH] Plan 9 from Bell Labs 1993-05-27 --- port/portdat.h | 11 ++- port/portfns.h | 3 +- port/qio.c | 228 +++++++++++++++++++++++++++++++++++-------------- 3 files changed, 176 insertions(+), 66 deletions(-) diff --git a/port/portdat.h b/port/portdat.h index 34d73315107931d27574541c634145b1d8a731fc..575791b40e64ba884ecfadc140d29795fd131e25 100644 --- a/port/portdat.h +++ b/port/portdat.h @@ -647,9 +647,11 @@ struct Block uchar *wp; /* first empty byte */ uchar *lim; /* 1 past the end of the buffer */ uchar *base; /* start of the buffer */ + uchar flag; Rendez r; /* waiting reader */ }; +#define BLEN(b) ((b)->wp - (b)->rp) struct Queue { @@ -675,10 +677,13 @@ struct Queue enum { - Qstarve=1, /* consumer starved */ - Qmsg=1, /* message oriented */ + /* Block.flag */ + Bfilled=1, /* block filled */ + + /* Queue.state */ + Qstarve=1, /* consumer starved */ + Qmsg=2, /* message stream */ }; -#define BLEN(b) ((b)->wp - (b)->rp) enum diff --git a/port/portfns.h b/port/portfns.h index ba16f4bf407b62df6bcf3ad8a010c0d4a24e8f64..444ff9266cc1e9df282c548164dbfc834450ec63 100644 --- a/port/portfns.h +++ b/port/portfns.h @@ -85,6 +85,7 @@ void gotolabel(Label*); int haswaitq(void*); int hwcursmove(int, int); int hwcursset(uchar*, uchar*, int, int); +void iallocinit(void); long ibrk(ulong, int); int incref(Ref*); void initq(IOQ*); @@ -173,7 +174,7 @@ int qconsume(Queue*, uchar*, int); void qlock(QLock*); Queue* qopen(int, void (*)(void*), void*); int qproduce(Queue*, uchar*, int); -long qread(Queue*, char*, int, int); +long qread(Queue*, char*, int); void qunlock(QLock*); long qwrite(Queue*, char*, int); int readnum(ulong, char*, ulong, ulong, int); diff --git a/port/qio.c b/port/qio.c index a0666b9e1017484713e6b2ede8f015287ad0a566..ca99719a791e671fd506f836e22ee531812a46f9 100644 --- a/port/qio.c +++ b/port/qio.c @@ -142,15 +142,18 @@ iallocb(int size) cl = &arena.alloc[pow]; lock(cl); p = cl->first; - if(p){ - cl->have--; - cl->first = p->next; + if(p == 0){ + unlock(cl); + return 0; } + cl->have--; + cl->first = p->next; unlock(cl); b = (Block *)p; b->base = (uchar*)(b+1); b->wp = b->rp = b->base; b->lim = b->base + (1<flag = 0; return b; } panic("iallocb %d\n", size); @@ -186,6 +189,7 @@ allocb(int size) b->base = (uchar*)(b+1); b->rp = b->wp = b->base; b->lim = b->base + size; + b->flag = 0; return b; } @@ -200,7 +204,9 @@ qconsume(Queue *q, uchar *p, int len) Block *b; int n; + /* sync with qwrite */ lock(q); + b = q->bfirst; if(b == 0){ q->state |= Qstarve; @@ -211,7 +217,7 @@ qconsume(Queue *q, uchar *p, int len) if(n < len) len = n; memmove(p, b->rp, len); - if((q->state&Qmsg) || len == n) + if((q->state & Qmsg) || len == n) q->bfirst = b->next; else b->rp += len; @@ -223,7 +229,7 @@ qconsume(Queue *q, uchar *p, int len) unlock(q); - if((q->state&Qmsg) || len == n) + if((q->state & Qmsg) || len == n) ifree(b); return len; @@ -235,18 +241,21 @@ qproduce0(Queue *q, uchar *p, int len) Block *b; int n; + /* sync with qread */ lock(q); + b = q->rfirst; if(b){ /* hand to waiting receiver */ + q->rfirst = b->next; + unlock(q); n = b->lim - b->wp; if(n < len) len = n; memmove(b->wp, p, len); b->wp += len; - q->rfirst = b->next; + b->flag |= Bfilled; wakeup(&b->r); - unlock(q); return len; } @@ -258,17 +267,19 @@ qproduce0(Queue *q, uchar *p, int len) /* save in buffer */ b = q->bfirst; - if((q->state&Qmsg)==0 && b && b->lim-b->wp <= len){ + if((q->state & Qmsg) == 0 && b && b->lim - b->wp <= len){ memmove(b->wp, p, len); b->wp += len; + b->flag |= Bfilled; } else { b = iallocb(len); if(b == 0){ unlock(q); return -1; } + memmove(b->wp, p, len); b->wp += len; - memmove(b->rp, p, len); + b->flag |= Bfilled; if(q->bfirst) q->blast->next = b; else @@ -277,6 +288,7 @@ qproduce0(Queue *q, uchar *p, int len) } q->len += len; unlock(q); + return len; } @@ -285,14 +297,13 @@ qproduce(Queue *q, uchar *p, int len) { int n, sofar; - if(q->state&Qmsg) - return qproduce0(q, p, len); - - for(sofar = 0; sofar < len; sofar += n){ - n = qproduce0(q, p+sofar, len-sofar); + sofar = 0; + do { + n = qproduce0(q, p + sofar, len - sofar); if(n < 0) break; - } + sofar += n; + } while(sofar < len && (q->state & Qmsg) == 0); return sofar; } @@ -315,30 +326,60 @@ qopen(int limit, void (*kick)(void*), void *arg) return q; } +ulong qrtoomany; +ulong qrtoofew; + static int bfilled(void *a) { Block *b = a; - return b->wp - b->rp; + return b->flag & Bfilled; } long qread(Queue *q, char *p, int len) { - Block *b, *bb; + Block *b, *bb, **l; int x, n; qlock(&q->rlock); + b = 0; + if(waserror()){ + qunlock(&q->rlock); + if(b) + free(b); + nexterror(); + } - /* ... to be replaced by a kmapping if need be */ - b = allocb(len); - + /* + * If there are no buffered blocks, allocate a block + * for the qproducer/qwrite to fill. This is + * optimistic and and we will + * sometimes be wrong: after locking we may either + * have to throw away or allocate one. + * + * We hope to replace the allocb with a kmap later on. + */ +retry: + if(q->bfirst == 0) + b = allocb(len); + + /* sync with qwrite/qproduce */ x = splhi(); lock(q); + bb = q->bfirst; if(bb == 0){ - /* wait for our block to be filled */ + if(b == 0){ + /* we guessed wrong, drop the locks and try again */ + unlock(q); + splx(x); + qrtoofew++; + goto retry; + } + + /* add ourselves to the list of readers */ if(q->rfirst) q->rlast->next = b; else @@ -347,10 +388,33 @@ qread(Queue *q, char *p, int len) unlock(q); splx(x); qunlock(&q->rlock); + poperror(); + + if(waserror()){ + /* on error, unlink us from the chain */ + x = splhi(); + lock(q); + l = &q->rfirst; + for(bb = q->rfirst; bb; bb = bb->next){ + if(b == bb){ + *l = bb->next; + break; + } else + l = &bb->next; + } + unlock(q); + splx(x); + free(b); + nexterror(); + } + + /* wait for the producer */ sleep(&b->r, bfilled, b); n = BLEN(b); memmove(p, b->rp, n); + poperror(); free(b); + return n; } @@ -362,11 +426,13 @@ qread(Queue *q, char *p, int len) q->len -= n; unlock(q); splx(x); + + /* do this outside of the lock(q)! */ memmove(p, bb->rp, n); bb->rp += n; - /* free it or put it back */ - if(drop || bb->rp == bb->wp) + /* free it or put it back on the queue */ + if(bb->rp >= bb->wp || (q->state&Qmsg)) free(bb); else { x = splhi(); @@ -376,53 +442,63 @@ qread(Queue *q, char *p, int len) unlock(q); splx(x); } + + poperror(); qunlock(&q->rlock); - free(b); + if(b){ + qrtoomany++; + free(b); + } return n; } -static int -qnotfull(void *a) -{ - Queue *q = a; - - return q->len < q->limit; -} +ulong qwtoomany; +ulong qwtoofew; static long -qwrite0(Queue *q, char *p, int len) +qwrite0(Queue *q, char *p, int len, Block *b) { - Block *b, *bb; - int x, n; - - b = allocb(len); + Block *bb; + int x, n, sofar; + /* sync with qconsume/qread */ x = splhi(); lock(q); - bb = q->rfirst; - if(bb){ + + sofar = 0; + while(bb = q->rfirst){ /* hand to waiting receiver */ - n = bb->lim - bb->wp; q->rfirst = bb->next; unlock(q); splx(x); - if(n < len) - len = n; - memmove(bb->wp, p, len); - bb->wp += len; + n = bb->lim - bb->wp; + if(n > len-sofar) + n = len - sofar; + memmove(bb->wp, p+sofar, n); + bb->wp += n; + bb->flag |= Bfilled; wakeup(&bb->r); - free(b); - return len; + sofar += n; + if(sofar == len){ + if(b){ + free(b); /* we were wrong to allocate */ + qwtoomany++; + } + return len; + } } - - memmove(b->rp, p, len); - b->wp += len; - /* flow control */ - if(!qnotfull(q)) - sleep(&q->r, qnotfull, q); + /* buffer what ever is left */ + if(b == 0){ + /* we should have alloc'd, return to qwrite and have it do it */ + unlock(q); + splx(x); + qwtoofew++; + return sofar; + } + b->rp += sofar; x = splhi(); lock(q); @@ -442,28 +518,56 @@ qwrite0(Queue *q, char *p, int len) return len; } +static int +qnotfull(void *a) +{ + Queue *q = a; + + return q->len < q->limit; +} + long qwrite(Queue *q, char *p, int len) { - int n, sofar; + int n, i; + Block *b; + + /* + * If there are no readers, grab a buffer and copy + * into it before locking anything down. This + * provides the highest concurrency but we will + * sometimes be wrong: after locking we may either + * have to throw away or allocate one. + */ + if(q->rfirst == 0){ + b = allocb(len); + memmove(b->wp, p, len); + b->wp += len; + } else + b = 0; + /* ensure atomic writes */ qlock(&q->wlock); if(waserror()){ qunlock(&q->wlock); nexterror(); } - if(q->state&Qmsg){ - sofar = qwrite0(q, p, len); - } else { - for(sofar = 0; sofar < len; sofar += n){ - n = qwrite0(q, p+sofar, len-sofar); - if(n < 0) - break; - } + /* flow control */ + sleep(&q->r, qnotfull, q); + + n = qwrite0(q, p, len, b); + if(n != len){ + /* no readers and we need a buffer */ + i = len - n; + b = allocb(i); + memmove(b->wp, p + n, i); + b->wp += n; + n += qwrite0(q, p + n, i, b); } - poperror(); qunlock(&q->wlock); - return sofar; + poperror(); + + return n; }