// 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 { 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 { 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 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 { 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) 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") }