package dispenser import ( "context" "errors" "io" "reflect" "testing" "time" ) func apReply(st []byte) []byte { frame := []byte{STX, 0x30, 0x30, 0, byte(len(st) + 2), 'S', 'F'} frame = append(frame, st...) frame = append(frame, ETX) return append(frame, calculateBCC(frame)) } func workerRequest(c *Client, ctx context.Context, typ cmdType) cmdResp { ch := make(chan cmdResp, 1) c.handle(cmdReq{typ: typ, ctx: ctx, respCh: ch}) return <-ch } func TestDeliveryDispatchBoundary(t *testing.T) { transportAddress(t) for _, tc := range []struct { name string short int badACK bool cancelAfter int pending bool wantErr bool }{ {name: "success", pending: true}, {name: "command short write", short: 1, wantErr: true}, {name: "bad ACK", badACK: true, wantErr: true}, {name: "ENQ short write", short: 2, pending: true, wantErr: true}, {name: "cancel before ENQ", cancelAfter: 1, wantErr: true}, {name: "cancel after ENQ", cancelAfter: 2, pending: true, wantErr: true}, } { t.Run(tc.name, func(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() ack := append([]byte(nil), vendorACK...) if tc.badACK { ack[0] = 0 } p := &scriptedTransport{chunks: [][]byte{ack}, shortWrite: tc.short} p.afterWrite = func() { if len(p.writes) == tc.cancelAfter { cancel() } } c := &Client{port: p} r := workerRequest(c, ctx, cmdOutOfMouth) if (r.err != nil) != tc.wantErr || c.deliveryPending != tc.pending { t.Errorf("FC0(%s): err=%v pending=%v, want error=%v pending=%v", tc.name, r.err, c.deliveryPending, tc.wantErr, tc.pending) } if len(p.writes) > 2 { t.Errorf("FC0 writes=%d, want no resend", len(p.writes)) } }) } t.Run("definite failure preserves previous pending", func(t *testing.T) { c := &Client{port: &scriptedTransport{writeErr: io.ErrClosedPipe}, deliveryPending: true} workerRequest(c, context.Background(), cmdOutOfMouth) if !c.deliveryPending { t.Error("failed FC0 erased previous pending") } }) t.Run("ENQ transport error remains pending", func(t *testing.T) { p := &scriptedTransport{chunks: [][]byte{vendorACK}} p.afterWrite = func() { if len(p.writes) == 2 { p.writeErr = io.ErrClosedPipe } } c := &Client{port: p} r := workerRequest(c, context.Background(), cmdOutOfMouth) if !errors.Is(r.err, io.ErrClosedPipe) || !c.deliveryPending { t.Errorf("ENQ error=%v pending=%v, want closed pipe and pending", r.err, c.deliveryPending) } }) } func TestWorkerClearanceValidation(t *testing.T) { transportAddress(t) for _, tc := range []struct { name string st []byte clear bool }{ {"clear", status(0x30), true}, {"low stock clear", []byte{0x30, 0x30, 0x31, 0x30}, true}, {"residual encoder", status(0x33), false}, {"all sensors", status(0x37), false}, {"ready", status(0x34), true}, {"unknown diagnostic", []byte{0x30, 0x30, 0xFF, 0x30}, true}, {"jam", []byte{0x30, 0x30, 0x32, 0x30}, true}, {"overlap", []byte{0x30, 0x30, 0x34, 0x30}, true}, {"rejection", []byte{0x36, 0x30, 0x30, 0x30}, true}, {"empty", status(0x38), false}, {"preparing", []byte{0x31, 0x30, 0x30, 0x30}, false}, {"dispensing", []byte{0x30, 0x38, 0x30, 0x34}, false}, {"capturing", []byte{0x30, 0x34, 0x30, 0x30}, false}, {"ready with stale errors", []byte{0x36, 0x32, 0x34, 0x34}, true}, {"short payload", []byte{0x30, 0x30, 0x30}, false}, {"unknown position", status(0x40), false}, } { t.Run(tc.name, func(t *testing.T) { p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(tc.st), vendorACK}} c := &Client{port: p, deliveryPending: true, deliveryStarted: time.Now().Add(-deliveryMinimumWait), lastStatus: status(0x30), lastStatusT: time.Now(), statusTTL: time.Hour} r := workerRequest(c, context.Background(), cmdStatus) if c.deliveryPending == tc.clear { t.Errorf("FC7(% X): err=%v pending=%v, want clear=%v", tc.st, r.err, c.deliveryPending, tc.clear) } want := []string{"AP"} if got := wireCommands(p); !reflect.DeepEqual(got, want) { t.Errorf("commands=%v, want %v", got, want) } }) } t.Run("bad BCC cannot clear", func(t *testing.T) { frame := apReply(status(0x30)) frame[len(frame)-1] ^= 1 p := &scriptedTransport{chunks: [][]byte{vendorACK, frame}} c := &Client{port: p, deliveryPending: true} r := workerRequest(c, context.Background(), cmdStatus) if r.err == nil || !c.deliveryPending { t.Errorf("bad BCC: err=%v pending=%v", r.err, c.deliveryPending) } }) } func wireCommands(p *scriptedTransport) []string { var commands []string for _, w := range p.writes { if len(w) > 6 && w[0] == STX { commands = append(commands, string(w[5:len(w)-2])) } } return commands } func deliveryWorkerClient(t *testing.T, p *scriptedTransport) (*Client, *time.Time) { t.Helper() transportAddress(t) now := time.Unix(0, 0) c := &Client{port: p, reqCh: make(chan cmdReq, 16), done: make(chan struct{})} c.sequenceTiming = sequenceTiming{now: func() time.Time { return now }, wait: func(ctx context.Context, d time.Duration) error { if err := ctx.Err(); err != nil { return err } now = now.Add(d) return nil }} stopped := make(chan struct{}) go func() { defer close(stopped); c.loop() }() t.Cleanup(func() { c.Close(); <-stopped }) return c, &now } func TestDeliveryIncidentReplay(t *testing.T) { for _, next := range []bool{false, true} { t.Run(map[bool]string{false: "current", true: "next"}[next], func(t *testing.T) { p := &scriptedTransport{chunks: [][]byte{ vendorACK, apReply(status(0x33)), // current card at encoder; encoding succeeds vendorACK, // FC0 vendorACK, apReply(status(0x37)), // best effort must defer vendorACK, apReply(status(0x33)), // previous encoder sensor must not succeed vendorACK, apReply(status(0x34)), vendorACK, // FC7 vendorACK, apReply(status(0x33)), }} c, now := deliveryWorkerClient(t, p) if _, err := c.PrepareCurrentCard(context.Background()); err != nil { t.Fatal(err) } if _, err := c.DeliverCurrentCard(context.Background()); err != nil { t.Fatal(err) } // A cached clear result must not authorize FC7. c.mu.Lock() c.lastStatus = status(0x30) c.lastStatusT = time.Now() c.statusTTL = time.Hour c.mu.Unlock() if err := c.BeginPrepareNextCard(context.Background()); err != nil { t.Fatal(err) } if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP", "FC0", "AP", "AP", "AP", "FC7"}) { t.Fatalf("prestaging commands=%v, want clearance APs then one FC7", got) } prepare := c.PrepareCurrentCard if next { prepare = c.PrepareNextCard } if _, err := prepare(context.Background()); err != nil { t.Fatal(err) } want := []string{"AP", "FC0", "AP", "AP", "AP", "FC7", "AP"} if got := wireCommands(p); !reflect.DeepEqual(got, want) { t.Errorf("incident commands=%v, want %v", got, want) } if elapsed := now.Sub(time.Unix(0, 0)); elapsed != 2*time.Second { t.Errorf("clearance waits=%s, want 2s", elapsed) } }) } } func TestDeliveryClearWaitsTwoSeconds(t *testing.T) { for _, position := range []byte{0x30, 0x34} { p := &scriptedTransport{chunks: [][]byte{vendorACK, vendorACK, apReply(status(position)), vendorACK, apReply(status(position)), vendorACK, apReply(status(position)), vendorACK, vendorACK, apReply(status(0x37))}} c, now := deliveryWorkerClient(t, p) if _, err := c.DeliverCurrentCard(context.Background()); err != nil { t.Fatal(err) } if _, err := c.PrepareCurrentCard(context.Background()); err != nil { t.Fatal(err) } want := []string{"FC0", "AP", "AP", "AP", "FC7", "AP"} if got := wireCommands(p); !reflect.DeepEqual(got, want) { t.Errorf("clearance %X commands=%v, want %v", position, got, want) } if elapsed := now.Sub(time.Unix(0, 0)); elapsed != deliveryMinimumWait { t.Errorf("clearance %X elapsed=%s, want 2s", position, elapsed) } } } func TestDeliveryPendingIsNotPersisted(t *testing.T) { c := NewClient(nil, 1) defer c.Close() if c.deliveryPending { t.Error("new Client has pending delivery") } p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x37))}} fresh, _ := deliveryWorkerClient(t, p) if _, err := fresh.PrepareCurrentCard(context.Background()); err != nil { t.Errorf("fresh client encoder status: %v", err) } if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP"}) { t.Errorf("restart commands=%v, want AP only", got) } } func TestDeliveryCancellationAfterACKDoesNotSetPending(t *testing.T) { transportAddress(t) ctx, cancel := context.WithCancel(context.Background()) defer cancel() p := &scriptedTransport{chunks: [][]byte{vendorACK}, afterRead: cancel} c := &Client{port: p} r := workerRequest(c, ctx, cmdOutOfMouth) if !errors.Is(r.err, context.Canceled) || c.deliveryPending || len(p.writes) != 1 { t.Errorf("cancel after ACK: err=%v pending=%v writes=%d, want canceled, false, 1", r.err, c.deliveryPending, len(p.writes)) } } func TestWorkerSerializesClearanceAndFC7(t *testing.T) { p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x30)), vendorACK, vendorACK, vendorACK, apReply(status(0x37))}} c, _ := deliveryWorkerClient(t, p) // Set initial state before any request can reach the worker. c.deliveryPending = true c.deliveryStarted = c.now().Add(-deliveryMinimumWait) reading := make(chan struct{}) resume := make(chan struct{}) p.afterWrite = func() { if len(p.writes) == 1 { close(reading) <-resume } } ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() fc7 := make(chan error, 1) go func() { fc7 <- c.BeginPrepareNextCard(ctx) }() select { case <-reading: case <-ctx.Done(): t.Fatal("worker did not begin AP") } fc0 := make(chan cmdResp, 1) // This FC0 is queued while AP is in progress. It must not slip between // the clearance observation and FC7 dispatch. c.reqCh <- cmdReq{typ: cmdOutOfMouth, ctx: ctx, respCh: fc0} close(resume) if err := <-fc7; err != nil { t.Fatal(err) } if r := <-fc0; r.err != nil { t.Fatal(r.err) } want := []string{"AP", "FC7", "FC0"} if got := wireCommands(p); !reflect.DeepEqual(got, want) { t.Errorf("serialized commands=%v, want %v", got, want) } } func TestDeliveryCancellationAfterFreshAPKeepsPending(t *testing.T) { transportAddress(t) ctx, cancel := context.WithCancel(context.Background()) defer cancel() p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x34))}} p.afterRead = func() { if len(p.chunks) == 0 { cancel() } } c := &Client{port: p, deliveryPending: true, deliveryStarted: time.Now().Add(-time.Minute)} r := workerRequest(c, ctx, cmdToEncoder) if !errors.Is(r.err, context.Canceled) || !c.deliveryPending { t.Errorf("AP cancellation err=%v pending=%t, want canceled and pending", r.err, c.deliveryPending) } if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP"}) { t.Errorf("cancelled clearance commands=%v, want AP only", got) } } func TestPreparationWireFailuresNeverShake(t *testing.T) { badBCC := apReply(status(0x30)) badBCC[len(badBCC)-1] ^= 1 for _, frame := range [][]byte{nil, apReply(status(0x30))[:8], badBCC, apReply(status(0x40))} { p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x30)), vendorACK, vendorACK}} if frame != nil { p.chunks = append(p.chunks, frame) } c, _ := deliveryWorkerClient(t, p) if _, err := c.PrepareCurrentCard(context.Background()); err != nil { t.Errorf("frame % X: err=%v, want encoder opportunity", frame, err) } want := []string{"AP", "FC7", "AP", "AP", "AP", "AP", "AP"} if got := wireCommands(p); !reflect.DeepEqual(got, want) { t.Errorf("frame % X commands=%v, want %v without RS or resends", frame, got, want) } } } func TestBeginPrepareNextCardWaitsForClearanceWithoutReadinessPolling(t *testing.T) { for _, position := range []byte{0x30, 0x34} { p := &scriptedTransport{chunks: [][]byte{vendorACK, vendorACK, apReply(status(position)), vendorACK, apReply(status(position)), vendorACK, apReply(status(position)), vendorACK}} c, now := deliveryWorkerClient(t, p) if _, err := c.DeliverCurrentCard(context.Background()); err != nil { t.Fatal(err) } if err := c.BeginPrepareNextCard(context.Background()); err != nil { t.Fatal(err) } want := []string{"FC0", "AP", "AP", "AP", "FC7"} if got := wireCommands(p); !reflect.DeepEqual(got, want) { t.Errorf("BeginPrepareNextCard(%X) commands = %v, want %v", position, got, want) } if elapsed := now.Sub(time.Unix(0, 0)); elapsed != deliveryMinimumWait { t.Errorf("BeginPrepareNextCard(%X) waited %s, want %s", position, elapsed, deliveryMinimumWait) } } }