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