diff --git a/Cargo.lock b/Cargo.lock index 3cae82f..9d1cfff 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -187,8 +187,10 @@ checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" name = "content-telemetry-core" version = "0.1.0" dependencies = [ + "base64", "chrono", "dashmap", + "rand", "serde", "serde_json", "sqlx", diff --git a/Cargo.toml b/Cargo.toml index 8b67f20..9334a47 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -15,8 +15,10 @@ repository = "https://github.com/openattribution-org/telemetry-server" content-telemetry-core = { path = "crates/content-telemetry-core" } axum = "0.8" +base64 = "0.22" chrono = { version = "0.4", features = ["serde"] } dashmap = "6" +rand = "0.8" serde = { version = "1", features = ["derive"] } serde_json = "1" sqlx = { version = "0.8", features = [ diff --git a/README.md b/README.md index c1dd4d1..9a82d77 100644 --- a/README.md +++ b/README.md @@ -52,13 +52,13 @@ curl -s -X POST localhost:8080/events \ -H 'content-type: application/json' -H "x-organization-id: $ORG" \ -d '{ "document_type": "event_batch", - "schema_version": "0.1", + "schema_version": "1.0", "session_id": "'"$SESSION"'", "events": [ {"type": "content_grounded", "timestamp": "2026-08-06T12:00:00Z", "content_url": "https://example.com/article-1", - "data": {"grounding_scope": "session"}} + "data": {"scope": "session"}} ]}' curl -s "localhost:8080/sessions/$SESSION/document" -H "x-organization-id: $ORG" @@ -78,7 +78,7 @@ disclosure and document materialisation — against a live build. | `GET` | `/sessions/{id}/document` | The session as a standard session document | | `POST` | `/events` | Ingest events, singly or in batches of up to 500 | | `POST` | `/click-tokens` | Mint a ctx token for a click-out | -| `GET` | `/ctx/{token}` | Resolve a ctx token to its click manifest | +| `GET` | `/ctx/{token}` | Resolve a ctx token to its click context | Events bind to a session in this order: the event's own `session_id`, the batch's, the event's `ctx_token`, the batch's, and otherwise nothing — @@ -103,8 +103,8 @@ No handler changes, because no handler knows how the answer was reached. The exception is `GET /ctx/{token}`, deliberately unauthenticated: the destination of a click-out has no account here. The token is the credential, -and what it discloses is bounded by two-sided consent and by a manifest shape -that never contains the session id. +and what it discloses is bounded by two-sided consent and by a click-context +shape that never contains the session id. ## The database @@ -123,12 +123,19 @@ others. ## Conformance -Ingest enforces the v1 structural rules: `content_cited`, `content_presented` -and `content_reproduced` carry an event id and an `output_id`; -`content_engaged` carries the `presentation_id` of the presentation it acted -on; content events carry a resolvable `content_url` or `content_id`. The -`content_displayed` type v1 withdrew is refused rather than rewritten into a -claim the emitter never made. +Ingest enforces the v1 structural rules on documents declaring +`schema_version` `"1.0"`: `content_cited` and `content_presented` carry an +event id and an `output_id`; `content_grounded` carries `data.scope`; +`content_cited` carries `data.citation_type`; `content_retrieved` carries +`source_role`; `content_engaged` carries the `presentation_id` of the +presentation it acted on; content events carry a resolvable `content_url` or +`content_id`. Documents still declaring `"0.1"` are accepted for the +transition and normalised under the specification's migration rules instead. +The `content_displayed` type v1 withdrew is refused rather than rewritten +into a claim the emitter never made. `content_reproduced`, which never made +it out of the pre-release draft, is no longer a core type: rows stored under +it are treated like any other extension event and quarantined under the +document's `extensions` member. Two things are stored as given rather than validated, because the specification says a consumer must not reject a document over either: diff --git a/crates/content-telemetry-core/Cargo.toml b/crates/content-telemetry-core/Cargo.toml index 916c739..4c61570 100644 --- a/crates/content-telemetry-core/Cargo.toml +++ b/crates/content-telemetry-core/Cargo.toml @@ -10,8 +10,10 @@ repository.workspace = true path = "src/lib.rs" [dependencies] +base64.workspace = true chrono.workspace = true dashmap.workspace = true +rand.workspace = true serde.workspace = true serde_json.workspace = true sqlx.workspace = true diff --git a/crates/content-telemetry-core/migrations/0003_ctx_token_terms_ref.sql b/crates/content-telemetry-core/migrations/0003_ctx_token_terms_ref.sql new file mode 100644 index 0000000..f3a1984 --- /dev/null +++ b/crates/content-telemetry-core/migrations/0003_ctx_token_terms_ref.sql @@ -0,0 +1,17 @@ +-- v1 event fields the original schema predates. +-- +-- ctx_token (specification section 5.2, 7.4.1): on an agent-reported +-- content_engaged event, the click token minted for that engagement's +-- presentation, recorded so destination reports can be joined to it at +-- resolution. The token-to-presentation binding is issuer state and this +-- column is where it survives. +-- +-- terms_ref (specification section 5.2.4): reference to the governing terms +-- the emitter associates with the event. A processor MUST preserve it +-- unchanged - it is stored and served byte-for-byte, never normalised. + +ALTER TABLE events ADD COLUMN IF NOT EXISTS ctx_token TEXT; +ALTER TABLE events ADD COLUMN IF NOT EXISTS terms_ref TEXT; + +-- Destination reports join to the agent-recorded engagement by token. +CREATE INDEX IF NOT EXISTS idx_events_ctx_token ON events(ctx_token) WHERE ctx_token IS NOT NULL; diff --git a/crates/content-telemetry-core/profiles/rsl-1.json b/crates/content-telemetry-core/profiles/rsl-1.json index b962490..0d8c1c9 100644 --- a/crates/content-telemetry-core/profiles/rsl-1.json +++ b/crates/content-telemetry-core/profiles/rsl-1.json @@ -51,7 +51,7 @@ { "id": "grounding_scope", "kind": "field_coverage", - "spec_ref": "standard s5.7.2", + "spec_ref": "standard s6.4", "data_field": "scope", "applies_to": ["content_grounded"], "description": "content_grounded events carry data.scope." @@ -59,7 +59,7 @@ { "id": "citation_type", "kind": "field_coverage", - "spec_ref": "standard s5.7.3", + "spec_ref": "standard s6.5", "data_field": "citation_type", "applies_to": ["content_cited"], "description": "content_cited events carry data.citation_type." diff --git a/crates/content-telemetry-core/src/conformance.rs b/crates/content-telemetry-core/src/conformance.rs index ab8fbd7..5c94cd8 100644 --- a/crates/content-telemetry-core/src/conformance.rs +++ b/crates/content-telemetry-core/src/conformance.rs @@ -7,9 +7,40 @@ use serde_json::Value; -/// Schema versions this consumer accepts. During the 0.x preview period a -/// consumer accepts the exact same minor version only (spec 5.7.4). -pub const ACCEPTED_SCHEMA_VERSIONS: &[&str] = &["0.1"]; +/// Schema versions this consumer accepts. `"1.0"` is the version this +/// implementation targets. `"0.1"` remains accepted for the transition: +/// live member edge workers still declare it, and spec 12.1 gives a +/// consumer explicit migration rules for reading preview documents +/// (`bot_category` as `purpose`, defaulted `scope` and `citation_type`). +/// Under the spec's own rule the two lines do not interoperate — a strict +/// v1 consumer rejects `"0.1"` — so accepting both is a deliberate, +/// temporary deployment choice, not the 5.7.4 default. Documents declaring +/// `"1.0"` get full v1 strictness; documents declaring `"0.1"` are +/// normalised per 12.1 where a rule exists and tolerated otherwise. +pub const ACCEPTED_SCHEMA_VERSIONS: &[&str] = &["0.1", "1.0"]; + +/// The two schema lines this consumer reads. Which line a document declares +/// decides how much strictness applies at ingest: `V1_0` documents are held +/// to the v1 structural rules, `V0_1` documents are migrated per spec 12.1. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SchemaLine { + /// The v0.1 preview line, read under the migration rules of spec 12.1. + V0_1, + /// The v1.0 line, held to the full v1 structural rules. + V1_0, +} + +/// Which schema line a declared `schema_version` selects, or None when the +/// version is not accepted at all. Absent versions read as the current line: +/// the field postdates the earliest envelopes, and every live v0.1 emitter +/// declares its version explicitly, so absence means a current emitter. +pub fn schema_line(version: Option<&str>) -> Option { + match version { + None | Some("1.0") => Some(SchemaLine::V1_0), + Some("0.1") => Some(SchemaLine::V0_1), + Some(_) => None, + } +} /// Conformance levels the standard defines (spec 5.7). pub const STANDARD_CONFORMANCE_LEVELS: &[&str] = &["retrieval", "grounding", "citation"]; @@ -19,7 +50,6 @@ pub const STANDARD_CONFORMANCE_LEVELS: &[&str] = &["retrieval", "grounding", "ci pub const CONTENT_EVENT_TYPES: &[&str] = &[ "content_retrieved", "content_grounded", - "content_reproduced", "content_cited", "content_presented", "content_engaged", @@ -38,12 +68,14 @@ pub const CITATION_TYPES: &[&str] = &[ ]; pub const CITATION_POSITIONS: &[&str] = &["primary", "supporting", "mentioned", "unclassified"]; pub const GROUNDING_SCOPES: &[&str] = &["session", "turn"]; -pub const REPRODUCTION_TYPES: &[&str] = &["verbatim", "near_verbatim", "unclassified"]; pub const PRESENTATION_KINDS: &[&str] = &["content", "source_reference"]; /// The event type v1 withdrew (spec 12.1). Emitters MUST NOT send it on the /// v1 integration line; stored v0.1 rows keep it and are quarantined under -/// `extensions.events` at materialisation. +/// `extensions.events` at materialisation. `content_reproduced`, the other +/// type 12.1 excludes, existed only on the pre-release v1-draft line: it is +/// simply not in the core sets here, so stored draft rows self-quarantine +/// the same way without a named constant. pub const WITHDRAWN_EVENT_TYPE_DISPLAYED: &str = "content_displayed"; /// V0.1 event `data` fields prohibited by the v1 migration rule (spec 9.1). @@ -72,11 +104,6 @@ pub fn normalise_event_data_enums(event_type: &str, data: &mut Value) -> Vec &[( - "reproduction_type", - REPRODUCTION_TYPES, - Some("unclassified"), - )], _ => return Vec::new(), }; @@ -110,10 +137,66 @@ pub fn normalise_event_data_enums(event_type: &str, data: &mut Value) -> Vec) -> bool { - match version { - None => true, - Some(v) => ACCEPTED_SCHEMA_VERSIONS.contains(&v), + schema_line(version).is_some() +} + +/// Apply the spec 12.1 migration rules for reading a v0.1 preview event as +/// v1, in place, returning the names of the fields written: +/// +/// - `content_cited` without `data.citation_type` reads as `unclassified`; +/// - `content_grounded` without `data.scope` reads as `turn` when the event +/// carries a `turn_id` and `session` otherwise; +/// - a `data.bot_category` value is read as `purpose` (the v1 name for the +/// field). The original member is kept so pre-rename readers still see +/// it during the transition; the read layer coalesces the two. +/// +/// Everything else the preview line tolerated stays tolerated: 12.1 defines +/// no other defaulting rule, and inventing one would manufacture claims the +/// emitter never made. Used both at ingest for documents declaring `"0.1"` +/// and at materialisation for stored rows written before versioned ingest. +pub fn apply_v0_migration( + event_type: &str, + turn_id: Option<&str>, + data: &mut Value, +) -> Vec { + let needs_default = matches!(event_type, "content_grounded" | "content_cited"); + if data.is_null() && needs_default { + *data = Value::Object(serde_json::Map::new()); + } + let Some(obj) = data.as_object_mut() else { + return Vec::new(); + }; + + let mut changed = Vec::new(); + let absent = |obj: &serde_json::Map, field: &str| { + matches!(obj.get(field), None | Some(Value::Null)) + }; + + match event_type { + "content_cited" if absent(obj, "citation_type") => { + obj.insert( + "citation_type".to_string(), + Value::String("unclassified".to_string()), + ); + changed.push("citation_type".to_string()); + } + "content_grounded" if absent(obj, "scope") => { + let scope = if turn_id.is_some() { "turn" } else { "session" }; + obj.insert("scope".to_string(), Value::String(scope.to_string())); + changed.push("scope".to_string()); + } + _ => {} + } + + if let Some(Value::String(category)) = obj.get("bot_category") + && absent(obj, "purpose") + { + let category = category.clone(); + obj.insert("purpose".to_string(), Value::String(category)); + changed.push("purpose".to_string()); } + + changed } /// Result of normalising an emitter-supplied conformance level. @@ -159,11 +242,16 @@ pub fn content_identifier_present( content_url.is_some_and(|v| !v.is_empty()) || content_id.is_some_and(|v| !v.is_empty()) } -/// Strip the event `data` fields v1 withdrew (spec 9.1) in place, returning -/// the names of the fields removed. A hashed IP address is a pseudonym, not -/// an anonymous value, so `ip_hash` is treated as personal data: the -/// consumer drops it rather than storing it, on ingest and again on read -/// for rows written before this rule existed. +/// Strip the event `data` fields v1 withdrew (spec 9.1, 12.1) in place, +/// returning the names of the fields removed. A hashed IP address is a +/// pseudonym, not an anonymous value, so `ip_hash` is treated as personal +/// data: the consumer drops it rather than storing it, on ingest and again +/// on read for rows written before this rule existed. +/// +/// `content_fingerprint.preserved_in_output` is also withdrawn (spec 6.4, +/// 12.1): v1 defines no output-side reuse reporting, and a grounding +/// fingerprint is a grounding-time claim only. It lives one level down, so +/// the top-level sweep cannot catch it. pub fn strip_withdrawn_data_fields(data: &mut Value) -> Vec { let Some(obj) = data.as_object_mut() else { return Vec::new(); @@ -174,14 +262,22 @@ pub fn strip_withdrawn_data_fields(data: &mut Value) -> Vec { stripped.push((*field).to_string()); } } + if let Some(fingerprint) = obj + .get_mut("content_fingerprint") + .and_then(Value::as_object_mut) + && fingerprint.remove("preserved_in_output").is_some() + { + stripped.push("content_fingerprint.preserved_in_output".to_string()); + } stripped } /// The v1 structural requirements the standard's JSON Schema enforces per -/// event type (spec 5.2, 6.5-6.8): reproduction, citation and presentation -/// events carry an emitter-assigned `id` and an `output_id`; reproduction -/// and citation carry a resolvable source reference; reproduction carries -/// `data.reproduction_type`; presentation carries `data.presentation_kind` +/// event type (spec 5.2, 6.4-6.8): grounding carries `data.scope` (closed) +/// with `provenance` and `cached` kept consistent; citation and +/// presentation events carry an emitter-assigned `id` and an `output_id`; +/// citation carries a resolvable source reference and +/// `data.citation_type`; presentation carries `data.presentation_kind` /// (closed) and `data.presentation_type`; engagement carries the /// `presentation_id` of the exact presentation occurrence acted upon. /// @@ -202,47 +298,69 @@ pub fn v1_structural_violation( let data_str = |field: &str| -> Option<&str> { data.get(field).and_then(Value::as_str) }; match event_type { - "content_reproduced" | "content_cited" | "content_presented" => { + "content_grounded" => { + // scope is required and schema-closed (spec 6.4): the occurrence + // boundary and every counting model depend on it. + if !data_str("scope").is_some_and(|v| GROUNDING_SCOPES.contains(&v)) { + return Some(format!( + "content_grounded events must carry data.scope of: {}", + GROUNDING_SCOPES.join(", ") + )); + } + // provenance and cached must agree (spec 6.4): agent_fetched + // asserts a live fetch this session, agent_cached asserts reuse. + // third_party_sourced leaves cached unconstrained. + let cached = data.get("cached").and_then(Value::as_bool); + match data_str("provenance") { + Some("agent_fetched") if cached != Some(false) => Some( + "content_grounded provenance 'agent_fetched' requires data.cached: false" + .to_string(), + ), + Some("agent_cached") if cached != Some(true) => Some( + "content_grounded provenance 'agent_cached' requires data.cached: true" + .to_string(), + ), + _ => None, + } + } + "content_cited" | "content_presented" => { if !id_present { return Some(format!("{event_type} events must carry an event id")); } if !output_id.is_some_and(|v| !v.is_empty()) { return Some(format!("{event_type} events must carry output_id")); } - match event_type { - "content_reproduced" | "content_cited" => { - // Schema-enforced for these two types (spec 6.5, 6.6), - // over and above the application-layer identifier rule. - if !(content_url.is_some_and(|v| !v.is_empty()) - || content_id.is_some_and(|v| !v.is_empty())) - { - return Some(format!( - "{event_type} events must carry a resolvable content_url or content_id" - )); - } - if event_type == "content_reproduced" - && !data_str("reproduction_type").is_some_and(|v| !v.is_empty()) - { - return Some( - "content_reproduced events must carry data.reproduction_type" - .to_string(), - ); - } + if event_type == "content_cited" { + // Schema-enforced for citations (spec 6.5), over and above + // the application-layer identifier rule. + if !(content_url.is_some_and(|v| !v.is_empty()) + || content_id.is_some_and(|v| !v.is_empty())) + { + return Some(format!( + "{event_type} events must carry a resolvable content_url or content_id" + )); + } + // citation_type is required and schema-enforced (spec 6.5); + // an emitter that cannot classify uses 'unclassified'. + if !data_str("citation_type").is_some_and(|v| !v.is_empty()) { + return Some( + "content_cited events must carry data.citation_type \ + ('unclassified' when the agent cannot classify)" + .to_string(), + ); } - _ => { - let kind = data_str("presentation_kind"); - if !kind.is_some_and(|v| PRESENTATION_KINDS.contains(&v)) { - return Some(format!( - "content_presented events must carry data.presentation_kind of: {}", - PRESENTATION_KINDS.join(", ") - )); - } - if !data_str("presentation_type").is_some_and(|v| !v.is_empty()) { - return Some( - "content_presented events must carry data.presentation_type" - .to_string(), - ); - } + } else { + let kind = data_str("presentation_kind"); + if !kind.is_some_and(|v| PRESENTATION_KINDS.contains(&v)) { + return Some(format!( + "content_presented events must carry data.presentation_kind of: {}", + PRESENTATION_KINDS.join(", ") + )); + } + if !data_str("presentation_type").is_some_and(|v| !v.is_empty()) { + return Some( + "content_presented events must carry data.presentation_type".to_string(), + ); } } None @@ -262,6 +380,57 @@ pub fn v1_structural_violation( } } +/// The v1 rule that `source_role` MUST be present on every +/// `content_retrieved` event (spec 5.2.2, 5.7.5): without it a consumer +/// cannot tell an agent-reported fetch from an origin- or edge-reported +/// one. Kept separate from `v1_structural_violation` because `source_role` +/// is an event-level field, not a `data` member. +pub fn source_role_violation(event_type: &str, source_role: Option<&str>) -> Option { + if event_type == "content_retrieved" && !source_role.is_some_and(|v| !v.is_empty()) { + Some("content_retrieved events must carry source_role".to_string()) + } else { + None + } +} + +/// The v1 field-placement rules (spec 5.7.5): fields scoped to one event +/// type MUST NOT appear on others. `presentation_id` and the event-level +/// `ctx_token` belong only on `content_engaged`; `citation_id` only on +/// `content_presented`; `turn` only on `turn_started` and `turn_completed`. +/// Returns a description of the first misplaced field, or None. +pub fn field_placement_violation( + event_type: &str, + presentation_id_present: bool, + ctx_token_present: bool, + citation_id_present: bool, + turn_present: bool, +) -> Option { + if event_type != "content_engaged" { + if presentation_id_present { + return Some(format!( + "presentation_id may only appear on content_engaged events, not {event_type}" + )); + } + if ctx_token_present { + return Some(format!( + "an event-level ctx_token may only appear on content_engaged events, not \ + {event_type}" + )); + } + } + if citation_id_present && event_type != "content_presented" { + return Some(format!( + "citation_id may only appear on content_presented events, not {event_type}" + )); + } + if turn_present && !matches!(event_type, "turn_started" | "turn_completed") { + return Some(format!( + "turn may only appear on turn_started and turn_completed events, not {event_type}" + )); + } + None +} + /// Conversation-turn fields that MUST NOT be present at each privacy level /// (spec 5.5). `minimal` keeps only token counts and content URL arrays; /// `intent` strips the raw query/response text. PrivacyLevel is a closed @@ -358,12 +527,7 @@ mod tests { fn media_type_is_an_open_vocabulary() { // media_type tolerates emitter-defined values beyond the core set // (spec Annex A): unknown values pass through untouched. - for event_type in [ - "content_retrieved", - "content_grounded", - "content_reproduced", - "content_cited", - ] { + for event_type in ["content_retrieved", "content_grounded", "content_cited"] { let mut data = json!({ "media_type": "dataset" }); assert!( normalise_event_data_enums(event_type, &mut data).is_empty(), @@ -403,11 +567,77 @@ mod tests { } #[test] - fn schema_version_accepts_exact_minor_and_absent() { + fn schema_version_accepts_v1_and_transitional_v0() { + // "1.0" is the implemented version; absent reads as current. assert!(schema_version_accepted(None)); + assert!(schema_version_accepted(Some("1.0"))); + // "0.1" stays accepted during the transition: live member edge + // workers still declare it, and spec 12.1 defines how a consumer + // reads preview documents. This is deliberately more lenient than + // the spec's non-interoperation rule and is removed once every + // emitter declares "1.0". assert!(schema_version_accepted(Some("0.1"))); assert!(!schema_version_accepted(Some("0.2"))); - assert!(!schema_version_accepted(Some("1.0"))); + assert!(!schema_version_accepted(Some("1.1"))); + assert!(!schema_version_accepted(Some("2.0"))); + } + + #[test] + fn schema_line_maps_versions_to_strictness() { + assert_eq!(schema_line(None), Some(SchemaLine::V1_0)); + assert_eq!(schema_line(Some("1.0")), Some(SchemaLine::V1_0)); + assert_eq!(schema_line(Some("0.1")), Some(SchemaLine::V0_1)); + assert_eq!(schema_line(Some("9.9")), None); + } + + #[test] + fn v0_migration_defaults_citation_type_to_unclassified() { + let mut data = json!({ "excerpt_chars": 90 }); + let changed = apply_v0_migration("content_cited", None, &mut data); + assert_eq!(changed, vec!["citation_type"]); + assert_eq!(data["citation_type"], "unclassified"); + assert_eq!(data["excerpt_chars"], 90); + + // A supplied value is never overwritten. + let mut data = json!({ "citation_type": "direct_quote" }); + assert!(apply_v0_migration("content_cited", None, &mut data).is_empty()); + assert_eq!(data["citation_type"], "direct_quote"); + } + + #[test] + fn v0_migration_defaults_grounding_scope_by_turn_presence() { + let mut data = json!({}); + apply_v0_migration("content_grounded", Some("turn-1"), &mut data); + assert_eq!(data["scope"], "turn"); + + let mut data = json!({}); + apply_v0_migration("content_grounded", None, &mut data); + assert_eq!(data["scope"], "session"); + + // A supplied scope is never overwritten. + let mut data = json!({ "scope": "session" }); + assert!(apply_v0_migration("content_grounded", Some("turn-1"), &mut data).is_empty()); + assert_eq!(data["scope"], "session"); + + // Null data still gains the required member. + let mut data = Value::Null; + apply_v0_migration("content_grounded", None, &mut data); + assert_eq!(data["scope"], "session"); + } + + #[test] + fn v0_migration_reads_bot_category_as_purpose() { + let mut data = json!({ "bot_category": "inference", "bot_name": "Claude-User" }); + let changed = apply_v0_migration("content_retrieved", None, &mut data); + assert_eq!(changed, vec!["purpose"]); + assert_eq!(data["purpose"], "inference"); + // The original member is kept for pre-rename readers. + assert_eq!(data["bot_category"], "inference"); + + // An explicit purpose wins over the legacy field. + let mut data = json!({ "bot_category": "training", "purpose": "search" }); + assert!(apply_v0_migration("content_retrieved", None, &mut data).is_empty()); + assert_eq!(data["purpose"], "search"); } #[test] @@ -434,12 +664,18 @@ mod tests { } #[test] - fn reproduction_type_normalises_to_unclassified() { - let mut data = json!({ "reproduction_type": "loose_paraphrase", "reproduced_chars": 90 }); - let changed = normalise_event_data_enums("content_reproduced", &mut data); - assert_eq!(changed, vec!["reproduction_type"]); - assert_eq!(data["reproduction_type"], "unclassified"); - assert_eq!(data["reproduced_chars"], 90); + fn withdrawn_reproduced_type_carries_no_rules() { + // content_reproduced existed only on the pre-release v1-draft line + // (spec 12.1). It is an extension type now: no enum normalisation + // applies, and no v1 structural requirement recognises it. + let mut data = json!({ "reproduction_type": "loose_paraphrase" }); + assert!(normalise_event_data_enums("content_reproduced", &mut data).is_empty()); + assert_eq!(data["reproduction_type"], "loose_paraphrase"); + assert!(!CONTENT_EVENT_TYPES.contains(&"content_reproduced")); + assert!( + v1_structural_violation("content_reproduced", false, None, false, None, None, &data) + .is_none() + ); } #[test] @@ -458,13 +694,12 @@ mod tests { #[test] fn v1_structure_requires_output_identity_on_response_layer_events() { - for event_type in ["content_reproduced", "content_cited", "content_presented"] { + for event_type in ["content_cited", "content_presented"] { let data = match event_type { - "content_reproduced" => json!({ "reproduction_type": "verbatim" }), "content_presented" => { json!({ "presentation_kind": "source_reference", "presentation_type": "link" }) } - _ => json!({}), + _ => json!({ "citation_type": "direct_quote" }), }; assert!( v1_structural_violation( @@ -509,39 +744,25 @@ mod tests { } #[test] - fn v1_structure_requires_source_reference_on_reproduced_and_cited() { - for event_type in ["content_reproduced", "content_cited"] { - let data = json!({ "reproduction_type": "verbatim" }); - assert!( - v1_structural_violation( - event_type, - true, - Some("response:1"), - false, - None, - None, - &data - ) - .is_some(), - "{event_type} accepted an event with no source reference" - ); - } - } - - #[test] - fn v1_structure_requires_typed_reproduction_and_presentation_data() { + fn v1_structure_requires_source_reference_on_cited() { + let data = json!({ "citation_type": "direct_quote" }); assert!( v1_structural_violation( - "content_reproduced", + "content_cited", true, Some("response:1"), false, - Some("https://example.com/a"), None, - &json!({}) + None, + &data ) - .is_some() + .is_some(), + "content_cited accepted an event with no source reference" ); + } + + #[test] + fn v1_structure_requires_typed_presentation_data() { // presentation_kind is closed with no fallback member, so an // out-of-set value is a violation, not a normalisation case. assert!( @@ -612,6 +833,153 @@ mod tests { ); } + #[test] + fn v1_structure_requires_grounding_scope() { + let ok = |data: &Value| { + v1_structural_violation( + "content_grounded", + false, + None, + false, + Some("https://example.com/a"), + None, + data, + ) + }; + assert!(ok(&json!({ "scope": "session" })).is_none()); + assert!(ok(&json!({ "scope": "turn" })).is_none()); + assert!(ok(&json!({})).is_some(), "missing scope must be rejected"); + assert!( + ok(&json!({ "scope": "paragraph" })).is_some(), + "out-of-set scope must be rejected" + ); + } + + #[test] + fn v1_structure_requires_provenance_cached_consistency() { + let check = |data: Value| { + v1_structural_violation( + "content_grounded", + false, + None, + false, + Some("https://example.com/a"), + None, + &data, + ) + }; + // agent_fetched asserts a live fetch: cached must be false. + assert!( + check(json!({ "scope": "turn", "provenance": "agent_fetched", "cached": false })) + .is_none() + ); + assert!( + check(json!({ "scope": "turn", "provenance": "agent_fetched", "cached": true })) + .is_some() + ); + assert!(check(json!({ "scope": "turn", "provenance": "agent_fetched" })).is_some()); + // agent_cached asserts reuse: cached must be true. + assert!( + check(json!({ "scope": "turn", "provenance": "agent_cached", "cached": true })) + .is_none() + ); + assert!( + check(json!({ "scope": "turn", "provenance": "agent_cached", "cached": false })) + .is_some() + ); + assert!(check(json!({ "scope": "turn", "provenance": "agent_cached" })).is_some()); + // third_party_sourced leaves cached unconstrained (spec 6.4). + assert!( + check(json!({ "scope": "turn", "provenance": "third_party_sourced", "cached": true })) + .is_none() + ); + assert!(check(json!({ "scope": "turn", "provenance": "third_party_sourced" })).is_none()); + // Absent provenance constrains nothing. + assert!(check(json!({ "scope": "turn", "cached": true })).is_none()); + } + + #[test] + fn v1_structure_requires_citation_type() { + let check = |data: Value| { + v1_structural_violation( + "content_cited", + true, + Some("response:1"), + false, + Some("https://example.com/a"), + None, + &data, + ) + }; + assert!(check(json!({ "citation_type": "unclassified" })).is_none()); + assert!( + check(json!({})).is_some(), + "missing citation_type must be rejected" + ); + } + + #[test] + fn source_role_required_on_retrieval_only() { + assert!(source_role_violation("content_retrieved", None).is_some()); + assert!(source_role_violation("content_retrieved", Some("")).is_some()); + assert!(source_role_violation("content_retrieved", Some("edge")).is_none()); + assert!(source_role_violation("content_grounded", None).is_none()); + assert!(source_role_violation("turn_started", None).is_none()); + } + + #[test] + fn scoped_fields_may_not_appear_on_other_types() { + // Conforming placements pass. + assert!(field_placement_violation("content_engaged", true, true, false, false).is_none()); + assert!( + field_placement_violation("content_presented", false, false, true, false).is_none() + ); + assert!(field_placement_violation("turn_started", false, false, false, true).is_none()); + assert!(field_placement_violation("turn_completed", false, false, false, true).is_none()); + + // Misplacements are violations (spec 5.7.5). + assert!( + field_placement_violation("content_presented", true, false, false, false).is_some() + ); + assert!( + field_placement_violation("content_retrieved", false, true, false, false).is_some() + ); + assert!(field_placement_violation("content_cited", false, false, true, false).is_some()); + assert!(field_placement_violation("content_engaged", true, false, true, false).is_some()); + assert!(field_placement_violation("content_grounded", false, false, false, true).is_some()); + assert!( + field_placement_violation("checkout_completed", false, false, false, true).is_some() + ); + } + + #[test] + fn withdrawn_preserved_in_output_is_stripped_from_fingerprints() { + let mut data = json!({ + "scope": "turn", + "content_fingerprint": { + "scheme": "example:watermark", + "detected": true, + "preserved_in_output": true + } + }); + assert_eq!( + strip_withdrawn_data_fields(&mut data), + vec!["content_fingerprint.preserved_in_output"] + ); + assert_eq!(data["content_fingerprint"]["detected"], true); + assert!( + data["content_fingerprint"] + .get("preserved_in_output") + .is_none() + ); + + // A conforming fingerprint is left alone. + let mut clean = json!({ + "content_fingerprint": { "scheme": "example:watermark", "detected": false } + }); + assert!(strip_withdrawn_data_fields(&mut clean).is_empty()); + } + #[test] fn content_identifier_required_on_content_events_only() { assert!(!content_identifier_present("content_grounded", None, None)); diff --git a/crates/content-telemetry-core/src/models/event.rs b/crates/content-telemetry-core/src/models/event.rs index 39558b6..86a6703 100644 --- a/crates/content-telemetry-core/src/models/event.rs +++ b/crates/content-telemetry-core/src/models/event.rs @@ -20,18 +20,21 @@ pub struct TelemetryEventInput { pub content_telemetry_id: Option, pub content_url: Option, pub content_id: Option, - /// Output-artifact identity (spec 5.2): required on content_reproduced, - /// content_cited and content_presented events. + /// Output-artifact identity (spec 5.2): required on content_cited and + /// content_presented events. pub output_id: Option, pub output_element_id: Option, - /// On content_presented or content_reproduced, the id of the associated - /// content_cited event; absent for uncited presentations and uncredited - /// reproductions. + /// On content_presented, the id of the associated content_cited event; + /// absent for uncited presentations. pub citation_id: Option, /// On content_engaged, the id of the exact content_presented event the /// action occurred on. Required (spec 6.8). pub presentation_id: Option, pub license_ref: Option, + /// Reference to the governing terms the emitter associates with this + /// event (spec 5.2.4). Opaque: stored and served byte-for-byte, never + /// resolved, validated or rewritten. + pub terms_ref: Option, pub product_id: Option, pub turn: Option, #[serde(default)] @@ -106,7 +109,13 @@ pub struct EventRow { pub output_element_id: Option, pub citation_id: Option, pub presentation_id: Option, + /// On content_engaged, the click token minted for this engagement's + /// presentation (spec 5.2, 7.4.1), recorded so destination reports can + /// be joined to it. + pub ctx_token: Option, pub license_ref: Option, + /// Governing-terms reference (spec 5.2.4), preserved byte-for-byte. + pub terms_ref: Option, pub product_id: Option, pub turn_data: Option, pub event_data: serde_json::Value, diff --git a/crates/content-telemetry-core/src/models/query.rs b/crates/content-telemetry-core/src/models/query.rs index 85136c6..998db21 100644 --- a/crates/content-telemetry-core/src/models/query.rs +++ b/crates/content-telemetry-core/src/models/query.rs @@ -46,7 +46,14 @@ pub struct AgentBreakdown { /// Bot name from edge detection (event_data.bot_name), e.g. "ClaudeBot". /// Set when agent_id is null and the edge worker identified a known bot UA. pub bot_name: Option, - /// Edge-detected bot purpose: "training" | "inference" | "search". + /// Access purpose as classified by the reporting party (spec 6.2): + /// "training" | "inference" | "search" | "advertising", open enum. + /// Read from event_data.purpose, falling back to the v0.1 name + /// bot_category (spec 12.1). + pub purpose: Option, + /// Deprecated alias for `purpose` (the v0.1 field name). Carries the + /// same value during the transition so existing dashboard readers keep + /// working; new readers use `purpose`. pub bot_category: Option, pub event_count: i64, pub session_count: i64, @@ -104,13 +111,14 @@ pub struct DayFunnelCount { pub date: NaiveDate, pub retrieved: i64, pub grounded: i64, - /// Reproduction and citation are sibling response-layer claims - /// (spec 4.3), not consecutive funnel stages. - pub reproduced: i64, pub cited: i64, /// Includes stored v0.1 `content_displayed` events, so charts stay /// continuous across the v1 rename. pub presented: i64, + /// Deprecated alias for `presented` (the v0.1 stage name). Carries the + /// same value during the transition - the website dashboard still reads + /// `displayed`; a coordinated rename retires it. + pub displayed: i64, pub engaged: i64, } @@ -119,7 +127,7 @@ pub struct DayFunnelCount { /// Edge counts come from the publisher's own emitter under its own key — /// zero trust in the agent required. Everything past retrieval is /// agent-attested: correlation operates at the retrieval level only (spec -/// s7.3), so grounded/cited/displayed/engaged counts are what the agent +/// s7.3), so grounded/cited/presented/engaged counts are what the agent /// *says*, deterrence-audited, never corroborated. The field names say so. #[derive(Debug, Clone, Serialize)] pub struct AgentReconciliation { @@ -148,6 +156,11 @@ pub struct AgentReconciliation { pub struct AgentAttestedCounts { pub grounded: i64, pub cited: i64, + /// Presented count; stored v0.1 `content_displayed` rows are counted + /// inside this stage so history stays continuous across the v1 rename. + pub presented: i64, + /// Deprecated alias for `presented` (the v0.1 stage name). Same value + /// during the transition; the website dashboard still reads it. pub displayed: i64, pub engaged: i64, } @@ -248,8 +261,11 @@ pub struct PublisherQueryParams { /// Filter to a specific bot name from edge detection /// (event_data.bot_name), e.g. "ClaudeBot". pub bot: Option, - /// Filter to a bot category ("training" | "inference" | "search"). - pub bot_category: Option, + /// Filter to an access purpose ("training" | "inference" | "search" | + /// "advertising", spec 6.2). `bot_category` is accepted as a deprecated + /// alias for the v0.1 field name (spec 12.1). + #[serde(alias = "bot_category")] + pub purpose: Option, } #[derive(Debug, Clone, Default, Deserialize)] @@ -287,7 +303,9 @@ pub struct PaginatedQueryParams { pub until: Option>, pub domain: Option, pub bot: Option, - pub bot_category: Option, + /// Access-purpose filter; `bot_category` accepted as a deprecated alias. + #[serde(alias = "bot_category")] + pub purpose: Option, #[serde(default = "default_limit")] pub limit: i64, #[serde(default)] diff --git a/crates/content-telemetry-core/src/services/click_tokens.rs b/crates/content-telemetry-core/src/services/click_tokens.rs index 531a4b5..6a11dd7 100644 --- a/crates/content-telemetry-core/src/services/click_tokens.rs +++ b/crates/content-telemetry-core/src/services/click_tokens.rs @@ -10,16 +10,54 @@ use crate::models::click_token::{ }; use crate::models::event::EventRow; +/// Whether a ctx token value satisfies the spec 7.4.1 pattern +/// `^ct_[A-Za-z0-9_-]{16,240}$`. Applied at mint — including to +/// caller-supplied overrides — and at ingest, so no other shape enters the +/// system. The pattern is ASCII-only, so byte-wise checks are exact. +pub fn ctx_token_well_formed(token: &str) -> bool { + let Some(suffix) = token.strip_prefix("ct_") else { + return false; + }; + (16..=240).contains(&suffix.len()) + && suffix + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-') +} + +/// Mint a well-formed ctx token: `ct_` plus 32 bytes from the operating +/// system's CSPRNG, base64url-encoded without padding. That is 256 bits of +/// randomness against the spec's 96-bit floor (7.4.1): holding one token +/// gives no way to derive or enumerate another, and the value encodes no +/// content, session or user identifier. +fn mint_ctx_token() -> String { + use base64::Engine as _; + use rand::RngCore as _; + + let mut bytes = [0u8; 32]; + rand::rngs::OsRng.fill_bytes(&mut bytes); + format!( + "ct_{}", + base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(bytes) + ) +} + /// Create a click token mapping a click-out event to a session. /// -/// If `token` is None, a random UUID-based token is generated. +/// If `token` is None, a fresh `ct_`-prefixed token is minted from the +/// system CSPRNG. Callers passing their own token are responsible for +/// checking [`ctx_token_well_formed`] first (the HTTP layer does); the +/// debug assertion catches library misuse. pub async fn create_click_token( pool: &PgPool, session_id: Uuid, content_url: &str, token: Option<&str>, ) -> Result { - let token_value = token.map_or_else(|| Uuid::new_v4().to_string(), String::from); + debug_assert!( + token.is_none_or(ctx_token_well_formed), + "caller-supplied ctx tokens must match ^ct_[A-Za-z0-9_-]{{16,240}}$" + ); + let token_value = token.map_or_else(mint_ctx_token, String::from); sqlx::query_as::<_, ClickTokenRow>( "INSERT INTO click_tokens (token, session_id, content_url) @@ -81,6 +119,20 @@ pub async fn resolve_session_id(pool: &PgPool, token: &str) -> Result Result { .await?; Ok(result.rows_affected()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn minted_tokens_match_the_spec_pattern() { + for _ in 0..16 { + let token = mint_ctx_token(); + assert!(ctx_token_well_formed(&token), "minted {token} malformed"); + // 32 bytes base64url without padding is 43 characters. + assert_eq!(token.len(), "ct_".len() + 43); + } + // Distinct mints never collide in practice; two in a row certainly + // must not (a collision here means the RNG is broken). + assert_ne!(mint_ctx_token(), mint_ctx_token()); + } + + #[test] + fn well_formedness_follows_the_spec_grammar() { + assert!(ctx_token_well_formed("ct_0123456789abcdef")); + assert!(ctx_token_well_formed(&format!("ct_{}", "a".repeat(240)))); + + // Wrong or missing prefix. + assert!(!ctx_token_well_formed("cx_0123456789abcdef")); + assert!(!ctx_token_well_formed("0123456789abcdef")); + // Legacy UUID tokens are not well-formed. + assert!(!ctx_token_well_formed( + "d76318b8-4a06-4c48-8929-0c2f9b59d0c8" + )); + // Suffix length bounds: 16..=240. + assert!(!ctx_token_well_formed("ct_012345678901234")); + assert!(!ctx_token_well_formed(&format!("ct_{}", "a".repeat(241)))); + // Characters outside [A-Za-z0-9_-]. + assert!(!ctx_token_well_formed("ct_0123456789abcde!")); + assert!(!ctx_token_well_formed("ct_0123456789abcdé")); + assert!(!ctx_token_well_formed("")); + } +} diff --git a/crates/content-telemetry-core/src/services/events.rs b/crates/content-telemetry-core/src/services/events.rs index 234d050..9458474 100644 --- a/crates/content-telemetry-core/src/services/events.rs +++ b/crates/content-telemetry-core/src/services/events.rs @@ -78,7 +78,9 @@ pub async fn create_events_with_context( let mut output_element_ids: Vec> = Vec::with_capacity(len); let mut citation_ids: Vec> = Vec::with_capacity(len); let mut presentation_ids: Vec> = Vec::with_capacity(len); + let mut ctx_tokens: Vec> = Vec::with_capacity(len); let mut license_refs: Vec> = Vec::with_capacity(len); + let mut terms_refs: Vec> = Vec::with_capacity(len); let mut product_ids: Vec> = Vec::with_capacity(len); let mut turn_datas: Vec> = Vec::with_capacity(len); let mut event_datas = Vec::with_capacity(len); @@ -98,7 +100,12 @@ pub async fn create_events_with_context( output_element_ids.push(event.event.output_element_id.clone()); citation_ids.push(event.event.citation_id); presentation_ids.push(event.event.presentation_id); + ctx_tokens.push(event.event.ctx_token.clone()); license_refs.push(event.event.license_ref.clone()); + // terms_ref is preserved byte-for-byte (spec 5.2.4): a processor + // MUST NOT rewrite it, so it is cloned and bound with no + // normalisation of any kind. + terms_refs.push(event.event.terms_ref.clone()); product_ids.push(event.event.product_id); // A consumer that receives a privacy-violating turn strips the // offending fields rather than rejecting the document (spec 5.7.5). @@ -124,14 +131,14 @@ pub async fn create_events_with_context( r"INSERT INTO events ( id, session_id, organization_id, event_type, source_role, content_telemetry_id, content_url, content_id, turn_id, - output_id, output_element_id, citation_id, presentation_id, license_ref, - product_id, turn_data, event_data, event_timestamp + output_id, output_element_id, citation_id, presentation_id, ctx_token, license_ref, + terms_ref, product_id, turn_data, event_data, event_timestamp ) SELECT * FROM UNNEST( $1::uuid[], $2::uuid[], $3::uuid[], $4::text[], $5::text[], $6::uuid[], $7::text[], $8::text[], $9::text[], - $10::text[], $11::text[], $12::uuid[], $13::uuid[], $14::text[], - $15::uuid[], $16::jsonb[], $17::jsonb[], $18::timestamptz[] + $10::text[], $11::text[], $12::uuid[], $13::uuid[], $14::text[], $15::text[], + $16::text[], $17::uuid[], $18::jsonb[], $19::jsonb[], $20::timestamptz[] ) ON CONFLICT (id) DO NOTHING RETURNING *", @@ -149,7 +156,9 @@ pub async fn create_events_with_context( .bind(&output_element_ids) .bind(&citation_ids) .bind(&presentation_ids) + .bind(&ctx_tokens) .bind(&license_refs) + .bind(&terms_refs) .bind(&product_ids) .bind(&turn_datas) .bind(&event_datas) diff --git a/crates/content-telemetry-core/src/services/queries.rs b/crates/content-telemetry-core/src/services/queries.rs index 0477c87..aebfbd8 100644 --- a/crates/content-telemetry-core/src/services/queries.rs +++ b/crates/content-telemetry-core/src/services/queries.rs @@ -20,7 +20,7 @@ pub async fn get_publisher_summary( until: Option>, domain_filter: Option<&str>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, ) -> Result { let effective = effective_domains(domains, domain_filter); let patterns: Vec = domain_like_patterns(&effective); @@ -33,7 +33,7 @@ pub async fn get_publisher_summary( // // LEFT JOIN sessions so a bot filter can reach session-attached events. // Edge events identify their bot via event_data->>'bot_name'; agent - // self-report events (grounded/reproduced/cited/presented/engaged) + // self-report events (grounded/cited/presented/engaged) // identify it via the session's agent_id, mirroring the agent // breakdown's identity model. // Without this join, any bot filter silently drops every non-edge funnel @@ -66,12 +66,14 @@ pub async fn get_publisher_summary( )); param_idx += 1; } - if bot_category.is_some() { - // Category lives on event_data only; agent self-report events carry it - // when the emitter stamps it (the demo generator does). There is no - // category column on sessions to fall back to. + if purpose.is_some() { + // Purpose lives on event_data only; agent self-report events carry + // it when the emitter stamps it (the demo generator does). There is + // no purpose column on sessions to fall back to. COALESCE reads the + // v1 'purpose' name first and falls back to the stored v0.1 + // 'bot_category' rows (spec 12.1). query.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); } query.push_str(" GROUP BY e.event_type ORDER BY count DESC"); @@ -89,7 +91,7 @@ pub async fn get_publisher_summary( if let Some(b) = bot { q = q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { q = q.bind(c); } @@ -97,9 +99,9 @@ pub async fn get_publisher_summary( // they scan the same data independently, so parallelising cuts dashboard latency. let (rows, source_rows, status_rows, agents) = tokio::try_join!( q.fetch_all(pool), - query_source_breakdown(pool, &patterns, since, until, bot, bot_category), - query_status_breakdown(pool, &patterns, since, until, bot, bot_category), - query_agent_breakdown(pool, &patterns, since, until, bot, bot_category), + query_source_breakdown(pool, &patterns, since, until, bot, purpose), + query_status_breakdown(pool, &patterns, since, until, bot, purpose), + query_agent_breakdown(pool, &patterns, since, until, bot, purpose), )?; let total_events: i64 = rows.iter().map(|r| r.count).sum(); @@ -143,7 +145,7 @@ pub async fn get_publisher_summary( } /// Get per-day funnel counts -/// (retrieved/grounded/reproduced/cited/presented/engaged) for a publisher. +/// (retrieved/grounded/cited/presented/engaged) for a publisher. /// /// One row per UTC day in the window. The dashboard chart previously /// bucketed a sampled events list client-side, which collapsed to "all @@ -161,7 +163,7 @@ pub async fn get_publisher_timeseries( until: Option>, domain_filter: Option<&str>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, ) -> Result, sqlx::Error> { let effective = effective_domains(domains, domain_filter); let patterns: Vec = domain_like_patterns(&effective); @@ -174,8 +176,9 @@ pub async fn get_publisher_timeseries( "SELECT date_trunc('day', e.event_timestamp)::date AS day, COUNT(*) FILTER (WHERE e.event_type = 'content_retrieved') AS retrieved, COUNT(*) FILTER (WHERE e.event_type = 'content_grounded') AS grounded, - COUNT(*) FILTER (WHERE e.event_type = 'content_reproduced') AS reproduced, COUNT(*) FILTER (WHERE e.event_type = 'content_cited') AS cited, + -- Legacy branch: stored v0.1 content_displayed rows count + -- inside the presented stage (spec 12.1 renamed the type). COUNT(*) FILTER (WHERE e.event_type IN ('content_presented','content_displayed')) AS presented, COUNT(*) FILTER (WHERE e.event_type = 'content_engaged') AS engaged FROM events e @@ -203,18 +206,18 @@ pub async fn get_publisher_timeseries( )); param_idx += 1; } - if bot_category.is_some() { - // See get_publisher_summary: category lives on event_data only. + if purpose.is_some() { + // See get_publisher_summary: purpose lives on event_data only. query.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); } query.push_str( - " AND e.event_type IN ('content_retrieved','content_grounded','content_reproduced','content_cited','content_presented','content_displayed','content_engaged') \ + " AND e.event_type IN ('content_retrieved','content_grounded','content_cited','content_presented','content_displayed','content_engaged') \ GROUP BY day ORDER BY day ASC", ); - let mut q = sqlx::query_as::<_, (NaiveDate, i64, i64, i64, i64, i64, i64)>(&query); + let mut q = sqlx::query_as::<_, (NaiveDate, i64, i64, i64, i64, i64)>(&query); for pattern in &patterns { q = q.bind(pattern); } @@ -227,7 +230,7 @@ pub async fn get_publisher_timeseries( if let Some(b) = bot { q = q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { q = q.bind(c); } @@ -235,13 +238,15 @@ pub async fn get_publisher_timeseries( Ok(rows .into_iter() .map( - |(date, retrieved, grounded, reproduced, cited, presented, engaged)| DayFunnelCount { + |(date, retrieved, grounded, cited, presented, engaged)| DayFunnelCount { date, retrieved, grounded, - reproduced, cited, presented, + // Same value under the deprecated v0.1 stage name during + // the transition; the dashboard still reads `displayed`. + displayed: presented, engaged, }, ) @@ -257,7 +262,7 @@ pub async fn get_publisher_events( until: Option>, domain_filter: Option<&str>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, limit: i64, offset: i64, ) -> Result, sqlx::Error> { @@ -275,7 +280,7 @@ pub async fn get_publisher_events( } let (count_sql, data_sql) = - build_publisher_event_queries(&patterns, since, until, bot, bot_category); + build_publisher_event_queries(&patterns, since, until, bot, purpose); let mut count_q = sqlx::query_scalar::<_, i64>(&count_sql); for p in &patterns { @@ -290,13 +295,13 @@ pub async fn get_publisher_events( if let Some(b) = bot { count_q = count_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { count_q = count_q.bind(c); } let optional_count = usize::from(since.is_some()) + usize::from(until.is_some()) + usize::from(bot.is_some()) - + usize::from(bot_category.is_some()); + + usize::from(purpose.is_some()); let full_data_sql = format!( "{data_sql} LIMIT ${} OFFSET ${}", patterns.len() + 1 + optional_count, @@ -315,7 +320,7 @@ pub async fn get_publisher_events( if let Some(b) = bot { data_q = data_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { data_q = data_q.bind(c); } data_q = data_q.bind(limit).bind(offset); @@ -355,7 +360,7 @@ pub async fn get_publisher_url_metrics( until: Option>, domain_filter: Option<&str>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, limit: i64, offset: i64, ) -> Result, sqlx::Error> { @@ -391,9 +396,9 @@ pub async fn get_publisher_url_metrics( time_filter.push_str(&format!(" AND e.event_data->>'bot_name' = ${param_idx}")); param_idx += 1; } - if bot_category.is_some() { + if purpose.is_some() { time_filter.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); param_idx += 1; } @@ -419,7 +424,7 @@ pub async fn get_publisher_url_metrics( if let Some(b) = bot { count_q = count_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { count_q = count_q.bind(c); } let data_sql = format!( @@ -447,7 +452,7 @@ pub async fn get_publisher_url_metrics( if let Some(b) = bot { data_q = data_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { data_q = data_q.bind(c); } data_q = data_q.bind(limit).bind(offset); @@ -479,9 +484,9 @@ pub async fn get_publisher_url_metrics( tb_sql.push_str(&format!(" AND event_data->>'bot_name' = ${tb_param_idx}")); tb_param_idx += 1; } - if bot_category.is_some() { + if purpose.is_some() { tb_sql.push_str(&format!( - " AND event_data->>'bot_category' = ${tb_param_idx}" + " AND COALESCE(event_data->>'purpose', event_data->>'bot_category') = ${tb_param_idx}" )); } tb_sql.push_str(&format!( @@ -498,7 +503,7 @@ pub async fn get_publisher_url_metrics( if let Some(b) = bot { tb_q = tb_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { tb_q = tb_q.bind(c); } tb_q.fetch_all(pool).await? @@ -645,8 +650,9 @@ pub async fn get_agent_timeseries( "SELECT date_trunc('day', e.event_timestamp)::date AS day, COUNT(*) FILTER (WHERE e.event_type = 'content_retrieved') AS retrieved, COUNT(*) FILTER (WHERE e.event_type = 'content_grounded') AS grounded, - COUNT(*) FILTER (WHERE e.event_type = 'content_reproduced') AS reproduced, COUNT(*) FILTER (WHERE e.event_type = 'content_cited') AS cited, + -- Legacy branch: stored v0.1 content_displayed rows count + -- inside the presented stage (spec 12.1 renamed the type). COUNT(*) FILTER (WHERE e.event_type IN ('content_presented','content_displayed')) AS presented, COUNT(*) FILTER (WHERE e.event_type = 'content_engaged') AS engaged FROM events e @@ -662,11 +668,11 @@ pub async fn get_agent_timeseries( query.push_str(&format!(" AND e.event_timestamp <= ${param_idx}")); } query.push_str( - " AND e.event_type IN ('content_retrieved','content_grounded','content_reproduced','content_cited','content_presented','content_displayed','content_engaged') \ + " AND e.event_type IN ('content_retrieved','content_grounded','content_cited','content_presented','content_displayed','content_engaged') \ GROUP BY day ORDER BY day ASC", ); - let mut q = sqlx::query_as::<_, (NaiveDate, i64, i64, i64, i64, i64, i64)>(&query).bind(org_id); + let mut q = sqlx::query_as::<_, (NaiveDate, i64, i64, i64, i64, i64)>(&query).bind(org_id); if let Some(ref s) = since { q = q.bind(s); } @@ -678,13 +684,15 @@ pub async fn get_agent_timeseries( Ok(rows .into_iter() .map( - |(date, retrieved, grounded, reproduced, cited, presented, engaged)| DayFunnelCount { + |(date, retrieved, grounded, cited, presented, engaged)| DayFunnelCount { date, retrieved, grounded, - reproduced, cited, presented, + // Same value under the deprecated v0.1 stage name during + // the transition; the dashboard still reads `displayed`. + displayed: presented, engaged, }, ) @@ -1134,7 +1142,7 @@ fn build_publisher_event_queries( since: Option>, until: Option>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, ) -> (String, String) { let like_clauses: Vec = (1..=patterns.len()) .map(|i| format!("e.content_url LIKE ${i}")) @@ -1155,9 +1163,9 @@ fn build_publisher_event_queries( time_filter.push_str(&format!(" AND e.event_data->>'bot_name' = ${param_idx}")); param_idx += 1; } - if bot_category.is_some() { + if purpose.is_some() { time_filter.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); } @@ -1181,7 +1189,7 @@ async fn query_source_breakdown( since: Option>, until: Option>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, ) -> Result, sqlx::Error> { let like_clauses: Vec = (1..=patterns.len()) .map(|i| format!("e.content_url LIKE ${i}")) @@ -1202,9 +1210,9 @@ async fn query_source_breakdown( time_filter.push_str(&format!(" AND e.event_data->>'bot_name' = ${param_idx}")); param_idx += 1; } - if bot_category.is_some() { + if purpose.is_some() { time_filter.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); } @@ -1231,7 +1239,7 @@ async fn query_source_breakdown( if let Some(b) = bot { q = q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { q = q.bind(c); } @@ -1244,7 +1252,7 @@ async fn query_status_breakdown( since: Option>, until: Option>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, ) -> Result, sqlx::Error> { let like_clauses: Vec = (1..=patterns.len()) .map(|i| format!("e.content_url LIKE ${i}")) @@ -1265,9 +1273,9 @@ async fn query_status_breakdown( time_filter.push_str(&format!(" AND e.event_data->>'bot_name' = ${param_idx}")); param_idx += 1; } - if bot_category.is_some() { + if purpose.is_some() { time_filter.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); } @@ -1298,7 +1306,7 @@ async fn query_status_breakdown( if let Some(b) = bot { q = q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { q = q.bind(c); } @@ -1311,7 +1319,7 @@ async fn query_agent_breakdown( since: Option>, until: Option>, bot: Option<&str>, - bot_category: Option<&str>, + purpose: Option<&str>, ) -> Result, sqlx::Error> { let like_clauses: Vec = (1..=patterns.len()) .map(|i| format!("e.content_url LIKE ${i}")) @@ -1332,9 +1340,9 @@ async fn query_agent_breakdown( time_filter.push_str(&format!(" AND e.event_data->>'bot_name' = ${param_idx}")); param_idx += 1; } - if bot_category.is_some() { + if purpose.is_some() { time_filter.push_str(&format!( - " AND e.event_data->>'bot_category' = ${param_idx}" + " AND COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category') = ${param_idx}" )); } @@ -1342,19 +1350,20 @@ async fn query_agent_breakdown( // fall back to the bot identity the edge worker stamps onto event_data. // Without this, every edge retrieval collapses into a single "Unknown" row. let bot_name_expr = "NULLIF(e.event_data->>'bot_name', '')"; - let bot_category_expr = "NULLIF(e.event_data->>'bot_category', '')"; + let purpose_expr = + "NULLIF(COALESCE(e.event_data->>'purpose', e.event_data->>'bot_category'), '')"; // Get totals per agent let totals_sql = format!( "SELECT s.platform_id, s.agent_id, {bot_name_expr} as bot_name, - {bot_category_expr} as bot_category, + {purpose_expr} as purpose, COUNT(*) as event_count, COUNT(DISTINCT e.session_id) as session_count FROM events e LEFT JOIN sessions s ON e.session_id = s.id WHERE ({where_like}){time_filter} - GROUP BY s.platform_id, s.agent_id, bot_name, bot_category + GROUP BY s.platform_id, s.agent_id, bot_name, purpose ORDER BY event_count DESC" ); @@ -1362,14 +1371,14 @@ async fn query_agent_breakdown( let source_sql = format!( "SELECT s.platform_id, s.agent_id, {bot_name_expr} as bot_name, - {bot_category_expr} as bot_category, + {purpose_expr} as purpose, e.source_role, COUNT(*) as count, COUNT(DISTINCT e.session_id) as sessions FROM events e LEFT JOIN sessions s ON e.session_id = s.id WHERE ({where_like}){time_filter} - GROUP BY s.platform_id, s.agent_id, bot_name, bot_category, e.source_role + GROUP BY s.platform_id, s.agent_id, bot_name, purpose, e.source_role ORDER BY s.platform_id, s.agent_id, count DESC" ); @@ -1377,13 +1386,13 @@ async fn query_agent_breakdown( let type_sql = format!( "SELECT s.platform_id, s.agent_id, {bot_name_expr} as bot_name, - {bot_category_expr} as bot_category, + {purpose_expr} as purpose, e.event_type, COUNT(*) as count FROM events e LEFT JOIN sessions s ON e.session_id = s.id WHERE ({where_like}){time_filter} - GROUP BY s.platform_id, s.agent_id, bot_name, bot_category, e.event_type + GROUP BY s.platform_id, s.agent_id, bot_name, purpose, e.event_type ORDER BY s.platform_id, s.agent_id, count DESC" ); @@ -1400,7 +1409,7 @@ async fn query_agent_breakdown( if let Some(b) = bot { totals_q = totals_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { totals_q = totals_q.bind(c); } @@ -1417,7 +1426,7 @@ async fn query_agent_breakdown( if let Some(b) = bot { source_q = source_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { source_q = source_q.bind(c); } @@ -1434,7 +1443,7 @@ async fn query_agent_breakdown( if let Some(b) = bot { type_q = type_q.bind(b); } - if let Some(c) = bot_category { + if let Some(c) = purpose { type_q = type_q.bind(c); } @@ -1451,12 +1460,12 @@ async fn query_agent_breakdown( Option, ); - // Group source breakdowns by (platform_id, agent_id, bot_name, bot_category) + // Group source breakdowns by (platform_id, agent_id, bot_name, purpose) let mut source_map: std::collections::HashMap> = std::collections::HashMap::new(); for s in sources { source_map - .entry((s.platform_id, s.agent_id, s.bot_name, s.bot_category)) + .entry((s.platform_id, s.agent_id, s.bot_name, s.purpose)) .or_default() .push(SourceRoleCount { source_role: s.source_role, @@ -1465,12 +1474,12 @@ async fn query_agent_breakdown( }); } - // Group event_type breakdowns by (platform_id, agent_id, bot_name, bot_category) + // Group event_type breakdowns by (platform_id, agent_id, bot_name, purpose) let mut type_map: std::collections::HashMap> = std::collections::HashMap::new(); for t in types { type_map - .entry((t.platform_id, t.agent_id, t.bot_name, t.bot_category)) + .entry((t.platform_id, t.agent_id, t.bot_name, t.purpose)) .or_default() .push(EventTypeCount { event_type: t.event_type, @@ -1485,7 +1494,7 @@ async fn query_agent_breakdown( r.platform_id.clone(), r.agent_id.clone(), r.bot_name.clone(), - r.bot_category.clone(), + r.purpose.clone(), ); let by_source = source_map.remove(&key).unwrap_or_default(); let by_event_type = type_map.remove(&key).unwrap_or_default(); @@ -1493,7 +1502,10 @@ async fn query_agent_breakdown( platform_id: r.platform_id, agent_id: r.agent_id, bot_name: r.bot_name, - bot_category: r.bot_category, + // The same value under both names during the bot_category → + // purpose transition (spec 12.1). + bot_category: r.purpose.clone(), + purpose: r.purpose, event_count: r.event_count, session_count: r.session_count, by_source, @@ -1564,7 +1576,7 @@ pub async fn get_publisher_reconciliation( COUNT(*) FILTER (WHERE e.event_type = 'content_retrieved') AS self_retrievals, COUNT(*) FILTER (WHERE e.event_type = 'content_grounded') AS grounded, COUNT(*) FILTER (WHERE e.event_type = 'content_cited') AS cited, - COUNT(*) FILTER (WHERE e.event_type = 'content_displayed') AS displayed, + COUNT(*) FILTER (WHERE e.event_type IN ('content_presented','content_displayed')) AS presented, COUNT(*) FILTER (WHERE e.event_type = 'content_engaged') AS engaged, COUNT(DISTINCT e.session_id) AS sessions_reporting FROM events e @@ -1578,7 +1590,7 @@ pub async fn get_publisher_reconciliation( COALESCE(self_report.self_retrievals, 0) AS self_reported_retrievals, COALESCE(self_report.grounded, 0) AS grounded, COALESCE(self_report.cited, 0) AS cited, - COALESCE(self_report.displayed, 0) AS displayed, + COALESCE(self_report.presented, 0) AS presented, COALESCE(self_report.engaged, 0) AS engaged, COALESCE(self_report.sessions_reporting, 0) AS sessions_reporting FROM edge @@ -1613,7 +1625,9 @@ pub async fn get_publisher_reconciliation( agent_attested: AgentAttestedCounts { grounded: r.grounded, cited: r.cited, - displayed: r.displayed, + presented: r.presented, + // Deprecated alias, same value during the transition. + displayed: r.presented, engaged: r.engaged, }, }) @@ -1806,7 +1820,7 @@ struct ReconciliationRow { self_reported_retrievals: i64, grounded: i64, cited: i64, - displayed: i64, + presented: i64, engaged: i64, sessions_reporting: i64, } @@ -1855,7 +1869,7 @@ struct AgentBreakdownRow { platform_id: Option, agent_id: Option, bot_name: Option, - bot_category: Option, + purpose: Option, event_count: i64, session_count: i64, } @@ -1878,7 +1892,7 @@ struct AgentSourceRow { platform_id: Option, agent_id: Option, bot_name: Option, - bot_category: Option, + purpose: Option, source_role: Option, count: i64, sessions: i64, @@ -1889,7 +1903,7 @@ struct AgentTypeRow { platform_id: Option, agent_id: Option, bot_name: Option, - bot_category: Option, + purpose: Option, event_type: String, count: i64, } diff --git a/crates/content-telemetry-core/src/standard.rs b/crates/content-telemetry-core/src/standard.rs index 3ebec75..b75271f 100644 --- a/crates/content-telemetry-core/src/standard.rs +++ b/crates/content-telemetry-core/src/standard.rs @@ -19,11 +19,13 @@ use crate::models::session::SessionWithEvents; /// extension event and moves to the document's `extensions` member - /// including `content_displayed`, which v1 withdrew (spec 12.1): stored /// v0.1 rows keep their type and are quarantined rather than rewritten -/// into presentation claims the emitter never made. +/// into presentation claims the emitter never made. The same applies to +/// `content_reproduced`, which existed only on the pre-release v1-draft +/// line and is not part of v1 (spec 12.1): stored draft rows self- +/// quarantine here rather than being rewritten. const CORE_EVENT_TYPES: &[&str] = &[ "content_retrieved", "content_grounded", - "content_reproduced", "content_cited", "content_presented", "content_engaged", @@ -31,13 +33,29 @@ const CORE_EVENT_TYPES: &[&str] = &[ "turn_completed", ]; +/// Prepare a stored row's `data` for the v1 document: normalise closed +/// enums (spec Annex A), strip the fields v1 withdrew (spec 9.1), and apply +/// the spec 12.1 migration defaults for rows written before versioned +/// ingest. Rows ingested on the v1 line come out unchanged; pre-v1 rows +/// come out reading as the v1 document the migration rules define. +fn prepared_event_data(row: &EventRow) -> Value { + let mut data = row.event_data.clone(); + if !data.is_null() { + crate::conformance::normalise_event_data_enums(&row.event_type, &mut data); + crate::conformance::strip_withdrawn_data_fields(&mut data); + } + crate::conformance::apply_v0_migration(&row.event_type, row.turn_id.as_deref(), &mut data); + data +} + /// Whether a stored row satisfies the v1 structural requirements for its -/// event type. Rows ingested on the v1 line always do (ingest rejects -/// violations); pre-v1 rows that do not - an engagement without a -/// `presentation_id`, a citation without an `output_id` - are quarantined -/// under `extensions.events` so the materialised document stays valid -/// against the v1 schema. -fn meets_v1_structure(row: &EventRow) -> bool { +/// event type, judged against its prepared (migration-normalised) `data`. +/// Rows ingested on the v1 line always do (ingest rejects violations); +/// pre-v1 rows that still do not after normalisation - an engagement +/// without a `presentation_id`, a citation without an `output_id` - are +/// quarantined under `extensions.events` so the materialised document stays +/// valid against the v1 schema. +fn meets_v1_structure(row: &EventRow, prepared_data: &Value) -> bool { crate::conformance::v1_structural_violation( &row.event_type, true, @@ -45,9 +63,19 @@ fn meets_v1_structure(row: &EventRow) -> bool { row.presentation_id.is_some(), row.content_url.as_deref(), row.content_id.as_deref(), - &row.event_data, + prepared_data, ) .is_none() + && crate::conformance::source_role_violation(&row.event_type, row.source_role.as_deref()) + .is_none() + && crate::conformance::field_placement_violation( + &row.event_type, + row.presentation_id.is_some(), + row.ctx_token.is_some(), + row.citation_id.is_some(), + row.turn_data.is_some(), + ) + .is_none() } fn insert_if_some(obj: &mut Map, key: &str, value: Option) { @@ -58,7 +86,7 @@ fn insert_if_some(obj: &mut Map, key: &str, value: Option) } } -fn standard_event(row: &EventRow) -> Value { +fn standard_event(row: &EventRow, prepared_data: &Value) -> Value { let mut event = Map::new(); event.insert("id".to_string(), json!(row.id)); event.insert("type".to_string(), json!(row.event_type)); @@ -80,6 +108,16 @@ fn standard_event(row: &EventRow) -> Value { "presentation_id", row.presentation_id.map(|v| json!(v)), ); + // The event-level ctx_token belongs on content_engaged only (spec 5.2, + // 5.7.5): the agent records the token minted for the engaged + // presentation so destination reports join to it. + if row.event_type == "content_engaged" { + insert_if_some( + &mut event, + "ctx_token", + row.ctx_token.clone().map(Value::from), + ); + } insert_if_some( &mut event, "source_role", @@ -105,17 +143,21 @@ fn standard_event(row: &EventRow) -> Value { "license_ref", row.license_ref.clone().map(Value::from), ); + // terms_ref passes through byte-for-byte (spec 5.2.4): a processor MUST + // preserve it unchanged and MUST NOT remove or rewrite it. + insert_if_some( + &mut event, + "terms_ref", + row.terms_ref.clone().map(Value::from), + ); insert_if_some(&mut event, "turn", row.turn_data.clone()); - if !row.event_data.is_null() { + if !prepared_data.is_null() { // Rows stored before ingest normalised closed enums can still carry - // out-of-set values; normalise again on read so the materialised - // document always validates (spec Annex A). The same applies to the - // fields v1 withdrew (spec 9.1): rows written before the ip_hash - // prohibition are stripped on read. - let mut data = row.event_data.clone(); - crate::conformance::normalise_event_data_enums(&row.event_type, &mut data); - crate::conformance::strip_withdrawn_data_fields(&mut data); - event.insert("data".to_string(), data); + // out-of-set values, and rows stored before versioned ingest can + // lack the members v1 requires; `prepared_event_data` normalised, + // migrated (spec 12.1) and stripped (spec 9.1) them on read so the + // materialised document always validates. + event.insert("data".to_string(), prepared_data.clone()); } // product_id is an extension field; carry it inside data where the // schema permits custom members, not as a top-level event field. @@ -142,7 +184,12 @@ pub fn standard_document(swe: &SessionWithEvents) -> Value { let mut doc = Map::new(); doc.insert("document_type".to_string(), json!("session")); - doc.insert("schema_version".to_string(), json!("0.1")); + // The document this module builds is a v1 session document: stored rows + // are read under the spec 12.1 migration rules (scope and citation_type + // defaults, bot_category as purpose) and rows that still do not satisfy + // the v1 shape are quarantined under `extensions.events`, so the stamp + // is truthful after normalisation. + doc.insert("schema_version".to_string(), json!("1.0")); doc.insert("session_id".to_string(), json!(session.id)); insert_if_some( &mut doc, @@ -184,20 +231,34 @@ pub fn standard_document(swe: &SessionWithEvents) -> Value { None => {} } - let (core_events, extension_events): (Vec<&EventRow>, Vec<&EventRow>) = swe + type Prepared<'a> = Vec<(&'a EventRow, Value)>; + let (core_events, extension_events): (Prepared<'_>, Prepared<'_>) = swe .events .iter() - .partition(|e| CORE_EVENT_TYPES.contains(&e.event_type.as_str()) && meets_v1_structure(e)); + .map(|e| (e, prepared_event_data(e))) + .partition(|(e, data)| { + CORE_EVENT_TYPES.contains(&e.event_type.as_str()) && meets_v1_structure(e, data) + }); doc.insert( "events".to_string(), - Value::Array(core_events.iter().map(|e| standard_event(e)).collect()), + Value::Array( + core_events + .iter() + .map(|(e, data)| standard_event(e, data)) + .collect(), + ), ); if !extension_events.is_empty() { extensions.insert( "events".to_string(), - Value::Array(extension_events.iter().map(|e| standard_event(e)).collect()), + Value::Array( + extension_events + .iter() + .map(|(e, data)| standard_event(e, data)) + .collect(), + ), ); } @@ -240,3 +301,187 @@ pub fn standard_document(swe: &SessionWithEvents) -> Value { Value::Object(doc) } + +#[cfg(test)] +mod tests { + use chrono::Utc; + use uuid::Uuid; + + use super::*; + use crate::models::session::SessionRow; + + fn session_row() -> SessionRow { + SessionRow { + id: Uuid::new_v4(), + organization_id: Uuid::new_v4(), + parent_session_id: None, + initiator_type: "user".to_string(), + initiator: None, + content_scope: None, + manifest_ref: None, + conformance_level: Some("citation".to_string()), + config_snapshot_hash: None, + agent_id: Some("test-agent".to_string()), + external_session_id: None, + prior_session_ids: None, + user_context: json!({}), + platform_id: None, + client_type: None, + client_info: None, + started_at: Utc::now(), + ended_at: None, + outcome_type: None, + outcome_value: None, + created_at: Utc::now(), + updated_at: Utc::now(), + } + } + + fn event_row(event_type: &str, data: Value) -> EventRow { + EventRow { + id: Uuid::new_v4(), + session_id: None, + organization_id: Uuid::new_v4(), + event_type: event_type.to_string(), + source_role: Some("agent".to_string()), + content_telemetry_id: None, + content_url: Some("https://example.com/a".to_string()), + content_id: None, + turn_id: None, + output_id: None, + output_element_id: None, + citation_id: None, + presentation_id: None, + ctx_token: None, + license_ref: None, + terms_ref: None, + product_id: None, + turn_data: None, + event_data: data, + event_timestamp: Utc::now(), + created_at: Utc::now(), + } + } + + fn event_types(doc: &Value, pointer: &str) -> Vec { + doc.pointer(pointer) + .and_then(Value::as_array) + .map(|events| { + events + .iter() + .filter_map(|e| e.get("type").and_then(Value::as_str)) + .map(ToString::to_string) + .collect() + }) + .unwrap_or_default() + } + + #[test] + fn withdrawn_types_self_quarantine_into_extensions() { + // A stored pre-release content_reproduced row and a stored v0.1 + // content_displayed row both fall outside CORE_EVENT_TYPES, so the + // partition moves them under extensions.events untouched (spec 12.1) + // while conforming core rows stay in the document body. + let swe = SessionWithEvents { + session: session_row(), + events: vec![ + event_row("content_grounded", json!({ "scope": "session" })), + event_row( + "content_reproduced", + json!({ "reproduction_type": "verbatim", "output_id": "out-1" }), + ), + event_row("content_displayed", json!({ "display_type": "link" })), + event_row("checkout_completed", json!({})), + ], + }; + + let doc = standard_document(&swe); + assert_eq!(doc["schema_version"], "1.0"); + assert_eq!(event_types(&doc, "/events"), vec!["content_grounded"]); + + let quarantined = event_types(&doc, "/extensions/events"); + assert_eq!( + quarantined, + vec![ + "content_reproduced", + "content_displayed", + "checkout_completed" + ] + ); + + // Quarantined rows keep the claims their emitters made. + let reproduced = &doc["extensions"]["events"][0]; + assert_eq!(reproduced["data"]["reproduction_type"], "verbatim"); + } + + #[test] + fn ctx_token_and_terms_ref_materialise_where_the_spec_places_them() { + let mut engaged = event_row( + "content_engaged", + json!({ "engagement_type": "link_click" }), + ); + engaged.presentation_id = Some(Uuid::new_v4()); + engaged.ctx_token = Some("ct_dGVzdHRva2VudmFsdWU".to_string()); + engaged.terms_ref = Some("https://example.com/terms/2026-01".to_string()); + + let mut grounded = event_row("content_grounded", json!({ "scope": "session" })); + grounded.terms_ref = Some("opaque:terms-77".to_string()); + + // A pre-v1-tolerated row with a ctx_token on the wrong type is + // quarantined by the field-placement rule, never emitted as core. + let mut misplaced = event_row("content_grounded", json!({ "scope": "session" })); + misplaced.ctx_token = Some("ct_bWlzcGxhY2VkdG9rZW4".to_string()); + + let swe = SessionWithEvents { + session: session_row(), + events: vec![engaged, grounded, misplaced], + }; + let doc = standard_document(&swe); + + assert_eq!( + event_types(&doc, "/events"), + vec!["content_engaged", "content_grounded"] + ); + assert_eq!(doc["events"][0]["ctx_token"], "ct_dGVzdHRva2VudmFsdWU"); + // terms_ref passes through byte-for-byte on any event (spec 5.2.4). + assert_eq!( + doc["events"][0]["terms_ref"], + "https://example.com/terms/2026-01" + ); + assert_eq!(doc["events"][1]["terms_ref"], "opaque:terms-77"); + assert_eq!( + event_types(&doc, "/extensions/events"), + vec!["content_grounded"] + ); + } + + #[test] + fn stored_v0_rows_materialise_under_the_migration_rules() { + // Rows written before versioned ingest lack the members v1 requires; + // the spec 12.1 defaults are applied on read so the document body + // keeps them rather than quarantining history (and the "1.0" stamp + // stays truthful). + let mut cited = event_row("content_cited", json!({})); + cited.output_id = Some("response:1".to_string()); + let mut grounded_turn = event_row("content_grounded", json!({})); + grounded_turn.turn_id = Some("turn-1".to_string()); + + let swe = SessionWithEvents { + session: session_row(), + events: vec![ + event_row("content_grounded", json!({})), + grounded_turn, + cited, + ], + }; + + let doc = standard_document(&swe); + assert_eq!( + event_types(&doc, "/events"), + vec!["content_grounded", "content_grounded", "content_cited"] + ); + assert_eq!(doc["events"][0]["data"]["scope"], "session"); + assert_eq!(doc["events"][1]["data"]["scope"], "turn"); + assert_eq!(doc["events"][2]["data"]["citation_type"], "unclassified"); + } +} diff --git a/crates/content-telemetry-core/tests/integration.rs b/crates/content-telemetry-core/tests/integration.rs index 012edae..d4106d5 100644 --- a/crates/content-telemetry-core/tests/integration.rs +++ b/crates/content-telemetry-core/tests/integration.rs @@ -177,6 +177,7 @@ async fn events_created_for_active_session(pool: PgPool) { citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data: serde_json::json!({}), @@ -197,6 +198,7 @@ async fn events_created_for_active_session(pool: PgPool) { citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data: serde_json::json!({}), @@ -270,6 +272,7 @@ async fn session_with_events_returns_both(pool: PgPool) { citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data: serde_json::json!({}), @@ -320,6 +323,7 @@ async fn click_token_created_and_looked_up(pool: PgPool) { citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data: serde_json::json!({}), @@ -340,6 +344,7 @@ async fn click_token_created_and_looked_up(pool: PgPool) { citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data: serde_json::json!({}), @@ -360,6 +365,7 @@ async fn click_token_created_and_looked_up(pool: PgPool) { citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data: serde_json::json!({ @@ -517,6 +523,7 @@ fn agent_event(event_type: &str, url: &str, data: serde_json::Value) -> Telemetr citation_id: None, presentation_id: None, license_ref: None, + terms_ref: None, product_id: None, turn: None, data, diff --git a/crates/content-telemetry-server/src/routes.rs b/crates/content-telemetry-server/src/routes.rs index 57ecb35..317dc5a 100644 --- a/crates/content-telemetry-server/src/routes.rs +++ b/crates/content-telemetry-server/src/routes.rs @@ -139,9 +139,9 @@ async fn bulk_session( )); } - if !conformance::schema_version_accepted(req.schema_version.as_deref()) { + let Some(line) = conformance::schema_line(req.schema_version.as_deref()) else { return Err(unsupported_schema_version(req.schema_version.as_deref())); - } + }; check_initiator_type(&req.initiator_type)?; if let Some(outcome) = req.outcome.as_ref() { @@ -154,8 +154,10 @@ async fn bulk_session( let mut events_in = req.events; for (index, event) in events_in.iter_mut().enumerate() { - validate::check_event(event, state.max_event_age_days, index)?; - log_notes(validate::normalise(event)); + // Normalise before checking: a "0.1" document's migration defaults + // (spec 12.1) must land before the structural rules judge it. + log_notes(validate::normalise(event, line)); + validate::check_event(event, line, state.max_event_age_days, index)?; } let level = conformance::normalise_conformance_level(req.conformance_level.as_deref()); @@ -274,9 +276,9 @@ async fn record_events( )); } - if !conformance::schema_version_accepted(req.schema_version.as_deref()) { + let Some(line) = conformance::schema_line(req.schema_version.as_deref()) else { return Err(unsupported_schema_version(req.schema_version.as_deref())); - } + }; // Both envelopes are accepted: `event` for a single standalone event, // `events` for a batch (specification section 7.1). @@ -295,6 +297,21 @@ async fn record_events( return Err(too_large(incoming.len())); } + // An envelope ctx_token accompanies content_engaged events only (spec + // 5.7.5). Checked here, not during binding: a batch that also presents + // a session_id binds through the session, but presenting the token + // alongside non-engagement claims is still malformed. + if req.ctx_token.is_some() + && let Some(index) = incoming + .iter() + .position(|e| e.event_type != "content_engaged") + { + return Err(ApiError::bad_request(format!( + "event {index}: an envelope ctx_token may only accompany content_engaged events; \ + supply a session_id instead" + ))); + } + let batch_session = req .session_id .as_deref() @@ -302,8 +319,10 @@ async fn record_events( let batch_session = batch_session.transpose()?; for (index, event) in incoming.iter_mut().enumerate() { - validate::check_event(event, state.max_event_age_days, index)?; - log_notes(validate::normalise(event)); + // Normalise before checking: a "0.1" document's migration defaults + // (spec 12.1) must land before the structural rules judge it. + log_notes(validate::normalise(event, line)); + validate::check_event(event, line, state.max_event_age_days, index)?; } let defaults = BatchDefaults { @@ -395,6 +414,16 @@ async fn resolve_binding( ))); } + // Well-formedness before lookup (spec 7.4.1): every token this + // server mints matches the pattern, so anything else can only be + // noise or probing and never reaches the database. + if !click_tokens::ctx_token_well_formed(token) { + return Err(ApiError::bad_request(format!( + "event {index}: malformed ctx_token; token values match \ + ^ct_[A-Za-z0-9_-]{{16,240}}$ (spec 7.4.1)" + ))); + } + let session_id = click_tokens::resolve_session_id(&state.pool, token) .await? .ok_or_else(|| { @@ -472,6 +501,16 @@ async fn create_click_token( OrgContext(org): OrgContext, Json(req): Json, ) -> Result { + // A caller-supplied token is held to the same spec 7.4.1 shape as a + // minted one; the unguessability of its suffix is the issuer's burden. + if let Some(token) = req.token.as_deref() + && !click_tokens::ctx_token_well_formed(token) + { + return Err(ApiError::bad_request( + "token must match ^ct_[A-Za-z0-9_-]{16,240}$ (spec 7.4.1)", + )); + } + let session = sessions::find_owned_session(&state.pool, org, req.session_id) .await? .ok_or(ApiError::NotFound)?; @@ -495,15 +534,21 @@ async fn create_click_token( )) } -/// Resolve a ctx token to its click manifest. +/// Resolve a ctx token to its click context. /// /// Deliberately unauthenticated: the destination of a click-out has no /// account here, and requiring one would defeat the mechanism. The token is /// the credential, and what it discloses is bounded twice over — by -/// two-sided consent inside the core lookup, and by the manifest shape, which +/// two-sided consent inside the core lookup, and by the response shape, which /// never contains the session id. A missing token, an expired one and a /// non-consenting one are all 404, so the endpoint cannot be used to probe /// which sessions exist. +/// +/// TODO(spec 7.4.4): the response is still the v0.1 click-manifest shape; +/// the v1 four-component click context (engagement, clicked-content +/// lineage, turn-scoped contributing sources, count-based session summary) +/// is a separate design change — see the note on +/// `click_tokens::lookup_by_token`. async fn lookup_ctx( State(state): State, Path(token): Path, diff --git a/crates/content-telemetry-server/src/validate.rs b/crates/content-telemetry-server/src/validate.rs index ecc7148..45f184f 100644 --- a/crates/content-telemetry-server/src/validate.rs +++ b/crates/content-telemetry-server/src/validate.rs @@ -11,6 +11,7 @@ use chrono::{DateTime, Duration, Utc}; use content_telemetry_core::conformance; use content_telemetry_core::models::event::TelemetryEventInput; +use content_telemetry_core::services::click_tokens; use crate::error::ApiError; @@ -67,9 +68,14 @@ pub fn check_timestamp( } /// Structural checks that apply to every event regardless of how it binds to -/// a session. +/// a session. `line` is the schema line the enclosing document declared: +/// the full v1 structural rules bind documents on the `"1.0"` line (and +/// undeclared documents, read as current), while `"0.1"` documents are +/// normalised per spec 12.1 by [`normalise`] and tolerated where no +/// migration rule exists — the transition posture for live v0.1 emitters. pub fn check_event( event: &TelemetryEventInput, + line: conformance::SchemaLine, max_event_age_days: i64, index: usize, ) -> Result<(), ApiError> { @@ -81,12 +87,12 @@ pub fn check_event( // v1 withdrew content_displayed (specification section 12.1). Stored v0.1 // rows keep their type, but nothing new is accepted under it: the - // replacement types distinguish presenting content from presenting a - // reference to it, and rewriting one into the other would manufacture a - // claim the emitter never made. + // replacement type distinguishes presenting content from presenting a + // reference to it via data.presentation_kind, and rewriting one into the + // other would manufacture a claim the emitter never made. if event.event_type == conformance::WITHDRAWN_EVENT_TYPE_DISPLAYED { return Err(ApiError::bad_request(format!( - "event {index}: '{}' was withdrawn in v1 — use content_presented or content_reproduced", + "event {index}: '{}' was withdrawn in v1 — use content_presented", conformance::WITHDRAWN_EVENT_TYPE_DISPLAYED ))); } @@ -113,16 +119,50 @@ pub fn check_event( ))); } - if let Some(violation) = conformance::v1_structural_violation( - &event.event_type, - event.id.is_some(), - event.output_id.as_deref(), - event.presentation_id.is_some(), - event.content_url.as_deref(), - event.content_id.as_deref(), - &event.data, - ) { - return Err(ApiError::bad_request(format!("event {index}: {violation}"))); + // The v1 structural rules bind the "1.0" line only. A "0.1" document + // has already been normalised per spec 12.1 where a rule exists; what + // the preview line tolerated beyond that stays tolerated at ingest and + // is quarantined at materialisation instead. + if line == conformance::SchemaLine::V1_0 { + if let Some(violation) = conformance::v1_structural_violation( + &event.event_type, + event.id.is_some(), + event.output_id.as_deref(), + event.presentation_id.is_some(), + event.content_url.as_deref(), + event.content_id.as_deref(), + &event.data, + ) { + return Err(ApiError::bad_request(format!("event {index}: {violation}"))); + } + + if let Some(violation) = + conformance::source_role_violation(&event.event_type, event.source_role.as_deref()) + { + return Err(ApiError::bad_request(format!("event {index}: {violation}"))); + } + + if let Some(violation) = conformance::field_placement_violation( + &event.event_type, + event.presentation_id.is_some(), + event.ctx_token.is_some(), + event.citation_id.is_some(), + event.turn.is_some(), + ) { + return Err(ApiError::bad_request(format!("event {index}: {violation}"))); + } + + // A recorded event-level ctx_token is a claim about a minted token + // (spec 7.4.1); a value outside the token grammar can never have + // been minted, so it is rejected rather than stored. + if let Some(token) = event.ctx_token.as_deref() + && !click_tokens::ctx_token_well_formed(token) + { + return Err(ApiError::bad_request(format!( + "event {index}: malformed ctx_token; token values match \ + ^ct_[A-Za-z0-9_-]{{16,240}}$ (spec 7.4.1)" + ))); + } } Ok(()) @@ -153,10 +193,23 @@ pub fn check_sessionless(event: &TelemetryEventInput, index: usize) -> Result<() } /// Apply the spec's normalisations in place: fold recognised enum synonyms, -/// and drop fields v1 withdrew on privacy grounds. Returns the notes worth -/// logging so an emitter can be told what was changed. -pub fn normalise(event: &mut TelemetryEventInput) -> Vec { +/// drop fields v1 withdrew on privacy grounds, and — for documents on the +/// `"0.1"` line — apply the spec 12.1 migration defaults so preview events +/// read as the v1 events the migration rules define. Runs before +/// [`check_event`], so a migrated `"0.1"` event passes the checks its +/// defaults satisfy. Returns the notes worth logging so an emitter can be +/// told what was changed. +pub fn normalise(event: &mut TelemetryEventInput, line: conformance::SchemaLine) -> Vec { let mut notes = conformance::normalise_event_data_enums(&event.event_type, &mut event.data); + + if line == conformance::SchemaLine::V0_1 { + notes.extend(conformance::apply_v0_migration( + &event.event_type, + event.turn_id.as_deref(), + &mut event.data, + )); + } + notes.extend(conformance::strip_withdrawn_data_fields(&mut event.data)); if let Some(turn) = event.turn.as_mut() { diff --git a/scripts/smoke.sh b/scripts/smoke.sh index 1fe0657..6e84b05 100755 --- a/scripts/smoke.sh +++ b/scripts/smoke.sh @@ -75,11 +75,11 @@ RESP=$(curl -s -X POST "$BASE/events" \ -H 'content-type: application/json' -H "x-organization-id: $AGENT" \ -d "{ \"document_type\":\"event_batch\", - \"schema_version\":\"0.1\", + \"schema_version\":\"1.0\", \"session_id\":\"$SESSION\", \"events\":[ - {\"type\":\"content_grounded\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"data\":{\"grounding_scope\":\"session\"}}, - {\"id\":\"$CITED\",\"type\":\"content_cited\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"output_id\":\"out-1\"}, + {\"type\":\"content_grounded\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"data\":{\"scope\":\"session\"}}, + {\"id\":\"$CITED\",\"type\":\"content_cited\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"output_id\":\"out-1\",\"data\":{\"citation_type\":\"direct_quote\"}}, {\"id\":\"$PRESENTED\",\"type\":\"content_presented\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"output_id\":\"out-1\",\"citation_id\":\"$CITED\",\"data\":{\"presentation_kind\":\"content\",\"presentation_type\":\"summary\"}} ]}") [ "$(echo "$RESP" | jget 'd.get("events_created","ERR")')" = "3" ] \ @@ -132,7 +132,7 @@ RESP=$(curl -s -X POST "$BASE/events" -H 'content-type: application/json' -H "x- RESP=$(curl -s -X POST "$BASE/events" -H 'content-type: application/json' -H "x-organization-id: $PUBLISHER" \ -d "{\"ctx_token\":\"$TOKEN\",\"events\":[{\"id\":\"$(uuid)\",\"type\":\"content_cited\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"output_id\":\"out-9\"}]}") -echo "$RESP" | grep -q 'only carry content_engaged' && ok "ctx token cannot write non-engagement claims" || bad "expected ctx restriction, got: $RESP" +echo "$RESP" | grep -q 'content_engaged' && ok "ctx token cannot write non-engagement claims" || bad "expected ctx restriction, got: $RESP" echo "== standard session document" DOC=$(curl -s "$BASE/sessions/$SESSION/document" -H "x-organization-id: $AGENT") @@ -146,13 +146,13 @@ RESP=$(curl -s -X POST "$BASE/sessions/end" -H 'content-type: application/json' [ "$(echo "$RESP" | jget 'd.get("status","ERR")')" = "ok" ] && ok "session ended with outcome" || bad "end failed: $RESP" RESP=$(curl -s -X POST "$BASE/events" -H 'content-type: application/json' -H "x-organization-id: $AGENT" \ - -d "{\"session_id\":\"$SESSION\",\"events\":[{\"type\":\"content_grounded\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\"}]}") + -d "{\"session_id\":\"$SESSION\",\"events\":[{\"type\":\"content_grounded\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"data\":{\"scope\":\"session\"}}]}") echo "$RESP" | grep -q 'session has ended' && ok "events rejected after session end" || bad "expected ended-session error, got: $RESP" echo "== bulk session document" RESP=$(curl -s -X POST "$BASE/sessions/bulk" -H 'content-type: application/json' -H "x-organization-id: $AGENT" \ - -d "{\"document_type\":\"session\",\"schema_version\":\"0.1\",\"session_id\":\"$(uuid)\",\"agent_id\":\"demo-agent\", - \"events\":[{\"type\":\"content_grounded\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\"}], + -d "{\"document_type\":\"session\",\"schema_version\":\"1.0\",\"session_id\":\"$(uuid)\",\"agent_id\":\"demo-agent\", + \"events\":[{\"type\":\"content_grounded\",\"timestamp\":\"$TS\",\"content_url\":\"$URL\",\"data\":{\"scope\":\"session\"}}], \"outcome\":{\"type\":\"browse\"}}") [ "$(echo "$RESP" | jget 'd.get("events_created","ERR")')" = "1" ] && ok "bulk document ingested" || bad "bulk failed: $RESP" [ "$(echo "$RESP" | jget 'str(d.get("outcome_recorded"))')" = "True" ] && ok "bulk outcome recorded" || bad "bulk outcome not recorded: $RESP"