From aef1bbbe576c3be2d35114b42894529958b6bd45 Mon Sep 17 00:00:00 2001 From: yurii Date: Tue, 29 Sep 2026 18:58:26 +0100 Subject: [PATCH] feat(dispenser): add worker-owned idle card prestaging --- cmd/hardlink/main.go | 3 +- internal/dispenser/K720_DISPENSER_CONTRACT.md | 335 ++++++++++++++++ internal/dispenser/dispenserclient.go | 100 +++-- internal/dispenser/maintenance.go | 180 +++++++++ internal/dispenser/maintenance_test.go | 362 ++++++++++++++++++ internal/handlers/doorcard_handlers_test.go | 44 ++- internal/handlers/handlers.go | 4 + internal/handlers/testhandlers.go | 3 + release notes.md | 3 + 9 files changed, 1005 insertions(+), 29 deletions(-) create mode 100644 internal/dispenser/K720_DISPENSER_CONTRACT.md create mode 100644 internal/dispenser/maintenance.go create mode 100644 internal/dispenser/maintenance_test.go diff --git a/cmd/hardlink/main.go b/cmd/hardlink/main.go index 8c68eef..fe6cfbf 100644 --- a/cmd/hardlink/main.go +++ b/cmd/hardlink/main.go @@ -33,7 +33,7 @@ import ( ) const ( - buildVersion = "v2.1.0" + buildVersion = "v2.1.1" serviceName = "hardlink" pollingFrequency = 8 * time.Second ) @@ -227,6 +227,7 @@ func main() { // Start polling for dispenser status every 10 seconds disp.StartPolling(pollingFrequency) + disp.StartMaintenance() } mux := http.NewServeMux() diff --git a/internal/dispenser/K720_DISPENSER_CONTRACT.md b/internal/dispenser/K720_DISPENSER_CONTRACT.md new file mode 100644 index 0000000..ac1f60c --- /dev/null +++ b/internal/dispenser/K720_DISPENSER_CONTRACT.md @@ -0,0 +1,335 @@ +# K720 dispenser issuance contract + +Hardlink treats the K720 as an unreliable mechanical peripheral with a strict transport +protocol and imperfect status telemetry. The implementation must keep transport +validation strict while making mechanical decisions from fresh, bounded observations. +Diagnostic bytes are useful evidence, but they must not override a valid physical +position or turn an uncertain observation into a fatal guest-facing error. + +This document records the dispenser invariants agreed for the clean implementation. +It is intended to be a design contract: future changes should preserve these rules +unless production evidence justifies changing them explicitly. + +## Ownership and concurrency + +- The dispenser worker is the single owner of serial commands and `deliveryPending`. + New code must not access the serial port directly from handlers or background + goroutines. +- Commands that can move a card are serialized through the worker. There must never + be concurrent AP/FC7/FC0/RS traffic from separate flows. +- `deliveryPending` is worker-owned state. It must not be duplicated or mutated by + HTTP handlers. +- A real `/issuedoorcard` operation always has priority over maintenance/prestage + work. +- Cancellation of a real caller must stop caller-owned work. Internal bounded + observation timeouts must remain distinguishable from caller cancellation. +- Do not create overlapping maintenance goroutines or timers per request. Any + background maintenance must have one clear owner and lifecycle. + +## Status and observation model + +AP transport parsing is strict. A response is usable only when its framing and +contents are valid, including ACK/address/header/length/ETX/BCC/type checks. + +Malformed, truncated, timed-out or otherwise invalid AP responses are unusable +observations. They are logged and discarded. They are never converted into a +fabricated physical state such as `0x30` or `0x38`. + +Every fresh valid AP observation reclassifies the current mechanical state. No +physical position is latched across later observations. + +The fourth AP status byte is the physical position used by the preparation logic: + +- Encoder confirmed: + `0x32, 0x33, 0x36, 0x37, 0x3A, 0x3B, 0x3E, 0x3F`. + A card may be handed to `LockSequence` immediately. +- Valid but uncertain: + `0x31, 0x34, 0x35, 0x39, 0x3C, 0x3D`. + Continue bounded fresh observation; if no stronger state appears, one encoding + opportunity may still be granted after the uncertainty window. +- No card on sensors: + exact `0x30`. + This may trigger the bounded mechanical recovery sequence. +- Card well empty: + exact `0x38`. + This is the only physical state that becomes `ErrCardWellEmpty`. + +The first three diagnostic bytes are advisory. Values such as `Preparing card fails`, +`Dispense card error`, `Card jammed`, or `Command cannot execute` may coexist with a +usable physical position. They may be logged, but they must not automatically block +encoding when the fourth byte confirms a usable card position. + +## Preparing the current card + +A `/issuedoorcard` request receives at most one `LockSequence` opportunity. + +Preparation uses bounded fresh observation: + +- Poll interval: 1 second. +- After a successful FC7, allow 3 seconds of observation before shake recovery is + eligible for a persistent exact `0x30`. +- For uncertain or unusable observations after a successful FC7, use one shared + 4-second uncertainty window. Switching between uncertain and unusable observations + does not restart that window. +- A later usable observation always takes precedence: + encoder confirmed -> handoff; + exact `0x38` -> empty; + exact `0x30` -> resume the remaining shake budget and reset uncertainty timing. +- Total preparation budget: 32 seconds. + +For persistent exact `0x30`, the request may perform at most three shake recoveries: + +`RS -> 2 second settle -> FC7` + +There is no additional plain FC7 retry. Including the initial dispatch, the maximum +per request is four FC7 commands and three RS commands. + +If all three shakes are exhausted while the latest usable physical state remains +exact `0x30`, preparation fails. If the internal preparation deadline is reached +while the latest state is uncertain or unusable after a successful FC7, the request +may hand off one encoding opportunity. Exact `0x38` remains empty. + +FC7 or RS dispatch failures are real transport/command failures. They are not +converted into status uncertainty and are not retried merely because their dispatch +failed. + +## Encoding and delivery + +`LockSequence` runs exactly once for a prepared card in a single `/issuedoorcard` +request. The dispenser layer does not implement a second encoding attempt on the same +physical card. + +After `LockSequence`, FC0 is dispatched exactly once to move the current card toward +the guest: + +- `LockSequence` success -> FC0 once -> successful issuance remains successful. +- `LockSequence` failure -> FC0 once -> return the original encoding error. +- FC0 errors are logged only. They do not replace the original encoding result and + they do not become `ErrCardWellEmpty`. + +A failed `LockSequence` does not start next-card prestaging. + +The design intentionally does not add retained-card state, per-card failure counters, +provider-specific encoder retries, special bad-card recovery, CP/capture recovery, or +multiple encoding attempts per physical card without production evidence requiring +them. + +## Delivery clearance + +After FC0, the worker marks the previous delivery as pending. A later FC7 must not be +sent until the previous delivery has been considered clear. + +Delivery clearance uses fresh AP observations and is owned by the worker. + +Current bounds: + +- Minimum clearance wait: 2 seconds. +- Clearance timeout: 6 seconds. +- Poll interval: 1 second. + +A valid physical position of `0x30` or `0x34` may clear `deliveryPending` after the +minimum wait unless the same fresh observation explicitly indicates active +preparing/dispensing/capturing movement. + +Exact `0x38` reports physical empty and does not silently clear the state. + +If the internal clearance timeout is reached while the caller context is still +valid, the worker may make the bounded assumption that delivery has cleared unless +the latest fresh usable observation establishes exact physical empty (`0x38`) or +explicit active movement. This fallback may therefore occur while the latest +physical position is `0x33`; `0x33` itself is not positive clearance evidence. +The fallback is an internal timeout policy, not a reclassification of the observed +position. + +If the latest fresh usable observation still explicitly indicates movement, that is +not treated as unknown. `deliveryPending` remains set and the clearance attempt +times out. + +Caller cancellation or deadline always takes precedence over the internal clearance +fallback. If the caller context expires, return the caller error and do not admit a +subsequent FC7 from that caller-owned operation. + +A later unusable observation supersedes older movement evidence; stale movement +information must not be latched indefinitely. + +## Next-card prestage and guest UX + +Preparing the next card is an optimization for the next guest, not part of the +business success of the current guest's issuance. + +After a successful `LockSequence` and FC0, Hardlink may make a short best-effort +attempt to prestage the next card: + +- wait for worker-owned delivery clearance; +- if clearance is obtained within the bounded prestage context, send exactly one FC7; +- do not run readiness polling, RS/shake recovery, or the full current-card + preparation flow; +- prestage failure is log-only and must not change the successful `/issuedoorcard` + result. + +Do not extend the current guest's screen by 15-20 seconds merely to guarantee that +the next card reaches the encoder. The guest-facing flow must remain bounded even +when the dispenser is slow to become ready for prestage. + +If the short prestage window expires, skip that prestage attempt. A later request or +maintenance cycle may prepare the next card. + +## HTTP contract + +Physical empty and operational failure are deliberately different outcomes. + +- HTTP 503 is reserved for `ErrCardWellEmpty`, derived only from an exact valid + physical `0x38`. +- All other dispenser preparation, transport, command and encoding failures return + HTTP 502. +- Normal request/protocol validation keeps its existing 400/405/415 behavior. +- Diagnostic text such as `Preparing card fails` must never by itself produce 503. + +Operafyne treats: + +- 502 as retryable while attempts remain; +- 503 as terminal `dispenser_failed`; +- a maximum of three total issue attempts: the initial attempt plus up to two user + retries. + +The dispenser layer must preserve this distinction. + +## Passive status and alerts + +Ordinary status queries remain strict and passive. They must not reuse tolerant +issuance semantics to fabricate a status or hide malformed transport. + +Passive polling captures the current foreground activity generation when the poll is +queued. The worker skips a passive poll if foreground activity is active when it is +dispatched or if the captured generation has become obsolete. Passive polling does +not reset the idle-maintenance clock. + +Status/diagnostic observations may be useful for logs and support alerts, but alerts +must not alter physical state classification or issuance success. + +Idle maintenance observations are intentionally quiet: they do not invoke the normal +stock-update callback or generate repeated support email such as `Preparing card +fails`. Maintenance failures are local diagnostic events only. + +## Idle prestage maintenance + +Idle prestaging recovers from a short post-FC0 prestage timeout without keeping the +current guest waiting. + +The implementation uses the existing dispenser worker as the single maintenance +owner: + +- one resettable maintenance timer is owned by the serial-worker loop; there is no + maintenance goroutine and no timer created per request; +- `/issuedoorcard` and `/testissuedoorcard` register foreground activity before + using the dispenser/encoder and release it when the handler finishes; +- foreground activity is reference-counted so overlapping requests suppress + maintenance until the last active request finishes; +- each foreground registration increments an activity generation and invalidates the + previous idle deadline; +- when the last foreground request finishes, the next maintenance attempt is + scheduled one minute later; +- enabling maintenance at startup schedules the first check one minute later when + no foreground request is active; +- a stale timer or stale generation cannot perform maintenance; +- after a maintenance attempt, the next check is scheduled one minute from completion + only if the same generation is still authoritative. A later foreground completion + therefore cannot have its newer deadline overwritten by an older maintenance + attempt. + +Foreground registration and maintenance transaction admission use the same activity +guard. The guard is used only for admission/state bookkeeping and is never held +across serial I/O, queue waits, sleeps, callbacks, or `LockSequence`. + +A maintenance attempt is deliberately weaker than foreground preparation: + +1. create a bounded 5-second maintenance context; +2. recheck maintenance admission before AP; +3. obtain one fresh AP using the existing strict transport parser; +4. classify the fresh status using the common physical classifier; +5. do nothing for encoder-present, exact `0x38`, unusable, explicit-movement, or + otherwise ineligible observations; +6. when `deliveryPending` is set, require fresh qualifying clearance evidence and the + existing minimum clearance wait before clearing it; +7. never use the foreground six-second assumed-clearance fallback; +8. recheck foreground count, activity generation, maintenance enabled state, and + maintenance context immediately before FC7 admission; +9. dispatch at most one FC7 and return immediately. + +If foreground activity begins while an already-admitted maintenance AP is in +progress, that AP may finish, but the second admission check prevents maintenance +from sending FC7 afterward. An FC7 that was already admitted and dispatched before +foreground registration cannot be recalled. + +Maintenance never performs: + +- RS or shake recovery; +- `LockSequence`; +- `PrepareCurrentCard`; +- `PrepareNextCard` or `BeginPrepareNextCard`; +- readiness polling after FC7; +- the 32-second foreground preparation flow; +- repeated FC7 attempts; +- the foreground assumed-clearance fallback. + +The 5-second maintenance context bounds new admissions and context-aware waits. An +already-admitted serial read may still finish according to the existing serial read +timeout, but no subsequent maintenance transaction is admitted after cancellation or +expiry. + +`StartMaintenance` and `StopMaintenance` are idempotent lifecycle operations. +`Client.Close` disables maintenance and stops the timer. +Stopping maintenance cancels the current maintenance context and waits only for an +already-admitted maintenance transaction to finish under the existing bounded serial +behavior. + +## Automated verification + +Tests should preserve the behavioural contract rather than only exercise individual +functions. + +At minimum, cover: + +- strict AP framing and malformed/truncated/timeout responses; +- tolerant issuance observations never fabricating `0x30` or `0x38`; +- every physical position class and transitions between classes; +- exact `0x38` as the only `ErrCardWellEmpty` path; +- persistent `0x30` using at most three `RS -> settle -> FC7` recoveries; +- no extra FC7 after the shake budget; +- uncertain/unusable shared timing and later usable-state precedence; +- caller cancellation versus internal preparation timeout; +- exactly one `LockSequence` opportunity per request; +- FC0 exactly once after encoding success or failure; +- FC0 failure remaining log-only; +- failed encoding never starting prestage; +- worker-owned `deliveryPending` and guarded FC7 dispatch; +- bounded delivery-clearance fallback and explicit-movement timeout; +- successful issuance remaining successful when next-card prestage fails; +- HTTP 503 only for exact physical empty and 502 for other issuance failures; +- Operafyne retry semantics remaining three total attempts. + +Idle maintenance verification additionally covers: + +- first attempt occurs one minute after the latest issuance finishes; +- a new issuance resets that idle interval; +- overlapping foreground requests suppress maintenance until the last request + finishes; +- both `/issuedoorcard` and `/testissuedoorcard` participate in foreground activity + registration; +- stale timer events and stale activity generations cannot perform maintenance; +- foreground registration before maintenance AP prevents AP admission; +- foreground registration during an admitted AP prevents the later FC7; +- passive polls are skipped while foreground activity is active or when their + captured generation is stale; +- maintenance performs one AP and at most one FC7, with no RS/shake/encoding; +- encoder-present, exact-empty, movement and unusable observations are no-ops; +- pending delivery requires fresh normal clearance plus the minimum wait and never + uses the foreground assumed-clearance fallback; +- maintenance errors do not produce guest-facing failures, stock callbacks or + repeated alert email; +- repeated maintenance start/stop calls are safe and do not accumulate work; +- cancellation, `StopMaintenance` and `Client.Close` prevent subsequent maintenance + commands after shutdown admission is revoked. + +For Hardlink verification, continue excluding the existing test that sends real +email when running the full suite. diff --git a/internal/dispenser/dispenserclient.go b/internal/dispenser/dispenserclient.go index a5f196f..9a72232 100644 --- a/internal/dispenser/dispenserclient.go +++ b/internal/dispenser/dispenserclient.go @@ -23,9 +23,11 @@ const ( ) type cmdReq struct { - typ cmdType - ctx context.Context - respCh chan cmdResp + passive bool + generation uint64 + typ cmdType + ctx context.Context + respCh chan cmdResp } type cmdResp struct { @@ -51,7 +53,10 @@ const ( ) type Client struct { - port serialTransport + activity activityGuard + activityWake chan struct{} + closeOnce sync.Once + port serialTransport reqCh chan cmdReq done chan struct{} @@ -80,9 +85,10 @@ func NewClient(port *serial.Port, queueSize int) *Client { queueSize = 16 } c := &Client{ - port: port, - reqCh: make(chan cmdReq, queueSize), - done: make(chan struct{}), + activityWake: make(chan struct{}, 1), + port: port, + reqCh: make(chan cmdReq, queueSize), + done: make(chan struct{}), sequenceTiming: sequenceTiming{ now: time.Now, wait: waitForSequence, @@ -106,12 +112,13 @@ func waitForSequence(ctx context.Context, duration time.Duration) error { } func (c *Client) Close() { - select { - case <-c.done: - return - default: + c.closeOnce.Do(func() { + c.activity.Lock() + c.activity.closed = true + c.activity.Unlock() + c.StopMaintenance() close(c.done) - } + }) } // SetStatusTTL sets the duration for which cached status is considered fresh. @@ -150,7 +157,7 @@ func (c *Client) setStock(statusBytes []byte) { } // StartPolling performs a periodic status refresh. -// It will NOT interrupt commands: it enqueues only when queue is idle. +// Passive requests are admitted by the worker only for the captured idle generation. func (c *Client) StartPolling(interval time.Duration) { if interval <= 0 { return @@ -168,8 +175,13 @@ func (c *Client) StartPolling(interval time.Duration) { if len(c.reqCh) != 0 { continue } + generation, idle := c.passiveGeneration() + if !idle { + continue + } ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) - _, err := c.CheckStatus(ctx) + response := c.sendRequest(cmdReq{typ: cmdStatus, ctx: ctx, passive: true, generation: generation}) + err := response.err if err != nil { log.Debugf("dispenser polling: %v", err) } @@ -180,12 +192,28 @@ func (c *Client) StartPolling(interval time.Duration) { } func (c *Client) loop() { + timer := time.NewTimer(time.Hour) + defer timer.Stop() for { + if !timer.Stop() { + select { + case <-timer.C: + default: + } + } + var tick <-chan time.Time + if delay, enabled := c.maintenanceDelay(); enabled { + timer.Reset(delay) + tick = timer.C + } select { case <-c.done: return + case <-c.activityWake: case req := <-c.reqCh: c.handle(req) + case <-tick: + c.maintainIdleCard() } } } @@ -198,6 +226,21 @@ func (c *Client) handle(req cmdReq) { default: } + if req.passive { + if !c.admitPassive(req.ctx, req.generation) { + req.respCh <- cmdResp{} + return + } + c.mu.RLock() + st := append([]byte(nil), c.lastStatus...) + fresh := len(st) == 4 && time.Since(c.lastStatusT) <= c.statusTTL + c.mu.RUnlock() + if fresh { + c.setStock(st) + req.respCh <- cmdResp{status: st} + return + } + } switch req.typ { case cmdStatus: st, err := c.readWorkerStatus(req.ctx) @@ -300,20 +343,33 @@ func (c *Client) do(ctx context.Context, typ cmdType) ([]byte, error) { } func (c *Client) doResponse(ctx context.Context, typ cmdType) cmdResp { - rch := make(chan cmdResp, 1) - req := cmdReq{typ: typ, ctx: ctx, respCh: rch} + return c.sendRequest(cmdReq{typ: typ, ctx: ctx}) +} +func (c *Client) sendRequest(req cmdReq) cmdResp { + if err := req.ctx.Err(); err != nil { + return cmdResp{err: err} + } + req.respCh = make(chan cmdResp, 1) + select { + case <-c.done: + return cmdResp{err: context.Canceled} + default: + } select { case c.reqCh <- req: - case <-ctx.Done(): - return cmdResp{err: ctx.Err()} + case <-c.done: + return cmdResp{err: context.Canceled} + case <-req.ctx.Done(): + return cmdResp{err: req.ctx.Err()} } - select { - case r := <-rch: + case r := <-req.respCh: return r - case <-ctx.Done(): - return cmdResp{err: ctx.Err()} + case <-c.done: + return cmdResp{err: context.Canceled} + case <-req.ctx.Done(): + return cmdResp{err: req.ctx.Err()} } } diff --git a/internal/dispenser/maintenance.go b/internal/dispenser/maintenance.go new file mode 100644 index 0000000..fe15271 --- /dev/null +++ b/internal/dispenser/maintenance.go @@ -0,0 +1,180 @@ +package dispenser + +import ( + "context" + "sync" + "time" + + log "github.com/sirupsen/logrus" +) + +const maintenanceInterval = time.Minute +const maintenanceTimeout = 5 * time.Second + +// activityGuard arbitrates admission, not the duration of serial transactions. +// The serial worker remains the sole owner of the port and deliveryPending. +type activityGuard struct { + sync.Mutex + foreground int + generation uint64 + enabled bool + closed bool + next time.Time + cancel context.CancelFunc + finished chan struct{} +} + +func (c *Client) wakeWorker() { + select { + case c.activityWake <- struct{}{}: + default: + } +} + +// BeginForeground suppresses idle work until the returned release function runs. +func (c *Client) BeginForeground() func() { + c.activity.Lock() + c.activity.foreground++ + c.activity.generation++ + c.activity.next = time.Time{} + c.activity.Unlock() + c.wakeWorker() + var once sync.Once + return func() { + once.Do(func() { + c.activity.Lock() + c.activity.foreground-- + if c.activity.foreground == 0 && c.activity.enabled && !c.activity.closed { + c.activity.next = c.now().Add(maintenanceInterval) + } + c.activity.Unlock() + c.wakeWorker() + }) + } +} + +// StartMaintenance enables the worker timer once, with an initial one-minute delay. +func (c *Client) StartMaintenance() { + c.activity.Lock() + if !c.activity.enabled && !c.activity.closed { + c.activity.enabled = true + c.activity.generation++ + if c.activity.foreground == 0 { + c.activity.next = c.now().Add(maintenanceInterval) + } + } + c.activity.Unlock() + c.wakeWorker() +} + +// StopMaintenance prevents new admissions and waits for an admitted attempt to exit. +// An underlying serial read can finish only according to its existing read timeout. +func (c *Client) StopMaintenance() { + c.activity.Lock() + c.activity.enabled = false + c.activity.generation++ + c.activity.next = time.Time{} + if c.activity.cancel != nil { + c.activity.cancel() + } + finished := c.activity.finished + c.activity.Unlock() + c.wakeWorker() + if finished != nil { + <-finished + } +} + +func (c *Client) maintenanceDelay() (time.Duration, bool) { + c.activity.Lock() + defer c.activity.Unlock() + if !c.activity.enabled || c.activity.closed || c.activity.foreground != 0 || c.activity.next.IsZero() { + return 0, false + } + remaining := c.activity.next.Sub(c.now()) + if remaining < 0 { + remaining = 0 + } + return remaining, true +} + +func (c *Client) passiveGeneration() (uint64, bool) { + c.activity.Lock() + defer c.activity.Unlock() + return c.activity.generation, c.activity.foreground == 0 && !c.activity.closed +} + +func (c *Client) admitPassive(ctx context.Context, generation uint64) bool { + c.activity.Lock() + defer c.activity.Unlock() + return ctx.Err() == nil && !c.activity.closed && c.activity.foreground == 0 && c.activity.generation == generation +} + +func (c *Client) admitMaintenance(ctx context.Context, generation uint64) bool { + c.activity.Lock() + defer c.activity.Unlock() + return ctx.Err() == nil && c.activity.enabled && !c.activity.closed && c.activity.foreground == 0 && c.activity.generation == generation +} + +// maintainIdleCard is called only by the serial worker, never through its queue. +func (c *Client) maintainIdleCard() { + c.activity.Lock() + if !c.activity.enabled || c.activity.closed || c.activity.foreground != 0 || c.activity.next.IsZero() || c.now().Before(c.activity.next) { + c.activity.Unlock() + return + } + ctx, cancel := context.WithTimeout(context.Background(), maintenanceTimeout) + generation := c.activity.generation + finished := make(chan struct{}) + c.activity.cancel, c.activity.finished = cancel, finished + c.activity.Unlock() + defer func() { + cancel() + c.activity.Lock() + c.activity.cancel, c.activity.finished = nil, nil + if c.activity.enabled && !c.activity.closed && c.activity.foreground == 0 && c.activity.generation == generation { + c.activity.next = c.now().Add(maintenanceInterval) + } + close(finished) + c.activity.Unlock() + }() + if !c.admitMaintenance(ctx, generation) { + return + } + // Quiet observation: do not publish stock callbacks or manufacture cached status. + st, err := checkDispenserStatus(ctx, c.port) + if err != nil { + log.Debugf("idle dispenser maintenance AP: %v", err) + return + } + class, err := classifyPreparationStatus(st) + if err != nil { + log.Debugf("idle dispenser maintenance unusable AP: %v", err) + return + } + log.Debugf("idle dispenser maintenance AP; class=%s raw status: % X", class, st) + if class == encoderConfirmed || class == positionWellEmpty { + return + } + clear, err := deliveryClearance(st) + if err != nil || !clear { + return + } + if c.deliveryPending && c.now().Sub(c.deliveryStarted) < deliveryMinimumWait { + return + } + // This is the admission boundary shared with foreground registration and stop. + if !c.admitMaintenance(ctx, generation) { + return + } + if c.deliveryPending { + c.deliveryPending = false + } + err = cardToEncoderPosition(ctx, c.port) + c.invalidateStatusCache() + if err != nil { + log.Warnf("idle dispenser maintenance FC7 dispatch: %v", err) + return + } + log.Info("idle dispenser maintenance FC7 dispatched") +} diff --git a/internal/dispenser/maintenance_test.go b/internal/dispenser/maintenance_test.go new file mode 100644 index 0000000..0dd9155 --- /dev/null +++ b/internal/dispenser/maintenance_test.go @@ -0,0 +1,362 @@ +package dispenser + +import ( + "context" + "errors" + "reflect" + "sync" + "testing" + "time" +) + +func maintenanceClient(t *testing.T, p *scriptedTransport) (*Client, *time.Time) { + t.Helper() + transportAddress(t) + now := time.Unix(100, 0) + c := &Client{port: p, done: make(chan struct{}), activityWake: make(chan struct{}, 1), sequenceTiming: sequenceTiming{now: func() time.Time { return now }, wait: waitForSequence}} + t.Cleanup(c.Close) + return c, &now +} + +func TestMaintenanceIdleScheduling(t *testing.T) { + c, now := maintenanceClient(t, &scriptedTransport{}) + c.StartMaintenance() + first := c.activity.next + c.StartMaintenance() + if c.activity.next != first { + t.Fatal("repeated start reset the timer") + } + *now = now.Add(59 * time.Second) + c.maintainIdleCard() + if len(c.port.(*scriptedTransport).writes) != 0 { + t.Fatal("maintenance ran before one minute") + } + release1 := c.BeginForeground() + release2 := c.BeginForeground() + release1() + if _, enabled := c.maintenanceDelay(); enabled { + t.Fatal("timer enabled while another issuance active") + } + *now = now.Add(20 * time.Second) + release2() + release2() + if delay, enabled := c.maintenanceDelay(); !enabled || delay != time.Minute { + t.Fatalf("after last release delay=%v enabled=%t", delay, enabled) + } + c.StopMaintenance() + if _, enabled := c.maintenanceDelay(); enabled { + t.Fatal("stop left timer enabled") + } +} + +func TestMaintenanceUsesCommonEligibility(t *testing.T) { + for _, tc := range []struct { + name string + st []byte + pending bool + age time.Duration + fc7 bool + }{ + {"clear", status(0x30), false, 0, true}, + {"staging", status(0x34), true, 3 * time.Second, true}, + {"too early", status(0x30), true, time.Second, false}, + {"encoder", status(0x33), true, time.Minute, false}, + {"empty", status(0x38), false, 0, false}, + {"uncertain", status(0x35), false, 0, false}, + {"movement", []byte{0x31, 0x30, 0x30, 0x34}, true, time.Minute, false}, + {"invalid", status(0x40), true, time.Minute, false}, + } { + t.Run(tc.name, func(t *testing.T) { + p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(tc.st), vendorACK}} + c, now := maintenanceClient(t, p) + c.StartMaintenance() + *now = now.Add(time.Minute) + c.deliveryPending = tc.pending + c.deliveryStarted = now.Add(-tc.age) + callbacks := 0 + c.OnStockUpdate(func(string) { callbacks++ }) + c.maintainIdleCard() + want := []string{"AP"} + if tc.fc7 { + want = append(want, "FC7") + } + if got := wireCommands(p); !reflect.DeepEqual(got, want) { + t.Errorf("commands=%v want %v", got, want) + } + if callbacks != 0 { + t.Errorf("maintenance published %d stock callbacks", callbacks) + } + if tc.pending && !tc.fc7 && !c.deliveryPending { + t.Error("ineligible maintenance cleared pending") + } + if delay, enabled := c.maintenanceDelay(); !enabled || delay != time.Minute { + t.Errorf("next maintenance delay=%v enabled=%t", delay, enabled) + } + }) + } +} + +func TestMaintenanceForegroundBetweenAPAndFC7(t *testing.T) { + p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x30)), vendorACK}} + c, now := maintenanceClient(t, p) + c.StartMaintenance() + *now = now.Add(time.Minute) + var release func() + p.afterRead = func() { + if len(p.chunks) == 1 && release == nil { + release = c.BeginForeground() + } + } + c.maintainIdleCard() + if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP"}) { + t.Errorf("foreground during AP commands=%v want AP only", got) + } + if release == nil { + t.Fatal("barrier not reached") + } + release() + if delay, _ := c.maintenanceDelay(); delay != time.Minute { + t.Errorf("foreground completion delay=%v", delay) + } +} + +func TestMaintenanceAdmissionConditions(t *testing.T) { + c, _ := maintenanceClient(t, &scriptedTransport{}) + c.StartMaintenance() + generation := c.activity.generation + ctx, cancel := context.WithCancel(context.Background()) + if !c.admitMaintenance(ctx, generation) { + t.Fatal("idle admission rejected") + } + cancel() + if c.admitMaintenance(ctx, generation) { + t.Fatal("canceled admission accepted") + } + expired, stop := context.WithDeadline(context.Background(), time.Now().Add(-time.Second)) + defer stop() + if c.admitMaintenance(expired, generation) { + t.Fatal("expired admission accepted") + } + release := c.BeginForeground() + if c.admitMaintenance(context.Background(), generation) { + t.Fatal("foreground/stale admission accepted") + } + release() + if c.admitMaintenance(context.Background(), generation) { + t.Fatal("obsolete generation accepted") + } + generation = c.activity.generation + c.StopMaintenance() + if c.admitMaintenance(context.Background(), generation) { + t.Fatal("stopped admission accepted") + } +} + +func TestPassivePollGenerationAtDispatch(t *testing.T) { + c, _ := maintenanceClient(t, &scriptedTransport{}) + generation, _ := c.passiveGeneration() + release := c.BeginForeground() + for _, active := range []bool{true, false} { + if !active { + release() + } + response := make(chan cmdResp, 1) + c.handle(cmdReq{ctx: context.Background(), typ: cmdStatus, passive: true, generation: generation, respCh: response}) + <-response + if len(c.port.(*scriptedTransport).writes) != 0 { + t.Errorf("active=%t stale passive request touched port", active) + } + } +} + +func TestMaintenanceStopWaitsForAdmittedRead(t *testing.T) { + p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x30)), vendorACK}} + c, now := maintenanceClient(t, p) + c.StartMaintenance() + *now = now.Add(time.Minute) + entered, resume := make(chan struct{}), make(chan struct{}) + var once sync.Once + p.afterRead = func() { once.Do(func() { close(entered); <-resume }) } + finished := make(chan struct{}) + go func() { c.maintainIdleCard(); close(finished) }() + <-entered + stopped := make(chan struct{}) + go func() { c.StopMaintenance(); close(stopped) }() + // Wait until stop has disabled admission. + for { + c.activity.Lock() + enabled := c.activity.enabled + c.activity.Unlock() + if !enabled { + break + } + time.Sleep(time.Millisecond) + } + select { + case <-stopped: + t.Fatal("stop returned before admitted read finished") + default: + } + close(resume) + <-finished + <-stopped + if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP"}) { + t.Errorf("stop commands=%v want AP only", got) + } + c.StopMaintenance() + c.Close() + c.Close() +} + +func TestMaintenanceWorkerTimerWake(t *testing.T) { + c := NewClient(nil, 1) + c.StartMaintenance() + release := c.BeginForeground() + // An immediate obsolete timer is harmless while foreground is registered. + c.activity.Lock() + c.activity.next = time.Now().Add(-time.Second) + c.activity.Unlock() + c.wakeWorker() + release() + if delay, enabled := c.maintenanceDelay(); !enabled || delay <= 59*time.Second { + t.Errorf("worker schedule delay=%v enabled=%t", delay, enabled) + } + c.Close() + if r := c.doResponse(context.Background(), cmdStatus); r.err == nil { + t.Fatal("closed client accepted a request") + } +} + +func TestMaintenanceWorkerTimerDispatch(t *testing.T) { + transportAddress(t) + p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x30)), vendorACK}} + c := &Client{port: p, reqCh: make(chan cmdReq, 1), done: make(chan struct{}), activityWake: make(chan struct{}, 1), sequenceTiming: sequenceTiming{now: time.Now, wait: waitForSequence}} + c.StartMaintenance() + c.activity.Lock() + c.activity.next = time.Now().Add(-time.Second) + c.activity.Unlock() + workerDone := make(chan struct{}) + go func() { defer close(workerDone); c.loop() }() + t.Cleanup(func() { c.Close(); <-workerDone }) + dispatched := make(chan struct{}) + // Observe completion through the activity guard, without touching the transport concurrently. + go func() { + for { + c.activity.Lock() + next := c.activity.next + c.activity.Unlock() + if time.Until(next) > 50*time.Second { + close(dispatched) + return + } + select { + case <-c.done: + return + case <-time.After(time.Millisecond): + } + } + }() + select { + case <-dispatched: + case <-time.After(4 * time.Second): + t.Fatal("worker timer did not finish maintenance") + } + c.StopMaintenance() + if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP", "FC7"}) { + t.Errorf("timer commands=%v", got) + } +} + +func TestForegroundEncoderClearanceCharacterization(t *testing.T) { + for _, movement := range []bool{false, true} { + st := status(0x33) + if movement { + st[0] = 0x31 + } + p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(st), vendorACK, apReply(st), vendorACK, apReply(st), vendorACK}} + c, now := maintenanceClient(t, p) + c.deliveryPending = true + c.deliveryStarted = now.Add(-time.Minute) + c.sequenceTiming.wait = func(context.Context, time.Duration) error { *now = now.Add(2 * time.Second); return nil } + r := workerRequest(c, context.Background(), cmdToEncoder) + want := []string{"AP", "AP", "AP"} + if movement { + if !errors.Is(r.err, context.DeadlineExceeded) || !c.deliveryPending { + t.Errorf("movement err=%v pending=%t", r.err, c.deliveryPending) + } + } else { + want = append(want, "FC7") + if r.err != nil || c.deliveryPending { + t.Errorf("internal fallback err=%v pending=%t", r.err, c.deliveryPending) + } + } + if got := wireCommands(p); !reflect.DeepEqual(got, want) { + t.Errorf("movement=%t commands=%v want %v", movement, got, want) + } + } +} + +func TestForegroundEncoderClearanceCallerDeadline(t *testing.T) { + p := &scriptedTransport{chunks: [][]byte{vendorACK, apReply(status(0x33)), vendorACK, apReply(status(0x33)), vendorACK, apReply(status(0x33))}} + c, now := maintenanceClient(t, p) + c.deliveryPending = true + c.deliveryStarted = *now + start := *now + samples := 0 + p.afterRead = func() { + if len(p.chunks)%2 == 0 { // Only complete frame reads advance the observation clock. + if len(p.chunks) == 4 || len(p.chunks) == 2 || len(p.chunks) == 0 { + samples++ + *now = now.Add(time.Second) + } + } + } + ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second) + defer cancel() + c.sequenceTiming.wait = func(context.Context, time.Duration) error { + if samples == 3 { + <-ctx.Done() + return ctx.Err() + } + *now = now.Add(time.Second) + return nil + } + r := workerRequest(c, ctx, cmdToEncoder) + if !errors.Is(r.err, context.DeadlineExceeded) || !c.deliveryPending { + t.Errorf("caller deadline err=%v pending=%t", r.err, c.deliveryPending) + } + if got := wireCommands(p); !reflect.DeepEqual(got, []string{"AP", "AP", "AP"}) { + t.Errorf("caller deadline commands=%v want AP only", got) + } + if elapsed := now.Sub(start); elapsed != 5*time.Second { + t.Errorf("last observation elapsed=%v want 5s", elapsed) + } +} + +func TestMaintenanceFailuresAreQuietAndRetryLater(t *testing.T) { + for _, mechanical := range []bool{false, true} { + p := &scriptedTransport{chunks: [][]byte{{0x10, 0x06, 0x30}}} + if mechanical { + p.chunks = [][]byte{vendorACK, apReply(status(0x30))} + } + c, now := maintenanceClient(t, p) + c.StartMaintenance() + *now = now.Add(time.Minute) + callbacks := 0 + c.OnStockUpdate(func(string) { callbacks++ }) + c.maintainIdleCard() + want := []string{"AP"} + if mechanical { + want = append(want, "FC7") + } + if got := wireCommands(p); !reflect.DeepEqual(got, want) { + t.Errorf("mechanical=%t commands=%v want %v", mechanical, got, want) + } + if callbacks != 0 { + t.Errorf("maintenance failure callbacks=%d", callbacks) + } + if delay, enabled := c.maintenanceDelay(); !enabled || delay != time.Minute { + t.Errorf("retry delay=%v enabled=%t", delay, enabled) + } + } +} diff --git a/internal/handlers/doorcard_handlers_test.go b/internal/handlers/doorcard_handlers_test.go index 7b13b63..e0fe71b 100644 --- a/internal/handlers/doorcard_handlers_test.go +++ b/internal/handlers/doorcard_handlers_test.go @@ -25,15 +25,22 @@ type dispenserCallResult struct { } type fakeDoorCardDispenser struct { - prepareCurrent dispenserCallResult - deliverCurrent dispenserCallResult - beginNextErr error - beginNext func() error - prepareNext dispenserCallResult - calls []string + prepareCurrent dispenserCallResult + deliverCurrent dispenserCallResult + beginNextErr error + beginNext func() error + prepareNext dispenserCallResult + activity int + registrations int + releases int + outsideActivity bool + calls []string } func (d *fakeDoorCardDispenser) PrepareCurrentCard(context.Context) (string, error) { + if d.activity == 0 { + d.outsideActivity = true + } d.calls = append(d.calls, "prepare current") return d.prepareCurrent.status, d.prepareCurrent.err } @@ -377,3 +384,28 @@ func TestIssueDoorCardOnlyEmptyPreparationReturns503(t *testing.T) { } } } + +func (d *fakeDoorCardDispenser) BeginForeground() func() { + d.activity++ + d.registrations++ + return func() { d.activity--; d.releases++ } +} + +func TestDoorCardForegroundLifecycle(t *testing.T) { + for _, testEndpoint := range []bool{false, true} { + d := &fakeDoorCardDispenser{prepareCurrent: dispenserCallResult{err: errors.New("stop before encoding")}} + lock := &fakeDoorCardLockServer{} + _, _, app := performIssueDoorCardRequest(t, d, lock) + if d.registrations != 1 || d.releases != 1 || d.activity != 0 || d.outsideActivity { + t.Fatalf("issue registration=%d release=%d active=%d outside=%t", d.registrations, d.releases, d.activity, d.outsideActivity) + } + if testEndpoint { + req := httptest.NewRequest(http.MethodPost, "/testissuedoorcard", strings.NewReader("{}")) + req.Header.Set("Content-Type", "application/json") + app.testIssueDoorCard(httptest.NewRecorder(), req) + if d.registrations != 2 || d.releases != 2 || d.activity != 0 || d.outsideActivity { + t.Errorf("test endpoint registration=%d release=%d active=%d outside=%t", d.registrations, d.releases, d.activity, d.outsideActivity) + } + } + } +} diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 354e3be..cbf75dd 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -27,6 +27,7 @@ import ( ) type doorCardDispenser interface { + BeginForeground() func() PrepareCurrentCard(context.Context) (string, error) DeliverCurrentCard(context.Context) (string, error) BeginPrepareNextCard(context.Context) error @@ -290,6 +291,9 @@ func (app *App) issueDoorCard(w http.ResponseWriter, r *http.Request) { return } + release := app.disp.BeginForeground() + defer release() + status, err := app.disp.PrepareCurrentCard(r.Context()) app.SetCardWellStatus(status) if err != nil { diff --git a/internal/handlers/testhandlers.go b/internal/handlers/testhandlers.go index c512b6f..a03bb84 100644 --- a/internal/handlers/testhandlers.go +++ b/internal/handlers/testhandlers.go @@ -59,6 +59,9 @@ func (app *App) testIssueDoorCard(w http.ResponseWriter, r *http.Request) { // Ensure dispenser ready (card at encoder) BEFORE we attempt encoding. // With queued dispenser ops, this will not clash with polling. + release := app.disp.BeginForeground() + defer release() + status, err := app.disp.PrepareCurrentCard(r.Context()) app.SetCardWellStatus(status) if err != nil { diff --git a/release notes.md b/release notes.md index c70d255..3485805 100644 --- a/release notes.md +++ b/release notes.md @@ -2,6 +2,9 @@ builtVersion is a const in main.go +#### v2.1.1 - 28 September 2026 +feat(dispenser): add worker-owned idle card prestaging + #### v2.1.0 - 28 September 2026 fix(dispenser): tolerate unusable AP observations during card preparation