3f84f82c3d
gates / gates (push) Successful in 56s
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
349 lines
12 KiB
Go
349 lines
12 KiB
Go
package report
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"sort"
|
|
"strconv"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Operator actions (R-314 / R-279 / R-177, `09` §3 decision 185 — D1, design option A of
|
|
// felhom.eu/documentation/audits/day-2026-10-08/design-R-314-279-177.md).
|
|
//
|
|
// THE DOOR IS THE REPORT REPLY, the shape selftail.go proved for log pulls: the operator presses a
|
|
// button on the hub, the hub stores a row and wakes this box's wait channel, the out-of-cycle report's
|
|
// reply carries `operator_actions: [{id, action, arg}]`, and the RUNNING controller acts on it — inside
|
|
// the process that owns the state, so there is no second process and no lost update (the CLI levers'
|
|
// problem, main.go). The result rides the next report as `operator_action_results: [{id, outcome,
|
|
// message}]`; the hub lists an action until its result has arrived, then stops.
|
|
//
|
|
// THE LIST IS CLOSED. Four actions, and run_job names one of four jobs. Nothing here deletes data,
|
|
// starts a countdown or shortens one — `03` §4 asks for a signing key only to destroy or overwrite the
|
|
// only copy, and none of these does. An unknown action, an unknown job or an argument out of range is
|
|
// answered `refused` and NOTHING is called. Pinned by TestOpActions_ClosedList (which fails when an
|
|
// entry is added, so a new action is a decision, not an edit) and TestOpActions_UnknownIsRefused.
|
|
//
|
|
// ONCE PER ID. An id is acted on the first time a reply carries it; every later reply that still lists
|
|
// it (the result has not reached the hub yet) only re-sends the result. The record is in memory: after a
|
|
// controller restart an action the hub still lists runs again. That is safe by construction — an
|
|
// off-site run or a check repeated adds nothing harmful, a second stop finds nothing to stop, and a
|
|
// second extension counts from the new „now" and can never shorten (ExtendAbandon refuses that).
|
|
|
|
const (
|
|
OpOffsiteBackupNow = "offsite_backup_now"
|
|
OpAbandonStop = "abandon_stop"
|
|
OpAbandonExtend = "abandon_extend"
|
|
OpRunJob = "run_job"
|
|
|
|
OutcomeDone = "done"
|
|
OutcomeRefused = "refused"
|
|
OutcomeFailed = "failed"
|
|
|
|
// AbandonExtendMinDays / AbandonExtendMaxDays bound abandon_extend's argument (decision 185).
|
|
AbandonExtendMinDays = 1
|
|
AbandonExtendMaxDays = 30
|
|
|
|
// opActionsPerReply caps how many NEW actions one reply may start — a defence against a hub bug
|
|
// flooding the box, never reached by a person pressing buttons.
|
|
opActionsPerReply = 10
|
|
// opActionForgetAfter drops a finished, no-longer-listed id from memory.
|
|
opActionForgetAfter = 24 * time.Hour
|
|
// offsiteRunTimeout matches the household's own „back up now" (offbox_handlers.go).
|
|
offsiteRunTimeout = 3 * time.Hour
|
|
)
|
|
|
|
// operatorActionNames is THE closed list. Adding an entry fails TestOpActions_ClosedList by design.
|
|
var operatorActionNames = []string{OpOffsiteBackupNow, OpAbandonStop, OpAbandonExtend, OpRunJob}
|
|
|
|
// operatorJobNames is the fixed set run_job may name — the scheduler's own job names (main.go).
|
|
var operatorJobNames = []string{"fill-watch", "offsite-integrity", "offsite-proof", "disk-health-check"}
|
|
|
|
// OperatorActionNames returns the closed action list (sorted copy).
|
|
func OperatorActionNames() []string { return sortedCopy(operatorActionNames) }
|
|
|
|
// OperatorJobNames returns the fixed run_job set (sorted copy).
|
|
func OperatorJobNames() []string { return sortedCopy(operatorJobNames) }
|
|
|
|
func sortedCopy(in []string) []string {
|
|
out := append([]string(nil), in...)
|
|
sort.Strings(out)
|
|
return out
|
|
}
|
|
|
|
func contains(list []string, v string) bool {
|
|
for _, x := range list {
|
|
if x == v {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// ErrActionRefused, wrapped by a handler, turns its error into a `refused` outcome (the action could
|
|
// not START — e.g. an off-site run already in flight) instead of `failed` (it started and went wrong).
|
|
var ErrActionRefused = errors.New("refused")
|
|
|
|
// OperatorAction is one entry of the reply's operator_actions list.
|
|
type OperatorAction struct {
|
|
ID int64 `json:"id"`
|
|
Action string `json:"action"`
|
|
Arg string `json:"arg,omitempty"`
|
|
}
|
|
|
|
// OperatorActionResult is one entry of the report's operator_action_results list.
|
|
type OperatorActionResult struct {
|
|
ID int64 `json:"id"`
|
|
Outcome string `json:"outcome"` // done | refused | failed
|
|
Message string `json:"message"`
|
|
}
|
|
|
|
// OperatorActionHandlers are the four doors, wired by main.go. A nil handler answers `refused`
|
|
// („not available on this box").
|
|
type OperatorActionHandlers struct {
|
|
// OffsiteBackupNow returns a refusal sentence when the run may not start (the household's
|
|
// „back up now" checks), else the run itself, which the executor starts in its own goroutine.
|
|
OffsiteBackupNow func() (refusal string, run func(ctx context.Context) error)
|
|
AbandonStop func() error
|
|
AbandonExtend func(days int) (time.Time, error)
|
|
// RunJob is Scheduler.RunNow: an error = not started (unknown/running); the channel = the job's end.
|
|
RunJob func(name string) (<-chan error, error)
|
|
// ResultReady sends an out-of-cycle report so a result reaches the hub in seconds (optional).
|
|
ResultReady func()
|
|
}
|
|
|
|
type opEntry struct {
|
|
finished bool
|
|
result OperatorActionResult
|
|
finishedAt time.Time
|
|
}
|
|
|
|
// OperatorActions is the executor. One per process.
|
|
type OperatorActions struct {
|
|
ctx context.Context
|
|
h OperatorActionHandlers
|
|
logger *log.Logger
|
|
now func() time.Time
|
|
|
|
mu sync.Mutex
|
|
seen map[int64]*opEntry
|
|
listed map[int64]bool // ids the LAST reply carried — results are sent only for these
|
|
running sync.WaitGroup // tests wait on async actions
|
|
}
|
|
|
|
// NewOperatorActions builds the executor. ctx bounds the asynchronous actions.
|
|
func NewOperatorActions(ctx context.Context, h OperatorActionHandlers, logger *log.Logger) *OperatorActions {
|
|
return &OperatorActions{ctx: ctx, h: h, logger: logger, now: time.Now,
|
|
seen: map[int64]*opEntry{}, listed: map[int64]bool{}}
|
|
}
|
|
|
|
func (o *OperatorActions) logf(format string, args ...interface{}) {
|
|
if o.logger != nil {
|
|
o.logger.Printf(format, args...)
|
|
}
|
|
}
|
|
|
|
// Reconcile takes one reply's operator_actions list: new ids are acted on (once), every listed id is
|
|
// remembered so its result keeps riding the reports until the hub stops listing it.
|
|
func (o *OperatorActions) Reconcile(list []OperatorAction) {
|
|
o.mu.Lock()
|
|
o.listed = make(map[int64]bool, len(list))
|
|
var fresh []OperatorAction
|
|
for _, a := range list {
|
|
o.listed[a.ID] = true
|
|
if _, ok := o.seen[a.ID]; ok {
|
|
continue // once per id: already acted on (or acting) — only its result is re-sent
|
|
}
|
|
if len(fresh) >= opActionsPerReply {
|
|
continue // not marked seen: the next reply offers it again
|
|
}
|
|
o.seen[a.ID] = &opEntry{}
|
|
fresh = append(fresh, a)
|
|
}
|
|
// Forget finished ids the hub no longer lists, after a day (memory bound; the hub never re-lists
|
|
// an id whose result it holds).
|
|
cutoff := o.now().Add(-opActionForgetAfter)
|
|
for id, e := range o.seen {
|
|
if e.finished && !o.listed[id] && e.finishedAt.Before(cutoff) {
|
|
delete(o.seen, id)
|
|
}
|
|
}
|
|
o.mu.Unlock()
|
|
|
|
if len(list) > 0 {
|
|
o.logf("[DEBUG] [opaction] reply listed %d action(s), %d new", len(list), len(fresh))
|
|
}
|
|
for _, a := range fresh {
|
|
o.dispatch(a)
|
|
}
|
|
}
|
|
|
|
// Results returns the finished results the hub is still waiting for (ids the last reply listed).
|
|
func (o *OperatorActions) Results() []OperatorActionResult {
|
|
o.mu.Lock()
|
|
defer o.mu.Unlock()
|
|
var out []OperatorActionResult
|
|
for id, e := range o.seen {
|
|
if e.finished && o.listed[id] {
|
|
out = append(out, e.result)
|
|
}
|
|
}
|
|
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
|
|
return out
|
|
}
|
|
|
|
func (o *OperatorActions) finish(a OperatorAction, outcome, msg string) {
|
|
o.mu.Lock()
|
|
if e, ok := o.seen[a.ID]; ok {
|
|
e.finished = true
|
|
e.finishedAt = o.now()
|
|
e.result = OperatorActionResult{ID: a.ID, Outcome: outcome, Message: msg}
|
|
}
|
|
o.mu.Unlock()
|
|
o.logf("[INFO] [opaction] operator action #%d %s: %s — %s", a.ID, a.Action, outcome, msg)
|
|
if o.h.ResultReady != nil {
|
|
o.h.ResultReady()
|
|
}
|
|
}
|
|
|
|
func outcomeOf(err error) string {
|
|
switch {
|
|
case err == nil:
|
|
return OutcomeDone
|
|
case errors.Is(err, ErrActionRefused):
|
|
return OutcomeRefused
|
|
default:
|
|
return OutcomeFailed
|
|
}
|
|
}
|
|
|
|
// dispatch validates against the closed list and acts. Validation happens BEFORE any handler is
|
|
// touched: a refusal calls nothing.
|
|
func (o *OperatorActions) dispatch(a OperatorAction) {
|
|
// Customer-visible transparency (the box's own log, like an operator log pull): never silent.
|
|
o.logf("[INFO] [opaction] operator action #%d received: %s", a.ID, a.Action)
|
|
switch a.Action {
|
|
case OpOffsiteBackupNow:
|
|
if a.Arg != "" {
|
|
o.finish(a, OutcomeRefused, "this action takes no argument")
|
|
return
|
|
}
|
|
if o.h.OffsiteBackupNow == nil {
|
|
o.finish(a, OutcomeRefused, "off-site backup is not available on this box")
|
|
return
|
|
}
|
|
refusal, run := o.h.OffsiteBackupNow()
|
|
if refusal != "" || run == nil {
|
|
if refusal == "" {
|
|
refusal = "the off-site backup could not be started"
|
|
}
|
|
o.finish(a, OutcomeRefused, refusal)
|
|
return
|
|
}
|
|
o.running.Add(1)
|
|
go func() {
|
|
defer o.running.Done()
|
|
ctx, cancel := context.WithTimeout(o.ctx, offsiteRunTimeout)
|
|
defer cancel()
|
|
err := run(ctx)
|
|
msg := "the off-site backup finished"
|
|
if err != nil {
|
|
msg = "the off-site backup: " + err.Error()
|
|
}
|
|
o.finish(a, outcomeOf(err), msg)
|
|
}()
|
|
case OpAbandonStop:
|
|
if a.Arg != "" {
|
|
o.finish(a, OutcomeRefused, "this action takes no argument")
|
|
return
|
|
}
|
|
if o.h.AbandonStop == nil {
|
|
o.finish(a, OutcomeRefused, "not available on this box")
|
|
return
|
|
}
|
|
if err := o.h.AbandonStop(); err != nil {
|
|
o.finish(a, outcomeOf(err), err.Error())
|
|
return
|
|
}
|
|
o.finish(a, OutcomeDone, "the deletion countdown is stopped; the set-aside history is kept")
|
|
case OpAbandonExtend:
|
|
days, err := strconv.Atoi(a.Arg)
|
|
if err != nil || days < AbandonExtendMinDays || days > AbandonExtendMaxDays {
|
|
o.finish(a, OutcomeRefused, fmt.Sprintf("the number of days must be %d-%d", AbandonExtendMinDays, AbandonExtendMaxDays))
|
|
return
|
|
}
|
|
if o.h.AbandonExtend == nil {
|
|
o.finish(a, OutcomeRefused, "not available on this box")
|
|
return
|
|
}
|
|
due, err := o.h.AbandonExtend(days)
|
|
if err != nil {
|
|
o.finish(a, outcomeOf(err), err.Error())
|
|
return
|
|
}
|
|
o.finish(a, OutcomeDone, "the set-aside history is now deleted on "+due.Format("2006-01-02"))
|
|
case OpRunJob:
|
|
if !contains(operatorJobNames, a.Arg) {
|
|
o.finish(a, OutcomeRefused, "unknown job name")
|
|
return
|
|
}
|
|
if o.h.RunJob == nil {
|
|
o.finish(a, OutcomeRefused, "not available on this box")
|
|
return
|
|
}
|
|
done, err := o.h.RunJob(a.Arg)
|
|
if err != nil {
|
|
o.finish(a, OutcomeRefused, "the job "+a.Arg+" was not started: "+err.Error())
|
|
return
|
|
}
|
|
o.running.Add(1)
|
|
go func() {
|
|
defer o.running.Done()
|
|
var jerr error
|
|
select {
|
|
case jerr = <-done:
|
|
case <-o.ctx.Done():
|
|
jerr = o.ctx.Err()
|
|
}
|
|
// „done" says the job RAN TO ITS END, never what it found (security review 2026-10-08,
|
|
// „presence is not success"): offsite-integrity and offsite-proof return nil whatever their
|
|
// verdict — the verdict travels in their own log line, event and alarm.
|
|
msg := "the job " + a.Arg + " ran to its end — this does not say what it found; its finding is in its own log line and alarms"
|
|
if jerr != nil {
|
|
msg = "the job " + a.Arg + ": " + jerr.Error()
|
|
}
|
|
o.finish(a, outcomeOf(jerr), msg)
|
|
}()
|
|
default:
|
|
o.finish(a, OutcomeRefused, "unknown action")
|
|
}
|
|
}
|
|
|
|
// ── package wiring (the selftail.go shape: BuildReport reads it, main.go sets it) ──
|
|
|
|
var (
|
|
activeOpActionsMu sync.Mutex
|
|
activeOpActions *OperatorActions
|
|
)
|
|
|
|
// SetOperatorActions wires the executor whose results BuildReport attaches. nil = none.
|
|
func SetOperatorActions(o *OperatorActions) {
|
|
activeOpActionsMu.Lock()
|
|
defer activeOpActionsMu.Unlock()
|
|
activeOpActions = o
|
|
}
|
|
|
|
// pendingOperatorActionResults is what BuildReport attaches (nil in the steady state).
|
|
func pendingOperatorActionResults() []OperatorActionResult {
|
|
activeOpActionsMu.Lock()
|
|
o := activeOpActions
|
|
activeOpActionsMu.Unlock()
|
|
if o == nil {
|
|
return nil
|
|
}
|
|
return o.Results()
|
|
}
|