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
20 changes: 16 additions & 4 deletions src/api/routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ import {
listPullRequestFiles,
listPullRequestReviews,
listRecentMergedPullRequests,
listLatestSignalSnapshotsByTarget,
listRepoLabels,
listRepoSyncSegments,
listRepoSyncStates,
Expand Down Expand Up @@ -97,7 +98,7 @@ import {
buildQueueHealth,
buildRegistryChangeReport,
} from "../signals/engine";
import { attachDataQuality, buildCoreSignalFidelity, buildRepoDataQuality, buildSignalFidelity } from "../signals/data-quality";
import { attachDataQuality, buildCoreSignalFidelity, buildFreshnessSloReport, buildRepoDataQuality, buildSignalFidelity } from "../signals/data-quality";
import { buildPullRequestReviewability } from "../signals/reward-risk";
import { buildLocalBranchAnalysis } from "../signals/local-branch";
import { buildRepoSettingsPreview } from "../signals/settings-preview";
Expand Down Expand Up @@ -397,20 +398,25 @@ export function createApp() {
});

app.get("/v1/sync/status", async (c) => {
const [snapshot, repositories, segments, totals, detailStates, installations, rateLimits] = await Promise.all([
const [snapshot, scoringSnapshot, repositories, segments, totals, detailStates, installations, rateLimits, signalSnapshots, bounties] = await Promise.all([
getLatestRegistrySnapshot(c.env),
getLatestScoringModelSnapshot(c.env),
listRepoSyncStates(c.env),
listRepoSyncSegments(c.env),
listLatestRepoGithubTotalsSnapshots(c.env),
listAllPullRequestDetailSyncStates(c.env),
listInstallationHealth(c.env),
listLatestGitHubRateLimitObservations(c.env, 20),
listLatestSignalSnapshotsByTarget(c.env),
listBounties(c.env),
]);
const repoCount = snapshot?.repoCount ?? repositories.length;
const coreSignalFidelity = buildCoreSignalFidelity(repoCount, repositories, segments, totals, detailStates);
const freshnessSlo = buildFreshnessSloReport({ registrySnapshot: snapshot, scoringSnapshot, repoCount, syncStates: repositories, totals, segments, signalSnapshots, bounties });
return c.json({
generatedAt: nowIso(),
signalFidelity: buildSignalFidelity(repoCount, repositories, segments),
freshnessSlo,
coreSignalFidelity,
historyCoverage: coreSignalFidelity.historyCoverage,
refreshingRepos: coreSignalFidelity.refreshingRepos,
Expand All @@ -425,7 +431,7 @@ export function createApp() {
});

app.get("/v1/readiness", async (c) => {
const [snapshot, scoringSnapshot, syncStates, syncSegments, totals, detailStates, installations, installationHealth, rateLimits] = await Promise.all([
const [snapshot, scoringSnapshot, syncStates, syncSegments, totals, detailStates, installations, installationHealth, rateLimits, signalSnapshots, bounties] = await Promise.all([
getLatestRegistrySnapshot(c.env),
getLatestScoringModelSnapshot(c.env),
listRepoSyncStates(c.env),
Expand All @@ -435,10 +441,13 @@ export function createApp() {
listInstallations(c.env),
listInstallationHealth(c.env),
listLatestGitHubRateLimitObservations(c.env, 20),
listLatestSignalSnapshotsByTarget(c.env),
listBounties(c.env),
]);
const repoCount = snapshot?.repoCount ?? syncStates.length;
const signalFidelity = buildSignalFidelity(repoCount, syncStates, syncSegments);
const coreSignalFidelity = buildCoreSignalFidelity(repoCount, syncStates, syncSegments, totals, detailStates);
const freshnessSlo = buildFreshnessSloReport({ registrySnapshot: snapshot, scoringSnapshot, repoCount, syncStates, totals, segments: syncSegments, signalSnapshots, bounties });
const statusCounts = syncStates.reduce<Record<string, number>>((counts, state) => {
counts[state.status] = (counts[state.status] ?? 0) + 1;
return counts;
Expand All @@ -459,6 +468,7 @@ export function createApp() {
...(signalFidelity.cappedRepos.length > 0 ? [`${signalFidelity.cappedRepos.length} repo sync(s) hit local pagination caps; signal fidelity is degraded.`] : []),
...(signalFidelity.rateLimitedRepos.length > 0 ? [`${signalFidelity.rateLimitedRepos.length} repo sync(s) encountered GitHub rate limiting.`] : []),
...(signalFidelity.staleRepos.length > 0 ? [`${signalFidelity.staleRepos.length} repo sync(s) are stale.`] : []),
...(freshnessSlo.status !== "fresh" ? [`Freshness SLO is ${freshnessSlo.status}; ${freshnessSlo.warnings.length} stale, missing, or blocked signal source(s) need repair.`] : []),
...(installationHealth.some((health) => health.status !== "healthy") ? ["One or more GitHub App installations need attention."] : []),
];
const ready = Boolean(snapshot) && Boolean(c.env.INTERNAL_JOB_TOKEN) && Boolean(c.env.GITTENSORY_API_TOKEN);
Expand All @@ -469,14 +479,16 @@ export function createApp() {
Boolean(c.env.GITHUB_PUBLIC_TOKEN) &&
missingSyncCount === 0 &&
failingSyncs.length === 0 &&
coreSignalFidelity.status === "complete"
coreSignalFidelity.status === "complete" &&
freshnessSlo.launchBlockingCount === 0
: false;
return c.json({
status: ready ? "ready" : "needs_attention",
generatedAt: nowIso(),
ready,
readyForPublicReview,
signalFidelity,
freshnessSlo,
coreSignalFidelity,
historyCoverage: coreSignalFidelity.historyCoverage,
partialRepos: signalFidelity.partialRepos,
Expand Down
32 changes: 32 additions & 0 deletions src/db/repositories.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1414,6 +1414,38 @@ export async function listSignalSnapshots(env: Env, signalType: string, targetKe
return rows.map(toSignalSnapshotRecord);
}

export async function listLatestSignalSnapshotsByTarget(env: Env): Promise<SignalSnapshotRecord[]> {
const { results } = await env.DB.prepare(
`
SELECT id, signal_type, target_key, repo_full_name, payload_json, generated_at
FROM (
SELECT
id,
signal_type,
target_key,
repo_full_name,
payload_json,
generated_at,
row_number() OVER (
PARTITION BY signal_type, target_key
ORDER BY generated_at DESC, id DESC
) AS snapshot_rank
FROM signal_snapshots
)
WHERE snapshot_rank = 1
ORDER BY signal_type, target_key
`,
).all<{ id: string; signal_type: string; target_key: string; repo_full_name: string | null; payload_json: string; generated_at: string }>();
return results.map((row) => ({
id: row.id,
signalType: row.signal_type,
targetKey: row.target_key,
repoFullName: row.repo_full_name,
payload: parseJson<Record<string, never>>(row.payload_json, {}),
generatedAt: row.generated_at,
}));
}

export async function createAgentRun(env: Env, run: AgentRunRecord): Promise<void> {
const db = getDb(env.DB);
await db.insert(agentRuns).values({
Expand Down
24 changes: 24 additions & 0 deletions src/openapi/schemas.ts
Original file line number Diff line number Diff line change
Expand Up @@ -665,6 +665,18 @@ export const SyncStatusSchema = z
.object({
generatedAt: z.string(),
signalFidelity: SignalFidelitySchema,
freshnessSlo: z.object({
status: z.enum(["fresh", "degraded", "blocked"]),
generatedAt: z.string(),
staleCount: z.number(),
degradedCount: z.number(),
blockedCount: z.number(),
missingCount: z.number(),
launchBlockingCount: z.number(),
repairRecommended: z.boolean(),
items: z.array(z.object({ area: z.string(), targetKey: z.string(), status: z.string(), launchBlocking: z.boolean(), ageSeconds: z.number().optional(), sloSeconds: z.number(), breachSeconds: z.number().optional(), observedAt: z.string().nullable().optional(), summary: z.string() })),
warnings: z.array(z.string()),
}),
coreSignalFidelity: CoreSignalFidelitySchema,
historyCoverage: z.enum(["sampled", "counts_only", "full"]),
refreshingRepos: z.array(z.string()),
Expand All @@ -685,6 +697,18 @@ export const ReadinessSchema = z
ready: z.boolean(),
readyForPublicReview: z.boolean(),
signalFidelity: SignalFidelitySchema,
freshnessSlo: z.object({
status: z.enum(["fresh", "degraded", "blocked"]),
generatedAt: z.string(),
staleCount: z.number(),
degradedCount: z.number(),
blockedCount: z.number(),
missingCount: z.number(),
launchBlockingCount: z.number(),
repairRecommended: z.boolean(),
items: z.array(z.object({ area: z.string(), targetKey: z.string(), status: z.string(), launchBlocking: z.boolean(), ageSeconds: z.number().optional(), sloSeconds: z.number(), breachSeconds: z.number().optional(), observedAt: z.string().nullable().optional(), summary: z.string() })),
warnings: z.array(z.string()),
}),
coreSignalFidelity: CoreSignalFidelitySchema,
historyCoverage: z.enum(["sampled", "counts_only", "full"]),
partialRepos: z.array(z.string()),
Expand Down
18 changes: 15 additions & 3 deletions src/queue/processors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
listContributorRepoStats,
listIssues,
listIssueSignalSample,
listLatestSignalSnapshotsByTarget,
listOtherOpenPullRequests,
listOpenPullRequests,
listPullRequests,
Expand Down Expand Up @@ -57,6 +58,10 @@ import { buildIssueAdvisory, buildPullRequestAdvisory } from "../rules/advisory"
import { getOrCreateScoringModelSnapshot, refreshScoringModelSnapshot } from "../scoring/model";
import { buildAndPersistContributorDecisionPack } from "../services/decision-pack";
import { executeAgentRun, explainBlockersWithAgent, planNextWork } from "../services/agent-orchestrator";
import {
buildFreshnessSloReport,
freshnessAuditMetadata,
} from "../signals/data-quality";
import {
buildBurdenForecast,
buildCollisionEdges,
Expand Down Expand Up @@ -202,7 +207,7 @@ async function fanOutRepoSignalSnapshotJobs(env: Env, requestedBy: "schedule" |
}

async function repairDataFidelity(env: Env, requestedBy: "schedule" | "api" | "test"): Promise<void> {
const [repositories, segments] = await Promise.all([listRepositories(env), listRepoSyncSegments(env)]);
const [repositories, segments, signalSnapshots] = await Promise.all([listRepositories(env), listRepoSyncSegments(env), listLatestSignalSnapshotsByTarget(env)]);
const requiredSegments = new Set(["labels", "open_issues", "open_pull_requests"]);
const segmentsByRepo = new Map<string, Set<string>>();
for (const segment of segments) {
Expand All @@ -213,6 +218,7 @@ async function repairDataFidelity(env: Env, requestedBy: "schedule" | "api" | "t
}
}
const registeredRepos = repositories.filter((repo) => repo.isRegistered);
const freshnessSlo = buildFreshnessSloReport({ repoCount: registeredRepos.length, segments, signalSnapshots });
const repairs = [];
const signalRefreshes = [];
for (const repo of registeredRepos) {
Expand Down Expand Up @@ -247,8 +253,14 @@ async function repairDataFidelity(env: Env, requestedBy: "schedule" | "api" | "t
]);
await recordAuditEvent(env, {
eventType: "sync.fidelity_repair",
outcome: repairs.length > 0 ? "queued" : "completed",
metadata: { requestedBy, repairCount: repairs.length, signalRefreshCount: signalRefreshes.length, repairs: repairs.slice(0, 25) },
outcome: repairs.length > 0 || freshnessSlo.repairRecommended ? "queued" : "completed",
metadata: { requestedBy, repairCount: repairs.length, signalRefreshCount: signalRefreshes.length, repairs: repairs.slice(0, 25), freshnessSlo: freshnessAuditMetadata(freshnessSlo) },
});
await recordAuditEvent(env, {
eventType: "signals.freshness_slo",
outcome: freshnessSlo.repairRecommended ? "queued" : "completed",
detail: freshnessSlo.status,
metadata: { requestedBy, ...freshnessAuditMetadata(freshnessSlo) },
});
}

Expand Down
103 changes: 102 additions & 1 deletion src/signals/data-quality.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,17 @@
import type { DataQuality, PullRequestDetailSyncStateRecord, RepoGithubTotalsSnapshotRecord, RepoSyncSegmentRecord, RepoSyncStateRecord } from "../types";
import type { BountyRecord, DataQuality, PullRequestDetailSyncStateRecord, RegistrySnapshot, RepoGithubTotalsSnapshotRecord, RepoSyncSegmentRecord, RepoSyncStateRecord, ScoringModelSnapshotRecord, SignalSnapshotRecord } from "../types";
import { nowIso } from "../utils/json";

const DEFAULT_STALE_MS = 7 * 24 * 60 * 60 * 1000;
const FRESHNESS_SLO_MS = {
registry: DEFAULT_STALE_MS,
scoring_model: DEFAULT_STALE_MS,
github_totals: DEFAULT_STALE_MS,
repo_segments: DEFAULT_STALE_MS,
decision_pack: 6 * 60 * 60 * 1000,
bounty_data: 24 * 60 * 60 * 1000,
signal_snapshot: 12 * 60 * 60 * 1000,
};
const LAUNCH_BLOCKING_FRESHNESS_AREAS = new Set<keyof typeof FRESHNESS_SLO_MS>(["registry", "scoring_model", "github_totals", "repo_segments"]);
const COMPLETE_SEGMENT_STATUSES = new Set<RepoSyncSegmentRecord["status"]>(["complete", "not_modified", "sampled"]);
const BLOCKING_SEGMENT_STATUSES = new Set<RepoSyncSegmentRecord["status"]>(["error", "rate_limited", "waiting_rate_limit", "skipped"]);
const REQUIRED_OPEN_SEGMENTS = new Set<RepoSyncSegmentRecord["segment"]>(["metadata", "labels", "open_issues", "open_pull_requests", "pull_request_files", "pull_request_reviews", "check_summaries"]);
Expand Down Expand Up @@ -31,6 +41,80 @@ export type CoreSignalFidelity = {
historyCoverage: "sampled" | "counts_only" | "full";
};

export type FreshnessSloReport = {
status: "fresh" | "degraded" | "blocked";
generatedAt: string;
staleCount: number;
degradedCount: number;
blockedCount: number;
missingCount: number;
launchBlockingCount: number;
repairRecommended: boolean;
items: Array<{ area: keyof typeof FRESHNESS_SLO_MS; targetKey: string; status: "fresh" | "stale" | "degraded" | "blocked" | "missing"; launchBlocking: boolean; ageSeconds?: number; sloSeconds: number; breachSeconds?: number; observedAt?: string | null; summary: string }>;
warnings: string[];
};

export function buildFreshnessSloReport(args: {
registrySnapshot?: RegistrySnapshot | null;
scoringSnapshot?: ScoringModelSnapshotRecord | null;
repoCount?: number;
syncStates?: RepoSyncStateRecord[];
totals?: RepoGithubTotalsSnapshotRecord[];
segments?: RepoSyncSegmentRecord[];
signalSnapshots?: SignalSnapshotRecord[];
bounties?: BountyRecord[];
expectedDecisionPackKeys?: string[];
nowMs?: number;
}): FreshnessSloReport {
const nowMs = args.nowMs ?? Date.now();
const items: FreshnessSloReport["items"] = [];
const add = (area: keyof typeof FRESHNESS_SLO_MS, targetKey: string, observedAt: string | null | undefined, forced?: "blocked" | "degraded" | "missing") => {
const observedMs = observedAt ? Date.parse(observedAt) : NaN;
const validObservedAt = observedAt && Number.isFinite(observedMs) ? observedAt : null;
const ageSeconds = validObservedAt ? Math.max(0, Math.floor((nowMs - observedMs) / 1000)) : undefined;
const status = forced ?? (!validObservedAt ? "missing" : ageSeconds !== undefined && ageSeconds * 1000 > FRESHNESS_SLO_MS[area] ? "stale" : "fresh");
const launchBlocking = status !== "fresh" && LAUNCH_BLOCKING_FRESHNESS_AREAS.has(area);
items.push({ area, targetKey, status, launchBlocking, ...(ageSeconds !== undefined ? { ageSeconds, breachSeconds: Math.max(0, ageSeconds - Math.floor(FRESHNESS_SLO_MS[area] / 1000)) } : {}), sloSeconds: Math.floor(FRESHNESS_SLO_MS[area] / 1000), observedAt: validObservedAt, summary: `${area}:${targetKey} is ${status}` });
};
if ("registrySnapshot" in args) add("registry", "latest", args.registrySnapshot?.fetchedAt);
if ("scoringSnapshot" in args) add("scoring_model", "latest", args.scoringSnapshot?.fetchedAt);
if ("totals" in args && (args.repoCount ?? 0) > 0) add("github_totals", "registered_repos", oldest(args.totals?.map((total) => total.fetchedAt)), args.totals?.length ? undefined : "missing");
if ((args.repoCount ?? 0) > 0) {
const segmentBlocked = args.segments?.some((segment) => BLOCKING_SEGMENT_STATUSES.has(segment.status) && !hasEffectiveSegmentCoverage(segment)) || args.syncStates?.some((state) => ["error", "skipped", "rate_limited"].includes(state.status));
const segmentDegraded = args.syncStates?.some((state) => !["success", "never_synced"].includes(state.status));
add("repo_segments", "registered_repos", oldest(args.segments?.map((segment) => segment.completedAt ?? segment.updatedAt)), segmentBlocked ? "blocked" : segmentDegraded ? "degraded" : args.segments?.length ? undefined : "missing");
}
for (const [key, snapshots] of groupBy(args.signalSnapshots ?? [], (snapshot) => `${snapshot.signalType}\0${snapshot.targetKey}`)) {
const type = snapshots[0]?.signalType ?? key;
const targetKey = snapshots[0]?.targetKey ?? type;
add(type === "contributor-decision-pack" ? "decision_pack" : "signal_snapshot", targetKey ?? type, newest(snapshots.map((snapshot) => snapshot.generatedAt)));
}
for (const key of args.expectedDecisionPackKeys ?? []) {
if (!items.some((item) => item.area === "decision_pack" && item.targetKey === key)) add("decision_pack", key, null, "missing");
}
if (args.bounties?.length) add("bounty_data", "all_bounties", oldest(args.bounties.map((bounty) => bounty.updatedAt ?? bounty.discoveredAt)));
const staleCount = items.filter((item) => item.status === "stale").length;
const degradedCount = items.filter((item) => item.status === "degraded").length;
const blockedCount = items.filter((item) => item.status === "blocked").length;
const missingCount = items.filter((item) => item.status === "missing").length;
const launchBlockingCount = items.filter((item) => item.launchBlocking).length;
const status = blockedCount > 0 ? "blocked" : staleCount + degradedCount + missingCount > 0 ? "degraded" : "fresh";
return { status, generatedAt: nowIso(), staleCount, degradedCount, blockedCount, missingCount, launchBlockingCount, repairRecommended: status !== "fresh", items, warnings: items.filter((item) => item.status !== "fresh").map((item) => item.summary) };
}

export function freshnessAuditMetadata(report: FreshnessSloReport) {
return {
status: report.status,
staleCount: report.staleCount,
degradedCount: report.degradedCount,
blockedCount: report.blockedCount,
missingCount: report.missingCount,
launchBlockingCount: report.launchBlockingCount,
repairRecommended: report.repairRecommended,
affectedAreas: [...new Set(report.items.filter((item) => item.status !== "fresh").map((item) => item.area))],
};
}

export function buildRepoDataQuality(
repoFullName: string,
syncState: RepoSyncStateRecord | null | undefined,
Expand Down Expand Up @@ -260,6 +344,23 @@ function groupByRepo<T extends { repoFullName: string }>(records: T[]): Map<stri
return grouped;
}

function groupBy<T>(records: T[], keyFor: (record: T) => string): Map<string, T[]> {
const grouped = new Map<string, T[]>();
for (const record of records) {
const key = keyFor(record);
grouped.set(key, [...(grouped.get(key) ?? []), record]);
}
return grouped;
}

function oldest(values: Array<string | null | undefined> | undefined): string | null | undefined {
return values?.filter((value): value is string => Boolean(value && Number.isFinite(Date.parse(value)))).sort()[0];
}

function newest(values: Array<string | null | undefined> | undefined): string | null | undefined {
return values?.filter((value): value is string => Boolean(value && Number.isFinite(Date.parse(value)))).sort().at(-1);
}

function isCompleteCount(segment: RepoSyncSegmentRecord | undefined, expected: number | null | undefined): boolean {
return Boolean(segment && hasCompleteCountCoverage(segment, expected) && hasUsableRequiredSegmentCoverage(segment, expected));
}
Expand Down
Loading