hardlink/internal/dispenser/delivery_test.go

364 lines
12 KiB
Go

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)
}
}
}