import { randomUUID } from "node:crypto"; import type { OneDevSessionContext } from "./context.js"; import { runTod, TodError, type CommandExecutor } from "./tod.js"; export type OneDevWatchKind = "build" | "pull_attention"; export interface OneDevWatchEmission { id: string; kind: OneDevWatchKind; ref: string; message: string; details: Readonly>; } export interface OneDevWatchInfo { id: string; kind: OneDevWatchKind; ref: string; intervalSeconds: number; expiresAt: string; failures: number; lastError?: string; } export interface AddOneDevWatchOptions { kind: OneDevWatchKind; ref: string; intervalSeconds?: number; ttlMinutes?: number; signal?: AbortSignal; } export interface AddOneDevWatchResult { watch?: OneDevWatchInfo; message: string; } interface BuildSnapshot { kind: "build"; fingerprint: string; status?: string; } interface PullSnapshot { kind: "pull_attention"; fingerprint: string; state?: string; head?: string; mergeable?: boolean; conflicted?: boolean; unresolvedComments: number; reviewSignals: number; buildSummary: string; } type WatchSnapshot = BuildSnapshot | PullSnapshot; interface WatchRecord { id: string; kind: OneDevWatchKind; ref: string; intervalMs: number; expiresAt: number; nextPollAt: number; failures: number; lastError?: string; baseline: WatchSnapshot; } interface WatchDeps { exec: CommandExecutor; context: () => OneDevSessionContext; } const DEFAULT_INTERVAL_SECONDS = 30; const DEFAULT_TTL_MINUTES = 60; const MIN_INTERVAL_SECONDS = 15; const MAX_INTERVAL_SECONDS = 300; const MAX_TTL_MINUTES = 24 * 60; const MAX_WATCHES = 10; const MAX_BACKOFF_MS = 5 * 60_000; const SNAPSHOT_TIMEOUT_MS = 10_000; const TERMINAL_BUILD_STATUSES = new Set([ "SUCCESSFUL", "FAILED", "CANCELLED", "TIMED_OUT", ]); type JsonPrimitive = string | number | boolean | null; type JsonValue = JsonPrimitive | JsonValue[] | { [key: string]: JsonValue }; function isJsonValue(value: unknown): value is JsonValue { if ( value === null || typeof value === "string" || typeof value === "number" || typeof value === "boolean" ) { return true; } if (Array.isArray(value)) return value.every(isJsonValue); if (typeof value !== "object") return false; return Object.values(value).every(isJsonValue); } function json(text: string, source: string): JsonValue { try { const parsed: unknown = JSON.parse(text); if (!isJsonValue(parsed)) throw new Error("response is not a JSON value"); return parsed; } catch (error) { throw new TodError( `tod returned invalid JSON for ${source}: ${error instanceof Error ? error.message : String(error)}`, ); } } function stable(value: JsonValue): string { if (Array.isArray(value)) return `[${value.map(stable).join(",")}]`; if (typeof value === "object" && value !== null) { return `{${Object.entries(value) .sort(([left], [right]) => left.localeCompare(right)) .map(([key, child]) => `${JSON.stringify(key)}:${stable(child)}`) .join(",")}}`; } return JSON.stringify(value) ?? "null"; } function visit( value: JsonValue, callback: (key: string, value: JsonValue) => void, ): void { if (Array.isArray(value)) { for (const item of value) visit(item, callback); return; } if (typeof value !== "object" || value === null) return; for (const [key, child] of Object.entries(value)) { callback(key, child); visit(child, callback); } } function firstString(value: JsonValue, keys: readonly string[]): string | undefined { let found: string | undefined; visit(value, (key, child) => { if ( found === undefined && keys.includes(key) && (typeof child === "string" || typeof child === "number") ) { found = String(child); } }); return found; } function firstBoolean(value: JsonValue, keys: readonly string[]): boolean | undefined { let found: boolean | undefined; visit(value, (key, child) => { if (found === undefined && keys.includes(key) && typeof child === "boolean") { found = child; } }); return found; } function topLevelString( value: JsonValue, keys: readonly string[], ): string | undefined { if (typeof value !== "object" || value === null || Array.isArray(value)) return undefined; for (const key of keys) { const child = value[key]; if (typeof child === "string" && child.trim() !== "") return child.trim(); if (typeof child === "number") return String(child); } return undefined; } function topLevelBoolean( value: JsonValue, keys: readonly string[], ): boolean | undefined { if (typeof value !== "object" || value === null || Array.isArray(value)) return undefined; for (const key of keys) { const child = value[key]; if (typeof child === "boolean") return child; } return undefined; } function unresolvedCodeComments( value: JsonValue, ): Array<{ id: string | number; replies: number }> { const comments = Array.isArray(value) ? value : []; const unresolved: Array<{ id: string | number; replies: number }> = []; for (let index = 0; index < comments.length; index += 1) { const comment = comments[index]; if (Array.isArray(comment) || typeof comment !== "object" || comment === null) continue; const record = comment; if (record.resolved === false) { const id = record.id; unresolved.push({ id: typeof id === "string" || typeof id === "number" ? id : index, replies: Array.isArray(record.replies) ? record.replies.length : 0, }); } } return unresolved; } function buildSignals( value: JsonValue, ): Array<{ id: string; status: string }> { const builds = Array.isArray(value) ? value : []; return builds .map((build, index) => { if (Array.isArray(build) || typeof build !== "object" || build === null) return undefined; const record = build; const rawStatus = record.status; if (typeof rawStatus !== "string") return undefined; const rawId = record.number ?? record.id ?? index; return { id: String(rawId), status: rawStatus.toUpperCase() }; }) .filter((build): build is { id: string; status: string } => build !== undefined) .sort((left, right) => left.id.localeCompare(right.id)); } function valuesAtKeys(value: JsonValue, keys: readonly string[]): JsonValue[] { const values: JsonValue[] = []; visit(value, (key, child) => { if (keys.includes(key)) values.push(child); }); return values; } function terminalBuild(status: string | undefined): boolean { return status !== undefined && TERMINAL_BUILD_STATUSES.has(status.toUpperCase()); } function assertRange( value: number, minimum: number, maximum: number, name: string, ): number { if (!Number.isFinite(value) || value < minimum || value > maximum) { throw new TodError(`${name} must be between ${minimum} and ${maximum}`); } return value; } export class OneDevSourceWatchManager { readonly #records = new Map(); #timer: ReturnType | undefined; #polling = false; #generation = 0; constructor( readonly deps: WatchDeps, readonly emit: (emission: OneDevWatchEmission) => void, readonly now: () => number = Date.now, ) {} list(): OneDevWatchInfo[] { return [...this.#records.values()] .sort((left, right) => left.expiresAt - right.expiresAt) .map((record) => this.#info(record)); } stop(id: string): boolean { const stopped = this.#records.delete(id); if (stopped) this.#schedule(); return stopped; } clear(): void { this.#generation += 1; this.#records.clear(); if (this.#timer) clearTimeout(this.#timer); this.#timer = undefined; } async add(options: AddOneDevWatchOptions): Promise { const ref = options.ref.trim(); if (!ref) throw new TodError("watch ref must not be empty"); const existing = this.#duplicate(options.kind, ref); if (existing) return this.#duplicateResult(existing); if (this.#records.size >= MAX_WATCHES) { throw new TodError(`at most ${MAX_WATCHES} OneDev watches may be active`); } const intervalSeconds = assertRange( options.intervalSeconds ?? DEFAULT_INTERVAL_SECONDS, MIN_INTERVAL_SECONDS, MAX_INTERVAL_SECONDS, "interval_seconds", ); const ttlMinutes = assertRange( options.ttlMinutes ?? DEFAULT_TTL_MINUTES, 1, MAX_TTL_MINUTES, "ttl_minutes", ); const generation = this.#generation; const baseline = await this.#snapshot(options.kind, ref, options.signal); if (generation !== this.#generation) { throw new TodError("watch cancelled because the OneDev session changed"); } if (baseline.kind === "build" && terminalBuild(baseline.status)) { return { message: `Build ${ref} is already ${baseline.status}; no watch was added.`, }; } const duplicate = this.#duplicate(options.kind, ref); if (duplicate) { return this.#duplicateResult(duplicate); } if (this.#records.size >= MAX_WATCHES) { throw new TodError(`at most ${MAX_WATCHES} OneDev watches may be active`); } const now = this.now(); const record: WatchRecord = { id: randomUUID().slice(0, 8), kind: options.kind, ref, intervalMs: intervalSeconds * 1_000, expiresAt: now + ttlMinutes * 60_000, nextPollAt: now + intervalSeconds * 1_000, failures: 0, baseline, }; this.#records.set(record.id, record); this.#schedule(); return { watch: this.#info(record), message: `Watch ${record.id} will notify once when ${options.kind} ${ref} needs attention.`, }; } #duplicate(kind: OneDevWatchKind, ref: string): WatchRecord | undefined { return [...this.#records.values()].find( (record) => record.kind === kind && record.ref === ref, ); } #duplicateResult(record: WatchRecord): AddOneDevWatchResult { return { watch: this.#info(record), message: `Watch ${record.id} already monitors ${record.kind} ${record.ref}.`, }; } #info(record: WatchRecord): OneDevWatchInfo { return { id: record.id, kind: record.kind, ref: record.ref, intervalSeconds: record.intervalMs / 1_000, expiresAt: new Date(record.expiresAt).toISOString(), failures: record.failures, lastError: record.lastError, }; } async #snapshot( kind: OneDevWatchKind, ref: string, signal?: AbortSignal, ): Promise { const context = this.deps.context(); if (context.status !== "ready" || !context.project) { throw new TodError( context.problem ?? "OneDev watch requires a ready project context", ); } if (kind === "build") { const output = await runTod(this.deps.exec, ["build", "get", ref], { cwd: context.cwd, signal, timeoutMs: SNAPSHOT_TIMEOUT_MS, maxOutputBytes: 512_000, }); const value = json(output.text, `build ${ref}`); const status = firstString(value, ["status"])?.toUpperCase(); return { kind, status, fingerprint: status ?? stable(value), }; } const [pull, comments, builds] = await Promise.all([ runTod(this.deps.exec, ["pr", "get", ref], { cwd: context.cwd, signal, timeoutMs: SNAPSHOT_TIMEOUT_MS, maxOutputBytes: 512_000, }), runTod(this.deps.exec, ["pr", "get-code-comments", ref], { cwd: context.cwd, signal, timeoutMs: SNAPSHOT_TIMEOUT_MS, maxOutputBytes: 512_000, }), runTod(this.deps.exec, ["pr", "get-builds", ref], { cwd: context.cwd, signal, timeoutMs: SNAPSHOT_TIMEOUT_MS, maxOutputBytes: 512_000, }), ]); const pullValue = json(pull.text, `pull request ${ref}`); const unresolved = unresolvedCodeComments( json(comments.text, `pull request ${ref} code comments`), ); const buildsValue = json(builds.text, `pull request ${ref} builds`); const buildSignal = buildSignals( buildsValue, ); const reviews = valuesAtKeys(pullValue, [ "reviewStatus", "reviewDecision", "pendingReviewers", "reviewers", "requestedChanges", ]); const selected = { state: topLevelString(pullValue, ["status", "state"]) ?? null, head: topLevelString(pullValue, [ "headCommitHash", "sourceCommitHash", "headCommit", "sourceCommit", ]) ?? null, mergeable: topLevelBoolean(pullValue, ["mergeable"]) ?? null, conflicted: topLevelBoolean(pullValue, ["hasConflicts", "conflicted"]) ?? null, unresolved, reviews, buildSignal, }; return { kind, fingerprint: stable(selected), state: selected.state ?? undefined, head: selected.head ?? undefined, mergeable: selected.mergeable ?? undefined, conflicted: selected.conflicted ?? undefined, unresolvedComments: unresolved.length, reviewSignals: reviews.length, buildSummary: buildSignal.map((build) => `${build.id}:${build.status}`).join(", ") || "none", }; } #changed(record: WatchRecord, snapshot: WatchSnapshot): boolean { if (record.kind === "build" && snapshot.kind === "build") { return ( terminalBuild(snapshot.status) && snapshot.fingerprint !== record.baseline.fingerprint ); } return snapshot.fingerprint !== record.baseline.fingerprint; } #emission(record: WatchRecord, snapshot: WatchSnapshot): OneDevWatchEmission { if (snapshot.kind === "build") { return { id: record.id, kind: record.kind, ref: record.ref, message: `OneDev build ${record.ref} is ${snapshot.status ?? "updated"}. Inspect it with onedev_build get/get_log.`, details: { status: snapshot.status ?? null }, }; } return { id: record.id, kind: record.kind, ref: record.ref, message: `OneDev pull request ${record.ref} needs attention. Inspect it with onedev_pull get/get_code_comments/get_builds.`, details: { state: snapshot.state ?? null, head: snapshot.head?.slice(0, 12) ?? null, mergeable: snapshot.mergeable ?? null, conflicted: snapshot.conflicted ?? null, unresolvedComments: snapshot.unresolvedComments, reviewSignals: snapshot.reviewSignals, builds: snapshot.buildSummary.slice(0, 200), }, }; } #schedule(): void { if (this.#timer) clearTimeout(this.#timer); this.#timer = undefined; if (this.#records.size === 0 || this.#polling) return; const next = Math.min( ...[...this.#records.values()].map((record) => Math.min(record.nextPollAt, record.expiresAt), ), ); this.#timer = setTimeout(() => { this.#timer = undefined; void this.#tick(); }, Math.max(0, next - this.now())); this.#timer.unref?.(); } async #tick(): Promise { if (this.#polling) return; this.#polling = true; const generation = this.#generation; try { const now = this.now(); const due = [...this.#records.values()].filter( (record) => record.nextPollAt <= now || record.expiresAt <= now, ); for (const record of due) { if (generation !== this.#generation || !this.#records.has(record.id)) break; try { const snapshot = await this.#snapshot(record.kind, record.ref); if (generation !== this.#generation || !this.#records.has(record.id)) break; if (this.#changed(record, snapshot)) { this.#records.delete(record.id); try { this.emit(this.#emission(record, snapshot)); } catch { // a throwing consumer must not re-enter the poll-failure path } continue; } record.failures = 0; record.lastError = undefined; record.baseline = snapshot; record.nextPollAt = this.now() + record.intervalMs; } catch (error) { record.failures += 1; record.lastError = error instanceof Error ? error.message : String(error); record.nextPollAt = this.now() + Math.min( record.intervalMs * 2 ** Math.min(record.failures, 6), MAX_BACKOFF_MS, ); } if (record.expiresAt <= now) this.#records.delete(record.id); } } finally { this.#polling = false; if (generation === this.#generation) this.#schedule(); } } }