From cc7c6340905c8929e7c40ae6723e279c9cc6adf0 Mon Sep 17 00:00:00 2001 From: David du Colombier <0intro@gmail.com> Date: Sat, 1 Aug 1998 00:00:00 +0000 Subject: [PATCH] Plan 9 from Bell Labs 1998-08-01 --- ip/rudp.c | 71 +++++++++++++++++++++++++++++++++++++++++++------------ ip/udp.c | 3 +-- 2 files changed, 57 insertions(+), 17 deletions(-) diff --git a/ip/rudp.c b/ip/rudp.c index 6f2981e82df0fe659c6be69d5130c834e2854c0a..2ec802699a986b034de3669ecd797007a59f2ddc 100644 --- a/ip/rudp.c +++ b/ip/rudp.c @@ -14,6 +14,15 @@ #define DEBUG 0 #define DPRINT if(DEBUG)print +#define SEQDIFF(a,b) ( (a)>=(b)?\ + (a)-(b):\ + 0xffffffffUL-((b)-(a)) ) +#define INSEQ(a,start,end) ( (start)<=(end)?\ + ((a)>(start)&&(a)<=(end)):\ + ((a)>(start)||(a)<=(end)) ) +#define UNACKED(r) SEQDIFF(r->sndseq, r->ackrcvd) +#define NEXTSEQ(a) ( (a)+1 == 0 ? 1 : (a)+1 ) + enum { UDP_HDRSIZE = 20, /* pseudo header + udp header */ @@ -27,7 +36,7 @@ enum Rudprxms = 200, Rudptickms = 100, Rudpmaxxmit = 10, - + Maxunacked = 100, }; typedef struct Udphdr Udphdr; @@ -93,7 +102,7 @@ struct Reliable { Reliable *next; - uchar addr[IPaddrlen]; /* always V6 when put here */ + uchar addr[IPaddrlen]; /* always V6 when put here */ ushort port; Block *unacked; /* unacked msg list */ @@ -110,6 +119,11 @@ struct Reliable ulong acksent; /* last ack sent */ ulong ackrcvd; /* last msg for which ack was rcvd */ + + /* flow control */ + QLock lock; + Rendez vous; + int blocked; }; @@ -185,8 +199,17 @@ rudpconnect(Conv *c, char **argv, int argc) static int rudpstate(Conv *c, char *state, int n) { - USED(c); - return snprint(state, n, "%s", "Reliable UDP V0.1"); + Rudpcb *ucb; + Reliable *r; + int m; + + m = snprint(state, n, "%s", c->inuse?"Open":"Closed"); + ucb = (Rudpcb*)c->ptcl; + qlock(ucb); + for(r = ucb->r; r; r = r->next) + m += snprint(state+m, n-m, " %I/%d", r->addr, UNACKED(r)); + qunlock(ucb); + return m; } static char* @@ -238,6 +261,12 @@ rudpclose(Conv *c) unlock(c); } +int +flow(Reliable *r) +{ + return UNACKED(r) <= Maxunacked; +} + void rudpkick(Conv *c, int) { @@ -339,14 +368,14 @@ rudpkick(Conv *c, int) qlock(ucb); r = relstate(ucb, raddr, rport, "kick"); - r->sndseq++; + r->sndseq = NEXTSEQ(r->sndseq); hnputl(rh->relseq, r->sndseq); hnputl(rh->relsgen, r->sndgen); hnputl(rh->relack, r->rcvseq); /* ACK last rcvd packet */ hnputl(rh->relagen, r->rcvgen); - if(r->rcvseq < r->acksent) + if(r->rcvseq != r->acksent) r->acksent = r->rcvseq; hnputs(uh->udpcksum, ptclcsum(bp, UDP_IPHDR, dlen+UDP_RHDRSIZE)); @@ -360,6 +389,15 @@ rudpkick(Conv *c, int) r->sndseq, r->sndgen, r->rcvseq, r->rcvgen, r->sndgen); ipoput(f, bp, 0, c->ttl); + + /* flow control of sorts */ + qlock(&r->lock); + if(UNACKED(r) > Maxunacked){ + r->blocked = 1; + sleep(&r->vous, flow, r); + r->blocked = 0; + } + qunlock(&r->lock); } void @@ -664,7 +702,7 @@ loop: if(r->timeout > Rudprxms*r->xmits) relrexmit(c, r); } - if(r->acksent < r->rcvseq) + if(r->acksent != r->rcvseq) relsendack(c, r); } qunlock(ucb); @@ -736,7 +774,6 @@ reliput(Conv *c, Block *bp, uchar *addr, ushort port) ack = nhgetl(rh->relack); agen = nhgetl(rh->relagen); - upriv = c->p->priv; ucb = (Rudpcb*)c->ptcl; r = relstate(ucb, addr, port, "input"); @@ -748,16 +785,20 @@ reliput(Conv *c, Block *bp, uchar *addr, ushort port) /* dequeue acked packets */ if(ack && agen == r->sndgen){ ackreal = 0; - while(r->unacked != nil && ack > r->ackrcvd){ + while(r->unacked != nil && INSEQ(ack, r->ackrcvd, r->sndseq)){ nbp = r->unacked; r->unacked = nbp->list; DPRINT("%d/%d acked, r->sndgen = %d\n", ack, agen, r->sndgen); freeb(nbp); - r->ackrcvd++; + r->ackrcvd = NEXTSEQ(r->ackrcvd); ackreal = 1; } + /* flow control */ + if(UNACKED(r) < Maxunacked/8 && r->blocked) + wakeup(&r->vous); + /* * retransmit next packet if the acked packet * was transmitted more than once @@ -777,7 +818,6 @@ reliput(Conv *c, Block *bp, uchar *addr, ushort port) if(seq != 1) return -1; - /* new connection */ if(r->rcvgen != 0){ DPRINT("new con r->rcvgen = %d, sgen = %d\n", r->rcvgen, sgen); @@ -786,8 +826,8 @@ reliput(Conv *c, Block *bp, uchar *addr, ushort port) r->rcvgen = sgen; } - /* no message */ - if(seq == 0) + /* no message or input queue full */ + if(seq == 0 || qfull(c->rq)) return -1; if(DEBUG && ++drop == drop_rate){ @@ -797,9 +837,9 @@ reliput(Conv *c, Block *bp, uchar *addr, ushort port) } /* refuse out of order delivery */ - if(seq != r->rcvseq + 1){ + if(seq != NEXTSEQ(r->rcvseq)){ upriv->orders++; - DPRINT("out of sequence %d not %d\n", seq, r->rcvseq + 1); + DPRINT("out of sequence %d not %d\n", seq, NEXTSEQ(r->rcvseq)); return -1; } r->rcvseq = seq; @@ -880,6 +920,7 @@ relhangup(Conv *, Reliable *r) r->ackrcvd = 0; r->xmits = 0; r->timeout = 0; + wakeup(&r->vous); } /* diff --git a/ip/udp.c b/ip/udp.c index 660d178af4fa6c086de58b6011ffbe726e788d6e..474241d97faaaea15799a36c59d1973b42a1fbc6 100644 --- a/ip/udp.c +++ b/ip/udp.c @@ -91,8 +91,7 @@ udpconnect(Conv *c, char **argv, int argc) static int udpstate(Conv *c, char *state, int n) { - USED(c); - return snprint(state, n, "%s", "Datagram"); + return snprint(state, n, "%s", c->inuse?"Open":"Closed"); } static char*