diff --git a/packages/daemon/src/config.ts b/packages/daemon/src/config.ts index 78a53a7..1f91ef3 100644 --- a/packages/daemon/src/config.ts +++ b/packages/daemon/src/config.ts @@ -64,6 +64,14 @@ export const configSchema = z.object({ infraRetryMax: z.coerce.number().int().nonnegative().default(3), infraRetryBaseMs: z.coerce.number().int().positive().default(15_000), recoveryReport: z.boolean().default(true), + // Stuck-badge sweep: how often to re-verify every "processing" badge. + badgeSweepIntervalMs: z.coerce.number().int().positive().default(60_000), + // A processing badge older than this with no live signal is stale. + badgeTtlMs: z.coerce.number().int().positive().default(10 * 60_000), + // Signal ②: session output must have happened within this window. + badgeOutputTtlMs: z.coerce.number().int().positive().default(10 * 60_000), + // Signal ③: model traffic (token delta) within this window. + badgeModelTtlMs: z.coerce.number().int().positive().default(15 * 60_000), }), db: z.object({ driver: z.enum(["sqlite", "mysql"]).default("sqlite"), @@ -129,6 +137,10 @@ function readWorkSection() { : 3, infraRetryBaseMs: process.env.WORK_INFRA_RETRY_BASE_MS ? Math.max(1, Number(process.env.WORK_INFRA_RETRY_BASE_MS)) : 15_000, recoveryReport: !(process.env.WORK_RECOVERY_REPORT === "false" || process.env.WORK_RECOVERY_REPORT === "0"), + badgeSweepIntervalMs: process.env.WORK_BADGE_SWEEP_MS ? Math.max(1_000, Number(process.env.WORK_BADGE_SWEEP_MS)) : 60_000, + badgeTtlMs: process.env.WORK_BADGE_TTL_MIN ? Math.max(1, Number(process.env.WORK_BADGE_TTL_MIN)) * 60_000 : 10 * 60_000, + badgeOutputTtlMs: process.env.WORK_BADGE_OUTPUT_TTL_MIN ? Math.max(1, Number(process.env.WORK_BADGE_OUTPUT_TTL_MIN)) * 60_000 : 10 * 60_000, + badgeModelTtlMs: process.env.WORK_BADGE_MODEL_TTL_MIN ? Math.max(1, Number(process.env.WORK_BADGE_MODEL_TTL_MIN)) * 60_000 : 15 * 60_000, }; } diff --git a/packages/daemon/src/db.ts b/packages/daemon/src/db.ts index eba4704..af1e0a3 100644 --- a/packages/daemon/src/db.ts +++ b/packages/daemon/src/db.ts @@ -550,6 +550,14 @@ async function runMigrations(db: AsyncDatabase): Promise { "generation", sqlite ? "generation INTEGER NOT NULL DEFAULT 0" : "generation INT NOT NULL DEFAULT 0" ); + // Processing-badge heartbeat deadline (ms epoch). NULL = session not + // currently holding a processing badge; past due = badge must be treated + // as stale by the stuck-badge sweep regardless of other signals. + await ensureColumn( + tSessions, + "expected_heartbeat_at", + sqlite ? "expected_heartbeat_at INTEGER" : "expected_heartbeat_at BIGINT" + ); // messages.model — per-message model override from the webhook payload. // Persisted so queued/nudged/recovered messages keep their model instead of diff --git a/packages/daemon/src/gitea.ts b/packages/daemon/src/gitea.ts index 7d2c8aa..3b58de4 100644 --- a/packages/daemon/src/gitea.ts +++ b/packages/daemon/src/gitea.ts @@ -101,6 +101,12 @@ export class GiteaClient { ); } + async listAiStatusBadges(status: string) { + return this.request<{ badges: Array<{ owner: string; repo: string; number: number; aiStatus: string; since: string | null }> }>( + "GET", `/ai-status?status=${encodeURIComponent(status)}` + ); + } + async listComments(owner: string, repo: string, issueNumber: number) { return this.request< Array<{ id: number; body: string; created_at: string; user: { login: string } }> diff --git a/packages/daemon/src/op.ts b/packages/daemon/src/op.ts index b547731..225b9ec 100644 --- a/packages/daemon/src/op.ts +++ b/packages/daemon/src/op.ts @@ -34,6 +34,7 @@ interface SessionRow { nudge_rounds: number; stuck_nudge_rounds: number; generation: number; + expected_heartbeat_at: number | null; } interface MessageRow { @@ -90,6 +91,7 @@ function rowToSession(row: SessionRow): OpSession { nudgeRounds: row.nudge_rounds ?? 0, stuckNudgeRounds: row.stuck_nudge_rounds ?? 0, generation: row.generation ?? 0, + expectedHeartbeatAt: row.expected_heartbeat_at ?? undefined, }; } @@ -548,6 +550,11 @@ export class Store { await getDB().run("UPDATE {{op_sessions}} SET opencode_session_id = NULL WHERE issue_id = ?", [issueUid]); } + /** Set or clear the processing-badge heartbeat deadline for a session. */ + async setExpectedHeartbeat(sessionId: string, ms: number | null): Promise { + await getDB().run("UPDATE {{op_sessions}} SET expected_heartbeat_at = ? WHERE uid = ?", [ms, sessionId]); + } + async markDaemonStatus(daemonId: number, status: "active" | "drained" | "dead"): Promise { await getDB().run( "UPDATE {{daemons}} SET status = ? WHERE id = ?", diff --git a/packages/daemon/src/opencode.ts b/packages/daemon/src/opencode.ts index 0908732..e1f9e36 100644 --- a/packages/daemon/src/opencode.ts +++ b/packages/daemon/src/opencode.ts @@ -8,6 +8,7 @@ import type { Config } from "./config"; import type { Store } from "./op"; import type { IssueTracker, TrackerRef, TrackerEvent, TrackerComment, Issue, OpSession, Message } from "./trackers/types"; import { formatKey, parseKey } from "./trackers/types"; +import type { BadgeEntry } from "./trackers/types"; import type { RuntimeBackend, RuntimeHandle } from "./runtime/types"; import { OpencodeBackend } from "./runtime/opencode-backend"; import { PiBackend } from "./runtime/pi-backend"; @@ -174,6 +175,80 @@ export function hasRecoveryDelivery( }); } +// ─── Stuck-badge sweep pure helpers (ework#12) ─── + +export interface BadgeSignals { + pidAlive: boolean; + /** ms since last session output; null = no output ever recorded */ + outputAgeMs: number | null; + /** ms since last model traffic (token delta); null = none observed */ + modelAgeMs: number | null; +} + +export interface BadgeSignalLimits { + outputTtlMs: number; + modelTtlMs: number; +} + +/** + * Three-signal liveness verdict (ework#12 spec item 1): a session counts as + * alive when ANY of pid-alive / fresh-output / fresh-model-traffic holds AND + * its heartbeat TTL has not expired. TTL expiry alone forces stale — that is + * what catches "wrote the badge and never managed it" paths where no signal + * source exists at all. Exported for unit testing. + */ +export function evaluateBadgeSignals( + signals: BadgeSignals, + limits: BadgeSignalLimits, + heartbeatExpired: boolean, +): { alive: boolean; reasons: string[] } { + const reasons: string[] = []; + const outputFresh = signals.outputAgeMs !== null && signals.outputAgeMs < limits.outputTtlMs; + const modelFresh = signals.modelAgeMs !== null && signals.modelAgeMs < limits.modelTtlMs; + if (!signals.pidAlive) reasons.push("pid-dead"); + if (!outputFresh) reasons.push("no-recent-output"); + if (!modelFresh) reasons.push("no-recent-model"); + if (heartbeatExpired) reasons.push("heartbeat-expired"); + return { alive: !heartbeatExpired && (signals.pidAlive || outputFresh || modelFresh), reasons }; +} + +export type BadgeAction = "keep" | "grace" | "complete" | "reset"; + +/** + * Arbitration for a stale processing badge (ework#12 spec items 2+3): + * younger than the badge TTL → grace period (re-check next cycle); stale + + * bot delivery found in the comment stream → completed; stale without + * delivery → reset to idle with a [system] notice. Exported for unit testing. + */ +export function decideBadgeAction(badgeAgeMs: number | null, badgeTtlMs: number, delivered: boolean): BadgeAction { + if (badgeAgeMs !== null && badgeAgeMs < badgeTtlMs) return "grace"; + return delivered ? "complete" : "reset"; +} + +/** Marker so the reset notice is posted at most once per badge generation. */ +export const BADGE_RESET_MARKER = ""; + +export function buildBadgeResetNotice(): string { + return `[system] 🏷 ⚠️ 状态已复位:AI 处理徽标长时间无响应(进程已死且无输出/模型流量),已重置为空闲。如需继续处理请回复本 issue。${BADGE_RESET_MARKER}`; +} + +export interface BadgeSweepEntry { + owner: string; + repo: string; + number: number; + aiStatus: string; + since: number | null; + /** Whether this daemon has an issue row for the badge (false = orphan). */ + engineRecord: boolean; + ownedByMe: boolean; + sessions: number; + pidAlive: boolean; + outputAgeSec: number | null; + modelAgeSec: number | null; + verdict: "alive" | "stale" | "orphan" | "skipped"; + action?: string; +} + /** * Query the opencode SQLite DB for a session's assistant-message output tokens. * Returns `{hasOutput: true}` (safe default) when the DB can't be opened or @@ -437,6 +512,8 @@ export interface EngineOptions { replyBurst?: { max: number; windowMs: number }; /** Set false in tests to drive recover() explicitly instead of fire-and-forget from the constructor. */ recoverOnBoot?: boolean; + /** Set false in tests that drive sweepStuckBadges() manually. */ + badgeSweep?: boolean; } function createDefaultBackend(cfg: Config): RuntimeBackend { @@ -587,6 +664,14 @@ export class Engine { private lastWorkdirGcAt = 0; private badgeWrites = new Map(); private observerTimer?: ReturnType; + // Stuck-badge sweep state (ework#12) + private badgeSweepTimer?: ReturnType; + private badgeSweepRunning = false; + private lastModelAt = new Map(); + private lastModelTokens = new Map(); + private lastHeartbeatRefreshAt = new Map(); + private lastBadgeSweepAt = 0; + private lastBadgeSweepEntries: BadgeSweepEntry[] = []; private groupConfigs = new Map(); private cloneUrls = new Map(); @@ -632,6 +717,7 @@ export class Engine { this.maxConcurrent = cfg.work.maxConcurrent; this.maxConcurrentExplicit = cfg.work.maxConcurrentExplicit; this.startGlobalObserver(); + if (opts.badgeSweep !== false) this.startBadgeSweep(); if (opts.recoverOnBoot !== false) void this.recover(); } @@ -1238,6 +1324,9 @@ export class Engine { () => {}, ); void tracker.updateStatus(ref, "processing"); + // Every processing write arms the badge heartbeat (ework#12 item 3): if no + // subsequent activity refreshes it, the sweep treats the badge as stale. + void this.refreshHeartbeat(session.id); const workdir = await this.resolveWorkdir(session, issue); log.info(`engine: session "${session.name}" created for ${k}, workdir=${workdir}`); @@ -1807,7 +1896,15 @@ export class Engine { env: childEnv, }, { - onOutput: () => { this.lastOutputAt.set(k, Date.now()); }, + onOutput: () => { + const n = Date.now(); + this.lastOutputAt.set(k, n); + // Throttled heartbeat refresh so long chatty runs never expire. + if (n - (this.lastHeartbeatRefreshAt.get(k) ?? 0) > 30_000) { + this.lastHeartbeatRefreshAt.set(k, n); + void this.refreshHeartbeat(session.id); + } + }, onSessionId: async (id: string) => { if (!session.opencodeSessionId) { await this.store.updateSession(session.id, { opencodeSessionId: id }); @@ -1848,6 +1945,8 @@ export class Engine { await this.store.updateSession(session.id, { opencodePid: handle.pid }); await this.persistRuntimeState(session.id); + this.lastHeartbeatRefreshAt.set(k, Date.now()); + void this.refreshHeartbeat(session.id); log.info(`engine: spawned pid=${handle.pid} for ${k} (backend=${backend.name})`); @@ -2317,6 +2416,7 @@ export class Engine { this.running.add(k); this.currentMessage.set(k, msg.id); await this.store.updateSession(session.id, { state: "running" }); + void this.refreshHeartbeat(session.id); void this.execProcess(k, session, issue, msg); } @@ -2614,7 +2714,12 @@ export class Engine { const ref = { trackerType: issue.trackerType, scope: { owner: scopeParts[0]!, repo: scopeParts[1]! }, issueId: String(issue.trackerIssueId) }; // cache only after success — a failed write must retry next cycle, not be skipped forever await this.getTracker(issue.trackerType).updateStatus(ref, desired).then( - () => { this.badgeWrites.set(issue.id, desired); }, + () => { + this.badgeWrites.set(issue.id, desired); + if (desired === "processing") { + for (const s of sessions) if (s.state === "running") void this.refreshHeartbeat(s.id); + } + }, () => { /* best-effort; retried next cycle */ }, ); } catch { /* badge convergence is best-effort */ } @@ -3225,10 +3330,218 @@ export class Engine { return !!proc; } + // ─── Stuck-Badge Sweep (ework#12) ─── + + private startBadgeSweep() { + this.badgeSweepTimer = setInterval(() => { void this.sweepStuckBadges(); }, this.cfg.work.badgeSweepIntervalMs); + } + + private allTrackers(): IssueTracker[] { + const reg = this.trackers; + if (reg instanceof Map) return [...reg.values()]; + return [reg.get("gitea")].filter((t): t is IssueTracker => !!t); + } + + private pidAlive(pid: number): boolean { + try { process.kill(pid, 0); return true; } catch (err) { + return (err as NodeJS.ErrnoException).code === "EPERM"; + } + } + + /** Arm/refresh a session's badge-heartbeat deadline (spec item 3). Best-effort. */ + private async refreshHeartbeat(sessionId: string): Promise { + try { await this.store.setExpectedHeartbeat(sessionId, Date.now() + this.cfg.work.badgeTtlMs); } catch { /* best-effort */ } + } + + private async fetchWebAiStatus(owner: string, repo: string, number: number): Promise { + const url = `${this.cfg.gitea.url}/api/v1/dispatch-state?owner=${encodeURIComponent(owner)}&repo=${encodeURIComponent(repo)}&number=${number}`; + const resp = await fetch(url, { signal: AbortSignal.timeout(5000), headers: { Authorization: `token ${this.cfg.gitea.token}` } }); + if (!resp.ok) throw new Error(`web returned ${resp.status}`); + const data = await resp.json() as { aiStatus?: string }; + return data.aiStatus ?? ""; + } + + /** + * Continuous stuck-badge detector (ework#12 spec items 2+3+4). Enumerates every + * fleet-wide "processing" badge through the web listing endpoint. Badges whose + * issue this daemon knows about are verified with three signals (pid alive / + * recent output / recent model traffic + heartbeat TTL); orphan badges — the + * engine has ZERO records — go straight to arbitration: bot delivery found in + * the comment stream → completed, otherwise → idle + [system] notice. Every + * action is logged and recorded for GET /api/badges. Public so tests can drive + * it directly. + */ + async sweepStuckBadges(): Promise { + if (this.destroyed || this.badgeSweepRunning) return; + this.badgeSweepRunning = true; + try { + try { await this.store.releaseDeadOwners(this.cfg.work.leaseTtlMs); } catch { /* lease reaping retried next cycle */ } + const entries: BadgeSweepEntry[] = []; + const seen = new Set(); + for (const tracker of this.allTrackers()) { + if (!tracker.listBadges) continue; + let badges: BadgeEntry[]; + try { badges = await tracker.listBadges("processing"); } catch { continue; } + for (const b of badges) { + const dedupeKey = `${b.owner}/${b.repo}#${b.number}`; + if (seen.has(dedupeKey)) continue; + seen.add(dedupeKey); + try { + entries.push(await this.sweepOneBadge(tracker, b)); + } catch (err) { + log.warn(`engine: badge-sweep failed for ${dedupeKey}: ${(err as Error).message}`); + } + } + } + this.lastBadgeSweepAt = Date.now(); + this.lastBadgeSweepEntries = entries.slice(0, 200); + const acted = entries.filter((e) => e.action && !e.action.startsWith("grace") && e.action !== "owned-by-other-daemon"); + if (acted.length > 0) { + log.info(`engine: badge-sweep: ${entries.length} badge(s) checked, actions: ${acted.map((e) => `${e.owner}/${e.repo}#${e.number}→${e.action}`).join("; ")}`); + } + } finally { + this.badgeSweepRunning = false; + } + } + + private async sweepOneBadge(tracker: IssueTracker, b: BadgeEntry): Promise { + const now = Date.now(); + const ref: TrackerRef = { trackerType: tracker.type, scope: { owner: b.owner, repo: b.repo }, issueId: String(b.number) }; + const entry: BadgeSweepEntry = { + owner: b.owner, repo: b.repo, number: b.number, aiStatus: b.aiStatus, since: b.since, + engineRecord: false, ownedByMe: false, sessions: 0, + pidAlive: false, outputAgeSec: null, modelAgeSec: null, verdict: "skipped", + }; + + let issue: Issue | undefined; + try { issue = await this.store.findIssue(tracker.type, `${b.owner}/${b.repo}`, String(b.number)); } catch { /* store blip — skip this round */ } + if (issue && issue.ownerDaemonId !== null && issue.ownerDaemonId !== this.daemonId) { + entry.action = "owned-by-other-daemon"; + log.info(`engine: badge-sweep skip ${b.owner}/${b.repo}#${b.number}: owned by daemon #${issue.ownerDaemonId}`); + return entry; + } + entry.engineRecord = !!issue; + entry.ownedByMe = !!issue && (issue.ownerDaemonId === null || issue.ownerDaemonId === this.daemonId); + + let sessions: OpSession[] = []; + if (issue) sessions = await this.store.getSessionsForIssue(issue.id).catch(() => [] as OpSession[]); + entry.sessions = sessions.length; + + // Three-signal verification over every engine-known session + let anyLive = false; + for (const s of sessions) { + const k = this.sessionKey(s, issue!); + const handle = this.processes.get(k); + const pid = handle?.pid ?? s.opencodePid; + const pidAlive = pid != null ? this.pidAlive(pid) : false; + const outTs = this.lastOutputAt.get(k) ?? s.lastOutputAt; + const outputAgeMs = outTs != null ? now - outTs : null; + let modelAgeMs = this.lastModelAt.has(k) ? now - this.lastModelAt.get(k)! : null; + // Model-traffic probe only when the other two signals are both dead (cheap-first) + if (!pidAlive && (outputAgeMs === null || outputAgeMs >= this.cfg.work.badgeOutputTtlMs) && s.opencodeSessionId) { + const backend = this.backendFor(k, s.opencodeSessionId); + try { + const res = await backend.getSessionOutputTokens(s.opencodeSessionId); + const prev = this.lastModelTokens.get(k); + if (prev !== undefined && res.tokenCount > prev) this.lastModelAt.set(k, now); + this.lastModelTokens.set(k, res.tokenCount); + modelAgeMs = this.lastModelAt.has(k) ? now - this.lastModelAt.get(k)! : null; + } catch { /* probe failed — signal stays unknown */ } + } + const heartbeatExpired = s.state === "running" && s.expectedHeartbeatAt != null && now > s.expectedHeartbeatAt; + const verdict = evaluateBadgeSignals( + { pidAlive, outputAgeMs, modelAgeMs }, + { outputTtlMs: this.cfg.work.badgeOutputTtlMs, modelTtlMs: this.cfg.work.badgeModelTtlMs }, + heartbeatExpired, + ); + if (verdict.alive) { + anyLive = true; + entry.pidAlive = true; + void this.refreshHeartbeat(s.id); + } + if (outputAgeMs !== null) { + const sec = Math.round(outputAgeMs / 1000); + entry.outputAgeSec = entry.outputAgeSec === null ? sec : Math.min(entry.outputAgeSec, sec); + } + if (modelAgeMs !== null) { + const sec = Math.round(modelAgeMs / 1000); + entry.modelAgeSec = entry.modelAgeSec === null ? sec : Math.min(entry.modelAgeSec, sec); + } + } + + if (sessions.length > 0 && anyLive) { + entry.verdict = "alive"; + return entry; + } + + // Stale (all known sessions dead) or orphan (zero engine records) + entry.verdict = sessions.length === 0 ? "orphan" : "stale"; + const age = b.since !== null ? now - b.since : null; + if (decideBadgeAction(age, this.cfg.work.badgeTtlMs, false) === "grace") { + entry.action = "grace period (badge younger than ttl)"; + return entry; + } + + // Re-verify the web still shows processing — never clobber a concurrent fixup + let current = "processing"; + try { current = await this.fetchWebAiStatus(b.owner, b.repo, b.number); } catch { current = ""; } + if (current !== "processing") { + entry.action = `web already shows "${current}" — no-op`; + return entry; + } + + // Arbitration: the web comment stream is the truth source for delivery + let comments: TrackerComment[]; + try { comments = await tracker.listComments(ref); } catch (err) { + entry.action = `listComments failed: ${(err as Error).message}`; + return entry; + } + const promptTime = await this.arbitrationPromptTime(issue, sessions, b.since); + const delivered = hasRecoveryDelivery(comments, (a) => tracker.isBotUser(a), promptTime); + const decision = decideBadgeAction(age, this.cfg.work.badgeTtlMs, delivered); + + if (decision === "complete") { + try { await tracker.updateStatus(ref, "completed"); } catch (err) { entry.action = `updateStatus(completed) failed: ${(err as Error).message}`; return entry; } + entry.action = "flipped→completed (bot delivery found)"; + log.info(`engine: badge-sweep ${b.owner}/${b.repo}#${b.number} flipped to completed: bot delivery found in comment stream`); + } else { + try { await tracker.updateStatus(ref, ""); } catch (err) { entry.action = `updateStatus(idle) failed: ${(err as Error).message}`; return entry; } + const alreadyNotified = comments.some((c) => c.body.includes(BADGE_RESET_MARKER)); + if (!alreadyNotified) void tracker.createComment(ref, buildBadgeResetNotice()).catch(() => {}); + entry.action = alreadyNotified ? "flipped→idle (notice already posted)" : "flipped→idle + [system] notice (no delivery)"; + log.warn(`engine: badge-sweep ${b.owner}/${b.repo}#${b.number} reset to idle: ${sessions.length === 0 ? "orphan badge, no engine record" : "all sessions dead"} and no bot delivery`); + } + return entry; + } + + /** Delivery-arbitration baseline: oldest non-terminal message (same rule as + * restart recovery), falling back to the badge write time, then epoch 0. */ + private async arbitrationPromptTime(issue: Issue | undefined, sessions: OpSession[], sinceMs: number | null): Promise { + let oldestMs: number | null = null; + if (issue) { + for (const s of sessions) { + const msgs = await this.store.getMessagesForSession(s.id).catch(() => [] as Message[]); + for (const m of msgs) { + if (m.status !== "pending" && m.status !== "running" && m.status !== "interrupted") continue; + const t = new Date(m.createdAt).getTime(); + if (!Number.isNaN(t) && (oldestMs === null || t < oldestMs)) oldestMs = t; + } + } + } + if (oldestMs !== null) return new Date(oldestMs); + if (sinceMs !== null) return new Date(sinceMs); + return new Date(0); + } + + getBadgeSweepState(): { checkedAt: number; intervalMs: number; entries: BadgeSweepEntry[] } { + return { checkedAt: this.lastBadgeSweepAt, intervalMs: this.cfg.work.badgeSweepIntervalMs, entries: this.lastBadgeSweepEntries }; + } + destroy() { this.destroyed = true; this.stopHeartbeat(); if (this.observerTimer) clearInterval(this.observerTimer); + if (this.badgeSweepTimer) clearInterval(this.badgeSweepTimer); this.observedIssues.clear(); for (const [, proc] of this.processes) { try { process.kill(proc.pid, "SIGKILL"); } catch { /* dead */ } @@ -3244,6 +3557,9 @@ export class Engine { this.stuckNudgeRounds.clear(); this.currentPrompt.clear(); this.generation.clear(); + this.lastModelAt.clear(); + this.lastModelTokens.clear(); + this.lastHeartbeatRefreshAt.clear(); this.groupConfigs.clear(); this.cloneUrls.clear(); this.senders.clear(); diff --git a/packages/daemon/src/server.ts b/packages/daemon/src/server.ts index 773fc1f..28c7ccb 100644 --- a/packages/daemon/src/server.ts +++ b/packages/daemon/src/server.ts @@ -120,6 +120,9 @@ export function createServer( if (pathname === "/api/processes") { return json(engine.getProcesses()); } + if (pathname === "/api/badges") { + return json(engine.getBadgeSweepState()); + } const sessionMsgsMatch = pathname.match(/^\/api\/sessions\/([0-9a-f-]+)\/messages$/); if (sessionMsgsMatch) { diff --git a/packages/daemon/src/trackers/gitea-tracker.ts b/packages/daemon/src/trackers/gitea-tracker.ts index c0a924a..ffea72a 100644 --- a/packages/daemon/src/trackers/gitea-tracker.ts +++ b/packages/daemon/src/trackers/gitea-tracker.ts @@ -6,6 +6,7 @@ import type { TrackerEvent, TrackerComment, TrackerInstructions, + BadgeEntry, } from "./types"; export class GiteaTracker implements IssueTracker { @@ -71,6 +72,22 @@ export class GiteaTracker implements IssueTracker { ); } + async listBadges(status: string): Promise { + try { + const res = await this.client.listAiStatusBadges(status); + return (res.badges ?? []).map((b) => ({ + owner: b.owner, + repo: b.repo, + number: b.number, + aiStatus: b.aiStatus, + since: b.since ? Date.parse(b.since) : null, + })).filter((b) => b.since === null || !Number.isNaN(b.since)); + } catch (e) { + console.warn(`[tracker] listBadges(${status}) failed:`, (e as Error).message); + return []; + } + } + async setCommentModel(ref: TrackerRef, commentId: string, model: string) { await this.client.setCommentModel( this.owner(ref), this.repo(ref), Number(commentId), model diff --git a/packages/daemon/src/trackers/types.ts b/packages/daemon/src/trackers/types.ts index f160e36..52bf2b1 100644 --- a/packages/daemon/src/trackers/types.ts +++ b/packages/daemon/src/trackers/types.ts @@ -108,6 +108,19 @@ export interface OpSession { nudgeRounds?: number; stuckNudgeRounds?: number; generation?: number; + // Processing-badge TTL: instant by which this session must refresh its + // heartbeat (spawn / output / model traffic). Past due = badge stale. + expectedHeartbeatAt?: number; +} + +/** A tracker-side issue currently carrying a given ai_status badge. */ +export interface BadgeEntry { + owner: string; + repo: string; + number: number; + aiStatus: string; + /** ISO timestamp of the status change, or null when unknown. */ + since: number | null; } /** Message = a prompt enqueued for a session */ @@ -147,6 +160,12 @@ export interface IssueTracker { listComments(ref: TrackerRef): Promise; closeIssue(ref: TrackerRef): Promise; updateStatus(ref: TrackerRef, status: string, detail?: string): Promise; + /** + * Enumerate every issue fleet-wide carrying the given ai_status. + * Optional: backends without a global listing endpoint return undefined and + * the stuck-badge sweep degrades to per-issue verification only. + */ + listBadges?(status: string): Promise; setReaction(ref: TrackerRef, commentId: string, content: string, remove?: boolean): Promise; setCommentModel(ref: TrackerRef, commentId: string, model: string): Promise; diff --git a/packages/daemon/tests/stuck-badge.test.ts b/packages/daemon/tests/stuck-badge.test.ts new file mode 100644 index 0000000..bbfd935 --- /dev/null +++ b/packages/daemon/tests/stuck-badge.test.ts @@ -0,0 +1,472 @@ +import { beforeAll, beforeEach, afterEach, describe, it, expect } from "bun:test"; +import { mkdirSync, rmSync } from "fs"; +import { join } from "path"; +import { tmpdir } from "os"; +import { Store } from "../src/op"; +import { + Engine, + type TakeoverStrategy, + evaluateBadgeSignals, + decideBadgeAction, + BADGE_RESET_MARKER, + buildBadgeResetNotice, +} from "../src/opencode"; +import { initDB, getDB } from "../src/db"; +import { loadConfig, type Config } from "../src/config"; +import type { + IssueTracker, + TrackerRef, + TrackerEvent, + TrackerComment, + TrackerInstructions, + Issue, + OpSession, + BadgeEntry, +} from "../src/trackers/types"; +import type { + RuntimeBackend, + RuntimeHandle, + RuntimeSpawnOpts, + RuntimeSpawnCallbacks, + LastModelResult, + SessionOutputResult, +} from "../src/runtime/types"; + +// Stuck-badge sweep harness (ework#12). Models the web side with an in-memory +// badge table + a local mock of GET /api/v1/dispatch-state, then drives the +// engine's sweep both manually and through its real 1-min-class timer. + +const LEASE_TTL_MS = 500; +const HEARTBEAT_MS = 100; +const BOT_USER = "ework-daemon"; +const HUMAN_USER = "tester"; +// Tiny TTLs so every scenario settles in well under the spec's 2-minute bound. +const SWEEP_INTERVAL_MS = 300; +const BADGE_TTL_MS = 500; +const OUTPUT_TTL_MS = 400; +const MODEL_TTL_MS = 400; + +interface BadgeRecord { status: string; since: number } +type BadgeTable = Map; // key: "owner/repo#number" + +let workdirBase: string; +let badgeTable: BadgeTable; +let fakeTracker: FakeTracker; +let trackerRegistry: Map; +let mockWeb: MockWeb | null = null; +const liveEngines: Engine[] = []; + +beforeAll(async () => { + await initDB(); +}); + +beforeEach(async () => { + const db = getDB(); + const mysql = db.dialect === "mysql"; + await db.exec(mysql ? "SET FOREIGN_KEY_CHECKS = 0" : "PRAGMA foreign_keys = OFF"); + for (const t of ["messages", "op_sessions", "issues", "daemons"]) { + await db.exec(`DELETE FROM {{${t}}}`); + } + await db.exec(mysql ? "SET FOREIGN_KEY_CHECKS = 1" : "PRAGMA foreign_keys = ON"); + + workdirBase = `${tmpdir()}/ework-daemon-badge-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`; + mkdirSync(workdirBase, { recursive: true }); + + badgeTable = new Map(); + fakeTracker = new FakeTracker(badgeTable); + trackerRegistry = new Map(); + trackerRegistry.set(fakeTracker.type, fakeTracker); +}); + +afterEach(async () => { + for (const e of liveEngines) { + try { e.destroy(); } catch { /* already destroyed */ } + } + liveEngines.length = 0; + mockWeb?.stop(); + mockWeb = null; + try { rmSync(workdirBase, { recursive: true, force: true }); } catch { /* gone */ } +}); + +/** Web-side badge table standing in for issues.ai_status (+ai_status_since). */ +class FakeTracker implements IssueTracker { + readonly type = "gitea"; + comments: TrackerComment[] = []; + statuses: Array<{ key: string; status: string }> = []; + private nextId = 1; + constructor(private badges: BadgeTable) {} + + formatScopeKey(scope: Record): string { + return `${scope.owner}/${scope.repo}`; + } + + private key(ref: TrackerRef): string { + return `${ref.scope["owner"]}/${ref.scope["repo"]}#${ref.issueId}`; + } + + async createComment(_ref: TrackerRef, body: string): Promise<{ id: string }> { + const id = `c${this.nextId++}`; + this.comments.push({ id, body, author: BOT_USER, createdAt: new Date().toISOString() }); + return { id }; + } + + async editComment(): Promise {} + async deleteComment(): Promise {} + async listComments(): Promise { return [...this.comments]; } + async closeIssue(): Promise {} + + async updateStatus(ref: TrackerRef, status: string): Promise { + const k = this.key(ref); + this.statuses.push({ key: k, status }); + if (status === "") this.badges.delete(k); + else this.badges.set(k, { status, since: Date.now() }); + } + + async listBadges(status: string): Promise { + const out: BadgeEntry[] = []; + for (const [k, rec] of this.badges) { + if (rec.status !== status) continue; + const [owner, rest] = k.split("/"); + const [repo, number] = rest!.split("#"); + out.push({ owner: owner!, repo: repo!, number: Number(number), aiStatus: rec.status, since: rec.since }); + } + return out; + } + + async setCommentModel(): Promise {} + async setReaction(): Promise {} + + getTrackerInstructions(_ref: TrackerRef): TrackerInstructions { + return { clone: "git clone fake", issueRef: "fake/ref" }; + } + + verifyWebhookSignature(): boolean { return true; } + parseWebhookEvent(): TrackerEvent | null { return null; } + isBotUser(author: string): boolean { return author === BOT_USER; } +} + +/** Minimal mock of the web machine endpoints the sweep depends on. */ +class MockWeb { + private server: ReturnType; + constructor(private badges: BadgeTable) { + this.server = Bun.serve({ + port: 0, + fetch: (req) => { + const url = new URL(req.url); + if (url.pathname === "/api/v1/dispatch-state") { + const key = `${url.searchParams.get("owner")}/${url.searchParams.get("repo")}#${url.searchParams.get("number")}`; + return Response.json({ dispatchOff: false, aiStatus: this.badges.get(key)?.status ?? "" }); + } + return Response.json({ error: "not found" }, 404); + }, + }); + } + get url(): string { return `http://127.0.0.1:${this.server.port}`; } + stop(): void { this.server.stop(true); } +} + +class TestTakeoverStrategy implements TakeoverStrategy { + constructor(private baseWorkdir: string) {} + async acquireWorkdir(session: OpSession, issue: Issue): Promise { + const dir = join(this.baseWorkdir, String(issue.trackerIssueId), session.name); + mkdirSync(dir, { recursive: true }); + return dir; + } + async resumeOpenCodeSession(): Promise { return null; } +} + +class FakeBackend implements RuntimeBackend { + readonly name = "fake"; + spawns = 0; + constructor(private comments: TrackerComment[]) {} + + async spawn(opts: RuntimeSpawnOpts, callbacks: RuntimeSpawnCallbacks): Promise { + this.spawns++; + const proc = Bun.spawn(["sh", "-c", "sleep 30"], { stdout: "ignore", stderr: "ignore" }); + callbacks.onOutput(); + const exited = proc.exited.then(() => 0); + return { + pid: proc.pid, + exited, + stderrText: Promise.resolve(""), + stderrPartial: () => "", + stderrCancel: () => {}, + }; + } + + async lastSessionModel(): Promise { return { model: "fake/model" }; } + async sessionExists(): Promise { return true; } + async getSessionOutputTokens(): Promise { return { hasOutput: true, tokenCount: 5 }; } +} + +function makeConfig(webUrl: string): Config { + const cfg = loadConfig(); + return { + ...cfg, + gitea: { ...cfg.gitea, url: webUrl }, + opencode: { ...cfg.opencode, binary: "fake-opencode", baseWorkdir: workdirBase, dbPath: join(workdirBase, "opencode.db") }, + work: { + ...cfg.work, + capacity: 4, + maxConcurrent: 4, + maxConcurrentExplicit: false, + heartbeatMs: HEARTBEAT_MS, + leaseTtlMs: LEASE_TTL_MS, + badgeSweepIntervalMs: SWEEP_INTERVAL_MS, + badgeTtlMs: BADGE_TTL_MS, + badgeOutputTtlMs: OUTPUT_TTL_MS, + badgeModelTtlMs: MODEL_TTL_MS, + }, + }; +} + +async function bootEngine(name: string, store: Store, cfg: Config, port: number, backend: FakeBackend, badgeSweep = false): Promise<{ engine: Engine; daemonId: number }> { + await store.releaseDeadOwners(cfg.work.leaseTtlMs); + const daemonId = await store.registerDaemon(`host-${name}`, `127.0.0.1:${port}`, cfg.work.capacity, cfg.work.leaseTtlMs); + await store.claimAllOwnerless(daemonId); + const engine = new Engine(cfg, store, trackerRegistry, { + daemonId, + gateChecker: async () => ({ allowed: true, reason: "test" }), + takeover: new TestTakeoverStrategy(workdirBase), + backend, + recoverOnBoot: false, + badgeSweep, + }); + engine.startHeartbeat(cfg.work.heartbeatMs); + liveEngines.push(engine); + return { engine, daemonId }; +} + +function commentEvent(issueId: string, body: string): TrackerEvent { + return { + type: "comment_created", + ref: { trackerType: "gitea", scope: { owner: "dog", repo: "repo" }, issueId }, + issue: { title: "stuck badge test", body: "b", state: "open", author: HUMAN_USER }, + comment: { id: `wc-${Date.now()}`, body, author: HUMAN_USER }, + }; +} + +const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); + +async function waitFor(fn: () => T | undefined | Promise, timeoutMs = 15_000): Promise { + const start = Date.now(); + for (;;) { + const result = await fn(); + if (result !== undefined) return result; + if (Date.now() - start >= timeoutMs) break; + await sleep(50); + } + return await fn(); +} + +async function messageRows(): Promise> { + const db = getDB(); + return db.all<{ status: string }>(`SELECT status FROM {{messages}} ORDER BY created_at ASC`); +} + +function badgeKey(issueId: string): string { return `dog/repo#${issueId}`; } + +/** Age a badge past its TTL as if the process that wrote it vanished. */ +function ageBadge(issueId: string): void { + const rec = badgeTable.get(badgeKey(issueId)); + if (rec) rec.since = Date.now() - BADGE_TTL_MS * 2; +} + +async function getRunningPid(): Promise { + const db = getDB(); + const row = await db.get<{ opencode_pid: number | null }>(`SELECT opencode_pid FROM {{op_sessions}} LIMIT 1`); + return row?.opencode_pid ?? null; +} + +describe("badge signal helpers (unit)", () => { + it("evaluateBadgeSignals: dead pid + no output + no model → stuck", () => { + const v = evaluateBadgeSignals({ pidAlive: false, outputAgeMs: 999_999, modelAgeMs: null }, { outputTtlMs: OUTPUT_TTL_MS, modelTtlMs: MODEL_TTL_MS }, false); + expect(v.alive).toBe(false); + expect(v.reasons).toEqual(expect.arrayContaining(["pid-dead", "no-recent-output", "no-recent-model"])); + }); + + it("evaluateBadgeSignals: live pid alone keeps the badge", () => { + const v = evaluateBadgeSignals({ pidAlive: true, outputAgeMs: null, modelAgeMs: null }, { outputTtlMs: OUTPUT_TTL_MS, modelTtlMs: MODEL_TTL_MS }, false); + expect(v.alive).toBe(true); + }); + + it("evaluateBadgeSignals: fresh output alone keeps the badge", () => { + const v = evaluateBadgeSignals({ pidAlive: false, outputAgeMs: 1000, modelAgeMs: null }, { outputTtlMs: OUTPUT_TTL_MS * 10, modelTtlMs: MODEL_TTL_MS }, false); + expect(v.alive).toBe(true); + }); + + it("evaluateBadgeSignals: expired heartbeat forces stale even with a live pid", () => { + const v = evaluateBadgeSignals({ pidAlive: true, outputAgeMs: 1000, modelAgeMs: 1000 }, { outputTtlMs: OUTPUT_TTL_MS * 10, modelTtlMs: MODEL_TTL_MS * 10 }, true); + expect(v.alive).toBe(false); + expect(v.reasons).toContain("heartbeat-expired"); + }); + + it("decideBadgeAction: young badge gets grace; stale decides by delivery", () => { + expect(decideBadgeAction(BADGE_TTL_MS / 2, BADGE_TTL_MS, false)).toBe("grace"); + expect(decideBadgeAction(BADGE_TTL_MS * 2, BADGE_TTL_MS, true)).toBe("complete"); + expect(decideBadgeAction(BADGE_TTL_MS * 2, BADGE_TTL_MS, false)).toBe("reset"); + // Unknown since → never grace (safe direction: arbitrate). + expect(decideBadgeAction(null, BADGE_TTL_MS, false)).toBe("reset"); + }); + + it("reset notice carries the dedup marker", () => { + expect(buildBadgeResetNotice()).toContain(BADGE_RESET_MARKER); + expect(buildBadgeResetNotice().startsWith("[system]")).toBe(true); + }); +}); + +describe("stuck-badge sweep (integration)", () => { + it("flips a dead-and-record-wiped processing badge to idle + [system] notice (acceptance)", async () => { + mockWeb = new MockWeb(badgeTable); + const store = new Store(); + const cfg = makeConfig(mockWeb.url); + const backend = new FakeBackend(fakeTracker.comments); + const { engine: engineA } = await bootEngine("A", store, cfg, 7421, backend); + + await engineA.handleEvent(commentEvent("911", "please fix")); + expect(await waitFor(async () => { + const r = await messageRows(); + return r.length === 1 && r[0]?.status === "running" ? true : undefined; + })).toBe(true); + expect(badgeTable.get(badgeKey("911"))?.status).toBe("processing"); + + // Hard crash mid-run (kills the child), then wipe the engine's records. + engineA.destroy(); + await sleep(LEASE_TTL_MS + 150); + const db = getDB(); + await db.exec(`DELETE FROM {{op_sessions}}`); + ageBadge("911"); + + const { engine: engineB } = await bootEngine("B", store, cfg, 7422, backend); + await engineB.sweepStuckBadges(); + + const flips = fakeTracker.statuses.filter((s) => s.key === badgeKey("911")); + expect(flips.some((s) => s.status === "")).toBe(true); + expect(flips.some((s) => s.status === "completed")).toBe(false); + const notices = fakeTracker.comments.filter((c) => c.body.includes(BADGE_RESET_MARKER)); + expect(notices.length).toBe(1); + // Idempotent: a second sweep must not re-notify (badge already gone). + await engineB.sweepStuckBadges(); + expect(fakeTracker.comments.filter((c) => c.body.includes(BADGE_RESET_MARKER)).length).toBe(1); + }); + + it("keeps a badge whose session process is genuinely alive", async () => { + mockWeb = new MockWeb(badgeTable); + const store = new Store(); + const cfg = makeConfig(mockWeb.url); + const backend = new FakeBackend(fakeTracker.comments); + const { engine } = await bootEngine("A", store, cfg, 7423, backend); + + await engine.handleEvent(commentEvent("912", "please fix")); + expect(await waitFor(async () => { + const r = await messageRows(); + return r.length === 1 && r[0]?.status === "running" ? true : undefined; + })).toBe(true); + ageBadge("912"); // badge is old, but the child is still sleeping + + await engine.sweepStuckBadges(); + + const flips = fakeTracker.statuses.filter((s) => s.key === badgeKey("912")); + expect(flips.every((s) => s.status === "processing")).toBe(true); + expect(fakeTracker.comments.some((c) => c.body.includes(BADGE_RESET_MARKER))).toBe(false); + // Live sessions get their heartbeat refreshed, not clobbered. + const hb = await getDB().get<{ expected_heartbeat_at: number | null }>(`SELECT expected_heartbeat_at FROM {{op_sessions}} LIMIT 1`); + expect(hb?.expected_heartbeat_at ?? 0).toBeGreaterThan(Date.now()); + }); + + it("completes an orphan badge when the comment stream shows a bot delivery", async () => { + mockWeb = new MockWeb(badgeTable); + const store = new Store(); + const cfg = makeConfig(mockWeb.url); + const backend = new FakeBackend(fakeTracker.comments); + const { engine } = await bootEngine("A", store, cfg, 7424, backend); + + // Pure orphan: the engine has NO issue/session/message rows for this badge. + badgeTable.set(badgeKey("913"), { status: "processing", since: Date.now() - BADGE_TTL_MS * 2 }); + fakeTracker.comments.push({ + id: "bot-delivery-913", + body: "[bot] 🏷 已完成,修复提交在 PR #12,测试全部通过。", + author: BOT_USER, + createdAt: new Date(Date.now() - BADGE_TTL_MS).toISOString(), + }); + + await engine.sweepStuckBadges(); + + const flips = fakeTracker.statuses.filter((s) => s.key === badgeKey("913")); + expect(flips.some((s) => s.status === "completed")).toBe(true); + expect(fakeTracker.comments.some((c) => c.body.includes(BADGE_RESET_MARKER))).toBe(false); + }); + + it("resets an orphan badge without any delivery and notifies once", async () => { + mockWeb = new MockWeb(badgeTable); + const store = new Store(); + const cfg = makeConfig(mockWeb.url); + const backend = new FakeBackend(fakeTracker.comments); + const { engine } = await bootEngine("A", store, cfg, 7425, backend); + + badgeTable.set(badgeKey("914"), { status: "processing", since: Date.now() - BADGE_TTL_MS * 2 }); + fakeTracker.comments.push({ id: "human-914", body: "hello?", author: HUMAN_USER, createdAt: new Date().toISOString() }); + + await engine.sweepStuckBadges(); + + const flips = fakeTracker.statuses.filter((s) => s.key === badgeKey("914")); + expect(flips.some((s) => s.status === "")).toBe(true); + expect(fakeTracker.comments.filter((c) => c.body.includes(BADGE_RESET_MARKER)).length).toBe(1); + }); + + it("skips badges owned by another live daemon", async () => { + mockWeb = new MockWeb(badgeTable); + const store = new Store(); + const cfg = makeConfig(mockWeb.url); + const backend = new FakeBackend(fakeTracker.comments); + const { engine } = await bootEngine("A", store, cfg, 7426, backend); + + await engine.handleEvent(commentEvent("915", "please fix")); + expect(await waitFor(async () => { + const r = await messageRows(); + return r.length === 1 && r[0]?.status === "running" ? true : undefined; + })).toBe(true); + + // Pretend a different (live, registered) daemon owns this issue. + const otherDaemon = await store.registerDaemon("host-other", "127.0.0.1:7499", 1, cfg.work.leaseTtlMs); + await getDB().run(`UPDATE {{issues}} SET owner_daemon_id = ? WHERE tracker_issue_id = '915'`, [otherDaemon]); + ageBadge("915"); + + await engine.sweepStuckBadges(); + + const flips = fakeTracker.statuses.filter((s) => s.key === badgeKey("915")); + expect(flips.every((s) => s.status === "processing")).toBe(true); + const state = engine.getBadgeSweepState(); + const entry = state.entries.find((e) => e.number === 915); + expect(entry?.verdict).toBe("skipped"); + }); + + it("auto-flips within the timer loop without manual driving (spec: < 2 min)", async () => { + mockWeb = new MockWeb(badgeTable); + const store = new Store(); + const cfg = makeConfig(mockWeb.url); + const backend = new FakeBackend(fakeTracker.comments); + const { engine: engineA } = await bootEngine("A", store, cfg, 7427, backend, true); + + await engineA.handleEvent(commentEvent("916", "please fix")); + expect(await waitFor(async () => { + const r = await messageRows(); + return r.length === 1 && r[0]?.status === "running" ? true : undefined; + })).toBe(true); + + // Crash + record wipe, exactly like the weekly recurrence class. + engineA.destroy(); + await sleep(LEASE_TTL_MS + 150); + await getDB().exec(`DELETE FROM {{op_sessions}}`); + ageBadge("916"); + + const { engine: engineB } = await bootEngine("B", store, cfg, 7428, backend, true); + const started = Date.now(); + const flipped = await waitFor(() => + fakeTracker.statuses.some((s) => s.key === badgeKey("916") && s.status === "") ? true : undefined, + ); + expect(flipped).toBe(true); + expect(Date.now() - started).toBeLessThan(2 * 60_000); + void engineB; + }); +}); diff --git a/packages/web/src/db.ts b/packages/web/src/db.ts index 9bddde3..4a60615 100644 --- a/packages/web/src/db.ts +++ b/packages/web/src/db.ts @@ -131,6 +131,9 @@ function migrateIssuesTable(db: Database): void { if (!have.has("ai_status")) { db.exec(applyPrefix("ALTER TABLE {{issues}} ADD COLUMN ai_status TEXT NOT NULL DEFAULT ''")); } + if (!have.has("ai_status_since")) { + db.exec(applyPrefix("ALTER TABLE {{issues}} ADD COLUMN ai_status_since TEXT NOT NULL DEFAULT ''")); + } if (!have.has("model")) { db.exec(applyPrefix("ALTER TABLE {{issues}} ADD COLUMN model TEXT NOT NULL DEFAULT ''")); } @@ -430,6 +433,7 @@ async function migrateMysqlColumn(pool: Pool, table: string, column: string, ddl async function migrateMysqlIssuesAiStatus(pool: Pool): Promise { await migrateMysqlColumn(pool, "issues", "ai_status", "ai_status VARCHAR(32) NOT NULL DEFAULT ''"); + await migrateMysqlColumn(pool, "issues", "ai_status_since", "ai_status_since VARCHAR(64) NOT NULL DEFAULT ''"); await migrateMysqlColumn(pool, "issues", "model", "model VARCHAR(128) NOT NULL DEFAULT ''"); await migrateMysqlColumn(pool, "issues", "runtime", "runtime VARCHAR(32) NOT NULL DEFAULT ''"); await migrateMysqlColumn(pool, "issues", "upstream_issue_number", "upstream_issue_number INT DEFAULT NULL"); diff --git a/packages/web/src/index.ts b/packages/web/src/index.ts index 853f0f5..de08565 100644 --- a/packages/web/src/index.ts +++ b/packages/web/src/index.ts @@ -40,6 +40,7 @@ import { postComment, setIssueState, updateIssueAiStatus, + listIssuesByAiStatus, updateIssueModel, updateIssueRuntime, listIssues, @@ -116,6 +117,7 @@ import { import { classifyActor, type CommentView } from "./render/components"; import { buildWebhooksPage } from "./views/webhooks"; import { buildWebhookDeliveriesPage } from "./views/webhookDeliveries"; +import { buildBadgeMonitorPage, type DaemonBadgeReport, type DaemonBadgeEntry } from "./views/badgeMonitor"; import { browseRemoteFile, proxyFileSince, RemoteFileError } from "./remote-file"; import { buildProjectMembersPage } from "./views/projectMembers"; import { @@ -632,6 +634,14 @@ async function handle(req: Request, url: URL, ip: string, ctx: { authed: boolean const concurrency = Number(cfgKv[`concurrency:${owner}/${repo}`]) || null; return json({ dispatchOff: globalOff || projectOff || issueOff, aiStatus, sessionResetMs, concurrency }); } + if (url.pathname === "/api/v1/ai-status" && req.method === "GET") { + const status = url.searchParams.get("status") ?? ""; + if (!status) return json({ error: "status required" }, 400); + // Machine-facing badge enumeration for the daemon's stuck-badge sweep: + // every issue currently carrying the given ai_status, fleet-wide. + const badges = await listIssuesByAiStatus(status); + return json({ badges }); + } if (url.pathname === "/api/v1/wake-logins" && req.method === "GET") { const owner = url.searchParams.get("owner") ?? ""; const repo = url.searchParams.get("repo") ?? ""; @@ -1722,6 +1732,29 @@ async function handle(req: Request, url: URL, ip: string, ctx: { authed: boolean return html(buildWebhookDeliveriesPage(ctx.user!, deliveries)); } + if (url.pathname === "/admin/badges") { + const webBadges = await listIssuesByAiStatus("processing"); + const active = await getActiveDaemons(); + const reports = await Promise.all( + active.map(async (d): Promise => { + try { + const res = await fetch(`${d.endpoint.replace(/\/$/, "")}/api/badges`, { + signal: AbortSignal.timeout(5_000), + }); + if (!res.ok) throw new Error(`HTTP ${res.status}`); + const body = (await res.json()) as { checkedAt: number; intervalMs: number; entries: DaemonBadgeEntry[] }; + return { daemonId: d.id, endpoint: d.endpoint, reachable: true, ...body }; + } catch { + return { + daemonId: d.id, endpoint: d.endpoint, reachable: false, + checkedAt: Date.now(), intervalMs: 0, entries: [], + }; + } + }), + ); + return html(buildBadgeMonitorPage(ctx.user!, webBadges, reports, Date.now())); + } + const adminPatRevoke = url.pathname.match(/^\/admin\/tokens\/(\d+)\/revoke$/); if (adminPatRevoke && req.method === "POST") { const id = Number(adminPatRevoke[1]); diff --git a/packages/web/src/schema-mysql.sql b/packages/web/src/schema-mysql.sql index 37bb3a8..87b66e1 100644 --- a/packages/web/src/schema-mysql.sql +++ b/packages/web/src/schema-mysql.sql @@ -54,6 +54,7 @@ CREATE TABLE IF NOT EXISTS {{issues}} ( updated_at VARCHAR(40) NOT NULL, closed_at VARCHAR(40) DEFAULT NULL, ai_status VARCHAR(32) NOT NULL DEFAULT '', + ai_status_since VARCHAR(64) NOT NULL DEFAULT '', model VARCHAR(128) NOT NULL DEFAULT '', runtime VARCHAR(32) NOT NULL DEFAULT '', upstream_issue_number INT DEFAULT NULL, diff --git a/packages/web/src/schema.sql b/packages/web/src/schema.sql index b00df7f..9888be1 100644 --- a/packages/web/src/schema.sql +++ b/packages/web/src/schema.sql @@ -65,6 +65,8 @@ CREATE TABLE IF NOT EXISTS {{issues}} ( closed_at TEXT, -- AI processing status: '' (none) | 'processing' | 'halted' | 'completed' | 'failed' ai_status TEXT NOT NULL DEFAULT '', + -- ISO timestamp of when ai_status last changed (TTL basis for badge liveness). + ai_status_since TEXT NOT NULL DEFAULT '', -- Resolved "provider/model" for this issue. Empty = inherit project/global default. model TEXT NOT NULL DEFAULT '', -- Runtime backend pinned for this issue ('' = daemon default, 'opencode', 'pi'). diff --git a/packages/web/src/store.ts b/packages/web/src/store.ts index b7e9d3d..31f76e1 100644 --- a/packages/web/src/store.ts +++ b/packages/web/src/store.ts @@ -525,7 +525,36 @@ export async function setIssueState( } export async function updateIssueAiStatus(issueId: number, status: string): Promise { - await getDB().run("UPDATE {{issues}} SET ai_status = ? WHERE id = ?", [status, issueId]); + const row = await getDB().get<{ ai_status: string }>("SELECT ai_status FROM {{issues}} WHERE id = ?", [issueId]); + if (!row) return; + if (row.ai_status === status) return; + // Stamp the change time only when the status actually flips — badge age + // (now - ai_status_since) drives the daemon's TTL arbitration. + await getDB().run("UPDATE {{issues}} SET ai_status = ?, ai_status_since = ? WHERE id = ?", [status, now(), issueId]); +} + +export interface AiStatusBadgeRow { + owner: string; + repo: string; + number: number; + aiStatus: string; + since: string | null; +} + +export async function listIssuesByAiStatus(status: string): Promise { + const rows = await getDB().all<{ + owner: string; name: string; number: number; ai_status: string; ai_status_since: string; + }>( + "SELECT p.owner AS owner, p.name AS name, i.number AS number, i.ai_status AS ai_status, i.ai_status_since AS ai_status_since FROM {{issues}} i JOIN {{projects}} p ON i.project_id = p.id WHERE i.ai_status = ? ORDER BY i.ai_status_since ASC", + [status], + ); + return rows.map((r) => ({ + owner: r.owner, + repo: r.name, + number: r.number, + aiStatus: r.ai_status, + since: r.ai_status_since || null, + })); } export async function getIssueAiStatusByNumber(owner: string, repo: string, number: number): Promise { diff --git a/packages/web/src/views/badgeMonitor.ts b/packages/web/src/views/badgeMonitor.ts new file mode 100644 index 0000000..7de8dbf --- /dev/null +++ b/packages/web/src/views/badgeMonitor.ts @@ -0,0 +1,150 @@ +import type { UserRow } from "../store"; +import type { AiStatusBadgeRow } from "../store"; +import { THEME_CSS, escapeHtml, escapeAttr, tabNavHTML } from "../render/layout"; + +export interface DaemonBadgeEntry { + owner: string; + repo: string; + number: number; + aiStatus: string; + since: number | null; + engineRecord: boolean; + ownedByMe: boolean; + sessions: number; + pidAlive: boolean; + outputAgeSec: number | null; + modelAgeSec: number | null; + verdict: "alive" | "stale" | "orphan" | "skipped"; + action?: string; +} + +export interface DaemonBadgeReport { + daemonId: number; + endpoint: string; + reachable: boolean; + checkedAt: number; + intervalMs: number; + entries: DaemonBadgeEntry[]; +} + +const VERDICT_LABEL: Record = { + alive: { text: "存活", cls: "ok" }, + stale: { text: "卡死", cls: "err" }, + orphan: { text: "孤儿", cls: "warn" }, + skipped: { text: "跳过", cls: "" }, +}; + +function ageText(sec: number | null): string { + if (sec === null) return "—"; + if (sec < 60) return `${Math.round(sec)}s`; + if (sec < 3600) return `${Math.round(sec / 60)}m`; + return `${(sec / 3600).toFixed(1)}h`; +} + +function signalCell(v: boolean | null): string { + if (v === null) return `—`; + return v ? `✓` : `✗`; +} + +function entryRow(e: DaemonBadgeEntry): string { + const v = VERDICT_LABEL[e.verdict] ?? VERDICT_LABEL.skipped; + const issue = `/issues/${escapeAttr(e.owner)}/${escapeAttr(e.repo)}/${e.number}`; + const action = e.action ? `
${escapeHtml(e.action)}
` : ""; + return ` + ${escapeHtml(e.owner + "/" + e.repo + "#" + e.number)} + ${escapeHtml(e.aiStatus || "(空)")} + ${e.engineRecord ? "有" : "无"} + ${e.sessions} + ${signalCell(e.pidAlive)} + ${ageText(e.outputAgeSec)} + ${ageText(e.modelAgeSec)} + ${v.text} + ${action} + `; +} + +function daemonSection(d: DaemonBadgeReport): string { + if (!d.reachable) { + return `
+

daemon #${d.daemonId} 不可达

+

GET ${escapeHtml(d.endpoint)}/api/badges 失败

+
`; + } + const rows = d.entries.length > 0 + ? d.entries.map(entryRow).join("") + : `该 daemon 视角下没有 processing 徽标`; + const stale = d.entries.filter((e) => e.verdict === "stale" || e.verdict === "orphan").length; + return `
+

daemon #${d.daemonId} + ${stale > 0 ? stale + " 异常" : "正常"} + 扫描间隔 ${Math.round(d.intervalMs / 1000)}s · 上次 ${new Date(d.checkedAt).toLocaleString()} +

+ + + ${rows} +
issue状态引擎记录会话数pid输出模型判定动作
+
`; +} + +export function buildBadgeMonitorPage( + viewer: UserRow, + webBadges: AiStatusBadgeRow[], + daemons: DaemonBadgeReport[], + nowMs: number, +): string { + const webRows = webBadges.length > 0 + ? webBadges.map((b) => { + const since = b.since ? Date.parse(b.since) : null; + const age = since ? ageText((nowMs - since) / 1000) : "未知"; + return ` + ${escapeHtml(b.owner + "/" + b.repo + "#" + b.number)} + ${escapeHtml(b.aiStatus)} + ${b.since ? escapeHtml(b.since) : "—"} + ${age} + `; + }).join("") + : `web 侧当前没有 processing 徽标`; + + const daemonSections = daemons.length > 0 + ? daemons.map(daemonSection).join("\n") + : `

没有活跃 daemon(心跳超过 2 分钟)

`; + + return ` + + + +ework-web · Processing 徽标监控 + + + +
+

Processing 徽标监控

+

web 侧全部 processing 徽标 + 各 daemon 卡死探测器的三信号状态(pid 存活 / 会话输出 / 模型流量)。卡死探测器每分钟扫描,孤儿徽标自动仲裁。

+

Web 侧 processing 列表(${webBadges.length})

+ + +${webRows} +
issue状态写入时间徽标年龄
+${daemonSections} +
+`; +} diff --git a/packages/web/src/views/users.ts b/packages/web/src/views/users.ts index d65f396..07df978 100644 --- a/packages/web/src/views/users.ts +++ b/packages/web/src/views/users.ts @@ -139,6 +139,7 @@ ${flashHtml} 用户 所有 Token Webhook 投递 + 徽标监控