// 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 ) type cmdReq struct { typ cmdType ctx context.Context respCh chan cmdResp } var errDeliveryPending = errors.New("next-card preparation deferred: previous delivery is not clear") 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 sequenceRetryAfter = 6 * time.Second sequenceResetWait = 2 * time.Second sequenceTimeout = 16 * time.Second ) type Client struct { 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{ 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 := c.readWorkerStatus(req.ctx) req.respCh <- cmdResp{status: st, err: err, deliveryPending: c.deliveryPending} case cmdToEncoder: if c.deliveryPending { // Stay inside port ownership: never enqueue a request from the worker. st, err := c.readWorkerStatus(req.ctx) if err == nil { _, err = deliveryClearance(st) } if err == nil && c.deliveryPending { err = errDeliveryPending log.Debugf("next-card preparation deferred; raw status: % X", st) } if err != nil { req.respCh <- cmdResp{err: err} return } } err := cardToEncoderPosition(req.ctx, c.port) c.invalidateStatusCache() req.respCh <- cmdResp{err: err} case cmdReset: err := resetDispenser(req.ctx, c.port) c.invalidateStatusCache() req.respCh <- cmdResp{err: err} case cmdOutOfMouth: attempted, err := cardOutOfMouth(req.ctx, c.port) if attempted { c.deliveryPending = true c.deliveryStarted = time.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) { if err := validateDispenserStatusData(status); err != nil { return false, err } if isCardWellEmpty(status) { return false, ErrCardWellEmpty } if isPreparationMoving(status) { return false, nil } if err := dispenserStatusError(status); err != nil { return false, err } return status[0] == 0x30 && status[1] == 0x30 && status[3] == 0x30, nil } // 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 c.deliveryPending { 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.deliveryPending = false log.Infof("previous delivery cleared after %s; FC7 permitted", time.Since(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 { rch := make(chan cmdResp, 1) req := cmdReq{typ: typ, ctx: ctx, respCh: rch} select { case c.reqCh <- req: case <-ctx.Done(): return cmdResp{err: ctx.Err()} } select { case r := <-rch: return r case <-ctx.Done(): return cmdResp{err: 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) status, err := response.status, response.err if err == nil && response.deliveryPending { err = errDeliveryPending } 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) } } } // waitForDeliveryClearance reuses the first clear AP as the initial preparation sample. // Its bounded window is separate from recovery, which starts after initial FC7. func (c *Client) waitForDeliveryClearance(ctx context.Context, operation string) (status []byte, stock string, resultErr error) { ctx, cancel := context.WithTimeout(ctx, sequenceTimeout) defer cancel() deadline := c.sequenceTiming.now().Add(sequenceTimeout) var lastStatus []byte defer func() { if resultErr != nil { log.Warnf("[%s] delivery clearance stopped: %v; %s raw status: % X", operation, resultErr, statusDescription(lastStatus), lastStatus) } }() for { if err := ctx.Err(); err != nil { return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, err) } if !c.sequenceTiming.now().Before(deadline) { return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, context.DeadlineExceeded) } r := c.doResponse(ctx, cmdStatus) if len(r.status) > 0 { lastStatus = r.status } if r.err != nil { return nil, stock, fmt.Errorf("[%s] delivery clearance status: %w", operation, r.err) } if len(r.status) == 4 { stock = stockTake(r.status) c.setStock(r.status) } if err := ctx.Err(); err != nil { return nil, stock, err } if !c.sequenceTiming.now().Before(deadline) { return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, context.DeadlineExceeded) } if !r.deliveryPending { return r.status, stock, nil } if _, err := deliveryClearance(r.status); err != nil { return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, err) } wait := sequencePollInterval if remaining := deadline.Sub(c.sequenceTiming.now()); remaining < wait { wait = remaining } if err := c.sequenceTiming.wait(ctx, wait); err != nil { return nil, stock, fmt.Errorf("[%s] delivery clearance: %w", operation, err) } } } func (c *Client) prepareCardAtEncoder(ctx context.Context, operation string) (string, error) { status, stockStatus, err := c.waitForDeliveryClearance(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") }