181 lines
5.0 KiB
Go
181 lines
5.0 KiB
Go
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")
|
|
}
|