hardlink/internal/dispenser/dispenserclient.go

562 lines
14 KiB
Go

// Package dispenser provides a queue-based client (single owner of port).
package dispenser
import (
"context"
"fmt"
"sync"
"time"
log "github.com/sirupsen/logrus"
"github.com/tarm/serial"
)
type cmdType int
const (
cmdStatus cmdType = iota
cmdToEncoder
cmdOutOfMouth
cmdReset
)
type cmdReq struct {
typ cmdType
ctx context.Context
respCh chan cmdResp
}
type cmdResp struct {
status []byte
err error
}
type sequenceTiming struct {
now func() time.Time
wait func(context.Context, time.Duration) error
}
const (
sequencePollInterval = time.Second
sequenceRetryAfter = 6 * time.Second
sequenceResetWait = 2 * time.Second
sequenceTimeout = 16 * time.Second
)
type Client struct {
port serialTransport
reqCh chan cmdReq
done chan struct{}
sequenceTiming sequenceTiming
// status cache
mu sync.RWMutex
lastStatus []byte
lastStatusT time.Time
statusTTL time.Duration
// published "stock/cardwell" cache + callback
lastStockMu sync.RWMutex
lastStock string
onStock func(string)
}
// NewClient starts the worker that owns the serial port.
func NewClient(port *serial.Port, queueSize int) *Client {
if queueSize <= 0 {
queueSize = 16
}
c := &Client{
port: port,
reqCh: make(chan cmdReq, queueSize),
done: make(chan struct{}),
sequenceTiming: sequenceTiming{
now: time.Now,
wait: waitForSequence,
},
statusTTL: defaultStatusTTL,
}
go c.loop()
return c
}
func waitForSequence(ctx context.Context, duration time.Duration) error {
timer := time.NewTimer(duration)
defer timer.Stop()
select {
case <-timer.C:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (c *Client) Close() {
select {
case <-c.done:
return
default:
close(c.done)
}
}
// SetStatusTTL sets the duration for which cached status is considered fresh.
func (c *Client) SetStatusTTL(d time.Duration) {
c.mu.Lock()
c.statusTTL = d
c.mu.Unlock()
}
// OnStockUpdate registers a callback called whenever polling (or status reads) produce a stock status string.
func (c *Client) OnStockUpdate(fn func(string)) {
c.lastStockMu.Lock()
c.onStock = fn
c.lastStockMu.Unlock()
}
// LastStock returns the most recently computed stock/card-well status string.
func (c *Client) LastStock() string {
c.lastStockMu.RLock()
defer c.lastStockMu.RUnlock()
return c.lastStock
}
func (c *Client) setStock(statusBytes []byte) {
stock := stockTake(statusBytes)
c.lastStockMu.Lock()
c.lastStock = stock
fn := c.onStock
c.lastStockMu.Unlock()
// call outside lock
if fn != nil {
fn(stock)
}
}
// StartPolling performs a periodic status refresh.
// It will NOT interrupt commands: it enqueues only when queue is idle.
func (c *Client) StartPolling(interval time.Duration) {
if interval <= 0 {
return
}
go func() {
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-c.done:
return
case <-t.C:
// enqueue only if idle to avoid delaying real commands
if len(c.reqCh) != 0 {
continue
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
_, err := c.CheckStatus(ctx)
if err != nil {
log.Debugf("dispenser polling: %v", err)
}
cancel()
}
}
}()
}
func (c *Client) loop() {
for {
select {
case <-c.done:
return
case req := <-c.reqCh:
c.handle(req)
}
}
}
func (c *Client) handle(req cmdReq) {
select {
case <-req.ctx.Done():
req.respCh <- cmdResp{err: req.ctx.Err()}
return
default:
}
switch req.typ {
case cmdStatus:
st, err := checkDispenserStatus(req.ctx, c.port)
if err == nil && len(st) == 4 {
c.mu.Lock()
c.lastStatus = append([]byte(nil), st...)
c.lastStatusT = time.Now()
c.mu.Unlock()
// publish stock/cardwell
c.setStock(st)
}
req.respCh <- cmdResp{status: st, err: err}
case cmdToEncoder:
err := cardToEncoderPosition(req.ctx, c.port)
// A movement command makes any previously cached position unreliable.
c.invalidateStatusCache()
req.respCh <- cmdResp{err: err}
case cmdReset:
err := resetDispenser(req.ctx, c.port)
c.invalidateStatusCache()
req.respCh <- cmdResp{err: err}
case cmdOutOfMouth:
err := cardOutOfMouth(req.ctx, c.port)
// A movement command makes any previously cached position unreliable.
c.invalidateStatusCache()
req.respCh <- cmdResp{err: err}
default:
req.respCh <- cmdResp{err: fmt.Errorf("unknown command")}
}
}
func (c *Client) do(ctx context.Context, typ cmdType) ([]byte, error) {
rch := make(chan cmdResp, 1)
req := cmdReq{typ: typ, ctx: ctx, respCh: rch}
select {
case c.reqCh <- req:
case <-ctx.Done():
return nil, ctx.Err()
}
select {
case r := <-rch:
return r.status, r.err
case <-ctx.Done():
return nil, ctx.Err()
}
}
// CheckStatus returns cached status if fresh, otherwise enqueues a device status read.
func (c *Client) CheckStatus(ctx context.Context) ([]byte, error) {
c.mu.RLock()
ttl := c.statusTTL
st := append([]byte(nil), c.lastStatus...)
ts := c.lastStatusT
c.mu.RUnlock()
if len(st) == 4 && time.Since(ts) <= ttl {
// even when returning cached, keep stock in sync
c.setStock(st)
return st, nil
}
return c.do(ctx, cmdStatus)
}
func (c *Client) invalidateStatusCache() {
c.mu.Lock()
c.lastStatus = nil
c.lastStatusT = time.Time{}
c.mu.Unlock()
}
func (c *Client) ToEncoder(ctx context.Context) error {
_, err := c.do(ctx, cmdToEncoder)
return err
}
// Reset dispatches RS through the port owner; acceptance does not imply mechanical completion.
func (c *Client) Reset(ctx context.Context) error {
_, err := c.do(ctx, cmdReset)
return err
}
func (c *Client) OutOfMouth(ctx context.Context) error {
_, err := c.do(ctx, cmdOutOfMouth)
return err
}
// --------------------
// Public sequences updated to use Client (queue)
// --------------------
// DispenserPrepare checks status; if empty => ok; else ensure at encoder.
func (c *Client) DispenserPrepare(ctx context.Context) (string, error) {
const funcName = "DispenserPrepare"
stockStatus := ""
status, err := c.CheckStatus(ctx)
if err != nil {
return stockStatus, fmt.Errorf("[%s] check status: %w", funcName, err)
}
logStatus(status)
stockStatus = stockTake(status)
c.setStock(status)
if isCardWellEmpty(status) {
return stockStatus, nil
}
if isAtEncoderPosition(status) {
return stockStatus, nil
}
if err := c.ToEncoder(ctx); err != nil {
return stockStatus, fmt.Errorf("[%s] to encoder: %w", funcName, err)
}
time.Sleep(delay)
status, err = c.CheckStatus(ctx)
if err != nil {
return stockStatus, fmt.Errorf("[%s] re-check status: %w", funcName, err)
}
logStatus(status)
stockStatus = stockTake(status)
c.setStock(status)
return stockStatus, nil
}
func (c *Client) readSequenceStatus(ctx context.Context, operation string) ([]byte, string, error) {
status, err := c.do(ctx, cmdStatus)
if err != nil {
if ctxErr := ctx.Err(); ctxErr != nil {
return nil, "", fmt.Errorf("[%s] read status: %w", operation, ctxErr)
}
return nil, "", fmt.Errorf("[%s] read status: %w", operation, err)
}
stockStatus := ""
if len(status) == 4 {
stockStatus = stockTake(status)
c.setStock(status)
logStatus(status)
}
return status, stockStatus, nil
}
func preparationStatus(operation string, status []byte, allowCombinedFailure bool) (bool, error) {
// A confirmed read sensor takes precedence over every other diagnostic.
if isAtEncoderPosition(status) {
if hasPreparationDiagnostics(status) {
log.Warnf(
"[%s] card confirmed at encoder with dispenser diagnostics: %s raw status: % X",
operation,
statusDescription(status),
status,
)
}
return true, nil
}
if err := validateDispenserStatusData(status); err != nil {
return false, fmt.Errorf("[%s] %w", operation, err)
}
if isCardWellEmpty(status) {
return false, fmt.Errorf("[%s] %w", operation, ErrCardWellEmpty)
}
if isPreparationMoving(status) {
return false, nil
}
// Some dispenser firmware briefly reports 0x32 ("Preparing card fails")
// while the card is still travelling to the encoder. Treat that one
// diagnostic as transient during an active preparation sequence and let
// pollForEncoderPosition decide encoder success or timeout. Do not mask
// independent hard errors reported in the other status bytes.
// Combined command rejection is also recoverable before the one retry.
if status[0] == 0x32 || (status[0] == 0x36 && allowCombinedFailure) {
statusWithoutPrepareFailure := append([]byte(nil), status...)
statusWithoutPrepareFailure[0] = 0x30
if err := dispenserStatusError(statusWithoutPrepareFailure); err != nil {
return false, fmt.Errorf("[%s] %w", operation, err)
}
log.Warnf(
"[%s] transient preparation failure (%s); waiting for encoder position, raw status: % X",
operation,
statusDescription(status),
status,
)
return false, nil
}
if err := dispenserStatusError(status); err != nil {
return false, fmt.Errorf("[%s] %w", operation, err)
}
return false, nil
}
func (c *Client) pollForEncoderPosition(ctx context.Context, operation string) (stockStatus string, resultErr error) {
started := c.sequenceTiming.now()
recoveryAt := started.Add(sequenceRetryAfter)
deadline := started.Add(sequenceTimeout)
ctx, cancel := context.WithTimeout(ctx, sequenceTimeout)
defer cancel()
retried := false
resetAttempted := false
retryAfterReset := false
consecutiveFailures := 0
stage := "initial polling"
var lastStatus []byte
defer func() {
if resultErr != nil {
log.Warnf("[%s] card preparation failed; stage=%s reset_attempted=%t raw status: % X: %v", operation, stage, resetAttempted, lastStatus, resultErr)
}
}()
checkDeadline := func() error {
if err := ctx.Err(); err != nil {
return fmt.Errorf("[%s] %s: %w", operation, stage, err)
}
if !c.sequenceTiming.now().Before(deadline) {
return fmt.Errorf("[%s] %s: timed out after %s", operation, stage, sequenceTimeout)
}
return nil
}
for {
if err := checkDeadline(); err != nil {
return stockStatus, err
}
status, currentStockStatus, err := c.readSequenceStatus(ctx, operation)
stockStatus = currentStockStatus
lastStatus = status
if err != nil {
return stockStatus, err
}
if err := checkDeadline(); err != nil {
return stockStatus, err
}
ready, err := preparationStatus(operation, status, !retried || retryAfterReset)
if err != nil {
return stockStatus, err
}
if ready {
if resetAttempted {
log.Infof("[%s] card preparation recovered after reset", operation)
}
return stockStatus, nil
}
moving := isPreparationMoving(status)
prepareFailed := status[0] == 0x32 || status[0] == 0x36
if prepareFailed && !moving {
consecutiveFailures++
} else {
consecutiveFailures = 0
}
if !moving && retryAfterReset {
stage = "post-reset FC7"
if err := checkDeadline(); err != nil {
return stockStatus, err
}
log.Infof("[%s] reset settle wait completed; retrying card preparation", operation)
if err := c.ToEncoder(ctx); err != nil {
return stockStatus, fmt.Errorf("[%s] post-reset retry command: %w", operation, err)
}
retryAfterReset = false
stage = "polling after reset and retry"
} else if !moving && !retried && !c.sequenceTiming.now().Before(recoveryAt) {
if prepareFailed && consecutiveFailures >= 2 {
stage = "RS recovery"
if err := checkDeadline(); err != nil {
return stockStatus, err
}
retried = true
resetAttempted = true
log.Warnf("[%s] persistent prepare failure; resetting dispenser, raw status: % X", operation, status)
if err := c.Reset(ctx); err != nil {
return stockStatus, fmt.Errorf("[%s] reset command: %w", operation, err)
}
stage = "reset settle"
if err := checkDeadline(); err != nil {
return stockStatus, err
}
log.Infof("[%s] reset dispatched; waiting for mechanical settling", operation)
wait := sequenceResetWait
if remaining := deadline.Sub(c.sequenceTiming.now()); remaining < wait {
wait = remaining
}
if err := c.sequenceTiming.wait(ctx, wait); err != nil {
return stockStatus, fmt.Errorf("[%s] reset settle: %w", operation, err)
}
retryAfterReset = true
stage = "post-reset status"
continue // Always read fresh physical status before the retry.
}
if !prepareFailed {
stage = "FC7 retry"
if err := checkDeadline(); err != nil {
return stockStatus, err
}
retried = true
if err := c.ToEncoder(ctx); err != nil {
return stockStatus, fmt.Errorf("[%s] retry command: %w", operation, err)
}
}
}
if err := checkDeadline(); err != nil {
return stockStatus, err
}
wait := sequencePollInterval
if remaining := deadline.Sub(c.sequenceTiming.now()); remaining < wait {
wait = remaining
}
if err := c.sequenceTiming.wait(ctx, wait); err != nil {
return stockStatus, fmt.Errorf("[%s] %s: %w", operation, stage, err)
}
}
}
func (c *Client) prepareCardAtEncoder(ctx context.Context, operation string) (string, error) {
status, stockStatus, err := c.readSequenceStatus(ctx, operation)
if err != nil {
return stockStatus, err
}
ready, err := preparationStatus(operation, status, true)
if err != nil {
return stockStatus, err
}
if ready {
return stockStatus, nil
}
if err := c.ToEncoder(ctx); err != nil {
return stockStatus, fmt.Errorf("[%s] to encoder: %w", operation, err)
}
return c.pollForEncoderPosition(ctx, operation)
}
// PrepareCurrentCard authoritatively places the card to be encoded at the encoder.
func (c *Client) PrepareCurrentCard(ctx context.Context) (string, error) {
return c.prepareCardAtEncoder(ctx, "PrepareCurrentCard")
}
// DeliverCurrentCard presents the encoded card and confirms command acceptance.
func (c *Client) DeliverCurrentCard(ctx context.Context) (string, error) {
const operation = "DeliverCurrentCard"
if err := c.OutOfMouth(ctx); err != nil {
return "", fmt.Errorf("[%s] out of mouth: %w", operation, err)
}
return "", nil
}
// BeginPrepareNextCard starts moving the next card to the encoder without waiting for readiness.
func (c *Client) BeginPrepareNextCard(ctx context.Context) error {
if err := c.ToEncoder(ctx); err != nil {
return fmt.Errorf("[BeginPrepareNextCard] to encoder: %w", err)
}
return nil
}
// PrepareNextCard places a new card at the encoder for a later issuance attempt.
func (c *Client) PrepareNextCard(ctx context.Context) (string, error) {
return c.prepareCardAtEncoder(ctx, "PrepareNextCard")
}