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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions crates/content-telemetry-core/src/models/query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ pub struct PublisherSummary {
pub total_sessions: i64,
pub events_by_type: Vec<EventTypeCount>,
pub events_by_source: Vec<SourceRoleCount>,
pub events_by_status: Vec<StatusCodeCount>,
pub agents: Vec<AgentBreakdown>,
pub period_start: Option<DateTime<Utc>>,
pub period_end: Option<DateTime<Utc>>,
Expand All @@ -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<i32>,
pub count: i64,
}

#[derive(Debug, Clone, Serialize)]
pub struct AgentBreakdown {
pub platform_id: Option<String>,
Expand Down
89 changes: 86 additions & 3 deletions crates/content-telemetry-core/src/services/queries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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),
)?;

Expand All @@ -118,6 +119,13 @@ pub async fn get_publisher_summary(
sessions: r.sessions,
})
.collect();
let events_by_status: Vec<StatusCodeCount> = status_rows
.into_iter()
.map(|r| StatusCodeCount {
status: r.status,
count: r.count,
})
.collect();

Ok(PublisherSummary {
organization_id: owner_id,
Expand All @@ -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,
Expand Down Expand Up @@ -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<DateTime<Utc>>,
until: Option<DateTime<Utc>>,
bot: Option<&str>,
bot_category: Option<&str>,
) -> Result<Vec<StatusCodeRow>, sqlx::Error> {
let like_clauses: Vec<String> = (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],
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -1790,6 +1867,12 @@ struct SourceRoleRow {
sessions: i64,
}

#[derive(Debug, sqlx::FromRow)]
struct StatusCodeRow {
status: Option<i32>,
count: i64,
}

#[derive(Debug, sqlx::FromRow)]
struct AgentSourceRow {
platform_id: Option<String>,
Expand Down
75 changes: 75 additions & 0 deletions crates/content-telemetry-core/tests/integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<i32>| {
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"
);
}
Loading