From 187d5227f5870bd640f1bc9d07cefa326762a934 Mon Sep 17 00:00:00 2001 From: David du Colombier <0intro@gmail.com> Date: Tue, 2 Mar 1999 00:00:00 +0000 Subject: [PATCH] Plan 9 from Bell Labs 1999-03-02 --- carrera/clock.c | 14 +++++----- ip/devip.c | 52 ++++++++++++++++++----------------- ip/ethermedium.c | 15 +++++++---- ip/gre.c | 7 ++++- ip/icmp.c | 3 ++- ip/il.c | 33 ++++++++++++++++++++--- ip/ip.c | 9 ++++--- ip/ip.h | 2 +- ip/ipifc.c | 10 +++---- ip/ipmux.c | 2 +- ip/rudp.c | 8 ++++-- ip/tcp.c | 70 +++++++++++++++++++++++++++--------------------- ip/udp.c | 29 +++++++++++++++----- port/devssl.c | 2 +- port/portfns.h | 1 + port/qlock.c | 16 +++++++++++ 16 files changed, 181 insertions(+), 92 deletions(-) diff --git a/carrera/clock.c b/carrera/clock.c index 7e0a2545dea0f2ea927bf698f1abe371acf9bb8d..4f4def6c985ef7fa7e8557679b5f8aa09f28a91e 100644 --- a/carrera/clock.c +++ b/carrera/clock.c @@ -78,15 +78,15 @@ clockinit(void) uvlong updatefastclock(ulong count) { - uvlong cyclecount; + vlong delta; /* keep track of higher precision time */ - cyclecount = count; - if(cyclecount < m->lastcyclecount) - m->fastclock += (0x100000000LL - m->lastcyclecount) + cyclecount; - else - m->fastclock += cyclecount - m->lastcyclecount; - m->lastcyclecount = cyclecount; + delta = count - m->lastcyclecount; + if(delta < 0) + delta += 0x100000000LL; + m->lastcyclecount = count; + m->fastclock += delta; + return m->fastclock; } diff --git a/ip/devip.c b/ip/devip.c index a3177f71cebcfdc6c4609f684a1918a12fd52ff4..f03fb89342c578f3fcbddeb63836a34975a3722a 100644 --- a/ip/devip.c +++ b/ip/devip.c @@ -367,7 +367,14 @@ ipopen(Chan* c, int omode) break; case Qclone: p = f->p[PROTO(c->qid)]; + qlock(p); + if(waserror()){ + qunlock(p); + nexterror(); + } cv = Fsprotoclone(p, commonuser()); + qunlock(p); + poperror(); if(cv == nil) { error(Enodev); break; @@ -380,9 +387,9 @@ ipopen(Chan* c, int omode) p = f->p[PROTO(c->qid)]; qlock(p); cv = p->conv[CONV(c->qid)]; - lock(cv); + qlock(cv); if(waserror()) { - unlock(cv); + qunlock(cv); qunlock(p); nexterror(); } @@ -398,7 +405,7 @@ ipopen(Chan* c, int omode) memmove(cv->owner, commonuser(), NAMELEN); cv->perm = 0660; } - unlock(cv); + qunlock(cv); qunlock(p); poperror(); break; @@ -422,14 +429,14 @@ ipopen(Chan* c, int omode) /* wait for a connect */ sleep(&cv->listenr, incoming, cv); - lock(cv); + qlock(cv); nc = cv->incall; if(nc != nil){ cv->incall = nc->next; c->qid = (Qid){QID(PROTO(c->qid), nc->x, Qctl), 0}; memmove(cv->owner, commonuser(), NAMELEN); } - unlock(cv); + qunlock(cv); qunlock(&cv->listenq); poperror(); @@ -466,10 +473,10 @@ closeconv(Conv *cv) Conv *nc; Ipmulti *mp; - lock(cv); + qlock(cv); if(--cv->inuse > 0) { - unlock(cv); + qunlock(cv); return; } @@ -1005,17 +1012,15 @@ Fsbuiltinproto(Fs* f, uchar proto) return f->t2p[proto] != nil; } +/* + * called with protocol locked + */ Conv* Fsprotoclone(Proto *p, char *user) { Conv *c, **pp, **ep; c = nil; - qlock(p); - if(waserror()) { - qunlock(p); - nexterror(); - } ep = &p->conv[p->nc]; for(pp = p->conv; pp < ep; pp++) { c = *pp; @@ -1023,7 +1028,7 @@ Fsprotoclone(Proto *p, char *user) c = malloc(sizeof(Conv)); if(c == nil) error(Enomem); - lock(c); + qlock(c); c->p = p; c->x = pp - p->conv; c->ptcl = malloc(p->ptclsize); @@ -1037,7 +1042,7 @@ Fsprotoclone(Proto *p, char *user) (*p->create)(c); break; } - if(canlock(c)){ + if(canqlock(c)){ /* * make sure both processes and protocol * are done with this Conv @@ -1045,12 +1050,10 @@ Fsprotoclone(Proto *p, char *user) if(c->inuse == 0 && (p->inuse == nil || (*p->inuse)(c) == 0)) break; - unlock(c); + qunlock(c); } } if(pp >= ep) { - qunlock(p); - poperror(); return nil; } @@ -1068,9 +1071,7 @@ Fsprotoclone(Proto *p, char *user) qreopen(c->wq); qreopen(c->eq); - unlock(c); - qunlock(p); - poperror(); + qunlock(c); return c; } @@ -1110,6 +1111,9 @@ Fsrcvpcolx(Fs *f, uchar proto) return f->t2p[proto]; } +/* + * called with protocol locked + */ Conv* Fsnewcall(Conv *c, uchar *raddr, ushort rport, uchar *laddr, ushort lport) { @@ -1117,19 +1121,19 @@ Fsnewcall(Conv *c, uchar *raddr, ushort rport, uchar *laddr, ushort lport) Conv **l; int i; - lock(c); + qlock(c); i = 0; for(l = &c->incall; *l; l = &(*l)->next) i++; if(i >= Maxincall) { - unlock(c); + qunlock(c); return nil; } /* find a free conversation */ nc = Fsprotoclone(c->p, network); if(nc == nil) { - unlock(c); + qunlock(c); return nil; } ipmove(nc->raddr, raddr); @@ -1138,7 +1142,7 @@ Fsnewcall(Conv *c, uchar *raddr, ushort rport, uchar *laddr, ushort lport) nc->lport = lport; nc->next = nil; *l = nc; - unlock(c); + qunlock(c); wakeup(&c->listenr); diff --git a/ip/ethermedium.c b/ip/ethermedium.c index 32b8b4944075e1eac92db36ff395ed3e5738591d..225934a128e8e728e0fac8f9be0b459ff771e57e 100644 --- a/ip/ethermedium.c +++ b/ip/ethermedium.c @@ -265,27 +265,32 @@ etherread(void *a) Ipifc *ifc; Block *bp; Etherrock *er; - int locked = 0; ifc = a; er = ifc->arg; er->readp = up; /* hide identity under a rock for unbind */ if(waserror()){ - if(locked) - runlock(ifc); er->readp = 0; pexit("hangup", 1); } for(;;){ bp = devtab[er->mchan->type]->bread(er->mchan, ifc->maxmtu, 0); - rlock(ifc); locked = 1; USED(locked); + if(!canrlock(ifc)){ + freeb(bp); + continue; + } + if(waserror()){ + runlock(ifc); + nexterror(); + } ifc->in++; bp->rp += ifc->m->hsize; if(ifc->lifc == nil) freeb(bp); else ipiput(er->f, ifc->lifc->local, bp); - runlock(ifc); locked = 0; USED(locked); + runlock(ifc); + poperror(); } } diff --git a/ip/gre.c b/ip/gre.c index f6fba2df977b54b75179f1ccfb16356066565add..e180b56da71ed0dc65580c9da6262c6649fa8592 100644 --- a/ip/gre.c +++ b/ip/gre.c @@ -115,7 +115,7 @@ greclose(Conv *c) c->lport = 0; c->rport = 0; - unlock(c); + qunlock(c); } int drop; @@ -179,6 +179,8 @@ greiput(Proto *gre, uchar*, Block *bp) v4tov6(raddr, ghp->src); eproto = nhgets(ghp->eproto); + qlock(gre); + /* Look for a conversation structure for this port and address */ c = nil; for(p = gre->conv; *p; p++) { @@ -190,10 +192,13 @@ greiput(Proto *gre, uchar*, Block *bp) } if(*p == nil) { + qunlock(gre); freeblist(bp); return; } + qunlock(gre); + /* * Trim the packet down to data size */ diff --git a/ip/icmp.c b/ip/icmp.c index c7af8f3120d0a9952f673a3810f048187dcbecde..c4fc421d931a71c73100e7356180bbeebe8a62a4 100644 --- a/ip/icmp.c +++ b/ip/icmp.c @@ -140,7 +140,7 @@ icmpclose(Conv *c) ipmove(c->laddr, IPnoaddr); ipmove(c->raddr, IPnoaddr); c->lport = 0; - unlock(c); + qunlock(c); } static void @@ -328,6 +328,7 @@ icmpiput(Proto *icmp, uchar*, Block *bp) } if(p->type <= Maxtype) ipriv->in[p->type]++; + switch(p->type) { case EchoRequest: r = mkechoreply(bp); diff --git a/ip/il.c b/ip/il.c index 680c6590e04db36bfa552643d4707d7647e4e8f7..2a7a459ecf946b831f17fa0c2a554c103a0e81d5 100644 --- a/ip/il.c +++ b/ip/il.c @@ -274,7 +274,7 @@ ilclose(Conv *c) break; } ilfreeq(ic); - unlock(c); + qunlock(c); } void @@ -521,17 +521,28 @@ iliput(Proto *il, uchar*, Block *bp) goto raise; } + qlock(il); + for(p = il->conv; *p; p++) { s = *p; if(s->lport == sp) if(s->rport == dp) if(ipcmp(s->raddr, raddr) == 0) { + qlock(s); + qunlock(il); + if(waserror()){ + qunlock(s); + nexterror(); + } ilprocess(s, ih, bp); + qunlock(s); + poperror(); return; } } if(ih->iltype != Ilsync){ + qunlock(il); if(ih->iltype < 0 || ih->iltype > Ilclose) st = "?"; else @@ -564,24 +575,36 @@ iliput(Proto *il, uchar*, Block *bp) else if(gen) s = gen; - else + else { + qunlock(il); goto raise; + } v4tov6(laddr, ih->dst); new = Fsnewcall(s, raddr, dp, laddr, sp); if(new == nil){ + qunlock(il); netlog(il->f, Logil, "il: bad newcall %I/%ud->%ud\n", raddr, sp, dp); ilsendctl(nil, ih, Ilclose, 0, nhgetl(ih->ilid), 0); goto raise; } ic = (Ilcb*)new->ptcl; + ic->conv = new; ic->state = Ilsyncee; ilcbinit(ic); ic->rstart = nhgetl(ih->ilid); + qlock(new); + qunlock(il); + if(waserror()){ + qunlock(new); + nexterror(); + } ilprocess(new, ih, bp); + qunlock(new); + poperror(); return; raise: @@ -1179,20 +1202,24 @@ iladvise(Proto *il, Block *bp, char *msg) /* Look for a connection, unfortunately the destination port is missing */ + qlock(il); for(p = il->conv; *p; p++) { s = *p; if(s->lport == psource) if(ipcmp(s->laddr, source) == 0) if(ipcmp(s->raddr, dest) == 0){ + qunlock(il); ic = (Ilcb*)s->ptcl; switch(ic->state){ case Ilsyncer: ilhangup(s, msg); break; } - break; + freeblist(bp); + return; } } + qunlock(il); freeblist(bp); } diff --git a/ip/ip.c b/ip/ip.c index 321419d5bdf6d64f14d1d90b445b0930c6366325..7c6e58d9c5702e644fdcc75ed46483e15a5281be 100644 --- a/ip/ip.c +++ b/ip/ip.c @@ -85,7 +85,7 @@ struct IP Fragment* flisthead; Fragment* fragfree; - ulong id; + Ref id; int iprouting; /* true if we route like a gateway */ void (*ipextprotoiput)(Block*); }; @@ -186,11 +186,12 @@ ipoput(Fs *f, Block *bp, int gating, int ttl) eh->ttl = ttl; } + if(!canrlock(ifc)) + goto free; if(waserror()){ runlock(ifc); nexterror(); } - rlock(ifc); if(ifc->m == nil) goto raise; @@ -198,7 +199,7 @@ ipoput(Fs *f, Block *bp, int gating, int ttl) medialen = ifc->m->maxmtu - ifc->m->hsize; if(len <= medialen) { if(!gating) - hnputs(eh->id, ip->id++); + hnputs(eh->id, incref(&ip->id)); hnputs(eh->length, len); eh->frag[0] = 0; eh->frag[1] = 0; @@ -231,7 +232,7 @@ ipoput(Fs *f, Block *bp, int gating, int ttl) if(gating) lid = nhgets(eh->id); else - lid = ip->id++; + lid = incref(&ip->id); offset = IPHDR; while(xp != nil && offset && offset >= BLEN(xp)) { diff --git a/ip/ip.h b/ip/ip.h index 1f7eacd25ddd2b6557baa47b7272a96207392e43..d5dc298d13a6053f9737fc035fcc1ffc43c76972 100644 --- a/ip/ip.h +++ b/ip/ip.h @@ -58,7 +58,7 @@ enum */ struct Conv { - Lock; + QLock; int x; /* conversation index */ Proto* p; diff --git a/ip/ipifc.c b/ip/ipifc.c index 573b6656e97e74ac4ab0329c2a0bdc41c561ca4f..21a9be760904baef7e3dda3c85332f30c3ef849d 100644 --- a/ip/ipifc.c +++ b/ip/ipifc.c @@ -136,9 +136,9 @@ ipifcbind(Conv *c, char **argv, int argc) ifc->minmtu = ifc->m->minmtu; ifc->maxmtu = ifc->m->maxmtu; if(ifc->m->unbindonclose == 0){ - lock(ifc->conv); + qlock(ifc->conv); ifc->conv->inuse++; - unlock(ifc->conv); + qunlock(ifc->conv); } ifc->ifcid++; @@ -166,9 +166,9 @@ ipifcunbind(Ipifc *ifc) /* dissociate routes */ if(ifc->m != nil && ifc->m->unbindonclose == 0){ - lock(ifc->conv); + qlock(ifc->conv); ifc->conv->inuse--; - unlock(ifc->conv); + qunlock(ifc->conv); } ifc->ifcid++; @@ -305,7 +305,7 @@ ipifcclose(Conv *c) m = ifc->m; if(m != nil && m->unbindonclose) ipifcunbind(ifc); - unlock(c); + qunlock(c); } /* diff --git a/ip/ipmux.c b/ip/ipmux.c index 2ef0a1cbf40f6c2b320e8c6b7c3eadfa3f246d03..159c638e5b6c860077f9c95002f950143ae45ef0 100644 --- a/ip/ipmux.c +++ b/ip/ipmux.c @@ -620,7 +620,7 @@ ipmuxclose(Conv *c) ipmuxtreefree(r->chain); r->chain = nil; - unlock(c); + qunlock(c); } /* diff --git a/ip/rudp.c b/ip/rudp.c index 1b2680044ee1e4c3e8465333ad755881c5738aa9..e2f95b4dec021a8293d8eaf120d908d6655bada6 100644 --- a/ip/rudp.c +++ b/ip/rudp.c @@ -289,7 +289,7 @@ rudpclose(Conv *c) qunlock(ucb); - unlock(c); + qunlock(c); } /* @@ -490,6 +490,8 @@ rudpiput(Proto *rudp, uchar *ia, Block *bp) } } + qlock(rudp); + /* Look for a conversation structure for this port */ c = nil; for(p = rudp->conv; *p; p++) { @@ -512,6 +514,7 @@ rudpiput(Proto *rudp, uchar *ia, Block *bp) } if(*p == nil) { + qunlock(rudp); upriv->ustats.rudpNoPorts++; netlog(f, Logrudp, "rudp: no conv %I!%d -> %I!%d\n", raddr, rport, laddr, lport); @@ -525,8 +528,9 @@ rudpiput(Proto *rudp, uchar *ia, Block *bp) } ucb = (Rudpcb*)c->ptcl; - qlock(ucb); + qunlock(rudp); + if(reliput(c, bp, raddr, rport) < 0){ qunlock(ucb); freeb(bp); diff --git a/ip/tcp.c b/ip/tcp.c index e5421ead375f51b93ce561f679656a8f4f5da06e..68d8f58e44552baf0e02b140ea1b2f4e060cb228 100644 --- a/ip/tcp.c +++ b/ip/tcp.c @@ -133,10 +133,12 @@ struct Reseq ushort length; }; +/* + * the qlock in the Conv locks this structure + */ typedef struct Tcpctl Tcpctl; struct Tcpctl { - QLock; uchar state; /* Connection state */ uchar type; /* Listening or active connection */ uchar code; /* Icmp code */ @@ -343,8 +345,6 @@ tcpclose(Conv *c) qhangup(c->wq, nil); qhangup(c->eq, nil); - unlock(c); - switch(tcb->state) { case Listen: /* @@ -352,31 +352,27 @@ tcpclose(Conv *c) */ Fsconnected(c, "Hangup"); - qlock(tcb); localclose(c, nil); break; case Closed: case Syn_sent: - qlock(tcb); localclose(c, nil); break; case Syn_received: case Established: - qlock(tcb); tcb->sndcnt++; tcb->snd.nxt++; tcpsetstate(c, Finwait1); tcpoutput(c); break; case Close_wait: - qlock(tcb); tcb->sndcnt++; tcb->snd.nxt++; tcpsetstate(c, Last_ack); tcpoutput(c); break; } - qunlock(tcb); + qunlock(c); } void @@ -400,11 +396,11 @@ tcpkick(Conv *s, int len) /* * Push data */ - qlock(tcb); + qlock(s); tcb->sndcnt += len; tcprcvwin(s); tcpoutput(s); - qunlock(tcb); + qunlock(s); break; default: localclose(s, "Hangup"); @@ -445,11 +441,11 @@ tcpacktimer(Conv *s) tcb = (Tcpctl*)s->ptcl; - qlock(tcb); + qlock(s); tcb->flags |= FORCE; tcprcvwin(s); tcpoutput(s); - qunlock(tcb); + qunlock(s); } static void @@ -667,12 +663,12 @@ tcpstart(Conv *s, int mode, ushort window) case TCP_CONNECT: tpriv->tstats.tcpActiveOpens++; /* Send SYN, go into SYN_SENT state */ - qlock(tcb); + qlock(s); tcb->flags |= ACTIVE; tcpsndsyn(tcb); tcpsetstate(s, Syn_sent); tcpoutput(s); - qunlock(tcb); + qunlock(s); break; } } @@ -893,10 +889,10 @@ tcphangup(Conv *s) tcb = (Tcpctl*)s->ptcl; if(waserror()){ - qunlock(tcb); + qunlock(s); return commonerror(); } - qlock(tcb); + qlock(s); if(s->raddr != 0) { seg.flags = RST | ACK; seg.ack = tcb->rcv.nxt; @@ -911,7 +907,7 @@ tcphangup(Conv *s) } localclose(s, nil); poperror(); - qunlock(tcb); + qunlock(s); return nil; } @@ -1159,6 +1155,8 @@ tcpiput(Proto *tcp, uchar*, Block *bp) return; } + /* lock protocol while searching for a conversation */ + qlock(tcp); /* Look for a connection. failing that look for a listener. */ for(p = tcp->conv; *p; p++) { @@ -1175,6 +1173,7 @@ tcpiput(Proto *tcp, uchar*, Block *bp) /* can't send packets to a listener */ tcb = (Tcpctl*)s->ptcl; if(tcb->state == Listen){ + qunlock(tcp); freeblist(bp); return; } @@ -1184,10 +1183,13 @@ tcpiput(Proto *tcp, uchar*, Block *bp) * dump packets with bogus flags */ if(seg.flags & RST){ + qunlock(tcp); freeblist(bp); return; } + if(seg.flags & ACK) { + qunlock(tcp); sndrst(tcp, source, dest, length, &seg); freeblist(bp); return; @@ -1223,8 +1225,9 @@ tcpiput(Proto *tcp, uchar*, Block *bp) s = tcpincoming(gen, &seg, source, dest); } if(s == nil) { - freeblist(bp); + qunlock(tcp); sndrst(tcp, source, dest, length, &seg); + freeblist(bp); return; } @@ -1233,7 +1236,8 @@ tcpiput(Proto *tcp, uchar*, Block *bp) * Out-of-band data is ignored - it was always a bad idea. */ tcb = (Tcpctl*)s->ptcl; - qlock(tcb); + qlock(s); + qunlock(tcp); if(tcb->kacounter > 0) tcb->kacounter = MAXBACKOFF; @@ -1283,7 +1287,7 @@ tcpiput(Proto *tcp, uchar*, Block *bp) else freeblist(bp); - qunlock(tcb); + qunlock(s); return; case Syn_received: /* doesn't matter if it's the correct ack, we're just trying to set timing */ @@ -1308,7 +1312,7 @@ tcpiput(Proto *tcp, uchar*, Block *bp) tcb->flags |= FORCE; goto output; } - qunlock(tcb); + qunlock(s); return; } @@ -1449,7 +1453,7 @@ tcpiput(Proto *tcp, uchar*, Block *bp) if(bp != nil) freeblist(bp); sndrst(tcp, source, dest, length, &seg); - qunlock(tcb); + qunlock(s); return; } } @@ -1514,10 +1518,10 @@ tcpiput(Proto *tcp, uchar*, Block *bp) } output: tcpoutput(s); - qunlock(tcb); + qunlock(s); return; raise: - qunlock(tcb); + qunlock(s); freeblist(bp); tcpkick(s, 0); } @@ -1773,9 +1777,9 @@ tcpkeepalive(Conv *s) if(--(tcb->kacounter) <= 0) localclose(s, Etimedout); else { - qlock(tcb); + qlock(s); tcpsendka(s); - qunlock(tcb); + qunlock(s); tcpgo(s->p->priv, &tcb->katimer); } } @@ -1810,7 +1814,7 @@ tcprxmit(Conv *s) tpriv = s->p->priv; tcb = (Tcpctl*)s->ptcl; - qlock(tcb); + qlock(s); tcb->flags |= RETRAN|FORCE; tcb->snd.ptr = tcb->snd.una; @@ -1826,7 +1830,7 @@ tcprxmit(Conv *s) tcpoutput(s); tpriv->tstats.tcpRetransSegs++; - qunlock(tcb); + qunlock(s); } void @@ -2035,6 +2039,7 @@ tcpadvise(Proto *tcp, Block *bp, char *msg) pdest = nhgets(h->tcpdport); /* Look for a connection */ + qlock(tcp); for(p = tcp->conv; *p; p++) { s = *p; tcb = (Tcpctl*)s->ptcl; @@ -2043,16 +2048,19 @@ tcpadvise(Proto *tcp, Block *bp, char *msg) if(tcb->state != Closed) if(ipcmp(s->raddr, dest) == 0) if(ipcmp(s->laddr, source) == 0){ - qlock(tcb); + qlock(s); + qunlock(tcp); switch(tcb->state){ case Syn_sent: localclose(s, msg); break; } - qunlock(tcb); - break; + qunlock(s); + freeblist(bp); + return; } } + qunlock(tcp); freeblist(bp); } diff --git a/ip/udp.c b/ip/udp.c index 5e7c42b492d47d69ba9e2670208f27c540798b0d..dff93c9ab5210614c155861e48ca0a2985d108ce 100644 --- a/ip/udp.c +++ b/ip/udp.c @@ -130,7 +130,7 @@ udpclose(Conv *c) ucb = (Udpcb*)c->ptcl; ucb->headers = 0; - unlock(c); + qunlock(c); } void @@ -270,6 +270,8 @@ udpiput(Proto *udp, uchar *ia, Block *bp) } } + qlock(udp); + /* Look for a conversation structure for this port */ c = nil; for(p = udp->conv; *p; p++) { @@ -293,6 +295,7 @@ udpiput(Proto *udp, uchar *ia, Block *bp) if(*p == nil) { upriv->ustats.udpNoPorts++; + qunlock(udp); netlog(f, Logudp, "udp: no conv %I!%d -> %I!%d\n", raddr, rport, laddr, lport); uh->Unused = ottl; @@ -302,12 +305,17 @@ udpiput(Proto *udp, uchar *ia, Block *bp) return; } + ucb = (Udpcb*)c->ptcl; + qlock(c); + qunlock(udp); + /* * Trim the packet down to data size */ len -= (UDP_HDRSIZE-UDP_PHDRSIZE); bp = trimblock(bp, UDP_IPHDR+UDP_HDRSIZE, len); if(bp == nil){ + qunlock(c); netlog(f, Logudp, "udp: len err %I.%d -> %I.%d\n", raddr, rport, laddr, lport); upriv->lenerr++; @@ -317,8 +325,6 @@ udpiput(Proto *udp, uchar *ia, Block *bp) netlog(f, Logudpmsg, "udp: %I.%d -> %I.%d l %d\n", raddr, rport, laddr, lport, len); - ucb = (Udpcb*)c->ptcl; - switch(ucb->headers){ case 6: /* pass the src address */ @@ -361,11 +367,16 @@ udpiput(Proto *udp, uchar *ia, Block *bp) bp = concatblock(bp); if(qfull(c->rq)){ + qunlock(c); netlog(f, Logudp, "udp: qfull %I.%d -> %I.%d\n", raddr, rport, laddr, lport); freeblist(bp); - }else - qpass(c->rq, bp); + return; + } + + qpass(c->rq, bp); + qunlock(c); + } char* @@ -402,17 +413,23 @@ udpadvise(Proto *udp, Block *bp, char *msg) pdest = nhgets(h->udpdport); /* Look for a connection */ + qlock(udp); for(p = udp->conv; *p; p++) { s = *p; if(s->rport == pdest) if(s->lport == psource) if(ipcmp(s->raddr, dest) == 0) if(ipcmp(s->laddr, source) == 0){ + qlock(s); + qunlock(udp); qhangup(s->rq, msg); qhangup(s->wq, msg); - break; + qunlock(s); + freeblist(bp); + return; } } + qunlock(udp); freeblist(bp); } diff --git a/port/devssl.c b/port/devssl.c index 537adfe6f9846a38ca3301361f25abeb1beb8c57..42a7af02f9972ad0cae8a4b7ea7835d82b737102 100644 --- a/port/devssl.c +++ b/port/devssl.c @@ -418,7 +418,7 @@ qremove(Block **l, int n, int discard) } else *l = b->next; b->next = 0; - break; + return first; } else if(i > n){ i -= n; if(discard){ diff --git a/port/portfns.h b/port/portfns.h index bd817fc6b894a0153de5161e1f55d9995e838d86..1fb78a2f5382e605eb01bd7a6277055c9359ba19 100644 --- a/port/portfns.h +++ b/port/portfns.h @@ -26,6 +26,7 @@ int canlock(Lock*); int canpage(Proc*); int canputc(void*); int canqlock(QLock*); +int canrlock(RWlock*); void chandevinit(void); void chandevreset(void); void chanfree(Chan*); diff --git a/port/qlock.c b/port/qlock.c index 8b8da34722383e09cff7d97aba724e0d4a946cb6..f26b1284266d183edd6347244744b1e9e8dd5fdd 100644 --- a/port/qlock.c +++ b/port/qlock.c @@ -196,3 +196,19 @@ wunlock(RWlock *q) q->writer = 0; unlock(&q->use); } + +/* same as rlock but punts if there are any readers waiting */ +int +canrlock(RWlock *q) +{ + lock(&q->use); +rwstats.rlock++; + if(q->writer == 0 && q->head == nil){ + /* no writer, go for it */ + q->readers++; + unlock(&q->use); + return 1; + } + unlock(&q->use); + return 0; +}