worker

package
v0.8.4 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 9, 2026 License: Apache-2.0 Imports: 31 Imported by: 0

Documentation

Overview

JSONB payload + HMAC signing for remediation jobs.

Mirrors the scan-job envelope (scheduler.JobPayload + hex HMAC tag), but the remediation payload carries its own load-bearing fields — request_id, rule_id, action, and (for rollback) txn_id — so it needs its own canonical encoding to HMAC-sign. We reuse the SAME queue key the scheduler derives (scheduler.DeriveQueueKey over the credential DEK); the encoding here is purpose-distinct from the scan encoding, so a scan tag can never validate a remediation payload and vice versa.

The worker verifies the HMAC BEFORE any host-mutating side effect, exactly like the scan path (system-worker-subcommand C-02).

Package worker — production remediation-job consumer.

RemediationWorker mirrors ScanWorker for the queued single-rule remediation path (Phase 7, Tier A free-core). It:

  1. Parses + HMAC-verifies the remediation payload BEFORE any side effect.
  2. Loads the approved (execute) / executed (rollback) request and guards its state via the remediation service's row-locked transitions.
  3. Calls executor.Remediate / executor.Rollback — which share the host's per-host inFlight guard with scans (a host is never scanned + remediated at the same instant).
  4. Writes the remediation_transactions journal and transitions the request to its terminal state (executed | failed | rolled_back).
  5. On a committed execute, flips THAT one rule to pass in host_rule_state via the transaction-log Writer (Kensa ran Validate before Commit, so the rule now passes — no full re-scan needed) so the compliance score moves.
  6. Publishes remediation.completed on the event bus, emits the audit event, and marks the queue row complete.

The worker owns the *kensa.Executor and the transaction-log Writer; the remediation service stays host-free and import-cycle-free (the worker maps kensa.RemediationTxn -> remediation.ExecTxn at this boundary).

Spec: api-remediation, system-worker-subcommand.

Package worker — production scan-job consumer.

ScanWorker is the long-lived loop that backs the `openwatch worker` subcommand. It:

  1. Claims one scan job at a time via queue.Dequeue (SKIP LOCKED).
  2. HMAC-verifies the payload via scheduler.Verify before any side effect.
  3. Takes a pg_advisory_xact_lock keyed on the host (per-host concurrency).
  4. Calls executor.Run with the host_id + policy_version (no framework arg — v2.0.0 framework-at-query-time architecture).
  5. On success: persists per-rule outcomes via transactionlog.Writer.Apply, resets host_backoff_state, queue.Complete.
  6. On transient executor errors: queue.Fail + UPSERT host_backoff_state with the ladder backoff (system-worker-subcommand C-05).
  7. On permanent executor errors: queue.Fail; the executor has already emitted scan.failed with the typed reason — the worker does NOT emit a second.
  8. On HMAC mismatch or malformed payload: queue.Fail + scheduler.job.hmac_rejected.

SIGTERM contract (C-07): ctx cancellation stops new Dequeue calls; the in-flight executor.Run completes (the executor owns its per-scan timeout); its result is persisted; then Run returns.

Spec: specs/system/worker-subcommand.spec.yaml

Package worker is the Stage-0 in-process job consumer. It polls the job_queue table on a short interval, drains diagnostics.test_job rows by emitting a diagnostics.test_job_completed audit event, and marks each job completed.

Stage 2 replaces this with a per-job-type dispatcher; the in-process loop here exists only to demonstrate the queue + correlation contract end-to-end (DoD step 16).

Spec: specs/release/stage-0-signoff.spec.yaml AC-10, C-02.

Index

Constants

View Source
const (
	RemediationActionExecute  = "execute"
	RemediationActionRollback = "rollback"
)

Remediation action discriminators carried in the payload.

View Source
const DefaultPollInterval = 1 * time.Second

DefaultPollInterval is the empty-queue sleep between Dequeue attempts. system-worker-subcommand C-10 / AC-11: 1s default, max 5s, no busy-spin.

View Source
const MaxBackoff = 24 * time.Hour

MaxBackoff is the suppress_until ceiling that kicks in at the 6th consecutive failure. After this point, the scheduler dispatcher will not enqueue a new job for the host for 24h — that is the de facto dead-letter.

View Source
const MaxPollInterval = 5 * time.Second

MaxPollInterval is the upper bound on the configurable poll_interval. Operators picking a longer value mistakenly increase scan-pickup latency without a corresponding benefit; we cap rather than let them pick 10m by accident.

View Source
const PollInterval = 200 * time.Millisecond

PollInterval is how often the loop checks for pending jobs when the queue is empty. Short enough that DoD step 16 sees a result within 2s.

View Source
const RemediationJobType = "remediation"

RemediationJobType is the queue.Job.JobType for a remediation execute or rollback. Distinct from ScanJobType so the dispatcher routes by type.

View Source
const ScanJobType = "scan"

ScanJobType is the queue.Job.JobType value scheduler.Service emits. Locked here so other job_types (diagnostics.test_job, future bulk ops) can coexist on the same queue without the scan worker grabbing them.

View Source
const TickInterval = 60 * time.Second

TickInterval is the nominal cadence of worker.loop.tick audit emission. System-worker-subcommand C-08 / AC-08: 55..65s range (60s + jitter).

Variables

This section is empty.

Functions

func MarshalRemediationJob

func MarshalRemediationJob(key []byte, p RemediationPayload) map[string]any

MarshalRemediationJob builds the signed JSONB body for queue.Enqueue. Used by the HTTP execute/rollback handlers so the worker can verify on claim.

Types

type Config

type Config struct {
	Pool         *pgxpool.Pool
	Executor     *kensa.Executor
	Writer       *transactionlog.Writer
	QueueKey     []byte // scheduler.DeriveQueueKey output
	PollInterval time.Duration
	Emit         EmitFunc

	// ScanResults, when non-nil, durably records every rule's outcome +
	// evidence for every scan (the /api/v1/scans audit memory). It is
	// written alongside Writer, never instead of it. nil disables the
	// durable write (legacy tests that only assert transaction-log
	// behavior leave it unset).
	ScanResults *scanresult.Writer

	// Bus, when non-nil, receives scan.completed events after outcomes
	// persist. The serve process passes its SSE bus; the dedicated
	// worker subcommand passes nil (its in-memory bus would have no
	// subscribers — cross-process delivery is a known non-goal).
	Bus *eventbus.Bus

	// Sched, when non-nil, receives PersistAfterScan after every
	// completed scan so host_compliance_schedule tracks the fresh
	// compliance state and next_scheduled_scan. Spec system-scheduler
	// v3.0.0: both scheduler-dispatched and on-demand scans update the
	// schedule on completion.
	Sched *scheduler.Service

	// RemediationProcessor, when non-nil, handles "remediation" job_type
	// rows this worker claims. queue.Dequeue is not type-filtered, so the
	// dedicated worker subcommand routes remediation jobs to it rather than
	// dead-lettering them. Spec api-remediation.
	RemediationProcessor *RemediationWorker

	// Regressions, when non-nil, projects each completed scan's
	// transaction-log regressions (a passing rule now fails) into a grouped
	// in-app notification. Best-effort: a projection error never fails the
	// scan. nil disables the in-app rule-regression bell (legacy tests, and
	// any boot path without the notification feed). Spec system-notifications.
	Regressions RegressionProjector

	// Drift, when non-nil, runs the compliance drift detector after every
	// completed scan's outcomes have committed and before ScanCompleted is
	// published. Production passes drift.NewService in both serve and the
	// worker subcommand. A detector error is logged and never fails the
	// scan. Spec system-drift-detector C-11, C-12.
	Drift DriftDetector

	// clock allows tests to inject a controllable time source.
	// Production passes time.Now.
	Clock func() time.Time
}

Config is the constructor argument bundle. All fields except clock are required.

type CredentialBridge

type CredentialBridge struct {
	// contains filtered or unexported fields
}

CredentialBridge wraps credential.Service to satisfy kensa.CredentialBridge. Constructor takes the live service; Resolve is invoked per scan by the executor.

func NewCredentialBridge

func NewCredentialBridge(svc *credential.Service) *CredentialBridge

NewCredentialBridge constructs a bridge backed by the live credential service. The same instance is shared across all scan jobs in the worker process.

func (*CredentialBridge) Resolve

func (b *CredentialBridge) Resolve(ctx context.Context, hostID uuid.UUID) ([]byte, func(), error)

Resolve fetches the host's credential (host-scoped first, system default fallback), extracts the SSH private key as a byte slice the executor can pass to Kensa, and returns a wipe function that zeros the slice on completion of the scan.

Returns kensa.ErrNoCredential when the underlying service reports no credential is available — the executor maps this to a short-circuit return without emitting scan.started (per kensa AC-09).

Returns kensa.ErrCredentialDecryption on any other resolve error; the executor maps this to scan.failed with reason credential_decryption_failed.

type DriftDetector added in v0.8.0

type DriftDetector interface {
	DetectForScan(ctx context.Context, hostID, scanID uuid.UUID) (drift.Report, error)
}

DriftDetector compares a completed scan against the host's prior state and emits compliance.drift.detected (and a DriftDetected bus event when the detector holds a bus) for a non-stable change. drift.Service implements it; the worker holds the interface so tests can stub a failing detector. Spec system-drift-detector C-11.

type EmitFunc

type EmitFunc func(ctx context.Context, code audit.Code, ev audit.Event)

EmitFunc matches audit.Emit's signature; production wires audit.Emit directly. Tests pass a recorder.

type GovernanceNotifier

type GovernanceNotifier interface {
	RemediationFailed(ctx context.Context, hostID uuid.UUID, ruleID, action, finalStatus string) error
}

GovernanceNotifier receives terminal remediation-failure signals for the in-app feed. notifyfeed.GovernanceProjector implements it; the worker holds the interface to avoid importing the feed package. Best-effort at the call site — a notify error never fails the job. Spec system-notifications (Slice 3).

type HostDiscoveryRunner

type HostDiscoveryRunner interface {
	RunDiscovery(ctx context.Context, hostID uuid.UUID) error
}

HostDiscoveryRunner is the seam the worker uses to invoke the OS Discovery flow when it drains a host.discovery job. The real implementation is internal/intelligence/discovery.Service.RunDiscovery (the error-only adapter over Discover). Interface lives here so the worker doesn't import intelligence (which would otherwise create a cycle via internal/credential's transitive imports).

type RegressionProjector

type RegressionProjector interface {
	ProjectScan(ctx context.Context, scanID, hostID uuid.UUID) error
}

RegressionProjector turns a completed scan's transaction-log changes into a grouped in-app notification. notifyfeed.Projector implements it; the worker holds the interface to avoid importing the feed package's concrete type and to let tests stub it. Spec system-notifications (Slice 2).

type RemediationConfig

type RemediationConfig struct {
	Pool     *pgxpool.Pool
	Executor *kensa.Executor
	Service  *remediation.Service
	Writer   *transactionlog.Writer
	QueueKey []byte
	Bus      *eventbus.Bus
	Emit     EmitFunc
	Clock    func() time.Time

	// Governance, when non-nil, receives a notification when a remediation
	// reaches a terminal FAILURE (an execute that failed, or a rollback that
	// did not restore). A successful user-initiated rollback is the intended
	// outcome and is NOT notified. Spec system-notifications (Slice 3).
	Governance GovernanceNotifier
}

RemediationConfig is the constructor bundle. All fields except Bus/Clock/Emit are required.

type RemediationPayload

type RemediationPayload struct {
	RequestID uuid.UUID
	HostID    uuid.UUID
	RuleID    string
	Action    string    // execute | rollback
	TxnID     uuid.UUID // rollback only; uuid.Nil for execute
	// ActorID is the user who invoked :execute or :rollback, so the worker's
	// terminal audit event names the same person the HTTP layer's intent
	// event did. uuid.Nil means no user initiated this job (system work),
	// and the worker then records a system actor rather than inventing one.
	// Signed when present: a job cannot be re-attributed after enqueue.
	ActorID uuid.UUID
	// ActorType is the audit actor type of ActorID: "user" for a session,
	// "api_key" for an API token (the token's own id, never its owner).
	// Empty on a payload signed before the field existed; the worker reads
	// that as "user", which is what those jobs meant. Signed when present,
	// so a token's job cannot be relabeled as a person's. bugs/OW-100.
	ActorType string
}

RemediationPayload is the typed payload of a remediation job.

type RemediationWorker

type RemediationWorker struct {
	// contains filtered or unexported fields
}

RemediationWorker processes "remediation" queue jobs. One per process, constructed at boot and held until shutdown. It is dispatched to by both the in-process Worker.process and the dedicated ScanWorker.ProcessJob (each routes by JobType so one worker drains the queue).

func NewRemediationWorker

func NewRemediationWorker(cfg RemediationConfig) *RemediationWorker

NewRemediationWorker wires a RemediationWorker.

func (*RemediationWorker) ProcessJob

func (w *RemediationWorker) ProcessJob(ctx context.Context, j *queue.Job)

ProcessJob runs the full remediation pipeline for one job. Recovers from panics so a rogue remediation does not take down the worker.

type ReportRenderer

type ReportRenderer interface {
	ProcessJob(ctx context.Context, j *queue.Job)
}

Worker drains pending jobs from job_queue. One Worker per process is enough for Stage 0; multi-worker setups are Stage 2.

The worker can run several claim/process loops concurrently (WithConcurrency) so a fleet of queued scans does not drain one host at a time; the queue's SKIP LOCKED claim and the scan path's per-host advisory lock keep concurrent draining safe (system-job-queue C-07). ReportRenderer renders a report's faces for a claimed "report.render" job and publishes the ready event. Implemented by report.RenderProcessor; defined as an interface here so the worker package does not import the report package (which would cycle through the queue/eventbus packages they share).

type ScanWorker

type ScanWorker struct {
	// contains filtered or unexported fields
}

ScanWorker is the production worker. One per process. Constructed at boot in cmd/openwatch/worker.go via NewScanWorker; held until SIGTERM.

func NewScanWorker

func NewScanWorker(cfg Config) *ScanWorker

NewScanWorker wires a ScanWorker. The dependencies are constructed at boot in cmd/openwatch/worker.go and held for the process lifetime.

func (*ScanWorker) ProcessJob

func (w *ScanWorker) ProcessJob(ctx context.Context, j *queue.Job)

ProcessJob runs the full per-job pipeline. Recovers from panics so a rogue scan does not take down the worker; emits a typed audit on the panic path.

func (*ScanWorker) Run

func (w *ScanWorker) Run(ctx context.Context) error

Run is the worker loop. Returns nil on clean shutdown (ctx canceled and any in-flight job completed). Returns a non-nil error only if the loop encountered a fatal problem (none today — every per-job error is classified and handled).

Spec C-07 / AC-06: ctx cancellation stops new Dequeue calls; an in-flight job is allowed to finish (the executor owns its per-scan timeout), its result is applied, then Run returns.

type Worker

type Worker struct {
	// contains filtered or unexported fields
}

func New

func New(pool *pgxpool.Pool) *Worker

New constructs a Worker bound to the given pool. Call Start to begin the drain loop and Stop to exit cleanly. Defaults to one (serial) loop; call WithConcurrency to fan out.

func (*Worker) HasScanProcessor added in v0.8.0

func (w *Worker) HasScanProcessor() bool

HasScanProcessor reports whether a scan processor is registered on this worker, so a caller can verify the wiring rather than read the source and hope. system-worker-subcommand AC-19 uses it to prove the serve process really runs scans: a source-inspection check passes on a registration guarded by an impossible condition, and this does not.

Read-only. It changes nothing about how jobs are claimed or run.

func (*Worker) Start

func (w *Worker) Start(ctx context.Context)

Start kicks off w.concurrency drain loops on background goroutines. Returns immediately. Safe to call once per Worker. Each loop claims jobs independently (SKIP LOCKED), so up to w.concurrency jobs run at once.

func (*Worker) Stop

func (w *Worker) Stop()

Stop signals the loop to exit and waits for the in-flight drain (if any) to complete.

func (*Worker) WithConcurrency

func (w *Worker) WithConcurrency(n int) *Worker

WithConcurrency sets how many claim/process loops run at once. A value < 1 clamps to 1 (strictly serial). Each loop independently claims jobs via SKIP LOCKED, so N loops process up to N distinct hosts in parallel while the per-host advisory lock still serializes same-host work. Spec system-job-queue C-07.

func (*Worker) WithDiscovery

func (w *Worker) WithDiscovery(d HostDiscoveryRunner) *Worker

WithDiscovery registers the OS Discovery runner. When set, the worker processes host.discovery jobs by calling Discover; nil keeps the legacy behavior (host.discovery fails as unsupported). Spec system-host-discovery C-05.

func (*Worker) WithRemediationProcessor

func (w *Worker) WithRemediationProcessor(rw *RemediationWorker) *Worker

WithRemediationProcessor registers a RemediationWorker whose ProcessJob handles "remediation" jobs claimed by THIS worker's loop. Like scan jobs, queue.Dequeue is not type-filtered, so the in-process worker must route remediation jobs rather than fail them as unsupported. Spec api-remediation.

func (*Worker) WithReportProcessor

func (w *Worker) WithReportProcessor(rp ReportRenderer) *Worker

WithReportProcessor registers a report RenderProcessor whose ProcessJob handles "report.render" jobs claimed by THIS worker's loop. Like scan and remediation jobs, queue.Dequeue is not type-filtered, so the in-process worker must route report renders rather than fail them as unsupported. Spec api-reports.

func (*Worker) WithScanProcessor

func (w *Worker) WithScanProcessor(sw *ScanWorker) *Worker

WithScanProcessor registers a ScanWorker whose ProcessJob handles "scan" jobs claimed by THIS worker's loop. queue.Dequeue is not type-filtered, so the in-process worker must route scan jobs rather than fail them as unsupported — otherwise a serve-process claim of an on-demand scan would dead-end the job. The single-binary deployment processes scans in-process; a dedicated `openwatch worker` process can run alongside (both are capable; first claim wins). Spec api-host-scan / system-scan-runs.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL