M port/devip.c => port/devip.c +54 -12
@@ 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:
M port/ipdat.h => port/ipdat.h +3 -2
@@ 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))
M port/net.c => port/net.c +1 -0
@@ 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;
}
M port/tcpif.c => port/tcpif.c +6 -3
@@ 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);
M port/tcpinput.c => port/tcpinput.c +50 -8
@@ 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;