feat(store): outbox consumers for TAXII, Clearfolio, and orchestrator - #106
Conversation
Enqueue operator-triggered TAXII polls, Clearfolio submits, and
contextual-orchestrator SOC analysis on the PostgreSQL leased outbox.
Request path returns 202; GET /api/outbox/{id} exposes receipt evidence.
Secrets stay in the credential registry. Startup postgres save advances
snapshot_version so the first management write cannot false-conflict.
|
Important Review skippedAuto reviews are disabled on base/target branches other than the default branch. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Plus Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
| let response = request.send().await.map_err(|err| { | ||
| outbox::DispatchError::Transient(format!("clearfolio request failed: {err}")) | ||
| })?; |
There was a problem hiding this comment.
🔴 Hung Clearfolio or LLM upstream stalls the entire outbox worker
The Clearfolio submit (request.send()) and the SOC LLM call run inside the single leased outbox worker on an HTTP client with no timeout. A hung upstream blocks the drain loop forever, halting all outbox dispatch, including security-event SIEM export (lib.rs).
Prompt for agents
The outbound HTTP client state.http is built in outbound_http_client (src/lib.rs:3909-3913) with no request timeout. Two new worker dispatch paths use it without a per-request timeout: execute_clearfolio_submit (request.send() at src/lib.rs:853) and execute_soc_analyze (send() at src/lib.rs:1058-1065). Both run inside the single leased outbox worker through drain_due_async, which processes claimed messages sequentially and awaits each dispatch. A slow or unresponsive Clearfolio or SOC LLM endpoint will block the drain loop indefinitely, stalling all outbox processing (including security-event stdout SIEM export). Add a bounded per-request .timeout(...) to these two sends (as fetch_taxii_objects already does at src/lib.rs:2218), or set a default timeout on the shared client, so a hung upstream fails as a Transient dispatch error instead of blocking the worker.
Was this helpful? React with 👍 or 👎 to provide feedback.
| pub async fn get_outbox_item( | ||
| &self, | ||
| message_id: &str, | ||
| ) -> Result<Option<(OutboxMessage, Option<String>)>, String> { | ||
| let mut client = self.client.lock().await; | ||
| let tx = client | ||
| .transaction() | ||
| .await | ||
| .map_err(|error| format!("control plane get-outbox transaction failed: {error}"))?; | ||
| tx.execute( | ||
| "SELECT set_config('wardnet.tenant_id', $1, true)", | ||
| &[&self.tenant_id], | ||
| ) | ||
| .await | ||
| .map_err(|error| format!("control plane tenant context failed: {error}"))?; | ||
| let row = tx | ||
| .query_opt( | ||
| "SELECT message_id, aggregate_id, aggregate_version, event_type, schema_version, | ||
| created_unix, payload_json, payload_hash, idempotency_key, message_status, | ||
| lease_owner, lease_expires_unix, attempt_count, first_attempt_unix, | ||
| last_attempt_unix, next_available_unix, terminal_reason | ||
| FROM outbox_message WHERE tenant_id = $1 AND message_id = $2", | ||
| &[&self.tenant_id, &message_id], | ||
| ) | ||
| .await | ||
| .map_err(|error| format!("control plane get outbox_message failed: {error}"))?; | ||
| let Some(row) = row else { | ||
| tx.commit() | ||
| .await | ||
| .map_err(|error| format!("control plane get-outbox commit failed: {error}"))?; | ||
| return Ok(None); | ||
| }; | ||
| let message = row_to_outbox(&row, &self.tenant_id); | ||
| let evidence = tx | ||
| .query_opt( | ||
| "SELECT receipt_evidence FROM outbox_receipt | ||
| WHERE tenant_id = $1 AND message_id = $2", | ||
| &[&self.tenant_id, &message_id], | ||
| ) | ||
| .await | ||
| .map_err(|error| format!("control plane get outbox_receipt failed: {error}"))? | ||
| .map(|row| row.get::<_, String>(0)); | ||
| tx.commit() | ||
| .await | ||
| .map_err(|error| format!("control plane get-outbox commit failed: {error}"))?; | ||
| Ok(Some((message, evidence))) | ||
| } |
There was a problem hiding this comment.
📝 Info: Pruned outbox rows can hide a successful receipt
get_outbox_item returns 404 once the outbox_message row is pruned to EVENT_LIMIT, even though the receipt row survives. The browser's pollOutboxReceipt treats 404 as keep-waiting and eventually times out. Under high throughput with a small EVENT_LIMIT, a succeeded effect can surface as a timeout.
Was this helpful? React with 👍 or 👎 to provide feedback.
| let _ = plane | ||
| .drain_once(&owner, now_unix() as i64, outbox::dispatch_stdout) | ||
| .drain_due_async(&owner, now_unix() as i64, move |message| { | ||
| let state = worker_state.clone(); | ||
| async move { dispatch_outbox_message(&state, &message).await } | ||
| }) | ||
| .await; |
There was a problem hiding this comment.
📝 Info: HTTP effects share the SIEM export worker
The single outbox worker drains a claimed batch sequentially. A slow TAXII poll (bounded at 15s) delays every other message in the batch, including security-event SIEM export. Folding outbound HTTP effects into the same at-least-once worker couples effect latency to SIEM timeliness.
Was this helpful? React with 👍 or 👎 to provide feedback.
| let status = response.status().as_u16(); | ||
| let json = response.json::<serde_json::Value>().await.map_err(|err| { | ||
| outbox::DispatchError::Transient(format!("llm response read failed: {err}")) | ||
| })?; | ||
| let preview = json.to_string(); | ||
| outbox::classify_http_status(status, &preview)?; |
There was a problem hiding this comment.
📝 Info: SOC LLM non-JSON 4xx retried instead of dead-lettered
execute_soc_analyze reads response.json() before classify_http_status. A 4xx with a non-JSON body fails the parse and becomes Transient, so the worker retries a permanent client error to MAX_ATTEMPTS instead of dead-lettering. Clearfolio avoids this by classifying on status before parsing.
Was this helpful? React with 👍 or 👎 to provide feedback.
| pub async fn enqueue_effect( | ||
| &self, | ||
| event_type: &'static str, | ||
| aggregate_id: &str, | ||
| payload_json: String, | ||
| ) -> Result<String, String> { | ||
| let created_unix = unix_now_i64(); | ||
| let hash = outbox::payload_hash(&payload_json); | ||
| let unique = format!("{created_unix}:{hash}"); | ||
| let (message_id, idempotency_key) = | ||
| outbox::effect_ids(event_type, &self.tenant_id, &unique); | ||
| let mut client = self.client.lock().await; | ||
| let tx = client | ||
| .transaction() | ||
| .await | ||
| .map_err(|error| format!("control plane effect transaction failed: {error}"))?; | ||
| tx.execute( | ||
| "SELECT set_config('wardnet.tenant_id', $1, true)", | ||
| &[&self.tenant_id], | ||
| ) | ||
| .await | ||
| .map_err(|error| format!("control plane tenant context failed: {error}"))?; | ||
| insert_outbox( | ||
| &tx, | ||
| &self.tenant_id, | ||
| &OutboxInsert { | ||
| message_id: message_id.clone(), | ||
| aggregate_id: aggregate_id.to_string(), | ||
| aggregate_version: created_unix, | ||
| event_type, | ||
| created_unix, | ||
| payload_json, | ||
| payload_hash: hash, | ||
| idempotency_key, | ||
| }, | ||
| ) | ||
| .await?; | ||
| prune_processed_outbox(&tx, &self.tenant_id, self.event_limit).await?; | ||
| tx.commit() | ||
| .await | ||
| .map_err(|error| format!("control plane effect commit failed: {error}"))?; | ||
| Ok(message_id) | ||
| } |
There was a problem hiding this comment.
📝 Info: Same-second duplicate enqueue silently dedups
enqueue_effect keys idempotency on created_unix plus payload hash and inserts ON CONFLICT DO NOTHING, yet returns message_id without checking whether a row was written. Two identical requests in the same second collapse to one message; both return 202 with that id. Idempotent by design, but the caller cannot tell a fresh enqueue from a dedup.
Was this helpful? React with 👍 or 👎 to provide feedback.
There was a problem hiding this comment.
Pull request overview
OpenCode cannot approve yet because required coverage evidence did not pass.
Review outcome
1. HIGH .github/workflows/opencode-review.yml:1 - Coverage evidence did not prove required test/docstring evidence
-
Problem: The required coverage-evidence job result was
failure, so OpenCode cannot establish approval sufficiency for this head. -
Root cause: Automated approval is only valid when the same-head coverage-evidence job proves supported repository test suites passed and configured docstring gates passed or were advisory, or reports not applicable because no supported source files or package manifests exist. Missing, failed, skipped, unavailable, or unsupported-tooling test evidence is a blocker.
-
Fix: Install or configure the repository test/docstring evidence tooling when source files or package manifests exist, rerun the current-head coverage-evidence job, and approve only after it reports
successwith required evidence or explicit no-source not-applicable evidence. -
Regression test: Keep the approval branch checking
needs.coverage-evidence.result == successbefore posting APPROVE, and publish REQUEST_CHANGES when coverage-evidence blocker states such as cancelled, skipped, failed, unsupported-tooling, or below-100 evidence are present. -
Result: REQUEST_CHANGES
-
Reason: coverage-evidence result was
failure, so required test/docstring evidence was not proven for current head18fd115b2c12d6c4956cc08b2a2f81c567529976. -
Head SHA:
18fd115b2c12d6c4956cc08b2a2f81c567529976 -
Workflow run: 32702428652
-
Workflow attempt: 1
Coverage evidence
Coverage Decision
- Result: FAIL
- Test evidence: not proven passing
- Docstring evidence: not proven passing when configured
- Failure count: 1
Changed-File Evidence Map
flowchart LR
PR["PR changed files"] --> Evidence["OpenCode bounded evidence"]
Evidence --> S1["Changed file (5 files)"]
S1 --> I1["repository behavior"]
I1 --> R1["Review risk: Changed file (5 files)"]
R1 --> V1["required checks"]
Evidence --> S2["Docs (3 files)"]
S2 --> I2["operator or user guidance"]
I2 --> R2["Review risk: Docs (3 files)"]
R2 --> V2["docs review"]
OpenCode Review Overview
Pull request overviewOpenCode cannot approve yet because required coverage evidence did not pass. Review outcome1. HIGH .github/workflows/opencode-review.yml:1 - Coverage evidence did not prove required test/docstring evidence
Coverage evidenceCoverage Decision
Changed-File Evidence Mapflowchart LR
PR["PR changed files"] --> Evidence["OpenCode bounded evidence"]
Evidence --> S1["Changed file (5 files)"]
S1 --> I1["repository behavior"]
I1 --> R1["Review risk: Changed file (5 files)"]
R1 --> V1["required checks"]
Evidence --> S2["Docs (3 files)"]
S2 --> I2["operator or user guidance"]
I2 --> R2["Review risk: Docs (3 files)"]
R2 --> V2["docs review"]
|
6831fbc
into
feat/issue-80-optimistic-concurrency
* feat(security): fail-closed destination policy for outbound HTTP One DestinationPolicy mediates gateway upstreams, threat-intel fetches, Clearfolio, SOC LLM, and the Coraza sidecar URL. Private, loopback, link-local, CGNAT, and metadata classes are denied unless DESTINATION_ALLOWLIST (or loopback development) permits them; DESTINATION_DENYLIST wins. Clients ignore ambient HTTP proxies and do not follow redirects. Refs #79. * docs: record PR #96 in the product-technical gap baseline * feat(security): harden destination policy per-IP CIDR and readiness order CIDR allowlist matches apply per resolved address, authorize non-default ports, and reject prefixes outside the address-family width. IPv6 site-local is a denied class. Hostnames that merely contain 0x are not hex IP literals. AppState constructors default to production policy; seeded fixtures opt into development. Blocking DNS runs on spawn_blocking with a timeout. Persistence and destination-list validation complete before the readiness line. * feat(security): pin outbound HTTP to evaluated destination addresses After destination policy allows a host, the reqwest client resolves only those IPs so a rebinding answer cannot reach loopback, private, or metadata classes. Host and SNI stay on the original name. Unpinned hostnames fail closed instead of falling back to OS DNS. * feat(waf): evaluate live gateway transactions with in-process libcoraza Issue #86 remainder: dlopen operator-supplied libcoraza and drive the C ABI on each /gateway request. Missing library or empty ruleset fail closed before bind. CI stays hermetic with a fixture cdylib that exports the same symbols. * docs: record PR #97 in the product-technical gap baseline * feat(store): require PostgreSQL as the production control plane Non-loopback binds fail closed without CONTROL_PLANE_DATABASE_URL. Migrations create 3NF two-word tables with default-deny RLS. Snapshot persist commits policy rows and audit records in one transaction. Loopback still uses the JSON file or memory adapter. * feat(store): transactional outbox and leased workers Issue #81 first slice on the PostgreSQL control plane. Security events append with an outbox row in one transaction instead of rewriting the snapshot. Workers claim with SKIP LOCKED, retry, dead-letter, and record unique receipts. Stdout SIEM is at-least-once; the receipt is the exactly-once ack. Also deterministic ORDER BY on postgres loads (still-valid #98 finding). Do not re-implement the postgres gate. * feat(store): rustls for production PostgreSQL sslmode=require (#100) * feat(store): rustls for production PostgreSQL sslmode=require Issue #80 remainder. sslmode=require/verify-ca/verify-full connect with rustls and Mozilla roots; certificates are always verified. allow/prefer are still rejected so the process cannot silently drop to plaintext. Live test against plaintext CI postgres proves fail-closed. Do not re-implement the postgres gate or the outbox. * docs: record PR #100 in the product-technical gap baseline * fix(store): rewrite verify-full sslmode for tokio-postgres 0.7 Still-valid #100 Devin finding. tokio-postgres 0.7 only parses disable/prefer/require. Map verify-ca/verify-full to require before connect; rustls still verifies certificates. Password query-lookalikes are left untouched. * feat(store): bound outbox listing and prune processed rows (#101) * feat(store): bound outbox listing and prune processed rows Still-valid #99 finding. GET /api/outbox returns at most EVENT_LIMIT rows (dead letters, then pending, then leased, then processed). Processed outbox_message rows prune to that cap; receipts and dead letters stay. Do not re-implement the outbox, postgres gate, or rustls. * docs: record PR #101 in the product-technical gap baseline * fix(store): prune processed outbox to EVENT_LIMIT on save and ack Still-valid #101 Devin finding. Snapshot save and worker ack used LIST_LIMIT (1000) while append used operator EVENT_LIMIT. Store the configured cap on PostgresPlane so all three paths retain the same processed-row bound. Receipts and dead letters stay. * feat(store): logical backup and isolated restore drill (#102) * feat(store): logical backup and isolated restore drill Issue #80 remainder stacked on #101. GET /api/backup exports a hashed tenant snapshot; POST /api/backup restores after schema and payload-hash checks; POST /api/backup/drill restores into an isolated tenant, compares unmasked invariants, and drops the drill rows. Declared RPO is last successful export; declared RTO is 60s. File/memory adapters report backup=disabled. Do not re-implement rustls, outbox, or retention. * docs: record PR #102 in the product-technical gap baseline * feat(store): non-owner PostgreSQL runtime role after migrate (#103) * feat(store): non-owner PostgreSQL runtime role after migrate Still-valid #98 finding. CI connects as a superuser, which bypasses FORCE RLS. Migrations stay on the login role, then SET ROLE wardnet_runtime (NOSUPERUSER, NOBYPASSRLS, not table owner). Missing tenant GUC yields no rows; DROP TABLE and DISABLE RLS are denied. Do not re-implement rustls, outbox, retention, or backup/restore. * docs: record PR #103 in the product-technical gap baseline * fix(store): restore logical backups across role-only schema versions v3 only provisions wardnet_runtime and does not change table shape. verify() accepts schema 2 through the current migration version so a role-only upgrade cannot void the last pre-upgrade logical backup. * feat(store): HASH-partition security_event by tenant (#104) * feat(store): HASH-partition security_event by tenant Convert unpartitioned security_event to PARTITION BY HASH (tenant_id) with eight children under pg_advisory_lock. Rows keep unmasked client IPs and paths. /healthz.event_partitions reports the child count. Logical restore still accepts schema 2 through the current version. * docs: record PR #104 in the product-technical gap baseline * feat(store): optimistic concurrency on postgres snapshots Issue #80 last remainder. tenant_account.snapshot_version must match the loaded token or persist fails closed (HTTP 409). Restores overwrite. File/memory stay single-writer. Do not re-implement rustls, outbox, runtime role, HASH, or backup/restore. * docs: record PR #105 in the product-technical gap baseline * fix(store): keep postgres snapshot_version aligned after startup save load_postgres was saving with OCC and leaving the in-memory token one behind the database, so every later management write returned HTTP 409. Advance the loaded snapshot_version to the value save() wrote. * feat(store): outbox consumers for TAXII, Clearfolio, and orchestrator Enqueue operator-triggered TAXII polls, Clearfolio submits, and contextual-orchestrator SOC analysis on the PostgreSQL leased outbox. Request path returns 202; GET /api/outbox/{id} exposes receipt evidence. Secrets stay in the credential registry. Startup postgres save advances snapshot_version so the first management write cannot false-conflict. * docs: record PR #106 in the product-technical gap baseline * feat(release): tagged GitHub Release with SHA-256 and immutable GHCR (#107) * feat(release): tagged GitHub Release with SHA-256 and immutable GHCR Issue #84 first slice. A vX.Y.Z tag builds a locked binary, checksums, a GitHub Release, and ghcr.io/contextualwisdomlab/waf-ids-ai-soc:vX.Y.Z with no moving latest tag. Promotion and rollback are tag-for-tag. Do not re-implement store slices or OCC. * docs: record PR #107 in the product-technical gap baseline * fix(release): basename checksums and serialize postgres GRANTs SHA256SUMS recorded dist/ prefixes so sha256sum -c failed next to the downloaded binary. Emit basenames. Parallel PostgresPlane connects raced HASH convert GRANT with SET ROLE GRANT (tuple concurrently updated); hold the advisory lock across both. Do not re-implement HASH layout. * feat(release): keyless cosign, SPDX SBOM, and SLSA on the same tag Issue #84 remainder. GitHub OIDC signs the binary, checksums, SBOMs, and the GHCR image by digest. Release is created only after signatures. Syft SPDX fails closed without syft or non-SPDX JSON. NIST SP 800-218 is attached. Do not re-implement checksums or store slices. * feat(release): refuse lightweight tags and pin k8s by digest Issue #84 remainder. Annotated vX.Y.Z tags only; lightweight tags fail closed before the release job builds. Kubernetes pin is the GHCR content digest; tag aliases are refused. Do not re-implement checksums or cosign/SBOM. * docs: record PR #109 in the product-technical gap baseline * feat(waf): detect OWASP CRS attack battery on the live binary Issue #11 first slice. The build-script libcoraza ABI stub gains a deterministic battery covering SQLi (942100), XSS (941100), path traversal (930100), Unix RCE (932100, with first-match ordering so '; cat /etc/passwd' attributes to RCE over traversal), and Log4j JNDI (944120) in raw and percent-encoded forms across URI and POST-body phases. tests/binary.rs now starts the real gateway with the stub engine, creates a block route through the admin API, fires nine cases over HTTP, and asserts each is 403-blocked citing the expected CRS rule id while a benign request still forwards; /api/events must record one event per attempt with the forwarded client IP kept unmasked. Doctoring: docs/doctoring/ci-attack-evidence-battery.md grounds the split between detection-path evidence (CI) and detection efficacy (operator-supplied libcoraza + Core Rule Set), APA 7th. * fix(control-plane): close OCC and credential race gaps * fix(release): capture pushed image digest
Summary
Issue #81 extra consumers stacked on #105.
PostgreSQL operator-triggered external effects now use the same leased outbox / receipt contract as security events:
taxii.collection_polled— TAXII 2.1 poll + STIX importclearfolio.document_submitted— Clearfolio convert jobsoc.analysis_requested— contextual-orchestrator chat completions (advisory only; never auto-enforces)Request path on PostgreSQL returns HTTP 202 with
message_id.GET /api/outbox/{id}returns receipt evidence. File/memory adapters keep the previous synchronous path.Secrets never enter outbox payloads. Inline TAXII Basic/Bearer is rejected on the durable path;
taxii_bearerandsoc_llm_tokenlive in the credential registry. Client IPs, paths, indicator values, and actor names stay unmasked.HTTP dispatch releases the PostgreSQL client lock during outbound I/O (
drain_due_async). 429/5xx retry; other 4xx dead-letter.Also lands on this stack: #105 startup
snapshot_versionalignment (Devin still-valid 409).Do not re-implement #78, sidecar #95, pin #96, libcoraza #97, postgres gate #98, outbox #99, rustls #100, retention #101, backup/restore #102, runtime role #103, HASH #104, or OCC #105.
Merge
Stacked on #105 (
feat/issue-80-optimistic-concurrency). Merge order: #94 independently; #95 then #96 then #97 then #98 then #99 then #105 then this. Do not--adminmerge. Org ruleset 18156473 (2 independent approvals) remains a policy blocker for this actor.Verification
cargo fmt --checkcargo test --locked --workspacewithCONTROL_PLANE_TEST_DATABASE_URLcargo clippy --locked --workspace --all-targets -- -D warningsscripts/smoke.shGET /healthzplus/adminand/api/commercial/readiness(2B KRW unchanged)