Files
admin cb692f8788 v0.83.0: observability pass — always-DEBUG capture ring + GET /debug/logs + heartbeat log-pull + gap-fill sweep
Capture layer: applog.New returns (logger, Ring) — slog fan-out, stderr at the
configured level, ~1000-entry ring fixed at LevelDebug (remote diagnostics
without a config flip). GET /debug/logs (token-authed, ?raw=1) + request-level
DEBUG middleware. Heartbeat log-pull mirrors the report logtail pattern:
envelope log_tail_requested -> next heartbeat carries log_tail (128KB cap,
consume-once, failed-push retry proven). Gap-fill sweep over netverify/
netstorage/netmount/signedjobs/selfupdate/disks/controller-swap/desired/loop.
Red-proofs: ring-at-emit-level FAILs capture test; drain removed FAILs
consume-once; dropped phase line FAILs the S7 log-sequence smoke.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PSK5g6qYLknKj8u3QAFEr6
2026-07-11 16:24:07 +02:00

230 lines
8.5 KiB
Go

// Package signedjobs is the slice-10B consumer of the hub's signed-jobs queue: it fetches each
// opaque queued blob, runs it through the reversibility gate (the LOCKED authz pipeline: pinned-key
// SSHSIG → namespace → allow-list-by-key-material → crypto → host target → time window → durable
// nonce-burn), and ONLY on all-pass hands the verified+bound op to an Executor. The order is
// verify → burn nonce (durable, inside Verify) → execute → report-completion, so no destructive op
// runs on an unsigned, non-pinned-signer, replayed, expired, retargeted, or path-only blob, and an
// interrupted op can't be re-authorized by replaying the same blob after a crash.
package signedjobs
import (
"context"
"encoding/base64"
"encoding/json"
"fmt"
"log/slog"
"sync"
"time"
"gitea.dooplex.hu/admin/felhom-agent/internal/authz"
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
"gitea.dooplex.hu/admin/felhom-agent/internal/reconcile"
)
// Envelope is the queued signed-op job payload (the opaque bytes the hub stores per job). The
// operator's signing CLI (cmd/felhom-opsign) produces it; this runner consumes it. The hub never
// forges one (it holds no signing key) and never opens it.
type Envelope struct {
OpBlobB64 string `json:"op_blob_b64"` // base64 of the canonical OpBlob JSON
SigArmored string `json:"sig_armored"` // armored SSHSIG over the op-blob (namespace felhom-op-v1)
}
// JobSource fetches + clears jobs. Satisfied by *hub.Client.
type JobSource interface {
Jobs(ctx context.Context) ([]hub.JobWire, error)
CompleteJob(ctx context.Context, jobID string) error
}
// Authorizer is the reversibility gate. Satisfied by *reconcile.Gate.
type Authorizer interface {
Authorize(intent reconcile.Intent, signed *reconcile.SignedOp) reconcile.Decision
}
// Executor runs a verified+bound destructive op. 10B wires the storage-wipe executor; other op
// classes (guest_destroy, decommission, restore_overwrite — 10D) plug in per-slice. The executor
// owns the RESOURCE-level anti-retarget (durable-id resolution + re-inspection at execution).
type Executor interface {
// Execute runs the op named by `op` with the VERIFIED canonical params. ErrNoExecutor means
// "this op has no executor in this build" (left in the queue for a later slice).
Execute(ctx context.Context, op string, params json.RawMessage) error
}
// ErrNoExecutor signals an op class with no executor wired in this build (don't clear the job).
var ErrNoExecutor = fmt.Errorf("signedjobs: no executor for this op class in this build")
// Runner polls the jobs queue and drives verify→execute for each job. Single-flight: a RunOnce
// already in progress short-circuits a concurrent trigger.
type Runner struct {
source JobSource
gate Authorizer
exec Executor
hostID string
logger *slog.Logger
mu sync.Mutex
running bool
}
// NewRunner builds a runner. All deps are required.
func NewRunner(source JobSource, gate Authorizer, exec Executor, hostID string, logger *slog.Logger) *Runner {
if logger == nil {
logger = slog.Default()
}
return &Runner{source: source, gate: gate, exec: exec, hostID: hostID, logger: logger}
}
// OnEnvelope implements hub.EnvelopeObserver: when the heartbeat flags pending signed ops, run a
// pass (off the heartbeat path so a long mkfs never blocks the loop). Single-flight.
func (r *Runner) OnEnvelope(ctx context.Context, env *hub.ControlEnvelope) {
if env == nil || !env.HasSignedOps {
return
}
go func() {
if _, err := r.RunOnce(context.Background()); err != nil {
r.logger.Warn("signedjobs: pass failed", "err", err)
}
}()
}
// RunOnce fetches the pending jobs and processes each. Returns the number of jobs PROCESSED
// (executed or terminally rejected/cleared). A fetch error returns early; per-job errors are
// logged and do not abort the pass.
func (r *Runner) RunOnce(ctx context.Context) (int, error) {
r.mu.Lock()
if r.running {
r.mu.Unlock()
return 0, nil // a pass is already in progress
}
r.running = true
r.mu.Unlock()
defer func() { r.mu.Lock(); r.running = false; r.mu.Unlock() }()
start := time.Now()
jobs, err := r.source.Jobs(ctx)
if err != nil {
return 0, fmt.Errorf("signedjobs: fetch jobs: %w", err)
}
ids := make([]string, len(jobs))
for i, j := range jobs {
ids[i] = j.JobID
}
r.logger.Debug("signedjobs: jobs fetched", "count", len(jobs), "ids", ids,
"duration_ms", time.Since(start).Milliseconds())
processed := 0
for _, j := range jobs {
if r.processJob(ctx, j) {
processed++
}
}
return processed, nil
}
// processJob handles one job; returns true if it was terminally handled (cleared from the queue).
func (r *Runner) processJob(ctx context.Context, j hub.JobWire) bool {
signed, ob, ok := r.decode(j)
if !ok {
// A malformed/forged envelope is permanently bad — clear it loudly (a queued garbage job
// must not wedge the queue). No nonce was consumed (we never reached Verify).
r.logger.Error("signedjobs: malformed job envelope — clearing", "job", j.JobID)
r.complete(ctx, j.JobID)
return true
}
// Op received — class/target/expiry only, NEVER the signature or nonce material.
r.logger.Debug("signedjobs: op received",
"job", j.JobID, "op", ob.Op, "target_host", ob.Target.HostID,
"issued_at", ob.IssuedAt.Format(time.RFC3339), "expires_at", ob.ExpiresAt.Format(time.RFC3339))
intent := reconcile.Intent{
Class: reconcile.OpClass(ob.Op),
HostID: ob.Target.HostID,
GuestID: ob.Target.GuestID,
ParamsJSON: ob.Params,
Source: reconcile.SourceOneShotJob,
// Provenance is the zero value (no agent-internal evidence) — a hub-sourced op is
// classified by its op class alone; a destructive class MUST carry a valid signature.
}
dec := r.gate.Authorize(intent, signed)
if !dec.Allowed {
// The signature/binding failed (forged, replayed, expired, retargeted, role-denied, …).
// The executor is NEVER called. This blob can never become valid, so clear it — but LOG
// it as a security event (a compromised hub queuing a forged op lands here).
r.logger.Warn("signedjobs: REJECTED signed op — executor not called",
"job", j.JobID, "op", ob.Op, "reason", dec.Reason, "err", errStr(dec.Err))
r.complete(ctx, j.JobID)
return true
}
// Allowed: the nonce is already durably burned (Verify, before this point). Execute.
r.logger.Warn("signedjobs: AUTHORIZED signed op — executing",
"job", j.JobID, "op", ob.Op, "key_id", dec.Verified.KeyID, "nonce", dec.Verified.Nonce)
err := r.exec.Execute(ctx, ob.Op, ob.Params)
switch {
case err == nil:
r.logger.Warn("signedjobs: signed op COMPLETED", "job", j.JobID, "op", ob.Op)
r.complete(ctx, j.JobID)
return true
case errorsIs(err, ErrNoExecutor):
// No executor for this op class in this build — leave it queued for the owning slice.
r.logger.Info("signedjobs: no executor for op (left queued)", "job", j.JobID, "op", ob.Op)
return false
default:
// Execution failed AFTER a passing verify — the nonce is spent, so re-running the same
// blob would be rejected as a replay anyway. Clear it (the operator re-signs with a fresh
// nonce to retry) and log the failure loudly.
r.logger.Error("signedjobs: signed op execution FAILED (nonce spent — clearing)",
"job", j.JobID, "op", ob.Op, "err", err)
r.complete(ctx, j.JobID)
return true
}
}
// decode unpacks a JobWire → SignedOp{op-blob, armored sig} + the parsed OpBlob (for routing).
func (r *Runner) decode(j hub.JobWire) (*reconcile.SignedOp, authz.OpBlob, bool) {
raw, err := base64.StdEncoding.DecodeString(j.BlobB64)
if err != nil {
return nil, authz.OpBlob{}, false
}
var env Envelope
if err := json.Unmarshal(raw, &env); err != nil || env.OpBlobB64 == "" || env.SigArmored == "" {
return nil, authz.OpBlob{}, false
}
opBlob, err := base64.StdEncoding.DecodeString(env.OpBlobB64)
if err != nil {
return nil, authz.OpBlob{}, false
}
var ob authz.OpBlob
if err := json.Unmarshal(opBlob, &ob); err != nil || ob.Op == "" {
return nil, authz.OpBlob{}, false
}
return &reconcile.SignedOp{Blob: opBlob, Sig: []byte(env.SigArmored)}, ob, true
}
func (r *Runner) complete(ctx context.Context, jobID string) {
if err := r.source.CompleteJob(ctx, jobID); err != nil {
r.logger.Warn("signedjobs: clearing job failed (will retry next pass)", "job", jobID, "err", err)
}
}
func errStr(e error) string {
if e == nil {
return ""
}
return e.Error()
}
func errorsIs(err, target error) bool {
for err != nil {
if err == target {
return true
}
type w interface{ Unwrap() error }
if u, ok := err.(w); ok {
err = u.Unwrap()
} else {
return false
}
}
return false
}