// 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 }