diff --git a/packages/backend/convex/_generated/api.d.ts b/packages/backend/convex/_generated/api.d.ts index 9bf908a..011af5d 100644 --- a/packages/backend/convex/_generated/api.d.ts +++ b/packages/backend/convex/_generated/api.d.ts @@ -32,6 +32,7 @@ import type * as lib_eventClaimCoverage from "../lib/eventClaimCoverage.js"; import type * as lib_openai from "../lib/openai.js"; import type * as mbfc from "../mbfc.js"; import type * as migrations from "../migrations.js"; +import type * as pipelineDiagnostics from "../pipelineDiagnostics.js"; import type * as privateData from "../privateData.js"; import type * as prompts from "../prompts.js"; import type * as seeds from "../seeds.js"; @@ -76,6 +77,7 @@ declare const fullApi: ApiFromModules<{ "lib/openai": typeof lib_openai; mbfc: typeof mbfc; migrations: typeof migrations; + pipelineDiagnostics: typeof pipelineDiagnostics; privateData: typeof privateData; prompts: typeof prompts; seeds: typeof seeds; diff --git a/packages/backend/convex/aiBudget.ts b/packages/backend/convex/aiBudget.ts index 98f5f59..71dfde4 100644 --- a/packages/backend/convex/aiBudget.ts +++ b/packages/backend/convex/aiBudget.ts @@ -28,6 +28,8 @@ import type { MutationCtx, QueryCtx } from "./_generated/server"; /** Default model pricing as of 2025. Add new models as needed. */ const DEFAULT_MODEL_RATES: Record = { + "gpt-5-nano": { input: 0.00000005, output: 0.0000004 }, + "gpt-5-mini": { input: 0.00000025, output: 0.000002 }, "gpt-4o-mini": { input: 0.00000015, output: 0.0000006 }, "gpt-4o": { input: 0.0000025, output: 0.00001 }, "gpt-4.1-nano": { input: 0.0000001, output: 0.0000004 }, @@ -67,14 +69,33 @@ export function calculateCost( model: string, inputTokens: number, outputTokens: number, +): number { + return calculateCostWithCachedInput(model, inputTokens, 0, outputTokens); +} + +export function calculateCostWithCachedInput( + model: string, + inputTokens: number, + cachedInputTokens: number, + outputTokens: number, ): number { const rates = MODEL_RATES[model]; + const cachedInputRate = rates ? rates.input * 0.1 : undefined; if (!rates) { // Unknown model — use gpt-4o-mini rates as conservative fallback const fallback = MODEL_RATES["gpt-4o-mini"]!; return inputTokens * fallback.input + outputTokens * fallback.output; } - return inputTokens * rates.input + outputTokens * rates.output; + const safeCachedInputTokens = Math.min( + Math.max(0, cachedInputTokens), + Math.max(0, inputTokens), + ); + const billableInputTokens = Math.max(0, inputTokens - safeCachedInputTokens); + return ( + billableInputTokens * rates.input + + safeCachedInputTokens * (cachedInputRate ?? rates.input) + + outputTokens * rates.output + ); } // --------------------------------------------------------------------------- @@ -329,6 +350,7 @@ export const logUsage = internalMutation({ callType: v.optional(v.string()), inputTokens: v.number(), outputTokens: v.number(), + cachedInputTokens: v.optional(v.number()), costUsd: v.number(), eventId: v.optional(v.id("events")), articleId: v.optional(v.id("articles")), @@ -350,6 +372,7 @@ async function recordUsageInternal( callType: string; inputTokens: number; outputTokens: number; + cachedInputTokens?: number; costUsd: number; eventId?: Id<"events">; articleId?: Id<"articles">; @@ -391,6 +414,7 @@ async function recordUsageInternal( callType: args.callType, inputTokens: args.inputTokens, outputTokens: args.outputTokens, + cachedInputTokens: args.cachedInputTokens, costUsd: args.costUsd, eventId: args.eventId, articleId: args.articleId, @@ -414,6 +438,7 @@ export const recordUsage = internalMutation({ model: v.string(), inputTokens: v.number(), outputTokens: v.number(), + cachedInputTokens: v.optional(v.number()), costUsd: v.number(), eventId: v.optional(v.id("events")), articleId: v.optional(v.id("articles")), @@ -449,6 +474,7 @@ export const getTodaysUsage = internalQuery({ calls: number; inputTokens: number; outputTokens: number; + cachedInputTokens: number; costUsd: number; latencyMs: number; }; @@ -464,12 +490,14 @@ export const getTodaysUsage = internalQuery({ calls: 0, inputTokens: 0, outputTokens: 0, + cachedInputTokens: 0, costUsd: 0, latencyMs: 0, }; existing.calls++; existing.inputTokens += row.inputTokens; existing.outputTokens += row.outputTokens; + existing.cachedInputTokens += row.cachedInputTokens ?? 0; existing.costUsd += row.costUsd; existing.latencyMs += row.latencyMs ?? 0; group[key] = existing; diff --git a/packages/backend/convex/claimDivergenceNode.ts b/packages/backend/convex/claimDivergenceNode.ts index 7c997c3..bde5dbc 100644 --- a/packages/backend/convex/claimDivergenceNode.ts +++ b/packages/backend/convex/claimDivergenceNode.ts @@ -14,7 +14,7 @@ import { type ClaimType, } from "./prompts"; -const DEFAULT_MODEL = "gpt-4o-mini"; +const DEFAULT_MODEL = "gpt-5-nano"; const DEFAULT_ENABLED = true; const DEFAULT_BATCH_SIZE = 4; const DEFAULT_SCAN_LIMIT = 60; diff --git a/packages/backend/convex/config.ts b/packages/backend/convex/config.ts index dc35ac6..b517db9 100644 --- a/packages/backend/convex/config.ts +++ b/packages/backend/convex/config.ts @@ -420,7 +420,7 @@ export const seedDefaults = internalMutation({ }, { key: "event_summary_model", - value: "gpt-4o-mini", + value: "gpt-5-nano", description: "OpenAI chat model used for event perspective summaries.", }, @@ -468,7 +468,7 @@ export const seedDefaults = internalMutation({ }, { key: "article_fact_extraction_model", - value: "gpt-4o-mini", + value: "gpt-5-nano", description: "OpenAI chat model used to extract atomic facts from articles during enrichment.", }, @@ -498,7 +498,7 @@ export const seedDefaults = internalMutation({ }, { key: "article_bias_detection_model", - value: "gpt-4o-mini", + value: "gpt-5-nano", description: "OpenAI chat model used for per-article bias component scoring during enrichment.", }, @@ -552,7 +552,7 @@ export const seedDefaults = internalMutation({ }, { key: "claim_analysis_model", - value: "gpt-4o-mini", + value: "gpt-5-nano", description: "OpenAI chat model used for event-level claim divergence analysis.", }, diff --git a/packages/backend/convex/enrichment.ts b/packages/backend/convex/enrichment.ts index 7f0d732..3e8e9f9 100644 --- a/packages/backend/convex/enrichment.ts +++ b/packages/backend/convex/enrichment.ts @@ -11,12 +11,25 @@ import type { MutationCtx } from "./_generated/server"; import type { Doc } from "./_generated/dataModel"; import { refreshEventClaimCoverage } from "./lib/eventClaimCoverage"; -const ARTICLE_AI_STATUS_VALIDATOR = v.union( +const ARTICLE_FACT_STATUS_VALIDATOR = v.union( + v.literal("pending"), + v.literal("deferred"), v.literal("succeeded"), + v.literal("succeeded_empty"), v.literal("failed"), v.literal("skipped"), ); +const ARTICLE_BIAS_STATUS_VALIDATOR = v.union( + v.literal("deferred"), + v.literal("succeeded"), + v.literal("failed"), + v.literal("skipped"), +); + +const MAX_FACT_EXTRACTION_ATTEMPTS = 3; +const MAX_BIAS_DETECTION_ATTEMPTS = 3; + function sourceBiasLabel(source: Doc<"sources">): string { const mbfcCategory = source.mbfcCategory?.toLowerCase(); if ( @@ -111,6 +124,7 @@ async function toClaimedArticle( sourceName: source.name, sourceLean: sourceBiasLabel(source), sourceReliability: source.reliabilityScore, + previousStatus: article.status, }; } @@ -172,18 +186,30 @@ function articleNeedsReenrichment( if ((article.summary ?? "").trim().length < 120) return true; if ( article.atomicFacts === undefined && - article.factExtractionStatus !== "skipped" + article.factExtractionStatus !== "skipped" && + article.factExtractionStatus !== "succeeded_empty" + ) { + return true; + } + if ( + article.factExtractionStatus === "failed" && + (article.factExtractionAttempts ?? 0) < MAX_FACT_EXTRACTION_ATTEMPTS ) { return true; } - if (article.factExtractionStatus === "failed") return true; + if (article.factExtractionStatus === "deferred") return true; if ( article.biasAnalyzedAt === undefined && article.biasDetectionStatus !== "skipped" ) { return true; } - if (article.biasDetectionStatus === "failed") return true; + if ( + article.biasDetectionStatus === "failed" || + article.biasDetectionStatus === "deferred" + ) { + return (article.biasDetectionAttempts ?? 0) < MAX_BIAS_DETECTION_ATTEMPTS; + } if (article.url.includes("news.google.com")) return true; if (article.canonicalUrl.includes("news.google.com")) return true; return false; @@ -242,6 +268,163 @@ export const claimArticlesForReenrichment = internalMutation({ }, }); +export const claimArticlesNeedingFactExtraction = internalMutation({ + args: { + limit: v.number(), + runId: v.string(), + leaseExpiresAt: v.number(), + beforePublishedAt: v.optional(v.number()), + includeFailed: v.optional(v.boolean()), + includeSucceededEmpty: v.optional(v.boolean()), + }, + handler: async ( + ctx, + { + limit, + runId, + leaseExpiresAt, + beforePublishedAt, + includeFailed, + includeSucceededEmpty, + }, + ) => { + const batchSize = Math.max(0, Math.floor(limit)); + if (batchSize === 0) return []; + + const now = Date.now(); + const candidates = await ctx.db + .query("articles") + .withIndex("by_published", (q) => + beforePublishedAt ? q.lt("publishedAt", beforePublishedAt) : q, + ) + .order("desc") + .take(Math.min(1000, Math.max(batchSize * 10, batchSize))); + + const claimed = []; + for (const article of candidates) { + if (claimed.length >= batchSize) break; + if (article.status === "discarded") continue; + if ( + article.status === "processing" && + (article.enrichmentLeaseExpiresAt ?? 0) > now + ) { + continue; + } + if ((article.atomicFacts ?? []).some((fact) => fact.trim().length > 0)) { + continue; + } + if (article.factExtractionStatus === "skipped") continue; + if (article.factExtractionStatus === "deferred") { + // Deferred rows are first-class retry candidates. + } else if (article.factExtractionStatus === "succeeded_empty") { + if (!includeSucceededEmpty) continue; + } else if (article.factExtractionStatus === "succeeded") { + if (!includeSucceededEmpty) continue; + } + if (article.factExtractionStatus === "failed" && !includeFailed) { + continue; + } + if ( + article.factExtractionStatus === "failed" && + (article.factExtractionAttempts ?? 0) >= MAX_FACT_EXTRACTION_ATTEMPTS + ) { + continue; + } + + const enriched = await toClaimedArticle(ctx, article); + if (!enriched) continue; + + await ctx.db.patch(article._id, { + status: "processing", + enrichmentRunId: runId, + enrichmentLeaseExpiresAt: leaseExpiresAt, + }); + claimed.push(enriched); + } + + return claimed; + }, +}); + +export const deferArticleFactExtraction = internalMutation({ + args: { + articleId: v.id("articles"), + runId: v.string(), + previousStatus: v.union( + v.literal("unprocessed"), + v.literal("processing"), + v.literal("enriched"), + v.literal("clustered"), + v.literal("discarded"), + ), + reason: v.string(), + attemptedAt: v.number(), + }, + handler: async (ctx, { articleId, runId, previousStatus, reason, attemptedAt }) => { + const article = await ctx.db.get(articleId); + if ( + !article || + article.status !== "processing" || + article.enrichmentRunId !== runId + ) { + return { updated: false, eventId: undefined }; + } + + await ctx.db.patch(articleId, { + status: previousStatus === "processing" ? "unprocessed" : previousStatus, + enrichmentRunId: undefined, + enrichmentLeaseExpiresAt: undefined, + factExtractionStatus: "deferred", + factExtractionError: reason.slice(0, 500), + factExtractionAttempts: (article.factExtractionAttempts ?? 0) + 1, + factExtractionLastAttemptAt: attemptedAt, + }); + + if (article.eventId) { + await refreshEventClaimCoverage(ctx, article.eventId); + } + + return { updated: true, eventId: article.eventId }; + }, +}); + +export const deferArticleBiasDetection = internalMutation({ + args: { + articleId: v.id("articles"), + runId: v.string(), + previousStatus: v.union( + v.literal("unprocessed"), + v.literal("processing"), + v.literal("enriched"), + v.literal("clustered"), + v.literal("discarded"), + ), + reason: v.string(), + }, + handler: async (ctx, { articleId, runId, previousStatus, reason }) => { + const article = await ctx.db.get(articleId); + if ( + !article || + article.status !== "processing" || + article.enrichmentRunId !== runId + ) { + return { updated: false, eventId: undefined }; + } + + await ctx.db.patch(articleId, { + status: previousStatus === "processing" ? "unprocessed" : previousStatus, + enrichmentRunId: undefined, + enrichmentLeaseExpiresAt: undefined, + biasDetectionStatus: "deferred", + biasDetectionError: reason.slice(0, 500), + biasDetectionAttempts: (article.biasDetectionAttempts ?? 0) + 1, + biasDetectionLastAttemptAt: Date.now(), + }); + + return { updated: true, eventId: article.eventId }; + }, +}); + export const claimEventArticlesForReenrichment = internalMutation({ args: { eventId: v.id("events"), @@ -298,11 +481,11 @@ export const markArticleEnriched = internalMutation({ sourceBiasDelta: v.optional(v.number()), sourceBiasOutlierFlag: v.optional(v.boolean()), biasAnalyzedAt: v.optional(v.number()), - biasDetectionStatus: v.optional(ARTICLE_AI_STATUS_VALIDATOR), + biasDetectionStatus: v.optional(ARTICLE_BIAS_STATUS_VALIDATOR), biasDetectionError: v.optional(v.string()), summary: v.optional(v.string()), atomicFacts: v.optional(v.array(v.string())), - factExtractionStatus: v.optional(ARTICLE_AI_STATUS_VALIDATOR), + factExtractionStatus: v.optional(ARTICLE_FACT_STATUS_VALIDATOR), factExtractionError: v.optional(v.string()), factExtractedAt: v.optional(v.number()), resolvedUrl: v.optional(v.string()), @@ -371,6 +554,16 @@ export const markArticleEnriched = internalMutation({ await ctx.db.delete(row._id); } + const shouldRecordFactAttempt = + factExtractionStatus !== undefined && + factExtractionStatus !== "skipped" && + !( + article.factExtractionStatus === "succeeded_empty" && + factExtractionStatus === "succeeded_empty" + ); + const shouldRecordBiasAttempt = + biasDetectionStatus !== undefined && biasDetectionStatus !== "skipped"; + // Store embedding in dedicated table (hot/cold split) await ctx.db.insert("articleEmbeddings", { articleId, @@ -390,6 +583,10 @@ export const markArticleEnriched = internalMutation({ ? { biasDetectionStatus, biasDetectionError, + biasDetectionAttempts: shouldRecordBiasAttempt + ? (article.biasDetectionAttempts ?? 0) + 1 + : article.biasDetectionAttempts, + biasDetectionLastAttemptAt: biasAnalyzedAt ?? Date.now(), } : {}), summary: summary ?? article.summary, @@ -399,6 +596,10 @@ export const markArticleEnriched = internalMutation({ factExtractionStatus, factExtractionError, factExtractedAt: factExtractedAt ?? article.factExtractedAt, + factExtractionAttempts: shouldRecordFactAttempt + ? (article.factExtractionAttempts ?? 0) + 1 + : article.factExtractionAttempts, + factExtractionLastAttemptAt: factExtractedAt ?? Date.now(), } : {}), url: resolvedUrl ?? article.url, diff --git a/packages/backend/convex/enrichmentNode.ts b/packages/backend/convex/enrichmentNode.ts index cd28261..0ec323b 100644 --- a/packages/backend/convex/enrichmentNode.ts +++ b/packages/backend/convex/enrichmentNode.ts @@ -39,13 +39,13 @@ const ARTICLE_LEASE_TTL_MS = 15 * 60 * 1000; /** OpenAI embedding model — cheap & effective for clustering */ const EMBEDDING_MODEL = "text-embedding-3-small"; -const DEFAULT_FACT_EXTRACTION_MODEL = "gpt-4o-mini"; +const DEFAULT_FACT_EXTRACTION_MODEL = "gpt-5-nano"; const DEFAULT_FACT_EXTRACTION_ENABLED = true; const DEFAULT_FACT_EXTRACTION_MAX_ARTICLES = 20; const DEFAULT_FACT_EXTRACTION_MAX_FACTS = 8; const DEFAULT_FACT_EXTRACTION_MAX_INPUT_CHARS = 2600; const DEFAULT_BIAS_DETECTION_ENABLED = true; -const DEFAULT_BIAS_DETECTION_MODEL = "gpt-4o-mini"; +const DEFAULT_BIAS_DETECTION_MODEL = "gpt-5-nano"; const DEFAULT_BIAS_DETECTION_MAX_ARTICLES = 20; const DEFAULT_BIAS_DETECTION_MAX_INPUT_CHARS = 6000; const DEFAULT_BIAS_SOURCE_DELTA_THRESHOLD = 2; @@ -138,11 +138,33 @@ type ArticleBiasResult = { }; type ArticleAiStatus = { - status: "succeeded" | "failed" | "skipped"; + status: + | "pending" + | "deferred" + | "succeeded" + | "succeeded_empty" + | "failed" + | "skipped"; error?: string; analyzedAt?: number; }; +type ArticleBiasStatus = "deferred" | "succeeded" | "failed" | "skipped"; + +function toArticleBiasStatus( + status: ArticleAiStatus["status"] | undefined, +): ArticleBiasStatus | undefined { + if ( + status === "deferred" || + status === "succeeded" || + status === "failed" || + status === "skipped" + ) { + return status; + } + return undefined; +} + type FactExtractionResult = { factsByArticleId: Map; statusByArticleId: Map, ArticleAiStatus>; @@ -172,6 +194,7 @@ type PreparedArticle = { resolvedUrl?: string; sourceLean: string; sourceReliability: number; + previousStatus: "unprocessed" | "processing" | "enriched" | "clustered" | "discarded"; imageUrl?: string; imageWidth?: number; imageHeight?: number; @@ -324,10 +347,6 @@ function parseAtomicFactsResponse( ); } - for (const id of articleIds) { - if (!result.has(id)) result.set(id, []); - } - return result; } @@ -430,6 +449,15 @@ function parseBiasScoringResponse( return results; } +function chunkArray(items: T[], chunkSize: number): T[][] { + const size = Math.max(1, Math.floor(chunkSize)); + const chunks: T[][] = []; + for (let i = 0; i < items.length; i += size) { + chunks.push(items.slice(i, i + size)); + } + return chunks; +} + async function scoreBiasForArticles( ctx: ActionCtx, articles: PreparedArticle[], @@ -441,10 +469,11 @@ async function scoreBiasForArticles( article.embeddingText.trim().length > 0 || article.extractedSummary?.trim() || article.rssSnippet?.trim(), - ) - .slice(0, settings.maxArticles); + ); const statusByArticleId = new Map, ArticleAiStatus>(); - const markSelected = ( + const biasByArticleId = new Map, ArticleBiasResult>(); + const markArticles = ( + targetArticles: PreparedArticle[], status: ArticleAiStatus["status"], error?: string, analyzedAt?: number, @@ -452,7 +481,7 @@ async function scoreBiasForArticles( const nextStatus: ArticleAiStatus = { status }; if (error !== undefined) nextStatus.error = error; if (analyzedAt !== undefined) nextStatus.analyzedAt = analyzedAt; - for (const article of selectedArticles) { + for (const article of targetArticles) { statusByArticleId.set(article._id, nextStatus); } }; @@ -462,90 +491,93 @@ async function scoreBiasForArticles( } if (!settings.enabled) { - markSelected("skipped", "Article bias detection is disabled"); + markArticles(selectedArticles, "skipped", "Article bias detection is disabled"); return { biasByArticleId: new Map(), statusByArticleId }; } - const budget = await ctx.runQuery(internal.aiBudget.checkBudget, {}); - if (!budget.allowed) { - console.warn( - `[enrichment] AI budget exhausted before bias scoring ($${budget.spentUsd}/$${budget.dailyLimitUsd}); skipping article bias detection`, - ); - markSelected("failed", "AI budget exhausted"); - return { biasByArticleId: new Map(), statusByArticleId }; - } - - const prompt = buildArticleBiasScoringPrompt({ - maxInputChars: settings.maxInputChars, - articles: selectedArticles.map((article) => ({ - id: article._id, - title: article.title, - sourceName: article.sourceName, - sourceLean: article.sourceLean, - sourceReliability: article.sourceReliability, - publishedAt: new Date(article.publishedAt).toISOString(), - summary: article.extractedSummary, - rssSnippet: article.rssSnippet ?? undefined, - bodyText: article.embeddingText.slice(0, settings.maxInputChars), - })), - }); + for (const chunk of chunkArray(selectedArticles, settings.maxArticles)) { + const budget = await ctx.runQuery(internal.aiBudget.checkBudget, {}); + if (!budget.allowed) { + console.warn( + `[enrichment] AI budget exhausted before bias scoring ($${budget.spentUsd}/$${budget.dailyLimitUsd}); skipping remaining article bias detection`, + ); + markArticles(chunk, "deferred", "AI budget exhausted", Date.now()); + continue; + } - try { - const response = await callOpenAI({ - kind: "chat", - model: settings.model, - temperature: 0, - maxTokens: Math.min(2500, 220 + selectedArticles.length * 120), - responseFormat: { - type: "json_schema", - json_schema: ARTICLE_BIAS_JSON_SCHEMA, - }, - messages: [ - { role: "system", content: prompt.system }, - { role: "user", content: prompt.user }, - ], - context: { callType: "bias_scoring" }, - runtime: ctx, + const prompt = buildArticleBiasScoringPrompt({ + maxInputChars: settings.maxInputChars, + articles: chunk.map((article) => ({ + id: article._id, + title: article.title, + sourceName: article.sourceName, + sourceLean: article.sourceLean, + sourceReliability: article.sourceReliability, + publishedAt: new Date(article.publishedAt).toISOString(), + summary: article.extractedSummary, + rssSnippet: article.rssSnippet ?? undefined, + bodyText: article.embeddingText.slice(0, settings.maxInputChars), + })), }); - const content = response.result; - if (!content) { - throw new Error( - response.error ?? "Bias scoring returned an empty response", + try { + const response = await callOpenAI({ + kind: "chat", + model: settings.model, + temperature: 0, + maxTokens: Math.min(2500, 220 + chunk.length * 120), + responseFormat: { + type: "json_schema", + json_schema: ARTICLE_BIAS_JSON_SCHEMA, + }, + messages: [ + { role: "system", content: prompt.system }, + { role: "user", content: prompt.user }, + ], + context: { callType: "bias_scoring" }, + runtime: ctx, + }); + + const content = response.result; + if (!content) { + throw new Error( + response.error ?? "Bias scoring returned an empty response", + ); + } + + const results = parseBiasScoringResponse( + content, + chunk, + settings.sourceDeltaThreshold, ); - } + const analyzedAt = Date.now(); + for (const [articleId, result] of results) { + biasByArticleId.set(articleId, result); + } + for (const article of chunk) { + statusByArticleId.set( + article._id, + results.has(article._id) + ? { status: "succeeded", analyzedAt } + : { + status: "deferred", + error: "Bias scoring omitted this article", + analyzedAt, + }, + ); + } - const results = parseBiasScoringResponse( - content, - selectedArticles, - settings.sourceDeltaThreshold, - ); - const analyzedAt = Date.now(); - for (const article of selectedArticles) { - statusByArticleId.set( - article._id, - results.has(article._id) - ? { status: "succeeded", analyzedAt } - : { status: "failed", error: "Bias scoring omitted this article" }, + console.log( + `[enrichment] Scored bias for ${results.size}/${chunk.length} articles (${response.usage.inputTokens}/${response.usage.outputTokens} tokens)`, ); + } catch (error) { + const message = error instanceof Error ? error.message : "Unknown error"; + console.error(`[enrichment] Article bias scoring failed: ${message}`); + markArticles(chunk, "failed", message.slice(0, 500)); } - - const inputTokens = response.usage.inputTokens; - const outputTokens = response.usage.outputTokens; - - console.log( - `[enrichment] Scored bias for ${results.size}/${selectedArticles.length} articles (${inputTokens}/${outputTokens} tokens)`, - ); - - return { biasByArticleId: results, statusByArticleId }; - } catch (error) { - const message = error instanceof Error ? error.message : "Unknown error"; - console.error( - `[enrichment] Article bias scoring failed: ${message}`, - ); - markSelected("failed", message.slice(0, 500)); - return { biasByArticleId: new Map(), statusByArticleId }; } + + return { biasByArticleId, statusByArticleId }; } async function extractAtomicFactsForArticles( @@ -559,10 +591,11 @@ async function extractAtomicFactsForArticles( article.embeddingText.trim().length > 0 || article.extractedSummary?.trim() || article.rssSnippet?.trim(), - ) - .slice(0, settings.maxArticles); + ); const statusByArticleId = new Map, ArticleAiStatus>(); - const markSelected = ( + const factsByArticleId = new Map(); + const markArticles = ( + targetArticles: PreparedArticle[], status: ArticleAiStatus["status"], error?: string, analyzedAt?: number, @@ -570,7 +603,7 @@ async function extractAtomicFactsForArticles( const nextStatus: ArticleAiStatus = { status }; if (error !== undefined) nextStatus.error = error; if (analyzedAt !== undefined) nextStatus.analyzedAt = analyzedAt; - for (const article of selectedArticles) { + for (const article of targetArticles) { statusByArticleId.set(article._id, nextStatus); } }; @@ -580,94 +613,102 @@ async function extractAtomicFactsForArticles( } if (!settings.enabled) { - markSelected("skipped", "Atomic fact extraction is disabled"); + markArticles(selectedArticles, "skipped", "Atomic fact extraction is disabled"); return { factsByArticleId: new Map(), statusByArticleId }; } - const budget = await ctx.runQuery(internal.aiBudget.checkBudget, {}); - if (!budget.allowed) { - console.warn( - `[enrichment] AI budget exhausted before fact extraction ($${budget.spentUsd}/$${budget.dailyLimitUsd}); skipping atomic facts`, - ); - markSelected("failed", "AI budget exhausted"); - return { factsByArticleId: new Map(), statusByArticleId }; - } - - const prompt = buildArticleFactExtractionPrompt({ - maxFactsPerArticle: settings.maxFactsPerArticle, - articles: selectedArticles.map((article) => ({ - id: article._id, - title: article.title, - sourceName: article.sourceName, - publishedAt: new Date(article.publishedAt).toISOString(), - entities: article.entities, - summary: article.extractedSummary, - rssSnippet: article.rssSnippet ?? undefined, - bodyText: article.embeddingText.slice(0, settings.maxInputChars), - })), - }); - - try { - const response = await callOpenAI({ - kind: "chat", - model: settings.model, - temperature: 0, - maxTokens: Math.min( - 3000, - 250 + selectedArticles.length * settings.maxFactsPerArticle * 32, - ), - responseFormat: { - type: "json_schema", - json_schema: ARTICLE_FACTS_JSON_SCHEMA, - }, - messages: [ - { role: "system", content: prompt.system }, - { role: "user", content: prompt.user }, - ], - context: { callType: "fact_extraction" }, - runtime: ctx, - }); - - const content = response.result; - if (!content) { - throw new Error( - response.error ?? "Fact extraction returned an empty response", + for (const chunk of chunkArray(selectedArticles, settings.maxArticles)) { + const budget = await ctx.runQuery(internal.aiBudget.checkBudget, {}); + if (!budget.allowed) { + console.warn( + `[enrichment] AI budget exhausted before fact extraction ($${budget.spentUsd}/$${budget.dailyLimitUsd}); skipping remaining atomic facts`, ); + markArticles(chunk, "deferred", "AI budget exhausted", Date.now()); + continue; } - const factsByArticleId = parseAtomicFactsResponse( - content, - new Set(selectedArticles.map((article) => article._id)), - settings.maxFactsPerArticle, - ); - const analyzedAt = Date.now(); - for (const article of selectedArticles) { - statusByArticleId.set(article._id, { - status: "succeeded", - analyzedAt, + const prompt = buildArticleFactExtractionPrompt({ + maxFactsPerArticle: settings.maxFactsPerArticle, + articles: chunk.map((article) => ({ + id: article._id, + title: article.title, + sourceName: article.sourceName, + publishedAt: new Date(article.publishedAt).toISOString(), + entities: article.entities, + summary: article.extractedSummary, + rssSnippet: article.rssSnippet ?? undefined, + bodyText: article.embeddingText.slice(0, settings.maxInputChars), + })), + }); + + try { + const response = await callOpenAI({ + kind: "chat", + model: settings.model, + temperature: 0, + maxTokens: Math.min( + 3000, + 250 + chunk.length * settings.maxFactsPerArticle * 32, + ), + responseFormat: { + type: "json_schema", + json_schema: ARTICLE_FACTS_JSON_SCHEMA, + }, + messages: [ + { role: "system", content: prompt.system }, + { role: "user", content: prompt.user }, + ], + context: { callType: "fact_extraction" }, + runtime: ctx, }); - } - const inputTokens = response.usage.inputTokens; - const outputTokens = response.usage.outputTokens; + const content = response.result; + if (!content) { + throw new Error( + response.error ?? "Fact extraction returned an empty response", + ); + } - const factCount = Array.from(factsByArticleId.values()).reduce( - (sum, facts) => sum + facts.length, - 0, - ); - console.log( - `[enrichment] Extracted ${factCount} atomic facts across ${factsByArticleId.size} articles (${inputTokens}/${outputTokens} tokens)`, - ); + const chunkFacts = parseAtomicFactsResponse( + content, + new Set(chunk.map((article) => article._id)), + settings.maxFactsPerArticle, + ); + for (const [articleId, facts] of chunkFacts) { + factsByArticleId.set(articleId, facts); + } + const analyzedAt = Date.now(); + for (const article of chunk) { + if (!chunkFacts.has(article._id)) { + statusByArticleId.set(article._id, { + status: "deferred", + error: "Fact extraction omitted this article", + analyzedAt, + }); + continue; + } + const facts = chunkFacts.get(article._id) ?? []; + statusByArticleId.set(article._id, { + status: facts.length > 0 ? "succeeded" : "succeeded_empty", + analyzedAt, + }); + } - return { factsByArticleId, statusByArticleId }; - } catch (error) { - const message = error instanceof Error ? error.message : "Unknown error"; - console.error( - `[enrichment] Atomic fact extraction failed: ${message}`, - ); - markSelected("failed", message.slice(0, 500)); - return { factsByArticleId: new Map(), statusByArticleId }; + const factCount = Array.from(chunkFacts.values()).reduce( + (sum, facts) => sum + facts.length, + 0, + ); + console.log( + `[enrichment] Extracted ${factCount} atomic facts across ${chunkFacts.size} articles (${response.usage.inputTokens}/${response.usage.outputTokens} tokens)`, + ); + } catch (error) { + const message = error instanceof Error ? error.message : "Unknown error"; + console.error(`[enrichment] Atomic fact extraction failed: ${message}`); + markArticles(chunk, "failed", message.slice(0, 500)); + } } + + return { factsByArticleId, statusByArticleId }; } async function mapWithConcurrency( @@ -708,6 +749,7 @@ async function runEnrichmentBatch( sourceName?: string; sourceLean: string; sourceReliability: number; + previousStatus: "unprocessed" | "processing" | "enriched" | "clustered" | "discarded"; }>, runId: string, ): Promise<{ @@ -835,10 +877,20 @@ async function runEnrichmentBatch( ), }; - const [factResult, biasResult] = await Promise.all([ - extractAtomicFactsForArticles(ctx, preparedArticles, factSettings), - scoreBiasForArticles(ctx, preparedArticles, biasSettings), - ]); + const factResult = await extractAtomicFactsForArticles( + ctx, + preparedArticles, + factSettings, + ); + const biasEligibleArticles = preparedArticles.filter( + (article) => + factResult.statusByArticleId.get(article._id)?.status !== "deferred", + ); + const biasResult = await scoreBiasForArticles( + ctx, + biasEligibleArticles, + biasSettings, + ); const factsByArticleId = factResult.factsByArticleId; const biasByArticleId = biasResult.biasByArticleId; @@ -855,6 +907,51 @@ async function runEnrichmentBatch( const bias = biasByArticleId.get(article._id); const factStatus = factResult.statusByArticleId.get(article._id); const biasStatus = biasResult.statusByArticleId.get(article._id); + if (factStatus?.status === "deferred") { + const result = await ctx.runMutation( + internal.enrichment.deferArticleFactExtraction, + { + articleId: article._id, + runId, + previousStatus: article.previousStatus, + reason: factStatus.error ?? "Atomic fact extraction deferred", + attemptedAt: factStatus.analyzedAt ?? Date.now(), + }, + ); + if (result.updated) { + failed++; + if (result.eventId) { + touchedEventIds.add(result.eventId); + } + } else { + console.warn( + `[enrichment] Article ${article._id} lease no longer belongs to run ${runId}; leaving it unchanged`, + ); + } + continue; + } + if (biasStatus?.status === "deferred") { + const result = await ctx.runMutation( + internal.enrichment.deferArticleBiasDetection, + { + articleId: article._id, + runId, + previousStatus: article.previousStatus, + reason: biasStatus.error ?? "Article bias detection deferred", + }, + ); + if (result.updated) { + failed++; + if (result.eventId) { + touchedEventIds.add(result.eventId); + } + } else { + console.warn( + `[enrichment] Article ${article._id} lease no longer belongs to run ${runId}; leaving it unchanged`, + ); + } + continue; + } const result = await ctx.runMutation( internal.enrichment.markArticleEnriched, { @@ -865,7 +962,7 @@ async function runEnrichmentBatch( sourceBiasDelta: bias?.sourceBiasDelta, sourceBiasOutlierFlag: bias?.sourceBiasOutlierFlag, biasAnalyzedAt: bias?.biasAnalyzedAt, - biasDetectionStatus: biasStatus?.status, + biasDetectionStatus: toArticleBiasStatus(biasStatus?.status), biasDetectionError: biasStatus?.error, summary: prepared.extractedSummary, atomicFacts: factsByArticleId.get(article._id), @@ -1056,6 +1153,114 @@ export const reenrichArticlesBackfill = internalAction({ }, }); +export const backfillAtomicFacts = internalAction({ + args: { + limit: v.optional(v.number()), + batchSize: v.optional(v.number()), + beforePublishedAt: v.optional(v.number()), + includeFailed: v.optional(v.boolean()), + includeSucceededEmpty: v.optional(v.boolean()), + force: v.optional(v.boolean()), + }, + handler: async ( + ctx, + { + limit, + batchSize, + beforePublishedAt, + includeFailed, + includeSucceededEmpty, + force, + }, + ): Promise<{ + processed: number; + enriched: number; + failed: number; + skipped: boolean; + budgetExhausted: boolean; + nextBeforePublishedAt?: number; + }> => { + const paused = await ctx.runQuery(internal.config.isPipelinePaused, {}); + if (paused && !force) { + console.log( + "[enrichment] Pipeline paused — skipping atomic facts backfill", + ); + return { + processed: 0, + enriched: 0, + failed: 0, + skipped: true, + budgetExhausted: false, + }; + } + + const totalLimit = safeInteger(limit, 100, 1, 500); + const perBatch = safeInteger(batchSize, 20, 1, BATCH_SIZE); + let processed = 0; + let enriched = 0; + let failed = 0; + let cursor = beforePublishedAt; + + try { + while (processed < totalLimit) { + const budget = await ctx.runQuery(internal.aiBudget.checkBudget, {}); + if (!budget.allowed) { + console.warn( + `[enrichment] AI budget exhausted ($${budget.spentUsd}/$${budget.dailyLimitUsd}). Stopping atomic facts backfill.`, + ); + return { + processed, + enriched, + failed, + skipped: false, + budgetExhausted: true, + nextBeforePublishedAt: cursor, + }; + } + + const runId = randomUUID(); + const articles = await ctx.runMutation( + internal.enrichment.claimArticlesNeedingFactExtraction, + { + limit: Math.min(perBatch, totalLimit - processed), + runId, + leaseExpiresAt: Date.now() + ARTICLE_LEASE_TTL_MS, + beforePublishedAt: cursor, + includeFailed: includeFailed ?? true, + includeSucceededEmpty: includeSucceededEmpty ?? false, + }, + ); + + if (articles.length === 0) break; + const oldestPublishedAt = articles.reduce( + (oldest: number | undefined, article) => + oldest === undefined + ? article.publishedAt + : Math.min(oldest, article.publishedAt), + undefined, + ); + + const result = await runEnrichmentBatch(ctx, articles, runId); + processed += articles.length; + enriched += result.enriched; + failed += result.failed; + cursor = oldestPublishedAt; + } + + return { + processed, + enriched, + failed, + skipped: false, + budgetExhausted: false, + nextBeforePublishedAt: cursor, + }; + } finally { + await shutdownPostHog(); + } + }, +}); + export const reenrichEventArticles = internalAction({ args: { eventId: v.id("events"), diff --git a/packages/backend/convex/lib/aiCall.ts b/packages/backend/convex/lib/aiCall.ts index 2fed3bd..4547816 100644 --- a/packages/backend/convex/lib/aiCall.ts +++ b/packages/backend/convex/lib/aiCall.ts @@ -3,7 +3,7 @@ import { internal } from "../_generated/api"; import type { Id } from "../_generated/dataModel"; import type { ActionCtx } from "../_generated/server"; -import { calculateCost } from "../aiBudget"; +import { calculateCost, calculateCostWithCachedInput } from "../aiBudget"; import { getOpenAI } from "./openai"; export type AICallType = @@ -27,6 +27,7 @@ type BudgetReservation = { type AICallUsage = { inputTokens: number; outputTokens: number; + cachedInputTokens: number; costUsd: number; latencyMs: number; }; @@ -44,6 +45,9 @@ type ChatArgs = { maxTokens?: number; temperature?: number; parseJson?: boolean; + promptCacheKey?: string; + promptCacheRetention?: "24h"; + reasoningEffort?: "minimal" | "low" | "medium" | "high"; context: AICallContext; runtime: ActionCtx; maxRetries?: number; @@ -69,6 +73,7 @@ function zeroUsage(): AICallUsage { return { inputTokens: 0, outputTokens: 0, + cachedInputTokens: 0, costUsd: 0, latencyMs: 0, }; @@ -101,6 +106,10 @@ function isFatalError(error: unknown): boolean { return status === 401 || status === 403 || status === 404; } +function isGpt5FamilyModel(model: string): boolean { + return model.startsWith("gpt-5"); +} + function retryDelayMs(attempt: number): number { const baseDelay = 500 * 3 ** attempt; const jitter = 0.75 + Math.random() * 0.5; @@ -161,6 +170,7 @@ async function logUsage( model, inputTokens: usage.inputTokens, outputTokens: usage.outputTokens, + cachedInputTokens: usage.cachedInputTokens, costUsd: usage.costUsd, latencyMs: usage.latencyMs, eventId: context.eventId, @@ -253,6 +263,7 @@ export async function callOpenAI( const usage = { inputTokens, outputTokens: 0, + cachedInputTokens: 0, costUsd: calculateCost(args.model, inputTokens, 0), latencyMs: Date.now() - startedAt, }; @@ -272,20 +283,43 @@ export async function callOpenAI( const response = await openai.chat.completions.create({ model: args.model, - temperature: args.temperature ?? 0.1, - ...(args.maxTokens ? { max_tokens: args.maxTokens } : {}), + ...(!isGpt5FamilyModel(args.model) + ? { temperature: args.temperature ?? 0.1 } + : { reasoning_effort: args.reasoningEffort ?? "minimal" }), + ...(args.maxTokens + ? isGpt5FamilyModel(args.model) + ? { max_completion_tokens: args.maxTokens } + : { max_tokens: args.maxTokens } + : {}), ...(args.responseFormat ? { response_format: args.responseFormat as never } : {}), + prompt_cache_key: + args.promptCacheKey ?? `biviant:${args.context.callType}`, + ...(args.promptCacheRetention + ? { prompt_cache_retention: args.promptCacheRetention } + : {}), messages: args.messages, - }); + } as never); const inputTokens = response.usage?.prompt_tokens ?? 0; const outputTokens = response.usage?.completion_tokens ?? 0; + const cachedInputTokens = + ( + response.usage as + | { prompt_tokens_details?: { cached_tokens?: number } } + | undefined + )?.prompt_tokens_details?.cached_tokens ?? 0; const usage = { inputTokens, outputTokens, - costUsd: calculateCost(args.model, inputTokens, outputTokens), + cachedInputTokens, + costUsd: calculateCostWithCachedInput( + args.model, + inputTokens, + cachedInputTokens, + outputTokens, + ), latencyMs: Date.now() - startedAt, }; const logged = await logUsage( diff --git a/packages/backend/convex/lib/articleExtraction.ts b/packages/backend/convex/lib/articleExtraction.ts index 48ce44c..a8ad475 100644 --- a/packages/backend/convex/lib/articleExtraction.ts +++ b/packages/backend/convex/lib/articleExtraction.ts @@ -99,6 +99,15 @@ const ENTITY_NOISE_TERMS = new Set([ "typically", "watch", ]); +const WEEKDAY_ENTITY_NOISE_TERMS = new Set([ + "monday", + "tuesday", + "wednesday", + "thursday", + "friday", + "saturday", + "sunday", +]); const ENTITY_ROLE_PREFIXES = [ "former", "president", @@ -204,18 +213,35 @@ function isUsefulEntityCandidate( count: number, titleEntities: Set, allUppercaseEntities: Set, + numericEntities: Set, ): boolean { if (entity.length < 3 || entity.length > 80) return false; if (!/[a-z0-9]/.test(entity)) return false; if (/^(?:[a-z]\s*)+$/.test(entity)) return false; const words = entity.split(/\s+/).filter(Boolean); - if (words.length === 0 || words.length > 5) return false; - if (words.some((word) => ENTITY_NOISE_TERMS.has(word))) return false; + if (words.length === 0 || words.length > 7) return false; + if ( + words.length === 1 && + (ENTITY_NOISE_TERMS.has(entity) || WEEKDAY_ENTITY_NOISE_TERMS.has(entity)) + ) { + return false; + } + if ( + words.length > 1 && + words.some( + (word) => + ENTITY_NOISE_TERMS.has(word) && + !WEEKDAY_ENTITY_NOISE_TERMS.has(word), + ) + ) { + return false; + } if (ENTITY_NOISE_TERMS.has(entity)) return false; if (words.length === 1) { return ( + numericEntities.has(entity) || titleEntities.has(entity) || count > 1 || (/^[a-z]{2,6}$/.test(entity.toLowerCase()) && @@ -226,12 +252,18 @@ function isUsefulEntityCandidate( return true; } -function addNumericEntities(text: string, scores: Map, weight: number) { +function addNumericEntities( + text: string, + scores: Map, + numericEntities: Set, + weight: number, +) { const numericMatches = text.match(NUMERIC_ENTITY_PATTERN) ?? []; for (const match of numericMatches) { const normalized = normalizeEntityCandidate(match).value; if (normalized.length >= 2) { scores.set(normalized, (scores.get(normalized) ?? 0) + weight); + numericEntities.add(normalized); } } } @@ -262,6 +294,8 @@ function collectProperNounCandidates( if (normalized.value) { scores.set(normalized.value, (scores.get(normalized.value) ?? 0) + weight); if (normalized.wasAllUppercase) allUppercaseEntities.add(normalized.value); + // Only proper nouns collected from the title get title-entity priority. + // Body matches use lower weights and must earn their way in by repetition. if (weight >= 3) titleEntities.add(normalized.value); } phrase = []; @@ -281,6 +315,7 @@ function extractEntityCandidates(title: string, ...texts: string[]): string[] { const scores = new Map(); const titleEntities = new Set(); const allUppercaseEntities = new Set(); + const numericEntities = new Set(); const cleanedTitle = stripTags(title); collectProperNounCandidates( @@ -290,7 +325,7 @@ function extractEntityCandidates(title: string, ...texts: string[]): string[] { allUppercaseEntities, 3, ); - addNumericEntities(cleanedTitle, scores, 3); + addNumericEntities(cleanedTitle, scores, numericEntities, 3); for (const rawText of texts) { const text = stripTags(rawText); @@ -302,12 +337,18 @@ function extractEntityCandidates(title: string, ...texts: string[]): string[] { allUppercaseEntities, 1, ); - addNumericEntities(text, scores, 1); + addNumericEntities(text, scores, numericEntities, 1); } return Array.from(scores.entries()) .filter(([entity, score]) => - isUsefulEntityCandidate(entity, score, titleEntities, allUppercaseEntities), + isUsefulEntityCandidate( + entity, + score, + titleEntities, + allUppercaseEntities, + numericEntities, + ), ) .sort((a, b) => b[1] - a[1] || a[0].localeCompare(b[0])) .map(([entity]) => entity) @@ -553,6 +594,19 @@ function absolutizeUrl(candidate: string | undefined, baseUrl: string): string | } } +function isAvatarImagePath(pathname: string): boolean { + const segments = pathname + .toLowerCase() + .split("/") + .map((segment) => segment.trim()) + .filter(Boolean); + return segments.some((segment, index) => { + if (segment === "avatar") return true; + if (index !== segments.length - 1) return false; + return /^avatar(?:[-_.]|$)/.test(segment); + }); +} + function isLikelyValidImageUrl(url: string | undefined): boolean { if (!url) return false; try { @@ -560,7 +614,7 @@ function isLikelyValidImageUrl(url: string | undefined): boolean { if (!/^https?:$/.test(parsed.protocol)) return false; const path = parsed.pathname.toLowerCase(); if ( - path.includes("avatar") || + isAvatarImagePath(path) || path.includes("author") || path.includes("logo") || path.includes("icon") || @@ -587,6 +641,12 @@ function normalizeEscapedUrl(candidate: string | undefined): string | undefined function scoreRawImageCandidate(url: string, hostname: string): number { let score = 0; const lower = url.toLowerCase(); + let pathname = ""; + try { + pathname = new URL(url).pathname; + } catch { + pathname = url; + } if (/\.(avif|jpe?g|png|webp)(?:$|\?)/i.test(lower)) score += 4; if (lower.includes("reutersmedia.net")) score += 6; @@ -607,7 +667,7 @@ function scoreRawImageCandidate(url: string, hostname: string): number { lower.includes("logo") || lower.includes("icon") || lower.includes("sprite") || - lower.includes("avatar") || + isAvatarImagePath(pathname) || lower.includes("author") || lower.includes("thumbnail") ) { diff --git a/packages/backend/convex/pipelineDiagnostics.ts b/packages/backend/convex/pipelineDiagnostics.ts new file mode 100644 index 0000000..02fb42e --- /dev/null +++ b/packages/backend/convex/pipelineDiagnostics.ts @@ -0,0 +1,345 @@ +import { v } from "convex/values"; +import { internalQuery } from "./_generated/server"; +import type { Doc, Id } from "./_generated/dataModel"; + +function hasText(value: string | undefined, minLength: number): boolean { + return (value ?? "").trim().length >= minLength; +} + +function hasAtomicFacts(article: Doc<"articles">): boolean { + return (article.atomicFacts ?? []).some((fact) => fact.trim().length > 0); +} + +function hasPerspectiveSummary(event: Doc<"events">): boolean { + return Boolean( + event.lastSummarizedAt && + event.perspectiveSummaries?.center?.trim() && + event.perspectiveSummaries?.left?.trim() && + event.perspectiveSummaries?.right?.trim() && + event.globalImpact?.trim(), + ); +} + +function summarizeDrop(from: number, to: number) { + return { + count: from - to, + pctOfPrevious: from === 0 ? 0 : Math.round(((from - to) / from) * 1000) / 10, + }; +} + +export const eventAiFunnel = internalQuery({ + args: { + days: v.optional(v.number()), + limit: v.optional(v.number()), + }, + handler: async (ctx, args) => { + const days = Math.max(1, Math.min(60, Math.floor(args.days ?? 14))); + const limit = Math.max(1, Math.min(2000, Math.floor(args.limit ?? 1000))); + const cutoff = Date.now() - days * 24 * 60 * 60 * 1000; + const events = ( + await ctx.db + .query("events") + .withIndex("by_status_recency", (q) => q.eq("status", "published")) + .order("desc") + .take(limit) + ).filter((event) => event.firstPublishedAt > cutoff); + + const stageEventIds = { + published: new Set>(), + article3: new Set>(), + source2: new Set>(), + factualArticle3: new Set>(), + factualSource2: new Set>(), + summarized: new Set>(), + claimAnalyzed: new Set>(), + hasClaims: new Set>(), + }; + + const samples = { + missingArticleCoverage: [] as Array<{ + eventId: Id<"events">; + title: string; + articleCount: number; + sourceCount: number; + }>, + missingFactualCoverage: [] as Array<{ + eventId: Id<"events">; + title: string; + articleCount: number; + sourceCount: number; + factualArticleCount: number; + factualSourceCount: number; + }>, + qualifiedNoSummary: [] as Array<{ + eventId: Id<"events">; + title: string; + factualArticleCount: number; + factualSourceCount: number; + latestSummaryJobStatus?: string; + latestSummaryJobReason?: string; + latestSummaryJobError?: string; + }>, + summarizedNoClaims: [] as Array<{ + eventId: Id<"events">; + title: string; + lastClaimAnalysisAt?: number; + factualArticleCount: number; + factualSourceCount: number; + }>, + }; + + for (const event of events) { + stageEventIds.published.add(event._id); + + const articles = await ctx.db + .query("articles") + .withIndex("by_event", (q) => q.eq("eventId", event._id)) + .collect(); + const sourceCount = new Set(articles.map((article) => article.sourceId)) + .size; + const factualArticles = articles.filter(hasAtomicFacts); + const factualSourceCount = new Set( + factualArticles.map((article) => article.sourceId), + ).size; + + if (articles.length >= 3) stageEventIds.article3.add(event._id); + if (articles.length >= 3 && sourceCount >= 2) { + stageEventIds.source2.add(event._id); + } else if (samples.missingArticleCoverage.length < 10) { + samples.missingArticleCoverage.push({ + eventId: event._id, + title: event.title, + articleCount: articles.length, + sourceCount, + }); + } + + if (articles.length >= 3 && sourceCount >= 2 && factualArticles.length >= 3) { + stageEventIds.factualArticle3.add(event._id); + } + if ( + articles.length >= 3 && + sourceCount >= 2 && + factualArticles.length >= 3 && + factualSourceCount >= 2 + ) { + stageEventIds.factualSource2.add(event._id); + } else if ( + articles.length >= 3 && + sourceCount >= 2 && + samples.missingFactualCoverage.length < 10 + ) { + samples.missingFactualCoverage.push({ + eventId: event._id, + title: event.title, + articleCount: articles.length, + sourceCount, + factualArticleCount: factualArticles.length, + factualSourceCount, + }); + } + + if (hasPerspectiveSummary(event)) { + stageEventIds.summarized.add(event._id); + } + if (event.lastClaimAnalysisAt) { + stageEventIds.claimAnalyzed.add(event._id); + } + + const hasClaimRows = Boolean( + await ctx.db + .query("eventClaims") + .withIndex("by_event", (q) => q.eq("eventId", event._id)) + .first(), + ); + if (hasClaimRows) stageEventIds.hasClaims.add(event._id); + + if ( + factualArticles.length >= 3 && + factualSourceCount >= 2 && + !hasPerspectiveSummary(event) && + samples.qualifiedNoSummary.length < 10 + ) { + const latestJob = await ctx.db + .query("eventSummaryJobs") + .withIndex("by_event_updatedAt", (q) => q.eq("eventId", event._id)) + .order("desc") + .first(); + samples.qualifiedNoSummary.push({ + eventId: event._id, + title: event.title, + factualArticleCount: factualArticles.length, + factualSourceCount, + latestSummaryJobStatus: latestJob?.status, + latestSummaryJobReason: latestJob?.reason, + latestSummaryJobError: latestJob?.lastError, + }); + } + + if ( + hasPerspectiveSummary(event) && + !hasClaimRows && + samples.summarizedNoClaims.length < 10 + ) { + samples.summarizedNoClaims.push({ + eventId: event._id, + title: event.title, + lastClaimAnalysisAt: event.lastClaimAnalysisAt, + factualArticleCount: factualArticles.length, + factualSourceCount, + }); + } + } + + const stages = { + published: stageEventIds.published.size, + articleCountAtLeast3: stageEventIds.article3.size, + sourceCountAtLeast2: stageEventIds.source2.size, + factualArticleCountAtLeast3: stageEventIds.factualArticle3.size, + factualSourceCountAtLeast2: stageEventIds.factualSource2.size, + summarized: stageEventIds.summarized.size, + claimAnalyzed: stageEventIds.claimAnalyzed.size, + hasClaimRows: stageEventIds.hasClaims.size, + }; + + return { + windowDays: days, + scannedEvents: events.length, + cutoff, + stages, + drops: { + publishedToArticle3: summarizeDrop( + stages.published, + stages.articleCountAtLeast3, + ), + article3ToSource2: summarizeDrop( + stages.articleCountAtLeast3, + stages.sourceCountAtLeast2, + ), + source2ToFactualArticle3: summarizeDrop( + stages.sourceCountAtLeast2, + stages.factualArticleCountAtLeast3, + ), + factualArticle3ToFactualSource2: summarizeDrop( + stages.factualArticleCountAtLeast3, + stages.factualSourceCountAtLeast2, + ), + factualSource2ToSummarized: summarizeDrop( + stages.factualSourceCountAtLeast2, + stages.summarized, + ), + summarizedToClaimAnalyzed: summarizeDrop( + stages.summarized, + stages.claimAnalyzed, + ), + summarizedToClaimRows: summarizeDrop( + stages.summarized, + stages.hasClaimRows, + ), + }, + samples, + }; + }, +}); + +export const articleFactExtractionFunnel = internalQuery({ + args: { + days: v.optional(v.number()), + limit: v.optional(v.number()), + }, + handler: async (ctx, args) => { + const days = Math.max(1, Math.min(60, Math.floor(args.days ?? 14))); + const limit = Math.max(1, Math.min(5000, Math.floor(args.limit ?? 2000))); + const cutoff = Date.now() - days * 24 * 60 * 60 * 1000; + const articles = ( + await ctx.db.query("articles").withIndex("by_published").order("desc").take(limit) + ).filter((article) => article.publishedAt > cutoff); + + const byStatus = new Map(); + const byFactStatus = new Map(); + const byExtractionQuality = new Map(); + for (const article of articles) { + byStatus.set(article.status, (byStatus.get(article.status) ?? 0) + 1); + byFactStatus.set( + article.factExtractionStatus ?? "unset", + (byFactStatus.get(article.factExtractionStatus ?? "unset") ?? 0) + 1, + ); + byExtractionQuality.set( + article.extractionQuality ?? "unset", + (byExtractionQuality.get(article.extractionQuality ?? "unset") ?? 0) + + 1, + ); + } + + return { + windowDays: days, + scannedArticles: articles.length, + cutoff, + total: articles.length, + withSummary: articles.filter((article) => hasText(article.summary, 80)) + .length, + withUsefulSummary: articles.filter((article) => + hasText(article.summary, 200), + ).length, + withUsefulRssSnippet: articles.filter((article) => + hasText(article.rssSnippet, 200), + ).length, + withEntities: articles.filter((article) => (article.entities ?? []).length > 0) + .length, + withAtomicFacts: articles.filter(hasAtomicFacts).length, + withBiasScore: articles.filter( + (article) => typeof article.aiBiasScore === "number", + ).length, + factExtractionSucceeded: articles.filter( + (article) => article.factExtractionStatus === "succeeded", + ).length, + factExtractionSucceededEmpty: articles.filter( + (article) => article.factExtractionStatus === "succeeded_empty", + ).length, + factExtractionDeferred: articles.filter( + (article) => article.factExtractionStatus === "deferred", + ).length, + factExtractionFailed: articles.filter( + (article) => article.factExtractionStatus === "failed", + ).length, + factExtractionSkipped: articles.filter( + (article) => article.factExtractionStatus === "skipped", + ).length, + neverFactAttempted: articles.filter( + (article) => article.factExtractionStatus === undefined, + ).length, + needsFactExtraction: articles.filter( + (article) => + article.status !== "processing" && + article.status !== "discarded" && + !hasAtomicFacts(article) && + article.factExtractionStatus !== "skipped" && + article.factExtractionStatus !== "succeeded_empty", + ).length, + byStatus: Object.fromEntries(byStatus.entries()), + byFactStatus: Object.fromEntries(byFactStatus.entries()), + byExtractionQuality: Object.fromEntries(byExtractionQuality.entries()), + samplesNeedingFacts: articles + .filter( + (article) => + article.status !== "processing" && + article.status !== "discarded" && + !hasAtomicFacts(article) && + article.factExtractionStatus !== "skipped" && + article.factExtractionStatus !== "succeeded_empty", + ) + .slice(0, 10) + .map((article) => ({ + articleId: article._id, + eventId: article.eventId, + title: article.title, + status: article.status, + factExtractionStatus: article.factExtractionStatus, + factExtractionError: article.factExtractionError, + summaryLength: (article.summary ?? "").length, + rssSnippetLength: (article.rssSnippet ?? "").length, + extractionQuality: article.extractionQuality, + })), + }; + }, +}); diff --git a/packages/backend/convex/schema.ts b/packages/backend/convex/schema.ts index 78ea0ad..e850d5f 100644 --- a/packages/backend/convex/schema.ts +++ b/packages/backend/convex/schema.ts @@ -251,13 +251,18 @@ export default defineSchema({ atomicFacts: v.optional(v.array(v.string())), // ["Vote count: 60-40", "Passed on: Tuesday", "Opposition: GOP"] factExtractionStatus: v.optional( v.union( + v.literal("pending"), + v.literal("deferred"), v.literal("succeeded"), + v.literal("succeeded_empty"), v.literal("failed"), v.literal("skipped"), ), ), factExtractionError: v.optional(v.string()), factExtractedAt: v.optional(v.number()), + factExtractionAttempts: v.optional(v.number()), + factExtractionLastAttemptAt: v.optional(v.number()), // Populated by enrichment pipeline (AI bias detection) aiBiasScore: v.optional(v.number()), @@ -276,12 +281,15 @@ export default defineSchema({ biasAnalyzedAt: v.optional(v.number()), biasDetectionStatus: v.optional( v.union( + v.literal("deferred"), v.literal("succeeded"), v.literal("failed"), v.literal("skipped"), ), ), biasDetectionError: v.optional(v.string()), + biasDetectionAttempts: v.optional(v.number()), + biasDetectionLastAttemptAt: v.optional(v.number()), status: v.union( v.literal("unprocessed"), @@ -503,6 +511,7 @@ export default defineSchema({ articleId: v.optional(v.id("articles")), inputTokens: v.number(), outputTokens: v.number(), + cachedInputTokens: v.optional(v.number()), costUsd: v.number(), // Pre-calculated latencyMs: v.optional(v.number()), timestamp: v.number(), diff --git a/packages/backend/convex/summarizationNode.ts b/packages/backend/convex/summarizationNode.ts index 9dcfa33..f2832de 100644 --- a/packages/backend/convex/summarizationNode.ts +++ b/packages/backend/convex/summarizationNode.ts @@ -10,7 +10,7 @@ import { shutdownPostHog } from "./lib/openai"; import { callOpenAI } from "./lib/aiCall"; import { buildEventSummaryPrompt, type EventSummaryOutput } from "./prompts"; -const DEFAULT_MODEL = "gpt-4o-mini"; +const DEFAULT_MODEL = "gpt-5-nano"; const DEFAULT_ENQUEUE_LIMIT = 40; const DEFAULT_BATCH_SIZE = 4; const DEFAULT_MAX_ATTEMPTS = 3;