From cf79dd1e29c30c219dafcb8d1dc719d18e3c1a3f Mon Sep 17 00:00:00 2001 From: mkdev11 Date: Fri, 29 May 2026 11:06:41 +0200 Subject: [PATCH 1/5] feat(signals): add burden forecast reader and cache-or-compute helper New getBurdenForecast reader on the existing burden_forecasts table plus a services/burden-forecast loader that returns an envelope with generatedAt, ageSeconds, and a freshness marker (fresh | stale, 6h window). Falls back to a computed forecast for known repos that have no cached snapshot yet, so callers see actionable data on first hit without waiting for the scheduled job. --- src/db/repositories.ts | 12 +++++++ src/services/burden-forecast.ts | 61 +++++++++++++++++++++++++++++++++ 2 files changed, 73 insertions(+) create mode 100644 src/services/burden-forecast.ts diff --git a/src/db/repositories.ts b/src/db/repositories.ts index 8de27569a3..8e1f9b9ee1 100644 --- a/src/db/repositories.ts +++ b/src/db/repositories.ts @@ -935,6 +935,18 @@ export async function upsertBurdenForecast(env: Env, forecast: BurdenForecastRec }); } +export async function getBurdenForecast(env: Env, repoFullName: string): Promise { + const db = getDb(env.DB); + const row = await db.select().from(burdenForecasts).where(eq(burdenForecasts.repoFullName, repoFullName)).limit(1); + const first = row[0]; + if (!first) return null; + return { + repoFullName: first.repoFullName, + payload: parseJson>(first.payloadJson, {}), + generatedAt: first.generatedAt, + }; +} + export async function persistRegistryDriftEvents(env: Env, events: RegistryDriftEventRecord[]): Promise { const db = getDb(env.DB); for (const event of events) { diff --git a/src/services/burden-forecast.ts b/src/services/burden-forecast.ts new file mode 100644 index 0000000000..afb375f03e --- /dev/null +++ b/src/services/burden-forecast.ts @@ -0,0 +1,61 @@ +import { + getBurdenForecast, + getRepository, + listIssueSignalSample, + listOpenPullRequests, + listRecentMergedPullRequests, +} from "../db/repositories"; +import { buildBurdenForecast, buildCollisionReport } from "../signals/engine"; + +export const BURDEN_FORECAST_MAX_AGE_MS = 6 * 60 * 60 * 1000; + +export type BurdenForecastFreshness = "fresh" | "stale"; + +export type BurdenForecastResponse = { + status: "ready"; + source: "snapshot" | "computed"; + repoFullName: string; + generatedAt: string; + ageSeconds: number; + freshness: BurdenForecastFreshness; + report: Record; +}; + +export async function loadOrComputeBurdenForecastResponse(env: Env, fullName: string): Promise { + const cached = await getBurdenForecast(env, fullName); + if (cached) { + const ageMs = forecastAgeMs(cached.generatedAt); + return { + status: "ready", + source: "snapshot", + repoFullName: fullName, + generatedAt: cached.generatedAt, + ageSeconds: Math.max(0, Math.floor(ageMs / 1000)), + freshness: ageMs > BURDEN_FORECAST_MAX_AGE_MS ? "stale" : "fresh", + report: cached.payload, + }; + } + const repo = await getRepository(env, fullName); + if (!repo) return null; + const [issues, pullRequests, recentMergedPullRequests] = await Promise.all([ + listIssueSignalSample(env, fullName), + listOpenPullRequests(env, fullName), + listRecentMergedPullRequests(env, fullName), + ]); + const collisions = buildCollisionReport(fullName, issues, pullRequests, recentMergedPullRequests); + const report = buildBurdenForecast(repo, issues, pullRequests, collisions, 30); + return { + status: "ready", + source: "computed", + repoFullName: fullName, + generatedAt: report.generatedAt, + ageSeconds: 0, + freshness: "fresh", + report: report as unknown as Record, + }; +} + +function forecastAgeMs(generatedAt: string): number { + const parsed = Date.parse(generatedAt); + return Number.isFinite(parsed) ? Date.now() - parsed : Number.POSITIVE_INFINITY; +} From c666f1d59cefae7b36a0bc6966ec370ea90d4e5b Mon Sep 17 00:00:00 2001 From: mkdev11 Date: Fri, 29 May 2026 11:06:43 +0200 Subject: [PATCH 2/5] feat(intelligence): include burden forecast in repo intelligence response buildRepoIntelligenceResponse now attaches the cached burden forecast plus a burdenForecastFreshness slice (source, generatedAt, ageSeconds, freshness) to both the snapshot and computed branches. Reads from the dedicated burden_forecasts table, so request-time scans are limited to the indexed PK lookup. --- src/api/routes.ts | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/src/api/routes.ts b/src/api/routes.ts index e99ebd1da7..ccbe51a352 100644 --- a/src/api/routes.ts +++ b/src/api/routes.ts @@ -78,6 +78,7 @@ import { loadContributorDecisionPackForServing, repoDecisionFromPack, } from "../services/decision-pack"; +import { loadOrComputeBurdenForecastResponse } from "../services/burden-forecast"; import { buildBountyAdvisory, buildBurdenForecast, @@ -1033,7 +1034,7 @@ export function createApp() { } async function buildRepoIntelligenceResponse(env: Env, fullName: string) { - const [repo, snapshots, dataQuality] = await Promise.all([ + const [repo, snapshots, dataQuality, burdenForecast] = await Promise.all([ getRepository(env, fullName), Promise.all( ["queue-health", "config-quality", "label-audit", "maintainer-lane", "maintainer-cut-readiness", "contributor-intake-health"].map(async (signalType) => [ @@ -1042,8 +1043,20 @@ async function buildRepoIntelligenceResponse(env: Env, fullName: string) { ]), ), loadRepoDataQuality(env, fullName), + loadOrComputeBurdenForecastResponse(env, fullName).catch(() => null), ]); const snapshotMap = Object.fromEntries(snapshots); + const burdenForecastSlice = burdenForecast + ? { + burdenForecast: burdenForecast.report, + burdenForecastFreshness: { + source: burdenForecast.source, + generatedAt: burdenForecast.generatedAt, + ageSeconds: burdenForecast.ageSeconds, + freshness: burdenForecast.freshness, + }, + } + : {}; if (snapshotMap["queue-health"] && snapshotMap["config-quality"] && snapshotMap["label-audit"]) { return { status: "ready", @@ -1059,6 +1072,7 @@ async function buildRepoIntelligenceResponse(env: Env, fullName: string) { maintainerCutReadiness: snapshotMap["maintainer-cut-readiness"], contributorIntakeHealth: snapshotMap["contributor-intake-health"], dataQuality, + ...burdenForecastSlice, }; } const [issues, pullRequests, recentMergedPullRequests, labels, queueCounts] = await Promise.all([ @@ -1090,6 +1104,7 @@ async function buildRepoIntelligenceResponse(env: Env, fullName: string) { maintainerCutReadiness, contributorIntakeHealth, dataQuality, + ...burdenForecastSlice, }; } From c408a141f89083c2ecd43bed797a979ae0a24fd5 Mon Sep 17 00:00:00 2001 From: mkdev11 Date: Fri, 29 May 2026 11:06:46 +0200 Subject: [PATCH 3/5] feat(mcp): add gittensory_get_burden_forecast tool Returns the cached or freshly-computed burden forecast for a repo along with its freshness marker so MCP consumers can decide whether to retry once the next scheduled rebuild lands. Falls back to a not_found payload when the repo is unknown. --- src/mcp/server.ts | 28 ++++++++++++++++++++++++++++ 1 file changed, 28 insertions(+) diff --git a/src/mcp/server.ts b/src/mcp/server.ts index f85179f4d5..1173fc800f 100644 --- a/src/mcp/server.ts +++ b/src/mcp/server.ts @@ -36,6 +36,7 @@ import { startAgentRun, } from "../services/agent-orchestrator"; import { loadContributorDecisionPackForServing, repoDecisionFromPack } from "../services/decision-pack"; +import { loadOrComputeBurdenForecastResponse } from "../services/burden-forecast"; import { buildBountyAdvisory, buildCollisionReport, @@ -244,6 +245,15 @@ export class GittensoryMcp { async (input) => this.toolResult(await this.getRepoContext(input)), ); + server.registerTool( + "gittensory_get_burden_forecast", + { + description: "Return the cached or freshly-computed maintainer burden forecast for a repo, including projected review load, queue growth risk, stale PR signals, and a freshness marker.", + inputSchema: ownerRepoShape, + }, + async (input) => this.toolResult(await this.getBurdenForecast(input)), + ); + server.registerTool( "gittensory_get_contributor_profile", { @@ -487,6 +497,24 @@ export class GittensoryMcp { }; } + private async getBurdenForecast(input: { owner: string; repo: string }): Promise { + const fullName = `${input.owner}/${input.repo}`; + const response = await loadOrComputeBurdenForecastResponse(this.env, fullName); + if (!response) { + return { + summary: `Gittensory has no cached burden forecast for ${fullName}.`, + data: { status: "not_found", repoFullName: fullName }, + }; + } + return { + summary: + response.source === "snapshot" + ? `Gittensory burden forecast for ${fullName} (cached, ${response.freshness}).` + : `Gittensory burden forecast for ${fullName} (computed from cached metadata).`, + data: response as unknown as Record, + }; + } + private async loadOpenQueueCounts(fullName: string): Promise<{ openIssues: number; openPullRequests: number }> { const [totals, openIssues, openPullRequests] = await Promise.all([ getLatestRepoGithubTotalsSnapshot(this.env, fullName), From 076b948314aeb71ac9776da97a83f6bb0a0ccd8a Mon Sep 17 00:00:00 2001 From: mkdev11 Date: Fri, 29 May 2026 11:06:49 +0200 Subject: [PATCH 4/5] test(signals): cover burden forecast fixtures and wiring - Builder fixtures: small clean queue, ragflow/sure-scale critical queue, duplicate-cluster trend, stale PR trend. - Service fixtures: snapshot freshness fresh/stale, computed fallback for known but uncached repos, broad-lister spy regression so the cached path never triggers request-time scans. - Queue.test.ts now asserts a burden forecast row was persisted after the build-burden-forecasts job. - api.test.ts asserts the intelligence response includes the cached forecast plus its freshness slice, and that the MCP tool round-trips through cache + missing-repo branches. --- test/integration/api.test.ts | 33 +++++- test/unit/burden-forecast.test.ts | 187 ++++++++++++++++++++++++++++++ test/unit/queue.test.ts | 4 + 3 files changed, 223 insertions(+), 1 deletion(-) create mode 100644 test/unit/burden-forecast.test.ts diff --git a/test/integration/api.test.ts b/test/integration/api.test.ts index edd5f32d37..4cf4ac3a28 100644 --- a/test/integration/api.test.ts +++ b/test/integration/api.test.ts @@ -1,6 +1,7 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import { upsertBounty, + upsertBurdenForecast, upsertCheckSummary, upsertInstallation, upsertInstallationHealth, @@ -23,6 +24,7 @@ import { createApp } from "../../src/api/routes"; import { normalizeRegistryPayload } from "../../src/registry/normalize"; import { persistRegistrySnapshot } from "../../src/registry/sync"; import { createTestEnv } from "../helpers/d1"; +import type { JsonValue } from "../../src/types"; describe("api routes", () => { afterEach(() => { @@ -720,9 +722,18 @@ describe("api routes", () => { generatedAt: "2026-05-25T00:00:00.000Z", }); } + await upsertBurdenForecast(env, { + repoFullName: "entrius/allways-ui", + payload: { repoFullName: "entrius/allways-ui", level: "medium", summary: "intelligence fixture" } as unknown as Record, + generatedAt: "2026-05-25T00:00:00.000Z", + }); const snapshotIntelligence = await app.request("/v1/repos/entrius/allways-ui/intelligence", { headers: apiHeaders(env) }, env); expect(snapshotIntelligence.status).toBe(200); - await expect(snapshotIntelligence.json()).resolves.toMatchObject({ source: "snapshot", queueHealth: { signals: { openPullRequests: 2 } } }); + const snapshotIntelligenceBody = (await snapshotIntelligence.json()) as Record & { burdenForecast?: Record; burdenForecastFreshness?: { freshness: string; source: string; ageSeconds: number } }; + expect(snapshotIntelligenceBody).toMatchObject({ source: "snapshot", queueHealth: { signals: { openPullRequests: 2 } } }); + expect(snapshotIntelligenceBody.burdenForecast).toMatchObject({ level: "medium" }); + expect(snapshotIntelligenceBody.burdenForecastFreshness).toMatchObject({ source: "snapshot", freshness: "stale" }); + expect(snapshotIntelligenceBody.burdenForecastFreshness?.ageSeconds).toBeGreaterThan(0); for (const path of [ "/v1/repos/entrius/allways-ui/issue-quality", @@ -1128,6 +1139,7 @@ describe("api routes", () => { const toolsPayload = (await mcpJson(toolsList)) as { result: { tools: Array<{ name: string }> } }; const toolNames = toolsPayload.result.tools.map((tool) => tool.name); expect(toolNames).toContain("gittensory_get_repo_context"); + expect(toolNames).toContain("gittensory_get_burden_forecast"); expect(toolNames).toContain("gittensory_get_contributor_profile"); expect(toolNames).toContain("gittensory_get_decision_pack"); expect(toolNames).toContain("gittensory_explain_repo_decision"); @@ -1245,8 +1257,27 @@ describe("api routes", () => { expect(missingRepoDecision.status).toBe(200); await expect(mcpJson(missingRepoDecision)).resolves.toMatchObject({ result: { structuredContent: { status: "not_found", decision: null } } }); + const missingBurdenForecast = await app.request( + "/mcp", + { + method: "POST", + headers: mcpHeaders(env), + body: JSON.stringify({ jsonrpc: "2.0", id: "missing-burden", method: "tools/call", params: { name: "gittensory_get_burden_forecast", arguments: { owner: "ghost", repo: "nothing" } } }), + }, + env, + ); + expect(missingBurdenForecast.status).toBe(200); + await expect(mcpJson(missingBurdenForecast)).resolves.toMatchObject({ result: { structuredContent: { status: "not_found", repoFullName: "ghost/nothing" } } }); + + await upsertBurdenForecast(env, { + repoFullName: "entrius/allways-ui", + payload: { repoFullName: "entrius/allways-ui", level: "low", summary: "mcp fixture", forecast: { projectedReviewLoad: 0, queueGrowthRisk: 0, stalePullRequests: 0, duplicateTrend: 0, reviewablePullRequests: 0 }, findings: [] } as unknown as Record, + generatedAt: new Date(Date.now() - 1000).toISOString(), + }); + for (const [name, args] of [ ["gittensory_get_repo_context", { owner: "entrius", repo: "allways-ui" }], + ["gittensory_get_burden_forecast", { owner: "entrius", repo: "allways-ui" }], ["gittensory_get_contributor_profile", { login: "oktofeesh1" }], ["gittensory_get_decision_pack", { login: "oktofeesh1" }], ["gittensory_explain_repo_decision", { login: "oktofeesh1", owner: "entrius", repo: "allways-ui" }], diff --git a/test/unit/burden-forecast.test.ts b/test/unit/burden-forecast.test.ts new file mode 100644 index 0000000000..cf19023ac4 --- /dev/null +++ b/test/unit/burden-forecast.test.ts @@ -0,0 +1,187 @@ +import { describe, expect, it, vi } from "vitest"; +import { getBurdenForecast, upsertBurdenForecast, upsertRepositoryFromGitHub } from "../../src/db/repositories"; +import { BURDEN_FORECAST_MAX_AGE_MS, loadOrComputeBurdenForecastResponse } from "../../src/services/burden-forecast"; +import { buildBurdenForecast, buildCollisionReport } from "../../src/signals/engine"; +import type { IssueRecord, JsonValue, PullRequestRecord, RepositoryRecord } from "../../src/types"; +import { createTestEnv } from "../helpers/d1"; + +describe("burden forecast builder", () => { + it("classifies a small clean queue as low burden with no findings", () => { + const repo = repoFixture("owner/small"); + const forecast = buildBurdenForecast(repo, [], [], buildCollisionReport(repo.fullName, [], []), 7); + expect(forecast.level).toBe("low"); + expect(forecast.findings).toEqual([]); + }); + + it("stays bounded on a ragflow/sure-style large queue and emits critical findings", () => { + const repo = repoFixture("ragflow/ragflow"); + const stalePr = pr(repo.fullName, 999, "stale", { updatedAt: "2025-01-01T00:00:00.000Z" }); + const open = Array.from({ length: 120 }, (_, index) => pr(repo.fullName, index + 1, `open ${index}`, { linkedIssues: [], updatedAt: new Date().toISOString() })); + const forecast = buildBurdenForecast(repo, [], [stalePr, ...open], buildCollisionReport(repo.fullName, [], [stalePr, ...open]), 30); + expect(forecast.level).toBe("critical"); + expect(forecast.findings.map((f) => f.code)).toEqual(expect.arrayContaining(["queue_growth_risk", "stale_review_load"])); + expect(forecast.forecast.stalePullRequests).toBeGreaterThan(0); + }); + + it("counts the duplicate cluster trend when multiple PRs reference the same issue", () => { + const repo = repoFixture("owner/duplicates"); + const issueRecord = issue(repo.fullName, 42, "Login flow broken"); + const a = pr(repo.fullName, 1, "Fix login flow", { linkedIssues: [42] }); + const b = pr(repo.fullName, 2, "Alternative login fix", { linkedIssues: [42] }); + const collisions = buildCollisionReport(repo.fullName, [issueRecord], [a, b]); + const forecast = buildBurdenForecast(repo, [issueRecord], [a, b], collisions, 7); + expect(collisions.summary.clusterCount).toBeGreaterThan(0); + expect(forecast.forecast.duplicateTrend).toBe(collisions.summary.clusterCount); + }); + + it("surfaces a stale PR trend in the forecast findings", () => { + const repo = repoFixture("owner/stale"); + const stalePrs = Array.from({ length: 4 }, (_, index) => pr(repo.fullName, index + 1, `stale ${index}`, { updatedAt: "2025-01-01T00:00:00.000Z", linkedIssues: [] })); + const forecast = buildBurdenForecast(repo, [], stalePrs, buildCollisionReport(repo.fullName, [], stalePrs), 30); + expect(forecast.forecast.stalePullRequests).toBe(4); + expect(forecast.findings.find((f) => f.code === "stale_review_load")?.detail).toContain("4 open PR(s)"); + }); +}); + +describe("loadOrComputeBurdenForecastResponse", () => { + it("returns null when the repo is unknown", async () => { + const env = createTestEnv(); + const response = await loadOrComputeBurdenForecastResponse(env, "ghost/missing"); + expect(response).toBeNull(); + }); + + it("returns a snapshot envelope with freshness:fresh for a recently persisted forecast", async () => { + const env = createTestEnv(); + await upsertRepositoryFromGitHub(env, { name: "fresh", full_name: "owner/fresh", private: false, owner: { login: "owner" }, default_branch: "main" }); + await upsertBurdenForecast(env, { + repoFullName: "owner/fresh", + payload: { repoFullName: "owner/fresh", level: "low", summary: "fresh fixture" } as unknown as Record, + generatedAt: new Date(Date.now() - 60_000).toISOString(), + }); + const response = await loadOrComputeBurdenForecastResponse(env, "owner/fresh"); + expect(response).toMatchObject({ + status: "ready", + source: "snapshot", + repoFullName: "owner/fresh", + freshness: "fresh", + report: { level: "low" }, + }); + expect(response?.ageSeconds).toBeGreaterThanOrEqual(0); + expect(response?.ageSeconds).toBeLessThan(BURDEN_FORECAST_MAX_AGE_MS / 1000); + }); + + it("surfaces freshness:stale when the cached forecast is older than the max age", async () => { + const env = createTestEnv(); + await upsertRepositoryFromGitHub(env, { name: "old", full_name: "owner/old", private: false, owner: { login: "owner" }, default_branch: "main" }); + await upsertBurdenForecast(env, { + repoFullName: "owner/old", + payload: { repoFullName: "owner/old", level: "high", summary: "stale fixture" } as unknown as Record, + generatedAt: "2025-01-01T00:00:00.000Z", + }); + const response = await loadOrComputeBurdenForecastResponse(env, "owner/old"); + expect(response).toMatchObject({ + status: "ready", + source: "snapshot", + freshness: "stale", + }); + expect(response?.ageSeconds).toBeGreaterThan(BURDEN_FORECAST_MAX_AGE_MS / 1000); + }); + + it("falls back to a computed forecast when no snapshot exists but the repo is known", async () => { + const env = createTestEnv(); + await upsertRepositoryFromGitHub(env, { name: "uncached", full_name: "owner/uncached", private: false, owner: { login: "owner" }, default_branch: "main" }); + const response = await loadOrComputeBurdenForecastResponse(env, "owner/uncached"); + expect(response).toMatchObject({ + status: "ready", + source: "computed", + freshness: "fresh", + ageSeconds: 0, + }); + expect(response?.report).toMatchObject({ repoFullName: "owner/uncached", level: "low" }); + }); + + it("does not call broad request-time listers when a cached forecast exists", async () => { + const env = createTestEnv(); + await upsertRepositoryFromGitHub(env, { name: "perf", full_name: "owner/perf", private: false, owner: { login: "owner" }, default_branch: "main" }); + await upsertBurdenForecast(env, { + repoFullName: "owner/perf", + payload: { repoFullName: "owner/perf", level: "low", summary: "fixture" } as unknown as Record, + generatedAt: new Date(Date.now() - 1000).toISOString(), + }); + const repositoriesModule = await import("../../src/db/repositories"); + const spies = [ + vi.spyOn(repositoriesModule, "listIssueSignalSample"), + vi.spyOn(repositoriesModule, "listOpenPullRequests"), + vi.spyOn(repositoriesModule, "listRecentMergedPullRequests"), + ]; + await loadOrComputeBurdenForecastResponse(env, "owner/perf"); + for (const spy of spies) { + expect(spy).not.toHaveBeenCalled(); + spy.mockRestore(); + } + }); + + it("getBurdenForecast round-trips through upsert", async () => { + const env = createTestEnv(); + await upsertBurdenForecast(env, { + repoFullName: "owner/round-trip", + payload: { level: "medium", summary: "round-trip" } as unknown as Record, + generatedAt: "2026-05-25T00:00:00.000Z", + }); + const row = await getBurdenForecast(env, "owner/round-trip"); + expect(row).toMatchObject({ repoFullName: "owner/round-trip", generatedAt: "2026-05-25T00:00:00.000Z" }); + expect(row?.payload).toMatchObject({ level: "medium", summary: "round-trip" }); + }); +}); + +function repoFixture(fullName: string): RepositoryRecord { + const [owner, name] = fullName.split("/"); + return { + fullName, + owner, + name, + isInstalled: true, + isRegistered: true, + isPrivate: false, + registryConfig: { + repo: fullName, + emissionShare: 0.02, + issueDiscoveryShare: 0, + maintainerCut: 0, + labelMultipliers: {}, + raw: {}, + }, + } as RepositoryRecord; +} + +function issue(repoFullName: string, number: number, title: string, overrides: Partial = {}): IssueRecord { + return { + repoFullName, + number, + title, + state: "open", + authorLogin: "reporter", + authorAssociation: "NONE", + labels: [], + linkedPrs: [], + body: "Detailed body for collision testing with enough content to be useful.", + updatedAt: new Date().toISOString(), + ...overrides, + } as IssueRecord; +} + +function pr(repoFullName: string, number: number, title: string, overrides: Partial = {}): PullRequestRecord { + return { + repoFullName, + number, + title, + state: "open", + authorLogin: "dev", + authorAssociation: "NONE", + labels: [], + linkedIssues: [], + body: "", + updatedAt: new Date().toISOString(), + ...overrides, + } as PullRequestRecord; +} diff --git a/test/unit/queue.test.ts b/test/unit/queue.test.ts index 7250fa161f..561c7b4c93 100644 --- a/test/unit/queue.test.ts +++ b/test/unit/queue.test.ts @@ -2,6 +2,7 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import { listCollisionEdges, createAgentRun, + getBurdenForecast, getContributorEvidence, getAgentRun, getContributorScoringProfile, @@ -100,6 +101,9 @@ describe("queue processors", () => { expect(await listSignalSnapshots(env, "contributor-decision-pack", "oktofeesh1")).not.toHaveLength(0); expect(await getContributorEvidence(env, "oktofeesh1")).toMatchObject({ login: "oktofeesh1" }); expect(await getContributorScoringProfile(env, "oktofeesh1")).toMatchObject({ login: "oktofeesh1" }); + const persistedBurden = await getBurdenForecast(env, "JSONbored/gittensory"); + expect(persistedBurden).toMatchObject({ repoFullName: "JSONbored/gittensory" }); + expect(persistedBurden?.payload).toMatchObject({ level: expect.any(String), summary: expect.any(String) }); }); it("runs queued agent jobs through the queue processor", async () => { From f5407dbb1e1a6d2100ea42c42b1880f2d3b08cdd Mon Sep 17 00:00:00 2001 From: mkdev11 Date: Fri, 29 May 2026 11:29:26 +0200 Subject: [PATCH 5/5] fix(signals): tighten burden forecast contract --- src/api/routes.ts | 24 +++++++-- src/openapi/schemas.ts | 9 ++++ src/services/burden-forecast.ts | 8 +-- test/integration/api.test.ts | 85 ++++++++++++++++++++++++++++++- test/unit/burden-forecast.test.ts | 45 ++++++++++++---- test/unit/openapi.test.ts | 1 + 6 files changed, 153 insertions(+), 19 deletions(-) diff --git a/src/api/routes.ts b/src/api/routes.ts index ccbe51a352..36bf985a89 100644 --- a/src/api/routes.ts +++ b/src/api/routes.ts @@ -103,7 +103,7 @@ import { attachDataQuality, buildCoreSignalFidelity, buildFreshnessSloReport, bu import { buildPullRequestReviewability } from "../signals/reward-risk"; import { buildLocalBranchAnalysis } from "../signals/local-branch"; import { buildRepoSettingsPreview } from "../signals/settings-preview"; -import type { ContributorEvidenceRecord, JobMessage, JsonValue, RepoSyncSegmentRecord } from "../types"; +import type { ContributorEvidenceRecord, DataQuality, JobMessage, JsonValue, RepoSyncSegmentRecord } from "../types"; import { errorMessage, nowIso } from "../utils/json"; type AppBindings = { Bindings: Env }; @@ -1034,6 +1034,7 @@ export function createApp() { } async function buildRepoIntelligenceResponse(env: Env, fullName: string) { + let burdenForecastError: unknown; const [repo, snapshots, dataQuality, burdenForecast] = await Promise.all([ getRepository(env, fullName), Promise.all( @@ -1043,8 +1044,14 @@ async function buildRepoIntelligenceResponse(env: Env, fullName: string) { ]), ), loadRepoDataQuality(env, fullName), - loadOrComputeBurdenForecastResponse(env, fullName).catch(() => null), + loadOrComputeBurdenForecastResponse(env, fullName).catch((error) => { + burdenForecastError = error; + return null; + }), ]); + const intelligenceDataQuality = burdenForecastError + ? withDataQualityWarning(dataQuality, `Burden forecast unavailable for ${fullName}: ${errorMessage(burdenForecastError)}`) + : dataQuality; const snapshotMap = Object.fromEntries(snapshots); const burdenForecastSlice = burdenForecast ? { @@ -1071,7 +1078,7 @@ async function buildRepoIntelligenceResponse(env: Env, fullName: string) { maintainerLane: snapshotMap["maintainer-lane"], maintainerCutReadiness: snapshotMap["maintainer-cut-readiness"], contributorIntakeHealth: snapshotMap["contributor-intake-health"], - dataQuality, + dataQuality: intelligenceDataQuality, ...burdenForecastSlice, }; } @@ -1103,11 +1110,20 @@ async function buildRepoIntelligenceResponse(env: Env, fullName: string) { maintainerLane, maintainerCutReadiness, contributorIntakeHealth, - dataQuality, + dataQuality: intelligenceDataQuality, ...burdenForecastSlice, }; } +function withDataQualityWarning(dataQuality: DataQuality, warning: string): DataQuality { + return { + ...dataQuality, + status: dataQuality.status === "complete" ? "degraded" : dataQuality.status, + partial: true, + warnings: [...new Set([...dataQuality.warnings, warning])], + }; +} + async function buildRegistrationReadinessResponse(env: Env, fullName: string) { const intelligence = await buildRepoIntelligenceResponse(env, fullName); const settings = await getRepositorySettings(env, fullName); diff --git a/src/openapi/schemas.ts b/src/openapi/schemas.ts index 12b5235c25..9a79b7e2af 100644 --- a/src/openapi/schemas.ts +++ b/src/openapi/schemas.ts @@ -1117,6 +1117,15 @@ export const RepoIntelligenceSchema = z maintainerCutReadiness: z.record(z.unknown()).nullable().optional(), contributorIntakeHealth: z.record(z.unknown()).nullable().optional(), dataQuality: z.record(z.unknown()), + burdenForecast: BurdenForecastSchema.optional(), + burdenForecastFreshness: z + .object({ + source: z.enum(["snapshot", "computed"]), + generatedAt: z.string(), + ageSeconds: z.number(), + freshness: z.enum(["fresh", "stale"]), + }) + .optional(), }) .openapi("RepoIntelligence"); diff --git a/src/services/burden-forecast.ts b/src/services/burden-forecast.ts index afb375f03e..09feb056fa 100644 --- a/src/services/burden-forecast.ts +++ b/src/services/burden-forecast.ts @@ -5,7 +5,7 @@ import { listOpenPullRequests, listRecentMergedPullRequests, } from "../db/repositories"; -import { buildBurdenForecast, buildCollisionReport } from "../signals/engine"; +import { buildBurdenForecast, buildCollisionReport, type BurdenForecast } from "../signals/engine"; export const BURDEN_FORECAST_MAX_AGE_MS = 6 * 60 * 60 * 1000; @@ -18,7 +18,7 @@ export type BurdenForecastResponse = { generatedAt: string; ageSeconds: number; freshness: BurdenForecastFreshness; - report: Record; + report: BurdenForecast; }; export async function loadOrComputeBurdenForecastResponse(env: Env, fullName: string): Promise { @@ -32,7 +32,7 @@ export async function loadOrComputeBurdenForecastResponse(env: Env, fullName: st generatedAt: cached.generatedAt, ageSeconds: Math.max(0, Math.floor(ageMs / 1000)), freshness: ageMs > BURDEN_FORECAST_MAX_AGE_MS ? "stale" : "fresh", - report: cached.payload, + report: cached.payload as unknown as BurdenForecast, }; } const repo = await getRepository(env, fullName); @@ -51,7 +51,7 @@ export async function loadOrComputeBurdenForecastResponse(env: Env, fullName: st generatedAt: report.generatedAt, ageSeconds: 0, freshness: "fresh", - report: report as unknown as Record, + report, }; } diff --git a/test/integration/api.test.ts b/test/integration/api.test.ts index 4cf4ac3a28..87c4ede6eb 100644 --- a/test/integration/api.test.ts +++ b/test/integration/api.test.ts @@ -21,6 +21,7 @@ import { upsertRepositorySettings, } from "../../src/db/repositories"; import { createApp } from "../../src/api/routes"; +import { BURDEN_FORECAST_MAX_AGE_MS } from "../../src/services/burden-forecast"; import { normalizeRegistryPayload } from "../../src/registry/normalize"; import { persistRegistrySnapshot } from "../../src/registry/sync"; import { createTestEnv } from "../helpers/d1"; @@ -722,10 +723,11 @@ describe("api routes", () => { generatedAt: "2026-05-25T00:00:00.000Z", }); } + const staleForecastGeneratedAt = new Date(Date.now() - BURDEN_FORECAST_MAX_AGE_MS - 60_000).toISOString(); await upsertBurdenForecast(env, { repoFullName: "entrius/allways-ui", payload: { repoFullName: "entrius/allways-ui", level: "medium", summary: "intelligence fixture" } as unknown as Record, - generatedAt: "2026-05-25T00:00:00.000Z", + generatedAt: staleForecastGeneratedAt, }); const snapshotIntelligence = await app.request("/v1/repos/entrius/allways-ui/intelligence", { headers: apiHeaders(env) }, env); expect(snapshotIntelligence.status).toBe(200); @@ -733,7 +735,25 @@ describe("api routes", () => { expect(snapshotIntelligenceBody).toMatchObject({ source: "snapshot", queueHealth: { signals: { openPullRequests: 2 } } }); expect(snapshotIntelligenceBody.burdenForecast).toMatchObject({ level: "medium" }); expect(snapshotIntelligenceBody.burdenForecastFreshness).toMatchObject({ source: "snapshot", freshness: "stale" }); - expect(snapshotIntelligenceBody.burdenForecastFreshness?.ageSeconds).toBeGreaterThan(0); + expect(snapshotIntelligenceBody.burdenForecastFreshness?.ageSeconds).toBeGreaterThanOrEqual(Math.floor((BURDEN_FORECAST_MAX_AGE_MS + 50_000) / 1000)); + expect(snapshotIntelligenceBody.burdenForecastFreshness?.ageSeconds).toBeLessThan(Math.floor((BURDEN_FORECAST_MAX_AGE_MS + 120_000) / 1000)); + + await upsertRepositoryFromGitHub(env, { name: "uncached-burden", full_name: "entrius/uncached-burden", private: false, owner: { login: "entrius" }, default_branch: "main" }); + const computedIntelligence = await app.request("/v1/repos/entrius/uncached-burden/intelligence", { headers: apiHeaders(env) }, env); + expect(computedIntelligence.status).toBe(200); + await expect(computedIntelligence.json()).resolves.toMatchObject({ + source: "computed", + burdenForecast: { repoFullName: "entrius/uncached-burden", level: "low" }, + burdenForecastFreshness: { source: "computed", freshness: "fresh", ageSeconds: 0 }, + }); + + const degradedForecastEnv = withBurdenForecastReadFailure(env); + const degradedIntelligence = await app.request("/v1/repos/entrius/allways-ui/intelligence", { headers: apiHeaders(env) }, degradedForecastEnv); + expect(degradedIntelligence.status).toBe(200); + const degradedBody = (await degradedIntelligence.json()) as Record & { dataQuality: { status: string; warnings: string[] }; burdenForecast?: unknown }; + expect(degradedBody.burdenForecast).toBeUndefined(); + expect(degradedBody.dataQuality.status).toBe("degraded"); + expect(degradedBody.dataQuality.warnings).toEqual(expect.arrayContaining([expect.stringMatching(/Burden forecast unavailable/i)])); for (const path of [ "/v1/repos/entrius/allways-ui/issue-quality", @@ -1275,6 +1295,51 @@ describe("api routes", () => { generatedAt: new Date(Date.now() - 1000).toISOString(), }); + const cachedBurdenForecast = await app.request( + "/mcp", + { + method: "POST", + headers: mcpHeaders(env), + body: JSON.stringify({ jsonrpc: "2.0", id: "cached-burden", method: "tools/call", params: { name: "gittensory_get_burden_forecast", arguments: { owner: "entrius", repo: "allways-ui" } } }), + }, + env, + ); + expect(cachedBurdenForecast.status).toBe(200); + await expect(mcpJson(cachedBurdenForecast)).resolves.toMatchObject({ + result: { + structuredContent: { + status: "ready", + source: "snapshot", + repoFullName: "entrius/allways-ui", + freshness: "fresh", + report: { level: "low" }, + }, + }, + }); + + await upsertRepositoryFromGitHub(env, { name: "mcp-computed-burden", full_name: "entrius/mcp-computed-burden", private: false, owner: { login: "entrius" }, default_branch: "main" }); + const computedBurdenForecast = await app.request( + "/mcp", + { + method: "POST", + headers: mcpHeaders(env), + body: JSON.stringify({ jsonrpc: "2.0", id: "computed-burden", method: "tools/call", params: { name: "gittensory_get_burden_forecast", arguments: { owner: "entrius", repo: "mcp-computed-burden" } } }), + }, + env, + ); + expect(computedBurdenForecast.status).toBe(200); + await expect(mcpJson(computedBurdenForecast)).resolves.toMatchObject({ + result: { + structuredContent: { + status: "ready", + source: "computed", + repoFullName: "entrius/mcp-computed-burden", + freshness: "fresh", + report: { repoFullName: "entrius/mcp-computed-burden", level: "low" }, + }, + }, + }); + for (const [name, args] of [ ["gittensory_get_repo_context", { owner: "entrius", repo: "allways-ui" }], ["gittensory_get_burden_forecast", { owner: "entrius", repo: "allways-ui" }], @@ -2094,6 +2159,22 @@ function apiHeaders(env: Env): Record { }; } +function withBurdenForecastReadFailure(env: Env): Env { + const db = env.DB as unknown as { prepare: (sql: string) => unknown; batch: (statements: unknown[]) => Promise }; + return { + ...env, + DB: { + prepare(sql: string) { + if (/burden_forecasts/i.test(sql)) throw new Error("forecast table unavailable"); + return db.prepare(sql); + }, + batch(statements: unknown[]) { + return db.batch(statements); + }, + } as unknown as D1Database, + }; +} + function stubOktofeeshFetch(): void { vi.stubGlobal("fetch", async (input: RequestInfo | URL) => { const url = input.toString(); diff --git a/test/unit/burden-forecast.test.ts b/test/unit/burden-forecast.test.ts index cf19023ac4..a5d45e13d4 100644 --- a/test/unit/burden-forecast.test.ts +++ b/test/unit/burden-forecast.test.ts @@ -15,7 +15,7 @@ describe("burden forecast builder", () => { it("stays bounded on a ragflow/sure-style large queue and emits critical findings", () => { const repo = repoFixture("ragflow/ragflow"); - const stalePr = pr(repo.fullName, 999, "stale", { updatedAt: "2025-01-01T00:00:00.000Z" }); + const stalePr = pr(repo.fullName, 999, "stale", { updatedAt: daysAgo(31) }); const open = Array.from({ length: 120 }, (_, index) => pr(repo.fullName, index + 1, `open ${index}`, { linkedIssues: [], updatedAt: new Date().toISOString() })); const forecast = buildBurdenForecast(repo, [], [stalePr, ...open], buildCollisionReport(repo.fullName, [], [stalePr, ...open]), 30); expect(forecast.level).toBe("critical"); @@ -25,18 +25,19 @@ describe("burden forecast builder", () => { it("counts the duplicate cluster trend when multiple PRs reference the same issue", () => { const repo = repoFixture("owner/duplicates"); - const issueRecord = issue(repo.fullName, 42, "Login flow broken"); - const a = pr(repo.fullName, 1, "Fix login flow", { linkedIssues: [42] }); - const b = pr(repo.fullName, 2, "Alternative login fix", { linkedIssues: [42] }); + const issueRecord = issue(repo.fullName, 42, "Auth failure after reconnect"); + const a = pr(repo.fullName, 1, "Token refresh", { linkedIssues: [42] }); + const b = pr(repo.fullName, 2, "Session restore", { linkedIssues: [42] }); const collisions = buildCollisionReport(repo.fullName, [issueRecord], [a, b]); const forecast = buildBurdenForecast(repo, [issueRecord], [a, b], collisions, 7); - expect(collisions.summary.clusterCount).toBeGreaterThan(0); - expect(forecast.forecast.duplicateTrend).toBe(collisions.summary.clusterCount); + expect(collisions.summary.clusterCount).toBe(4); + expect(collisions.clusters.some((cluster) => cluster.items.map((item) => `${item.type}:${item.number}`).sort().join("|") === "issue:42|pull_request:1|pull_request:2")).toBe(true); + expect(forecast.forecast.duplicateTrend).toBe(4); }); it("surfaces a stale PR trend in the forecast findings", () => { const repo = repoFixture("owner/stale"); - const stalePrs = Array.from({ length: 4 }, (_, index) => pr(repo.fullName, index + 1, `stale ${index}`, { updatedAt: "2025-01-01T00:00:00.000Z", linkedIssues: [] })); + const stalePrs = Array.from({ length: 4 }, (_, index) => pr(repo.fullName, index + 1, `stale ${index}`, { updatedAt: daysAgo(31), linkedIssues: [] })); const forecast = buildBurdenForecast(repo, [], stalePrs, buildCollisionReport(repo.fullName, [], stalePrs), 30); expect(forecast.forecast.stalePullRequests).toBe(4); expect(forecast.findings.find((f) => f.code === "stale_review_load")?.detail).toContain("4 open PR(s)"); @@ -73,10 +74,11 @@ describe("loadOrComputeBurdenForecastResponse", () => { it("surfaces freshness:stale when the cached forecast is older than the max age", async () => { const env = createTestEnv(); await upsertRepositoryFromGitHub(env, { name: "old", full_name: "owner/old", private: false, owner: { login: "owner" }, default_branch: "main" }); + const generatedAt = new Date(Date.now() - BURDEN_FORECAST_MAX_AGE_MS - 60_000).toISOString(); await upsertBurdenForecast(env, { repoFullName: "owner/old", payload: { repoFullName: "owner/old", level: "high", summary: "stale fixture" } as unknown as Record, - generatedAt: "2025-01-01T00:00:00.000Z", + generatedAt, }); const response = await loadOrComputeBurdenForecastResponse(env, "owner/old"); expect(response).toMatchObject({ @@ -84,7 +86,8 @@ describe("loadOrComputeBurdenForecastResponse", () => { source: "snapshot", freshness: "stale", }); - expect(response?.ageSeconds).toBeGreaterThan(BURDEN_FORECAST_MAX_AGE_MS / 1000); + expect(response?.ageSeconds).toBeGreaterThanOrEqual(Math.floor((BURDEN_FORECAST_MAX_AGE_MS + 50_000) / 1000)); + expect(response?.ageSeconds).toBeLessThan(Math.floor((BURDEN_FORECAST_MAX_AGE_MS + 120_000) / 1000)); }); it("falls back to a computed forecast when no snapshot exists but the repo is known", async () => { @@ -121,6 +124,26 @@ describe("loadOrComputeBurdenForecastResponse", () => { } }); + it("uses only the bounded per-repo listers when computing a missing forecast", async () => { + const env = createTestEnv(); + await upsertRepositoryFromGitHub(env, { name: "computed-perf", full_name: "owner/computed-perf", private: false, owner: { login: "owner" }, default_branch: "main" }); + const repositoriesModule = await import("../../src/db/repositories"); + const spies = [ + vi.spyOn(repositoriesModule, "listIssueSignalSample"), + vi.spyOn(repositoriesModule, "listOpenPullRequests"), + vi.spyOn(repositoriesModule, "listRecentMergedPullRequests"), + ]; + + const response = await loadOrComputeBurdenForecastResponse(env, "owner/computed-perf"); + + expect(response).toMatchObject({ source: "computed", report: { repoFullName: "owner/computed-perf" } }); + for (const spy of spies) { + expect(spy).toHaveBeenCalledTimes(1); + expect(spy).toHaveBeenCalledWith(env, "owner/computed-perf"); + spy.mockRestore(); + } + }); + it("getBurdenForecast round-trips through upsert", async () => { const env = createTestEnv(); await upsertBurdenForecast(env, { @@ -154,6 +177,10 @@ function repoFixture(fullName: string): RepositoryRecord { } as RepositoryRecord; } +function daysAgo(days: number): string { + return new Date(Date.now() - days * 86_400_000).toISOString(); +} + function issue(repoFullName: string, number: number, title: string, overrides: Partial = {}): IssueRecord { return { repoFullName, diff --git a/test/unit/openapi.test.ts b/test/unit/openapi.test.ts index f1d656f100..ea6e675258 100644 --- a/test/unit/openapi.test.ts +++ b/test/unit/openapi.test.ts @@ -73,6 +73,7 @@ describe("OpenAPI contract", () => { expect(spec.components?.schemas?.AgentRunBundle).toBeDefined(); expect(spec.components?.schemas?.AgentAction).toBeDefined(); expect(JSON.stringify(spec.components?.schemas?.ScorePreviewResult)).toContain("scenarioPreviews"); + expect(JSON.stringify(spec.components?.schemas?.RepoIntelligence)).toContain("burdenForecastFreshness"); expect(JSON.stringify(spec.components?.schemas?.LocalBranchAnalysis)).toContain("baseFreshness"); expect(JSON.stringify(spec.components?.schemas?.LocalBranchAnalysis)).toContain("recommendedRerunCondition"); });