Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
3 changes: 3 additions & 0 deletions crates/corpus-core/src/dto.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
Expand Down
25 changes: 25 additions & 0 deletions crates/corpus-core/src/merlin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
})
Expand Down Expand Up @@ -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");
}
5 changes: 5 additions & 0 deletions crates/corpus-core/tests/merlin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,13 +51,18 @@ 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();
assert_eq!(replay.segment_id, first.segment_id);
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();
Expand Down
23 changes: 22 additions & 1 deletion docs/openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
Loading