diff --git a/internal/dispenser/delivery_test.go b/internal/dispenser/delivery_test.go new file mode 100644 index 0000000..80288b6 --- /dev/null +++ b/internal/dispenser/delivery_test.go @@ -0,0 +1,391 @@ +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}, + {"mouth", status(0x34), false}, + {"unknown diagnostic", []byte{0x30, 0x30, 0xFF, 0x30}, false}, + {"jam", []byte{0x30, 0x30, 0x32, 0x30}, false}, + {"overlap", []byte{0x30, 0x30, 0x34, 0x30}, false}, + {"rejection", []byte{0x36, 0x30, 0x30, 0x30}, false}, + {"empty", status(0x38), false}, + {"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(), lastStatus: status(0x30), lastStatusT: time.Now(), statusTTL: time.Hour} + r := workerRequest(c, context.Background(), cmdToEncoder) + if (r.err == nil) != tc.clear || 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 tc.clear { + want = append(want, "FC7") + } + 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(), cmdToEncoder) + 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, apReply(status(0x37)), + vendorACK, apReply(status(0x30)), // clearance + 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()); !errors.Is(err, errDeliveryPending) { + t.Fatalf("BeginPrepareNextCard=%v, want deferred", err) + } + if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP", "FC0", "AP"}) { + t.Fatalf("before clearance commands=%v", 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", "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 != 3*time.Second { + t.Errorf("clearance waits=%s, want 3s", elapsed) + } + }) + } +} + +func TestDeliveryClearImmediately(t *testing.T) { + p := &scriptedTransport{chunks: [][]byte{vendorACK, vendorACK, apReply(status(0x30)), vendorACK, vendorACK, apReply(status(0x33))}} + 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) + } + if got := wireCommands(p); !reflect.DeepEqual(got, []string{"FC0", "AP", "FC7", "AP"}) { + t.Errorf("commands=%v, want FC0 AP FC7 AP", got) + } + if !now.Equal(time.Unix(0, 0)) { + t.Errorf("immediate clearance waited until %v", now) + } +} + +func TestDeliveryClearanceWaitStops(t *testing.T) { + for _, tc := range []struct { + name string + response cmdResp + cancel bool + timeout bool + }{ + {name: "cancel", response: cmdResp{status: status(0x37), deliveryPending: true}, cancel: true}, + {name: "timeout", response: cmdResp{status: status(0x37), deliveryPending: true}, timeout: true}, + {name: "transport", response: cmdResp{err: io.ErrUnexpectedEOF}}, + {name: "jam", response: cmdResp{status: []byte{0x30, 0x30, 0x32, 0x30}, deliveryPending: true}}, + {name: "empty", response: cmdResp{status: status(0x38), deliveryPending: true}}, + } { + t.Run(tc.name, func(t *testing.T) { + responses := make([]cmdResp, 20) + for i := range responses { + responses[i] = tc.response + } + c, d := newSequenceTestClient(t, responses...) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + if tc.cancel { + c.sequenceTiming.wait = func(context.Context, time.Duration) error { cancel(); return ctx.Err() } + } + started := c.sequenceTiming.now() + _, err := c.PrepareCurrentCard(ctx) + if err == nil { + t.Fatal("preparation succeeded while clearance unavailable") + } + if tc.cancel && !errors.Is(err, context.Canceled) { + t.Errorf("error=%v, want canceled", err) + } + if tc.timeout { + if !errors.Is(err, context.DeadlineExceeded) || c.sequenceTiming.now().Sub(started) != sequenceTimeout { + t.Errorf("timeout err=%v elapsed=%s", err, c.sequenceTiming.now().Sub(started)) + } + } + for _, cmd := range d.commands { + if cmd != cmdStatus { + t.Errorf("clearance dispatched command %v, want only AP", cmd) + } + } + }) + } +} + +func TestDeliveryClearanceDoesNotConsumeRecoveryWindow(t *testing.T) { + responses := []cmdResp{{status: status(0x37), deliveryPending: true}, {status: status(0x33), deliveryPending: true}, {status: status(0x30)}} + responses = append(responses, prepareFailureResponses(0x32, 7)...) + responses = append(responses, cmdResp{status: status(0x30)}, cmdResp{status: status(0x33)}) + c, d := newSequenceTestClient(t, responses...) + if _, err := c.PrepareCurrentCard(context.Background()); err != nil { + t.Fatal(err) + } + var initial, reset time.Time + for i, cmd := range d.commands { + if cmd == cmdToEncoder && initial.IsZero() { + initial = d.commandTimes[i] + } + if cmd == cmdReset { + reset = d.commandTimes[i] + } + } + if reset.Sub(initial) != sequenceRetryAfter || commandCount(d.commands, cmdReset) != 1 || commandCount(d.commands, cmdToEncoder) != 2 { + t.Errorf("recovery commands=%v reset delay=%s, want two FC7 and one RS at 6s", d.commands, reset.Sub(initial)) + } + if got := c.sequenceTiming.now().Sub(reset); got != sequenceResetWait+sequencePollInterval { + t.Errorf("post-reset elapsed=%s, want settle plus poll=3s", got) + } +} + +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 TestDeliveryClearanceRetainsFullPreparationTimeout(t *testing.T) { + responses := make([]cmdResp, 8) + for i := range responses { + responses[i] = cmdResp{status: status(0x37), deliveryPending: true} + } + responses = append(responses, cmdResp{status: status(0x30)}) + responses = append(responses, prepareFailureResponses(0x32, 30)...) + c, d := newSequenceTestClient(t, responses...) + started := c.sequenceTiming.now() + if _, err := c.PrepareCurrentCard(context.Background()); err == nil { + t.Fatal("persistent failure succeeded") + } + if got := c.sequenceTiming.now().Sub(started); got != 8*time.Second+sequenceTimeout { + t.Errorf("total duration=%s, want clearance 8s + preparation 16s", got) + } + if commandCount(d.commands, cmdReset) != 1 || commandCount(d.commands, cmdToEncoder) != 2 { + t.Errorf("commands=%v, want one RS and two FC7", d.commands) + } +} + +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 = time.Now() + 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) + } + if err := c.ToEncoder(ctx); !errors.Is(err, errDeliveryPending) { + t.Errorf("FC7 after queued FC0=%v, want deferred", err) + } + want := []string{"AP", "FC7", "FC0", "AP"} + if got := wireCommands(p); !reflect.DeepEqual(got, want) { + t.Errorf("serialized commands=%v, want %v", got, want) + } +} diff --git a/internal/dispenser/dispenser.go b/internal/dispenser/dispenser.go index 8d29c44..9a5c3c3 100644 --- a/internal/dispenser/dispenser.go +++ b/internal/dispenser/dispenser.go @@ -254,17 +254,23 @@ type serialTransport interface { } func writePacket(ctx context.Context, port serialTransport, packet []byte) error { + _, err := writePacketAttempt(ctx, port, packet) + return err +} + +// writePacketAttempt distinguishes cancellation before Write from an ambiguous write. +func writePacketAttempt(ctx context.Context, port serialTransport, packet []byte) (bool, error) { if err := ctx.Err(); err != nil { - return err + return false, err } n, err := port.Write(packet) if err != nil { - return fmt.Errorf("write dispenser packet: %w", err) + return true, fmt.Errorf("write dispenser packet: %w", err) } if n != len(packet) { - return fmt.Errorf("write dispenser packet (%d/%d bytes): %w", n, len(packet), io.ErrShortWrite) + return true, fmt.Errorf("write dispenser packet (%d/%d bytes): %w", n, len(packet), io.ErrShortWrite) } - return ctx.Err() + return true, ctx.Err() } func readExact(ctx context.Context, port serialTransport, data []byte) error { @@ -410,7 +416,10 @@ func resetDispenser(ctx context.Context, port serialTransport) error { return dispatchCommand(ctx, port, commandRS, delay) } -func cardOutOfMouth(ctx context.Context, port serialTransport) error { +func cardOutOfMouth(ctx context.Context, port serialTransport) (bool, error) { log.Println("Send card to out mouth position") - return dispatchCommand(ctx, port, commandFC0, delay) + if err := sendAndReadACK(ctx, port, createPacket(Address, commandFC0), delay); err != nil { + return false, err + } + return writePacketAttempt(ctx, port, append([]byte{ENQ}, Address...)) } diff --git a/internal/dispenser/dispenserclient.go b/internal/dispenser/dispenserclient.go index fa8ef26..117e74e 100644 --- a/internal/dispenser/dispenserclient.go +++ b/internal/dispenser/dispenserclient.go @@ -3,6 +3,7 @@ package dispenser import ( "context" + "errors" "fmt" "sync" "time" @@ -26,9 +27,12 @@ type cmdReq struct { respCh chan cmdResp } +var errDeliveryPending = errors.New("next-card preparation deferred: previous delivery is not clear") + type cmdResp struct { - status []byte - err error + deliveryPending bool + status []byte + err error } type sequenceTiming struct { @@ -49,6 +53,10 @@ type Client struct { reqCh chan cmdReq done chan struct{} + // Owned exclusively by the serial worker. + deliveryPending bool + deliveryStarted time.Time + sequenceTiming sequenceTiming // status cache @@ -189,21 +197,26 @@ func (c *Client) handle(req cmdReq) { switch req.typ { case cmdStatus: - st, err := checkDispenserStatus(req.ctx, c.port) - if err == nil && len(st) == 4 { - c.mu.Lock() - c.lastStatus = append([]byte(nil), st...) - c.lastStatusT = time.Now() - c.mu.Unlock() - - // publish stock/cardwell - c.setStock(st) - } - req.respCh <- cmdResp{status: st, err: err} + st, err := c.readWorkerStatus(req.ctx) + req.respCh <- cmdResp{status: st, err: err, deliveryPending: c.deliveryPending} case cmdToEncoder: + if c.deliveryPending { + // Stay inside port ownership: never enqueue a request from the worker. + st, err := c.readWorkerStatus(req.ctx) + if err == nil { + _, err = deliveryClearance(st) + } + if err == nil && c.deliveryPending { + err = errDeliveryPending + log.Debugf("next-card preparation deferred; raw status: % X", st) + } + if err != nil { + req.respCh <- cmdResp{err: err} + return + } + } err := cardToEncoderPosition(req.ctx, c.port) - // A movement command makes any previously cached position unreliable. c.invalidateStatusCache() req.respCh <- cmdResp{err: err} @@ -213,7 +226,12 @@ func (c *Client) handle(req cmdReq) { req.respCh <- cmdResp{err: err} case cmdOutOfMouth: - err := cardOutOfMouth(req.ctx, c.port) + attempted, err := cardOutOfMouth(req.ctx, c.port) + if attempted { + c.deliveryPending = true + c.deliveryStarted = time.Now() + log.Info("delivery ENQ attempted; awaiting mechanical clearance") + } // A movement command makes any previously cached position unreliable. c.invalidateStatusCache() req.respCh <- cmdResp{err: err} @@ -223,21 +241,68 @@ func (c *Client) handle(req cmdReq) { } } +// deliveryClearance deliberately does not apply encoder-success precedence. +func deliveryClearance(status []byte) (bool, error) { + if err := validateDispenserStatusData(status); err != nil { + return false, err + } + if isCardWellEmpty(status) { + return false, ErrCardWellEmpty + } + if isPreparationMoving(status) { + return false, nil + } + if err := dispenserStatusError(status); err != nil { + return false, err + } + return status[0] == 0x30 && status[1] == 0x30 && status[3] == 0x30, nil +} + +// readWorkerStatus is called only by the serial worker and always reads fresh AP. +func (c *Client) readWorkerStatus(ctx context.Context) ([]byte, error) { + st, err := checkDispenserStatus(ctx, c.port) + if err != nil { + return st, err + } + if c.deliveryPending { + clear, clearanceErr := deliveryClearance(st) + if clearanceErr != nil { + log.Warnf("delivery clearance failed: %v; %s raw status: % X", clearanceErr, statusDescription(st), st) + } else if clear { + c.deliveryPending = false + log.Infof("previous delivery cleared after %s; FC7 permitted", time.Since(c.deliveryStarted)) + } + } + if len(st) == 4 { + c.mu.Lock() + c.lastStatus = append([]byte(nil), st...) + c.lastStatusT = time.Now() + c.mu.Unlock() + c.setStock(st) + } + return st, nil +} + func (c *Client) do(ctx context.Context, typ cmdType) ([]byte, error) { + r := c.doResponse(ctx, typ) + return r.status, r.err +} + +func (c *Client) doResponse(ctx context.Context, typ cmdType) cmdResp { rch := make(chan cmdResp, 1) req := cmdReq{typ: typ, ctx: ctx, respCh: rch} select { case c.reqCh <- req: case <-ctx.Done(): - return nil, ctx.Err() + return cmdResp{err: ctx.Err()} } select { case r := <-rch: - return r.status, r.err + return r case <-ctx.Done(): - return nil, ctx.Err() + return cmdResp{err: ctx.Err()} } } @@ -322,7 +387,11 @@ func (c *Client) DispenserPrepare(ctx context.Context) (string, error) { } func (c *Client) readSequenceStatus(ctx context.Context, operation string) ([]byte, string, error) { - status, err := c.do(ctx, cmdStatus) + response := c.doResponse(ctx, cmdStatus) + status, err := response.status, response.err + if err == nil && response.deliveryPending { + err = errDeliveryPending + } if err != nil { if ctxErr := ctx.Err(); ctxErr != nil { return nil, "", fmt.Errorf("[%s] read status: %w", operation, ctxErr) @@ -513,8 +582,60 @@ func (c *Client) pollForEncoderPosition(ctx context.Context, operation string) ( } } +// waitForDeliveryClearance reuses the first clear AP as the initial preparation sample. +// Its bounded window is separate from recovery, which starts after initial FC7. +func (c *Client) waitForDeliveryClearance(ctx context.Context, operation string) (status []byte, stock string, resultErr error) { + ctx, cancel := context.WithTimeout(ctx, sequenceTimeout) + defer cancel() + deadline := c.sequenceTiming.now().Add(sequenceTimeout) + var lastStatus []byte + defer func() { + if resultErr != nil { + log.Warnf("[%s] delivery clearance stopped: %v; %s raw status: % X", operation, resultErr, statusDescription(lastStatus), lastStatus) + } + }() + for { + if err := ctx.Err(); err != nil { + return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, err) + } + if !c.sequenceTiming.now().Before(deadline) { + return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, context.DeadlineExceeded) + } + r := c.doResponse(ctx, cmdStatus) + if len(r.status) > 0 { + lastStatus = r.status + } + if r.err != nil { + return nil, stock, fmt.Errorf("[%s] delivery clearance status: %w", operation, r.err) + } + if len(r.status) == 4 { + stock = stockTake(r.status) + c.setStock(r.status) + } + if err := ctx.Err(); err != nil { + return nil, stock, err + } + if !c.sequenceTiming.now().Before(deadline) { + return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, context.DeadlineExceeded) + } + if !r.deliveryPending { + return r.status, stock, nil + } + if _, err := deliveryClearance(r.status); err != nil { + return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, err) + } + wait := sequencePollInterval + if remaining := deadline.Sub(c.sequenceTiming.now()); remaining < wait { + wait = remaining + } + if err := c.sequenceTiming.wait(ctx, wait); err != nil { + return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, err) + } + } +} + func (c *Client) prepareCardAtEncoder(ctx context.Context, operation string) (string, error) { - status, stockStatus, err := c.readSequenceStatus(ctx, operation) + status, stockStatus, err := c.waitForDeliveryClearance(ctx, operation) if err != nil { return stockStatus, err } diff --git a/internal/handlers/doorcard_handlers_test.go b/internal/handlers/doorcard_handlers_test.go index 2b399b0..ff217a9 100644 --- a/internal/handlers/doorcard_handlers_test.go +++ b/internal/handlers/doorcard_handlers_test.go @@ -168,6 +168,16 @@ func TestIssueDoorCardPhysicalOutcomeContract(t *testing.T) { wantCalls: []string{"prepare current", "deliver current", "begin prepare next"}, wantLockSequence: 1, }, + { + name: "delivery clearance deferral does not undo successful issuance", + dispenser: fakeDoorCardDispenser{ + beginNextErr: errors.New("next-card preparation deferred: previous delivery is not clear"), + }, + wantHTTP: http.StatusOK, + wantMessage: "Card issued successfully", + wantCalls: []string{"prepare current", "deliver current", "begin prepare next"}, + wantLockSequence: 1, + }, { name: "encoding failure with safe recovery remains retryable", lockErr: encodingErr,