diff --git a/src/ai_transcript_ingest/legacy_receipt.rs b/src/ai_transcript_ingest/legacy_receipt.rs index 142a6c9ee..c646a573c 100644 --- a/src/ai_transcript_ingest/legacy_receipt.rs +++ b/src/ai_transcript_ingest/legacy_receipt.rs @@ -14,15 +14,64 @@ impl std::fmt::Display for IdempotencyConflict { impl std::error::Error for IdempotencyConflict {} -fn envelope_fingerprint(envelope: &EvidenceEnvelope) -> anyhow::Result { +fn v2_fingerprint_with_locator( + envelope: &EvidenceEnvelope, + locator: &str, +) -> anyhow::Result { // Supplemental display metadata can change without a transcript revision. + // Keep the persisted v2 format rollback-compatible with older Cortex + // releases; locator mutability is handled during replay verification. let mut evidence = envelope.clone(); evidence.source.title = None; evidence.source.title_provenance = None; + evidence.source.locator = locator.to_string(); let encoded = serde_json::to_vec(&evidence)?; Ok(format!("evidence-v2:sha256:{:x}", Sha256::digest(encoded))) } +fn envelope_fingerprint(envelope: &EvidenceEnvelope) -> anyhow::Result { + v2_fingerprint_with_locator(envelope, &envelope.source.locator) +} + +fn transient_v3_fingerprint(envelope: &EvidenceEnvelope) -> anyhow::Result { + // A pre-review patched deployment briefly wrote v3 receipts whose only + // semantic difference was excluding the movable locator. Continue to read + // those receipts and lazily rebind them to the durable v2 format. + let mut evidence = envelope.clone(); + evidence.source.title = None; + evidence.source.title_provenance = None; + evidence.source.locator.clear(); + Ok(format!( + "evidence-v3:sha256:{:x}", + Sha256::digest(serde_json::to_vec(&evidence)?) + )) +} + +fn canonical_v2_fingerprint( + envelope: &EvidenceEnvelope, + stored_locator: Option<&str>, +) -> anyhow::Result> { + let Some(locator) = stored_locator else { + return Ok(None); + }; + Ok(Some(v2_fingerprint_with_locator(envelope, locator)?)) +} + +fn v2_fingerprint_matches( + envelope: &EvidenceEnvelope, + stored_locator: Option<&str>, + previous: &str, +) -> anyhow::Result { + Ok(canonical_v2_fingerprint(envelope, stored_locator)?.as_deref() == Some(previous)) +} + +fn transient_v3_fingerprint_matches( + envelope: &EvidenceEnvelope, + previous: &str, +) -> anyhow::Result { + Ok(previous == transient_v3_fingerprint(envelope)?) +} + fn receipt_key(forwarder_identity: &str, source_record_id: &str, shared_bearer: bool) -> String { if shared_bearer { source_record_id.to_owned() @@ -34,9 +83,10 @@ fn receipt_key(forwarder_identity: &str, source_record_id: &str, shared_bearer: } } -/// Reconstruct only mutable titles from the stored row, then check the old -/// full-envelope hash. This also preserves timestamp-less exact replays: the -/// hash, unlike a canonical log timestamp, retains their original `None`. +/// Reconstruct mutable display/location metadata from the stored row, then +/// check the old full-envelope hash. This also preserves timestamp-less exact +/// replays: the hash, unlike a canonical log timestamp, retains their original +/// `None`. fn old_fingerprint_matches( tx: &rusqlite::Transaction<'_>, key: &str, @@ -44,7 +94,8 @@ fn old_fingerprint_matches( previous: &str, ) -> anyhow::Result { // An exact old hash is sufficient even if bounded log metadata omitted - // its source fields. Canonical metadata is needed only for a title change. + // its source fields. Canonical metadata is needed only when mutable title + // or locator metadata changed. if previous == format!("sha256:{:x}", Sha256::digest(serde_json::to_vec(envelope)?)) { return Ok(true); } @@ -65,6 +116,7 @@ fn old_fingerprint_matches( let mut original = envelope.clone(); original.source.title = source.title; original.source.title_provenance = source.title_provenance; + original.source.locator = source.locator; Ok(previous == format!( "sha256:{:x}", @@ -86,7 +138,7 @@ fn legacy_receipt_matches( "SELECT r.envelope_version, r.provider, r.source_identity, r.source_epoch, r.source_revision, l.timestamp, l.message, l.ai_project, l.ai_session_id, - l.ai_transcript_path, l.metadata_json + l.metadata_json FROM ai_transcript_forward_receipts r JOIN logs l ON l.id = r.log_id WHERE r.source_record_id = ?1", @@ -103,7 +155,6 @@ fn legacy_receipt_matches( row.get::<_, Option>(7)?, row.get::<_, Option>(8)?, row.get::<_, Option>(9)?, - row.get::<_, Option>(10)?, )) }, ) @@ -118,7 +169,6 @@ fn legacy_receipt_matches( message, ai_project, ai_session_id, - locator, metadata_json, )) = stored else { @@ -136,9 +186,11 @@ fn legacy_receipt_matches( let mut incoming_source = envelope.source.clone(); incoming_source.title = None; incoming_source.title_provenance = None; + incoming_source.locator.clear(); if let Some(source) = &mut stored_source { source.title = None; source.title_provenance = None; + source.locator.clear(); } let stored_capabilities = metadata .get("capabilities") @@ -171,7 +223,6 @@ fn legacy_receipt_matches( && message == envelope.message && ai_project == envelope.ai_project && ai_session_id == envelope.ai_session_id - && locator.as_deref() == Some(envelope.source.locator.as_str()) // A legacy canonical row does not record whether its timestamp came // from the source envelope or the receiver clock. Requiring the replay // to supply the stored value avoids silently binding an ambiguous @@ -246,32 +297,71 @@ fn insert_envelopes_with_identity( )?; let already_accepted = tx .query_row( - "SELECT request_fingerprint FROM ai_transcript_forward_receipts WHERE source_record_id = ?1", + "SELECT r.request_fingerprint, l.ai_transcript_path + FROM ai_transcript_forward_receipts r + JOIN logs l ON l.id = r.log_id + WHERE r.source_record_id = ?1", [&stored_receipt_key], - |row| row.get::<_, Option>(0), + |row| { + Ok(( + row.get::<_, Option>(0)?, + row.get::<_, Option>(1)?, + )) + }, ) .optional()?; - if let Some(previous_fingerprint) = already_accepted { + if let Some((previous_fingerprint, stored_locator)) = already_accepted { if previous_fingerprint.as_deref() != Some(request_fingerprint.as_str()) { - // Old hashes included titles. Rebind only after comparing all - // immutable evidence against the canonical row, never merely - // because the caller reused an existing source-record ID. - let matches = match previous_fingerprint.as_deref() { + // Compatibility fingerprints are validated against the + // canonical row before any rebinding. Persist v2 so a rollback + // still recognizes the durable receipt format; locator-move + // tolerance itself is provided by this version's verifier. + let (matches, replacement_fingerprint) = match previous_fingerprint.as_deref() { + Some(previous) if previous.starts_with("evidence-v3:sha256:") => { + let matches = transient_v3_fingerprint_matches(&envelope, previous)?; + let replacement = if matches { + canonical_v2_fingerprint(&envelope, stored_locator.as_deref())? + } else { + None + }; + (matches, replacement) + } + Some(previous) if previous.starts_with("evidence-v2:sha256:") => ( + v2_fingerprint_matches(&envelope, stored_locator.as_deref(), previous)?, + None, + ), Some(previous) if previous.starts_with("sha256:") => { - old_fingerprint_matches(&tx, &stored_receipt_key, &envelope, previous)? + let matches = + old_fingerprint_matches(&tx, &stored_receipt_key, &envelope, previous)?; + let replacement = if matches { + canonical_v2_fingerprint(&envelope, stored_locator.as_deref())? + } else { + None + }; + (matches, replacement) + } + Some(_) => (false, None), + None => { + let matches = legacy_receipt_matches(&tx, &stored_receipt_key, &envelope)?; + let replacement = if matches { + canonical_v2_fingerprint(&envelope, stored_locator.as_deref())? + } else { + None + }; + (matches, replacement) } - Some(_) => false, - None => legacy_receipt_matches(&tx, &stored_receipt_key, &envelope)?, }; if !matches { return Err(IdempotencyConflict.into()); } - tx.execute( - "UPDATE ai_transcript_forward_receipts - SET request_fingerprint = ?2 - WHERE source_record_id = ?1", - rusqlite::params![stored_receipt_key, request_fingerprint], - )?; + if let Some(replacement_fingerprint) = replacement_fingerprint { + tx.execute( + "UPDATE ai_transcript_forward_receipts + SET request_fingerprint = ?2 + WHERE source_record_id = ?1", + rusqlite::params![stored_receipt_key, replacement_fingerprint], + )?; + } } receipts.push(AiTranscriptReceipt { source_record_id: envelope.source_record_id, diff --git a/src/ai_transcript_replay_tests.rs b/src/ai_transcript_replay_tests.rs index a37462f83..11230a1c7 100644 --- a/src/ai_transcript_replay_tests.rs +++ b/src/ai_transcript_replay_tests.rs @@ -1,5 +1,222 @@ use super::*; +fn receipt_v2_fingerprint(record: &serde_json::Value) -> String { + use sha2::{Digest, Sha256}; + + let envelope: EvidenceEnvelope = serde_json::from_value(record["envelope"].clone()).unwrap(); + let mut evidence = scrub_envelope(envelope).unwrap(); + evidence.source.title = None; + evidence.source.title_provenance = None; + format!( + "evidence-v2:sha256:{:x}", + Sha256::digest(serde_json::to_vec(&evidence).unwrap()) + ) +} + +fn transient_receipt_v3_fingerprint(record: &serde_json::Value) -> String { + use sha2::{Digest, Sha256}; + + let envelope: EvidenceEnvelope = serde_json::from_value(record["envelope"].clone()).unwrap(); + let mut evidence = scrub_envelope(envelope).unwrap(); + evidence.source.title = None; + evidence.source.title_provenance = None; + evidence.source.locator.clear(); + format!( + "evidence-v3:sha256:{:x}", + Sha256::digest(serde_json::to_vec(&evidence).unwrap()) + ) +} + +#[tokio::test] +async fn archived_transcript_replay_accepts_changed_locator_and_preserves_v2_receipt() { + let (app, dir) = test_app(Some("secret")); + let original = sample_record(); + let original_v2 = receipt_v2_fingerprint(&original); + + let first = app + .clone() + .oneshot(transcript_request( + json!({"records": [original]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(first.status(), StatusCode::OK); + + let conn = rusqlite::Connection::open(dir.path().join("ai-transcript-ingest-test.db")).unwrap(); + let fingerprint: String = conn + .query_row( + "SELECT request_fingerprint FROM ai_transcript_forward_receipts", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(fingerprint, original_v2); + + for locator in ["f", "1"] { + let mut archived = sample_record(); + archived["envelope"]["source"]["locator"] = json!(format!("sha256:{}", locator.repeat(64))); + let replay = app + .clone() + .oneshot(transcript_request( + json!({"records": [archived]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(replay.status(), StatusCode::OK); + let body = axum::body::to_bytes(replay.into_body(), usize::MAX) + .await + .unwrap(); + let body: serde_json::Value = serde_json::from_slice(&body).unwrap(); + assert_eq!(body["receipts"][0]["disposition"], "duplicate"); + } + + let fingerprint: String = conn + .query_row( + "SELECT request_fingerprint FROM ai_transcript_forward_receipts", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(fingerprint, original_v2); + + let mut changed_evidence = sample_record(); + changed_evidence["envelope"]["source"]["locator"] = json!(format!("sha256:{}", "1".repeat(64))); + changed_evidence["envelope"]["message"] = json!("different evidence"); + let conflict = app + .oneshot(transcript_request( + json!({"records": [changed_evidence]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(conflict.status(), StatusCode::CONFLICT); + let (log_count, fingerprint): (i64, String) = conn + .query_row( + "SELECT (SELECT COUNT(*) FROM logs), request_fingerprint + FROM ai_transcript_forward_receipts", + [], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .unwrap(); + assert_eq!(log_count, 1); + assert_eq!(fingerprint, original_v2); +} + +#[tokio::test] +async fn transient_v3_receipt_rebinds_to_rollback_safe_v2() { + let (app, dir) = test_app(Some("secret")); + let original = sample_record(); + let original_v2 = receipt_v2_fingerprint(&original); + let transient_v3 = transient_receipt_v3_fingerprint(&original); + + let first = app + .clone() + .oneshot(transcript_request( + json!({"records": [original]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(first.status(), StatusCode::OK); + + let conn = rusqlite::Connection::open(dir.path().join("ai-transcript-ingest-test.db")).unwrap(); + conn.execute( + "UPDATE ai_transcript_forward_receipts SET request_fingerprint = ?1", + [&transient_v3], + ) + .unwrap(); + + let mut changed = sample_record(); + changed["envelope"]["message"] = json!("changed transient v3 evidence"); + let conflict = app + .clone() + .oneshot(transcript_request( + json!({"records": [changed]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(conflict.status(), StatusCode::CONFLICT); + let fingerprint: String = conn + .query_row( + "SELECT request_fingerprint FROM ai_transcript_forward_receipts", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(fingerprint, transient_v3); + + let mut archived = sample_record(); + archived["envelope"]["source"]["locator"] = json!(format!("sha256:{}", "f".repeat(64))); + let replay = app + .clone() + .oneshot(transcript_request( + json!({"records": [archived]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(replay.status(), StatusCode::OK); + + let fingerprint: String = conn + .query_row( + "SELECT request_fingerprint FROM ai_transcript_forward_receipts", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(fingerprint, original_v2); + + let mut moved_again = sample_record(); + moved_again["envelope"]["source"]["locator"] = json!(format!("sha256:{}", "1".repeat(64))); + let replay = app + .oneshot(transcript_request( + json!({"records": [moved_again]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(replay.status(), StatusCode::OK); +} + +#[tokio::test] +async fn legacy_null_fingerprint_receipt_accepts_locator_move_and_rebinds_v2() { + let (app, dir) = test_app(Some("secret")); + let original = sample_record(); + let original_v2 = receipt_v2_fingerprint(&original); + + let first = app + .clone() + .oneshot(transcript_request( + json!({"records": [original]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(first.status(), StatusCode::OK); + + let conn = rusqlite::Connection::open(dir.path().join("ai-transcript-ingest-test.db")).unwrap(); + conn.execute( + "UPDATE ai_transcript_forward_receipts SET request_fingerprint = NULL", + [], + ) + .unwrap(); + + let mut archived = sample_record(); + archived["envelope"]["source"]["locator"] = json!(format!("sha256:{}", "f".repeat(64))); + let replay = app + .clone() + .oneshot(transcript_request( + json!({"records": [archived]}).to_string(), + )) + .await + .unwrap(); + assert_eq!(replay.status(), StatusCode::OK); + + let fingerprint: String = conn + .query_row( + "SELECT request_fingerprint FROM ai_transcript_forward_receipts", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(fingerprint, original_v2); +} + #[tokio::test] async fn title_changes_and_missing_metadata_are_duplicate_replays() { let (app, dir) = test_app(Some("secret")); @@ -30,7 +247,7 @@ async fn title_changes_and_missing_metadata_are_duplicate_replays() { } #[tokio::test] -async fn old_full_envelope_receipt_accepts_title_only_changes_but_not_evidence_changes() { +async fn old_full_envelope_receipt_accepts_title_and_locator_changes_but_not_evidence_changes() { check_old_receipt_replay(false).await; check_old_receipt_replay(true).await; } @@ -42,6 +259,7 @@ async fn check_old_receipt_replay(timestamp_absent: bool) { if timestamp_absent { original["envelope"]["timestamp"] = serde_json::Value::Null; } + let expected_v2 = receipt_v2_fingerprint(&original); let envelope: EvidenceEnvelope = serde_json::from_value(original["envelope"].clone()).unwrap(); let scrubbed = scrub_envelope(envelope).unwrap(); let old_hash = format!( @@ -93,6 +311,7 @@ async fn check_old_receipt_replay(timestamp_absent: bool) { replay["envelope"]["timestamp"] = serde_json::Value::Null; } replay["envelope"]["source"]["title"] = json!("Renamed"); + replay["envelope"]["source"]["locator"] = json!(format!("sha256:{}", "f".repeat(64))); let renamed = app .clone() .oneshot(transcript_request( @@ -101,6 +320,14 @@ async fn check_old_receipt_replay(timestamp_absent: bool) { .await .unwrap(); assert_eq!(renamed.status(), StatusCode::OK); + let rebound: String = conn + .query_row( + "SELECT request_fingerprint FROM ai_transcript_forward_receipts", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(rebound, expected_v2); replay["envelope"]["message"] = json!("different evidence"); let changed = app .oneshot(transcript_request(json!({"records": [replay]}).to_string()))