From d428e767dd2429a5e9a89c9e3db122cd8840cd51 Mon Sep 17 00:00:00 2001 From: David du Colombier <0intro@gmail.com> Date: Thu, 16 Apr 1992 00:00:00 +0000 Subject: [PATCH] Plan 9 from Bell Labs 1992-04-16 --- port/devip.c | 66 ++++++++++++++++++++++++++++++++++++++++--------- port/ipdat.h | 5 ++-- port/net.c | 1 + port/tcpif.c | 9 ++++--- port/tcpinput.c | 58 +++++++++++++++++++++++++++++++++++++------ 5 files changed, 114 insertions(+), 25 deletions(-) diff --git a/port/devip.c b/port/devip.c index 9e21cbe2a3523dccf8c2fa878095d10d27f32e8e..5366b827c83d10175a3d7349404224467e18aa86 100644 --- a/port/devip.c +++ b/port/devip.c @@ -21,6 +21,7 @@ Ipifc *ipifc; /* IP protocol interfaces for stip */ Ipconv *ipconv[Nrprotocol]; /* Connections for each protocol */ Network *ipnet[Nrprotocol]; /* User level interface for protocol */ QLock ipalloc; /* Protocol port allocation lock */ +Ipconv *tcpbase; /* Base tcp connection */ /* ARPA User Datagram Protocol */ void udpstiput(Queue *, Block *); @@ -171,7 +172,6 @@ ipincoming(Ipconv *base, Ipconv *from) netown(new->net, new->index, u->p->user, 0); new->ref = 1; - new->newcon = 0; qunlock(new); return new; } @@ -262,10 +262,11 @@ ipwrite(Chan *c, char *a, long n, ulong offset) cp->pdst = atoi(ctlarg[1]); /* If we have no local port assign one */ - qlock(&ipalloc); - if(cp->psrc == 0) + if(cp->psrc == 0){ + qlock(&ipalloc); cp->psrc = nextport(ipconv[c->dev], priv); - qunlock(&ipalloc); + qunlock(&ipalloc); + } if(cp->stproto == &tcpinfo) tcpstart(cp, TCP_ACTIVE, Streamhi, 0); @@ -289,13 +290,20 @@ ipwrite(Chan *c, char *a, long n, ulong offset) port = atoi(field[1]); - qlock(&ipalloc); - if(portused(ipconv[c->dev], port)) { - qunlock(&ipalloc); - error(Einuse); - } - cp->psrc = port; - qunlock(&ipalloc); + if(port){ + qlock(&ipalloc); + if(portused(ipconv[c->dev], port)) { + qunlock(&ipalloc); + error(Einuse); + } + cp->psrc = port; + qunlock(&ipalloc); + } else if(*field[1] != '*'){ + qlock(&ipalloc); + cp->psrc = nextport(ipconv[c->dev], 0); + qunlock(&ipalloc); + } else + cp->psrc = 0; if(cp->stproto == &tcpinfo) tcpstart(cp, TCP_PASSIVE, Streamhi, 0); @@ -557,6 +565,8 @@ void tcpstopen(Queue *q, Stream *s) { Ipconv *ipc; + Tcpctl *tcb; + Block *bp; static int tcpkprocs; if(!Ipoutput) { @@ -573,6 +583,8 @@ tcpstopen(Queue *q, Stream *s) } + if(tcpbase == 0) + tcpbase = ipconv[s->dev]; ipc = &ipconv[s->dev][s->id]; ipc->ipinterface = newipifc(IP_TCPPROTO, tcp_input, ipconv[s->dev], 1500, 512, ETHER_HDR, "TCP"); @@ -584,6 +596,13 @@ tcpstopen(Queue *q, Stream *s) RD(q)->ptr = (void *)ipc; WR(q)->next->ptr = (void *)ipc->ipinterface; WR(q)->ptr = (void *)ipc; + + /* pass any waiting data upstream */ + tcb = &ipc->tcpctl; + qlock(tcb); + while(bp = getb(&tcb->rcvq)) + PUTNEXT(ipc->readq, bp); + qunlock(tcb); } void @@ -682,6 +701,7 @@ void tcpstclose(Queue *q) { Ipconv *s; + Ipconv *etab, *ifc; Tcpctl *tcb; s = (Ipconv *)(q->ptr); @@ -693,10 +713,32 @@ tcpstclose(Queue *q) qunlock(s); switch(tcb->state){ - case Closed: case Listen: + /* + * reset any incoming calls to this listener + */ + qlock(s); + s->backlog = 0; + s->curlog = 0; + etab = &tcpbase[conf.ip]; + for(ifc = tcpbase; ifc < etab; ifc++){ + if(ifc->newcon == s) { + ifc->newcon = 0; + tcpflushincoming(ifc); + } + } + qunlock(s); + + qlock(tcb); + close_self(s, 0); + qunlock(tcb); + break; + + case Closed: case Syn_sent: + qlock(tcb); close_self(s, 0); + qunlock(tcb); break; case Syn_received: diff --git a/port/ipdat.h b/port/ipdat.h index 3e138b1677275d6b79617eb45e78b3ef45fc4ddb..4b4d15ac3ef91576301ea41fe5f4e86b1f190178 100644 --- a/port/ipdat.h +++ b/port/ipdat.h @@ -231,7 +231,7 @@ struct Tctl char flags; char tos; - Block *rcvq; + Blist rcvq; ulong rcvcnt; Block *sndq; /* List of data going out */ @@ -249,8 +249,8 @@ struct Tctl struct Tcpctl { QLock; - struct Tctl; Rendez syner; + struct Tctl; }; struct Tcp @@ -499,6 +499,7 @@ void ipstatusfill(Chan*, char*, int); int ipforme(uchar*); void ipsetaddrs(void); int ipconbusy(Ipconv*); +void tcpflushincoming(Ipconv*); #define fmtaddr(xx) (xx>>24)&0xff,(xx>>16)&0xff,(xx>>8)&0xff,xx&0xff #define MIN(a, b) ((a) < (b) ? (a) : (b)) diff --git a/port/net.c b/port/net.c index 011fdc267cd22a0d2a07d433e2ecefe66ffd91a0..194fd987618cd1c3ec5060f5011206de89ad128f 100644 --- a/port/net.c +++ b/port/net.c @@ -268,5 +268,6 @@ netown(Network *np, int id, char *o, int omode) void netdisown(Network *np, int id) { +if(np == 0) panic("np == 0"); *np->prot[id].owner = 0; } diff --git a/port/tcpif.c b/port/tcpif.c index 830c6657b01de8138866c4f6fb7e91abeaf62f3c..ca3856aee6d8b6ec6a5cc56f8e6d6fdd777e08b5 100644 --- a/port/tcpif.c +++ b/port/tcpif.c @@ -46,9 +46,12 @@ state_upcall(Ipconv *s, char oldstate, char newstate) qunlock(s); nexterror(); } - if(s->readq == 0) - freeb(bp); - else + if(s->readq == 0){ + if(newstate == Close_wait) + putb(&tcb->rcvq, bp); + else + freeb(bp); + } else PUTNEXT(s->readq, bp); poperror(); qunlock(s); diff --git a/port/tcpinput.c b/port/tcpinput.c index d0e9770522b5df19053af5cc21880dc0895f0db1..4a9c25aa24cfa28332749d6ad73f39077b642c6b 100644 --- a/port/tcpinput.c +++ b/port/tcpinput.c @@ -79,25 +79,56 @@ reset(Ipaddr source, Ipaddr dest, char tos, ushort length, Tcp *seg) PUTNEXT(Ipoutput, hbp); } +/* + * flush an incoming call; send a reset to the remote side and close the + * conversation + */ +void +tcpflushincoming(Ipconv *s) +{ + Tcp seg; + Tcpctl *tcb; + + tcb = &s->tcpctl; + seg.source = s->pdst; + seg.dest = s->psrc; + seg.flags = ACK; + seg.seq = tcb->snd.ptr; + seg.ack = tcb->last_ack = tcb->rcv.nxt; + + reset(s->dst, Myip[Myself], 0, 0, &seg); + close_self(s, 0); +} + +static void +tcpmove(struct Tctl *to, struct Tctl *from) +{ + memmove(to, from, sizeof(struct Tctl)); +} + Ipconv* tcpincoming(Ipconv *ipc, Ipconv *s, Tcp *segp, Ipaddr source) { Ipconv *new; - if(s->curlog >= s->backlog) + qlock(s); + if(s->curlog >= s->backlog){ + qunlock(s); return 0; + } new = ipincoming(ipc, s); - if(new == 0) + if(new == 0){ + qunlock(s); return 0; + } - qlock(s); s->curlog++; qunlock(s); new->psrc = segp->dest; new->pdst = segp->source; new->dst = source; - memmove(&new->tcpctl, &s->tcpctl, sizeof(Tcpctl)); + tcpmove(&new->tcpctl, &s->tcpctl); new->tcpctl.flags &= ~CLONE; new->tcpctl.timer.arg = new; new->tcpctl.timer.state = TIMER_STOP; @@ -269,7 +300,8 @@ tcp_input(Ipconv *ipc, Block *bp) /* If we dont understand answer with a rst */ if(length) - if(s->readq == 0) { + if(s->readq == 0) + if(tcb->state == Closed) { freeb(bp); reset(source, dest, tos, length, &seg); goto done; @@ -370,9 +402,11 @@ tcp_input(Ipconv *ipc, Block *bp) case Finwait2: /* Place on receive queue */ tcb->rcvcnt += blen(bp); - if(bp) - if(s->readq) { - PUTNEXT(s->readq, bp); + if(bp){ + if(s->readq) + PUTNEXT(s->readq, bp); + else + putb(&tcb->rcvq, bp); bp = 0; } tcb->rcv.nxt += length; @@ -811,16 +845,24 @@ init_tcpctl(Ipconv *s) tcb->acktimer.arg = (void *)s; } +/* + * called with tcb locked + */ void close_self(Ipconv *s, char reason[]) { Reseq *rp,*rp1; Tcpctl *tcb = &s->tcpctl; + Block *bp; stop_timer(&tcb->timer); stop_timer(&tcb->rtt_timer); s->err = reason; + /* flush receive queue */ + while(bp = getb(&tcb->rcvq)) + freeb(bp); + /* Flush reassembly queue; nothing more can arrive */ for(rp = tcb->reseq;rp != 0;rp = rp1){ rp1 = rp->next;