Task — engineering-spec@1

"Layer-2 ingest pipeline (binlog → pgvector + Apache AGE) — chain-anchor verification + 1s p95 lag + per-tenant cursor + idempotent UPSERT"

doneTASK-MEMORY-101
module memory · class product · priority p0 · created 2026-05-15 · shipped 2026-05-23
depends on TASK-AI-019, TASK-AUTH-003 · blocks TASK-MEMORY-102, TASK-MEMORY-103, TASK-MEMORY-106, TASK-MEMORY-108, TASK-MEMORY-105, TASK-PROJ-008

§1 — Description (BCP-14 normative)

A long-lived Rust process MUST tail every <memory-root>/audit/*.binlog and ingest each row into pgvector (embeddings) + Apache AGE (graph) within ≤ 1 second of binlog append. Each piece:

  1. MUST consume from binlog via per-tenant offset cursor; persist last-consumed seq to Postgres for restart resume. Cursor table layer2_ingest_cursor (tenant_id PK, last_seq, updated_at).
  2. MUST compute BGE-M3 embedding via TASK-AI-019 sidecar for every memory-body row; insert into pgvector with (tenant_id, seq, id, kind, path, ts_ns, embedding 1024-dim, chain_anchor 32-bytes).
  3. MUST extract entities (PERSON, ORG, PLACE, CONCEPT) + relations via spaCy NER + custom recognizers; insert into AGE graph as nodes + edges. Failure to extract entities does NOT block pgvector insertion (per §1 #11 fallback).
  4. MUST tag every Layer 2 row with chain_anchor = SHA-256(canonical(Layer 1 row at seq N)) for tamper detection. Read paths (TASK-MEMORY-108) verify chain_anchor matches Layer 1 before returning results; mismatch → sev-1 + drop the result.
  5. MUST be idempotent — re-ingesting the same seq is a no-op (UPSERT on (tenant_id, seq)). Restart-mid-batch can re-process some rows; the UPSERT prevents duplicates.
  6. MUST target ingest staleness ≤ 1 second p95 from binlog append to Layer 2 visibility (DEC-074). Steady-state SLO measured via memory_layer2_ingest_lag_seconds histogram.
  7. MUST emit OTel metrics:
  1. MUST support tenant isolation via tenant_id column + RLS (TASK-AUTH-003 pattern). The layer2_memories table is in TENANT_SCOPED_TABLES registry; RLS USING + WITH CHECK clauses applied per TASK-AUTH-003 §1 #2.
  2. MUST NOT be the source of truth — DEC-070 invariant: Layer 1 wins on any conflict. Read paths comparing Layer 1 vs Layer 2 (e.g., on chain_anchor mismatch) MUST trust Layer 1 + flag Layer 2 row as stale.
  3. MUST retry BGE-M3 sidecar failures with exponential backoff (100ms, 250ms, 500ms, 1s, 2s) up to 5 attempts before marking the row as pending_embed_retry. Retry job (slice 2 follow-up) reprocesses such rows.
  4. MUST apply graceful AGE-entity-extraction fallback: if spaCy/AGE call fails, write the pgvector row (with embedding) AND insert an AGE-pending row marked state: pending_retry. The pgvector row is queryable; AGE backfill runs in a separate job.
  5. MUST validate Layer 1 row's structural integrity before ingestion: parse the binlog frame; check seq monotonicity (next seq = previous + 1 per tenant); chain_anchor recompute (frames carry their own hash). Frame-level validation failure → log sev-2 + skip + advance cursor (alternative would block ingest indefinitely on a corrupted row).
  6. MUST support concurrent multi-tenant ingestion via tokio task per tenant. Each tenant's ingest runs independently; one tenant's BGE saturation doesn't starve another. Task scheduling round-robin per CPU core.
  7. MUST support graceful shutdown: on SIGTERM, current in-flight transaction completes; cursor saved; tasks drain; service exits cleanly. No partial commits.
  8. SHOULD support a --rebuild flag for TASK-MEMORY-102 CI gate: starts from seq 0 for all tenants; ignores existing cursor; re-ingests everything from Layer 1.

§2 — Why this design (rationale for humans)

Why per-tenant cursor (DEC-073)? Global cursor would let one tenant's slow ingest stall everyone. Per-tenant + tokio task per tenant means tenant A can be 10s behind without affecting tenant B. Also: tenant-scoped restart resume is more granular.

Why chain_anchor (DEC-072)? Layer 1 is the source of truth, but read paths hit Layer 2 for performance. If Layer 1 is corrupted (e.g., disk bit-flip, malicious mutation), Layer 2 contains the pre-corruption state. The chain_anchor lets read paths detect drift: recompute SHA-256 of Layer 1 row at seq N; compare to stored chain_anchor; mismatch → flag as sev-1 + drop the result.

Why 1s p95 lag (DEC-074)? Users expect "I just saved this; let me search for it" to work. 1s is below human perception of staleness; longer windows produce "where did my data go?" complaints. The 1s budget covers BGE embedding (~50ms on GPU), entity extraction (~30ms), pgvector + AGE writes (~10ms) with healthy margin.

Why graceful AGE fallback (§1 #11)? Entity extraction is the slowest part of the pipeline; a spaCy hang would block ALL ingestion. Writing pgvector + marking AGE pending preserves the searchability (vector search works) while AGE catches up via background job.

Why idempotent UPSERT (§1 #5)? Restart-mid-batch is the recovery scenario. After SIGTERM, the cursor is saved at seq N; some rows up to N+1 might have been written with the old transaction not yet committed. On restart, ingest re-processes from cursor; UPSERT ensures no duplicates.

Why frame-level validation skip on corruption (§1 #12)? A single corrupted frame in binlog should NOT block ingest indefinitely. Skip + sev-2 + advance cursor lets ingest continue; the corrupted row is investigated separately. Without skip, one bad frame stalls the entire tenant's ingest.

Why concurrent multi-tenant (§1 #13)? Single-threaded ingest scales linearly with tenant count. At 100 tenants × 100 rows/sec, single thread saturates. Per-tenant tasks distribute across CPU cores.

Why DEC-070 Layer 1 wins (§1 #9)? Layer 2 is derived. If Layer 2 says X but Layer 1 says Y, Y is correct (Layer 2 derivation must have a bug). Read paths trust Layer 1 + flag Layer 2 for re-ingest.

Why per-tenant tokio task (§1 #13)? Better than one big task with internal queue: per-tenant task has natural backpressure (the task either keeps up or doesn't); no need for a separate scheduler. Failure mode is contained: one task panicking only affects one tenant.

Why graceful shutdown (§1 #14)? Mid-transaction abort would leave pgvector + AGE inconsistent. Drain + commit ensures atomicity — either all 3 inserts (pgvector, AGE, cursor) succeed or none.


§3 — API contract

-- services/memory/migrations/0001_layer2.sql
CREATE EXTENSION IF NOT EXISTS vector;
CREATE EXTENSION IF NOT EXISTS age;
LOAD 'age';
SET search_path = ag_catalog, "$user", public;
SELECT create_graph('cyberos_layer2');

CREATE TABLE layer2_memories (
    tenant_id     UUID NOT NULL,
    seq           BIGINT NOT NULL,
    id            UUID NOT NULL,
    kind          TEXT NOT NULL,
    path          TEXT NOT NULL,
    ts_ns         BIGINT NOT NULL,
    embedding     vector(1024) NOT NULL,
    chain_anchor  BYTEA NOT NULL,
    age_state     TEXT NOT NULL DEFAULT 'complete' CHECK (age_state IN ('complete', 'pending_retry')),
    created_at    TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    PRIMARY KEY (tenant_id, seq)
);

CREATE INDEX layer2_memories_hnsw ON layer2_memories USING hnsw (embedding vector_cosine_ops);
CREATE INDEX layer2_memories_kind_idx ON layer2_memories (tenant_id, kind);
CREATE INDEX layer2_memories_age_pending_idx ON layer2_memories (tenant_id) WHERE age_state = 'pending_retry';

ALTER TABLE layer2_memories ENABLE ROW LEVEL SECURITY;
ALTER TABLE layer2_memories FORCE ROW LEVEL SECURITY;
CREATE POLICY layer2_isolation ON layer2_memories
    USING      (tenant_id = current_setting('app.tenant_id', true)::uuid)
    WITH CHECK (tenant_id = current_setting('app.tenant_id', true)::uuid);
-- services/memory/migrations/0002_layer2_cursor.sql
CREATE TABLE layer2_ingest_cursor (
    tenant_id  UUID PRIMARY KEY,
    last_seq   BIGINT NOT NULL,
    updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
// services/memory/src/layer2/ingest.rs
use std::sync::Arc;

pub struct IngestConfig {
    pub bge_url: String,
    pub spacy_url: String,
    pub poll_interval: Duration,
}

pub async fn run_ingest_loop(pool: Arc<PgPool>, config: IngestConfig) -> anyhow::Result<()> {
    loop {
        let tenants = active_tenants(&pool).await?;
        let mut joinset = tokio::task::JoinSet::new();
        for tenant_id in tenants {
            let pool = pool.clone();
            let config = config.clone();
            joinset.spawn(async move {
                ingest_one_tenant(tenant_id, &pool, &config).await
            });
        }
        while let Some(_) = joinset.join_next().await {}
        tokio::time::sleep(config.poll_interval).await;
    }
}

async fn ingest_one_tenant(tenant_id: Uuid, pool: &PgPool, config: &IngestConfig) -> anyhow::Result<()> {
    let cursor = cursor::get(pool, tenant_id).await?;
    let frames = binlog_tail::read_frames_after(tenant_id, cursor).await?;

    for frame in frames {
        // §1 #12: structural validation
        if !chain_anchor::validate_frame(&frame) {
            tracing::warn!(tenant_id = %tenant_id, seq = frame.seq, "corrupted frame; skipping");
            cursor::advance(pool, tenant_id, frame.seq).await?;
            metrics::ingest_failure(tenant_id, "frame_corrupted");
            continue;
        }

        let embedding_result = bge::embed_with_retry(&config.bge_url, &frame.body).await;
        let entity_result = entity_extract::run(&config.spacy_url, &frame.body).await;

        let mut tx = pool.begin().await?;
        rls_set_tenant(&mut tx, tenant_id).await?;

        let embedding = match embedding_result {
            Ok(e) => e,
            Err(_) => {
                metrics::ingest_failure(tenant_id, "bge_down");
                continue;   // skip; retry job picks up later
            }
        };

        let chain_anchor = chain_anchor::compute(&frame);
        sqlx::query("INSERT INTO layer2_memories (tenant_id, seq, id, kind, path, ts_ns, embedding, chain_anchor, age_state)
                     VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
                     ON CONFLICT (tenant_id, seq) DO NOTHING")
            .bind(tenant_id).bind(frame.seq).bind(frame.id)
            .bind(&frame.kind).bind(&frame.path).bind(frame.ts_ns)
            .bind(embedding).bind(&chain_anchor)
            .bind(if entity_result.is_ok() { "complete" } else { "pending_retry" })
            .execute(&mut *tx).await?;

        if let Ok(entities) = entity_result {
            age::insert_entities_and_relations(&mut tx, tenant_id, frame.seq, &entities).await?;
        } else {
            metrics::ingest_failure(tenant_id, "age_fail");
        }

        cursor::advance_in_tx(&mut tx, tenant_id, frame.seq).await?;
        tx.commit().await?;

        metrics::lag_observed(tenant_id, (Utc::now().timestamp_nanos() - frame.ts_ns) / 1_000_000_000);
        metrics::cursor_advanced(tenant_id);
    }
    Ok(())
}
// services/memory/src/layer2/chain_anchor.rs
pub fn compute(frame: &BinlogFrame) -> [u8; 32] {
    let canonical = canonicalise_layer1_row(frame);
    sha256(&canonical)
}

pub fn verify(stored_anchor: &[u8; 32], current_layer1_row: &Layer1Row) -> bool {
    let recomputed = sha256(&canonicalise_layer1_row_from(current_layer1_row));
    recomputed == *stored_anchor
}

pub fn validate_frame(frame: &BinlogFrame) -> bool {
    // CRC check + format validation
    let computed_crc = crc32c::crc32c(&frame.payload);
    computed_crc == frame.crc
}

§4 — Acceptance criteria

  1. Append a Layer 1 row → within 1s, Layer 2 row exists — synthetic write to binlog; query layer2_memories within 1s; row present.
  2. Restart ingest process → resumes from cursor (no re-ingest of prior rows; no duplicate INSERT).
  3. Cross-tenant query: tenant A's ingest doesn't write to tenant B's index — RLS enforces; tenant B SELECT from layer2_memories returns 0 of A's rows.
  4. Idempotent: re-ingest same seq → UPSERT, single row in layer2_memories.
  5. Lag metric < 1s p95 under steady-state load — emit 1000 frames over 100s; histogram p95 < 1s.
  6. Embedding sidecar (TASK-AI-019) failure → ingest retries with backoff; doesn't crash — kill BGE; ingest queues retries; continues.
  7. AGE entity-extraction failure → pgvector row written; AGE row marked state: pending_retry — kill spaCy; pgvector still inserts; age_state='pending_retry'.
  8. CI rebuild test (TASK-MEMORY-102) recreates Layer 2 from Layer 1 in <30min for slice-1 data volume.
  9. Chain_anchor verifies post-write — Layer 1 row at seq N; recompute chain_anchor; equals stored value.
  10. Chain_anchor mismatch flagged sev-1 — manually corrupt Layer 1; read path detects mismatch; sev-1 OBS event.
  11. Frame-level validation skip on CRC failure — inject CRC-corrupt frame; ingest skips + advances cursor + sev-2 log.
  12. Concurrent multi-tenant ingest — 10 tenants × 100 rows; all 10 progress concurrently; per-tenant lag < 2s.
  13. Graceful shutdown — SIGTERM; in-flight tx commits; cursor saved; restart resumes correctly.
  14. --rebuild flag re-ingests from seq 0 — flag passed; cursor reset; full re-ingest.
  15. OTel metrics emit per ingest event.
  16. RLS on layer2_memories — INSERT with wrong tenant_id rejected (42501).
  17. HNSW index used for vector search — EXPLAIN shows Index Scan using layer2_memories_hnsw.

§5 — Verification

#[tokio::test]
async fn append_to_layer1_visible_in_layer2_within_1s() {
    let pool = test_pool().await;
    let tenant = test_tenant().await;
    spawn_ingest_loop(pool.clone()).await;

    let seq = test_helper::append_layer1_row(tenant, "decisions/test.md", "test body").await;
    let t0 = std::time::Instant::now();

    while t0.elapsed() < Duration::from_secs(2) {
        let row = sqlx::query("SELECT * FROM layer2_memories WHERE tenant_id = $1 AND seq = $2")
            .bind(tenant).bind(seq).fetch_optional(&pool).await.unwrap();
        if row.is_some() { return; }
        tokio::time::sleep(Duration::from_millis(50)).await;
    }
    panic!("Layer 2 row not visible within 2s");
}

#[tokio::test]
async fn restart_resumes_from_cursor() {
    let pool = test_pool().await;
    let tenant = test_tenant().await;
    let handle = spawn_ingest_loop(pool.clone()).await;
    test_helper::append_layer1_rows(tenant, 100).await;
    tokio::time::sleep(Duration::from_secs(2)).await;

    let count_before: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM layer2_memories").fetch_one(&pool).await.unwrap();

    handle.shutdown().await;
    let _ = spawn_ingest_loop(pool.clone()).await;
    tokio::time::sleep(Duration::from_secs(2)).await;

    let count_after: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM layer2_memories").fetch_one(&pool).await.unwrap();
    assert_eq!(count_before, count_after, "no duplicate rows after restart");
}

#[tokio::test]
async fn cross_tenant_isolation() {
    let pool = test_pool().await;
    let a = test_tenant().await;
    let b = test_tenant().await;
    test_helper::append_layer1_row(a, "decisions/a.md", "tenant A").await;
    test_helper::append_layer1_row(b, "decisions/b.md", "tenant B").await;
    tokio::time::sleep(Duration::from_secs(2)).await;

    rls::with_tenant(&pool, a, |tx| async move {
        let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM layer2_memories").fetch_one(&mut **tx).await.unwrap();
        assert_eq!(count, 1, "tenant A should see only their row");
    }).await;
}

#[tokio::test]
async fn bge_down_triggers_retry_no_crash() {
    bge_test_helper::stop();
    let pool = test_pool().await;
    let tenant = test_tenant().await;
    spawn_ingest_loop(pool.clone()).await;
    test_helper::append_layer1_row(tenant, "decisions/x.md", "body").await;
    tokio::time::sleep(Duration::from_secs(3)).await;
    let metric: u64 = otel_test_helper::counter_value("memory_layer2_ingest_failures_total", &[("reason", "bge_down")]);
    assert!(metric > 0);
    // ingest still alive
    bge_test_helper::start();
    tokio::time::sleep(Duration::from_secs(2)).await;
    let row = sqlx::query("SELECT 1 FROM layer2_memories WHERE tenant_id = $1").bind(tenant).fetch_optional(&pool).await.unwrap();
    assert!(row.is_some(), "ingest recovered after BGE restart");
}

#[tokio::test]
async fn age_failure_pgvector_still_writes() {
    spacy_test_helper::stop();
    let pool = test_pool().await;
    let tenant = test_tenant().await;
    spawn_ingest_loop(pool.clone()).await;
    test_helper::append_layer1_row(tenant, "decisions/x.md", "body").await;
    tokio::time::sleep(Duration::from_secs(2)).await;
    let row: (String,) = sqlx::query_as("SELECT age_state FROM layer2_memories WHERE tenant_id = $1").bind(tenant).fetch_one(&pool).await.unwrap();
    assert_eq!(row.0, "pending_retry");
}

#[tokio::test]
async fn chain_anchor_mismatch_sev1() {
    let pool = test_pool().await;
    let tenant = test_tenant().await;
    test_helper::append_layer1_row(tenant, "decisions/x.md", "body").await;
    tokio::time::sleep(Duration::from_secs(2)).await;

    test_helper::corrupt_layer1_row_at_seq(tenant, 1, "TAMPERED").await;

    // Read path queries chain_anchor
    let result = layer2_search::query_with_anchor_check(tenant, "body").await;
    assert!(matches!(result, Err(SearchError::ChainAnchorMismatch)));
    let metric: u64 = otel_test_helper::counter_value("memory_layer2_ingest_failures_total", &[("reason", "chain_anchor_mismatch")]);
    assert!(metric > 0);
}

#[tokio::test]
async fn idempotent_upsert() {
    let pool = test_pool().await;
    let tenant = test_tenant().await;
    spawn_ingest_loop(pool.clone()).await;
    let seq = test_helper::append_layer1_row(tenant, "decisions/x.md", "body").await;
    tokio::time::sleep(Duration::from_secs(2)).await;

    test_helper::reset_cursor(tenant, seq - 1).await;   // simulate restart-mid-batch
    tokio::time::sleep(Duration::from_secs(2)).await;

    let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM layer2_memories WHERE tenant_id = $1 AND seq = $2")
        .bind(tenant).bind(seq).fetch_one(&pool).await.unwrap();
    assert_eq!(count, 1);
}

#[tokio::test]
async fn rls_blocks_wrong_tenant_insert() {
    let pool = test_pool_as_cyberos_app().await;
    let a = test_tenant().await;
    let b = test_tenant().await;
    let err = rls::with_tenant(&pool, a, |tx| async move {
        sqlx::query("INSERT INTO layer2_memories (tenant_id, seq, id, kind, path, ts_ns, embedding, chain_anchor) VALUES ($1, 1, gen_random_uuid(), 'k', 'p', 0, $2, $3)")
            .bind(b)   // wrong tenant
            .bind(vec![0.0f32; 1024])
            .bind(vec![0u8; 32])
            .execute(&mut **tx).await
    }).await.expect_err("expected RLS violation");
    let violation = rls::classify_pg_error(&err).expect("expected RlsCheckViolation");
}

§6 — Implementation skeleton

See §3.


§7 — Dependencies


§8 — Example payloads

Layer 2 row (selected fields)

tenant_id:    550e8400-...
seq:          12345
id:           7e57c0de-...
kind:         decisions
path:         memories/decisions/2026-05-15-revoke-policy.md
ts_ns:        1747526400000000000
embedding:    [0.012, -0.034, ...]   (1024 dims)
chain_anchor: a3f9c8d7e6b5a4f3...
age_state:    complete

Cursor row

tenant_id:  550e8400-...
last_seq:   12345
updated_at: 2026-05-15T14:00:00.500Z

Sev-1 chain_anchor mismatch event

sev-1  memory_layer2_chain_anchor_mismatch
       tenant_id=550e... seq=12345 stored=a3f9c8d7... recomputed=ff00aa11...
       Layer 1 may be tampered; halting reads against this seq

§9 — Open questions

All resolved. Deferred:


§10 — Failure modes inventory

FailureDetectionOutcomeRecovery
BGE sidecar downreqwest errorIngest retries with backoff (5 attempts) then pending_embed_retry; metric bge_down incrementsSelf-heals when sidecar up; backfill job processes pending
AGE entity extraction failsspaCy errorpgvector row written with age_state=pending_retry; AGE backfill runs separatelySelf-heals; backfill job
Postgres deadlocksqlx errortx rollback + retry onceSelf-heals
Layer 1 row corruption (frame CRC)validate_frame failsSkip + advance cursor + sev-2 logEngineer investigates
Layer 1 row corruption (chain_anchor mismatch on read)Read-path verifysev-1 + drop resultEngineer investigates
Cursor table corruptionsqlx errorRestart failsOperator restores from backup
Concurrent ingest race (same tenant, two processes)UPSERT idempotentNo duplicatesBy design
Postgres unavailablesqlx connectIngest waits + retriesSelf-heals
RLS violation (cross-tenant write)postgres 42501Sev-1 alarmInvestigate ingest code
Restart mid-batchUPSERT idempotentNo duplicates; cursor advances correctlyBy design
HNSW index fragmentationSlow vector queriessev-3 alarm; REINDEXOperator action
Tenant deleted while ingest runningtenants table FK violationTenant ingest task exits cleanly; metricBy design
Embedding dim mismatch (BGE returns wrong dim)sqlx schema checkINSERT fails; sev-2Investigate TASK-AI-019
Frame deserialise fails (msgpack error)parse errorSkip + sev-2Engineer investigates Layer 1 writer
Disk full on layer2 partitionINSERT failssev-1; ingest pausesOperator extends disk
AGE extension missingstartup CREATE GRAPH failsService refuses to startOperator runs migration
pgvector extension missingstartup CREATE EXTENSION failsService refuses to startOperator runs migration
Tokio task panicobservability via tracingOther tenants unaffected; restart taskInvestigate
Lag > 1s p95 sustainedOTel histogram alarmsev-3Investigate BGE OR DB load
--rebuild on productionflag check + confirmationOperator must explicitly confirm in productionBy design (slice 2)

§11 — Notes


End of TASK-MEMORY-101. Status: done (implemented 2026-05-23).

As built (2026-07-02)

Apache AGE was removed; the layer-2 graph is the relational l2_edge table + pgvector. Any AGE/CREATE GRAPH references above are historical.