@@ 268,7 268,7 @@ ipwrite(Chan *c, char *a, long n, ulong offset)
if(cp->stproto == &tcpinfo)
tcpstart(cp, TCP_ACTIVE, Streamhi, 0);
else if(cp->stproto == &ilinfo)
- ilstart(cp, IL_ACTIVE, 10);
+ ilstart(cp, IL_ACTIVE, 20);
}
else if(strcmp(field[0], "announce") == 0) {
@@ 557,6 557,7 @@ ipremotefill(Chan *c, char *buf, int len)
cp = &ipconv[c->dev][connection];
sprint(buf, "%d.%d.%d.%d %d\n", fmtaddr(cp->dst), cp->pdst);
}
+
void
iplocalfill(Chan *c, char *buf, int len)
{
@@ 567,6 568,7 @@ iplocalfill(Chan *c, char *buf, int len)
cp = &ipconv[c->dev][connection];
sprint(buf, "%d.%d.%d.%d %d\n", fmtaddr(Myip), cp->psrc);
}
+
void
ipstatusfill(Chan *c, char *buf, int len)
{
@@ 580,8 582,8 @@ ipstatusfill(Chan *c, char *buf, int len)
tcpstate[cp->tcpctl.state],
cp->tcpctl.flags & CLONE ? "listen" : "connect");
else if(cp->stproto == &ilinfo)
- sprint(buf, "il/%d %d %s\n", connection, cp->ref,
- ilstate[cp->ilctl.state]);
+ sprint(buf, "il/%d %d %s rtt %d ms\n", connection, cp->ref,
+ ilstate[cp->ilctl.state], cp->ilctl.rtt);
else
sprint(buf, "%s/%d %d\n", cp->stproto->name, connection, cp->ref);
}
@@ 616,7 618,6 @@ iplisten(Chan *c)
for(;;) {
sleep(&s->listenr, iphavecon, s);
- print("listen wakes\n");
poperror();
new = base;
for(etab = &base[conf.ip]; new < etab; new++) {
@@ 18,26 18,31 @@ static Rendez ilackr;
char *ilstate[] = { "Closed", "Syncer", "Syncee", "Established", "Listening", "Closing" };
char *iltype[] = { "sync", "data", "dataquerey", "ack", "querey", "state", "close" };
+static char *etime = "connection timed out";
+/* Always Acktime < Fasttime < Slowtime << Ackkeepalive */
enum
{
- Mstime = 200,
- Slowtime = Mstime*20,
- Fasttime = Mstime,
+ Iltickms = 100,
+ Slowtime = 20*Iltickms,
+ Fasttime = 5*Iltickms,
+ Acktime = 3*Iltickms,
+ Ackkeepalive = 1000*Iltickms,
+ Defaultwin = 20,
};
+#define Backoff(s) (s)*=2
+#define Starttimer(s) {(s)->timeout = 0; (s)->fasttime = Fasttime;}
+
void ilrcvmsg(Ipconv*, Block*);
void ilackproc(void*);
-void ilsendctl(Ipconv*, Ilhdr*, int);
+void ilsendctl(Ipconv*, Ilhdr*, int, ulong, ulong);
void ilackq(Ilcb*, Block*);
void ilprocess(Ipconv*, Ilhdr*, Block*);
void ilpullup(Ipconv*);
-void ilhangup(Ipconv*, char*);
+void ilhangup(Ipconv*, char *);
void ilfreeq(Ilcb*);
-
-char Crefused[] = "connection refused";
-char Ctimedout[] = "connection timed out";
-char Creset[] = "connection reset by peer";
+void ilrexmit(Ilcb*);
void
ilopen(Queue *q, Stream *s)
@@ 88,11 93,13 @@ ilclose(Queue *q)
case Ilestablished:
ilfreeq(ic);
ic->state = Ilclosing;
- ilsendctl(s, 0, Ilclose);
+ ilsendctl(s, 0, Ilclose, ic->next, ic->recvd);
break;
case Illistening:
ic->state = Ilclosed;
s->psrc = 0;
+ s->pdst = 0;
+ s->dst = 0;
break;
}
netdisown(&s->ipinterface->net, s->index);
@@ 160,8 167,9 @@ iloput(Queue *q, Block *bp)
/* Checksum of ilheader plus data (not ip & no pseudo header) */
if(ilcksum)
hnputs(ih->ilsum, ptcl_csum(bp, IL_EHSIZE, dlen+IL_HDRSIZE));
-
ilackq(ic, bp);
+ ic->acktime = Ackkeepalive;
+
PUTNEXT(q, bp);
}
@@ 172,11 180,13 @@ ilackq(Ilcb *ic, Block *bp)
/* Enqueue a copy on the unacked queue in case this one gets lost */
np = copyb(bp, blen(bp));
+ qlock(&ic->ackq);
if(ic->unacked)
ic->unackedtail->list = np;
else
ic->unacked = np;
ic->unackedtail = np;
+ qunlock(&ic->ackq);
np->list = 0;
}
@@ 185,19 195,21 @@ ilackto(Ilcb *ic, ulong ackto)
{
Ilhdr *h;
Block *bp;
- ulong ack;
+ ulong id;
+ qlock(&ic->ackq);
while(ic->unacked) {
h = (Ilhdr *)ic->unacked->rptr;
- ack = nhgetl(h->ilack);
- if(ackto < ack)
+ id = nhgetl(h->ilid);
+ if(ackto < id)
break;
- ic->lastack = ackto;
+
bp = ic->unacked;
ic->unacked = bp->list;
bp->list = 0;
freeb(bp);
}
+ qunlock(&ic->ackq);
}
void
@@ 239,7 251,6 @@ ilrcvmsg(Ipconv *ipc, Block *bp)
etab = &ipc[conf.ip];
for(s = ipc; s < etab; s++)
- if(s->ilctl.state != Ilclosed)
if(s->psrc == sp)
if(s->pdst == dp)
if(s->dst == dst) {
@@ 257,11 268,10 @@ ilrcvmsg(Ipconv *ipc, Block *bp)
if(s->dst == 0) {
if(s->curlog > s->backlog)
goto reset;
+
new = ipincoming(ipc, s);
- if(new == 0) {
- print("incoming\n");
+ if(new == 0)
goto reset;
- }
new->newcon = 1;
new->ipinterface = s->ipinterface;
@@ 271,13 281,13 @@ ilrcvmsg(Ipconv *ipc, Block *bp)
ic = &new->ilctl;
ic->state = Ilsyncee;
-/*
initseq += TK2MS(MACHP(0)->ticks);
-*/
-initseq =1; ic->next = initseq;
- ic->start = ic->next;
+ ic->start = initseq;
+ ic->next = ic->start+1;
ic->recvd = 0;
ic->rstart = nhgetl(ih->ilid);
+ ic->slowtime = Slowtime;
+ ic->window = Defaultwin;
ilprocess(new, ih, bp);
s->ipinterface->ref++;
@@ 287,12 297,10 @@ initseq =1; ic->next = initseq;
}
}
drop:
- print("drop\n");
freeb(bp);
return;
reset:
- print("reset\n");
- ilsendctl(0, ih, Ilclose);
+ ilsendctl(0, ih, Ilclose, 0, 0);
freeb(bp);
}
@@ 307,7 315,6 @@ _ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
ack = nhgetl(h->ilack);
ic = &s->ilctl;
- ic->timeout = 0;
switch(ic->state) {
default:
panic("il unknown state");
@@ 321,20 328,21 @@ _ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
case Ilsync:
if(ack != ic->start) {
ic->state = Ilclosed;
- ilhangup(s, Crefused);
+ ilhangup(s, "connection rejected");
}
else {
ic->recvd = id;
ic->rstart = id;
- ilsendctl(s, 0, Ilack);
+ ilsendctl(s, 0, Ilack, ic->next, ic->recvd);
ic->state = Ilestablished;
ilpullup(s);
+ Starttimer(ic);
}
break;
case Ilclose:
if(ack == ic->start) {
ic->state = Ilclosed;
- ilhangup(s, Crefused);
+ ilhangup(s, "remote close");
}
break;
}
@@ 349,19 357,21 @@ _ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
ic->state = Ilclosed;
else {
ic->recvd = id;
- ilsendctl(s, 0, Ilsync);
+ ilsendctl(s, 0, Ilsync, ic->start, ic->recvd);
+ Starttimer(ic);
}
break;
case Ilack:
if(ack == ic->start) {
ic->state = Ilestablished;
ilpullup(s);
+ Starttimer(ic);
}
break;
case Ilclose:
- if(ack == ic->start) {
+ if(id == ic->next) {
ic->state = Ilclosed;
- ilhangup(s, Crefused);
+ ilhangup(s, "remote close");
}
break;
}
@@ 372,56 382,54 @@ _ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
case Ilsync:
if(id != ic->start) {
ic->state = Ilclosed;
- ilhangup(s, Creset);
+ ilhangup(s, "remote close");
+ }
+ else {
+ ilsendctl(s, 0, Ilack, ic->next, ic->rstart);
+ Starttimer(ic);
}
- else
- ilsendctl(s, 0, Ilack);
freeb(bp);
break;
case Ildata:
+ Starttimer(ic);
+ ilackto(ic, ack);
+ ic->acktime = Acktime;
+ iloutoforder(s, h, bp);
+ ilpullup(s);
+ break;
case Ildataquery:
- if(id < ic->recvd) {
- freeb(bp);
- break;
- }
- if(ack >= ic->recvd)
- ilackto(ic, ack);
+ Starttimer(ic);
+ ilackto(ic, ack);
+ ic->acktime = Acktime;
iloutoforder(s, h, bp);
ilpullup(s);
- if(h->iltype == Ildataquery)
- ilsendctl(s, 0, Ilstate);
+ ilsendctl(s, 0, Ilstate, ic->next, ic->recvd);
break;
case Ilack:
ilackto(ic, ack);
+ Starttimer(ic);
freeb(bp);
break;
case Ilquerey:
ilackto(ic, ack);
- ilsendctl(s, 0, Ilstate);
+ ilsendctl(s, 0, Ilstate, ic->next, ic->recvd);
+ Starttimer(ic);
freeb(bp);
break;
case Ilstate:
ilackto(ic, ack);
- if(ic->unacked) {
- nb = copyb(ic->unacked, blen(ic->unacked));
- h = (Ilhdr*)nb;
- h->iltype = Ildataquery;
- hnputl(h->ilack, ic->recvd);
- h->ilsum[0] = 0;
- h->ilsum[1] = 0;
- if(ilcksum)
- hnputs(h->ilsum, ptcl_csum(nb, IL_EHSIZE, nhgets(h->illen)));
- PUTNEXT(Ipoutput, nb);
- }
+ ilrexmit(ic);
+ Starttimer(ic);
freeb(bp);
break;
case Ilclose:
freeb(bp);
- if(id != ic->recvd)
+ if(ack < ic->start || ack > ic->next)
break;
- ilsendctl(s, 0, Ilclose);
+ ilsendctl(s, 0, Ilclose, ic->next, ic->recvd);
ic->state = Ilclosing;
ilfreeq(ic);
+ Starttimer(ic);
break;
}
break;
@@ 432,11 440,12 @@ _ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
switch(h->iltype) {
case Ilclose:
ic->recvd = id;
+ ilsendctl(s, 0, Ilclose, ic->next, ic->recvd);
if(ack == ic->next) {
ic->state = Ilclosed;
ilhangup(s, 0);
}
- ilsendctl(s, 0, Ilclose);
+ Starttimer(ic);
break;
default:
break;
@@ 446,6 455,29 @@ _ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
}
}
+void
+ilrexmit(Ilcb *ic)
+{
+ Block *nb;
+ Ilhdr *h;
+
+ if(ic->unacked == 0)
+ return;
+
+ nb = copyb(ic->unacked, blen(ic->unacked));
+ h = (Ilhdr*)nb->rptr;
+ DBG("rxmit %d.", nhgetl(h->ilid));
+
+ h->iltype = Ildataquery;
+ hnputl(h->ilack, ic->recvd);
+ h->ilsum[0] = 0;
+ h->ilsum[1] = 0;
+ if(ilcksum)
+ hnputs(h->ilsum, ptcl_csum(nb, IL_EHSIZE, nhgets(h->illen)));
+
+ PUTNEXT(Ipoutput, nb);
+}
+
/* DEBUG */
void
ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
@@ 453,22 485,23 @@ ilprocess(Ipconv *s, Ilhdr *h, Block *bp)
Ilcb *ic = &s->ilctl;
USED(ic);
- DBG("%-11s rcv %d/%d snt %d/%d pkt(%-6s id %d ack %d %d->%d) ",
+ DBG("%11s rcv %d/%d snt %d/%d pkt(%s id %d ack %d %d->%d) ",
ilstate[ic->state], ic->rstart, ic->recvd, ic->start, ic->next,
iltype[h->iltype], nhgetl(h->ilid), nhgetl(h->ilack),
nhgets(h->ilsrc), nhgets(h->ildst));
_ilprocess(s, h, bp);
- DBG("%-11s rcv %d snt %d\n", ilstate[ic->state], ic->recvd, ic->next);
+ DBG("%11s rcv %d snt %d\n", ilstate[ic->state], ic->recvd, ic->next);
}
void
ilhangup(Ipconv *s, char *msg)
{
Block *nb;
- ulong l;
+ int l;
+ DBG("hangup! %s %d/%d\n", msg ? msg : "??", s->psrc, s->pdst);
if(s->readq) {
if(msg) {
l = strlen(msg);
@@ 482,6 515,9 @@ ilhangup(Ipconv *s, char *msg)
nb->flags |= S_DELIM;
PUTNEXT(s->readq, nb);
}
+ s->psrc = 0;
+ s->pdst = 0;
+ s->dst = 0;
}
void
@@ 499,25 535,31 @@ ilpullup(Ipconv *s)
if(ic->state != Ilestablished)
return;
+ qlock(&ic->outo);
while(ic->outoforder) {
bp = ic->outoforder;
oh = (Ilhdr*)bp->rptr;
oid = nhgetl(oh->ilid);
- if(oid > ic->recvd)
- break;
- if(oid < ic->recvd) {
+ if(oid <= ic->recvd) {
ic->outoforder = bp->list;
freeb(bp);
+ continue;
}
- if(oid == ic->recvd) {
- ic->recvd++;
- ic->outoforder = bp->list;
- bp->list = 0;
- dlen = nhgets(oh->illen)-IL_HDRSIZE;
- bp = btrim(bp, IL_EHSIZE+IL_HDRSIZE, dlen);
- PUTNEXT(s->readq, bp);
- }
+ if(oid != ic->recvd+1)
+ break;
+
+ ic->recvd = oid;
+ ic->outoforder = bp->list;
+ ic->oblks--;
+
+ qunlock(&ic->outo);
+ bp->list = 0;
+ dlen = nhgets(oh->illen)-IL_HDRSIZE;
+ bp = btrim(bp, IL_EHSIZE+IL_HDRSIZE, dlen);
+ PUTNEXT(s->readq, bp);
+ qlock(&ic->outo);
}
+ qunlock(&ic->outo);
}
void
@@ 530,30 572,43 @@ iloutoforder(Ipconv *s, Ilhdr *h, Block *bp)
ic = &s->ilctl;
bp->list = 0;
- if(ic->outoforder == 0) {
- ic->outoforder = bp;
+
+
+ id = nhgetl(h->ilid);
+ /* Window checks */
+ if(id <= ic->recvd || ic->oblks > ic->window) {
+ freeb(bp);
return;
}
- id = nhgetl(h->id);
- l = &ic->outoforder;
- for(f = *l; f; f = f->list) {
- lid = ((Ilhdr*)(bp->rptr))->ilid;
- if(id > nhgetl(lid))
- break;
- l = &f->list;
+ /* Packet is acceptable so sort onto receive queue for pullup */
+ qlock(&ic->outo);
+ ic->oblks++;
+ if(ic->outoforder == 0)
+ ic->outoforder = bp;
+ else {
+ l = &ic->outoforder;
+ for(f = *l; f; f = f->list) {
+ lid = ((Ilhdr*)(f->rptr))->ilid;
+ if(id < nhgetl(lid)) {
+ bp->list = f;
+ *l = bp;
+ qunlock(&ic->outo);
+ return;
+ }
+ l = &f->list;
+ }
+ *l = bp;
}
- bp->list = *l;
- *l = bp;
+ qunlock(&ic->outo);
}
void
-ilsendctl(Ipconv *ipc, Ilhdr *inih, int type)
+ilsendctl(Ipconv *ipc, Ilhdr *inih, int type, ulong id, ulong ack)
{
Ilhdr *ih;
Ilcb *ic;
Block *bp;
- ulong id;
bp = allocb(IL_EHSIZE+IL_HDRSIZE);
bp->wptr += IL_EHSIZE+IL_HDRSIZE;
@@ 577,11 632,9 @@ ilsendctl(Ipconv *ipc, Ilhdr *inih, int type)
hnputl(ih->dst, ipc->dst);
hnputs(ih->ilsrc, ipc->psrc);
hnputs(ih->ildst, ipc->pdst);
- id = ic->next;
- if(type == Ilsync)
- id = ic->start;
hnputl(ih->ilid, id);
- hnputl(ih->ilack, ic->recvd);
+ hnputl(ih->ilack, ack);
+ ic->acktime = Ackkeepalive;
}
ih->iltype = type;
ih->ilspec = 0;
@@ 607,51 660,54 @@ ilackproc(void *a)
base = (Ipconv*)a;
end = &base[conf.ip];
+
for(;;) {
- tsleep(&ilackr, return0, 0, Mstime);
+ tsleep(&ilackr, return0, 0, Iltickms);
for(s = base; s < end; s++) {
ic = &s->ilctl;
+ ic->timeout += Iltickms;
switch(ic->state) {
- default:
+ case Ilclosed:
+ case Illistening:
break;
case Ilclosing:
- ic->timeout++;
- if(ic->timeout >= Slowtime) {
+ if(ic->timeout >= ic->fasttime) {
+ ilsendctl(s, 0, Ilclose, ic->next, ic->recvd);
+ Backoff(ic->fasttime);
+ }
+ if(ic->timeout >= ic->slowtime) {
ic->state = Ilclosed;
ilhangup(s, 0);
}
break;
case Ilsyncee:
- ic->timeout++;
- if(ic->timeout >= Slowtime) {
- ic->state = Ilclosed;
- ilhangup(s, Ctimedout);
- break;
- }
-print("Rxmit %d/%d %s", s->psrc, s->pdst, ilstate[ic->state]);
- ilsendctl(s, 0, Ilsync);
- break;
case Ilsyncer:
- ic->timeout++;
- if(ic->timeout >= Slowtime) {
+ if(ic->timeout >= ic->fasttime) {
+ ilsendctl(s, 0, Ilsync, ic->start, ic->recvd);
+ Backoff(ic->fasttime);
+ }
+ if(ic->timeout >= ic->slowtime) {
ic->state = Ilclosed;
- ilhangup(s, Ctimedout);
- break;
+ ilhangup(s, etime);
}
-print("Rxmit %d/%d %s", s->psrc, s->pdst, ilstate[ic->state]);
- ilsendctl(s, 0, Ilsync);
break;
case Ilestablished:
- ic->timeout++;
- if(ic->unacked == 0)
+ ic->acktime -= Iltickms;
+ if(ic->acktime <= 0)
+ ilsendctl(s, 0, Ilack, ic->next, ic->recvd);
+ if(ic->unacked == 0) {
+ ic->timeout = 0;
break;
- if(ic->timeout >= Slowtime) {
+ }
+ if(ic->timeout >= ic->fasttime) {
+ ilrexmit(ic);
+ Backoff(ic->fasttime);
+ }
+ if(ic->timeout >= ic->slowtime) {
ic->state = Ilclosed;
- ilhangup(s, Ctimedout);
+ ilhangup(s, etime);
break;
}
-print("Rxmit %d/%d %s", s->psrc, s->pdst, ilstate[ic->state]);
- ilsendctl(s, 0, Ilstate);
break;
}
}
@@ 669,14 725,13 @@ ilstart(Ipconv *ipc, int type, int window)
ic->timeout = 0;
ic->unacked = 0;
ic->outoforder = 0;
-/*
+ ic->slowtime = Slowtime;
+
+
initseq += TK2MS(MACHP(0)->ticks);
-*/
-initseq = 1;
- ic->next = initseq;
- ic->start = ic->next;
+ ic->start = initseq;
+ ic->next = ic->start+1;
ic->recvd = 0;
- ic->lastack = ic->next;
ic->window = window;
switch(type) {
@@ 685,7 740,7 @@ initseq = 1;
break;
case IL_ACTIVE:
ic->state = Ilsyncer;
- ilsendctl(ipc, 0, Ilsync);
+ ilsendctl(ipc, 0, Ilsync, ic->start, ic->recvd);
break;
}
}
@@ 695,14 750,19 @@ ilfreeq(Ilcb *ic)
{
Block *bp, *next;
+ qlock(&ic->ackq);
for(bp = ic->unacked; bp; bp = next) {
next = bp->list;
freeb(bp);
}
+ ic->unacked = 0;
+ qunlock(&ic->ackq);
+
+ qlock(&ic->outo);
for(bp = ic->outoforder; bp; bp = next) {
next = bp->list;
freeb(bp);
}
- ic->unacked = 0;
ic->outoforder = 0;
+ qunlock(&ic->outo);
}