diff --git a/crates/content-telemetry-core/src/models/query.rs b/crates/content-telemetry-core/src/models/query.rs index 2e8e66e..85136c6 100644 --- a/crates/content-telemetry-core/src/models/query.rs +++ b/crates/content-telemetry-core/src/models/query.rs @@ -17,6 +17,7 @@ pub struct PublisherSummary { pub total_sessions: i64, pub events_by_type: Vec, pub events_by_source: Vec, + pub events_by_status: Vec, pub agents: Vec, pub period_start: Option>, pub period_end: Option>, @@ -29,6 +30,15 @@ pub struct SourceRoleCount { pub sessions: i64, } +#[derive(Debug, Clone, Serialize)] +pub struct StatusCodeCount { + /// HTTP status observed at the edge (event_data.response_status). None + /// bucket counts events whose source recorded no status — only the edge + /// enrichment profile stamps one, so self-reported events land there. + pub status: Option, + pub count: i64, +} + #[derive(Debug, Clone, Serialize)] pub struct AgentBreakdown { pub platform_id: Option, diff --git a/crates/content-telemetry-core/src/services/queries.rs b/crates/content-telemetry-core/src/services/queries.rs index 46c238d..0477c87 100644 --- a/crates/content-telemetry-core/src/services/queries.rs +++ b/crates/content-telemetry-core/src/services/queries.rs @@ -7,7 +7,7 @@ use crate::models::query::{ AgentAttestedCounts, AgentBreakdown, AgentDomainBreakdown, AgentDomainMetric, AgentReconciliation, AgentSummary, DayFunnelCount, DomainPreview, DomainPreviewUrl, EventTypeCount, Paginated, PublisherEvent, PublisherSummary, PublisherUrlMetric, - SourceRoleCount, UnclaimedDomain, + SourceRoleCount, StatusCodeCount, UnclaimedDomain, }; /// Get publisher summary with event counts filtered by their domains. @@ -93,11 +93,12 @@ pub async fn get_publisher_summary( q = q.bind(c); } - // Run event counts, source breakdown, and agent breakdown concurrently — + // Run event counts, source, status, and agent breakdowns concurrently — // they scan the same data independently, so parallelising cuts dashboard latency. - let (rows, source_rows, agents) = tokio::try_join!( + 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), )?; @@ -118,6 +119,13 @@ pub async fn get_publisher_summary( sessions: r.sessions, }) .collect(); + let events_by_status: Vec = status_rows + .into_iter() + .map(|r| StatusCodeCount { + status: r.status, + count: r.count, + }) + .collect(); Ok(PublisherSummary { organization_id: owner_id, @@ -127,6 +135,7 @@ pub async fn get_publisher_summary( total_sessions, events_by_type, events_by_source, + events_by_status, agents, period_start: since, period_end: until, @@ -1229,6 +1238,73 @@ async fn query_source_breakdown( q.fetch_all(pool).await } +async fn query_status_breakdown( + pool: &PgPool, + patterns: &[String], + since: Option>, + until: Option>, + bot: Option<&str>, + bot_category: Option<&str>, +) -> Result, sqlx::Error> { + let like_clauses: Vec = (1..=patterns.len()) + .map(|i| format!("e.content_url LIKE ${i}")) + .collect(); + let where_like = like_clauses.join(" OR "); + + let mut time_filter = String::new(); + let mut param_idx = patterns.len() + 1; + if since.is_some() { + time_filter.push_str(&format!(" AND e.event_timestamp >= ${param_idx}")); + param_idx += 1; + } + if until.is_some() { + time_filter.push_str(&format!(" AND e.event_timestamp <= ${param_idx}")); + param_idx += 1; + } + if bot.is_some() { + time_filter.push_str(&format!(" AND e.event_data->>'bot_name' = ${param_idx}")); + param_idx += 1; + } + if bot_category.is_some() { + time_filter.push_str(&format!( + " AND e.event_data->>'bot_category' = ${param_idx}" + )); + } + + // Only the edge enrichment profile stamps response_status; the guarded + // cast folds any event without a numeric status into the NULL bucket + // rather than failing the whole aggregation on malformed data. + let sql = format!( + "SELECT CASE WHEN e.event_data->>'response_status' ~ '^[0-9]+$' + THEN (e.event_data->>'response_status')::int + END as status, + COUNT(*) as count + FROM events e + WHERE ({where_like}){time_filter} + GROUP BY status + ORDER BY count DESC" + ); + + let mut q = sqlx::query_as::<_, StatusCodeRow>(&sql); + for p in patterns { + q = q.bind(p); + } + if let Some(ref s) = since { + q = q.bind(s); + } + if let Some(ref u) = until { + q = q.bind(u); + } + if let Some(b) = bot { + q = q.bind(b); + } + if let Some(c) = bot_category { + q = q.bind(c); + } + + q.fetch_all(pool).await +} + async fn query_agent_breakdown( pool: &PgPool, patterns: &[String], @@ -1715,6 +1791,7 @@ fn empty_summary(owner_id: Uuid, domains: &[String]) -> PublisherSummary { total_sessions: 0, events_by_type: vec![], events_by_source: vec![], + events_by_status: vec![], agents: vec![], period_start: None, period_end: None, @@ -1790,6 +1867,12 @@ struct SourceRoleRow { sessions: i64, } +#[derive(Debug, sqlx::FromRow)] +struct StatusCodeRow { + status: Option, + count: i64, +} + #[derive(Debug, sqlx::FromRow)] struct AgentSourceRow { platform_id: Option, diff --git a/crates/content-telemetry-core/tests/integration.rs b/crates/content-telemetry-core/tests/integration.rs index ac07a17..012edae 100644 --- a/crates/content-telemetry-core/tests/integration.rs +++ b/crates/content-telemetry-core/tests/integration.rs @@ -687,3 +687,78 @@ async fn publisher_category_filter_includes_stamped_agent_citations(pool: PgPool "inference category should include the stamped citation" ); } + +// =================================================================== +// Publisher summary: HTTP status-code breakdown. Only edge enrichment +// stamps response_status; everything else must land in the NULL bucket +// rather than being dropped or failing the aggregation. +// =================================================================== + +#[sqlx::test(migrations = "./migrations", fixtures("setup"))] +async fn publisher_summary_breaks_down_status_codes(pool: PgPool) { + let mk_edge = |status: serde_json::Value| EdgeEventInput { + id: None, + event_type: "content_retrieved".to_string(), + timestamp: Utc::now(), + source_role: Some("edge".to_string()), + content_telemetry_id: None, + content_url: Some("https://example.com/a".to_string()), + content_id: None, + license_ref: None, + data: serde_json::json!({"bot_name": "GPTBot", "bot_category": "training", "response_status": status}), + }; + let edge = vec![ + mk_edge(serde_json::json!(200)), + mk_edge(serde_json::json!(200)), + mk_edge(serde_json::json!(404)), + // Malformed status must fold into the NULL bucket, not error the query. + mk_edge(serde_json::json!("cached")), + ]; + events::create_edge_events(&pool, publisher_org_id(), &edge) + .await + .unwrap(); + + // A self-report event with no response_status also lands in NULL. + let agent_org = agent_org_id(); + let session = sessions::create_session(&pool, agent_org, &minimal_session_request()) + .await + .unwrap(); + let inputs = vec![agent_event( + "content_cited", + "https://example.com/a", + serde_json::json!({}), + )]; + events::create_events(&pool, session.id, agent_org, &inputs) + .await + .unwrap(); + + let domains = vec!["example.com".to_string()]; + let summary = queries::get_publisher_summary( + &pool, + publisher_org_id(), + &domains, + None, + None, + None, + None, + None, + ) + .await + .unwrap(); + + let count_for = |status: Option| { + summary + .events_by_status + .iter() + .find(|r| r.status == status) + .map(|r| r.count) + .unwrap_or(0) + }; + assert_eq!(count_for(Some(200)), 2, "two 200 edge retrievals"); + assert_eq!(count_for(Some(404)), 1, "one 404 edge retrieval"); + assert_eq!( + count_for(None), + 2, + "malformed and statusless events fold into the NULL bucket" + ); +}