hardlink/internal/dispenser/dispenserclient.go

732 lines
19 KiB
Go

// Package dispenser provides a queue-based client (single owner of port).
package dispenser
import (
"context"
"errors"
"fmt"
"sync"
"time"
log "github.com/sirupsen/logrus"
"github.com/tarm/serial"
)
type cmdType int
const (
cmdStatus cmdType = iota
cmdToEncoder
cmdOutOfMouth
cmdReset
cmdDeliveryClearance
)
type cmdReq struct {
passive bool
generation uint64
typ cmdType
ctx context.Context
respCh chan cmdResp
}
type cmdResp struct {
deliveryPending bool
status []byte
err error
}
type sequenceTiming struct {
now func() time.Time
wait func(context.Context, time.Duration) error
}
const (
sequencePollInterval = time.Second
sequenceShakeAfter = 3 * time.Second
sequenceResetWait = 2 * time.Second
sequenceTimeout = 32 * time.Second
sequenceUncertainWait = 4 * time.Second
sequenceMaxShakes = 3
deliveryClearanceTimeout = 6 * time.Second
deliveryMinimumWait = 2 * time.Second
)
type Client struct {
activity activityGuard
activityWake chan struct{}
closeOnce sync.Once
port serialTransport
reqCh chan cmdReq
done chan struct{}
// Owned exclusively by the serial worker.
deliveryPending bool
deliveryStarted time.Time
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{
activityWake: make(chan struct{}, 1),
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() {
c.closeOnce.Do(func() {
c.activity.Lock()
c.activity.closed = true
c.activity.Unlock()
c.StopMaintenance()
close(c.done)
})
}
// SetStatusTTL sets the duration for which cached status is considered fresh.
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.
// Passive requests are admitted by the worker only for the captured idle generation.
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
}
generation, idle := c.passiveGeneration()
if !idle {
continue
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
response := c.sendRequest(cmdReq{typ: cmdStatus, ctx: ctx, passive: true, generation: generation})
err := response.err
if err != nil {
log.Debugf("dispenser polling: %v", err)
}
cancel()
}
}
}()
}
func (c *Client) loop() {
timer := time.NewTimer(time.Hour)
defer timer.Stop()
for {
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
var tick <-chan time.Time
if delay, enabled := c.maintenanceDelay(); enabled {
timer.Reset(delay)
tick = timer.C
}
select {
case <-c.done:
return
case <-c.activityWake:
case req := <-c.reqCh:
c.handle(req)
case <-tick:
c.maintainIdleCard()
}
}
}
func (c *Client) handle(req cmdReq) {
select {
case <-req.ctx.Done():
req.respCh <- cmdResp{err: req.ctx.Err()}
return
default:
}
if req.passive {
if !c.admitPassive(req.ctx, req.generation) {
req.respCh <- cmdResp{}
return
}
c.mu.RLock()
st := append([]byte(nil), c.lastStatus...)
fresh := len(st) == 4 && time.Since(c.lastStatusT) <= c.statusTTL
c.mu.RUnlock()
if fresh {
c.setStock(st)
req.respCh <- cmdResp{status: st}
return
}
}
switch req.typ {
case cmdStatus:
st, err := c.readWorkerStatus(req.ctx)
req.respCh <- cmdResp{status: st, err: err, deliveryPending: c.deliveryPending}
case cmdDeliveryClearance:
st, err := c.waitWorkerDeliveryClearance(req.ctx)
req.respCh <- cmdResp{status: st, err: err, deliveryPending: c.deliveryPending}
case cmdToEncoder:
if c.deliveryPending {
if _, err := c.waitWorkerDeliveryClearance(req.ctx); err != nil {
req.respCh <- cmdResp{err: err}
return
}
}
err := cardToEncoderPosition(req.ctx, c.port)
log.Infof("FC7 dispatch finished; dispatched=%t error=%v", err == nil, err)
c.invalidateStatusCache()
req.respCh <- cmdResp{err: err}
case cmdReset:
err := resetDispenser(req.ctx, c.port)
log.Infof("RS dispatch finished; dispatched=%t error=%v", err == nil, err)
c.invalidateStatusCache()
req.respCh <- cmdResp{err: err}
case cmdOutOfMouth:
attempted, err := cardOutOfMouth(req.ctx, c.port)
log.Infof("FC0 dispatch finished; ENQ_attempted=%t error=%v", attempted, err)
if attempted {
c.deliveryPending = true
c.deliveryStarted = c.now()
log.Info("delivery ENQ attempted; awaiting mechanical clearance")
}
// A movement command makes any previously cached position unreliable.
c.invalidateStatusCache()
req.respCh <- cmdResp{err: err}
default:
req.respCh <- cmdResp{err: fmt.Errorf("unknown command")}
}
}
// deliveryClearance deliberately does not apply encoder-success precedence.
func deliveryClearance(status []byte) (bool, error) {
class, err := classifyPreparationStatus(status)
if err != nil {
return false, err
}
if class == positionWellEmpty {
return false, ErrCardWellEmpty
}
if isPreparationMoving(status) || status[1] == 0x34 {
return false, nil
}
return status[3] == 0x30 || status[3] == 0x34, nil
}
func (c *Client) now() time.Time {
if c.sequenceTiming.now != nil {
return c.sequenceTiming.now()
}
return time.Now()
}
// readWorkerStatus is called only by the serial worker and always reads fresh AP.
func (c *Client) readWorkerStatus(ctx context.Context) ([]byte, error) {
st, err := checkDispenserStatus(ctx, c.port)
if err != nil {
return st, err
}
if err := ctx.Err(); err != nil {
return st, err
}
if c.deliveryPending {
class, _ := classifyPreparationStatus(st)
log.Infof("delivery clearance AP; class=%s elapsed=%s raw status: % X", class, c.now().Sub(c.deliveryStarted), st)
clear, clearanceErr := deliveryClearance(st)
if clearanceErr != nil {
log.Warnf("delivery clearance failed: %v; %s raw status: % X", clearanceErr, statusDescription(st), st)
} else if clear && c.now().Sub(c.deliveryStarted) >= deliveryMinimumWait {
c.deliveryPending = false
log.Infof("previous delivery cleared after %s; FC7 permitted", c.now().Sub(c.deliveryStarted))
}
}
if len(st) == 4 {
c.mu.Lock()
c.lastStatus = append([]byte(nil), st...)
c.lastStatusT = time.Now()
c.mu.Unlock()
c.setStock(st)
}
return st, nil
}
func (c *Client) do(ctx context.Context, typ cmdType) ([]byte, error) {
r := c.doResponse(ctx, typ)
return r.status, r.err
}
func (c *Client) doResponse(ctx context.Context, typ cmdType) cmdResp {
return c.sendRequest(cmdReq{typ: typ, ctx: ctx})
}
func (c *Client) sendRequest(req cmdReq) cmdResp {
if err := req.ctx.Err(); err != nil {
return cmdResp{err: err}
}
req.respCh = make(chan cmdResp, 1)
select {
case <-c.done:
return cmdResp{err: context.Canceled}
default:
}
select {
case c.reqCh <- req:
case <-c.done:
return cmdResp{err: context.Canceled}
case <-req.ctx.Done():
return cmdResp{err: req.ctx.Err()}
}
select {
case r := <-req.respCh:
return r
case <-c.done:
return cmdResp{err: context.Canceled}
case <-req.ctx.Done():
return cmdResp{err: req.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) {
response := c.doResponse(ctx, cmdStatus)
if err := ctx.Err(); err != nil {
return nil, "", err
}
if response.deliveryPending {
return c.waitForDeliveryClearance(ctx, operation)
}
status := usableObservation(response.status, response.err, operation)
stockStatus := ""
if len(status) == 4 {
stockStatus = stockTake(status)
c.setStock(status)
logStatus(status)
}
return status, stockStatus, nil
}
// waitWorkerDeliveryClearance runs only inside the serial worker.
// An unusable observation never becomes a fabricated physical position.
func (c *Client) waitWorkerDeliveryClearance(parent context.Context) ([]byte, error) {
ctx, cancel := context.WithTimeout(parent, deliveryClearanceTimeout)
defer cancel()
deadline := c.now().Add(deliveryClearanceTimeout)
var latest []byte
for {
if err := parent.Err(); err != nil {
return nil, err
}
if !c.now().Before(deadline) || ctx.Err() != nil {
class, err := classifyPreparationStatus(latest)
if err == nil && class == positionWellEmpty {
return latest, ErrCardWellEmpty
}
if err == nil && (isPreparationMoving(latest) || latest[1] == 0x34) {
return latest, context.DeadlineExceeded
}
c.deliveryPending = false
log.Warn("delivery clearance assumed after bounded observation fallback")
return latest, nil
}
st, err := c.readWorkerStatus(ctx)
if parent.Err() != nil {
return nil, parent.Err()
}
latest = usableObservation(st, err, "delivery clearance")
if latest != nil && latest[3] == 0x38 {
return latest, ErrCardWellEmpty
}
if !c.deliveryPending {
return latest, nil
}
remaining := deadline.Sub(c.now())
if remaining <= 0 || ctx.Err() != nil {
continue
}
wait := sequencePollInterval
if remaining < wait {
wait = remaining
}
if err := c.sequenceTiming.wait(ctx, wait); err != nil && parent.Err() != nil {
return nil, parent.Err()
}
}
}
// usableObservation preserves strict validation while discarding unusable telemetry.
func usableObservation(status []byte, err error, operation string) []byte {
if err == nil {
_, err = classifyPreparationStatus(status)
}
if err != nil {
log.Warnf("[%s] unusable AP observation; raw status: % X error=%v", operation, status, err)
return nil
}
return status
}
// waitForDeliveryClearance uses a worker request; no worker recursively enqueues.
func (c *Client) waitForDeliveryClearance(ctx context.Context, operation string) ([]byte, string, error) {
r := c.doResponse(ctx, cmdDeliveryClearance)
stock := ""
if len(r.status) == 4 {
stock = stockTake(r.status)
}
if r.err != nil {
return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, r.err)
}
return r.status, stock, nil
}
func (c *Client) prepareCardAtEncoder(parent context.Context, operation string) (stock string, resultErr error) {
status, stock, err := c.waitForDeliveryClearance(parent, operation)
if err != nil {
return stock, err
}
status = usableObservation(status, nil, operation)
started := c.now()
deadline := started.Add(sequenceTimeout)
ctx, cancel := context.WithTimeoutCause(parent, sequenceTimeout, ErrPreparationExhausted)
defer cancel()
shakes := 0
stage := "initial status"
var lastFC7, uncertainSince time.Time
var previous preparationClass
defer func() {
log.Infof("[%s] preparation finished; stage=%s shakes=%d elapsed=%s raw status: % X error=%v", operation, stage, shakes, c.now().Sub(started), status, resultErr)
}()
checkDeadline := func() error {
if err := parent.Err(); err != nil {
return err
}
if !c.now().Before(deadline) || errors.Is(context.Cause(ctx), ErrPreparationExhausted) {
return ErrPreparationExhausted
}
return ctx.Err()
}
// Only observation exhaustion can grant a deadline handoff. Command and
// reset-settle failures never pass through this policy.
observationResult := func(err error) error {
if parent.Err() != nil {
return parent.Err()
}
if !errors.Is(err, ErrPreparationExhausted) || lastFC7.IsZero() {
return err
}
class, _ := classifyPreparationStatus(status)
switch class {
case positionWellEmpty:
return ErrCardWellEmpty
case positionNoCard:
return err
default:
stage = "observation deadline encoder handoff"
return nil
}
}
wait := func(duration time.Duration) error {
if err := checkDeadline(); err != nil {
return err
}
if remaining := deadline.Sub(c.now()); remaining < duration {
duration = remaining
}
if err := c.sequenceTiming.wait(ctx, duration); err != nil {
if deadlineErr := checkDeadline(); deadlineErr != nil {
return deadlineErr
}
return err
}
return checkDeadline()
}
sendFC7 := func() error {
stage = "FC7 dispatch"
if err := checkDeadline(); err != nil {
return err
}
err := c.ToEncoder(ctx)
if deadlineErr := checkDeadline(); deadlineErr != nil {
return deadlineErr
}
if err != nil {
return fmt.Errorf("[%s] FC7 dispatch: %w", operation, err)
}
lastFC7 = c.now()
log.Infof("[%s] FC7 dispatched; shake=%d elapsed=%s", operation, shakes, lastFC7.Sub(started))
return nil
}
for {
if err := checkDeadline(); err != nil {
return stock, observationResult(err)
}
class, _ := classifyPreparationStatus(status)
// No class means unusable telemetry, sharing the uncertainty timer.
log.Infof("[%s] fresh AP; class=%s previous=%s elapsed=%s shake=%d raw status: % X", operation, class, previous, c.now().Sub(started), shakes, status)
previous = class
switch class {
case encoderConfirmed:
stage = "encoder sensor handoff"
return stock, observationResult(checkDeadline())
case positionWellEmpty:
stage = "empty"
return stock, ErrCardWellEmpty
}
if lastFC7.IsZero() {
if err := sendFC7(); err != nil {
return stock, err
}
} else if class == positionNoCard && c.now().Sub(lastFC7) >= sequenceShakeAfter {
uncertainSince = time.Time{}
if shakes == sequenceMaxShakes {
stage = "three shakes exhausted"
return stock, ErrPreparationExhausted
}
shakes++
stage = "RS dispatch"
if err := checkDeadline(); err != nil {
return stock, err
}
err := c.Reset(ctx)
if deadlineErr := checkDeadline(); deadlineErr != nil {
return stock, deadlineErr
}
if err != nil {
return stock, fmt.Errorf("[%s] RS dispatch: %w", operation, err)
}
stage = "reset settling"
log.Infof("[%s] RS dispatched; shake=%d settling=%s", operation, shakes, sequenceResetWait)
if err := wait(sequenceResetWait); err != nil {
return stock, err
}
log.Infof("[%s] reset settle wait completed; shake=%d", operation, shakes)
if err := sendFC7(); err != nil {
return stock, err
}
} else {
if class == positionUncertain || class == "" {
if uncertainSince.IsZero() {
uncertainSince = c.now()
}
if c.now().Sub(uncertainSince) >= sequenceUncertainWait {
stage = "uncertain position handoff"
return stock, observationResult(checkDeadline())
}
} else {
uncertainSince = time.Time{}
}
stage = "polling"
if err := wait(sequencePollInterval); err != nil {
return stock, observationResult(err)
}
}
stage = "fresh AP"
if err := checkDeadline(); err != nil {
return stock, observationResult(err)
}
status, stock, err = c.readSequenceStatus(ctx, operation)
if deadlineErr := checkDeadline(); deadlineErr != nil {
return stock, observationResult(deadlineErr)
}
if err != nil {
return stock, err
}
}
}
// PrepareCurrentCard grants one encoder opportunity after bounded physical preparation.
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 waits for worker-owned clearance and dispatches one FC7.
// It does not wait for encoder 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 runs the same bounded physical preparation for a later issuance attempt.
func (c *Client) PrepareNextCard(ctx context.Context) (string, error) {
return c.prepareCardAtEncoder(ctx, "PrepareNextCard")
}