Files

196 lines
5.0 KiB
Go

package ctl
import (
"sync"
"time"
)
// CommandTimeout is the deadline on a pending sid/cid correlation, measured
// from registration time (milestone plan §10).
const CommandTimeout = 30 * time.Second
// Trace is one correlation event. Phase is "ack" for the stanza receipt and
// "result" for the ctl command result or the terminal failure of a command
// that expected one. Ret carries the wire ret value or the failure reason.
type Trace struct {
SID string `json:"sid"`
CID string `json:"cid,omitempty"`
Command string `json:"command"`
Phase string `json:"phase"`
Ret string `json:"ret,omitempty"`
Errno *string `json:"errno,omitempty"`
At time.Time `json:"timestamp"`
}
// TraceSink receives completed traces. A nil sink discards them.
type TraceSink func(serial string, trace Trace)
type record struct {
cid string
command string
expectResult bool
sids map[string]struct{}
lastSID string
registered time.Time
}
// Correlator tracks outstanding sid and cid waiters. Registering a cid that is
// still outstanding attaches the new sid to the single existing record; every
// sid may then ack but the cid completes exactly once. A cid is reusable after
// completion.
type Correlator struct {
mu sync.Mutex
now func() time.Time
bySID map[string]*record
byCID map[string]*record
records map[*record]struct{}
}
// NewCorrelator creates a Correlator. A nil now uses time.Now.
func NewCorrelator(now func() time.Time) *Correlator {
if now == nil {
now = time.Now
}
return &Correlator{
now: now,
bySID: map[string]*record{},
byCID: map[string]*record{},
records: map[*record]struct{}{},
}
}
// Register adds sid as an ack waiter and, when cid is nonempty, a result
// waiter. A duplicate registration of a live cid does not allocate a second
// waiter; the sid joins the existing record.
func (c *Correlator) Register(sid, cid, command string, expectResult bool) {
c.mu.Lock()
defer c.mu.Unlock()
if cid != "" {
if r := c.byCID[cid]; r != nil {
r.sids[sid] = struct{}{}
r.lastSID = sid
c.bySID[sid] = r
return
}
}
r := &record{
cid: cid,
command: command,
expectResult: expectResult,
sids: map[string]struct{}{sid: {}},
lastSID: sid,
registered: c.now(),
}
c.records[r] = struct{}{}
c.bySID[sid] = r
if cid != "" {
c.byCID[cid] = r
}
}
// CompleteSID completes the stanza ack phase for sid. A sid-only record such
// as Move is fully completed by its ack.
func (c *Correlator) CompleteSID(sid string) (Trace, bool) {
c.mu.Lock()
defer c.mu.Unlock()
r := c.bySID[sid]
if r == nil {
return Trace{}, false
}
delete(c.bySID, sid)
delete(r.sids, sid)
tr := Trace{SID: sid, CID: r.cid, Command: r.command, Phase: "ack", At: c.now()}
if !r.expectResult && len(r.sids) == 0 {
c.remove(r)
}
return tr, true
}
// CompleteCID completes the result phase for cid. It returns true only the
// first time; the cid is free for reuse afterwards.
func (c *Correlator) CompleteCID(cid string) (Trace, bool) {
c.mu.Lock()
defer c.mu.Unlock()
r := c.byCID[cid]
if r == nil {
return Trace{}, false
}
c.remove(r)
return Trace{SID: r.lastSID, CID: r.cid, Command: r.command, Phase: "result", At: c.now()}, true
}
// FailGeneration emits one terminal trace per outstanding logical command and
// clears every waiter. The trace phase is "result" for commands expecting a
// ctl result and "ack" for sid-only commands.
func (c *Correlator) FailGeneration(reason string) []Trace {
c.mu.Lock()
defer c.mu.Unlock()
at := c.now()
traces := make([]Trace, 0, len(c.records))
for r := range c.records {
traces = append(traces, c.terminal(r, reason, at))
}
c.bySID = map[string]*record{}
c.byCID = map[string]*record{}
c.records = map[*record]struct{}{}
return traces
}
// NextDeadline returns the earliest expiry across outstanding records: each
// record's registration time plus CommandTimeout. It reports false when no
// records are outstanding.
func (c *Correlator) NextDeadline() (time.Time, bool) {
c.mu.Lock()
defer c.mu.Unlock()
var min time.Time
for r := range c.records {
d := r.registered.Add(CommandTimeout)
if min.IsZero() || d.Before(min) {
min = d
}
}
if min.IsZero() {
return time.Time{}, false
}
return min, true
}
// Expire emits one terminal trace per record registered at least
// CommandTimeout before now and removes it.
func (c *Correlator) Expire(now time.Time) []Trace {
c.mu.Lock()
defer c.mu.Unlock()
var traces []Trace
for r := range c.records {
if now.Sub(r.registered) >= CommandTimeout {
traces = append(traces, c.terminal(r, "timeout", now))
c.remove(r)
}
}
return traces
}
func (c *Correlator) terminal(r *record, reason string, at time.Time) Trace {
phase := "ack"
if r.expectResult {
phase = "result"
}
return Trace{SID: r.lastSID, CID: r.cid, Command: r.command, Phase: phase, Ret: reason, At: at}
}
func (c *Correlator) remove(r *record) {
delete(c.records, r)
for sid := range r.sids {
delete(c.bySID, sid)
}
if r.cid != "" {
delete(c.byCID, r.cid)
}
}