diff --git a/apps/mobile/src/features/threads/git/GitOverviewSheet.tsx b/apps/mobile/src/features/threads/git/GitOverviewSheet.tsx index 17b1b3ea377c..828c8838e31a 100644 --- a/apps/mobile/src/features/threads/git/GitOverviewSheet.tsx +++ b/apps/mobile/src/features/threads/git/GitOverviewSheet.tsx @@ -339,7 +339,7 @@ export function GitOverviewSheet(props: GitOverviewSheetProps) { { void tryOpenExternalUrl(link.url, "pull-request").then((opened) => { if (!opened) diff --git a/apps/mobile/src/features/threads/thread-work-log.tsx b/apps/mobile/src/features/threads/thread-work-log.tsx index 826b9205ac75..2f06b9462b4d 100644 --- a/apps/mobile/src/features/threads/thread-work-log.tsx +++ b/apps/mobile/src/features/threads/thread-work-log.tsx @@ -1368,6 +1368,8 @@ function toolGroupSummarySymbolName(kind: ToolGroupSummaryKind): AppSymbolName { case "link-pr": case "unlink-pr": case "list-prs": + case "watch-pr": + case "unwatch-pr": return "arrow.triangle.pull"; case "read": return { ios: "eye", android: "visibility" }; diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index 91ab19b55e01..5381645e9fcf 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -240,6 +240,7 @@ export const make = Effect.gen(function* () { threadTitleRegeneration: true, threadVisitedTracking: true, threadPullRequests: true, + threadPullRequestWatch: true, pullRequestStackActions: true, threadPullRequestLinking: true, serverResolvedCommandContext: true, diff --git a/apps/server/src/mcp/toolkits/pullRequests/handlers.test.ts b/apps/server/src/mcp/toolkits/pullRequests/handlers.test.ts index dab6e7176498..e329ade90d9e 100644 --- a/apps/server/src/mcp/toolkits/pullRequests/handlers.test.ts +++ b/apps/server/src/mcp/toolkits/pullRequests/handlers.test.ts @@ -222,6 +222,58 @@ describe("pull request toolkit handlers", () => { }), ); + it.effect("watching an unlinked pull request links it first", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + const result = yield* harness.call("watch_pull_request", { + url: "https://github.com/t3tools/t3code/pull/9", + }); + // The harness thread never changes, so the result reports what it still holds. + expect(result).toMatchObject({ number: 9, watching: false, wasWatching: false }); + expect(yield* Ref.get(harness.commands)).toMatchObject([ + { + type: "thread.pull-request.watch", + number: 9, + watching: true, + link: { url: "https://github.com/t3tools/t3code/pull/9", source: "agent" }, + }, + ]); + }), + ); + + it.effect("refuses to watch a merged pull request and stops an existing watch", () => + Effect.gen(function* () { + const watch = { + startedAt: "2026-08-20T00:00:00.000Z", + headSha: null, + failedChecks: [], + passed: false, + remarksThrough: "2026-08-20T00:00:00.000Z", + remarkIds: [], + conflicting: false, + wakes: 0, + }; + const merged = makeLink(1, { headBranch: "done" }); + const harness = yield* makeHarness({ + thread: makeThread([ + { ...merged, snapshot: merged.snapshot && { ...merged.snapshot, state: "merged" } }, + makeLink(2, { headBranch: "idle" }), + makeLink(3, { headBranch: "watched", watch }), + ]), + }); + const error = yield* harness + .call("watch_pull_request", { repository: "t3tools/t3code", number: 1 }) + .pipe(Effect.flip); + expect(error).toMatchObject({ _tag: "PullRequestNotOpenError", state: "merged" }); + expect( + yield* harness.call("unwatch_pull_request", { repository: "t3tools/t3code", number: 3 }), + ).toMatchObject({ wasWatching: true }); + expect(yield* Ref.get(harness.commands)).toMatchObject([ + { type: "thread.pull-request.watch", number: 3, watching: false }, + ]); + }), + ); + it.effect("links by repository and number, defaulting the host to the project's", () => Effect.gen(function* () { const harness = yield* makeHarness(); @@ -396,6 +448,7 @@ describe("pull request toolkit handlers", () => { number: 3, url: "https://github.com/t3tools/t3code/pull/3", source: "agent", + watching: false, state: "open", title: "PR 3", headBranch: "feat-c", diff --git a/apps/server/src/mcp/toolkits/pullRequests/handlers.ts b/apps/server/src/mcp/toolkits/pullRequests/handlers.ts index 19677e1bd212..97af65e1a7bc 100644 --- a/apps/server/src/mcp/toolkits/pullRequests/handlers.ts +++ b/apps/server/src/mcp/toolkits/pullRequests/handlers.ts @@ -34,7 +34,9 @@ import { PullRequestHostRequiredError, PullRequestUnlinkFailedError, PullRequestListFailedError, + PullRequestNotOpenError, type PullRequestTargetInput, + PullRequestWatchFailedError, PullRequestThreadNotFoundError, PullRequestsToolkit, type ThreadPullRequestEntry, @@ -122,6 +124,7 @@ function entryOf( number: link.number, url: link.url, source: link.source, + watching: link.watch !== undefined, state: link.snapshot?.state ?? null, title: link.snapshot?.title ?? null, headBranch: link.snapshot?.headBranch ?? null, @@ -163,7 +166,8 @@ const make = Effect.gen(function* () { Failure: | typeof PullRequestLinkFailedError | typeof PullRequestUnlinkFailedError - | typeof PullRequestListFailedError, + | typeof PullRequestListFailedError + | typeof PullRequestWatchFailedError, ) { const scope = yield* McpInvocationContext.requireMcpCapability("pull-requests"); const thread = yield* engine @@ -178,7 +182,10 @@ const make = Effect.gen(function* () { const projectOf = ( thread: OrchestrationV2ThreadShell, - Failure: typeof PullRequestLinkFailedError | typeof PullRequestUnlinkFailedError, + Failure: + | typeof PullRequestLinkFailedError + | typeof PullRequestUnlinkFailedError + | typeof PullRequestWatchFailedError, ) => projects.getShell(thread.projectId).pipe( Effect.map(Option.getOrUndefined), @@ -186,14 +193,65 @@ const make = Effect.gen(function* () { ); const dispatchFailure = - (Failure: typeof PullRequestLinkFailedError | typeof PullRequestUnlinkFailedError) => + ( + Failure: + | typeof PullRequestLinkFailedError + | typeof PullRequestUnlinkFailedError + | typeof PullRequestWatchFailedError, + ) => ( cause: Cause.Cause, - ): Effect.Effect => + ): Effect.Effect< + never, + PullRequestLinkFailedError | PullRequestUnlinkFailedError | PullRequestWatchFailedError + > => Cause.hasInterruptsOnly(cause) ? Effect.failCause(cause as Cause.Cause) : Effect.fail(new Failure({ cause })); + /** + * Starts or stops a watch. One command links an unlinked pull request and watches it, and the + * result reports the state the thread holds afterwards. + */ + const setWatching = Effect.fn("PullRequestsToolkit.setWatching")(function* ( + input: PullRequestTargetInput, + watching: boolean, + ) { + const thread = yield* requireThread(PullRequestWatchFailedError); + const project = yield* projectOf(thread, PullRequestWatchFailedError); + const target = yield* resolveTarget(input, project); + const watchedLink = (shell: OrchestrationV2ThreadShell) => + threadPullRequestsOf(shell).find( + (link) => link.source !== "stack-dismissed" && threadPullRequestKeysEqual(link, target), + ); + const before = watchedLink(thread); + const state = before?.snapshot?.state; + if (watching && state !== undefined && state !== "open") { + return yield* new PullRequestNotOpenError({ state }); + } + yield* engine + .dispatch({ + type: "thread.pull-request.watch", + commandId: yield* commandId("mcp-pr-watch", thread.id), + threadId: thread.id, + host: target.host, + repository: target.repository, + number: target.number, + watching, + ...(watching ? { link: { url: target.url, source: "agent" as const } } : {}), + }) + .pipe(Effect.catchCause(dispatchFailure(PullRequestWatchFailedError))); + const after = yield* requireThread(PullRequestWatchFailedError); + return { + host: target.host, + repository: target.repository, + number: target.number, + url: target.url, + watching: watchedLink(after)?.watch !== undefined, + wasWatching: before?.watch !== undefined, + }; + }); + return PullRequestsToolkit.of({ link_pull_request: (input) => Effect.gen(function* () { @@ -260,6 +318,8 @@ const make = Effect.gen(function* () { }), list_thread_pull_requests: () => requireThread(PullRequestListFailedError).pipe(Effect.map(listThreadPullRequests)), + watch_pull_request: (input) => setWatching(input, true), + unwatch_pull_request: (input) => setWatching(input, false), }); }); diff --git a/apps/server/src/mcp/toolkits/pullRequests/tools.ts b/apps/server/src/mcp/toolkits/pullRequests/tools.ts index 90f22cea2615..52ac473a3d50 100644 --- a/apps/server/src/mcp/toolkits/pullRequests/tools.ts +++ b/apps/server/src/mcp/toolkits/pullRequests/tools.ts @@ -108,6 +108,24 @@ export class PullRequestUnlinkFailedError extends Schema.TaggedError()( + "PullRequestWatchFailedError", + { cause: Schema.Defect() }, +) { + override get message(): string { + return "Could not change whether the pull request is watched."; + } +} + +export class PullRequestNotOpenError extends Schema.TaggedError()( + "PullRequestNotOpenError", + { state: Schema.String }, +) { + override get message(): string { + return `The pull request is ${this.state}, so there is nothing to watch.`; + } +} + export class PullRequestListFailedError extends Schema.TaggedError()( "PullRequestListFailedError", { cause: Schema.Defect() }, @@ -126,6 +144,8 @@ export const PullRequestToolError = Schema.Union([ PullRequestLinkFailedError, PullRequestUnlinkFailedError, PullRequestListFailedError, + PullRequestWatchFailedError, + PullRequestNotOpenError, ]); export type PullRequestToolError = typeof PullRequestToolError.Type; @@ -154,9 +174,21 @@ export const UnlinkPullRequestResult = Schema.Struct({ }); export type UnlinkPullRequestResult = typeof UnlinkPullRequestResult.Type; +export const WatchPullRequestResult = Schema.Struct({ + ...PullRequestIdentity, + watching: Schema.Boolean.annotate({ + description: "Whether T3 Code now watches the pull request for this thread.", + }), + wasWatching: Schema.Boolean.annotate({ + description: "Whether it was already watched before the call.", + }), +}); +export type WatchPullRequestResult = typeof WatchPullRequestResult.Type; + export const ThreadPullRequestEntry = Schema.Struct({ ...PullRequestIdentity, source: ThreadPullRequestLinkSource, + watching: Schema.Boolean, state: Schema.NullOr(PullRequestState), title: Schema.NullOr(Schema.String), headBranch: Schema.NullOr(Schema.String), @@ -224,8 +256,38 @@ const ListThreadPullRequestsTool = Tool.make("list_thread_pull_requests", { .annotate(Tool.Idempotent, true) .annotate(Tool.OpenWorld, false); +const WatchPullRequestTool = Tool.make("watch_pull_request", { + description: + "Have T3 Code watch an open pull request for this thread, linking it first if needed. T3 Code checks it every minute and wakes you with a message when a check fails, the required checks pass, someone else comments or reviews, or the branch starts to conflict with its base. Use this to monitor or babysit a pull request instead of polling, sleeping, or running a watcher. Only comments posted after this call wake you, so handle the existing ones first, then end your turn. A wake is news, not a merge decision: check readiness yourself before merging. Watching ends when the pull request merges or closes, when T3 Code cannot read it for 15 minutes, or when you call unwatch_pull_request.", + parameters: PullRequestTargetInput, + success: WatchPullRequestResult, + failure: PullRequestToolError, + dependencies, +}) + .annotate(Tool.Title, "Watch pull request") + .annotate(Tool.Readonly, false) + .annotate(Tool.Destructive, false) + .annotate(Tool.Idempotent, true) + .annotate(Tool.OpenWorld, false); + +const UnwatchPullRequestTool = Tool.make("unwatch_pull_request", { + description: + "Stop T3 Code from watching a pull request for this thread. The pull request stays linked. Pass the URL, or repository plus number.", + parameters: PullRequestTargetInput, + success: WatchPullRequestResult, + failure: PullRequestToolError, + dependencies, +}) + .annotate(Tool.Title, "Stop watching pull request") + .annotate(Tool.Readonly, false) + .annotate(Tool.Destructive, false) + .annotate(Tool.Idempotent, true) + .annotate(Tool.OpenWorld, false); + export const PullRequestsToolkit = Toolkit.make( LinkPullRequestTool, UnlinkPullRequestTool, ListThreadPullRequestsTool, + WatchPullRequestTool, + UnwatchPullRequestTool, ); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 28296f9d7b78..2707fe5b8344 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -19,6 +19,8 @@ import { OrchestrationV2Command, type OrchestrationV2InternalCommand, type OrchestrationV2ServerCommand, + type ThreadPullRequestLink, + type ThreadPullRequestWatch, type OrchestrationV2AppThread, type OrchestrationV2ContextHandoff, type OrchestrationV2ContextSourcePoint, @@ -371,6 +373,8 @@ function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { case "thread.pull-request.link": case "thread.pull-request.unlink": case "thread.pull-request-link.sync": + case "thread.pull-request.watch": + case "thread.pull-request-watch.sync": case "thread.pull-request.sync": case "thread.title.regeneration.complete": case "thread.runtime-mode.set": @@ -444,6 +448,35 @@ function hasLiveRun(projection: Pick): ); } +/** The link with its watch replaced, or removed when `watch` is undefined. */ +function withPullRequestWatch( + link: ThreadPullRequestLink, + watch: ThreadPullRequestWatch | undefined, +): ThreadPullRequestLink { + const { watch: _previous, ...rest } = link; + return watch === undefined ? rest : { ...rest, watch }; +} + +/** A legacy single-PR link as a link entry. Re-linking a pull request keeps its watch. */ +function legacyPullRequestLink( + thread: OrchestrationV2AppThread, + linked: ThreadLinkedPullRequest, + now: DateTime.Utc, +): ThreadPullRequestLink { + const key = legacyThreadPullRequestKey(linked); + return withPullRequestWatch( + { + ...key, + url: linked.url, + source: "manual", + linkedAt: DateTime.formatIso(now), + snapshot: null, + stack: null, + }, + threadPullRequestsOf(thread).find((link) => threadPullRequestKeysEqual(link, key))?.watch, + ); +} + function delegatedCompletionWakeDetail(taskIds: ReadonlyArray): string { const taskList = taskIds.join(", "); return taskIds.length === 1 @@ -2188,9 +2221,66 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); }); + // Checked under the thread lock: the watch or the thread can change while the host is read. + // The watch is recorded first so the wake's own thread events carry it. + const dispatchPullRequestWatchSync = Effect.fn("orchestrationV2.dispatch.pullRequestWatchSync")( + function* ( + command: Extract< + OrchestrationV2ServerCommand, + { readonly type: "thread.pull-request-watch.sync" } + >, + events: Ref.Ref>, + effects: Ref.Ref>, + ) { + const thread = yield* projectionStore + .getThread(command.threadId) + .pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), + ), + ); + const key = normalizeThreadPullRequestKey(command); + const link = threadPullRequestsOf(thread).find( + (candidate) => + candidate.source !== "stack-dismissed" && threadPullRequestKeysEqual(candidate, key), + ); + // Same rule as a direct message.dispatch: a provider-native subagent takes no messages. + const inactive = + thread.archivedAt !== null || + thread.settledOverride === "settled" || + thread.settledAt !== null || + isProviderNativeSubagentThread(thread); + if (link?.watch?.startedAt !== command.startedAt || (command.wake && inactive)) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: "The pull request watch ended or its thread settled while it was read.", + }); + } + yield* dispatchThreadMutation(command, events, effects); + if (command.wake === undefined) return; + yield* dispatchMessage( + { + type: "message.dispatch", + commandId: command.commandId, + threadId: command.threadId, + messageId: command.wake.messageId, + text: command.wake.text, + notification: command.wake.notification, + attachments: [], + dispatchMode: { type: "queue_after_active" }, + createdBy: "agent", + creationSource: "server", + }, + events, + effects, + ); + }, + ); + const dispatchThreadMutation = Effect.fn("orchestrationV2.dispatch.threadMutation")(function* ( command: Extract< - OrchestrationV2Command, + OrchestrationV2ServerCommand, { readonly type: | "thread.archive" @@ -2209,6 +2299,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio | "thread.pull-request.link" | "thread.pull-request.unlink" | "thread.pull-request-link.sync" + | "thread.pull-request.watch" + | "thread.pull-request-watch.sync" | "thread.pull-request.sync" | "thread.title.regeneration.complete" | "thread.runtime-mode.set" @@ -2236,6 +2328,16 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio cause: `Thread ${command.threadId} is deleted.`, }); } + if ( + command.type === "thread.pull-request.watch" && + command.watching && + isProviderNativeSubagentThread(thread) + ) { + return yield* new OrchestratorSubagentThreadReadOnlyError({ + commandId: command.commandId, + threadId: command.threadId, + }); + } if ( command.type === "thread.metadata.update" && command.expectedWorktreePath !== undefined && @@ -2719,16 +2821,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ), ...(command.linkedPullRequest - ? [ - { - ...legacyThreadPullRequestKey(command.linkedPullRequest), - url: command.linkedPullRequest.url, - source: "manual" as const, - linkedAt: DateTime.formatIso(now), - snapshot: null, - stack: null, - }, - ] + ? [legacyPullRequestLink(thread, command.linkedPullRequest, now)] : []), ], }), @@ -2801,7 +2894,12 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ); pullRequests = belongsToStack ? links.map((link) => - link === existing ? { ...link, source: "stack-dismissed" as const } : link, + link === existing + ? { + ...withPullRequestWatch(link, undefined), + source: "stack-dismissed" as const, + } + : link, ) : links.filter((link) => link !== existing); } else { @@ -2830,6 +2928,61 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio updatedAt: command.type === "thread.pull-request-link.sync" ? thread.updatedAt : now, }; } + case "thread.pull-request.watch": + case "thread.pull-request-watch.sync": { + const key = normalizeThreadPullRequestKey(command); + const startedAt = DateTime.formatIso(now); + const linked = threadPullRequestsOf(thread); + const visible = (link: ThreadPullRequestLink) => + link.source !== "stack-dismissed" && threadPullRequestKeysEqual(link, key); + // A watch started on an unlinked (or dismissed) pull request links it in the same step. + const links = + command.type === "thread.pull-request.watch" && + command.watching && + command.link !== undefined && + !linked.some(visible) + ? [ + ...linked.filter((link) => !threadPullRequestKeysEqual(link, key)), + { + ...key, + url: command.link.url, + source: command.link.source, + linkedAt: startedAt, + snapshot: null, + stack: null, + }, + ] + : linked; + const existing = links.find(visible); + if (existing === undefined) return thread; + const watch = + command.type === "thread.pull-request-watch.sync" + ? // Progress read before a stop or restart must not bring the old watch back. + existing.watch?.startedAt === command.startedAt + ? (command.watch ?? undefined) + : existing.watch + : !command.watching + ? undefined + : (existing.watch ?? { + startedAt, + headSha: null, + failedChecks: [], + passed: false, + remarksThrough: startedAt, + remarkIds: [], + conflicting: false, + wakes: 0, + }); + if (watch === existing.watch && links === linked) return thread; + return { + ...thread, + pullRequests: links.map((link) => + link === existing ? withPullRequestWatch(link, watch) : link, + ), + // A user or agent starting or stopping a watch is activity; recorded progress is not. + updatedAt: command.type === "thread.pull-request.watch" ? now : thread.updatedAt, + }; + } case "thread.pull-request.sync": return { ...thread, @@ -2857,16 +3010,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ), ...(command.linkedPullRequest - ? [ - { - ...legacyThreadPullRequestKey(command.linkedPullRequest), - url: command.linkedPullRequest.url, - source: "manual" as const, - linkedAt: DateTime.formatIso(now), - snapshot: null, - stack: null, - }, - ] + ? [legacyPullRequestLink(thread, command.linkedPullRequest, now)] : []), ], }), @@ -2927,6 +3071,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "thread.pull-request.link": case "thread.pull-request.unlink": case "thread.pull-request-link.sync": + case "thread.pull-request.watch": + case "thread.pull-request-watch.sync": case "thread.pull-request.sync": return "thread.pull-request-synced" as const; case "thread.runtime-mode.set": @@ -9130,6 +9276,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "thread.pull-request.link": case "thread.pull-request.unlink": case "thread.pull-request-link.sync": + case "thread.pull-request.watch": case "thread.pull-request.sync": case "thread.title.regeneration.complete": case "thread.runtime-mode.set": @@ -9138,6 +9285,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "provider.switch": yield* dispatchThreadMutation(command, events, effects); break; + case "thread.pull-request-watch.sync": + yield* dispatchPullRequestWatchSync(command, events, effects); + break; case "provider-session.detach": yield* dispatchProviderSessionDetach(command, events, effects); break; diff --git a/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts new file mode 100644 index 000000000000..4eb6ec80ce17 --- /dev/null +++ b/apps/server/src/orchestration-v2/PullRequestWatchReactor.ts @@ -0,0 +1,219 @@ +import { + CommandId, + MessageId, + type OrchestrationV2Notification, + type ThreadPullRequestLink, + type ThreadPullRequestWatch, +} from "@t3tools/contracts"; +import { + normalizeThreadPullRequestKey, + threadPullRequestKeyOf, + visibleThreadPullRequests, +} from "@t3tools/shared/threadPullRequests"; +import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Layer from "effect/Layer"; +import * as Schedule from "effect/Schedule"; +import type * as Scope from "effect/Scope"; + +import * as PullRequestService from "../pullRequest/PullRequestService.ts"; +import { forkParked } from "../serverActivation.ts"; +import * as Orchestrator from "./Orchestrator.ts"; +import * as ProjectionStore from "./ProjectionStore.ts"; +import { evaluatePullRequestWatch, pullRequestWatchMessage } from "./pullRequestWatch.ts"; + +/** Passes in a row that could not read a pull request before its watch ends (one a minute). */ +const READ_FAILURE_LIMIT = 15; + +const logFailure = + (message: string, fields: Record) => + (cause: Cause.Cause): Effect.Effect => + Cause.hasInterruptsOnly(cause) + ? Effect.interrupt + : Effect.logWarning(message, { ...fields, cause }); + +interface WatchTarget { + readonly thread: ProjectionStore.ProjectionThreadPullRequests; + readonly link: ThreadPullRequestLink; + readonly watch: ThreadPullRequestWatch; +} + +const failureKey = ({ thread, link, watch }: WatchTarget) => + `${thread.id} ${threadPullRequestKeyOf(link)} ${watch.startedAt}`; + +function watchesEqual(left: ThreadPullRequestWatch, right: ThreadPullRequestWatch): boolean { + return ( + left.startedAt === right.startedAt && + left.headSha === right.headSha && + left.failedChecks.join("\n") === right.failedChecks.join("\n") && + left.passed === right.passed && + left.remarksThrough === right.remarksThrough && + left.remarkIds.join("\n") === right.remarkIds.join("\n") && + left.conflicting === right.conflicting && + left.wakes === right.wakes + ); +} + +/** + * Wakes a thread's agent when a pull request it watches (`watch_pull_request`) needs a look: + * checks finished on the head commit, someone else commented, or the branch started to + * conflict. One pass a minute reads each watched pull request; settled threads wait until + * they are active again, and a merged or closed pull request ends its watch. + */ +export class PullRequestWatchReactor extends Context.Service< + PullRequestWatchReactor, + { + readonly start: () => Effect.Effect; + /** One pass over every watched pull request. */ + readonly sweep: Effect.Effect; + } +>()("t3/orchestration-v2/PullRequestWatchReactor") {} + +/** @public Service construction is part of the canonical Effect module API. */ +export const make = Effect.gen(function* () { + const engine = yield* Orchestrator.OrchestratorV2; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const pullRequests = yield* PullRequestService.PullRequestService; + const crypto = yield* Crypto.Crypto; + + // Passes in a row that failed, per watch. Kept in memory: a restart only delays the stop. + const readFailures = new Map(); + + // Host-level identity, with the repository as linked, the way pull request sync reads it. + const identityOf = (link: ThreadPullRequestLink) => ({ + host: normalizeThreadPullRequestKey(link).host, + repository: link.repository, + number: link.number, + }); + + /** + * Records what a pass saw, and wakes the agent with it. The orchestrator applies this only + * while the same watch is on, so a stop or restart that lands during the host read wins. + */ + const record = ( + target: WatchTarget, + next: ThreadPullRequestWatch | null, + wake?: { readonly text: string; readonly notification: OrchestrationV2Notification }, + ) => + Effect.gen(function* () { + const uuid = yield* crypto.randomUUIDv4; + yield* engine.dispatch({ + type: "thread.pull-request-watch.sync", + commandId: CommandId.make(`server:pr-watch:${target.thread.id}:${uuid}`), + threadId: target.thread.id, + ...identityOf(target.link), + startedAt: target.watch.startedAt, + watch: next, + ...(wake === undefined + ? {} + : { wake: { ...wake, messageId: MessageId.make(`message:pr-watch:${uuid}`) } }), + }); + }); + + // A watch that cannot read its pull request ends with a wake saying so, rather than showing + // "Watching" while it learns nothing. + const giveUp = (target: WatchTarget) => + record(target, null, { + text: `T3 Code stopped watching pull request #${target.link.number} (${target.link.url}) because it could not read it from the host for ${READ_FAILURE_LIMIT} minutes. Check it yourself, and call watch_pull_request to watch it again.`, + notification: { + source: { kind: "monitor" }, + outcome: "failed", + summary: `#${target.link.number}: stopped watching, could not read it`, + }, + }).pipe(Effect.catch(() => record(target, null))); + + const check = Effect.fn("PullRequestWatchReactor.check")(function* (target: WatchTarget) { + const { thread, link, watch } = target; + const pullRequest = identityOf(link); + // A merged pull request cannot reopen, so its watch ends without a host read, even on a + // settled thread. A closed one can, so the host decides below. + if (link.snapshot?.state === "merged") return yield* record(target, null); + if (thread.settledOverride === "settled" || thread.settledAt !== null) return; + + const reference = { projectId: thread.projectId, ...pullRequest }; + const read = yield* Effect.exit( + Effect.all( + [ + pullRequests.detail({ ...reference, allowStale: false }), + pullRequests.activity(reference), + ], + { concurrency: 2 }, + ), + ); + // Only host reads count towards giving up; a refused wake is not the host's fault. + const key = failureKey(target); + if (Exit.isFailure(read)) { + if (Cause.hasInterruptsOnly(read.cause)) return yield* Effect.failCause(read.cause); + const failures = (readFailures.get(key) ?? 0) + 1; + readFailures.set(key, failures); + // The count stays until the stop lands, so a failed stop is tried again next pass. + if (failures >= READ_FAILURE_LIMIT) { + yield* giveUp(target); + readFailures.delete(key); + } + return yield* Effect.failCause(read.cause); + } + readFailures.delete(key); + const [detail, activity] = read.value; + if (detail.state !== "open") return yield* record(target, null); + + // A degraded read (GitHub's review thread query failed) is truncated with no long thread to + // explain it, and would skip review comments, so remarks wait for a later pass. Replies past + // the first ten of a long review thread are not read. + const degraded = + activity.commentsTruncated && + !activity.reviewThreads.some((reviewThread) => reviewThread.nextCommentsCursor !== undefined); + const report = evaluatePullRequestWatch(watch, detail, degraded ? null : activity.comments); + if (report.changes.length > 0) { + return yield* record( + target, + report.exhausted ? null : report.next, + pullRequestWatchMessage({ + number: link.number, + url: link.url, + baseBranch: detail.baseBranch, + headSha: report.next.headSha, + report, + }), + ); + } + if (!watchesEqual(report.next, watch)) yield* record(target, report.next); + }); + + const sweep = Effect.gen(function* () { + const threads = yield* projections.getThreadsWithPullRequests(); + const targets = threads.flatMap((thread) => + visibleThreadPullRequests(thread.pullRequests ?? []).flatMap((link) => + link.watch === undefined ? [] : [{ thread, link, watch: link.watch }], + ), + ); + const keys = new Set(targets.map(failureKey)); + for (const key of readFailures.keys()) if (!keys.has(key)) readFailures.delete(key); + yield* Effect.forEach( + targets, + (target) => + check(target).pipe( + Effect.catchCause( + logFailure("pull request watch check failed", { + threadId: target.thread.id, + pullRequest: threadPullRequestKeyOf(target.link), + }), + ), + ), + { concurrency: 4, discard: true }, + ); + }).pipe( + Effect.catchCause(logFailure("pull request watch sweep failed", {})), + Effect.withSpan("PullRequestWatchReactor.sweep"), + ); + + const start: PullRequestWatchReactor["Service"]["start"] = () => + forkParked(sweep.pipe(Effect.repeat(Schedule.spaced("1 minute")), Effect.asVoid)); + + return { start, sweep } satisfies PullRequestWatchReactor["Service"]; +}); + +export const layer = Layer.effect(PullRequestWatchReactor, make); diff --git a/apps/server/src/orchestration-v2/pullRequestWatch.test.ts b/apps/server/src/orchestration-v2/pullRequestWatch.test.ts new file mode 100644 index 000000000000..c247d78fd332 --- /dev/null +++ b/apps/server/src/orchestration-v2/pullRequestWatch.test.ts @@ -0,0 +1,217 @@ +import type { + PullRequestCheck, + PullRequestComment, + PullRequestDetail, + ThreadPullRequestWatch, +} from "@t3tools/contracts"; +import { assert, describe, it } from "@effect/vitest"; + +import { + PULL_REQUEST_WATCH_WAKE_LIMIT, + evaluatePullRequestWatch, + pullRequestWatchMessage, +} from "./pullRequestWatch.ts"; + +const STARTED = "2026-10-02T12:00:00.000Z"; + +const watch = (overrides: Partial = {}): ThreadPullRequestWatch => ({ + startedAt: STARTED, + headSha: null, + failedChecks: [], + passed: false, + remarksThrough: STARTED, + remarkIds: [], + conflicting: false, + wakes: 0, + ...overrides, +}); + +const check = (name: string, status: PullRequestCheck["status"]): PullRequestCheck => ({ + name, + status, + description: null, + url: `https://ci.example/${name}`, +}); + +type Detail = Parameters[1]; + +const detail = (overrides: Partial = {}): Detail => ({ + headSha: "aaaaaaaaaa", + checks: [check("lint", "success"), check("test", "pending")], + mergeability: "mergeable", + viewer: "agent-user", + author: { login: "agent-user", name: null, avatarUrl: null }, + ...overrides, +}); + +const remark = ( + login: string, + createdAt: string, + body = "Please rename this.", +): PullRequestComment => ({ + id: `${login}-${createdAt}`, + kind: "review-comment", + author: { login, name: null, avatarUrl: null }, + body, + createdAt, + url: `https://github.com/o/r/pull/1#${login}`, + path: "src/index.ts", + reviewState: null, +}); + +const noRemarks: ReadonlyArray = []; + +describe("evaluatePullRequestWatch", () => { + it("reports each failure at once, even while another check never finishes", () => { + const bot = check("CodeRabbit", "pending"); + const first = detail({ checks: [check("lint", "failure"), check("test", "pending"), bot] }); + const lint = evaluatePullRequestWatch(watch(), first, noRemarks); + assert.deepEqual(lint.changes, [{ kind: "checks-failed", failed: [check("lint", "failure")] }]); + assert.deepEqual(evaluatePullRequestWatch(lint.next, first, noRemarks).changes, []); + + // A different job failing later is news of its own. + const second = detail({ checks: [check("lint", "failure"), check("test", "failure"), bot] }); + const test = evaluatePullRequestWatch(lint.next, second, noRemarks); + assert.deepEqual(test.changes, [{ kind: "checks-failed", failed: [check("test", "failure")] }]); + + // A rerun leaves the list while it runs, so failing again is reported again. + const rerun = evaluatePullRequestWatch(test.next, first, noRemarks); + assert.equal(evaluatePullRequestWatch(rerun.next, second, noRemarks).changes.length, 1); + + // A push reports its failures, even ones that failed between two passes. + const pushed = detail({ ...second, headSha: "bbbbbbbbbb" }); + assert.equal(evaluatePullRequestWatch(test.next, pushed, noRemarks).changes.length, 1); + }); + + it("reports passed once the required checks pass, whatever the others do", () => { + const required = (name: string, status: PullRequestCheck["status"]) => ({ + ...check(name, status), + required: true, + }); + const green = detail({ + checks: [required("test", "success"), required("lint", "success"), check("bot", "pending")], + }); + const passed = evaluatePullRequestWatch(watch(), green, noRemarks); + assert.deepEqual(passed.changes, [{ kind: "checks-passed", count: 2, required: true }]); + assert.deepEqual(evaluatePullRequestWatch(passed.next, green, noRemarks).changes, []); + + // Where nothing is marked required, every check has to pass. + const plain = detail({ checks: [check("test", "success"), check("bot", "pending")] }); + assert.deepEqual(evaluatePullRequestWatch(watch(), plain, noRemarks).changes, []); + }); + + it("keeps remarks for a later pass when the conversation was not read whole", () => { + const comments = [remark("reviewer", "2026-10-02T12:06:00Z")]; + const partial = evaluatePullRequestWatch(watch(), detail(), null); + assert.deepEqual(partial.changes, []); + assert.equal( + evaluatePullRequestWatch(partial.next, detail(), comments).changes[0]?.kind, + "remarks", + ); + }); + + it("reports a remark that shows up late with the same time as a reported one", () => { + const first = remark("reviewer", "2026-10-02T12:06:00Z"); + const late = { ...remark("bot", "2026-10-02T12:06:00Z"), id: "late" }; + const reported = evaluatePullRequestWatch(watch(), detail(), [first]); + const again = evaluatePullRequestWatch(reported.next, detail(), [first, late]); + assert.deepEqual(again.changes, [{ kind: "remarks", remarks: [late] }]); + assert.deepEqual(again.next.remarkIds, [first.id, "late"]); + }); + + it("does not treat a failed check read as a rerun", () => { + const failed = detail({ checks: [check("lint", "failure")] }); + const reported = evaluatePullRequestWatch(watch(), failed, noRemarks); + assert.equal(reported.changes.length, 1); + const unreadable = evaluatePullRequestWatch(reported.next, detail({ checks: [] }), noRemarks); + assert.deepEqual(evaluatePullRequestWatch(unreadable.next, failed, noRemarks).changes, []); + }); + + it("wakes for the pull request's author when the agent is someone else", () => { + const contributor = detail({ author: { login: "contributor", name: null, avatarUrl: null } }); + const reply = remark("contributor", "2026-10-02T12:06:00Z"); + assert.deepEqual(evaluatePullRequestWatch(watch(), contributor, [reply]).changes, [ + { kind: "remarks", remarks: [reply] }, + ]); + // Without a viewer, the author is taken to be the agent. + const noViewer = detail({ viewer: undefined, author: contributor.author }); + assert.deepEqual(evaluatePullRequestWatch(watch(), noViewer, [reply]).changes, []); + }); + + it("reports remarks from others once and never the agent's own", () => { + const comments = [ + remark("agent-user", "2026-10-02T12:05:00Z", "Fixed in the latest push."), + remark("macroscope-app[bot]", "2026-10-02T12:06:00Z"), + remark("reviewer", "2026-10-02T11:00:00Z", "Older than the watch."), + ]; + const report = evaluatePullRequestWatch(watch(), detail(), comments); + assert.deepEqual(report.changes, [{ kind: "remarks", remarks: [comments[1]!] }]); + assert.equal(report.next.remarksThrough, "2026-10-02T12:06:00Z"); + assert.deepEqual(evaluatePullRequestWatch(report.next, detail(), comments).changes, []); + }); + + it("reports a conflict once, until the branch is clean again", () => { + const conflicting = detail({ mergeability: "conflicting" }); + const first = evaluatePullRequestWatch(watch(), conflicting, noRemarks); + assert.deepEqual(first.changes, [{ kind: "conflicting" }]); + // GitHub answers "unknown" while it recomputes after a push; that is not a resolution. + const recomputing = evaluatePullRequestWatch( + first.next, + detail({ mergeability: "unknown" }), + noRemarks, + ); + assert.deepEqual( + evaluatePullRequestWatch(recomputing.next, conflicting, noRemarks).changes, + [], + ); + const clean = evaluatePullRequestWatch(first.next, detail(), noRemarks); + assert.deepEqual(evaluatePullRequestWatch(clean.next, conflicting, noRemarks).changes, [ + { kind: "conflicting" }, + ]); + }); + + it("does not spend the comment wake limit on check results", () => { + const tired = watch({ headSha: "aaaaaaaaaa", wakes: PULL_REQUEST_WATCH_WAKE_LIMIT - 1 }); + const result = evaluatePullRequestWatch( + tired, + detail({ checks: [check("lint", "failure")] }), + noRemarks, + ); + assert.isFalse(result.exhausted); + assert.equal(result.next.wakes, 0); + }); + + it("stops after the wake limit unless the head moves", () => { + const comments = [remark("reviewer", "2026-10-02T12:10:00Z")]; + const tired = watch({ headSha: "aaaaaaaaaa", wakes: PULL_REQUEST_WATCH_WAKE_LIMIT - 1 }); + assert.isTrue(evaluatePullRequestWatch(tired, detail(), comments).exhausted); + const pushed = evaluatePullRequestWatch(tired, detail({ headSha: "cccccccccc" }), comments); + assert.isFalse(pushed.exhausted); + assert.equal(pushed.next.wakes, 1); + }); +}); + +describe("pullRequestWatchMessage", () => { + it("tells the agent what changed and marks failures for the timeline", () => { + const report = evaluatePullRequestWatch( + watch(), + detail({ checks: [check("lint", "failure")] }), + [remark("reviewer", "2026-10-02T12:10:00Z", "Needs a test.")], + ); + const message = pullRequestWatchMessage({ + number: 12, + url: "https://github.com/o/r/pull/12", + baseBranch: "main", + headSha: report.next.headSha, + report, + }); + assert.include(message.text, "- Checks failed on aaaaaaa:\n - lint https://ci.example/lint"); + assert.include(message.text, ' - reviewer on src/index.ts: "Needs a test."'); + assert.include(message.text, "unwatch_pull_request"); + assert.deepEqual(message.notification, { + source: { kind: "monitor" }, + outcome: "failed", + summary: "#12: checks failed, new comments", + }); + }); +}); diff --git a/apps/server/src/orchestration-v2/pullRequestWatch.ts b/apps/server/src/orchestration-v2/pullRequestWatch.ts new file mode 100644 index 000000000000..af9ce27636ee --- /dev/null +++ b/apps/server/src/orchestration-v2/pullRequestWatch.ts @@ -0,0 +1,209 @@ +import type { + OrchestrationV2Notification, + PullRequestCheck, + PullRequestComment, + PullRequestDetail, + ThreadPullRequestWatch, +} from "@t3tools/contracts"; + +/** + * Wakes in a row that bring only comments. Check, conflict, or push news resets the count, so + * this only stops a chatty bot looping an agent that is replying to it. + */ +export const PULL_REQUEST_WATCH_WAKE_LIMIT = 10; +const LISTED_ITEMS = 10; +const SNIPPET_LENGTH = 200; + +export type PullRequestWatchChange = + | { readonly kind: "checks-failed"; readonly failed: ReadonlyArray } + | { readonly kind: "checks-passed"; readonly count: number; readonly required: boolean } + | { readonly kind: "remarks"; readonly remarks: ReadonlyArray } + | { readonly kind: "conflicting" }; + +export interface PullRequestWatchReport { + /** What the agent has not been told yet. Empty means no wake. */ + readonly changes: ReadonlyArray; + /** The watch to record, whether or not anything is reported. */ + readonly next: ThreadPullRequestWatch; + /** This report spends the last wake before the limit, so watching stops after it. */ + readonly exhausted: boolean; +} + +// "action-required" is a finished check that needs someone, so the agent hears about it. +const isFailedCheck = (check: PullRequestCheck) => + check.status === "failure" || check.status === "cancelled" || check.status === "action-required"; + +/** + * Compares a watched pull request with what its agent was last told. Each check is reported as + * soon as it fails, so a check that never finishes (an advisory review bot) cannot hold the + * news back. "Passed" is reported once the checks the base branch requires all passed, or all + * checks where the host marks none required. Remarks count when someone other than the agent's + * own account wrote them, so its own replies never wake it. `remarks` is null when the + * conversation could not be read; remarks then wait for a later pass. + */ +export function evaluatePullRequestWatch( + watch: ThreadPullRequestWatch, + detail: Pick, + remarks: ReadonlyArray | null, +): PullRequestWatchReport { + const changes: Array = []; + const headSha = detail.headSha ?? null; + const headMoved = headSha !== watch.headSha; + + // An empty list keeps the last state: a host can answer with one when its check read fails. + let failedChecks = headMoved ? [] : watch.failedChecks; + let passed = headMoved ? false : watch.passed; + if (detail.checks.length > 0) { + const failed = detail.checks.filter(isFailedCheck); + const newlyFailed = failed.filter((check) => !failedChecks.includes(check.name)); + if (newlyFailed.length > 0) changes.push({ kind: "checks-failed", failed: newlyFailed }); + // A check that runs again leaves the list, so a rerun that fails again is reported. + failedChecks = failed.map((check) => check.name); + + const required = detail.checks.filter((check) => check.required === true); + const gate = required.length > 0 ? required : detail.checks; + const passedNow = gate.every((check) => check.status !== "pending" && !isFailedCheck(check)); + if (passedNow && !passed) { + changes.push({ kind: "checks-passed", count: gate.length, required: required.length > 0 }); + } + passed = passedNow; + } + + const own = (detail.viewer ?? detail.author?.login)?.toLowerCase(); + const through = Date.parse(watch.remarksThrough); + // GitHub times are per second, so remarks at the boundary time are told apart by ID. + const fresh = (remarks ?? []).filter((remark) => { + const at = Date.parse(remark.createdAt); + return ( + (at > through || (at === through && !watch.remarkIds.includes(remark.id))) && + remark.author?.login.toLowerCase() !== own + ); + }); + if (fresh.length > 0) changes.push({ kind: "remarks", remarks: fresh }); + const latest = Math.max(through, ...fresh.map((remark) => Date.parse(remark.createdAt))); + const atLatest = fresh.filter((remark) => Date.parse(remark.createdAt) === latest); + const remarksThrough = latest === through ? watch.remarksThrough : atLatest[0]!.createdAt; + const remarkIds = [ + ...(latest === through ? watch.remarkIds : []), + ...atLatest.map((remark) => remark.id), + ]; + + if (detail.mergeability === "conflicting" && !watch.conflicting) { + changes.push({ kind: "conflicting" }); + } + // "unknown" is GitHub still computing after a push; only a clean answer clears a conflict. + const conflicting = + detail.mergeability === "unknown" ? watch.conflicting : detail.mergeability === "conflicting"; + + const commentsOnly = changes.length > 0 && changes.every((change) => change.kind === "remarks"); + const progress = headMoved || (changes.length > 0 && !commentsOnly); + const wakes = (progress ? 0 : watch.wakes) + (commentsOnly ? 1 : 0); + return { + changes, + next: { + startedAt: watch.startedAt, + headSha, + failedChecks, + passed, + remarksThrough, + remarkIds, + conflicting, + wakes, + }, + exhausted: commentsOnly && wakes >= PULL_REQUEST_WATCH_WAKE_LIMIT, + }; +} + +function snippet(body: string): string { + const text = body + .replaceAll(//g, " ") + .replaceAll(/\s+/g, " ") + .trim(); + return text.length <= SNIPPET_LENGTH ? text : `${text.slice(0, SNIPPET_LENGTH - 3)}...`; +} + +function listed(items: ReadonlyArray, line: (item: T) => string): Array { + const lines = items.slice(0, LISTED_ITEMS).map(line); + if (items.length > LISTED_ITEMS) lines.push(` - and ${items.length - LISTED_ITEMS} more`); + return lines; +} + +function changeLines( + change: PullRequestWatchChange, + context: { readonly baseBranch: string; readonly commit: string }, +): Array { + switch (change.kind) { + case "checks-failed": + return [ + `- Checks failed${context.commit}:`, + ...listed( + change.failed, + (check) => + ` - ${check.name}${check.status === "failure" ? "" : ` (${check.status})`}${check.url ? ` ${check.url}` : ""}`, + ), + ]; + case "checks-passed": + return [ + `- All ${change.count} ${change.required ? "required " : ""}${change.count === 1 ? "check" : "checks"} passed${context.commit}.`, + ]; + case "remarks": + return [ + `- ${change.remarks.length} new ${change.remarks.length === 1 ? "comment" : "comments"}:`, + ...listed(change.remarks, (remark) => { + const where = remark.path === null ? "" : ` on ${remark.path}`; + const body = snippet(remark.body); + const said = body.length === 0 ? (remark.reviewState ?? "reviewed") : `"${body}"`; + return ` - ${remark.author?.login ?? "someone"}${where}: ${said}${remark.url ? ` ${remark.url}` : ""}`; + }), + ]; + case "conflicting": + return [`- The branch now conflicts with ${context.baseBranch}.`]; + } +} + +const SUMMARY: Record = { + "checks-failed": "checks failed", + "checks-passed": "checks passed", + remarks: "new comments", + conflicting: "merge conflict", +}; + +/** The wake the agent reads and the timeline notification the user sees. */ +export function pullRequestWatchMessage(input: { + readonly number: number; + readonly url: string; + readonly baseBranch: string; + readonly headSha: string | null; + readonly report: PullRequestWatchReport; +}): { readonly text: string; readonly notification: OrchestrationV2Notification } { + const { changes, exhausted } = input.report; + const context = { + baseBranch: input.baseBranch, + commit: input.headSha === null ? "" : ` on ${input.headSha.slice(0, 7)}`, + }; + const text = [ + `Update on pull request #${input.number} (${input.url}), which T3 Code is watching for you:`, + ...changes.flatMap((change) => changeLines(change, context)), + "", + exhausted + ? `T3 Code stopped watching after ${PULL_REQUEST_WATCH_WAKE_LIMIT} comment-only updates in a row. Call watch_pull_request to watch it again.` + : "Look into each item and act on it as your task requires. T3 Code keeps watching and wakes you on the next change, so end your turn when you are done. Call unwatch_pull_request when you no longer need updates.", + ].join("\n"); + const failed = changes.some( + (change) => change.kind === "checks-failed" || change.kind === "conflicting", + ); + const summary = changes.map((change) => SUMMARY[change.kind]); + if (exhausted) summary.push("stopped watching"); + return { + text, + notification: { + source: { kind: "monitor" }, + outcome: failed + ? "failed" + : changes.every((change) => change.kind === "checks-passed") + ? "completed" + : "updated", + summary: `#${input.number}: ${summary.join(", ")}`, + }, + }; +} diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index d077f7e20910..2c5235c767d7 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -17,6 +17,7 @@ import { type ModelSelection, type OrchestrationV2Run, ProjectId, + type PullRequestDetail, ProviderDriverKind, ProviderInstanceId, ProviderThreadId, @@ -58,6 +59,9 @@ import * as EffectOutbox from "./EffectOutbox.ts"; import * as EventSink from "./EventSink.ts"; import * as ProviderRuntimeRecoveryService from "./ProviderRuntimeRecoveryService.ts"; import * as ProjectionMaintenance from "./ProjectionMaintenance.ts"; +import * as ProjectionStore from "./ProjectionStore.ts"; +import * as PullRequestWatchReactor from "./PullRequestWatchReactor.ts"; +import * as PullRequestService from "../pullRequest/PullRequestService.ts"; import * as ProjectStore from "./ProjectStore.ts"; import type { ProviderAdapterV2SessionRuntime, ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import * as ProviderSessionManager from "./ProviderSessionManager.ts"; @@ -212,6 +216,7 @@ const TestLayer = Layer.mergeAll( OrchestrationV2LayerLive, OrchestrationV2EventSinkLayerLive, ProjectStore.layer, + ProjectionStore.layer, EffectOutbox.layer, ThreadCommandExecutor.layer, ).pipe( @@ -2148,6 +2153,298 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }), ); + it.effect("starts, records, and stops a pull request watch", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const maintenance = yield* ProjectionMaintenance.ProjectionMaintenanceV2; + const threadId = ThreadId.make("runtime-pull-request-watch"); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("pr-watch-create"), + threadId, + projectId: ProjectId.make("pr-watch-project"), + title: "Watch", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + }); + const key = { host: "github.com", repository: "pingdotgg/t3code", number: 7 }; + const url = "https://github.com/pingdotgg/t3code/pull/7"; + const watchOf = Effect.map( + orchestrator.getThreadShell(threadId), + (thread) => thread?.pullRequests?.[0]?.watch, + ); + + // Watching an unlinked pull request links it in the same command. + yield* orchestrator.dispatch({ + type: "thread.pull-request.watch", + commandId: CommandId.make("pr-watch-start"), + threadId, + ...key, + watching: true, + link: { url, source: "agent" }, + }); + assert.equal( + (yield* orchestrator.getThreadShell(threadId))?.pullRequests?.[0]?.source, + "agent", + ); + const started = yield* watchOf; + assert.isDefined(started); + if (started === undefined) return; + + // A legacy client re-linking the same pull request keeps its watch. + yield* orchestrator.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make("pr-watch-legacy-relink"), + threadId, + linkedPullRequest: { projectId: ProjectId.make("pr-watch-project"), ...key, url }, + }); + assert.deepEqual(yield* watchOf, started); + + const recorded = { ...started, headSha: "abc123", failedChecks: ["lint"], wakes: 1 }; + yield* orchestrator.dispatch({ + type: "thread.pull-request-watch.sync", + commandId: CommandId.make("pr-watch-record"), + threadId, + ...key, + startedAt: started.startedAt, + watch: recorded, + }); + assert.deepEqual(yield* watchOf, recorded); + assert.isTrue((yield* maintenance.rebuild).valid); + assert.deepEqual(yield* watchOf, recorded); + + yield* orchestrator.dispatch({ + type: "thread.pull-request.watch", + commandId: CommandId.make("pr-watch-stop"), + threadId, + ...key, + watching: false, + }); + // A wake read before the stop must neither wake the agent nor bring the watch back. + const late = yield* orchestrator + .dispatch({ + type: "thread.pull-request-watch.sync", + commandId: CommandId.make("pr-watch-late-record"), + threadId, + ...key, + startedAt: started.startedAt, + watch: { ...recorded, wakes: 2 }, + wake: { + messageId: MessageId.make("pr-watch-late-wake"), + text: "Update", + notification: { source: { kind: "monitor" }, outcome: "updated", summary: "#7" }, + }, + }) + .pipe(Effect.flip); + assert.equal(late._tag, "OrchestratorDispatchError"); + assert.isUndefined(yield* watchOf); + const { messages } = yield* orchestrator.getThreadRecords(threadId, ["messages"]); + assert.deepEqual(messages, []); + }), + ); + + it.effect("ends a watch it cannot read, and tells the agent", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const threadId = ThreadId.make("runtime-pull-request-watch-unreadable"); + const projectId = ProjectId.make("pr-watch-unreadable-project"); + yield* seedProject({ + projectId, + title: "Watch unreadable", + workspaceRoot: "/workspace/watch-unreadable", + defaultModelSelection: null, + createdAt: "2026-10-01T00:00:00.000Z", + }); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("pr-watch-unreadable-create"), + threadId, + projectId, + title: "Watch unreadable", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + }); + yield* orchestrator.dispatch({ + type: "thread.pull-request.watch", + commandId: CommandId.make("pr-watch-unreadable-start"), + threadId, + host: "github.com", + repository: "pingdotgg/t3code", + number: 8, + watching: true, + link: { url: "https://github.com/pingdotgg/t3code/pull/8", source: "agent" }, + }); + const reactor = yield* PullRequestWatchReactor.make.pipe( + Effect.provide( + Layer.mergeAll( + NodeServices.layer, + Layer.mock(PullRequestService.PullRequestService)({ + detail: () => Effect.die("host unreachable"), + activity: () => Effect.die("host unreachable"), + }), + ), + ), + ); + for (let pass = 0; pass < 15; pass += 1) yield* reactor.sweep; + + const thread = yield* orchestrator.getThreadShell(threadId); + assert.isUndefined(thread?.pullRequests?.[0]?.watch); + const { messages } = yield* orchestrator.getThreadRecords(threadId, ["messages"]); + assert.deepEqual( + messages.flatMap((message) => message.notification?.summary ?? []), + ["#8: stopped watching, could not read it"], + ); + }), + ); + + it.effect("wakes a watched thread once for failed checks and a review comment", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const threadId = ThreadId.make("runtime-pull-request-watch-wake"); + const projectId = ProjectId.make("pr-watch-wake-project"); + yield* seedProject({ + projectId, + title: "Watch wake", + workspaceRoot: "/workspace/watch", + defaultModelSelection: null, + createdAt: "2026-10-01T00:00:00.000Z", + }); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("pr-watch-wake-create"), + threadId, + projectId, + title: "Watch wake", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + }); + const key = { host: "github.com", repository: "pingdotgg/t3code", number: 7 }; + const url = "https://github.com/pingdotgg/t3code/pull/7"; + yield* orchestrator.dispatch({ + type: "thread.pull-request.link", + commandId: CommandId.make("pr-watch-wake-link"), + threadId, + ...key, + url, + source: "agent", + }); + yield* orchestrator.dispatch({ + type: "thread.pull-request.watch", + commandId: CommandId.make("pr-watch-wake-start"), + threadId, + ...key, + watching: true, + }); + + const at = "2026-10-02T12:00:00.000Z"; + const detail: PullRequestDetail = { + provider: "github", + capabilities: { + diff: true, + comment: true, + actions: [], + mergeMethods: [], + search: false, + review: { inlineComment: false, reply: false, resolve: false, verdicts: [] }, + reviewers: { request: false, listCandidates: false }, + }, + viewerPermissions: { + actions: [], + comment: true, + resolve: true, + verdicts: [], + requestReviewers: false, + }, + projectId, + projectTitle: "Watch wake", + workspaceRoot: "/workspace/watch", + repository: key.repository, + number: key.number, + title: "Watched pull request", + body: "", + url, + author: { login: "agent-user", name: null, avatarUrl: null }, + state: "open", + isDraft: false, + mergeability: "mergeable", + additions: 1, + deletions: 0, + changedFiles: 1, + headBranch: "feature", + headSha: "abc1234def", + baseBranch: "main", + createdAt: at, + updatedAt: at, + mergedAt: null, + closedAt: null, + reviewers: [], + labels: [], + checks: [{ name: "lint", status: "failure", description: null, url: null }], + mergeCapabilities: { merge: true, squash: true, rebase: true }, + viewer: "agent-user", + }; + const reactor = yield* PullRequestWatchReactor.make.pipe( + Effect.provide( + Layer.mergeAll( + NodeServices.layer, + Layer.mock(PullRequestService.PullRequestService)({ + detail: () => Effect.succeed(detail), + activity: () => + Effect.succeed({ + comments: [ + { + id: "review-1", + kind: "review-comment", + author: { login: "reviewer", name: null, avatarUrl: null }, + body: "One more thing.", + createdAt: "2999-01-01T00:00:00.000Z", + url: null, + path: "src/index.ts", + reviewState: null, + }, + ], + commentCount: 1, + commentsTruncated: false, + reviewThreads: [], + commits: [], + }), + }), + ), + ), + ); + yield* reactor.sweep; + yield* reactor.sweep; + + const { messages } = yield* orchestrator.getThreadRecords(threadId, ["messages"]); + assert.deepEqual( + messages.flatMap((message) => + message.notification === undefined ? [] : [message.notification.summary], + ), + ["#7: checks failed, new comments"], + ); + const watch = (yield* orchestrator.getThreadShell(threadId))?.pullRequests?.[0]?.watch; + assert.deepEqual( + { headSha: watch?.headSha, failedChecks: watch?.failedChecks, wakes: watch?.wakes }, + { headSha: "abc1234def", failedChecks: ["lint"], wakes: 0 }, + ); + }), + ); + it.effect("persists rejected command receipts across retries", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; diff --git a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts index 4cc834dddd62..b9cad107409a 100644 --- a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts +++ b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts @@ -153,6 +153,7 @@ function commandThreadIds(command: OrchestrationV2Command): ReadonlyArray -When the t3-code MCP server exposes link_pull_request, you must use it to register every pull request you create or work on for this thread. Call link_pull_request with the full PR URL immediately after creating a PR or starting work on an existing PR. For a stack, call it for every layer, not just the current branch or the top PR. This applies when creating or updating PRs through gh, gh stack, another CLI, or the host API: those operations do not register the PRs with this thread. Linking an already-linked PR is safe. Before finishing PR work, call list_thread_pull_requests and link any PR from your work that is missing. Do not link unrelated PRs mentioned only as background. If a linking call fails, report that failure instead of claiming the PR is linked. +When the t3-code MCP server exposes link_pull_request, you must use it to register every pull request you create or work on for this thread. Call link_pull_request with the full PR URL immediately after creating a PR or starting work on an existing PR. For a stack, call it for every layer, not just the current branch or the top PR. This applies when creating or updating PRs through gh, gh stack, another CLI, or the host API: those operations do not register the PRs with this thread. Linking an already-linked PR is safe. Before finishing PR work, call list_thread_pull_requests and link any PR from your work that is missing. Do not link unrelated PRs mentioned only as background. If a linking call fails, report that failure instead of claiming the PR is linked. When asked to monitor, watch, or babysit a PR and watch_pull_request is available, call it and end your turn: T3 Code wakes you when checks finish, someone else comments, or the branch conflicts, so do not poll or run your own watcher. `; /** diff --git a/apps/server/src/pullRequest/GitHubPullRequestCli.ts b/apps/server/src/pullRequest/GitHubPullRequestCli.ts index 020b571fe248..f40cff8f50bf 100644 --- a/apps/server/src/pullRequest/GitHubPullRequestCli.ts +++ b/apps/server/src/pullRequest/GitHubPullRequestCli.ts @@ -52,7 +52,7 @@ import { decodePullRequestActivityJson, decodePullRequestDetailJson, decodePullRequestCoreJson, - PULL_REQUEST_CORE_GRAPHQL_QUERY, + pullRequestCoreGraphQlQuery, type GitHubPullRequestCore, type GitHubPullRequestSummary, decodePullRequestPreviewJson, @@ -1595,7 +1595,7 @@ export const make = Effect.gen(function* () { ["-F", `number=${input.number}`], ["-f", `headRef=refs/pull/${input.number}/head`], ], - query: PULL_REQUEST_CORE_GRAPHQL_QUERY, + query: pullRequestCoreGraphQlQuery(input.host), decode: decodePullRequestCoreJson, }), ), diff --git a/apps/server/src/pullRequest/PullRequestProvider.ts b/apps/server/src/pullRequest/PullRequestProvider.ts index ed4312c125c7..cc370637c601 100644 --- a/apps/server/src/pullRequest/PullRequestProvider.ts +++ b/apps/server/src/pullRequest/PullRequestProvider.ts @@ -213,6 +213,8 @@ export interface ProviderChangeRequestStat { } export interface ProviderChangeRequestDetail extends ProviderChangeRequest { + /** The head commit, where the host's detail read reports it. */ + readonly headSha?: string | null; readonly body: string; readonly changedFiles: number; readonly mergedAt: string | null; diff --git a/apps/server/src/pullRequest/PullRequestService.ts b/apps/server/src/pullRequest/PullRequestService.ts index 44c175ab652f..293318f91c4b 100644 --- a/apps/server/src/pullRequest/PullRequestService.ts +++ b/apps/server/src/pullRequest/PullRequestService.ts @@ -1697,6 +1697,7 @@ export const make = Effect.gen(function* () { ...(changeRequest.headRepositoryNameWithOwner === undefined ? {} : { headRepositoryNameWithOwner: changeRequest.headRepositoryNameWithOwner }), + ...(changeRequest.headSha ? { headSha: changeRequest.headSha } : {}), baseBranch: changeRequest.baseBranch, createdAt: changeRequest.createdAt, updatedAt: changeRequest.updatedAt, diff --git a/apps/server/src/pullRequest/gitHubPullRequestJson.test.ts b/apps/server/src/pullRequest/gitHubPullRequestJson.test.ts index 12239a2af0ee..1393ab6e6b6a 100644 --- a/apps/server/src/pullRequest/gitHubPullRequestJson.test.ts +++ b/apps/server/src/pullRequest/gitHubPullRequestJson.test.ts @@ -26,6 +26,7 @@ import { decodeWorkflowRunApprovalsJson, reviewThreadConversation, REVIEW_THREADS_GRAPHQL_QUERY, + pullRequestCoreGraphQlQuery, pullRequestSearchGraphQlQuery, } from "./gitHubPullRequestJson.ts"; @@ -285,6 +286,25 @@ describe("pull request detail decoding", () => { ]); }); + it("keeps what branch protection requires, and asks for it on github.com only", () => { + const raw = JSON.parse(detailJson) as Record; + const detail = expectSuccess( + decodePullRequestDetailJson( + JSON.stringify({ + ...raw, + statusCheckRollup: [ + { __typename: "CheckRun", name: "test", status: "IN_PROGRESS", isRequired: true }, + { __typename: "StatusContext", context: "bot", state: "PENDING", isRequired: false }, + { __typename: "StatusContext", context: "legacy", state: "SUCCESS" }, + ], + }), + ), + ); + expect(detail.checks.map((check) => check.required)).toEqual([true, false, undefined]); + expect(pullRequestCoreGraphQlQuery("github.com")).toContain("isRequired"); + expect(pullRequestCoreGraphQlQuery("github.example.com")).not.toContain("isRequired"); + }); + it("keeps a workflow waiting for approval out of the passing state", () => { const raw = JSON.parse(detailJson) as Record; const detail = expectSuccess( diff --git a/apps/server/src/pullRequest/gitHubPullRequestJson.ts b/apps/server/src/pullRequest/gitHubPullRequestJson.ts index 8e3a9ee330b8..d1600754b50b 100644 --- a/apps/server/src/pullRequest/gitHubPullRequestJson.ts +++ b/apps/server/src/pullRequest/gitHubPullRequestJson.ts @@ -93,6 +93,8 @@ const RawCheckSchema = Schema.Struct({ workflowName: Schema.optional(Schema.NullOr(Schema.String)), startedAt: Schema.optional(Schema.NullOr(Schema.String)), completedAt: Schema.optional(Schema.NullOr(Schema.String)), + /** Branch protection requires this check; read by the detail query on github.com only. */ + isRequired: Schema.optional(Schema.NullOr(Schema.Boolean)), }); const RawListItemSchema = Schema.Struct({ @@ -706,8 +708,15 @@ export const PULL_REQUEST_LIST_JSON_FIELDS = export const PULL_REQUEST_DETAIL_JSON_FIELDS = `${PULL_REQUEST_LIST_JSON_FIELDS},body,changedFiles,closedAt,isCrossRepository,headRepositoryOwner,headRefOid,autoMergeRequest`; -/** Pull refs let the comparison share the detail read without first resolving a fork branch. */ -export const PULL_REQUEST_CORE_GRAPHQL_QUERY = `query($owner: String!, $name: String!, $number: Int!, $headRef: String!) { +/** + * Pull refs let the comparison share the detail read without first resolving a fork branch. + * `isRequired` is asked for on github.com only: an older Enterprise server may not know it, and + * an unknown field fails the whole read. + */ +export const pullRequestCoreGraphQlQuery = (host: string) => { + const required = + host.toLowerCase() === "github.com" ? " isRequired(pullRequestNumber: $number)" : ""; + return `query($owner: String!, $name: String!, $number: Int!, $headRef: String!) { repository(owner: $owner, name: $name) { mergeCommitAllowed squashMergeAllowed rebaseMergeAllowed viewerPermission pullRequest(number: $number) { @@ -727,9 +736,9 @@ export const PULL_REQUEST_CORE_GRAPHQL_QUERY = `query($owner: String!, $name: St nodes { commit { statusCheckRollup { contexts(first: 100) { nodes { __typename - ... on StatusContext { context state targetUrl createdAt description } + ... on StatusContext { context state targetUrl createdAt description${required} } ... on CheckRun { - name status conclusion startedAt completedAt detailsUrl + name status conclusion startedAt completedAt detailsUrl${required} checkSuite { workflowRun { workflow { name } } } } } @@ -739,6 +748,7 @@ export const PULL_REQUEST_CORE_GRAPHQL_QUERY = `query($owner: String!, $name: St } } }`; +}; export const PULL_REQUEST_PREVIEW_GRAPHQL_QUERY = `query($owner: String!, $name: String!, $number: Int!) { repository(owner: $owner, name: $name) { @@ -1468,6 +1478,7 @@ function toCheckEntries( status: toCheckStatus(check), description: trimmed(check.description), url: trimmed(check.detailsUrl) ?? trimmed(check.targetUrl), + ...(typeof check.isRequired === "boolean" ? { required: check.isRequired } : {}), }, workflowName: trimmed(check.workflowName), at: realTimestamp(check.completedAt) ?? realTimestamp(check.startedAt), diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 98039e56642a..9e9d5c5e3064 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -4,6 +4,7 @@ import * as Random from "effect/Random"; import * as Semaphore from "effect/Semaphore"; import * as StorageCleanup from "./storageCleanup.ts"; import * as PullRequestSyncReactor from "./orchestration-v2/PullRequestSyncReactor.ts"; +import * as PullRequestWatchReactor from "./orchestration-v2/PullRequestWatchReactor.ts"; // @effect-diagnostics nodeBuiltinImport:off import * as NodeHttp from "node:http"; @@ -521,6 +522,16 @@ const RuntimeCoreDependenciesBaseLive = Layer.mergeAll( Layer.provide(PullRequestServiceLive), Layer.provide(ProjectionStoreV2.layer), ), + Layer.effectDiscard( + Effect.gen(function* () { + const service = yield* PullRequestWatchReactor.PullRequestWatchReactor; + yield* service.start(); + }), + ).pipe( + Layer.provide(PullRequestWatchReactor.layer), + Layer.provide(PullRequestServiceLive), + Layer.provide(ProjectionStoreV2.layer), + ), // Subscribes to `account.rate-limits.updated` so usage bars track live // telemetry instead of waiting for the next status probe. ProviderUsageLimitsIngestionLive, diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 5dec1b6d411d..b5113b5f41a3 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -3630,6 +3630,8 @@ function toolGroupSummaryIconName( case "link-pr": case "unlink-pr": case "list-prs": + case "watch-pr": + case "unwatch-pr": return "pull-request"; case "read": return "eye"; diff --git a/apps/web/src/components/pullRequest/ThreadPullRequestsPanel.tsx b/apps/web/src/components/pullRequest/ThreadPullRequestsPanel.tsx index 3dd1bcbefaff..643c17d4b775 100644 --- a/apps/web/src/components/pullRequest/ThreadPullRequestsPanel.tsx +++ b/apps/web/src/components/pullRequest/ThreadPullRequestsPanel.tsx @@ -3,7 +3,14 @@ import { resolveThreadPullRequestChains, visibleThreadPullRequests, } from "@t3tools/shared/threadPullRequests"; -import { ArrowUpRightIcon, LinkIcon, MoreHorizontalIcon, PlusIcon } from "lucide-react"; +import { + ArrowUpRightIcon, + EyeIcon, + EyeOffIcon, + LinkIcon, + MoreHorizontalIcon, + PlusIcon, +} from "lucide-react"; import { useCallback, useMemo } from "react"; import { writeTextToClipboard } from "~/hooks/useCopyToClipboard"; @@ -68,14 +75,19 @@ function LinkRow({ line, threadRef, onUnlink, + onSetWatching, }: { line: PullRequestListLine; threadRef: ScopedThreadRef; onUnlink: (link: ThreadPullRequestLink) => void; + /** Null when the environment cannot watch pull requests. */ + onSetWatching: ((link: ThreadPullRequestLink, watching: boolean) => void) | null; }) { const openPrLink = useOpenPrLink(threadRef); const { link, depth, stack } = line; const snapshot = link.snapshot; + const open = snapshot === null || snapshot.state === "open"; + const watching = link.watch !== undefined; return (
- {snapshot.checksState ? : null} - {snapshot.reviewDecision ? ( + {watching ? ( + + }> + + + + Watching: the agent wakes when checks finish, someone comments, or the branch + conflicts + + + ) : null} + {snapshot?.checksState ? : null} + {snapshot?.reviewDecision ? ( ) : null} @@ -220,6 +243,12 @@ function LinkRow({ Open + {onSetWatching !== null && open ? ( + onSetWatching(link, !watching)}> + {watching ? : } + {watching ? "Stop watching" : "Watch for changes"} + + ) : null} onUnlink(link)}> {link.source === "stack" ? "Dismiss from thread" : "Unlink from thread"} @@ -248,6 +277,10 @@ function EnabledThreadPullRequestsPanel({ threadRef }: { threadRef: ScopedThread const thread = useThreadShell(threadRef); const openLinkDialog = useCallback(() => openLinkPullRequestDialog(threadRef), [threadRef]); const unlink = useAtomCommand(threadEnvironment.unlinkPullRequest, { reportFailure: true }); + const watch = useAtomCommand(threadEnvironment.watchPullRequest, { reportFailure: true }); + const supportsWatch = + useServerConfigs().get(threadRef.environmentId)?.environment.capabilities + .threadPullRequestWatch === true; const links = useMemo(() => visibleThreadPullRequests(thread?.pullRequests ?? []), [thread]); const lines = useMemo(() => pullRequestListLines(resolveThreadPullRequestChains(links)), [links]); const handleUnlink = useCallback( @@ -264,6 +297,21 @@ function EnabledThreadPullRequestsPanel({ threadRef }: { threadRef: ScopedThread }, [threadRef, unlink], ); + const handleSetWatching = useCallback( + (link: ThreadPullRequestLink, watching: boolean) => { + void watch({ + environmentId: threadRef.environmentId, + input: { + threadId: threadRef.threadId, + host: link.host, + repository: link.repository, + number: link.number, + watching, + }, + }); + }, + [threadRef, watch], + ); const openCount = useMemo( () => links.filter((link) => link.snapshot === null || link.snapshot.state === "open").length, [links], @@ -304,6 +352,7 @@ function EnabledThreadPullRequestsPanel({ threadRef }: { threadRef: ScopedThread line={line} threadRef={threadRef} onUnlink={handleUnlink} + onSetWatching={supportsWatch ? handleSetWatching : null} /> ))}
diff --git a/docs/user/source-control.md b/docs/user/source-control.md index 3d068a232be7..64beb946c408 100644 --- a/docs/user/source-control.md +++ b/docs/user/source-control.md @@ -181,6 +181,13 @@ closed reviews refresh periodically so reopening one on the host is detected. Me when requested. With **Auto-settle merged threads** enabled, a thread can settle after every linked review is terminal. An open or unsynced link keeps it active. +Ask the agent to watch, monitor, or babysit a pull request and it calls `watch_pull_request`. While +the thread is active, the server checks the pull request every minute and wakes the agent when a check +fails, the required checks pass, someone else comments or reviews, or the branch starts to conflict. +Comments from your own account do not wake it. Watching ends when the pull request merges or closes, +after 10 wakes in a row that bring only comments, or when the server cannot read the pull request for +15 minutes. To start or stop it yourself, use the row menu in the **Linked pull requests** panel. + Cross-repository links use a project on the same host. Azure DevOps reviews require a project checked out from the matching organization and repository. diff --git a/packages/client-runtime/src/operations/commands.ts b/packages/client-runtime/src/operations/commands.ts index ded1ff733e22..c0cb8e67504c 100644 --- a/packages/client-runtime/src/operations/commands.ts +++ b/packages/client-runtime/src/operations/commands.ts @@ -1027,6 +1027,20 @@ export const linkThreadPullRequest = Effect.fn("EnvironmentCommands.linkThreadPu }); }, ); +export type WatchThreadPullRequestInput = Omit< + Extract, + "type" | "commandId" +> & + CommandMetadata; +export const watchThreadPullRequest = Effect.fn("EnvironmentCommands.watchThreadPullRequest")( + function* (input: WatchThreadPullRequestInput) { + return yield* dispatch({ + ...input, + type: "thread.pull-request.watch", + commandId: yield* allocateCommandId(input), + }); + }, +); export const unlinkThreadPullRequest = Effect.fn("EnvironmentCommands.unlinkThreadPullRequest")( function* (input: UnlinkThreadPullRequestInput) { return yield* dispatch({ diff --git a/packages/client-runtime/src/state/threadCommands.ts b/packages/client-runtime/src/state/threadCommands.ts index 685ccc454a08..e3ad5c07f583 100644 --- a/packages/client-runtime/src/state/threadCommands.ts +++ b/packages/client-runtime/src/state/threadCommands.ts @@ -48,6 +48,7 @@ import { type UnarchiveThreadInput, type UnlinkThreadPullRequestInput, type UnpinThreadInput, + type WatchThreadPullRequestInput, type UnsettleThreadInput, type UnsnoozeThreadInput, type UpdateThreadMetadataInput, @@ -86,6 +87,7 @@ import { unsnoozeThread, updateThreadMetadata, visitThread, + watchThreadPullRequest, } from "../operations/commands.ts"; import type { EnvironmentRegistry } from "../connection/registry.ts"; import * as EnvironmentSupervisor from "../connection/supervisor.ts"; @@ -130,6 +132,7 @@ export type { UnsnoozeThreadInput, UpdateThreadMetadataInput, VisitThreadInput, + WatchThreadPullRequestInput, } from "../operations/commands.ts"; export function createThreadEnvironmentAtoms( @@ -251,6 +254,12 @@ export function createThreadEnvironmentAtoms( scheduler, concurrency, }), + watchPullRequest: createEnvironmentCommand(runtime, { + label: "environment-data:commands:thread:watch-pull-request", + execute: (input: WatchThreadPullRequestInput) => watchThreadPullRequest(input), + scheduler, + concurrency, + }), setRuntimeMode: createEnvironmentCommand(runtime, { label: "environment-data:commands:thread:set-runtime-mode", execute: (input: SetThreadRuntimeModeInput) => setThreadRuntimeMode(input), diff --git a/packages/client-runtime/src/t3ToolSummary.ts b/packages/client-runtime/src/t3ToolSummary.ts index 79577892a697..c6b40e5a62ff 100644 --- a/packages/client-runtime/src/t3ToolSummary.ts +++ b/packages/client-runtime/src/t3ToolSummary.ts @@ -372,6 +372,16 @@ export function summarizeT3ToolCalls( case "unlink-pr": label = phrase("Unlinked", "unlink", quantity(selected.length, "pull request")); break; + case "watch-pr": + label = phrase("Watching", "watch", quantity(selected.length, "pull request")); + break; + case "unwatch-pr": + label = phrase( + "Stopped watching", + "stop watching", + quantity(selected.length, "pull request"), + ); + break; case "list-prs": label = phrase( "Checked", diff --git a/packages/client-runtime/src/work-log/presentation.ts b/packages/client-runtime/src/work-log/presentation.ts index fdaf28dcf3b7..84229d3bac7f 100644 --- a/packages/client-runtime/src/work-log/presentation.ts +++ b/packages/client-runtime/src/work-log/presentation.ts @@ -85,6 +85,8 @@ export type ToolGroupAction = | "link-pr" | "unlink-pr" | "list-prs" + | "watch-pr" + | "unwatch-pr" | "read" | "edit" | "command" @@ -156,7 +158,9 @@ function resolveT3McpToolPresentation( const actionKind = definition.summaryAction === "link-pr" || definition.summaryAction === "unlink-pr" || - definition.summaryAction === "list-prs" + definition.summaryAction === "list-prs" || + definition.summaryAction === "watch-pr" || + definition.summaryAction === "unwatch-pr" ? definition.summaryAction : undefined; const payload = asRecord(data); @@ -539,6 +543,10 @@ function toolGroupActionLabel(action: ToolGroupAction, count: number): string { return `Linked ${count} ${count === 1 ? "pull request" : "pull requests"}`; case "unlink-pr": return `Unlinked ${count} ${count === 1 ? "pull request" : "pull requests"}`; + case "watch-pr": + return `Watching ${count} ${count === 1 ? "pull request" : "pull requests"}`; + case "unwatch-pr": + return `Stopped watching ${count} ${count === 1 ? "pull request" : "pull requests"}`; case "list-prs": return count === 1 ? "Checked linked pull requests" diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index a649cfe38de0..6752fa2a62ee 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -162,6 +162,8 @@ export const ExecutionEnvironmentCapabilities = Schema.Struct({ shaping and validation when this is absent. */ serverResolvedCommandContext: Schema.optionalKey(Schema.Boolean), threadPullRequests: Schema.optionalKey(Schema.Boolean), + /** Server understands thread.pull-request.watch and wakes agents on pull request changes. */ + threadPullRequestWatch: Schema.optionalKey(Schema.Boolean), pullRequestStackActions: Schema.optionalKey(Schema.Boolean), /** The update path clients should offer for this server. Absent on servers that must be relaunched manually (dev checkouts, Windows diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 1820008b77dd..6004ab119917 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -46,6 +46,7 @@ import { ThreadPullRequestLinkSource, ThreadPullRequestSnapshot, ThreadPullRequestStack, + ThreadPullRequestWatch, } from "./threadPullRequest.ts"; import { ProviderApprovalDecision, @@ -2600,6 +2601,18 @@ export const OrchestrationV2Command = Schema.Union([ snapshot: ThreadPullRequestSnapshot, stack: Schema.NullOr(ThreadPullRequestStack), }), + /** Start or stop the server watching a linked pull request for this thread's agent. */ + Schema.Struct({ + type: Schema.Literal("thread.pull-request.watch"), + commandId: CommandId, + threadId: ThreadId, + ...ThreadPullRequestKey.fields, + watching: Schema.Boolean, + /** Links the pull request first when starting a watch on one the thread has not linked. */ + link: Schema.optional( + Schema.Struct({ url: TrimmedNonEmptyString, source: ThreadPullRequestLinkSource }), + ), + }), Schema.Struct({ type: Schema.Literal("thread.pull-request.sync"), commandId: CommandId, @@ -2859,6 +2872,27 @@ export type OrchestrationV2Command = typeof OrchestrationV2Command.Type; * send them. */ const OrchestrationV2InternalCommand = Schema.Union([ + /** + * Records what a pull request watch saw, and wakes the agent in the same transaction when + * `wake` is set. Rejected once the watch started at `startedAt` has ended, and a wake is + * rejected on a settled or archived thread, so a read that raced either changes nothing. + */ + Schema.Struct({ + type: Schema.Literal("thread.pull-request-watch.sync"), + commandId: CommandId, + threadId: ThreadId, + ...ThreadPullRequestKey.fields, + startedAt: IsoDateTime, + /** The watch to record, or null to end it. */ + watch: Schema.NullOr(ThreadPullRequestWatch), + wake: Schema.optional( + Schema.Struct({ + messageId: MessageId, + text: Schema.String, + notification: OrchestrationV2Notification, + }), + ), + }), /** Records that the provider rollback `requestId` failed for good. */ Schema.Struct({ type: Schema.Literal("checkpoint.rollback.fail"), diff --git a/packages/contracts/src/pullRequest.ts b/packages/contracts/src/pullRequest.ts index cf8c00db0a88..c8feaad9a41f 100644 --- a/packages/contracts/src/pullRequest.ts +++ b/packages/contracts/src/pullRequest.ts @@ -152,6 +152,8 @@ export const PullRequestCheck = Schema.Struct({ status: PullRequestCheckStatus, description: Schema.NullOr(Schema.String), url: Schema.NullOr(Schema.String), + /** The base branch requires this check to merge. Absent where the host does not say. */ + required: Schema.optional(Schema.Boolean), }); export type PullRequestCheck = typeof PullRequestCheck.Type; @@ -844,6 +846,8 @@ export const PullRequestDetail = Schema.Struct({ changedFiles: NonNegativeInt, headBranch: TrimmedNonEmptyString, headRepositoryNameWithOwner: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), + /** The head commit, where the host reports it with the detail. */ + headSha: Schema.optional(TrimmedNonEmptyString), baseBranch: TrimmedNonEmptyString, createdAt: IsoDateTime, updatedAt: IsoDateTime, diff --git a/packages/contracts/src/threadPullRequest.ts b/packages/contracts/src/threadPullRequest.ts index 119503c2a85d..9df5cbe9670a 100644 --- a/packages/contracts/src/threadPullRequest.ts +++ b/packages/contracts/src/threadPullRequest.ts @@ -92,6 +92,30 @@ export const ThreadPullRequestKey = Schema.Struct({ }); export type ThreadPullRequestKey = typeof ThreadPullRequestKey.Type; +/** + * Present while the server watches the pull request for its thread. The server wakes the + * thread's agent when checks finish on the head commit, someone else comments, or the branch + * starts to conflict. The other fields record what the agent was last told, so each change is + * reported once. + */ +export const ThreadPullRequestWatch = Schema.Struct({ + startedAt: IsoDateTime, + /** Head commit at the last pass; null where the host does not report one. */ + headSha: Schema.NullOr(TrimmedNonEmptyString), + /** Failed checks on that commit the agent was told about; a rerun that fails again is news. */ + failedChecks: Schema.Array(TrimmedNonEmptyString), + /** The agent was told the required checks on that commit passed. */ + passed: Schema.Boolean, + /** Remarks from others created up to this host time were reported. */ + remarksThrough: IsoDateTime, + /** Remarks created exactly at `remarksThrough` that were reported, so a late one still counts. */ + remarkIds: Schema.Array(TrimmedNonEmptyString), + conflicting: Schema.Boolean, + /** Comment-only wakes in a row. Watching stops at a limit, so bots cannot loop it. */ + wakes: NonNegativeInt, +}); +export type ThreadPullRequestWatch = typeof ThreadPullRequestWatch.Type; + export const ThreadPullRequestLink = Schema.Struct({ ...ThreadPullRequestKey.fields, url: TrimmedNonEmptyString, @@ -99,5 +123,6 @@ export const ThreadPullRequestLink = Schema.Struct({ linkedAt: IsoDateTime, snapshot: Schema.NullOr(ThreadPullRequestSnapshot), stack: Schema.NullOr(ThreadPullRequestStack), + watch: Schema.optional(ThreadPullRequestWatch), }); export type ThreadPullRequestLink = typeof ThreadPullRequestLink.Type; diff --git a/packages/shared/src/t3McpToolPresentation.ts b/packages/shared/src/t3McpToolPresentation.ts index dd38709bd0e6..d85104cf4bc1 100644 --- a/packages/shared/src/t3McpToolPresentation.ts +++ b/packages/shared/src/t3McpToolPresentation.ts @@ -55,6 +55,8 @@ export type T3McpToolSummaryAction = | "link-pr" | "unlink-pr" | "list-prs" + | "watch-pr" + | "unwatch-pr" | "browser" | "device"; @@ -93,6 +95,16 @@ const T3_MCP_TOOLS: Readonly> = { "list-prs", "pull-request", ), + watch_pull_request: tool( + ["Watch", "Watching", "Watching", "a pull request"], + "watch-pr", + "pull-request", + ), + unwatch_pull_request: tool( + ["Stop watching", "Stopping watching", "Stopped watching", "a pull request"], + "unwatch-pr", + "pull-request", + ), orchestrator_capabilities: tool( ["Get", "Getting", "Got", "orchestration capabilities"], "capabilities",