Compare commits
No commits in common. "development" and "v2.1.0" have entirely different histories.
developmen
...
v2.1.0
@ -33,7 +33,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
buildVersion = "v2.1.3"
|
buildVersion = "v2.1.0"
|
||||||
serviceName = "hardlink"
|
serviceName = "hardlink"
|
||||||
pollingFrequency = 8 * time.Second
|
pollingFrequency = 8 * time.Second
|
||||||
)
|
)
|
||||||
@ -227,7 +227,6 @@ func main() {
|
|||||||
|
|
||||||
// Start polling for dispenser status every 10 seconds
|
// Start polling for dispenser status every 10 seconds
|
||||||
disp.StartPolling(pollingFrequency)
|
disp.StartPolling(pollingFrequency)
|
||||||
disp.StartMaintenance()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
|
|||||||
@ -1,335 +0,0 @@
|
|||||||
# 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.
|
|
||||||
@ -295,60 +295,11 @@ func sendAndReadACK(ctx context.Context, port serialTransport, packet []byte, pr
|
|||||||
if err := waitForSequence(ctx, processingDelay); err != nil {
|
if err := waitForSequence(ctx, processingDelay); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return readACK(ctx, port)
|
response := make([]byte, 3)
|
||||||
}
|
if err := readExact(ctx, port, response); err != nil {
|
||||||
|
return fmt.Errorf("read ACK: %w", err)
|
||||||
// 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)
|
|
||||||
}
|
}
|
||||||
}()
|
return checkACK(response)
|
||||||
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 fmt.Errorf("read ACK: scan limit of %d bytes exhausted", scanLimit)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// queryStatus accepts only the fixed RF/AP payload sizes, before reading a body.
|
// queryStatus accepts only the fixed RF/AP payload sizes, before reading a body.
|
||||||
|
|||||||
@ -23,8 +23,6 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type cmdReq struct {
|
type cmdReq struct {
|
||||||
passive bool
|
|
||||||
generation uint64
|
|
||||||
typ cmdType
|
typ cmdType
|
||||||
ctx context.Context
|
ctx context.Context
|
||||||
respCh chan cmdResp
|
respCh chan cmdResp
|
||||||
@ -53,9 +51,6 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type Client struct {
|
type Client struct {
|
||||||
activity activityGuard
|
|
||||||
activityWake chan struct{}
|
|
||||||
closeOnce sync.Once
|
|
||||||
port serialTransport
|
port serialTransport
|
||||||
|
|
||||||
reqCh chan cmdReq
|
reqCh chan cmdReq
|
||||||
@ -85,7 +80,6 @@ func NewClient(port *serial.Port, queueSize int) *Client {
|
|||||||
queueSize = 16
|
queueSize = 16
|
||||||
}
|
}
|
||||||
c := &Client{
|
c := &Client{
|
||||||
activityWake: make(chan struct{}, 1),
|
|
||||||
port: port,
|
port: port,
|
||||||
reqCh: make(chan cmdReq, queueSize),
|
reqCh: make(chan cmdReq, queueSize),
|
||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
@ -112,13 +106,12 @@ func waitForSequence(ctx context.Context, duration time.Duration) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *Client) Close() {
|
func (c *Client) Close() {
|
||||||
c.closeOnce.Do(func() {
|
select {
|
||||||
c.activity.Lock()
|
case <-c.done:
|
||||||
c.activity.closed = true
|
return
|
||||||
c.activity.Unlock()
|
default:
|
||||||
c.StopMaintenance()
|
|
||||||
close(c.done)
|
close(c.done)
|
||||||
})
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// SetStatusTTL sets the duration for which cached status is considered fresh.
|
// SetStatusTTL sets the duration for which cached status is considered fresh.
|
||||||
@ -157,7 +150,7 @@ func (c *Client) setStock(statusBytes []byte) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// StartPolling performs a periodic status refresh.
|
// StartPolling performs a periodic status refresh.
|
||||||
// Passive requests are admitted by the worker only for the captured idle generation.
|
// It will NOT interrupt commands: it enqueues only when queue is idle.
|
||||||
func (c *Client) StartPolling(interval time.Duration) {
|
func (c *Client) StartPolling(interval time.Duration) {
|
||||||
if interval <= 0 {
|
if interval <= 0 {
|
||||||
return
|
return
|
||||||
@ -175,13 +168,8 @@ func (c *Client) StartPolling(interval time.Duration) {
|
|||||||
if len(c.reqCh) != 0 {
|
if len(c.reqCh) != 0 {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
generation, idle := c.passiveGeneration()
|
|
||||||
if !idle {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
response := c.sendRequest(cmdReq{typ: cmdStatus, ctx: ctx, passive: true, generation: generation})
|
_, err := c.CheckStatus(ctx)
|
||||||
err := response.err
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Debugf("dispenser polling: %v", err)
|
log.Debugf("dispenser polling: %v", err)
|
||||||
}
|
}
|
||||||
@ -192,28 +180,12 @@ func (c *Client) StartPolling(interval time.Duration) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *Client) loop() {
|
func (c *Client) loop() {
|
||||||
timer := time.NewTimer(time.Hour)
|
|
||||||
defer timer.Stop()
|
|
||||||
for {
|
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 {
|
select {
|
||||||
case <-c.done:
|
case <-c.done:
|
||||||
return
|
return
|
||||||
case <-c.activityWake:
|
|
||||||
case req := <-c.reqCh:
|
case req := <-c.reqCh:
|
||||||
c.handle(req)
|
c.handle(req)
|
||||||
case <-tick:
|
|
||||||
c.maintainIdleCard()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -226,21 +198,6 @@ func (c *Client) handle(req cmdReq) {
|
|||||||
default:
|
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 {
|
switch req.typ {
|
||||||
case cmdStatus:
|
case cmdStatus:
|
||||||
st, err := c.readWorkerStatus(req.ctx)
|
st, err := c.readWorkerStatus(req.ctx)
|
||||||
@ -343,33 +300,20 @@ func (c *Client) do(ctx context.Context, typ cmdType) ([]byte, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *Client) doResponse(ctx context.Context, typ cmdType) cmdResp {
|
func (c *Client) doResponse(ctx context.Context, typ cmdType) cmdResp {
|
||||||
return c.sendRequest(cmdReq{typ: typ, ctx: ctx})
|
rch := make(chan cmdResp, 1)
|
||||||
}
|
req := cmdReq{typ: typ, ctx: ctx, respCh: rch}
|
||||||
|
|
||||||
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 {
|
select {
|
||||||
case c.reqCh <- req:
|
case c.reqCh <- req:
|
||||||
case <-c.done:
|
case <-ctx.Done():
|
||||||
return cmdResp{err: context.Canceled}
|
return cmdResp{err: ctx.Err()}
|
||||||
case <-req.ctx.Done():
|
|
||||||
return cmdResp{err: req.ctx.Err()}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
select {
|
select {
|
||||||
case r := <-req.respCh:
|
case r := <-rch:
|
||||||
return r
|
return r
|
||||||
case <-c.done:
|
case <-ctx.Done():
|
||||||
return cmdResp{err: context.Canceled}
|
return cmdResp{err: ctx.Err()}
|
||||||
case <-req.ctx.Done():
|
|
||||||
return cmdResp{err: req.ctx.Err()}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -1,180 +0,0 @@
|
|||||||
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")
|
|
||||||
}
|
|
||||||
@ -1,362 +0,0 @@
|
|||||||
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@ -6,7 +6,6 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"os"
|
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@ -25,146 +24,6 @@ var vendorFrames = []struct {
|
|||||||
{"RS", commandRS, []byte{2, 0x30, 0x30, 0, 2, 0x52, 0x53, 3, 2}},
|
{"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) {
|
func TestVendorOutboundFrames(t *testing.T) {
|
||||||
for _, tc := range vendorFrames {
|
for _, tc := range vendorFrames {
|
||||||
t.Run(tc.name, func(t *testing.T) {
|
t.Run(tc.name, func(t *testing.T) {
|
||||||
@ -357,11 +216,7 @@ func TestTransportCancellation(t *testing.T) {
|
|||||||
case "processing wait":
|
case "processing wait":
|
||||||
p.afterWrite = cancel
|
p.afterWrite = cancel
|
||||||
case "after ACK":
|
case "after ACK":
|
||||||
p.afterRead = func() {
|
p.afterRead = cancel
|
||||||
if len(p.chunks) == 1 {
|
|
||||||
cancel()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
case "after ENQ":
|
case "after ENQ":
|
||||||
wantWrites = 2
|
wantWrites = 2
|
||||||
p.afterWrite = func() {
|
p.afterWrite = func() {
|
||||||
@ -374,7 +229,7 @@ func TestTransportCancellation(t *testing.T) {
|
|||||||
reads := 0
|
reads := 0
|
||||||
p.afterRead = func() {
|
p.afterRead = func() {
|
||||||
reads++
|
reads++
|
||||||
if reads == 3 {
|
if reads == 2 {
|
||||||
cancel()
|
cancel()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -30,17 +30,10 @@ type fakeDoorCardDispenser struct {
|
|||||||
beginNextErr error
|
beginNextErr error
|
||||||
beginNext func() error
|
beginNext func() error
|
||||||
prepareNext dispenserCallResult
|
prepareNext dispenserCallResult
|
||||||
activity int
|
|
||||||
registrations int
|
|
||||||
releases int
|
|
||||||
outsideActivity bool
|
|
||||||
calls []string
|
calls []string
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d *fakeDoorCardDispenser) PrepareCurrentCard(context.Context) (string, error) {
|
func (d *fakeDoorCardDispenser) PrepareCurrentCard(context.Context) (string, error) {
|
||||||
if d.activity == 0 {
|
|
||||||
d.outsideActivity = true
|
|
||||||
}
|
|
||||||
d.calls = append(d.calls, "prepare current")
|
d.calls = append(d.calls, "prepare current")
|
||||||
return d.prepareCurrent.status, d.prepareCurrent.err
|
return d.prepareCurrent.status, d.prepareCurrent.err
|
||||||
}
|
}
|
||||||
@ -384,28 +377,3 @@ 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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@ -27,7 +27,6 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type doorCardDispenser interface {
|
type doorCardDispenser interface {
|
||||||
BeginForeground() func()
|
|
||||||
PrepareCurrentCard(context.Context) (string, error)
|
PrepareCurrentCard(context.Context) (string, error)
|
||||||
DeliverCurrentCard(context.Context) (string, error)
|
DeliverCurrentCard(context.Context) (string, error)
|
||||||
BeginPrepareNextCard(context.Context) error
|
BeginPrepareNextCard(context.Context) error
|
||||||
@ -291,9 +290,6 @@ func (app *App) issueDoorCard(w http.ResponseWriter, r *http.Request) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
release := app.disp.BeginForeground()
|
|
||||||
defer release()
|
|
||||||
|
|
||||||
status, err := app.disp.PrepareCurrentCard(r.Context())
|
status, err := app.disp.PrepareCurrentCard(r.Context())
|
||||||
app.SetCardWellStatus(status)
|
app.SetCardWellStatus(status)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@ -59,9 +59,6 @@ func (app *App) testIssueDoorCard(w http.ResponseWriter, r *http.Request) {
|
|||||||
|
|
||||||
// Ensure dispenser ready (card at encoder) BEFORE we attempt encoding.
|
// Ensure dispenser ready (card at encoder) BEFORE we attempt encoding.
|
||||||
// With queued dispenser ops, this will not clash with polling.
|
// With queued dispenser ops, this will not clash with polling.
|
||||||
release := app.disp.BeginForeground()
|
|
||||||
defer release()
|
|
||||||
|
|
||||||
status, err := app.disp.PrepareCurrentCard(r.Context())
|
status, err := app.disp.PrepareCurrentCard(r.Context())
|
||||||
app.SetCardWellStatus(status)
|
app.SetCardWellStatus(status)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@ -2,22 +2,17 @@ package logging
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
log "github.com/sirupsen/logrus"
|
log "github.com/sirupsen/logrus"
|
||||||
)
|
)
|
||||||
|
|
||||||
// SetupLogging ensures the log directory, opens the rotating log writer, and configures logrus.
|
// setupLogging ensures log directory, opens log file, and configures logrus.
|
||||||
func SetupLogging(logDir, serviceName, buildVersion string) (io.WriteCloser, error) {
|
// Returns the *os.File so caller can defer its Close().
|
||||||
if err := os.MkdirAll(logDir, 0o755); err != nil {
|
func SetupLogging(logDir, serviceName, buildVersion string) (*os.File, error) {
|
||||||
return nil, fmt.Errorf("create log directory: %w", err)
|
fileName := logDir + serviceName + ".log"
|
||||||
}
|
f, err := os.OpenFile(fileName, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0666)
|
||||||
|
|
||||||
fileName := filepath.Join(logDir, serviceName+".log")
|
|
||||||
f, err := newWeeklyLogWriter(fileName, time.Local, defaultWeeklyLogRuntime())
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("open log file: %w", err)
|
return nil, fmt.Errorf("open log file: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@ -1,319 +0,0 @@
|
|||||||
package logging
|
|
||||||
|
|
||||||
import (
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"io"
|
|
||||||
"io/fs"
|
|
||||||
"log/slog"
|
|
||||||
"os"
|
|
||||||
"path/filepath"
|
|
||||||
"strconv"
|
|
||||||
"strings"
|
|
||||||
"sync"
|
|
||||||
"time"
|
|
||||||
_ "time/tzdata"
|
|
||||||
)
|
|
||||||
|
|
||||||
const weeklyLogRetention = 13
|
|
||||||
|
|
||||||
type logRotationTimer interface {
|
|
||||||
C() <-chan time.Time
|
|
||||||
Stop() bool
|
|
||||||
}
|
|
||||||
|
|
||||||
type realLogRotationTimer struct {
|
|
||||||
*time.Timer
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t realLogRotationTimer) C() <-chan time.Time {
|
|
||||||
return t.Timer.C
|
|
||||||
}
|
|
||||||
|
|
||||||
type weeklyLogRuntime struct {
|
|
||||||
now func() time.Time
|
|
||||||
newTimer func(time.Duration) logRotationTimer
|
|
||||||
rename func(string, string) error
|
|
||||||
diagnostic io.Writer
|
|
||||||
}
|
|
||||||
|
|
||||||
func defaultWeeklyLogRuntime() weeklyLogRuntime {
|
|
||||||
return weeklyLogRuntime{
|
|
||||||
now: time.Now,
|
|
||||||
newTimer: func(duration time.Duration) logRotationTimer {
|
|
||||||
return realLogRotationTimer{Timer: time.NewTimer(duration)}
|
|
||||||
},
|
|
||||||
rename: os.Rename,
|
|
||||||
diagnostic: os.Stderr,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type weeklyLogWriter struct {
|
|
||||||
path string
|
|
||||||
location *time.Location
|
|
||||||
now func() time.Time
|
|
||||||
newTimer func(time.Duration) logRotationTimer
|
|
||||||
rename func(string, string) error
|
|
||||||
diagnostic io.Writer
|
|
||||||
|
|
||||||
mu sync.Mutex
|
|
||||||
file *os.File
|
|
||||||
weekStart time.Time
|
|
||||||
closed bool
|
|
||||||
stop chan struct{}
|
|
||||||
done chan struct{}
|
|
||||||
closeOnce sync.Once
|
|
||||||
closeErr error
|
|
||||||
}
|
|
||||||
|
|
||||||
func newWeeklyLogWriter(path string, location *time.Location, runtime weeklyLogRuntime) (*weeklyLogWriter, error) {
|
|
||||||
if location == nil {
|
|
||||||
location = time.Local
|
|
||||||
}
|
|
||||||
if runtime.now == nil {
|
|
||||||
runtime.now = time.Now
|
|
||||||
}
|
|
||||||
if runtime.newTimer == nil {
|
|
||||||
runtime.newTimer = defaultWeeklyLogRuntime().newTimer
|
|
||||||
}
|
|
||||||
if runtime.rename == nil {
|
|
||||||
runtime.rename = os.Rename
|
|
||||||
}
|
|
||||||
if runtime.diagnostic == nil {
|
|
||||||
runtime.diagnostic = os.Stderr
|
|
||||||
}
|
|
||||||
|
|
||||||
now := runtime.now().In(location)
|
|
||||||
weekStart := logWeekStart(now, location)
|
|
||||||
if info, err := os.Stat(path); err == nil {
|
|
||||||
weekStart = logWeekStart(info.ModTime(), location)
|
|
||||||
} else if !errors.Is(err, os.ErrNotExist) {
|
|
||||||
return nil, fmt.Errorf("inspect active log: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
file, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o666)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
writer := &weeklyLogWriter{
|
|
||||||
path: path,
|
|
||||||
location: location,
|
|
||||||
now: runtime.now,
|
|
||||||
newTimer: runtime.newTimer,
|
|
||||||
rename: runtime.rename,
|
|
||||||
diagnostic: runtime.diagnostic,
|
|
||||||
file: file,
|
|
||||||
weekStart: weekStart,
|
|
||||||
stop: make(chan struct{}),
|
|
||||||
done: make(chan struct{}),
|
|
||||||
}
|
|
||||||
writer.rotateIfNeeded(now)
|
|
||||||
go writer.run()
|
|
||||||
return writer, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) Write(data []byte) (int, error) {
|
|
||||||
now := w.now().In(w.location)
|
|
||||||
w.mu.Lock()
|
|
||||||
rotationErr := w.rotateIfNeededLocked(now)
|
|
||||||
if w.closed || w.file == nil {
|
|
||||||
w.mu.Unlock()
|
|
||||||
w.reportRotationError(rotationErr)
|
|
||||||
return 0, os.ErrClosed
|
|
||||||
}
|
|
||||||
n, writeErr := w.file.Write(data)
|
|
||||||
w.mu.Unlock()
|
|
||||||
w.reportRotationError(rotationErr)
|
|
||||||
return n, writeErr
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) Close() error {
|
|
||||||
w.closeOnce.Do(func() {
|
|
||||||
close(w.stop)
|
|
||||||
<-w.done
|
|
||||||
|
|
||||||
w.mu.Lock()
|
|
||||||
w.closed = true
|
|
||||||
if w.file != nil {
|
|
||||||
w.closeErr = w.file.Close()
|
|
||||||
w.file = nil
|
|
||||||
}
|
|
||||||
w.mu.Unlock()
|
|
||||||
})
|
|
||||||
return w.closeErr
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) run() {
|
|
||||||
defer close(w.done)
|
|
||||||
for {
|
|
||||||
now := w.now().In(w.location)
|
|
||||||
next := nextLogWeekStart(now, w.location)
|
|
||||||
duration := next.Sub(now)
|
|
||||||
if duration <= 0 {
|
|
||||||
duration = time.Nanosecond
|
|
||||||
}
|
|
||||||
timer := w.newTimer(duration)
|
|
||||||
select {
|
|
||||||
case <-timer.C():
|
|
||||||
w.rotateIfNeeded(w.now().In(w.location))
|
|
||||||
case <-w.stop:
|
|
||||||
timer.Stop()
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) rotateIfNeeded(now time.Time) {
|
|
||||||
w.mu.Lock()
|
|
||||||
err := w.rotateIfNeededLocked(now)
|
|
||||||
w.mu.Unlock()
|
|
||||||
w.reportRotationError(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) rotateIfNeededLocked(now time.Time) error {
|
|
||||||
if w.closed {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
targetWeek := logWeekStart(now, w.location)
|
|
||||||
if !targetWeek.After(w.weekStart) {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Record the attempted week even if rotation fails so every write in a
|
|
||||||
// broken environment does not retry and emit another diagnostic.
|
|
||||||
w.weekStart = targetWeek
|
|
||||||
return w.rotateLocked()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) rotateLocked() error {
|
|
||||||
if w.file == nil {
|
|
||||||
return fmt.Errorf("active log file is unavailable")
|
|
||||||
}
|
|
||||||
if err := w.file.Close(); err != nil {
|
|
||||||
w.file = nil
|
|
||||||
return errors.Join(fmt.Errorf("close active log: %w", err), w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
w.file = nil
|
|
||||||
|
|
||||||
if err := w.cleanupOlderGenerationsLocked(); err != nil {
|
|
||||||
return errors.Join(err, w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
if err := removeIfExists(w.generationPath(weeklyLogRetention)); err != nil {
|
|
||||||
return errors.Join(fmt.Errorf("remove oldest weekly log: %w", err), w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
for generation := weeklyLogRetention - 1; generation >= 1; generation-- {
|
|
||||||
source := w.generationPath(generation)
|
|
||||||
if _, err := os.Stat(source); err != nil {
|
|
||||||
if errors.Is(err, os.ErrNotExist) {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
return errors.Join(fmt.Errorf("inspect weekly log generation %d: %w", generation, err), w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
if err := w.rename(source, w.generationPath(generation+1)); err != nil {
|
|
||||||
return errors.Join(fmt.Errorf("shift weekly log generation %d: %w", generation, err), w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
activeRenamed := false
|
|
||||||
if _, err := os.Stat(w.path); err == nil {
|
|
||||||
if err := w.rename(w.path, w.generationPath(1)); err != nil {
|
|
||||||
return errors.Join(fmt.Errorf("archive active log: %w", err), w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
activeRenamed = true
|
|
||||||
} else if !errors.Is(err, os.ErrNotExist) {
|
|
||||||
return errors.Join(fmt.Errorf("inspect active log: %w", err), w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
|
|
||||||
file, err := os.OpenFile(w.path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o666)
|
|
||||||
if err == nil {
|
|
||||||
w.file = file
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
rotationErr := fmt.Errorf("open new active log: %w", err)
|
|
||||||
if activeRenamed {
|
|
||||||
if rollbackErr := w.rename(w.generationPath(1), w.path); rollbackErr != nil {
|
|
||||||
rotationErr = errors.Join(rotationErr, fmt.Errorf("restore archived active log: %w", rollbackErr))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return errors.Join(rotationErr, w.reopenActiveLocked())
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) cleanupOlderGenerationsLocked() error {
|
|
||||||
directory := filepath.Dir(w.path)
|
|
||||||
entries, err := os.ReadDir(directory)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("list weekly log directory: %w", err)
|
|
||||||
}
|
|
||||||
base := filepath.Base(w.path)
|
|
||||||
extension := filepath.Ext(base)
|
|
||||||
stem := strings.TrimSuffix(base, extension)
|
|
||||||
prefix := stem + "."
|
|
||||||
|
|
||||||
for _, entry := range entries {
|
|
||||||
if entry.IsDir() {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
name := entry.Name()
|
|
||||||
if !strings.HasPrefix(name, prefix) || !strings.HasSuffix(name, extension) {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
generationText := strings.TrimSuffix(strings.TrimPrefix(name, prefix), extension)
|
|
||||||
generation, err := strconv.Atoi(generationText)
|
|
||||||
if err != nil || generation <= weeklyLogRetention {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
if err := os.Remove(filepath.Join(directory, name)); err != nil && !errors.Is(err, os.ErrNotExist) {
|
|
||||||
return fmt.Errorf("remove weekly log generation %d: %w", generation, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) reopenActiveLocked() error {
|
|
||||||
file, err := os.OpenFile(w.path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o666)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("reopen active log: %w", err)
|
|
||||||
}
|
|
||||||
w.file = file
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) generationPath(generation int) string {
|
|
||||||
extension := filepath.Ext(w.path)
|
|
||||||
stem := strings.TrimSuffix(w.path, extension)
|
|
||||||
return fmt.Sprintf("%s.%d%s", stem, generation, extension)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (w *weeklyLogWriter) reportRotationError(err error) {
|
|
||||||
if err == nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
slog.New(slog.NewJSONHandler(w.diagnostic, nil)).Warn(
|
|
||||||
"weekly log rotation failed",
|
|
||||||
"module", "logging",
|
|
||||||
"action", "weekly_log_rotation_failed",
|
|
||||||
"path", w.path,
|
|
||||||
"error.message", err.Error(),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
func logWeekStart(value time.Time, location *time.Location) time.Time {
|
|
||||||
local := value.In(location)
|
|
||||||
daysSinceMonday := (int(local.Weekday()) + 6) % 7
|
|
||||||
monday := local.AddDate(0, 0, -daysSinceMonday)
|
|
||||||
return time.Date(monday.Year(), monday.Month(), monday.Day(), 0, 0, 0, 0, location)
|
|
||||||
}
|
|
||||||
|
|
||||||
func nextLogWeekStart(value time.Time, location *time.Location) time.Time {
|
|
||||||
return logWeekStart(value, location).AddDate(0, 0, 7)
|
|
||||||
}
|
|
||||||
|
|
||||||
func removeIfExists(path string) error {
|
|
||||||
err := os.Remove(path)
|
|
||||||
if errors.Is(err, fs.ErrNotExist) {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
@ -1,387 +0,0 @@
|
|||||||
package logging
|
|
||||||
|
|
||||||
import (
|
|
||||||
"bytes"
|
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"io"
|
|
||||||
"os"
|
|
||||||
"path/filepath"
|
|
||||||
"strings"
|
|
||||||
"sync"
|
|
||||||
"sync/atomic"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
log "github.com/sirupsen/logrus"
|
|
||||||
)
|
|
||||||
|
|
||||||
func TestWeeklyLogRotationShiftsAndRetainsExactGenerations(t *testing.T) {
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
writeLogTestFile(t, path, "active")
|
|
||||||
for generation := 1; generation <= 14; generation++ {
|
|
||||||
writeLogTestFile(t, weeklyLogTestGeneration(path, generation), fmt.Sprintf("generation-%d", generation))
|
|
||||||
}
|
|
||||||
writeLogTestFile(t, weeklyLogTestGeneration(path, 19), "generation-19")
|
|
||||||
for name, contents := range map[string]string{
|
|
||||||
"hardlink.backup.14.log": "related-looking backup",
|
|
||||||
"another.14.log": "unrelated log",
|
|
||||||
"hardlink.14.log.bak": "different extension",
|
|
||||||
} {
|
|
||||||
writeLogTestFile(t, filepath.Join(directory, name), contents)
|
|
||||||
}
|
|
||||||
|
|
||||||
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
|
|
||||||
setLogTestModTime(t, path, sunday)
|
|
||||||
writer := newLogTestWriter(t, path, sunday, weeklyLogRuntime{})
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
requireLogTestContents(t, path, "")
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "active")
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 2), "generation-1")
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 13), "generation-12")
|
|
||||||
for _, generation := range []int{14, 19} {
|
|
||||||
if _, err := os.Stat(weeklyLogTestGeneration(path, generation)); !errors.Is(err, os.ErrNotExist) {
|
|
||||||
t.Errorf("generation %d exists after rotation, want removed: %v", generation, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
for name, contents := range map[string]string{
|
|
||||||
"hardlink.backup.14.log": "related-looking backup",
|
|
||||||
"another.14.log": "unrelated log",
|
|
||||||
"hardlink.14.log.bak": "different extension",
|
|
||||||
} {
|
|
||||||
requireLogTestContents(t, filepath.Join(directory, name), contents)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWeeklyLogRotationOccursOncePerWeek(t *testing.T) {
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
writeLogTestFile(t, path, "week-zero\n")
|
|
||||||
sunday := time.Date(2026, time.January, 4, 18, 0, 0, 0, time.UTC)
|
|
||||||
setLogTestModTime(t, path, sunday)
|
|
||||||
writer := newLogTestWriter(t, path, sunday, weeklyLogRuntime{})
|
|
||||||
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.January, 4, 23, 59, 59, 0, time.UTC))
|
|
||||||
if _, err := os.Stat(weeklyLogTestGeneration(path, 1)); !errors.Is(err, os.ErrNotExist) {
|
|
||||||
t.Fatalf("rotation before Monday: %v, want no generation", err)
|
|
||||||
}
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.January, 5, 0, 0, 0, 0, time.UTC))
|
|
||||||
if _, err := writer.Write([]byte("week-one\n")); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.January, 6, 12, 0, 0, 0, time.UTC))
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "week-zero\n")
|
|
||||||
requireLogTestContents(t, path, "week-one\n")
|
|
||||||
if _, err := os.Stat(weeklyLogTestGeneration(path, 2)); !errors.Is(err, os.ErrNotExist) {
|
|
||||||
t.Errorf("same-week rotation created generation 2: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWeeklyLogStartupRotatesOnceAfterSeveralMissedWeeks(t *testing.T) {
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
writeLogTestFile(t, path, "old active")
|
|
||||||
writeLogTestFile(t, weeklyLogTestGeneration(path, 1), "older history")
|
|
||||||
setLogTestModTime(t, path, time.Date(2026, time.January, 5, 12, 0, 0, 0, time.UTC))
|
|
||||||
|
|
||||||
now := time.Date(2026, time.February, 23, 9, 0, 0, 0, time.UTC)
|
|
||||||
writer := newLogTestWriter(t, path, now, weeklyLogRuntime{})
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "old active")
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 2), "older history")
|
|
||||||
if _, err := os.Stat(weeklyLogTestGeneration(path, 3)); !errors.Is(err, os.ErrNotExist) {
|
|
||||||
t.Errorf("missed weeks created generation 3: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWeeklyLogSchedulerUsesLocalMondayBoundary(t *testing.T) {
|
|
||||||
location, err := time.LoadLocation("Europe/London")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
writeLogTestFile(t, path, "sunday")
|
|
||||||
sunday := time.Date(2026, time.March, 29, 0, 0, 0, 0, location)
|
|
||||||
setLogTestModTime(t, path, sunday)
|
|
||||||
|
|
||||||
clock := &atomic.Pointer[time.Time]{}
|
|
||||||
clock.Store(&sunday)
|
|
||||||
created := make(chan *logTestTimer, 2)
|
|
||||||
rotated := make(chan struct{})
|
|
||||||
runtime := weeklyLogRuntime{
|
|
||||||
now: func() time.Time { return *clock.Load() },
|
|
||||||
newTimer: func(duration time.Duration) logRotationTimer {
|
|
||||||
timer := &logTestTimer{duration: duration, ch: make(chan time.Time, 1)}
|
|
||||||
created <- timer
|
|
||||||
return timer
|
|
||||||
},
|
|
||||||
rename: func(oldPath, newPath string) error {
|
|
||||||
err := os.Rename(oldPath, newPath)
|
|
||||||
if oldPath == path && err == nil {
|
|
||||||
close(rotated)
|
|
||||||
}
|
|
||||||
return err
|
|
||||||
},
|
|
||||||
}
|
|
||||||
writer, err := newWeeklyLogWriter(path, location, runtime)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
timer := <-created
|
|
||||||
if timer.duration != 23*time.Hour {
|
|
||||||
t.Fatalf("timer duration = %v, want 23h across DST boundary", timer.duration)
|
|
||||||
}
|
|
||||||
monday := time.Date(2026, time.March, 30, 0, 0, 0, 0, location)
|
|
||||||
clock.Store(&monday)
|
|
||||||
timer.ch <- monday
|
|
||||||
<-rotated
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
requireLogTestContents(t, weeklyLogTestGeneration(path, 1), "sunday")
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWeeklyLogConcurrentWritesRemainCompleteAcrossRotation(t *testing.T) {
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
|
|
||||||
writer := newLogTestWriter(t, path, sunday, weeklyLogRuntime{})
|
|
||||||
|
|
||||||
const goroutines = 8
|
|
||||||
const recordsPerGoroutine = 100
|
|
||||||
writeErrors := make(chan error, goroutines)
|
|
||||||
var writers sync.WaitGroup
|
|
||||||
writers.Add(goroutines)
|
|
||||||
for worker := 0; worker < goroutines; worker++ {
|
|
||||||
go func() {
|
|
||||||
defer writers.Done()
|
|
||||||
for record := 0; record < recordsPerGoroutine; record++ {
|
|
||||||
if _, err := writer.Write([]byte(fmt.Sprintf("%d:%d\n", worker, record))); err != nil {
|
|
||||||
writeErrors <- err
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
|
|
||||||
writers.Wait()
|
|
||||||
close(writeErrors)
|
|
||||||
for err := range writeErrors {
|
|
||||||
t.Errorf("weeklyLogWriter.Write() error = %v, want nil", err)
|
|
||||||
}
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
contents := readLogTestFile(t, weeklyLogTestGeneration(path, 1)) + readLogTestFile(t, path)
|
|
||||||
lines := strings.Split(strings.TrimSpace(contents), "\n")
|
|
||||||
if len(lines) != goroutines*recordsPerGoroutine {
|
|
||||||
t.Fatalf("record count = %d, want %d", len(lines), goroutines*recordsPerGoroutine)
|
|
||||||
}
|
|
||||||
seen := make(map[string]bool, len(lines))
|
|
||||||
for _, line := range lines {
|
|
||||||
if _, _, ok := strings.Cut(line, ":"); !ok {
|
|
||||||
t.Fatalf("record %q has no separator", line)
|
|
||||||
}
|
|
||||||
seen[line] = true
|
|
||||||
}
|
|
||||||
if len(seen) != goroutines*recordsPerGoroutine {
|
|
||||||
t.Errorf("unique record count = %d, want %d", len(seen), goroutines*recordsPerGoroutine)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWeeklyLogCloseWaitsForRotationAndIsIdempotent(t *testing.T) {
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
writeLogTestFile(t, path, "active")
|
|
||||||
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
|
|
||||||
setLogTestModTime(t, path, sunday)
|
|
||||||
started := make(chan struct{})
|
|
||||||
release := make(chan struct{})
|
|
||||||
runtime := weeklyLogRuntime{rename: func(oldPath, newPath string) error {
|
|
||||||
if oldPath == path {
|
|
||||||
close(started)
|
|
||||||
<-release
|
|
||||||
}
|
|
||||||
return os.Rename(oldPath, newPath)
|
|
||||||
}}
|
|
||||||
writer := newLogTestWriter(t, path, sunday, runtime)
|
|
||||||
rotationDone := make(chan struct{})
|
|
||||||
go func() {
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
|
|
||||||
close(rotationDone)
|
|
||||||
}()
|
|
||||||
<-started
|
|
||||||
|
|
||||||
closeDone := make(chan error, 1)
|
|
||||||
go func() { closeDone <- writer.Close() }()
|
|
||||||
select {
|
|
||||||
case err := <-closeDone:
|
|
||||||
t.Fatalf("weeklyLogWriter.Close() during rotation returned %v, want blocked", err)
|
|
||||||
case <-time.After(20 * time.Millisecond):
|
|
||||||
}
|
|
||||||
close(release)
|
|
||||||
<-rotationDone
|
|
||||||
if err := <-closeDone; err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Errorf("second weeklyLogWriter.Close() error = %v, want nil", err)
|
|
||||||
}
|
|
||||||
if _, err := writer.Write([]byte("after close")); !errors.Is(err, os.ErrClosed) {
|
|
||||||
t.Errorf("weeklyLogWriter.Write() after close error = %v, want os.ErrClosed", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestWeeklyLogRotationFailureIsNonRecursiveAndPreservesWrites(t *testing.T) {
|
|
||||||
directory := t.TempDir()
|
|
||||||
path := filepath.Join(directory, "hardlink.log")
|
|
||||||
writeLogTestFile(t, path, "before\n")
|
|
||||||
sunday := time.Date(2026, time.August, 9, 12, 0, 0, 0, time.UTC)
|
|
||||||
setLogTestModTime(t, path, sunday)
|
|
||||||
var diagnostics bytes.Buffer
|
|
||||||
runtime := weeklyLogRuntime{
|
|
||||||
rename: func(oldPath, newPath string) error {
|
|
||||||
if oldPath == path {
|
|
||||||
return errors.New("injected rename failure")
|
|
||||||
}
|
|
||||||
return os.Rename(oldPath, newPath)
|
|
||||||
},
|
|
||||||
diagnostic: &diagnostics,
|
|
||||||
}
|
|
||||||
writer := newLogTestWriter(t, path, sunday, runtime)
|
|
||||||
writer.rotateIfNeeded(time.Date(2026, time.August, 10, 0, 0, 0, 0, time.UTC))
|
|
||||||
if _, err := writer.Write([]byte("after\n")); err != nil {
|
|
||||||
t.Fatalf("weeklyLogWriter.Write() after failed rotation error = %v, want nil", err)
|
|
||||||
}
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
requireLogTestContents(t, path, "before\nafter\n")
|
|
||||||
if !strings.Contains(diagnostics.String(), "weekly_log_rotation_failed") || !strings.Contains(diagnostics.String(), "injected rename failure") {
|
|
||||||
t.Errorf("rotation diagnostic = %q, want action and injected error", diagnostics.String())
|
|
||||||
}
|
|
||||||
if strings.Contains(readLogTestFile(t, path), "weekly log rotation failed") {
|
|
||||||
t.Error("rotation diagnostic recursively entered active log")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestSetupLoggingKeepsFilenameAndJSONFormat(t *testing.T) {
|
|
||||||
logger := log.StandardLogger()
|
|
||||||
previousOutput := logger.Out
|
|
||||||
previousFormatter := logger.Formatter
|
|
||||||
previousLevel := logger.Level
|
|
||||||
t.Cleanup(func() {
|
|
||||||
logger.SetOutput(previousOutput)
|
|
||||||
logger.SetFormatter(previousFormatter)
|
|
||||||
logger.SetLevel(previousLevel)
|
|
||||||
})
|
|
||||||
|
|
||||||
directory := filepath.Join(t.TempDir(), "nested", "logs")
|
|
||||||
writer, err := SetupLogging(directory, "hardlink", "v-test")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
log.WithField("probe", "value").Info("test record")
|
|
||||||
if err := writer.Close(); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
contents := readLogTestFile(t, filepath.Join(directory, "hardlink.log"))
|
|
||||||
for _, want := range []string{`"buildVersion":"v-test"`, `"msg":"Logging initialized"`, `"probe":"value"`, `"msg":"test record"`} {
|
|
||||||
if !strings.Contains(contents, want) {
|
|
||||||
t.Errorf("hardlink.log contents missing %q: %s", want, contents)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if _, err := os.Stat(filepath.Join(directory, "hardlink.1.log")); !errors.Is(err, os.ErrNotExist) {
|
|
||||||
t.Errorf("SetupLogging() created archive during current week: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type logTestTimer struct {
|
|
||||||
duration time.Duration
|
|
||||||
ch chan time.Time
|
|
||||||
stopped atomic.Bool
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *logTestTimer) C() <-chan time.Time {
|
|
||||||
return t.ch
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *logTestTimer) Stop() bool {
|
|
||||||
return !t.stopped.Swap(true)
|
|
||||||
}
|
|
||||||
|
|
||||||
func newLogTestWriter(t *testing.T, path string, now time.Time, overrides weeklyLogRuntime) *weeklyLogWriter {
|
|
||||||
t.Helper()
|
|
||||||
runtime := defaultWeeklyLogRuntime()
|
|
||||||
runtime.now = func() time.Time { return now }
|
|
||||||
if overrides.now != nil {
|
|
||||||
runtime.now = overrides.now
|
|
||||||
}
|
|
||||||
if overrides.newTimer != nil {
|
|
||||||
runtime.newTimer = overrides.newTimer
|
|
||||||
}
|
|
||||||
if overrides.rename != nil {
|
|
||||||
runtime.rename = overrides.rename
|
|
||||||
}
|
|
||||||
if overrides.diagnostic != nil {
|
|
||||||
runtime.diagnostic = overrides.diagnostic
|
|
||||||
}
|
|
||||||
writer, err := newWeeklyLogWriter(path, now.Location(), runtime)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
return writer
|
|
||||||
}
|
|
||||||
|
|
||||||
func weeklyLogTestGeneration(path string, generation int) string {
|
|
||||||
extension := filepath.Ext(path)
|
|
||||||
return fmt.Sprintf("%s.%d%s", strings.TrimSuffix(path, extension), generation, extension)
|
|
||||||
}
|
|
||||||
|
|
||||||
func writeLogTestFile(t *testing.T, path, contents string) {
|
|
||||||
t.Helper()
|
|
||||||
if err := os.WriteFile(path, []byte(contents), 0o600); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func setLogTestModTime(t *testing.T, path string, value time.Time) {
|
|
||||||
t.Helper()
|
|
||||||
if err := os.Chtimes(path, value, value); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func readLogTestFile(t *testing.T, path string) string {
|
|
||||||
t.Helper()
|
|
||||||
contents, err := os.ReadFile(path)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
return string(contents)
|
|
||||||
}
|
|
||||||
|
|
||||||
func requireLogTestContents(t *testing.T, path, want string) {
|
|
||||||
t.Helper()
|
|
||||||
if got := readLogTestFile(t, path); got != want {
|
|
||||||
t.Errorf("%s contents = %q, want %q", filepath.Base(path), got, want)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
var _ io.WriteCloser = (*weeklyLogWriter)(nil)
|
|
||||||
@ -2,15 +2,6 @@
|
|||||||
|
|
||||||
builtVersion is a const in main.go
|
builtVersion is a const in main.go
|
||||||
|
|
||||||
#### v2.1.3 - 30 September 2026
|
|
||||||
feat(logging): add weekly log rotation and retention
|
|
||||||
|
|
||||||
#### 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
|
|
||||||
|
|
||||||
#### v2.1.0 - 28 September 2026
|
#### v2.1.0 - 28 September 2026
|
||||||
fix(dispenser): tolerate unusable AP observations during card preparation
|
fix(dispenser): tolerate unusable AP observations during card preparation
|
||||||
|
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user