diff --git a/cmd/hardlink/main.go b/cmd/hardlink/main.go index fe6cfbf..131253e 100644 --- a/cmd/hardlink/main.go +++ b/cmd/hardlink/main.go @@ -33,7 +33,7 @@ import ( ) const ( - buildVersion = "v2.1.1" + buildVersion = "v2.1.2" serviceName = "hardlink" pollingFrequency = 8 * time.Second ) diff --git a/internal/dispenser/dispenser.go b/internal/dispenser/dispenser.go index ec56c13..90673ff 100644 --- a/internal/dispenser/dispenser.go +++ b/internal/dispenser/dispenser.go @@ -295,11 +295,60 @@ func sendAndReadACK(ctx context.Context, port serialTransport, packet []byte, pr if err := waitForSequence(ctx, processingDelay); err != nil { return err } - response := make([]byte, 3) - if err := readExact(ctx, port, response); err != nil { - return fmt.Errorf("read ACK: %w", err) + return readACK(ctx, port) +} + +// readACK scans sequentially without reading beyond a complete candidate. +// Wrong-address candidates are consumed in full; no bytes persist across calls. +func readACK(ctx context.Context, port serialTransport) (err error) { + const scanLimit = 64 + seen := make([]byte, 0, scanLimit) + defer func() { + if err != nil { + log.Warnf("dispenser ACK acquisition failed: address=% X seen=% X error=%v", Address, seen, err) + } + }() + var candidate [3]byte + used := 0 + for len(seen) < scanLimit { + if err := ctx.Err(); err != nil { + return err + } + need := 1 + if used > 0 { + need = len(candidate) - used + } + need = min(need, scanLimit-len(seen)) + chunk := candidate[used : used+need] + n, readErr := port.Read(chunk) + seen = append(seen, chunk[:n]...) + log.Debugf("dispenser ACK RX: address=% X n=%d bytes=% X error=%v", Address, n, chunk[:n], readErr) + if err := ctx.Err(); err != nil { + return err + } + used += n + if used == 1 && candidate[0] != ACK && candidate[0] != NAK { + log.Debugf("dispenser ACK ignored leading byte: % X", candidate[:1]) + used = 0 + } else if used == len(candidate) { + if len(Address) >= 2 && candidate[1] == Address[0] && candidate[2] == Address[1] { + log.Debugf("dispenser ACK/NAK accepted: address=% X token=% X", Address, candidate[:]) + if candidate[0] == ACK && len(seen) > len(candidate) { + log.Warnf("dispenser ACK resynchronized: address=% X ignored=% X accepted=% X", Address, seen[:len(seen)-len(candidate)], candidate[:]) + } + return checkACK(candidate[:]) + } + log.Debugf("dispenser ACK ignored wrong-address candidate: address=% X candidate=% X", Address, candidate[:]) + used = 0 + } + if readErr != nil { + return fmt.Errorf("read ACK: %w", readErr) + } + if n == 0 { + return fmt.Errorf("read ACK: %w", io.ErrNoProgress) + } } - return checkACK(response) + return fmt.Errorf("read ACK: scan limit of %d bytes exhausted", scanLimit) } // queryStatus accepts only the fixed RF/AP payload sizes, before reading a body. diff --git a/internal/dispenser/transport_test.go b/internal/dispenser/transport_test.go index 69ee972..6598094 100644 --- a/internal/dispenser/transport_test.go +++ b/internal/dispenser/transport_test.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "io" + "os" "testing" "time" ) @@ -24,6 +25,146 @@ var vendorFrames = []struct { {"RS", commandRS, []byte{2, 0x30, 0x30, 0, 2, 0x52, 0x53, 3, 2}}, } +func TestACKScannerFragmentation(t *testing.T) { + transportAddress(t) + for _, address := range []string{"00", "15"} { + Address = []byte(address) + wrong := []byte("15") + if address == "15" { + wrong = []byte("00") + } + for _, prefix := range [][]byte{nil, {0x30}, {0x10}, {address[1]}, {ACK, wrong[0], wrong[1], 0x30}, {NAK, wrong[0], wrong[1]}} { + for _, token := range []byte{ACK, NAK} { + wire := append(append([]byte(nil), prefix...), token, Address[0], Address[1]) + for mask := 0; mask < 1<<(len(wire)-1); mask++ { + t.Run(fmt.Sprintf("%s/%X/%d", address, wire, mask), func(t *testing.T) { + p := &scriptedTransport{} + start := 0 + for i := 1; i < len(wire); i++ { + if mask&(1<<(i-1)) != 0 { + p.chunks = append(p.chunks, wire[start:i]) + start = i + } + } + p.chunks = append(p.chunks, wire[start:]) + err := dispatchCommand(context.Background(), p, commandFC7, 0) + if (err != nil) != (token == NAK) { + t.Fatalf("dispatchCommand(% X)=%v", wire, err) + } + want := [][]byte{createPacket(Address, commandFC7)} + if token == ACK { + want = append(want, append([]byte{ENQ}, Address...)) + } + if len(p.writes) != len(want) { + t.Fatalf("writes=% X want=% X", p.writes, want) + } + for i := range want { + if !bytes.Equal(p.writes[i], want[i]) { + t.Errorf("write=% X want=% X", p.writes[i], want[i]) + } + } + if len(p.chunks) != 0 { + t.Errorf("unread token bytes=% X", p.chunks) + } + }) + } + } + } + } +} + +func TestACKScannerBoundsAndFollowingTransaction(t *testing.T) { + transportAddress(t) + for _, address := range []string{"00", "15"} { + Address = []byte(address) + ack := append([]byte{ACK}, Address...) + for _, leading := range []int{61, 62, 64} { + wire := append(bytes.Repeat([]byte{0xff}, leading), ack...) + p := &scriptedTransport{chunks: [][]byte{wire}} + err := dispatchCommand(context.Background(), p, commandRS, 0) + if (err == nil) != (leading == 61) { + t.Errorf("leading=%d error=%v", leading, err) + } + if leading != 61 && (len(p.writes) != 1 || len(bytes.Join(p.chunks, nil)) != len(wire)-64) { + t.Fatalf("scan exceeded bound: writes=% X remaining=% X", p.writes, p.chunks) + } + } + wire := append(append(append([]byte{0x10}, ack...), ack...), 0xfe) + p := &scriptedTransport{chunks: [][]byte{wire}} + for i := 0; i < 2; i++ { + if err := dispatchCommand(context.Background(), p, commandFC7, 0); err != nil { + t.Fatal(err) + } + } + if len(p.writes) != 4 || !bytes.Equal(bytes.Join(p.chunks, nil), []byte{0xfe}) { + t.Fatalf("transaction boundary lost: writes=% X remaining=% X", p.writes, p.chunks) + } + } +} + +type ackReadErrorTransport struct { + *scriptedTransport + reads, failAt int + err error +} + +func (p *ackReadErrorTransport) Read(b []byte) (int, error) { + p.reads++ + n, err := p.scriptedTransport.Read(b) + if p.reads == p.failAt { + return n, p.err + } + return n, err +} + +func TestACKScannerReadErrors(t *testing.T) { + transportAddress(t) + for _, address := range []string{"00", "15"} { + Address = []byte(address) + for _, token := range []byte{ACK, NAK} { + for _, readErr := range []error{io.EOF, os.ErrDeadlineExceeded, io.ErrClosedPipe} { + for _, failAt := range []int{1, 2} { + p := &ackReadErrorTransport{scriptedTransport: &scriptedTransport{chunks: [][]byte{{token, Address[0], Address[1]}}}, failAt: failAt, err: readErr} + err := dispatchCommand(context.Background(), p, commandFC7, 0) + if failAt == 1 && !errors.Is(err, readErr) { + t.Fatalf("incomplete token error=%v want=%v", err, readErr) + } + if failAt == 2 && ((err == nil) != (token == ACK) || errors.Is(err, readErr)) { + t.Fatalf("complete token %02X with read error returned %v", token, err) + } + wantWrites := 1 + if failAt == 2 && token == ACK { + wantWrites = 2 + } + if len(p.writes) != wantWrites || p.reads != failAt { + t.Fatalf("writes=%d reads=%d want=%d/%d", len(p.writes), p.reads, wantWrites, failAt) + } + } + } + } + for _, wire := range [][]byte{nil, {0x30, ACK, Address[0]}, {0xff, 0xfe}, {ACK, Address[0]}} { + p := &scriptedTransport{chunks: [][]byte{wire}} + if err := dispatchCommand(context.Background(), p, commandFC7, 0); err == nil || len(p.writes) != 1 { + t.Fatalf("incomplete/garbage % X: err=%v writes=% X", wire, err, p.writes) + } + } + for _, cancelAt := range []int{1, 2} { + ctx, cancel := context.WithCancel(context.Background()) + p := &ackReadErrorTransport{scriptedTransport: &scriptedTransport{chunks: [][]byte{{ACK, Address[0], Address[1]}}}, failAt: 2, err: io.EOF} + p.afterRead = func() { + if p.reads == cancelAt { + cancel() + } + } + err := dispatchCommand(ctx, p, commandFC7, 0) + cancel() + if !errors.Is(err, context.Canceled) || len(p.writes) != 1 || p.reads != cancelAt { + t.Fatalf("cancel read %d: err=%v writes=%d reads=%d", cancelAt, err, len(p.writes), p.reads) + } + } + } +} + func TestVendorOutboundFrames(t *testing.T) { for _, tc := range vendorFrames { t.Run(tc.name, func(t *testing.T) { @@ -216,7 +357,11 @@ func TestTransportCancellation(t *testing.T) { case "processing wait": p.afterWrite = cancel case "after ACK": - p.afterRead = cancel + p.afterRead = func() { + if len(p.chunks) == 1 { + cancel() + } + } case "after ENQ": wantWrites = 2 p.afterWrite = func() { @@ -229,7 +374,7 @@ func TestTransportCancellation(t *testing.T) { reads := 0 p.afterRead = func() { reads++ - if reads == 2 { + if reads == 3 { cancel() } } diff --git a/release notes.md b/release notes.md index 3485805..8d662ba 100644 --- a/release notes.md +++ b/release notes.md @@ -2,6 +2,9 @@ builtVersion is a const in main.go +#### v2.1.2 - 30 September 2026 +fix(dispenser): scan ACK responses without assuming 3-byte reads + #### v2.1.1 - 28 September 2026 feat(dispenser): add worker-owned idle card prestaging