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