From 23041b0556a696ce1ea56feee4e80f5d29a02187 Mon Sep 17 00:00:00 2001 From: Jonathan Haas Date: Tue, 4 Aug 2026 18:42:40 -0700 Subject: [PATCH] feat(merlin): add delivery receipts --- README.md | 4 ++++ crates/corpus-core/src/dto.rs | 3 +++ crates/corpus-core/src/merlin.rs | 25 +++++++++++++++++++++++++ crates/corpus-core/tests/merlin.rs | 5 +++++ docs/openapi.json | 23 ++++++++++++++++++++++- 5 files changed, 59 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 5d30078..f38f498 100644 --- a/README.md +++ b/README.md @@ -45,6 +45,10 @@ telemetry separately from verified artifact occurrences, so a process event can be joined to later byte capture without treating a path or event as a file hash. Use a dedicated `CORPUS_MERLIN_INGEST_TOKEN`; do not share the Corpus admin token or Merlin's sensor sync key. +Each batch returns a versioned delivery receipt containing the canonical segment +digest, accepted/duplicate counts, and a stable segment ID. Merlin validates +that receipt before marking its local delivery durable; malformed or stale +receipts are retried rather than silently acknowledged. ## What it does diff --git a/crates/corpus-core/src/dto.rs b/crates/corpus-core/src/dto.rs index a036371..d6a17cb 100644 --- a/crates/corpus-core/src/dto.rs +++ b/crates/corpus-core/src/dto.rs @@ -65,7 +65,10 @@ pub struct MerlinSegmentRequest { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct MerlinSegmentResponse { pub schema_version: i32, + pub receipt_version: i32, pub segment_id: Uuid, + pub segment_sha256: String, + pub status: String, pub accepted_events: usize, pub duplicate_events: usize, } diff --git a/crates/corpus-core/src/merlin.rs b/crates/corpus-core/src/merlin.rs index 1fa4009..ab72d94 100644 --- a/crates/corpus-core/src/merlin.rs +++ b/crates/corpus-core/src/merlin.rs @@ -148,7 +148,15 @@ pub async fn ingest_segment( Ok(MerlinSegmentResponse { schema_version: MERLIN_SCHEMA_VERSION, + receipt_version: 1, segment_id, + segment_sha256, + status: if accepted_events == 0 { + "duplicate" + } else { + "accepted" + } + .into(), accepted_events, duplicate_events: req.events.len() - accepted_events, }) @@ -349,3 +357,20 @@ mod tests { assert!(validate_request(&object).is_err()); } } + +#[test] +fn merlin_receipt_serializes_delivery_contract() { + let receipt = MerlinSegmentResponse { + schema_version: MERLIN_SCHEMA_VERSION, + receipt_version: 1, + segment_id: Uuid::nil(), + segment_sha256: "a".repeat(64), + status: "accepted".into(), + accepted_events: 2, + duplicate_events: 1, + }; + let value = serde_json::to_value(receipt).unwrap(); + assert_eq!(value["receipt_version"], 1); + assert_eq!(value["segment_sha256"], "a".repeat(64)); + assert_eq!(value["status"], "accepted"); +} diff --git a/crates/corpus-core/tests/merlin.rs b/crates/corpus-core/tests/merlin.rs index 149c53f..ac136be 100644 --- a/crates/corpus-core/tests/merlin.rs +++ b/crates/corpus-core/tests/merlin.rs @@ -51,6 +51,9 @@ async fn merlin_segment_ingest_is_idempotent_and_rejects_digest_conflicts() { assert_eq!(first.accepted_events, 2); assert_eq!(first.duplicate_events, 0); + assert_eq!(first.receipt_version, 1); + assert_eq!(first.segment_sha256, "a".repeat(64)); + assert_eq!(first.status, "accepted"); let replay = merlin::ingest_segment(&pool, tenant_id, &req) .await .unwrap(); @@ -58,6 +61,8 @@ async fn merlin_segment_ingest_is_idempotent_and_rejects_digest_conflicts() { assert_eq!(replay.accepted_events, 0); assert_eq!(replay.duplicate_events, 2); + assert_eq!(replay.receipt_version, 1); + assert_eq!(replay.status, "duplicate"); let observations = merlin::list_observations(&pool, tenant_id, Some("merlin-test-host"), 100) .await .unwrap(); diff --git a/docs/openapi.json b/docs/openapi.json index 3bce808..199cc81 100644 --- a/docs/openapi.json +++ b/docs/openapi.json @@ -105,7 +105,28 @@ "summary": "Ingest replay-safe raw Merlin telemetry", "security": [{ "merlinBearer": [] }, { "adminBearer": [] }], "parameters": [{ "$ref": "#/components/parameters/Tenant" }], - "responses": { "200": { "description": "Accepted and duplicate event counts" } } + "responses": { + "200": { + "description": "Delivery receipt with accepted and duplicate event counts", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["schema_version", "receipt_version", "segment_id", "segment_sha256", "status", "accepted_events", "duplicate_events"], + "properties": { + "schema_version": { "type": "integer", "example": 1 }, + "receipt_version": { "type": "integer", "example": 1 }, + "segment_id": { "type": "string", "format": "uuid" }, + "segment_sha256": { "type": "string", "pattern": "^[0-9a-f]{64}$" }, + "status": { "type": "string", "enum": ["accepted", "duplicate"] }, + "accepted_events": { "type": "integer", "minimum": 0 }, + "duplicate_events": { "type": "integer", "minimum": 0 } + } + } + } + } + } + } } }, "/api/v1/integrations/merlin/observations": {