@@ 117,17 117,20 @@ struct Line {
*/
struct Dk {
Lock;
+ int opened;
char name[64]; /* dk name */
Queue *wq; /* dk output queue */
Stream *s;
int lines; /* number of lines */
int ncsc; /* csc line number */
- Chan *csc; /* common signalling line */
+ Chan *csc;
Line line[Nline];
int restart;
int urpwindow;
+ Rendez timer;
};
static Dk dk[Ndk];
+static Lock dklock;
/*
* conversation states (for Line.state)
@@ 178,15 181,16 @@ extern Qinfo urpinfo;
* predeclared
*/
Chan* dkattach(char*);
-static void dkmuxconfig(Dk*, Block*);
-static int dkmesg(Dk*, int, int, int, int);
+static void dkmuxconfig(Queue*, Block*);
+static Chan* dkopenline(Dk*, int);
+static int dkmesg(Chan*, int, int, int, int);
static void dkcsckproc(void*);
static int dklisten(Chan*);
static void dkanswer(Chan*, int, int);
static void dkwindow(Chan*);
static void dkcall(int, Chan*, char*, char*, char*);
static void dktimer(void*);
-static void dkchgmesg(Dk*, Dkmsg*, int);
+static void dkchgmesg(Chan*, Dk*, Dkmsg*, int);
static void dkreplymesg(Dk*, Dkmsg*, int);
Chan* dkopen(Chan*, int);
static void dkhangup(Line*);
@@ 201,35 205,46 @@ static void dkmuxiput(Queue *, Block *);
Qinfo dkmuxinfo = { dkmuxiput, dkmuxoput, dkmuxopen, dkmuxclose, "dkmux" };
/*
- * a new dkmux. find a free dk structure and assign it to this queue.
- * when we get though here dp->s is meaningful and the name is set to "/".
- */
-static void
-dkmuxopen(Queue *q, Stream *s)
+ * Look for a dk struct with a name. If none exists, create one.
+ */
+static Dk *
+dkalloc(char *name)
{
Dk *dp;
- int i;
+ Dk *freep;
+ lock(&dklock);
+ freep = 0;
for(dp = dk; dp < &dk[Ndk]; dp++){
- if(dp->name[0]==0){
- lock(dp);
- if(dp->name[0]){
- /* someone was faster than us */
- unlock(dp);
- continue;
- }
- q->ptr = q->other->ptr = (void *)dp;
- dp->csc = 0;
- dp->ncsc = 4;
- dp->lines = 16;
- strcpy(dp->name, "/");
- dp->wq = WR(q);
- dp->s = s;
- unlock(dp);
- return;
+ if(strcmp(name, dp->name) == 0){
+ unlock(&dklock);
+ return dp;
}
+ if(dp->name[0] == 0)
+ freep = dp;
+ }
+ if(freep == 0){
+ unlock(&dklock);
+ error(0, Enoifc);
}
- error(0, Enoifc);
+ dp = freep;
+ dp->opened = 0;
+ dp->s = 0;
+ dp->ncsc = 1;
+ strncpy(dp->name, name, sizeof(freep->name));
+ unlock(&dklock);
+ return dp;
+}
+
+/*
+ * a new dkmux. find a free dk structure and assign it to this queue.
+ * when we get though here dp->s is meaningful and the name is set to "/".
+ */
+static void
+dkmuxopen(Queue *q, Stream *s)
+{
+ RD(q)->ptr = s;
+ WR(q)->ptr = 0;
}
/*
@@ 241,20 256,28 @@ dkmuxclose(Queue *q)
Dk *dp;
int i;
- dp = (Dk *)q->ptr;
+ dp = WR(q)->ptr;
+ if(dp == 0)
+ return;
/*
- * if we're the last user of the stream,
- * free the Dk structure
+ * disallow new dkstopens() on this line.
+ * the lock syncs with dkstopen().
*/
- if(dp->s->inuse == 1)
- dp->name[0] = 0;
+ lock(dp);
+ dp->opened = 0;
+ unlock(dp);
/*
* hang up all datakit connections
*/
for(i=dp->ncsc; i < dp->lines; i++)
dkhangup(&dp->line[i]);
+
+ /*
+ * wakeup the timer so it can die
+ */
+ wakeup(&dp->timer);
}
/*
@@ 263,12 286,9 @@ dkmuxclose(Queue *q)
static void
dkmuxoput(Queue *q, Block *bp)
{
- Dk *dp;
-
- dp = (Dk *)q->ptr;
if(bp->type != M_DATA){
if(streamparse("config", bp))
- dkmuxconfig(dp, bp);
+ dkmuxconfig(q, bp);
else
PUTNEXT(q, bp);
return;
@@ 293,6 313,14 @@ dkmuxiput(Queue *q, Block *bp)
Line *lp;
int line;
+ /*
+ * not configured yet
+ */
+ if(q->other->ptr == 0){
+ freeb(bp);
+ return;
+ }
+
dp = (Dk *)q->ptr;
if(bp->type != M_DATA){
PUTNEXT(q, bp);
@@ 341,6 369,12 @@ dkstopen(Queue *q, Stream *s)
Line *lp;
dp = &dk[s->dev];
+ lock(dp);
+ if(dp->opened==0 || streamenter(dp->s)<0){
+ unlock(dp);
+ error(0, Ehungup);
+ }
+ unlock(dp);
q->other->ptr = q->ptr = lp = &dp->line[s->id];
lp->dp = dp;
lp->rq = q;
@@ 356,11 390,32 @@ dkstclose(Queue *q)
{
Dk *dp;
Line *lp;
+ Chan *c;
lp = (Line *)q->ptr;
dp = lp->dp;
/*
+ * decrement ref count on mux'd line
+ */
+ streamexit(dp->s, 0);
+
+ /*
+ * don't tell controller about closing down the csc
+ */
+ if(lp - dp->line == dp->ncsc)
+ goto out;
+
+ c = 0;
+ if(waserror()){
+ if(c)
+ close(c);
+ lp->state = Lclosed;
+ goto out;
+ }
+ c = dkopenline(dp, dp->ncsc);
+
+ /*
* shake hands with dk
*/
switch(lp->state){
@@ 369,29 424,32 @@ dkstclose(Queue *q)
break;
case Lrclose:
- dkmesg(dp, T_CHG, D_CLOSE, lp - dp->line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, lp - dp->line, 0);
lp->state = Lclosed;
break;
case Lackwait:
- dkmesg(dp, T_CHG, D_CLOSE, lp - dp->line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, lp - dp->line, 0);
lp->state = Llclose;
break;
case Llistening:
- dkmesg(dp, T_CHG, D_CLOSE, lp - dp->line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, lp - dp->line, 0);
lp->state = Llclose;
break;
case Lconnected:
- dkmesg(dp, T_CHG, D_CLOSE, lp - dp->line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, lp - dp->line, 0);
lp->state = Llclose;
break;
case Lopened:
lp->state = Lclosed;
}
+ poperror();
+ close(c);
+out:
qlock(lp);
lp->rq = 0;
qunlock(lp);
@@ 448,87 506,102 @@ dkoput(Queue *q, Block *bp)
* we can configure only once
*/
static void
-dkmuxconfig(Dk *dp, Block *bp)
+dkmuxconfig(Queue *q, Block *bp)
{
- Chan *c;
+ Dk *dp;
char *fields[5];
int n;
char buf[64];
- static int dktimeron;
+ char name[NAMELEN];
+ int lines;
+ int ncsc;
+ int restart;
+ int window;
- if(dp->csc != 0){
+ if(WR(q)->ptr){
freeb(bp);
- error(0, Ebadarg);
+ error(0, Egreg);
}
/*
+ * defaults
+ */
+ ncsc = 1;
+ restart = 1;
+ lines = 16;
+ window = WS_2K;
+ strcpy(name, "dk");
+
+ /*
* parse
*/
- dp->restart = 1;
n = getfields((char *)bp->rptr, fields, 5, ' ');
- strcpy(dp->name, "dk");
- dp->urpwindow = WS_2K;
switch(n){
case 5:
- dp->urpwindow = strtoul(fields[4], 0, 0);
+ window = strtoul(fields[4], 0, 0);
case 4:
- strncpy(dp->name, fields[3], sizeof(dp->name));
+ strncpy(name, fields[3], sizeof(name));
case 3:
if(strcmp(fields[2], "restart")!=0)
- dp->restart = 0;
+ restart = 0;
case 2:
- dp->lines = strtoul(fields[1], 0, 0);
+ lines = strtoul(fields[1], 0, 0);
case 1:
- dp->ncsc = strtoul(fields[0], 0, 0);
+ ncsc = strtoul(fields[0], 0, 0);
break;
default:
freeb(bp);
error(0, Ebadarg);
}
freeb(bp);
- if(dp->ncsc <= 0 || dp->lines <= dp->ncsc){
- dp->lines = 16;
+ if(ncsc <= 0 || lines <= ncsc)
error(0, Ebadarg);
- }
- DPRINT("dkmuxconfig: ncsc=%d, lines=%d, restart=%d, name=\"%s\"\n",
- dp->ncsc, dp->lines, dp->restart, dp->name);
/*
- * open a stream for the csc and push urp onto it
+ * set up
*/
- c = 0;
- if(waserror()){
- if(c)
- close(c);
- nexterror();
+ dp = dkalloc(name);
+ lock(dp);
+ if(dp->opened){
+ unlock(dp);
+ error(0, Ebadarg);
}
- c = dkattach(dp->name);
- c->qid = STREAMQID(dp->ncsc, Sdataqid);
- dkopen(c, ORDWR);
- dp->csc = c;
+ dp->ncsc = ncsc;
+ dp->lines = lines;
+ dp->restart = restart;
+ dp->urpwindow = window;
+ dp->s = RD(q)->ptr;
+ q->ptr = q->other->ptr = dp;
+ dp->opened = 1;
+ dp->wq = WR(q);
+ unlock(dp);
+
+ /*
+ * open csc here so that boot, dktimer, and dkcsckproc aren't
+ * all fighting for it at once.
+ */
+ dp->csc = dkopenline(dp, dp->ncsc);
/*
* tell datakit we've rebooted. It should close all channels.
+ * do this here to get it done before trying to open a channel.
*/
if(dp->restart) {
- DPRINT("dkmuxconfig: restart %s\n", dp->name);
- dkmesg(dp, T_ALIVE, D_RESTART, 0, 0);
+ DPRINT("dktimer: restart %s\n", dp->name);
+ dkmesg(dp->csc, T_ALIVE, D_RESTART, 0, 0);
}
/*
- * start a process to deal with it
+ * start a process to listen to csc messages
*/
sprint(buf, "csc.%s.%d", dp->name, dp->ncsc);
kproc(buf, dkcsckproc, dp);
- poperror();
/*
- * start a keepalive process if one doesn't exist
+ * start a keepalive process
*/
- if(dktimeron == 0){
- dktimeron = 1;
- kproc("dktimer", dktimer, 0);
- }
+ sprint(buf, "timer.%s.%d", dp->name, dp->ncsc);
+ kproc(buf, dktimer, dp);
}
/*
@@ 619,54 692,19 @@ dkattach(char *spec)
*/
if(*spec == 0)
spec = "dk";
- for(dp = dk; dp < &dk[Ndk]; dp++){
- lock(dp);
- if(strcmp(spec, dp->name)==0)
- break;
- unlock(dp);
- }
- if(dp == &dk[Ndk])
- error(0, Enoifc);
-
- /*
- * don't let the multiplexed stream disappear under us
- */
- if(streamenter(dp->s) < 0){
- /*
- * it's closing down, forget it
- */
- unlock(dp);
- error(0, Ehungup);
- }
+ dp = dkalloc(spec);
/*
* return the new channel
*/
- if(waserror()){
- if(streamexit(dp->s, 0) == 0)
- dp->name[0] = 0;
- unlock(dp);
- nexterror();
- }
c = devattach('k', spec);
c->dev = dp - dk;
- unlock(dp);
- poperror();
return c;
}
-/*
- * clone as long as the multiplexed channel is not closing
- * down
- */
Chan*
dkclone(Chan *c, Chan *nc)
{
- Dk *dp;
-
- dp = &dk[c->dev];
- if(streamenter(dp->s) < 0)
- error(0, Ehungup);
return devclone(c, nc);
}
@@ 783,17 821,8 @@ dkclose(Chan *c)
{
Dk *dp;
- /* real closing happens in dkstclose */
if(c->stream)
streamclose(c);
-
- /*
- * Let go of the mulitplexed stream. If we're the last out,
- * free dp.
- */
- dp = &dk[c->dev];
- if(streamexit(dp->s, 0) == 0)
- dp->name[0] = 0;
}
long
@@ 900,16 929,35 @@ dkuserstr(Error *e, char *buf)
}
/*
+ * open the common signalling channel
+ */
+static Chan*
+dkopenline(Dk *dp, int line)
+{
+ Chan *c;
+
+ c = 0;
+ if(waserror()){
+ if(c)
+ close(c);
+ nexterror();
+ }
+ c = dkattach(dp->name);
+ c->qid = STREAMQID(line, Sdataqid);
+ dkopen(c, ORDWR);
+ poperror();
+
+ return c;
+}
+
+/*
* send a message to the datakit on the common signaling line
*/
static int
-dkmesg(Dk *dp, int type, int srv, int p0, int p1)
+dkmesg(Chan *c, int type, int srv, int p0, int p1)
{
Dkmsg d;
- Block *bp;
- if(dp->csc == 0)
- return -1;
if(waserror()){
print("dkmesg: error\n");
return -1;
@@ 926,7 974,7 @@ dkmesg(Dk *dp, int type, int srv, int p0, int p1)
d.param3h = 0;
d.param4l = 0;
d.param4h = 0;
- streamwrite(dp->csc, (char *)&d, sizeof(Dkmsg), 1);
+ streamwrite(c, (char *)&d, sizeof(Dkmsg), 1);
poperror();
return 0;
}
@@ 952,8 1000,9 @@ dkcall(int type, Chan *c, char *addr, char *nuser, char *machine)
Dk *dp;
Line *lp;
Chan *dc;
+ Chan *csc;
char *bang, *dot;
-
+
line = STREAMID(c->qid);
dp = &dk[c->dev];
lp = &dp->line[line];
@@ 1006,16 1055,22 @@ dkcall(int type, Chan *c, char *addr, char *nuser, char *machine)
}
/*
- * open the data file
+ * close temporary channels on error
*/
- dc = dkattach(dp->name);
+ dc = 0;
+ csc = 0;
if(waserror()){
- close(dc);
+ if(csc)
+ close(csc);
+ if(dc)
+ close(dc);
nexterror();
}
- dc->qid = STREAMQID(line, Sdataqid);
- dkopen(dc, ORDWR);
+ /*
+ * open the data file
+ */
+ dc = dkopenline(dp, line);
lp->calltolive = 4;
lp->state = Ldialing;
@@ 1023,7 1078,10 @@ dkcall(int type, Chan *c, char *addr, char *nuser, char *machine)
* tell the controller we want to make a call
*/
DPRINT("dialout\n");
- dkmesg(dp, t_val, d_val, line, W_WINDOW(dp->urpwindow,dp->urpwindow,2));
+ csc = dkopenline(dp, dp->ncsc);
+ dkmesg(csc, t_val, d_val, line, W_WINDOW(dp->urpwindow,dp->urpwindow,2));
+ close(csc);
+ csc = 0;
/*
* if redial, wait for a dial tone (otherwise we might send
@@ 1354,18 1412,11 @@ dkcsckproc(void *a)
Dk *dp;
Dkmsg d;
int line;
- int i;
- dp = (Dk *)a;
+ dp = a;
if(waserror()){
- Chan *csc;
-
- csc = dp->csc;
- lock(dp);
- dp->csc = 0;
- unlock(dp);
- close(csc);
+ close(dp->csc);
return;
}
@@ 1385,7 1436,7 @@ dkcsckproc(void *a)
switch (d.type) {
case T_CHG: /* controller wants to close a line */
- dkchgmesg(dp, &d, line);
+ dkchgmesg(dp->csc, dp, &d, line);
break;
case T_REPLY: /* reply to a dial request */
@@ 1415,7 1466,7 @@ dkcsckproc(void *a)
* datakit requests or confirms closing a line
*/
static void
-dkchgmesg(Dk *dp, Dkmsg *dialp, int line)
+dkchgmesg(Chan *c, Dk *dp, Dkmsg *dialp, int line)
{
Line *lp;
@@ 1424,7 1475,7 @@ dkchgmesg(Dk *dp, Dkmsg *dialp, int line)
case D_CLOSE: /* remote shutdown */
if (line <= 0 || line >= dp->lines) {
/* tell controller this line is not in use */
- dkmesg(dp, T_CHG, D_CLOSE, line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, line, 0);
return;
}
lp = &dp->line[line];
@@ 1445,13 1496,13 @@ dkchgmesg(Dk *dp, Dkmsg *dialp, int line)
break;
case Lopened:
- dkmesg(dp, T_CHG, D_CLOSE, line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, line, 0);
break;
case Llclose:
case Lclosed:
dkhangup(lp);
- dkmesg(dp, T_CHG, D_CLOSE, line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, line, 0);
lp->state = Lclosed;
break;
}
@@ 1460,7 1511,7 @@ dkchgmesg(Dk *dp, Dkmsg *dialp, int line)
case D_ISCLOSED: /* acknowledging a local shutdown */
if (line <= 0 || line >= dp->lines) {
/* tell controller this line is not in use */
- dkmesg(dp, T_CHG, D_CLOSE, line, 0);
+ dkmesg(c, T_CHG, D_CLOSE, line, 0);
return;
}
lp = &dp->line[line];
@@ 1554,62 1605,57 @@ dkreplymesg(Dk *dp, Dkmsg *dialp, int line)
}
/*
- * 15-second timer for all interfaces
+ * send a I'm alive message every 7.5 seconds and remind the dk of
+ * any closed channels it hasn't acknowledged.
*/
-static Rendez dkt;
-static int
-fuckit(void *a)
-{
- return 0;
-}
static void
dktimer(void *a)
{
int dki, i;
Dk *dp;
Line *lp;
+ Chan *c;
+
+ c = 0;
+ if(waserror()){
+ if(c)
+ close(c);
+ return;
+ }
- while(waserror())
- print("dktimer: error\n");
+ /*
+ * open csc
+ */
+ dp = (Dk *)a;
+ c = dkopenline(dp, dp->ncsc);
for(;;){
+ if(dp->opened==0)
+ error(0, Ehungup);
+
/*
- * loop through the active dks
+ * send keep alive
*/
- for(dki=0; dki<Ndk; dki++){
- dp = &dk[dki];
- if(!canlock(dp))
- continue;
- if(dp->csc==0){
- unlock(dp);
- continue;
- }
-
- /*
- * send keep alive
- */
- dkmesg(dp, T_ALIVE, D_CONTINUE, 0, 0);
+ dkmesg(c, T_ALIVE, D_CONTINUE, 0, 0);
- /*
- * remind controller of dead lines and
- * timeout calls that take to long
- */
- for (i=dp->ncsc+1; i<dp->lines; i++){
- lp = &dp->line[i];
- switch(lp->state){
- case Llclose:
- dkmesg(dp, T_CHG, D_CLOSE, i, 0);
- break;
+ /*
+ * remind controller of dead lines and
+ * timeout calls that take to long
+ */
+ for (i=dp->ncsc+1; i<dp->lines; i++){
+ lp = &dp->line[i];
+ switch(lp->state){
+ case Llclose:
+ dkmesg(c, T_CHG, D_CLOSE, i, 0);
+ break;
- case Ldialing:
- if(lp->calltolive==0 || --lp->calltolive!=0)
- break;
- dkreplymesg(dp, (Dkmsg *)0, i);
+ case Ldialing:
+ if(lp->calltolive==0 || --lp->calltolive!=0)
break;
- }
+ dkreplymesg(dp, (Dkmsg *)0, i);
+ break;
}
- unlock(dp);
}
- tsleep(&dkt, fuckit, 0, 7500);
+ tsleep(&dp->timer, return0, 0, 7500);
}
}