diff --git a/AGENTS.md b/AGENTS.md index 2e56829..8ae211c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -14,7 +14,9 @@ For features with background jobs, polling, retries, cancellation, or shutdown: - Validate retry context against the current snapshot, plan revision, assignment, and referenced code. If any context is stale, disable retry and require a new request. - Preserve the original timeout, cancellation, and shutdown reason through every layer. Do not replace actionable errors with generic cancellation text. - Begin shutdown by rejecting new work at the outer admission boundary. Drain already-admitted HTTP requests, then cancel and await jobs, then close storage. +- Bound the HTTP drain during shutdown. After its grace period, abort and await owned work before awaiting server closure so an admitted poll cannot deadlock teardown. - Polling endpoints should read only the state they need. Do not rebuild Git history or the full review merely to retrieve background-job status. +- Back off recurring external-status polling to a bounded cap. Reset the interval only after a meaningful lifecycle change or explicit user action. ## Async review UI @@ -74,6 +76,10 @@ Every reproduced race requires a failing-before and passing-after regression. As - A deadline must abort and await the underlying operation before releasing its in-flight ownership; rejecting only the caller can leave untracked work running. - Invalidate pre-action status caches after both successful and refused external mutations before rendering or fetching status again. Use a generation guard so reads started before or during the mutation cannot repopulate the cache afterward. - Keep irreversible integrations disabled in demo mode even when configuration or an injected dependency is present. After a stale or refused irreversible action, keep its control disabled until fresh state is loaded. +- When a durable external-action attempt is bound to an older snapshot, require approvals or evidence recorded against the replacement generation before another action, whether or not the prior outcome explicitly requested fresh review. A mismatch with the old context is not itself fresh review. +- If the external lifecycle mechanism or mode changes between validation passes, abort before the irreversible command. Create durable lifecycle ownership from the final stable mode, never from an earlier observation. +- When an irreversible command has an ambiguous timeout, cancellation, transport, or unknown outcome, retain durable in-flight ownership and reconcile external state before enabling retry. Only a confirmed refusal may become retryable failure. +- Correlate retry observations to the current attempt with an immutable external identity or event boundary, and fail closed when multiple post-boundary action sequences appear. Matching only the resource or commit identity can replay another attempt's terminal event. ## Blinded experiments diff --git a/docs/implementation/merge-queue.md b/docs/implementation/merge-queue.md new file mode 100644 index 0000000..d7c83ef --- /dev/null +++ b/docs/implementation/merge-queue.md @@ -0,0 +1,31 @@ +# Merge-queue lifecycle (#24) + +Queue-backed merging is enabled only when the configured GitHub adapter implements the recorded K1 queue-observation contract. The reviewed head remains the authority throughout the lifecycle; browser requests never supply a head, attempt ID, or queue conclusion. + +## Durable lifecycle + +SQLite owns the merge-attempt record and retains its plan revision, snapshot, review version, and exact reviewed head. Legal transitions are: + +```text +submitting -> queued -> merged + -> removed + -> failed +``` + +An observation may settle `submitting` after a process restart, because the enqueue command may have completed before the local queued update. Enqueue command success is never recorded as merged. Once submission begins, cancellation, timeout, transport failure, or another ambiguous command outcome leaves the persisted `submitting` record disabled, retains its diagnostic, and remains recoverable through queue inspection. Direct merges use the same durable ownership: a later pull-request read can reconcile an ambiguous direct command to merged, while any non-merged state remains disabled rather than being treated as proof that retry is safe. Only a confirmed GitHub refusal becomes a retryable `failed` attempt. Once the external command succeeds, a later local refresh failure likewise does not turn that committed action into a command failure. + +Every update compares the current attempt ID and legal source state. A delayed poll for an older attempt therefore cannot overwrite a retry. Immediately before enqueue, the coordinator persists GitHub's opaque connection cursor for the latest queue-timeline event; later inspection paginates forward from that cursor, with a bounded fail-closed page limit. Correlated state requires exactly one post-cursor add sequence and its matching active or terminal shape; multiple add sequences fail closed. An older or later attempt for the same reviewed head therefore cannot settle this attempt, and long-running restart recovery does not depend on the boundary remaining in a recent-event window. Removed and failed attempts retain GitHub's terminal reason. Retry creates a new attempt only when the same revision, snapshot, review version, and reviewed head are still current. Any snapshot replacement or plan amendment after a stored attempt requires every plan item to have an approval bound to the current snapshot and revision. When GitHub explicitly requires fresh review without a locally observed replacement snapshot, those approvals must come from a review generation at or after the attempted generation. A distinct replacement snapshot can satisfy the gate after complete re-review even if its commit SHA was restored to the original value. After those approvals are recorded, the old attempt is historical rather than retryable. + +## Runtime ownership + +The coordinator owns at most one enqueue and one shared queue inspection. Cancellation does not release either operation until its GitHub call settles. One fourteen-second deadline covers every sequential validation, queue-correlation, and submission stage rather than restarting for each remote call. Shutdown rejects every new API request at admission and rechecks irreversible work after partially received request bodies, then gives admitted requests a bounded fourteen-and-a-half-second drain. At the deadline it aborts request-scoped status inspections, destroys requests still blocked on partial bodies, and cancels coordinator work before waiting for server closure and closing SQLite, so an admitted request or queue poll cannot deadlock teardown. + +`GET /api/merge` is the narrow polling path for queue-backed attempts. It reads only the persisted attempt and the K1 queue observation; it does not reload Git history or reconstruct the review. Ambiguous direct submissions remain disabled and reconcile through a full review refresh instead of repeatedly calling the queue endpoint. The browser patches only the merge control and status banner, so current selection, scroll, code attachment, and composer drafts remain unchanged. Queue polling starts at two seconds, doubles while the same active state persists, and caps at thirty seconds; a lifecycle transition or explicit user action resets the interval. A review action advances the merge-poll generation, preventing an older response from re-enabling retry against pre-action review state. The action stays disabled while submitting or queued and after confirmed merge. Removal or failure exposes the reason but keeps retry disabled until a full review refresh revalidates current GitHub blockers. Retry is then offered only for an unchanged reviewed context, and the full gate is revalidated twice before another enqueue. + +Transient or incomplete GitHub observations leave the attempt active, surface an observation error, and continue polling rather than clearing the browser's active attempt. Submitting attempts use unknown-outcome wording until GitHub proves that the attempt is queued. Only a validated queued, merged, removed, or failed observation changes durable state. Queue rules count as the server-side current-base guard; adapters without queue inspection continue to fail closed. + +Both pre-action validation reads must agree on whether merge queues apply. Queue mode performs one final fresh validation after capturing its timeline cursor, immediately before durable ownership and `gh pr merge`. Any mode change aborts before the command; the coordinator never decides whether to create durable queue ownership from an earlier observation. + +## Acceptance evidence + +Controlled unit and browser regressions cover enqueue success, delayed merge, queue removal, unmergeable failure, head replacement, retry, stale observation publication, active-attempt retry refusal, shutdown settlement, durable restart recovery, and preservation of UI input while polling. Final evidence is recorded against the exact pushed PR head. diff --git a/github/merge.ts b/github/merge.ts index deab86c..08b5faa 100644 --- a/github/merge.ts +++ b/github/merge.ts @@ -22,6 +22,14 @@ export interface RemoteMergeState { } export interface MergeResult { url: string; } +export class MergeSubmissionError extends Error { + readonly outcome: 'refused' | 'unknown'; + constructor(message: string, outcome: 'refused' | 'unknown', options?: ErrorOptions) { + super(message, options); + this.name = 'MergeSubmissionError'; + this.outcome = outcome; + } +} export type MergeQueueEntryPhase = 'AWAITING_CHECKS' | 'LOCKED' | 'MERGEABLE' | 'QUEUED'; export type MergeQueueObservation = | { state: 'queued'; reviewedHead: string; entryId: string; phase: MergeQueueEntryPhase; position: number; enqueuedAt: string; queueHead: string } @@ -29,10 +37,11 @@ export type MergeQueueObservation = | { state: 'failed'; reviewedHead: string; entryId: string; reason: string } | { state: 'merged'; reviewedHead: string; mergedAt: string }; export interface MergeQueueGateway { - inspectQueue(expectedHead: string, options?: { signal?: AbortSignal; timeoutMs?: number }): Promise; + queueWatermark(expectedHead: string, options?: { signal?: AbortSignal; timeoutMs?: number }): Promise; + inspectQueue(expectedHead: string, options?: { signal?: AbortSignal; timeoutMs?: number; afterCursor?: string | null }): Promise; } export interface MergeGateway { - inspect(options?: { fresh?: boolean; timeoutMs?: number }): Promise; + inspect(options?: { fresh?: boolean; timeoutMs?: number; signal?: AbortSignal }): Promise; merge(expectedHead: string, options?: { signal?: AbortSignal }): Promise; } @@ -45,6 +54,10 @@ export interface GhMergeConfig { type RunGh = (args: readonly string[], options?: { signal?: AbortSignal }) => Promise; +function confirmedMergeRefusal(message: string): boolean { + return /required (?:approving )?review|required status check|branch protection|merge conflict|not mergeable|head (?:branch |commit )?(?:was )?(?:modified|changed)|does not match.*head|pull request.*(?:closed|draft)|merge method.*not allowed/i.test(message); +} + function fullSha(value: unknown, label: string): string { if (typeof value !== 'string' || !/^[a-f0-9]{40}$/.test(value)) throw new Error(`GitHub returned an invalid ${label}.`); return value; @@ -156,7 +169,7 @@ export class GhMergeGateway implements MergeGateway, MergeQueueGateway { if (!rule || typeof rule !== 'object') { rulesKnown = false; break; } const value = rule as { type?: unknown; parameters?: Record }; if (typeof value.type !== 'string') { rulesKnown = false; break; } - if (value.type === 'merge_queue') mergeQueue = true; + if (value.type === 'merge_queue') { mergeQueue = true; atomicBaseGuard = true; } if (value.type !== 'required_status_checks') continue; const parameters = value.parameters; if (!parameters || !Array.isArray(parameters.required_status_checks)) { rulesKnown = false; break; } @@ -223,16 +236,18 @@ export class GhMergeGateway implements MergeGateway, MergeQueueGateway { }; } - async inspect(options: { fresh?: boolean; timeoutMs?: number } = {}): Promise { + async inspect(options: { fresh?: boolean; timeoutMs?: number; signal?: AbortSignal } = {}): Promise { const timeoutMs = options.timeoutMs ?? 12_000; if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 12_000) throw new Error('Invalid GitHub inspection timeout.'); if (!options.fresh && this.#cache && this.#cache.expiresAt > Date.now()) return this.#cache.state; const generation = this.#generation; if (!options.fresh && this.#inflight?.generation === generation) return this.#inflight.promise; - const controller = new AbortController(); - const timer = setTimeout(() => controller.abort(new Error('GitHub merge-state inspection timed out.')), timeoutMs); - const attempt = this.#inspectNow(controller.signal).catch(error => { - if (controller.signal.aborted) throw new Error('GitHub merge-state inspection timed out.'); + const timeout = new AbortController(); + const timer = setTimeout(() => timeout.abort(new Error('GitHub merge-state inspection timed out.')), timeoutMs); + const signal = options.signal ? AbortSignal.any([options.signal, timeout.signal]) : timeout.signal; + const attempt = this.#inspectNow(signal).catch(error => { + if (timeout.signal.aborted) throw timeout.signal.reason; + if (options.signal?.aborted) throw options.signal.reason; throw error; }).then(state => { if (this.#generation === generation) this.#cache = { expiresAt: Date.now() + 5_000, state }; @@ -242,42 +257,103 @@ export class GhMergeGateway implements MergeGateway, MergeQueueGateway { return attempt; } - async inspectQueue(expectedHead: string, options: { signal?: AbortSignal; timeoutMs?: number } = {}): Promise { + async queueWatermark(expectedHead: string, options: { signal?: AbortSignal; timeoutMs?: number } = {}): Promise { fullSha(expectedHead, 'expected head SHA'); - const timeoutMs = options.timeoutMs ?? 12_000; - if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 12_000) throw new Error('Invalid GitHub queue inspection timeout.'); + const timeoutMs = options.timeoutMs ?? 6_000; + if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 12_000) throw new Error('Invalid GitHub queue watermark timeout.'); const [owner, name] = this.config.repository.split('/') as [string, string]; - const query = `query($owner:String!,$name:String!,$number:Int!){repository(owner:$owner,name:$name){pullRequest(number:$number){number headRefOid state mergedAt mergeQueueEntry{id state position enqueuedAt headCommit{oid} pullRequest{number headRefOid}} timelineItems(last:20,itemTypes:[ADDED_TO_MERGE_QUEUE_EVENT,REMOVED_FROM_MERGE_QUEUE_EVENT]){nodes{__typename ... on AddedToMergeQueueEvent{createdAt} ... on RemovedFromMergeQueueEvent{createdAt reason beforeCommit{oid}}}}}}`; + const query = `query($owner:String!,$name:String!,$number:Int!){repository(owner:$owner,name:$name){pullRequest(number:$number){number headRefOid timelineItems(last:1,itemTypes:[ADDED_TO_MERGE_QUEUE_EVENT,REMOVED_FROM_MERGE_QUEUE_EVENT]){edges{cursor node{id}}}}}}`; const timeout = new AbortController(); - const timer = setTimeout(() => timeout.abort(new Error('GitHub merge-queue inspection timed out.')), timeoutMs); + const timer = setTimeout(() => timeout.abort(new Error('GitHub merge-queue watermark timed out.')), timeoutMs); const signal = options.signal ? AbortSignal.any([options.signal, timeout.signal]) : timeout.signal; try { if (options.signal?.aborted) throw options.signal.reason; const response = await this.#json(['api','graphql','-f',`query=${query}`,'-f',`owner=${owner}`,'-f',`name=${name}`,'-F',`number=${this.config.pullRequest}`], signal) as { - data?: { repository?: { pullRequest?: Record | null } | null }; + data?: { repository?: { pullRequest?: { number?: unknown; headRefOid?: unknown; timelineItems?: { edges?: unknown } } | null } | null }; errors?: unknown; }; if (signal.aborted) throw signal.reason; - if (Object.hasOwn(response, 'errors') && (!Array.isArray(response.errors) || response.errors.length > 0)) throw new Error('GitHub returned merge-queue data with errors.'); + if (Object.hasOwn(response, 'errors') && (!Array.isArray(response.errors) || response.errors.length > 0)) throw new Error('GitHub returned merge-queue watermark data with errors.'); const pull = response.data?.repository?.pullRequest; + if (!pull || pull.number !== this.config.pullRequest || fullSha(pull.headRefOid, 'queue watermark head SHA') !== expectedHead || !Array.isArray(pull.timelineItems?.edges)) + throw new Error('GitHub returned an incomplete merge-queue watermark.'); + const edges = pull.timelineItems.edges; + if (edges.length > 1) throw new Error('GitHub returned an invalid merge-queue watermark.'); + const edge = edges[0]; + if (edge === undefined) return null; + if (!edge || typeof edge !== 'object' || Array.isArray(edge) || typeof (edge as { cursor?: unknown }).cursor !== 'string' || !(edge as { cursor: string }).cursor || + !(edge as { node?: unknown }).node || typeof (edge as { node: unknown }).node !== 'object' || Array.isArray((edge as { node: unknown }).node) || + typeof ((edge as { node: { id?: unknown } }).node.id) !== 'string' || !(edge as { node: { id: string } }).node.id) + throw new Error('GitHub returned an invalid merge-queue watermark.'); + return (edge as { cursor: string }).cursor; + } catch (error) { + if (options.signal?.aborted) throw options.signal.reason; + if (timeout.signal.aborted && !options.signal?.aborted) throw new Error('GitHub merge-queue watermark timed out.'); + throw error; + } finally { clearTimeout(timer); } + } + + async inspectQueue(expectedHead: string, options: { signal?: AbortSignal; timeoutMs?: number; afterCursor?: string | null } = {}): Promise { + fullSha(expectedHead, 'expected head SHA'); + const timeoutMs = options.timeoutMs ?? 12_000; + if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 12_000) throw new Error('Invalid GitHub queue inspection timeout.'); + const correlated = Object.hasOwn(options, 'afterCursor'); + if (options.afterCursor !== undefined && options.afterCursor !== null && (typeof options.afterCursor !== 'string' || !options.afterCursor || options.afterCursor.length > 512)) + throw new Error('Invalid merge-queue event cursor.'); + const [owner, name] = this.config.repository.split('/') as [string, string]; + const query = `query($owner:String!,$name:String!,$number:Int!,$after:String){repository(owner:$owner,name:$name){pullRequest(number:$number){number headRefOid state mergedAt mergeQueueEntry{id state position enqueuedAt headCommit{oid} pullRequest{number headRefOid}} timelineItems(first:100,after:$after,itemTypes:[ADDED_TO_MERGE_QUEUE_EVENT,REMOVED_FROM_MERGE_QUEUE_EVENT]){edges{cursor node{id __typename ... on AddedToMergeQueueEvent{createdAt} ... on RemovedFromMergeQueueEvent{createdAt reason beforeCommit{oid}}}} pageInfo{hasNextPage endCursor}}}}}`; + const timeout = new AbortController(); + const timer = setTimeout(() => timeout.abort(new Error('GitHub merge-queue inspection timed out.')), timeoutMs); + const signal = options.signal ? AbortSignal.any([options.signal, timeout.signal]) : timeout.signal; + try { + if (options.signal?.aborted) throw options.signal.reason; + let cursor = options.afterCursor ?? null, pull: Record | null = null, stable = '', page = 0; + const events: Array<{ id: string; type: unknown; createdAt: string; reason: string | undefined; beforeHead: string | undefined }> = []; + while (page++ < 10) { + const args = ['api','graphql','-f',`query=${query}`,'-f',`owner=${owner}`,'-f',`name=${name}`,'-F',`number=${this.config.pullRequest}`]; + if (cursor !== null) args.push('-f', `after=${cursor}`); + const response = await this.#json(args, signal) as { data?: { repository?: { pullRequest?: Record | null } | null }; errors?: unknown }; + if (signal.aborted) throw signal.reason; + if (Object.hasOwn(response, 'errors') && (!Array.isArray(response.errors) || response.errors.length > 0)) throw new Error('GitHub returned merge-queue data with errors.'); + const nextPull = response.data?.repository?.pullRequest; + if (!nextPull || nextPull.number !== this.config.pullRequest) throw new Error('GitHub returned an incomplete merge-queue pull request.'); + const current = JSON.stringify({ number: nextPull.number, headRefOid: nextPull.headRefOid, state: nextPull.state, mergedAt: nextPull.mergedAt, mergeQueueEntry: nextPull.mergeQueueEntry }); + if (stable && current !== stable) throw new Error('GitHub merge-queue state changed during paginated inspection.'); + stable = current; pull = nextPull; + const timeline = nextPull.timelineItems; + if (!timeline || typeof timeline !== 'object' || Array.isArray(timeline)) throw new Error('GitHub returned incomplete merge-queue history.'); + const value = timeline as { edges?: unknown; pageInfo?: { hasNextPage?: unknown; endCursor?: unknown } }; + if (!Array.isArray(value.edges) || !value.pageInfo || typeof value.pageInfo.hasNextPage !== 'boolean' || + (value.pageInfo.endCursor !== null && typeof value.pageInfo.endCursor !== 'string')) throw new Error('GitHub returned incomplete merge-queue pagination.'); + for (const edge of value.edges) { + if (!edge || typeof edge !== 'object' || Array.isArray(edge) || typeof (edge as { cursor?: unknown }).cursor !== 'string' || !(edge as { cursor: string }).cursor) + throw new Error('GitHub returned a malformed merge-queue edge.'); + const event = (edge as { node?: unknown }).node; + if (!event || typeof event !== 'object' || Array.isArray(event)) throw new Error('GitHub returned a malformed merge-queue event.'); + const node = event as { id?: unknown; __typename?: unknown; createdAt?: unknown; reason?: unknown; beforeCommit?: { oid?: unknown } | null }; + if (typeof node.id !== 'string' || !node.id || node.id.length > 512) throw new Error('GitHub returned a merge-queue event without stable identity.'); + if (!['AddedToMergeQueueEvent','RemovedFromMergeQueueEvent'].includes(String(node.__typename))) throw new Error('GitHub returned an unknown merge-queue event.'); + const createdAt = timestamp(node.createdAt, 'merge-queue event time'); + if (node.__typename === 'RemovedFromMergeQueueEvent' && (typeof node.reason !== 'string' || !node.reason.trim())) throw new Error('GitHub returned a merge-queue removal without a reason.'); + const beforeHead = node.__typename === 'RemovedFromMergeQueueEvent' ? fullSha(node.beforeCommit?.oid, 'removed merge-queue head SHA') : undefined; + events.push({ id: node.id, type: node.__typename, createdAt, reason: node.reason as string | undefined, beforeHead }); + } + if (!value.pageInfo.hasNextPage) break; + if (page === 10 || typeof value.pageInfo.endCursor !== 'string' || !value.pageInfo.endCursor) throw new Error('GitHub merge-queue history exceeded the inspection limit.'); + cursor = value.pageInfo.endCursor; + } if (!pull || pull.number !== this.config.pullRequest) throw new Error('GitHub returned an incomplete merge-queue pull request.'); const reviewedHead = fullSha(pull.headRefOid, 'queue pull request head SHA'); if (reviewedHead !== expectedHead) throw new Error('The pull request head changed after review.'); if (!['OPEN','CLOSED','MERGED'].includes(String(pull.state))) throw new Error('GitHub returned an invalid queue pull request state.'); if (!Object.hasOwn(pull, 'mergedAt') || (pull.mergedAt !== null && typeof pull.mergedAt !== 'string')) throw new Error('GitHub returned invalid merge completion data.'); - if (!Object.hasOwn(pull, 'mergeQueueEntry') || !Object.hasOwn(pull, 'timelineItems')) throw new Error('GitHub returned incomplete merge-queue data.'); + if (!Object.hasOwn(pull, 'mergeQueueEntry')) throw new Error('GitHub returned incomplete merge-queue data.'); const entry = pull.mergeQueueEntry; - const timeline = pull.timelineItems; - if (!timeline || typeof timeline !== 'object' || Array.isArray(timeline) || !Array.isArray((timeline as { nodes?: unknown }).nodes)) throw new Error('GitHub returned incomplete merge-queue history.'); - const events = (timeline as { nodes: unknown[] }).nodes.map(event => { - if (!event || typeof event !== 'object' || Array.isArray(event)) throw new Error('GitHub returned a malformed merge-queue event.'); - const value = event as { __typename?: unknown; createdAt?: unknown; reason?: unknown; beforeCommit?: { oid?: unknown } | null }; - if (!['AddedToMergeQueueEvent','RemovedFromMergeQueueEvent'].includes(String(value.__typename))) throw new Error('GitHub returned an unknown merge-queue event.'); - const createdAt = timestamp(value.createdAt, 'merge-queue event time'); - if (value.__typename === 'RemovedFromMergeQueueEvent' && (typeof value.reason !== 'string' || !value.reason.trim())) throw new Error('GitHub returned a merge-queue removal without a reason.'); - const beforeHead = value.__typename === 'RemovedFromMergeQueueEvent' ? fullSha(value.beforeCommit?.oid, 'removed merge-queue head SHA') : undefined; - return { type: value.__typename, createdAt, reason: value.reason as string | undefined, beforeHead }; - }); + const attemptEvents = events; + const addCount = attemptEvents.filter(event => event.type === 'AddedToMergeQueueEvent').length; + const hasCurrentAdd = addCount === 1 && attemptEvents[0]?.type === 'AddedToMergeQueueEvent'; + const singleSequence = hasCurrentAdd && (attemptEvents.length === 1 || (attemptEvents.length === 2 && attemptEvents[1]?.type === 'RemovedFromMergeQueueEvent')); + if (correlated && !singleSequence) throw new Error(addCount > 1 ? 'GitHub returned multiple enqueue sequences after the current attempt cursor.' : 'GitHub did not return one complete enqueue sequence for the current attempt.'); if (pull.state === 'MERGED') { if (entry !== null) throw new Error('GitHub returned an active queue entry for a merged pull request.'); return { state: 'merged', reviewedHead, mergedAt: timestamp(pull.mergedAt, 'merge completion time') }; @@ -289,19 +365,22 @@ export class GhMergeGateway implements MergeGateway, MergeQueueGateway { const value = entry as { id?: unknown; state?: unknown; position?: unknown; enqueuedAt?: unknown; headCommit?: { oid?: unknown } | null; pullRequest?: { number?: unknown; headRefOid?: unknown } | null }; if (typeof value.id !== 'string' || !value.id || !['AWAITING_CHECKS','LOCKED','MERGEABLE','QUEUED','UNMERGEABLE'].includes(String(value.state)) || !Number.isSafeInteger(value.position) || (value.position as number) < 0) throw new Error('GitHub returned an invalid merge-queue entry.'); - timestamp(value.enqueuedAt, 'merge-queue entry time'); + const enqueuedAt = timestamp(value.enqueuedAt, 'merge-queue entry time'); + if (correlated && attemptEvents.length !== 1) throw new Error('GitHub returned an ambiguous active event sequence for the current enqueue attempt.'); const queueHead = fullSha(value.headCommit?.oid, 'merge-queue head SHA'); const entryHead = fullSha(value.pullRequest?.headRefOid, 'merge-queue entry pull request head SHA'); if (value.pullRequest?.number !== this.config.pullRequest || queueHead !== expectedHead || entryHead !== expectedHead) throw new Error('The merge-queue entry does not match the reviewed pull request head.'); if (value.state === 'UNMERGEABLE') return { state: 'failed', reviewedHead, entryId: value.id, reason: 'GitHub reported the merge queue entry as unmergeable.' }; return { state: 'queued', reviewedHead, entryId: value.id, phase: value.state as MergeQueueEntryPhase, - position: value.position as number, enqueuedAt: value.enqueuedAt as string, queueHead, + position: value.position as number, enqueuedAt, queueHead, }; } - const last = events.at(-1); + if (correlated && attemptEvents.length === 0) throw new Error('GitHub did not return an event for the current enqueue attempt.'); + const last = (correlated ? attemptEvents : events).at(-1); if (!last || last.type !== 'RemovedFromMergeQueueEvent') throw new Error('GitHub did not confirm a queued, removed, failed, or merged state.'); if (last.beforeHead !== expectedHead) throw new Error('The merge-queue removal does not match the reviewed pull request head.'); + if (correlated && (attemptEvents.length !== 2 || attemptEvents[1]?.type !== 'RemovedFromMergeQueueEvent')) throw new Error('GitHub returned an ambiguous terminal event sequence for the current enqueue attempt.'); return { state: 'removed', reviewedHead, removedAt: last.createdAt, reason: last.reason! }; } catch (error) { if (options.signal?.aborted) throw options.signal.reason; @@ -319,6 +398,11 @@ export class GhMergeGateway implements MergeGateway, MergeQueueGateway { if (options.signal?.aborted) throw options.signal.reason; await this.run(['pr','merge',String(this.config.pullRequest),'--repo',this.config.repository,flag,'--match-head-commit',expectedHead], { signal: options.signal }); return { url: `https://github.com/${this.config.repository}/pull/${this.config.pullRequest}` }; + } catch (error) { + if (error instanceof MergeSubmissionError) throw error; + const message = error instanceof Error ? error.message : 'GitHub merge submission failed with an unknown outcome.'; + const outcome = !options.signal?.aborted && confirmedMergeRefusal(message) ? 'refused' : 'unknown'; + throw new MergeSubmissionError(message, outcome, { cause: error }); } finally { this.#generation++; this.#cache = null; } } } diff --git a/runner/merge.ts b/runner/merge.ts index 228cce7..2d675f9 100644 --- a/runner/merge.ts +++ b/runner/merge.ts @@ -1,20 +1,72 @@ import type { ReviewService } from './review.ts'; -import type { MergeGateway, MergeResult, RemoteMergeState } from '../github/merge.ts'; +import type { MergeAttempt } from './store.ts'; +import { MergeSubmissionError, type MergeGateway, type MergeQueueGateway, type MergeQueueObservation, type MergeResult, type RemoteMergeState } from '../github/merge.ts'; type ReviewView = ReturnType; +type QueueGateway = MergeGateway & MergeQueueGateway; export interface MergeBlocker { code: string; message: string; } -export interface MergeStatus { available: true; ready: boolean; blockers: MergeBlocker[]; remote: RemoteMergeState; } -export interface MergeUnavailableStatus { available: true; ready: false; blockers: MergeBlocker[]; remote: null; } +export interface MergeQueueStatus { + kind: MergeAttempt['kind']; state: MergeAttempt['state']; reviewedHead: string; url: string | null; reason: string | null; + phase: MergeAttempt['phase']; position: number | null; occurredAt: string | null; retryable: boolean; + observationError?: string; +} +export interface MergeStatus { + available: true; ready: boolean; action: 'merge' | 'retry' | null; blockers: MergeBlocker[]; + remote: RemoteMergeState; queue: MergeQueueStatus | null; +} +export interface MergeUnavailableStatus { available: true; ready: false; action: null; blockers: MergeBlocker[]; remote: null; queue: MergeQueueStatus | null; } + +function queueGateway(gateway: MergeGateway): gateway is QueueGateway { + const queue = gateway as Partial; + return typeof queue.inspectQueue === 'function' && typeof queue.queueWatermark === 'function'; +} export class MergeCoordinator { #active: Promise<{ status: MergeStatus; result: MergeResult }> | null = null; #abort: AbortController | null = null; + #queuePoll: Promise | null = null; + #queueAbort: AbortController | null = null; #closing = false; readonly service: ReviewService; readonly gateway: MergeGateway; - constructor(service: ReviewService, gateway: MergeGateway) { this.service = service; this.gateway = gateway; } + readonly operationTimeoutMs: number; + constructor(service: ReviewService, gateway: MergeGateway, operationTimeoutMs = 14_000) { + if (!Number.isSafeInteger(operationTimeoutMs) || operationTimeoutMs < 1 || operationTimeoutMs > 14_000) throw new Error('Invalid merge operation deadline.'); + this.service = service; this.gateway = gateway; this.operationTimeoutMs = operationTimeoutMs; + } + + #attempt(): MergeAttempt | null { + return this.service.store?.getMergeAttempt(this.service.config.identity) ?? null; + } + + #current(attempt: MergeAttempt): boolean { + const { store, config } = this.service; + if (!store || !config) return false; + const snapshot = store.getSnapshot(config.identity); + return store.getPlan(config.identity).revision === attempt.revision && snapshot.id === attempt.snapshotId && + snapshot.head === attempt.reviewedHead && store.reviewVersion(config.identity) === attempt.reviewVersion; + } - async status(view = this.service.load(), fresh = false): Promise { + #queueStatus(attempt = this.#attempt(), observationError?: string): MergeQueueStatus | null { + if (!attempt) return null; + return { + kind: attempt.kind, state: attempt.state, reviewedHead: attempt.reviewedHead, url: attempt.url, reason: attempt.reason, + phase: attempt.phase, position: attempt.position, occurredAt: attempt.occurredAt, + retryable: (attempt.state === 'removed' || attempt.state === 'failed') && !attempt.requiresFreshReview && this.#current(attempt), + ...(observationError ? { observationError } : {}), + }; + } + + #freshReviewComplete(attempt: MergeAttempt): boolean { + const { store, config } = this.service; + const plan = store.getPlan(config.identity), snapshot = store.getSnapshot(config.identity); + if (snapshot.id === attempt.snapshotId && plan.revision === attempt.revision && !attempt.requiresFreshReview) return true; + const approvals = store.getReview(config.identity).approvals; + return store.reviewVersion(config.identity) > attempt.reviewVersion && plan.items.every(item => approvals.some(approval => approval.item === item.id && + approval.revision === plan.revision && approval.snapshotId === snapshot.id && approval.reviewVersion !== undefined && approval.reviewVersion >= attempt.reviewVersion)); + } + + async status(view = this.service.load(), fresh = false, signal?: AbortSignal): Promise { const blockers: MergeBlocker[] = []; for (const item of view.items) { if (item.state !== 'approved') blockers.push({ code: 'approval', message: `${item.id} is ${item.state}.` }); @@ -27,23 +79,47 @@ export class MergeCoordinator { if (unplanned) blockers.push({ code: 'unplanned', message: `${unplanned} unplanned change${unplanned === 1 ? '' : 's'} remain.` }); const changes = view.notes.filter(note => note.kind === 'change' && note.revision === view.plan.revision && note.snapshotId === view.snapshot.id).length; if (changes) blockers.push({ code: 'changes', message: `${changes} change request${changes === 1 ? '' : 's'} remain open.` }); - const remote = await this.gateway.inspect({ fresh, timeoutMs: fresh ? 6_000 : undefined }); + const remote = await this.gateway.inspect({ fresh, timeoutMs: fresh ? 6_000 : undefined, signal }); + if (signal?.aborted) throw signal.reason; if (remote.pullRequestState !== 'OPEN') blockers.push({ code: 'pr-state', message: `Pull request is ${remote.pullRequestState.toLowerCase()}.` }); if (remote.base !== view.snapshot.base) blockers.push({ code: 'base', message: 'The base branch moved. Rebase and review the resulting snapshot.' }); if (remote.head !== view.snapshot.head) blockers.push({ code: 'head', message: 'The pull request head moved. Refresh the review.' }); if (remote.mergeable !== 'MERGEABLE') blockers.push({ code: 'mergeable', message: remote.mergeable === 'CONFLICTING' ? 'The pull request has merge conflicts.' : 'GitHub has not determined mergeability.' }); if (!remote.rulesKnown) blockers.push({ code: 'rules', message: 'Required branch checks could not be read.' }); - if (remote.mergeQueue) blockers.push({ code: 'merge-queue', message: 'Merge queues are not supported by this merge action yet.' }); + if (remote.mergeQueue && !queueGateway(this.gateway)) blockers.push({ code: 'merge-queue', message: 'This GitHub adapter cannot verify the merge-queue lifecycle.' }); if (!remote.atomicBaseGuard) blockers.push({ code: 'base-guard', message: 'GitHub does not expose a server-enforced guard for the validated base.' }); for (const check of remote.requiredChecks) if (check.state !== 'success') blockers.push({ code: 'check', message: `${check.context} is ${check.state}.` }); if (remote.alreadyFixed === 'found') blockers.push({ code: 'already-fixed', message: 'Another open or merged pull request references this issue.' }); if (remote.alreadyFixed === 'unknown') blockers.push({ code: 'already-fixed', message: 'The already-fixed check could not be completed.' }); - return { available: true, ready: blockers.length === 0, blockers, remote }; + + let attempt = this.#attempt(); + if (attempt?.kind === 'direct' && attempt.state === 'submitting') { + if (remote.pullRequestState === 'MERGED' && remote.head === attempt.reviewedHead) + this.service.store.finishMergeAttempt(this.service.config.identity, attempt.id, { state: 'merged' }); + attempt = this.#attempt(); + } + const queue = this.#queueStatus(attempt); + if (attempt?.state === 'submitting' || attempt?.state === 'queued') { + blockers.unshift({ code: 'queue-active', message: attempt.state === 'submitting' + ? attempt.reason ? `The merge submission outcome is unknown. ${attempt.reason} Waiting for GitHub reconciliation.` : attempt.kind === 'queue' ? 'The reviewed head is being submitted to the merge queue.' : 'The reviewed head is being submitted for direct merge.' + : 'The reviewed head is queued. Waiting for GitHub to confirm the outcome.' }); + } else if (attempt?.state === 'merged') { + blockers.unshift({ code: 'queue-merged', message: 'GitHub confirmed that the reviewed head was merged.' }); + } else if (attempt && !this.#freshReviewComplete(attempt)) { + blockers.unshift({ code: 'queue-head', message: attempt.reason ?? 'The pull request snapshot changed after the queue attempt. Review the replacement snapshot.' }); + } + const active = attempt?.state === 'submitting' || attempt?.state === 'queued' || attempt?.state === 'merged'; + const retry = !!queue?.retryable; + const ready = !active && blockers.length === 0; + return { available: true, ready, action: ready ? (retry ? 'retry' : 'merge') : null, blockers, remote, queue }; } - async displayStatus(view = this.service.load()): Promise { - try { return await this.status(view); } - catch (error) { return { available: true, ready: false, blockers: [{ code: 'github', message: `Could not read GitHub merge state. ${error instanceof Error ? error.message : 'Unknown error.'}` }], remote: null }; } + async displayStatus(view = this.service.load(), signal?: AbortSignal): Promise { + try { return await this.status(view, false, signal); } + catch (error) { + if (signal?.aborted) throw signal.reason; + return { available: true, ready: false, action: null, blockers: [{ code: 'github', message: `Could not read GitHub merge state. ${error instanceof Error ? error.message : 'Unknown error.'}` }], remote: null, queue: this.#queueStatus() }; + } } async merge(token: unknown): Promise<{ status: MergeStatus; result: MergeResult }> { @@ -51,8 +127,10 @@ export class MergeCoordinator { if (this.#active) throw new Error('A merge attempt is already running.'); if (typeof token !== 'string') throw new Error('Stale review state. Refresh before merging.'); const abort = new AbortController(); + const timer = setTimeout(() => abort.abort(new Error('Merge request deadline exceeded.')), this.operationTimeoutMs); this.#abort = abort; const attempt = this.#merge(token, abort.signal).finally(() => { + clearTimeout(timer); if (this.#active === attempt) this.#active = null; if (this.#abort === abort) this.#abort = null; }); @@ -61,31 +139,118 @@ export class MergeCoordinator { } async #merge(token: string, signal: AbortSignal): Promise<{ status: MergeStatus; result: MergeResult }> { + let queueAttempt: MergeAttempt | null = null; try { let view = this.service.load(); if (view.token !== token) throw new Error('Stale review state. Refresh before merging.'); - const status = await this.status(view, true); + const status = await this.#statusForMerge(view, signal); if (!status.ready) throw new Error(status.blockers[0]?.message ?? 'Merge is blocked.'); view = this.service.load(); if (view.token !== token) throw new Error('Review changed during merge validation. Refresh before merging.'); - const finalStatus = await this.status(view, true); + const finalStatus = await this.#statusForMerge(view, signal); if (finalStatus.remote.base !== status.remote.base || finalStatus.remote.head !== status.remote.head) throw new Error('The pull request changed during merge validation. Refresh before merging.'); + if (finalStatus.remote.mergeQueue !== status.remote.mergeQueue) throw new Error('Merge-queue requirements changed during validation. Refresh before merging.'); if (!finalStatus.ready) throw new Error(`Merge requirements changed during validation. ${finalStatus.blockers[0]!.message}`); + let commandStatus = finalStatus; + let queueWatermark: string | null = null; + if (status.remote.mergeQueue) { + if (!queueGateway(this.gateway)) throw new Error('This GitHub adapter cannot verify the merge-queue lifecycle.'); + if (view.expected.reviewVersion === undefined) throw new Error('A current review version is required for merging.'); + queueWatermark = await this.gateway.queueWatermark(status.remote.head, { signal, timeoutMs: 6_000 }); + commandStatus = await this.#statusForMerge(view, signal); + if (commandStatus.remote.base !== finalStatus.remote.base || commandStatus.remote.head !== finalStatus.remote.head || commandStatus.remote.mergeQueue !== finalStatus.remote.mergeQueue) + throw new Error('Merge-queue requirements changed after queue correlation. Refresh before merging.'); + if (!commandStatus.ready) throw new Error(`Merge requirements changed after queue correlation. ${commandStatus.blockers[0]!.message}`); + } if (this.service.load().token !== token) throw new Error('Review changed during merge validation. Refresh before merging.'); if (signal.aborted) throw signal.reason; - const result = await this.gateway.merge(status.remote.head, { signal }); - return { status, result }; + if (this.service.store && this.service.config && view.expected.reviewVersion !== undefined) + queueAttempt = this.service.store.beginMergeAttempt(this.service.config.identity, { ...view.expected, reviewVersion: view.expected.reviewVersion }, commandStatus.remote.head, queueWatermark, commandStatus.remote.mergeQueue ? 'queue' : 'direct'); + const result = await this.gateway.merge(commandStatus.remote.head, { signal }); + if (queueAttempt) { + // The enqueue command has already committed externally. A local refresh failure must not + // report that action as failed; the durable submitting record is recoverable by polling. + try { + if (queueAttempt.kind === 'queue') this.service.store.queueMergeAttempt(this.service.config.identity, queueAttempt.id, result.url); + else this.service.store.finishMergeAttempt(this.service.config.identity, queueAttempt.id, { state: 'merged' }); + } catch {} + } + return { status: commandStatus, result }; } catch (error) { + if (queueAttempt) try { + const message = error instanceof Error ? error.message : 'GitHub merge submission outcome is unknown.'; + if (error instanceof MergeSubmissionError && error.outcome === 'refused') { + this.service.store.finishMergeAttempt(this.service.config.identity, queueAttempt.id, { + state: 'failed', reason: message, + requiresFreshReview: /head (?:branch |commit )?(?:was )?(?:modified|changed)|does not match.*head|stale review/i.test(message), + }); + } else this.service.store.recordMergeAttemptDiagnostic(this.service.config.identity, queueAttempt.id, message); + } catch {} if (signal.aborted && signal.reason instanceof Error) throw signal.reason; throw error; } } + #statusForMerge(view: ReviewView, signal: AbortSignal): Promise { + if (signal.aborted) return Promise.reject(signal.reason); + return this.status(view, true, signal); + } + + async pollQueue(): Promise { + if (this.#closing) throw new Error('Merge coordinator is shutting down.'); + if (this.#queuePoll) return this.#queuePoll; + const attempt = this.#attempt(); + if (!attempt || (attempt.state !== 'submitting' && attempt.state !== 'queued')) return this.#queueStatus(attempt); + if (attempt.kind === 'direct') return this.#queueStatus(attempt); + if (!queueGateway(this.gateway)) return this.#queueStatus(attempt, 'This GitHub adapter cannot verify the merge-queue lifecycle.'); + const abort = new AbortController(); + this.#queueAbort = abort; + const poll = this.#pollQueue(attempt, abort.signal).finally(() => { + if (this.#queuePoll === poll) this.#queuePoll = null; + if (this.#queueAbort === abort) this.#queueAbort = null; + }); + this.#queuePoll = poll; + return poll; + } + + queueSnapshot(): MergeQueueStatus | null { return this.#queueStatus(); } + + async #pollQueue(attempt: MergeAttempt, signal: AbortSignal): Promise { + try { + const observation = await (this.gateway as QueueGateway).inspectQueue(attempt.reviewedHead, { signal, timeoutMs: 12_000, afterCursor: attempt.queueWatermark ?? null }); + this.#publishQueueObservation(attempt, observation); + return this.#queueStatus(); + } catch (error) { + if (signal.aborted) throw signal.reason; + const message = error instanceof Error ? error.message : 'Could not read the merge queue.'; + if (/head changed after review/i.test(message)) { + this.service.store.finishMergeAttempt(this.service.config.identity, attempt.id, { state: 'failed', reason: message, requiresFreshReview: true }); + return this.#queueStatus(); + } + return this.#queueStatus(attempt, message); + } + } + + #publishQueueObservation(attempt: MergeAttempt, observation: MergeQueueObservation): void { + const { store } = this.service, { identity } = this.service.config; + if (observation.reviewedHead !== attempt.reviewedHead) throw new Error('GitHub returned a merge-queue observation for a different reviewed head.'); + if (observation.state === 'queued') { + store.observeQueuedMerge(identity, attempt.id, { entryId: observation.entryId, phase: observation.phase, position: observation.position }); + } else if (observation.state === 'merged') { + store.finishMergeAttempt(identity, attempt.id, { state: 'merged', occurredAt: observation.mergedAt }); + } else if (observation.state === 'removed') { + store.finishMergeAttempt(identity, attempt.id, { state: 'removed', reason: observation.reason, occurredAt: observation.removedAt }); + } else { + store.finishMergeAttempt(identity, attempt.id, { state: 'failed', reason: observation.reason }); + } + } + async close(): Promise { this.#closing = true; - const active = this.#active; - if (!active) return; + const active = this.#active, poll = this.#queuePoll; this.#abort?.abort(new Error('Merge cancelled during shutdown.')); - try { await active; } catch {} + this.#queueAbort?.abort(new Error('Merge-queue inspection cancelled during shutdown.')); + if (active) try { await active; } catch {} + if (poll) try { await poll; } catch {} } } diff --git a/runner/review.ts b/runner/review.ts index ae0efa5..3425aa5 100644 --- a/runner/review.ts +++ b/runner/review.ts @@ -49,6 +49,13 @@ export class ReviewService { return { ...segment, file, key: createHash('sha256').update(keys[index]!).digest('hex'), originalRow: raw[index]!.row }; }); const states = approvalStates(plan, segments, saved.approvals, identity); + const mergeAttempt = this.store.getMergeAttempt(identity); + const replacementReview = !!mergeAttempt && (mergeAttempt.snapshotId !== snapshot.id || mergeAttempt.revision !== plan.revision || mergeAttempt.requiresFreshReview); + if (replacementReview) for (const item of plan.items) { + const approval = saved.approvals.find(value => value.item === item.id); + if (approval && (approval.revision !== plan.revision || approval.snapshotId !== snapshot.id || + (mergeAttempt.requiresFreshReview && (approval.reviewVersion === undefined || approval.reviewVersion < mergeAttempt.reviewVersion)))) states[item.id] = 'stale'; + } for (const item of plan.items) { if (states[item.id] === 'approved' && (segments.some(segment => segment.row === 'Ambiguous' && segment.owners.includes(item.id)) || item.depends_on.some(id => states[id] === 'stale'))) states[item.id] = 'stale'; } @@ -64,6 +71,10 @@ export class ReviewService { const before = prior ? JSON.parse(prior.fingerprint) : null; const reasons: string[] = []; if (states[item.id] === 'stale') { + const approval = saved.approvals.find(value => value.item === item.id); + if (replacementReview && approval?.snapshotId !== snapshot.id) reasons.push('Pull request snapshot changed after the queue attempt'); + if (replacementReview && approval?.revision !== plan.revision) reasons.push('Plan revision changed after the queue attempt'); + if (mergeAttempt?.requiresFreshReview && approval && (approval.reviewVersion === undefined || approval.reviewVersion < mergeAttempt.reviewVersion)) reasons.push('GitHub requires a fresh review after the queue attempt'); if (before && !isDeepStrictEqual(before.item.acceptance, item.acceptance)) reasons.push('Acceptance checks changed'); if (before && owned.some(segment => before.segments.some((old: { path: string; content: string; context: string }) => old.path === segment.path && old.content === segment.content && old.context !== segment.context))) reasons.push('Moved to another function'); for (const dep of item.depends_on) if (states[dep] === 'stale') reasons.push(`Depends on ${dep}, which changed`); diff --git a/runner/store.ts b/runner/store.ts index a838568..1762b42 100644 --- a/runner/store.ts +++ b/runner/store.ts @@ -17,6 +17,14 @@ export interface SuggestionRequest { state: SuggestionState; revision: number; s export interface SnippetReference { key: string; path: string; side: 'old' | 'new'; start: number; end: number; text: string; head: string; base: string } export interface QuestionAnswer { provider?: 'claude' | 'codex'; attempt: string; contextId?: string; status: 'pending' | 'complete' | 'failed'; expiresAt: number; text?: string; error?: string } export interface ReviewNote { id: string; item: string; kind: 'question' | 'change'; text: string; reference?: SnippetReference; answer?: QuestionAnswer; createdAt: string; revision: number; snapshotId: string } +export type MergeAttemptState = 'submitting' | 'queued' | 'merged' | 'removed' | 'failed'; +export interface MergeAttempt { + id: string; kind: 'queue' | 'direct'; state: MergeAttemptState; revision: number; snapshotId: string; reviewVersion: number; reviewedHead: string; + queueWatermark?: string | null; + url: string | null; reason: string | null; requiresFreshReview: boolean; entryId: string | null; + phase: 'AWAITING_CHECKS' | 'LOCKED' | 'MERGEABLE' | 'QUEUED' | null; position: number | null; + occurredAt: string | null; createdAt: string; updatedAt: string; +} export interface LedgerEntry { sha: string; owner: string | null; origin: 'owned' | 'foreign'; sourceSha: string | null } export interface Checkpoint { id: string; revision: number; snapshotId: string; item: string; @@ -42,8 +50,8 @@ export class Store { this.#db.exec('PRAGMA foreign_keys=ON; PRAGMA journal_mode=WAL; PRAGMA synchronous=FULL;'); this.#transaction(() => { const version = this.#get('PRAGMA user_version')!.user_version; - if (version !== 0 && version !== 1 && version !== 2 && version !== 3 && version !== 4) throw new Error('Unsupported store schema version.'); - if (version === 4) return; + if (version !== 0 && version !== 1 && version !== 2 && version !== 3 && version !== 4 && version !== 5) throw new Error('Unsupported store schema version.'); + if (version === 5) return; if (version === 0) this.#db.exec(` CREATE TABLE plans (key TEXT PRIMARY KEY, issue INTEGER NOT NULL, revision INTEGER NOT NULL, snapshot_id TEXT); CREATE TABLE revisions (key TEXT NOT NULL REFERENCES plans(key), revision INTEGER NOT NULL, data TEXT NOT NULL, PRIMARY KEY(key,revision)); @@ -65,6 +73,12 @@ export class Store { ALTER TABLE requests ADD COLUMN reason TEXT; UPDATE requests SET state='invalidated', reason='Request predates snapshot binding.' WHERE state IN ('pending','ready'); PRAGMA user_version=4;`); + if (version < 5) this.#db.exec(`CREATE TABLE IF NOT EXISTS merge_attempts ( + id TEXT PRIMARY KEY, + key TEXT NOT NULL REFERENCES plans(key), + data TEXT NOT NULL + ); + PRAGMA user_version=5;`); }); } catch (error) { this.#db.close(); throw error; } } @@ -179,6 +193,72 @@ export class Store { if (typeof reason !== 'string' || !reason.trim() || reason.length > 4000) throw new Error('Invalid cancellation reason.'); this.#run("UPDATE requests SET state='cancelled',reason=? WHERE id=? AND key=? AND state IN ('pending','ready')", reason.trim(), id, identityKey(identity)); } + beginMergeAttempt(identity: PlanIdentity, expected: ReviewState & { reviewVersion: number }, reviewedHead: string, queueWatermark: string | null = null, kind: MergeAttempt['kind'] = 'queue'): MergeAttempt { + sha(reviewedHead); + if (!['queue','direct'].includes(kind)) throw new Error('Invalid merge attempt kind.'); + if (queueWatermark !== null && (typeof queueWatermark !== 'string' || !queueWatermark || queueWatermark.length > 512)) throw new Error('Invalid merge-queue event cursor.'); + if (!Number.isSafeInteger(expected.reviewVersion) || expected.reviewVersion < 0) throw new Error('A current review version is required for merging.'); + const key = identityKey(identity); + return this.#transaction(() => { + this.#expect(key, expected); + const current = this.getMergeAttempt(identity); + if (current?.state === 'submitting' || current?.state === 'queued') throw new Error('A merge-queue attempt is already active.'); + if (current?.state === 'merged') throw new Error('The reviewed pull request is already merged.'); + const now = new Date().toISOString(); + const attempt: MergeAttempt = { + id: randomUUID(), kind, state: 'submitting', revision: expected.revision, snapshotId: expected.snapshotId, + reviewVersion: expected.reviewVersion, reviewedHead, queueWatermark, url: null, reason: null, requiresFreshReview: false, + entryId: null, phase: null, position: null, occurredAt: null, createdAt: now, updatedAt: now, + }; + this.#run('INSERT INTO merge_attempts VALUES (?,?,?)', attempt.id, key, encode(attempt)); + return attempt; + }); + } + #changeMergeAttempt(identity: PlanIdentity, id: string, allowed: readonly MergeAttemptState[], change: (attempt: MergeAttempt) => MergeAttempt): boolean { + const key = identityKey(identity); + return this.#transaction(() => { + const latest = this.#get('SELECT id,data FROM merge_attempts WHERE key=? ORDER BY rowid DESC LIMIT 1', key); + if (!latest || latest.id !== id) return false; + const attempt = decode(latest.data); + if (!allowed.includes(attempt.state)) return false; + const next = change(attempt); + return this.#run('UPDATE merge_attempts SET data=? WHERE id=? AND key=?', encode({ ...next, updatedAt: new Date().toISOString() }), id, key).changes === 1; + }); + } + queueMergeAttempt(identity: PlanIdentity, id: string, url: string): boolean { + if (typeof url !== 'string' || url.length > 2048 || !/^https:\/\//.test(url)) throw new Error('Invalid merge result URL.'); + return this.#changeMergeAttempt(identity, id, ['submitting'], attempt => ({ ...attempt, state: 'queued', url, reason: null })); + } + recordMergeAttemptDiagnostic(identity: PlanIdentity, id: string, reason: string): boolean { + if (typeof reason !== 'string' || !reason.trim() || reason.length > 4000) throw new Error('A bounded merge diagnostic is required.'); + return this.#changeMergeAttempt(identity, id, ['submitting'], attempt => ({ ...attempt, reason: reason.trim() })); + } + observeQueuedMerge(identity: PlanIdentity, id: string, observation: { entryId: string; phase: MergeAttempt['phase']; position: number }): boolean { + if (typeof observation.entryId !== 'string' || !observation.entryId || observation.entryId.length > 512 || + !['AWAITING_CHECKS','LOCKED','MERGEABLE','QUEUED'].includes(String(observation.phase)) || + !Number.isSafeInteger(observation.position) || observation.position < 0) throw new Error('Invalid merge-queue observation.'); + return this.#changeMergeAttempt(identity, id, ['submitting','queued'], attempt => ({ + ...attempt, state: 'queued', reason: null, entryId: observation.entryId, phase: observation.phase, position: observation.position, + })); + } + finishMergeAttempt(identity: PlanIdentity, id: string, outcome: { + state: 'merged' | 'removed' | 'failed'; reason?: string; occurredAt?: string; requiresFreshReview?: boolean; + }): boolean { + if (!['merged','removed','failed'].includes(outcome.state)) throw new Error('Invalid merge-queue outcome.'); + if (outcome.state !== 'merged' && (typeof outcome.reason !== 'string' || !outcome.reason.trim() || outcome.reason.length > 4000)) throw new Error('A bounded terminal merge reason is required.'); + if (outcome.occurredAt !== undefined && (!Number.isFinite(Date.parse(outcome.occurredAt)) || outcome.occurredAt.length > 64)) throw new Error('Invalid merge-queue timestamp.'); + return this.#changeMergeAttempt(identity, id, ['submitting','queued'], attempt => ({ + ...attempt, state: outcome.state, reason: outcome.state === 'merged' ? null : outcome.reason!.trim(), + occurredAt: outcome.occurredAt ?? null, requiresFreshReview: outcome.requiresFreshReview === true, + })); + } + getMergeAttempt(identity: PlanIdentity): MergeAttempt | null { + this.#current(identityKey(identity)); + const row = this.#get('SELECT data FROM merge_attempts WHERE key=? ORDER BY rowid DESC LIMIT 1', identityKey(identity)); + if (!row) return null; + const attempt = decode(row.data); + return { ...attempt, kind: attempt.kind ?? 'queue' }; + } getSuggestions(identity: PlanIdentity, id: string): SuggestionRequest { const row = this.#get('SELECT * FROM requests WHERE key=? AND id=?', identityKey(identity), id); if (!row) throw new Error('Unknown suggestion request.'); diff --git a/test/browser/review.spec.ts b/test/browser/review.spec.ts index 182ac06..0ab1032 100644 --- a/test/browser/review.spec.ts +++ b/test/browser/review.spec.ts @@ -8,7 +8,7 @@ import { createDemo } from '../../scripts/demo.ts'; import { choiceKeys } from '../../core/approvals.ts'; import { ReviewService } from '../../runner/review.ts'; import { startServer } from '../../web/server.ts'; -import type { MergeGateway } from '../../github/merge.ts'; +import type { MergeGateway, MergeQueueGateway } from '../../github/merge.ts'; let root: string, app: Awaited>; function removeDemoOutOfScope(config: typeof app.service.config) { chmodSync(join(config.repository,'run.sh'),0o644);execFileSync('git',['-c','core.hooksPath=/dev/null','commit','-am','Restore declared scope'],{cwd:config.repository,stdio:'pipe'}); } test.beforeEach(async () => { root=mkdtempSync(join(tmpdir(),'codeboost-browser-'));app=await startServer(createDemo(join(root,'demo')),0); }); @@ -49,7 +49,7 @@ test('shows merge blockers and submits one exact-head merge',async({page})=>{ const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let mergeCalls:string[]=[];let release!:()=>void;const held=new Promise(resolve=>{release=resolve;}); const gateway:MergeGateway={ inspect:async()=>{const snapshot=app.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:false,requiredChecks:[],alreadyFixed:'clear'};}, - merge:async head=>{mergeCalls.push(head);await held;app.service.load=()=>{throw new Error('post-command reload failed');};return {url:'https://github.com/example/repo/pull/21'};}, + merge:async head=>{mergeCalls.push(head);await held;app.service.load=()=>{throw new Error('post-command reload failed');};app.service.store.getMergeAttempt=()=>{throw new Error('post-command queue read failed');};return {url:'https://github.com/example/repo/pull/21'};}, }; app=await startServer(config,0,undefined,gateway);await page.goto(app.url); await expect(page.locator('#merge')).toBeDisabled();await page.getByRole('button',{name:'Review blockers'}).click();await expect(page.getByRole('heading',{name:'Merge blockers'})).toBeVisible();await expect(page.getByText(/is unreviewed/).first()).toBeVisible();await page.screenshot({path:'test-results/merge-blockers.png',fullPage:true});await page.getByRole('button',{name:'Close',exact:true}).click(); @@ -57,6 +57,71 @@ test('shows merge blockers and submits one exact-head merge',async({page})=>{ await page.getByRole('button',{name:'Refresh',exact:true}).click();const expectedHead=view.snapshot.head;await expect(page.getByRole('button',{name:'Merge PR',exact:true})).toBeEnabled();page.on('dialog',dialog=>dialog.accept()); await page.locator('#merge').evaluate((button:HTMLButtonElement)=>{button.click();button.click();});await expect.poll(()=>mergeCalls.length).toBe(1);await expect(page.locator('#merge')).toBeDisabled();release();await expect(page.locator('#banner')).toContainText('Merge submitted.');await expect(page.locator('#merge')).toBeDisabled();await page.locator('#merge').evaluate((button:HTMLButtonElement)=>button.click());expect(mergeCalls).toEqual([expectedHead]); }); +test('keeps the reviewed head queued until confirmed merged and preserves current input',async({page})=>{ + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app,phase:'queued'|'merged'='queued',queueReads=0; + const gateway:MergeGateway&MergeQueueGateway={ + inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:true,requiredChecks:[],alreadyFixed:'clear'};}, + queueWatermark:async()=>null, + merge:async()=>({url:'https://github.com/example/repo/pull/24'}), + inspectQueue:async head=>{queueReads++;return phase==='queued'?{state:'queued',reviewedHead:head,entryId:'MQE_1',phase:'AWAITING_CHECKS',position:2,enqueuedAt:'2026-09-24T08:00:00Z',queueHead:head}:{state:'merged',reviewedHead:head,mergedAt:'2026-09-24T08:10:00Z'};}, + }; + app=appRef=await startServer(config,0,undefined,gateway);let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); + await page.goto(app.url);await page.getByRole('button',{name:/P2 Document retry behavior/}).click();await page.getByLabel('Question about this item').fill('Keep this draft while polling.');page.on('dialog',dialog=>dialog.accept());await page.getByRole('button',{name:'Merge PR',exact:true}).click(); + await expect(page.getByRole('button',{name:'Merge queued',exact:true})).toBeDisabled();await expect(page.getByLabel('Question about this item')).toHaveValue('Keep this draft while polling.');await expect.poll(()=>queueReads).toBeGreaterThan(0); + phase='merged';await expect(page.getByRole('button',{name:'Merged',exact:true})).toBeDisabled({timeout:10000});await expect(page.locator('#banner')).toContainText('GitHub confirmed the reviewed head was merged.');await expect(page.getByRole('heading',{name:'Document retry behavior'})).toBeVisible();await expect(page.getByLabel('Question about this item')).toHaveValue('Keep this draft while polling.'); +}); +test('backs off repeated merge-queue polling',async({page})=>{ + await page.addInitScript(()=>{const delays:number[]=[];(window as typeof window&{__mergePollDelays:number[]}).__mergePollDelays=delays;const native=window.setTimeout.bind(window);window.setTimeout=((handler:TimerHandler,timeout?:number,...args:unknown[])=>{if(typeof timeout==='number')delays.push(timeout);return native(handler,timeout,...args);}) as typeof window.setTimeout;}); + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app,queueReads=0; + const gateway:MergeGateway&MergeQueueGateway={inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:true,requiredChecks:[],alreadyFixed:'clear'};},queueWatermark:async()=>null,merge:async()=>({url:'https://github.com/example/repo/pull/24'}),inspectQueue:async head=>{queueReads++;return{state:'queued',reviewedHead:head,entryId:'MQE_1',phase:'QUEUED',position:1,enqueuedAt:'2026-09-24T08:00:00Z',queueHead:head};}}; + app=appRef=await startServer(config,0,undefined,gateway);let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); + await page.goto(app.url);page.on('dialog',dialog=>dialog.accept());await page.getByRole('button',{name:'Merge PR',exact:true}).click();await expect.poll(()=>queueReads,{timeout:10000}).toBeGreaterThanOrEqual(2); + const delays=await page.evaluate(()=>(window as typeof window&{__mergePollDelays:number[]}).__mergePollDelays.filter(value=>value>=500));expect(delays.slice(0,2)).toEqual([2000,4000]); +}); +test('does not queue-poll an ambiguous direct merge attempt',async({page})=>{ + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app,polls=0; + const gateway:MergeGateway={inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:false,requiredChecks:[],alreadyFixed:'clear'};},merge:async()=>({url:''})}; + app=appRef=await startServer(config,0,undefined,gateway);const view=app.service.load();app.service.store.beginMergeAttempt(config.identity,{...view.expected,reviewVersion:view.expected.reviewVersion!},view.snapshot.head,null,'direct'); + await page.route('**/api/merge',async route=>{polls++;await route.continue();});await page.goto(app.url);await expect(page.getByRole('button',{name:'Submitting…',exact:true})).toBeDisabled();await page.waitForTimeout(2500);expect(polls).toBe(0); +}); +test('surfaces queue removal and retries only the same reviewed head',async({page})=>{ + test.slow(); + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app,mergeCalls=0; + const gateway:MergeGateway&MergeQueueGateway={ + inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:true,requiredChecks:[],alreadyFixed:'clear'};}, + queueWatermark:async()=>null, + merge:async()=>{mergeCalls++;return {url:'https://github.com/example/repo/pull/24'};}, + inspectQueue:async head=>mergeCalls===1?{state:'removed',reviewedHead:head,removedAt:'2026-09-24T08:05:00Z',reason:'Required check failed.'}:{state:'queued',reviewedHead:head,entryId:'MQE_2',phase:'QUEUED',position:1,enqueuedAt:'2026-09-24T08:06:00Z',queueHead:head}, + }; + app=appRef=await startServer(config,0,undefined,gateway);let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); + await page.goto(app.url);page.on('dialog',dialog=>dialog.accept());await page.getByRole('button',{name:'Merge PR',exact:true}).click();await expect(page.locator('#banner')).toContainText('Refresh to verify retry readiness.',{timeout:10000});await expect(page.getByRole('button',{name:'Retry merge',exact:true})).toHaveCount(0);await expect(page.locator('#merge')).toBeDisabled();await page.getByRole('button',{name:'Refresh',exact:true}).click();await expect(page.getByRole('button',{name:'Retry merge',exact:true})).toBeEnabled(); + await page.getByRole('button',{name:'Retry merge',exact:true}).click();await expect.poll(()=>mergeCalls).toBe(2);await expect(page.getByRole('button',{name:'Merge queued',exact:true})).toBeDisabled();expect(app.service.store.getMergeAttempt(config.identity)).toMatchObject({state:'queued',reviewedHead:view.snapshot.head}); +}); +test('ignores a merge poll started before a newer review action',async({page})=>{ + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app; + const gateway:MergeGateway&MergeQueueGateway={ + inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:true,requiredChecks:[],alreadyFixed:'clear'};}, + queueWatermark:async()=>null, + merge:async()=>({url:'https://github.com/example/repo/pull/24'}), + inspectQueue:async head=>({state:'removed',reviewedHead:head,removedAt:'2026-09-24T08:05:00Z',reason:'Old poll result.'}), + }; + app=appRef=await startServer(config,0,undefined,gateway);let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); + let release!:()=>void,releaseNew!:()=>void,arrived!:()=>void;const held=new Promise(resolve=>release=resolve),heldNew=new Promise(resolve=>releaseNew=resolve),started=new Promise(resolve=>arrived=resolve);let polls=0; + await page.route('**/api/merge',async route=>{polls++;if(polls===1){arrived();await held;await route.fulfill({status:200,contentType:'application/json',body:JSON.stringify({queue:{state:'removed',reviewedHead:view.snapshot.head,url:null,reason:'Old poll result.',phase:null,position:null,occurredAt:'2026-09-24T08:05:00Z',retryable:true}})});return;}await heldNew;await route.continue();}); + await page.goto(app.url);page.on('dialog',dialog=>dialog.accept());await page.getByRole('button',{name:'Merge PR',exact:true}).click();await started; + await page.getByRole('button',{name:'Request change',exact:true}).click();await page.getByLabel('Change to request').fill('New review decision.');await page.getByRole('button',{name:'Save change request'}).click();await expect(page.getByText('New review decision.',{exact:true})).toBeVisible(); + const staleResponse=page.waitForResponse(response=>response.url().endsWith('/api/merge'));release();try{await staleResponse;await page.waitForTimeout(100);await expect(page.locator('#banner')).not.toContainText('Old poll result.');await expect(page.getByRole('button',{name:'Retry merge',exact:true})).toHaveCount(0);await expect(page.locator('#merge')).toBeDisabled();}finally{releaseNew();} +}); +test('requires fresh review instead of retry when GitHub replaces the queued head',async({page})=>{ + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app; + const gateway:MergeGateway&MergeQueueGateway={ + inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:true,requiredChecks:[],alreadyFixed:'clear'};}, + queueWatermark:async()=>null, + merge:async()=>({url:'https://github.com/example/repo/pull/24'}),inspectQueue:async()=>{throw new Error('The pull request head changed after review.');}, + }; + app=appRef=await startServer(config,0,undefined,gateway);let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); + await page.goto(app.url);page.on('dialog',dialog=>dialog.accept());await page.getByRole('button',{name:'Merge PR',exact:true}).click();await expect(page.locator('#banner')).toContainText('The pull request head changed after review.',{timeout:10000});await expect(page.getByRole('button',{name:'Retry merge',exact:true})).toHaveCount(0);await expect(page.locator('#merge')).toBeDisabled(); +}); test('keeps stale merge failures disabled until refresh',async({page})=>{ const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app;const gateway:MergeGateway={inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:false,requiredChecks:[],alreadyFixed:'clear'};},merge:async()=>{throw new Error('head changed');}};app=appRef=await startServer(config,0,undefined,gateway); let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); @@ -383,6 +448,37 @@ test('drains an in-flight question request before closing its agent manager',asy expect(JSON.parse(response).notes.some((candidate:{text:string})=>candidate.text==='Question during shutdown')).toBe(true); } finally {reopened.close();app=await startServer(config,0);} }); +test('drains an admitted merge request before closing its coordinator',async()=>{ + const config={...app.service.config,demo:false};await app.close();removeDemoOutOfScope(config);let appRef:typeof app,started!:(value?:void)=>void,release!:(value?:void)=>void,commandSignal:AbortSignal|undefined; + const commandStarted=new Promise(resolve=>{started=resolve;}),held=new Promise(resolve=>{release=resolve;}); + const gateway:MergeGateway={inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:false,requiredChecks:[],alreadyFixed:'clear'};},merge:async(_head,options)=>{commandSignal=options?.signal;started();await held;return {url:'https://github.com/example/repo/pull/24'};}}; + app=appRef=await startServer(config,0,undefined,gateway);let view=app.service.load();for(const segment of view.segments.filter(value=>value.row==='Unplanned'||value.row==='Ambiguous'))view=app.service.act({action:'accept',key:segment.key,token:view.token});for(const item of view.items)view=app.service.act({action:'approve',item:item.id,confirmNoChange:item.count===0,token:view.token}); + const endpoint=new URL('/api/action',app.url),body=JSON.stringify({action:'merge',token:view.token});let status=0; + const completed=new Promise((resolve,reject)=>{const req=httpRequest(endpoint,{method:'POST',headers:{'x-codeboost-token':app.token,'content-type':'application/json','content-length':Buffer.byteLength(body)}},res=>{status=res.statusCode??0;res.resume();res.on('end',resolve);});req.on('error',reject);req.end(body);}); + await commandStarted;const closing=app.close();await new Promise(resolve=>setTimeout(resolve,25));expect(commandSignal?.aborted).toBe(false);release();await Promise.all([closing,completed]);expect(status).toBe(200);app=await startServer(config,0); +}); +test('bounds shutdown draining before aborting active queue polling',async()=>{ + const config={...app.service.config,demo:false};await app.close();let appRef:typeof app,started!:(value?:void)=>void,settled=false; + const pollingStarted=new Promise(resolve=>{started=resolve;}); + const gateway:MergeGateway&MergeQueueGateway={inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:true,requiredChecks:[],alreadyFixed:'clear'};},queueWatermark:async()=>null,merge:async()=>({url:'https://github.com/example/repo/pull/24'}),inspectQueue:async(_head,options)=>new Promise((_resolve,reject)=>{started();options?.signal?.addEventListener('abort',()=>{settled=true;reject(options.signal?.reason);},{once:true});})}; + app=appRef=await startServer(config,0,undefined,gateway,50);const view=app.service.load(),attempt=app.service.store.beginMergeAttempt(config.identity,{...view.expected,reviewVersion:view.expected.reviewVersion!},view.snapshot.head);app.service.store.queueMergeAttempt(config.identity,attempt.id,'https://github.com/example/repo/pull/24'); + const response=fetch(new URL('/api/merge',app.url),{headers:{'x-codeboost-token':app.token}});await pollingStarted;await app.close();expect(settled).toBe(true);expect((await response).status).toBe(409);app=await startServer(config,0); +}); +test('aborts an admitted review status inspection after the shutdown drain',async()=>{ + const config={...app.service.config,demo:false};await app.close();let started!:(value?:void)=>void,settled=false;const inspectionStarted=new Promise(resolve=>{started=resolve;}); + const gateway:MergeGateway={inspect:async options=>new Promise((_resolve,reject)=>{started();options?.signal?.addEventListener('abort',()=>{settled=true;reject(options.signal?.reason);},{once:true});}),merge:async()=>({url:''})}; + app=await startServer(config,0,undefined,gateway,50);const response=fetch(new URL('/api/review',app.url),{headers:{'x-codeboost-token':app.token}});await inspectionStarted;await app.close();expect(settled).toBe(true);expect((await response).status).toBe(409);app=await startServer(config,0); +}); +test('destroys a partial request body after the shutdown drain',async()=>{ + const config=app.service.config;await app.close();app=await startServer(config,0,undefined,undefined,50); + const endpoint=new URL('/api/action',app.url),body=JSON.stringify({action:'note'}); + let admitted!:(value?:void)=>void;const requestAdmitted=new Promise(resolve=>{admitted=resolve;}); + const completed=new Promise<'response'|'error'>(resolve=>{ + const req=httpRequest(endpoint,{method:'POST',headers:{'x-codeboost-token':app.token,'content-type':'application/json','content-length':Buffer.byteLength(body)}},res=>{res.resume();res.on('end',()=>resolve('response'));}); + req.on('error',()=>resolve('error'));req.write(body.slice(0,1),()=>admitted()); + }); + await requestAdmitted;const started=Date.now();await app.close();expect(Date.now()-started).toBeLessThan(1000);expect(await completed).toBe('error');app=await startServer(config,0); +}); test('blocks a partially received merge request when shutdown starts',async()=>{ const config={...app.service.config,demo:false};await app.close();let mergeCalls=0;let appRef:typeof app; const gateway:MergeGateway={inspect:async()=>{const snapshot=appRef.service.load().snapshot;return {base:snapshot.base,head:snapshot.head,pullRequestState:'OPEN',mergeable:'MERGEABLE',rulesKnown:true,atomicBaseGuard:true,mergeQueue:false,requiredChecks:[],alreadyFixed:'clear'};},merge:async()=>{mergeCalls++;return {url:''};}}; diff --git a/test/merge.test.ts b/test/merge.test.ts index 908a573..4509db3 100644 --- a/test/merge.test.ts +++ b/test/merge.test.ts @@ -1,7 +1,8 @@ import { expect, it, vi } from 'vitest'; import { ReviewService } from '../runner/review.ts'; import { MergeCoordinator } from '../runner/merge.ts'; -import { GhMergeGateway, type MergeGateway, type RemoteMergeState } from '../github/merge.ts'; +import { GhMergeGateway, MergeSubmissionError, type MergeGateway, type MergeQueueGateway, type MergeQueueObservation, type RemoteMergeState } from '../github/merge.ts'; +import { Store } from '../runner/store.ts'; type ReviewView = ReturnType; const sha = (digit: string) => digit.repeat(40); @@ -22,6 +23,31 @@ function gateway(states: RemoteMergeState[]): MergeGateway & { heads: string[] } return { heads, inspect: vi.fn(async () => states.shift() ?? states.at(-1)!), merge: vi.fn(async head => { heads.push(head); return { url: 'https://github.example/pr/1' }; }) }; } +function queueHarness(observations: Array) { + const identity = { repositoryId: 'repo', taskId: 'task', planId: 'plan' }; + const store = new Store(':memory:'); + const plan = { schema_version: 1 as const, revision: 1, issue: 24, summary: 'Queue', questions: [], items: [{ id: 'P1', title: 'Queue', intent: 'Queue safely', files: [{ path: 'a', kind: 'edit' as const, renamed_from: null, change: 'Change' }], acceptance: [{ type: 'check' as const, text: 'Works' }], depends_on: [] }] }; + const context = { identity, issue: 24, baseEntries: [{ path: 'a', kind: 'file' as const }], pathKey: (path: string) => path, allowedCommands: [] }; + store.createPlan(JSON.stringify(plan), 'json', context, sha('a'), sha('b')); + let view = { ...readyView(), expected: { revision: 1, snapshotId: store.getSnapshot(identity).id, reviewVersion: store.reviewVersion(identity) } } as ReviewView; + const service = { store, config: { identity }, load: vi.fn(() => view) } as unknown as ReviewService; + const merges: string[] = []; + const client: MergeGateway & MergeQueueGateway = { + inspect: vi.fn(async () => remote(view, { mergeQueue: true })), + queueWatermark: vi.fn(async () => 'CURSOR_before'), + merge: vi.fn(async head => { merges.push(head); return { url: 'https://github.example/pr/1' }; }), + inspectQueue: vi.fn(async () => { const next = observations.shift(); if (next instanceof Error) throw next; if (!next) throw new Error('No queue observation.'); return next; }), + }; + return { store, identity, service, client, merges, coordinator: new MergeCoordinator(service, client), view: () => view, amendPlan() { + const amended = store.importRevision(JSON.stringify({ ...plan, summary: 'Amended queue plan' }), 'json', context, 1); + view = { ...view, plan: amended, expected: { revision: amended.revision, snapshotId: view.expected.snapshotId, reviewVersion: store.reviewVersion(identity) }, token: 'review-amended-plan' } as ReviewView; + }, replaceHead(head: string) { + const expected = view.expected; + const snapshot = store.recordHistory(identity, expected, sha('a'), head, []); + view = { ...view, snapshot, expected: { revision: 1, snapshotId: snapshot.id, reviewVersion: store.reviewVersion(identity) }, token: `review-${head}` } as ReviewView; + }, changeToken(token: string) { view = { ...view, token } as ReviewView; } }; +} + it('lists every local review blocker before merge', async () => { const view = { ...readyView(), items: [{ ...readyView().items[0]!, state: 'stale', outside: ['undeclared.ts'] }], segments: [{ row: 'Unplanned' }], notes: [{ kind: 'change', revision: 1, snapshotId: 'snapshot' }] } as unknown as ReviewView; const service = serviceFor(view); @@ -70,11 +96,22 @@ it('rechecks the exact base and head, then invokes the guarded head merge', asyn expect(merged.result.url).toContain('/pr/1'); }); -it('budgets both fresh validation passes below the serving deadline', async () => { - const view = readyView(), service = serviceFor(view), options: Array<{ fresh?: boolean; timeoutMs?: number } | undefined> = []; - const client: MergeGateway = { inspect: vi.fn(async value => { options.push(value); return remote(view); }), merge: vi.fn(async () => ({ url: 'https://github.example/pr/1' })) }; - await new MergeCoordinator(service, client).merge(view.token); - expect(options).toEqual([{ fresh: true, timeoutMs: 6_000 }, { fresh: true, timeoutMs: 6_000 }]); +it('enforces one deadline across every merge validation stage', async () => { + const h = queueHarness([]); + h.client.inspect = vi.fn(async value => { + await new Promise((resolve, reject) => { + const timer = setTimeout(resolve, 100); + const signal = (value as { signal?: AbortSignal } | undefined)?.signal; + signal?.addEventListener('abort', () => { clearTimeout(timer); reject(signal.reason); }, { once: true }); + }); + return remote(h.view(), { mergeQueue: true }); + }); + try { + const started = Date.now(); + await expect(new MergeCoordinator(h.service, h.client, 250).merge(h.view().token)).rejects.toThrow(/deadline/i); + expect(Date.now() - started).toBeLessThan(500); + expect(h.client.merge).not.toHaveBeenCalled(); + } finally { h.store.close(); } }); it('refuses a base or head race after the initial validation', async () => { @@ -106,6 +143,20 @@ it('preserves the GitHub merge refusal', async () => { await expect(new MergeCoordinator(service, client).merge(view.token)).rejects.toThrow('Required review is missing.'); }); +it('persists and reconciles an ambiguous direct merge outcome', async () => { + const h = queueHarness([]);let state=remote(h.view(),{mergeQueue:false}); + h.client.inspect=vi.fn(async()=>state); + h.client.merge=vi.fn(async()=>{throw new MergeSubmissionError('Direct merge response was lost.','unknown');}); + try { + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow(/response was lost/i); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({kind:'direct',state:'submitting',reason:'Direct merge response was lost.'}); + expect(await h.coordinator.status(h.view())).toMatchObject({ready:false,action:null,queue:{kind:'direct',state:'submitting'}}); + state={...state,pullRequestState:'MERGED'}; + expect(await h.coordinator.status(h.view())).toMatchObject({ready:false,action:null,queue:{state:'merged'}}); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({kind:'direct',state:'merged'}); + } finally {await h.coordinator.close();h.store.close();} +}); + it('aborts and awaits an active merge command during shutdown', async () => { const view = readyView(), service = serviceFor(view); let commandStarted!: () => void, commandSettled = false; @@ -126,6 +177,224 @@ it('aborts and awaits an active merge command during shutdown', async () => { await expect(coordinator.merge(view.token)).rejects.toThrow(/shutting down/i); }); +it('persists enqueue success as queued and waits for a separate confirmed merge', async () => { + const h = queueHarness([ + { state: 'queued', reviewedHead: sha('b'), entryId: 'MQE_1', phase: 'AWAITING_CHECKS', position: 2, enqueuedAt: '2026-09-24T08:00:00Z', queueHead: sha('b') }, + { state: 'merged', reviewedHead: sha('b'), mergedAt: '2026-09-24T08:10:00Z' }, + ]); + try { + await h.coordinator.merge(h.view().token); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'queued', reviewedHead: sha('b') }); + expect((await h.coordinator.pollQueue())).toMatchObject({ state: 'queued', phase: 'AWAITING_CHECKS', position: 2 }); + expect(h.client.inspectQueue).toHaveBeenCalledWith(sha('b'), expect.objectContaining({ afterCursor: 'CURSOR_before' })); + expect((await h.coordinator.pollQueue())).toMatchObject({ state: 'merged', occurredAt: '2026-09-24T08:10:00Z', retryable: false }); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('keeps an externally successful enqueue committed when the queued-state refresh fails', async () => { + const h = queueHarness([{ state: 'queued', reviewedHead: sha('b'), entryId: 'MQE_1', phase: 'QUEUED', position: 1, enqueuedAt: '2026-09-24T08:00:00Z', queueHead: sha('b') }]); + const queue = h.store.queueMergeAttempt.bind(h.store); + h.store.queueMergeAttempt = vi.fn(() => { throw new Error('local refresh failed'); }); + try { + await expect(h.coordinator.merge(h.view().token)).resolves.toMatchObject({ result: { url: 'https://github.example/pr/1' } }); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'submitting', reviewedHead: sha('b') }); + h.store.queueMergeAttempt = queue; + expect(await h.coordinator.pollQueue()).toMatchObject({ state: 'queued', phase: 'QUEUED' }); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it.each([ + [{ state: 'removed', reviewedHead: sha('b'), removedAt: '2026-09-24T08:05:00Z', reason: 'Checks failed.' } as const, 'removed', 'Checks failed.'], + [{ state: 'failed', reviewedHead: sha('b'), entryId: 'MQE_1', reason: 'GitHub reported the merge queue entry as unmergeable.' } as const, 'failed', 'GitHub reported the merge queue entry as unmergeable.'], +])('persists a terminal %s queue result and safely enables retry', async (observation, state, reason) => { + const h = queueHarness([observation]); + try { + await h.coordinator.merge(h.view().token); + expect(await h.coordinator.pollQueue()).toMatchObject({ state, reason, retryable: true }); + expect((await h.coordinator.status(h.view())).action).toBe('retry'); + await h.coordinator.merge(h.view().token); + expect(h.merges).toEqual([sha('b'), sha('b')]); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'queued', reviewedHead: sha('b') }); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('requires current-snapshot approvals after an ordinary queue removal and force-push', async () => { + const h = queueHarness([{ state: 'removed', reviewedHead: sha('b'), removedAt: '2026-09-24T08:05:00Z', reason: 'Checks failed.' }]); + try { + await h.coordinator.merge(h.view().token); + await h.coordinator.pollQueue(); + h.replaceHead(sha('c')); + expect((await h.coordinator.status(h.view())).action).toBeNull(); + h.store.saveReview(h.identity, h.view().expected, [{ item: 'P1', fingerprint: 'replacement-review' }], []); + expect((await h.coordinator.status(h.view())).action).toBe('merge'); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('requires current-revision approvals after a same-snapshot plan amendment', async () => { + const h = queueHarness([{ state: 'removed', reviewedHead: sha('b'), removedAt: '2026-09-24T08:05:00Z', reason: 'Checks failed.' }]); + try { + await h.coordinator.merge(h.view().token); + await h.coordinator.pollQueue(); + h.amendPlan(); + expect((await h.coordinator.status(h.view())).action).toBeNull(); + h.store.saveReview(h.identity, h.view().expected, [{ item: 'P1', fingerprint: 'amended-plan-review' }], []); + expect((await h.coordinator.status(h.view())).action).toBe('merge'); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('requires a fresh review after the queued head is replaced', async () => { + const h = queueHarness([new Error('The pull request head changed after review.')]); + try { + await h.coordinator.merge(h.view().token); + expect(await h.coordinator.pollQueue()).toMatchObject({ state: 'failed', retryable: false, reason: 'The pull request head changed after review.' }); + expect((await h.coordinator.status(h.view())).action).toBeNull(); + h.replaceHead(sha('c')); + expect((await h.coordinator.status(h.view())).action).toBeNull(); + h.store.saveReview(h.identity, h.view().expected, [{ item: 'P1', fingerprint: 'fresh-review' }], []); + expect((await h.coordinator.status(h.view())).action).toBe('merge'); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('accepts post-attempt approvals when GitHub reports a transient head replacement', async () => { + const h = queueHarness([new Error('The pull request head changed after review.')]); + try { + await h.coordinator.merge(h.view().token); + await h.coordinator.pollQueue(); + expect((await h.coordinator.status(h.view())).action).toBeNull(); + h.store.saveReview(h.identity, h.view().expected, [{ item: 'P1', fingerprint: 'post-attempt-review' }], []); + expect((await h.coordinator.status(h.view())).action).toBe('merge'); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('accepts a fully re-reviewed replacement snapshot restored to the original head', async () => { + const h = queueHarness([new Error('The pull request head changed after review.')]); + try { + await h.coordinator.merge(h.view().token); + await h.coordinator.pollQueue(); + h.replaceHead(sha('c')); + h.replaceHead(sha('b')); + h.store.saveReview(h.identity, h.view().expected, [{ item: 'P1', fingerprint: 'restored-head-review' }], []); + expect((await h.coordinator.status(h.view())).action).toBe('merge'); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('refuses a merge-queue mode change between validation passes', async () => { + const h = queueHarness([]); + h.client.inspect = vi.fn() + .mockResolvedValueOnce(remote(h.view(), { mergeQueue: false })) + .mockResolvedValueOnce(remote(h.view(), { mergeQueue: true })); + try { + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow(/queue|requirements changed/i); + expect(h.merges).toEqual([]); + expect(h.store.getMergeAttempt(h.identity)).toBeNull(); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('revalidates the review generation after reading the queue watermark', async () => { + const h = queueHarness([]); + h.client.queueWatermark = vi.fn(async () => { h.changeToken('new-review-token'); return 'CURSOR_before'; }); + try { + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow(/review changed during merge validation/i); + expect(h.merges).toEqual([]); + expect(h.store.getMergeAttempt(h.identity)).toBeNull(); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('refuses a queue-mode change observed after reading the queue watermark', async () => { + const h = queueHarness([]); + h.client.inspect = vi.fn() + .mockResolvedValueOnce(remote(h.view(), { mergeQueue: true })) + .mockResolvedValueOnce(remote(h.view(), { mergeQueue: true })) + .mockResolvedValueOnce(remote(h.view(), { mergeQueue: false })); + try { + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow(/queue|requirements changed/i); + expect(h.merges).toEqual([]); + expect(h.store.getMergeAttempt(h.identity)).toBeNull(); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('keeps retry disabled while a prior queue attempt is active', async () => { + const h = queueHarness([]); + try { + await h.coordinator.merge(h.view().token); + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow(/queued|active/i); + expect(h.merges).toEqual([sha('b')]); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('keeps an active queue attempt visible when inspection fails', async () => { + const h = queueHarness([new Error('GitHub queue status unavailable.')]); + try { + await h.coordinator.merge(h.view().token); + await expect(h.coordinator.pollQueue()).resolves.toMatchObject({ + state: 'queued', + reviewedHead: sha('b'), + observationError: 'GitHub queue status unavailable.', + }); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'queued', reviewedHead: sha('b') }); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('aborts and awaits an active queue inspection during shutdown', async () => { + const h = queueHarness([]); + let settled = false; + h.client.inspectQueue = vi.fn(async (_head, options) => new Promise((_resolve, reject) => { + options?.signal?.addEventListener('abort', () => { settled = true; reject(options.signal?.reason); }, { once: true }); + })); + try { + await h.coordinator.merge(h.view().token); + const polling = h.coordinator.pollQueue(); + await h.coordinator.close(); + await expect(polling).rejects.toThrow(/shutdown/i); + expect(settled).toBe(true); + } finally { h.store.close(); } +}); + +it('keeps an aborted enqueue submitting until queue inspection recovers its outcome', async () => { + const h = queueHarness([{ state: 'queued', reviewedHead: sha('b'), entryId: 'MQE_1', phase: 'QUEUED', position: 1, enqueuedAt: '2026-09-24T08:00:00Z', queueHead: sha('b') }]); + let started!: () => void; + const commandStarted = new Promise(resolve => { started = resolve; }); + h.client.merge = vi.fn(async (_head, options) => new Promise((_resolve, reject) => { + started(); options?.signal?.addEventListener('abort', () => reject(options.signal?.reason), { once: true }); + })); + const merging = h.coordinator.merge(h.view().token); + await commandStarted; + await h.coordinator.close(); + await expect(merging).rejects.toThrow(/shutdown/i); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'submitting', reviewedHead: sha('b') }); + const recovered = new MergeCoordinator(h.service, h.client); + try { expect(await recovered.pollQueue()).toMatchObject({ state: 'queued', phase: 'QUEUED', retryable: false }); } + finally { await recovered.close(); h.store.close(); } +}); + +it('records only a confirmed enqueue refusal as retryable failure', async () => { + const h = queueHarness([]); + h.client.merge = vi.fn(async () => { throw new MergeSubmissionError('Required review is missing.', 'refused'); }); + try { + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow('Required review is missing.'); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'failed', reason: 'Required review is missing.' }); + expect(await h.coordinator.status()).toMatchObject({ ready: true, action: 'retry', queue: { state: 'failed', retryable: true } }); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('keeps an unknown enqueue failure submitting until external reconciliation', async () => { + const h = queueHarness([]); + h.client.merge = vi.fn(async () => { throw new MergeSubmissionError('GitHub merge submission timed out.', 'unknown'); }); + try { + await expect(h.coordinator.merge(h.view().token)).rejects.toThrow(/timed out/i); + expect(h.store.getMergeAttempt(h.identity)).toMatchObject({ state: 'submitting', reason: 'GitHub merge submission timed out.' }); + expect(await h.coordinator.status()).toMatchObject({ ready: false, action: null, queue: { state: 'submitting', retryable: false, reason: 'GitHub merge submission timed out.' } }); + } finally { await h.coordinator.close(); h.store.close(); } +}); + +it('classifies only explicit GitHub merge refusals as confirmed', async () => { + const config = { repository: 'owner/repo', pullRequest: 7, issue: 24 }; + const refusal = new GhMergeGateway(config, async () => { throw new Error('Required review is missing.'); }); + const unknown = new GhMergeGateway(config, async () => { throw new Error('request timed out'); }); + await expect(refusal.merge(sha('b'))).rejects.toMatchObject({ name: 'MergeSubmissionError', outcome: 'refused' }); + await expect(unknown.merge(sha('b'))).rejects.toMatchObject({ name: 'MergeSubmissionError', outcome: 'unknown' }); +}); + it('parses required checks from both rule sources and pins the gh merge head', async () => { const calls: string[][] = []; let pullReads = 0; @@ -162,8 +431,25 @@ it('rejects an unsupported runtime merge method', () => { expect(() => new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 21, method: 'typo' as 'merge' })).toThrow(/merge method/i); }); -const queueFixture = (pullRequest: Record) => JSON.stringify({ - data: { repository: { pullRequest: { number: 7, headRefOid: sha('b'), ...pullRequest } } }, +const queueFixture = (pullRequest: Record) => { + const timeline = pullRequest.timelineItems as { nodes?: unknown[] } | undefined; + const normalized = timeline && Array.isArray(timeline.nodes) ? { + ...pullRequest, + timelineItems: { + edges: timeline.nodes.map((node, index) => ({ + cursor: node && typeof node === 'object' && !Array.isArray(node) && typeof (node as { cursor?: unknown }).cursor === 'string' + ? (node as { cursor: string }).cursor : `CURSOR_${index}`, + node, + })), + pageInfo: { hasNextPage: false, endCursor: timeline.nodes.length ? `CURSOR_${timeline.nodes.length - 1}` : null }, + }, + } : pullRequest; + return JSON.stringify({ data: { repository: { pullRequest: { number: 7, headRefOid: sha('b'), ...normalized } } } }); +}; + +it('captures the stable queue timeline cursor before enqueue', async () => { + const run = async () => queueFixture({ timelineItems: { nodes: [{ id: 'MQEV_before', cursor: 'CURSOR_before' }] } }); + await expect(new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run).queueWatermark(sha('b'))).resolves.toBe('CURSOR_before'); }); it.each([ @@ -173,7 +459,7 @@ it.each([ ['MERGEABLE', 'queued'], ] as const)('maps the recorded %s merge-queue entry to %s', async (entryState, expectedState) => { const run = async () => queueFixture({ - state: 'OPEN', mergedAt: null, timelineItems: { nodes: [{ __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T08:00:00Z' }] }, + state: 'OPEN', mergedAt: null, timelineItems: { nodes: [{ id: 'MQEV_1', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T08:00:00Z' }] }, mergeQueueEntry: { id: 'MQE_1', state: entryState, position: 2, enqueuedAt: '2026-09-24T08:00:00Z', headCommit: { oid: sha('b') }, pullRequest: { number: 7, headRefOid: sha('b') } }, }); const observation = await new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run).inspectQueue(sha('b')); @@ -182,7 +468,7 @@ it.each([ it('maps a recorded unmergeable queue entry to a failed terminal state', async () => { const run = async () => queueFixture({ - state: 'OPEN', mergedAt: null, timelineItems: { nodes: [{ __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T08:00:00Z' }] }, + state: 'OPEN', mergedAt: null, timelineItems: { nodes: [{ id: 'MQEV_1', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T08:00:00Z' }] }, mergeQueueEntry: { id: 'MQE_1', state: 'UNMERGEABLE', position: 1, enqueuedAt: '2026-09-24T08:00:00Z', headCommit: { oid: sha('b') }, pullRequest: { number: 7, headRefOid: sha('b') } }, }); await expect(new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run).inspectQueue(sha('b'))).resolves.toEqual({ @@ -193,8 +479,8 @@ it('maps a recorded unmergeable queue entry to a failed terminal state', async ( it('preserves the recorded reason when GitHub removes a pull request from the merge queue', async () => { const run = async () => queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [ - { __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T08:00:00Z' }, - { __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T08:05:00Z', reason: 'Checks failed', beforeCommit: { oid: sha('b') } }, + { id: 'MQEV_1', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T08:00:00Z' }, + { id: 'MQEV_2', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T08:05:00Z', reason: 'Checks failed', beforeCommit: { oid: sha('b') } }, ] }, }); await expect(new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run).inspectQueue(sha('b'))).resolves.toEqual({ @@ -202,6 +488,79 @@ it('preserves the recorded reason when GitHub removes a pull request from the me }); }); +it('rejects a removal event at the pre-enqueue timeline cursor', async () => { + const run = async () => queueFixture({ + state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [] }, + }); + const client = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run); + await expect(client.inspectQueue(sha('b'), { afterCursor: 'CURSOR_removed_old' })).rejects.toThrow(/current attempt/i); +}); + +it('recovers terminal state after the stored cursor falls outside the recent event window', async () => { + const run = async () => queueFixture({ + state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [ + { id: 'MQEV_added_current', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T09:00:00Z' }, + { id: 'MQEV_removed_current', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T09:05:00Z', reason: 'Current attempt failed', beforeCommit: { oid: sha('b') } }, + ] }, + }); + const client = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run); + await expect(client.inspectQueue(sha('b'), { afterCursor: 'CURSOR_before_outside_window' })).resolves.toMatchObject({ + state: 'removed', reason: 'Current attempt failed', removedAt: '2026-09-24T09:05:00Z', + }); +}); + +it('paginates forward from the stored cursor to the terminal event', async () => { + let calls = 0; + const response = (edge: Record, hasNextPage: boolean, endCursor: string | null) => JSON.stringify({ data: { repository: { pullRequest: { + number: 7, headRefOid: sha('b'), state: 'OPEN', mergedAt: null, mergeQueueEntry: null, + timelineItems: { edges: [edge], pageInfo: { hasNextPage, endCursor } }, + } } } }); + const run = async (args: readonly string[]) => { + calls++; + if (calls === 1) { + expect(args).toContain('after=CURSOR_before'); + return response({ cursor: 'CURSOR_added', node: { id: 'MQEV_added', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T09:00:00Z' } }, true, 'CURSOR_added'); + } + expect(args).toContain('after=CURSOR_added'); + return response({ cursor: 'CURSOR_removed', node: { id: 'MQEV_removed', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T09:05:00Z', reason: 'Checks failed', beforeCommit: { oid: sha('b') } } }, false, 'CURSOR_removed'); + }; + const client = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run); + await expect(client.inspectQueue(sha('b'), { afterCursor: 'CURSOR_before' })).resolves.toMatchObject({ state: 'removed', reason: 'Checks failed' }); + expect(calls).toBe(2); +}); + +it('reads the tenth queue-history page and fails closed beyond it', async () => { + const response = (calls: number, hasNextPage: boolean) => JSON.stringify({ data: { repository: { pullRequest: { + number: 7, headRefOid: sha('b'), state: 'OPEN', mergedAt: null, mergeQueueEntry: null, + timelineItems: { edges: calls === 10 ? [ + { cursor: 'CURSOR_added', node: { id: 'MQEV_added', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T09:00:00Z' } }, + { cursor: 'CURSOR_removed', node: { id: 'MQEV_removed', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T09:05:00Z', reason: 'Checks failed', beforeCommit: { oid: sha('b') } } }, + ] : [], pageInfo: { hasNextPage, endCursor: `CURSOR_${calls}` } }, + } } } }); + let calls = 0; + const client = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, async () => response(++calls, calls < 10)); + await expect(client.inspectQueue(sha('b'), { afterCursor: 'CURSOR_before' })).resolves.toMatchObject({ state: 'removed', reason: 'Checks failed' }); + expect(calls).toBe(10); + + calls = 0; + const overLimit = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, async () => response(++calls, true)); + await expect(overLimit.inspectQueue(sha('b'), { afterCursor: 'CURSOR_before' })).rejects.toThrow(/exceeded the inspection limit/i); + expect(calls).toBe(10); +}); + +it('fails closed when multiple enqueue sequences follow the stored cursor', async () => { + const run = async () => queueFixture({ + state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [ + { id: 'MQEV_add_1', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T09:00:00Z' }, + { id: 'MQEV_remove_1', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T09:01:00Z', reason: 'First removal', beforeCommit: { oid: sha('b') } }, + { id: 'MQEV_add_2', __typename: 'AddedToMergeQueueEvent', createdAt: '2026-09-24T09:02:00Z' }, + { id: 'MQEV_remove_2', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T09:03:00Z', reason: 'Second removal', beforeCommit: { oid: sha('b') } }, + ] }, + }); + const client = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 24 }, run); + await expect(client.inspectQueue(sha('b'), { afterCursor: 'CURSOR_before' })).rejects.toThrow(/multiple|current enqueue attempt/i); +}); + it('reports merged only when GitHub confirms the reviewed head was merged', async () => { let call: readonly string[] = []; const run = async (args: readonly string[]) => { @@ -213,7 +572,7 @@ it('reports merged only when GitHub confirms the reviewed head was merged', asyn }); expect(call).toEqual([ 'api', 'graphql', '-f', - 'query=query($owner:String!,$name:String!,$number:Int!){repository(owner:$owner,name:$name){pullRequest(number:$number){number headRefOid state mergedAt mergeQueueEntry{id state position enqueuedAt headCommit{oid} pullRequest{number headRefOid}} timelineItems(last:20,itemTypes:[ADDED_TO_MERGE_QUEUE_EVENT,REMOVED_FROM_MERGE_QUEUE_EVENT]){nodes{__typename ... on AddedToMergeQueueEvent{createdAt} ... on RemovedFromMergeQueueEvent{createdAt reason beforeCommit{oid}}}}}}', + 'query=query($owner:String!,$name:String!,$number:Int!,$after:String){repository(owner:$owner,name:$name){pullRequest(number:$number){number headRefOid state mergedAt mergeQueueEntry{id state position enqueuedAt headCommit{oid} pullRequest{number headRefOid}} timelineItems(first:100,after:$after,itemTypes:[ADDED_TO_MERGE_QUEUE_EVENT,REMOVED_FROM_MERGE_QUEUE_EVENT]){edges{cursor node{id __typename ... on AddedToMergeQueueEvent{createdAt} ... on RemovedFromMergeQueueEvent{createdAt reason beforeCommit{oid}}}} pageInfo{hasNextPage endCursor}}}}}', '-f', 'owner=owner', '-f', 'name=repo', '-F', 'number=7', ]); }); @@ -221,9 +580,10 @@ it('reports merged only when GitHub confirms the reviewed head was merged', asyn it.each([ ['a replaced reviewed head', queueFixture({ state: 'OPEN', mergedAt: null, headRefOid: sha('d'), mergeQueueEntry: { id: 'MQE_1', state: 'QUEUED', position: 1, enqueuedAt: '2026-09-24T08:00:00Z', headCommit: { oid: sha('d') }, pullRequest: { number: 7, headRefOid: sha('d') } }, timelineItems: { nodes: [] } })], ['a queue entry for another head', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: { id: 'MQE_1', state: 'QUEUED', position: 1, enqueuedAt: '2026-09-24T08:00:00Z', headCommit: { oid: sha('d') }, pullRequest: { number: 7, headRefOid: sha('b') } }, timelineItems: { nodes: [] } })], - ['a stale removal from another head', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [{ __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T08:05:00Z', reason: 'Checks failed', beforeCommit: { oid: sha('d') } }] } })], + ['an unmergeable queue entry for another head', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: { id: 'MQE_1', state: 'UNMERGEABLE', position: 1, enqueuedAt: '2026-09-24T08:00:00Z', headCommit: { oid: sha('d') }, pullRequest: { number: 7, headRefOid: sha('b') } }, timelineItems: { nodes: [] } })], + ['a stale removal from another head', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [{ id: 'MQEV_1', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T08:05:00Z', reason: 'Checks failed', beforeCommit: { oid: sha('d') } }] } })], ['an absent queue entry without a removal event', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [] } })], - ['a removal without a reason', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [{ __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T08:05:00Z', reason: null, beforeCommit: { oid: sha('b') } }] } })], + ['a removal without a reason', queueFixture({ state: 'OPEN', mergedAt: null, mergeQueueEntry: null, timelineItems: { nodes: [{ id: 'MQEV_1', __typename: 'RemovedFromMergeQueueEvent', createdAt: '2026-09-24T08:05:00Z', reason: null, beforeCommit: { oid: sha('b') } }] } })], ['an omitted mergeQueueEntry field', queueFixture({ state: 'MERGED', mergedAt: '2026-09-24T08:10:00Z', timelineItems: { nodes: [] } })], ['an omitted timelineItems field', queueFixture({ state: 'MERGED', mergedAt: '2026-09-24T08:10:00Z', mergeQueueEntry: null })], ['a malformed timeline node on a merged response', queueFixture({ state: 'MERGED', mergedAt: '2026-09-24T08:10:00Z', mergeQueueEntry: null, timelineItems: { nodes: [null] } })], @@ -305,6 +665,16 @@ it('aborts one shared status inspection at its overall deadline', async () => { } finally { vi.useRealTimers(); } }); +it('preserves caller cancellation while reading merge status', async () => { + const controller = new AbortController(); + const run = async (_args: readonly string[], options?: { signal?: AbortSignal }) => new Promise((_resolve, reject) => { + options?.signal?.addEventListener('abort', () => reject(new Error('generic runner abort')), { once: true }); + }); + const pending = new GhMergeGateway({ repository: 'owner/repo', pullRequest: 7, issue: 21 }, run).inspect({ fresh: true, signal: controller.signal }); + controller.abort(new Error('merge request deadline exceeded')); + await expect(pending).rejects.toThrow('merge request deadline exceeded'); +}); + it.each([[false, true], [true, false]])('treats a protection 404 with protected=%s as rulesKnown=%s', async (protectedBranch, expectedKnown) => { const run = async (args: readonly string[]) => { const joined = args.join(' '); diff --git a/test/questions.test.ts b/test/questions.test.ts index ab6a25e..5e3692e 100644 --- a/test/questions.test.ts +++ b/test/questions.test.ts @@ -21,7 +21,7 @@ it('persists answers with plan, code, selected snippet and prior conversation co expect(prompt).toContain('Why cap the retry delay?');expect(prompt).toContain('Math.min');expect(prompt).toContain('Keep the signature.'); const after=service.load();expect(after.plan.revision).toBe(asked.plan.revision);expect(after.token).toBe(asked.token);expect(after.approved).toBe(0); const reopened=new ReviewService(service.config);services.push(reopened);expect(reopened.load().notes.at(-1)?.answer?.text).toContain('bounds retry latency'); -}); +},30_000); it('fails visibly and retries without duplicating the question or accepting stale completions',async()=>{ const service=fixture(),asked=question(service);let calls=0; const manager=new Questions(service,async()=>{if(++calls===1)throw new Error('Login required');return 'Recovered answer';});managers.push(manager);manager.start(asked.createdNoteId!,asked); diff --git a/test/review.test.ts b/test/review.test.ts index 944aba8..4a98a4a 100644 --- a/test/review.test.ts +++ b/test/review.test.ts @@ -7,6 +7,7 @@ import { DatabaseSync } from 'node:sqlite'; import { createDemo } from '../scripts/demo.ts'; import { ReviewService } from '../runner/review.ts'; import { startServer } from '../web/server.ts'; +import { approveItem } from '../core/approvals.ts'; // Each integration case performs several bounded real-Git reads. vi.setConfig({ testTimeout: 15000 }); const roots:string[]=[];const services:ReviewService[]=[]; @@ -57,6 +58,24 @@ it('rejects a concurrent store review edit through the atomic review counter',() other.store.addReviewNote(config.identity,view.expected,'P1','change','Please explain'); expect(()=>service.store.saveReview(config.identity,view.expected,[],[])).toThrow(/Stale/); }); +it('requires every item to be reviewed again after a queued head is replaced',()=>{ + const {service,config}=fixture();let view=service.load(); + service.store.saveReview(config.identity,view.expected,view.items.map(item=>approveItem(view.plan,view.segments,item.id,config.identity,item.count===0)),[]);view=service.load(); + const attempt=service.store.beginMergeAttempt(config.identity,{...view.expected,reviewVersion:view.expected.reviewVersion!},view.snapshot.head);service.store.queueMergeAttempt(config.identity,attempt.id,'https://github.example/pr/24');service.store.finishMergeAttempt(config.identity,attempt.id,{state:'failed',reason:'The pull request head changed after review.',requiresFreshReview:true}); + execFileSync('git',['-c','core.hooksPath=/dev/null','commit','--allow-empty','-m','Replace reviewed head'],{cwd:config.repository,stdio:'pipe'}); + view=service.load();expect(view.items.every(item=>item.state==='stale'&&item.reasons.includes('Pull request snapshot changed after the queue attempt'))).toBe(true); + const [first,...remaining]=view.items;service.store.saveReview(config.identity,view.expected,[approveItem(view.plan,view.segments,first!.id,config.identity,first!.count===0)],[]);view=service.load();expect(view.items.find(item=>item.id===first!.id)?.state).toBe('approved');expect(view.items.filter(item=>item.id!==first!.id).every(item=>item.state==='stale')).toBe(true); + service.store.saveReview(config.identity,view.expected,remaining.map(item=>approveItem(view.plan,view.segments,item.id,config.identity,item.count===0)),[]);view=service.load(); + expect(view.items.every(item=>item.state==='approved')).toBe(true); +},30000); +it('marks unchanged items stale after a same-snapshot plan amendment',()=>{ + const {service,config}=fixture();let view=service.load(); + service.store.saveReview(config.identity,view.expected,view.items.map(item=>approveItem(view.plan,view.segments,item.id,config.identity,item.count===0)),[]);view=service.load(); + const attempt=service.store.beginMergeAttempt(config.identity,{...view.expected,reviewVersion:view.expected.reviewVersion!},view.snapshot.head);service.store.queueMergeAttempt(config.identity,attempt.id,'https://github.example/pr/24');service.store.finishMergeAttempt(config.identity,attempt.id,{state:'removed',reason:'Checks failed.'}); + const amended={...view.plan,summary:'Amended review requirements'};service.store.importRevision(JSON.stringify(amended),'json',{identity:config.identity,issue:amended.issue,baseEntries:['retry.ts','README.md','run.sh'].map(path=>({path,kind:'file' as const})),pathKey:path=>path,allowedCommands:[]},amended.revision); + view=service.load();expect(view.items.every(item=>item.state==='stale'&&item.reasons.includes('Plan revision changed after the queue attempt'))).toBe(true); + const first=view.items[0]!;view=service.act({action:'approve',item:first.id,confirmNoChange:first.count===0,token:view.token});expect(view.items.find(item=>item.id===first.id)?.state).toBe('approved'); +}); it('refuses no-change confirmation while the item still owns ambiguous changes',async()=>{ const {service,config}=fixture();const {writeFileSync}=await import('node:fs');const {execFileSync}=await import('node:child_process'); writeFileSync(join(config.repository,'retry.ts'),'export function delay(attempt: number) {\n return Math.min(10000, 200 * 2 ** attempt);\n}\n'); diff --git a/test/store.test.ts b/test/store.test.ts index 5f76ec9..14b39cd 100644 --- a/test/store.test.ts +++ b/test/store.test.ts @@ -74,11 +74,37 @@ it('retains cancellation reasons across restart and never revives cancelled work expect(recovered.getSuggestions(identity, id)).toMatchObject({ state: 'cancelled', reason: 'Runner shut down.', snapshotId: expected.snapshotId }); expect(() => recovered.completeSuggestions(identity, id, reply())).toThrow(/stale|cancelled|complete/); }); +it('persists the merge-queue lifecycle and terminal reason across restart', () => { + const { store, path } = fixture(); + const expected = { ...state(store), reviewVersion: store.reviewVersion(identity) }; + const attempt = store.beginMergeAttempt(identity, expected, oid(2), 'MQEV_before'); + expect(store.recordMergeAttemptDiagnostic(identity, attempt.id, 'Submission timed out after GitHub may have accepted it.')).toBe(true); + expect(store.getMergeAttempt(identity)).toMatchObject({ state: 'submitting', reason: 'Submission timed out after GitHub may have accepted it.' }); + expect(store.queueMergeAttempt(identity, attempt.id, 'https://github.example/pr/1')).toBe(true); + expect(store.getMergeAttempt(identity)).toMatchObject({ state: 'queued', reason: null }); + expect(store.observeQueuedMerge(identity, attempt.id, { entryId: 'MQE_1', phase: 'AWAITING_CHECKS', position: 2 })).toBe(true); + expect(store.finishMergeAttempt(identity, attempt.id, { state: 'removed', reason: 'Checks failed.', occurredAt: '2026-09-24T08:05:00Z' })).toBe(true); + close(store); + expect(open(path).getMergeAttempt(identity)).toMatchObject({ + id: attempt.id, state: 'removed', reviewedHead: oid(2), queueWatermark: 'MQEV_before', reason: 'Checks failed.', entryId: 'MQE_1', phase: 'AWAITING_CHECKS', position: 2, + }); +}); +it('prevents a stale queue observation from overwriting a retry attempt', () => { + const { store } = fixture(); + const expected = { ...state(store), reviewVersion: store.reviewVersion(identity) }; + const first = store.beginMergeAttempt(identity, expected, oid(2)); + store.queueMergeAttempt(identity, first.id, 'https://github.example/pr/1'); + store.finishMergeAttempt(identity, first.id, { state: 'failed', reason: 'Queue failed.' }); + const retry = store.beginMergeAttempt(identity, expected, oid(2)); + store.queueMergeAttempt(identity, retry.id, 'https://github.example/pr/1'); + expect(store.finishMergeAttempt(identity, first.id, { state: 'merged', occurredAt: '2026-09-24T08:10:00Z' })).toBe(false); + expect(store.getMergeAttempt(identity)).toMatchObject({ id: retry.id, state: 'queued', reviewedHead: oid(2) }); +}); it('migrates unbound active requests to terminal history instead of reviving them', () => { const { store, path } = fixture(); const id = ready(store); close(store); const legacy = new DatabaseSync(path); - legacy.exec('ALTER TABLE requests DROP COLUMN reason; ALTER TABLE requests DROP COLUMN snapshot_id; PRAGMA user_version=3;'); + legacy.exec('DROP TABLE merge_attempts; ALTER TABLE requests DROP COLUMN reason; ALTER TABLE requests DROP COLUMN snapshot_id; PRAGMA user_version=3;'); legacy.close(); const recovered = open(path); expect(recovered.getSuggestions(identity, id)).toEqual({ state: 'invalidated', revision: 1, snapshotId: null, reply: reply(), reason: 'Request predates snapshot binding.' }); diff --git a/web/public/app.js b/web/public/app.js index c555030..51e6512 100644 --- a/web/public/app.js +++ b/web/public/app.js @@ -20,6 +20,11 @@ let data, since = false, busy = false; let reviewGeneration = 0; +let mergeGeneration = 0, + mergePollTimer = null, + mergePollState = null, + mergePollDelay = 2000; +const mergePollMaximumDelay = 30000; const drafts = new Map(); const attachments = new Map(); let snippetSelection = null; @@ -48,6 +53,11 @@ function rememberDraft() { if (selected) drafts.set(`${selected}:${mode}`, $("message").value); } function showFailure(message) { + mergeGeneration++; + if (mergePollTimer) clearTimeout(mergePollTimer); + mergePollTimer = null; + mergePollState = null; + mergePollDelay = 2000; data = null; snippetSelection = null; $("selection-actions").hidden = true; @@ -79,6 +89,9 @@ async function refresh() { if (busy) return; busy = true; reviewGeneration++; + mergeGeneration++; + mergePollState = null; + mergePollDelay = 2000; renderAttachment(); $("banner").textContent = "Linking changes to plan items…"; try { @@ -100,10 +113,79 @@ async function refresh() { renderAttachment(); } } +function renderMerge() { + const merge = data?.merge; + $("merge").hidden = !merge?.available; + $("merge-details").hidden = !merge?.available || (!merge.blockers.length && !merge.queue); + if (!merge?.available) return; + $("merge").disabled = !merge.ready; + $("merge").textContent = merge.action === "retry" + ? "Retry merge" + : merge.queue?.state === "submitting" + ? "Submitting…" + : merge.queue?.state === "queued" + ? "Merge queued" + : merge.queue?.state === "merged" + ? "Merged" + : merge.ready + ? "Merge PR" + : `${merge.blockers.length} blocker${merge.blockers.length === 1 ? "" : "s"}`; + $("merge-details").textContent = merge.queue ? "Merge status" : "Review blockers"; + scheduleMergePoll(); +} +function scheduleMergePoll() { + if (mergePollTimer) clearTimeout(mergePollTimer); + mergePollTimer = null; + const queue = data?.merge?.queue; + const state = queue?.state; + if (queue?.kind === "queue" && ["submitting", "queued"].includes(state)) { + if (state !== mergePollState) { + mergePollState = state; + mergePollDelay = 2000; + } + const generation = mergeGeneration; + mergePollTimer = setTimeout(() => pollMergeQueue(generation), mergePollDelay); + } else { + mergePollState = null; + mergePollDelay = 2000; + } +} +function backOffMergePoll() { + mergePollDelay = Math.min(mergePollDelay * 2, mergePollMaximumDelay); +} +async function pollMergeQueue(generation) { + try { + const update = await api("/api/merge"); + if (generation !== mergeGeneration || !data?.merge?.available) return; + data = { ...data, merge: { ...data.merge, queue: update.queue } }; + const queue = update.queue; + if (queue?.state === "merged") { + data.merge = { ...data.merge, ready: false, action: null, blockers: [{ code: "queue-merged", message: "GitHub confirmed that the reviewed head was merged." }] }; + $("banner").textContent = "GitHub confirmed the reviewed head was merged."; + } else if (queue?.state === "removed" || queue?.state === "failed") { + data.merge = { ...data.merge, ready: false, action: null, blockers: [{ code: "queue-refresh", message: `${queue.reason} Refresh to verify retry readiness.` }] }; + $("banner").textContent = `${queue.state === "removed" ? "Removed from merge queue" : "Merge queue failed"}. ${queue.reason} Refresh to verify retry readiness.`; + } else if (queue?.observationError) { + $("banner").textContent = `${queue.state === "submitting" ? "Merge submission status is unknown" : "Merge remains queued"}. ${queue.observationError}`; + } else if (queue?.state === "queued") { + $("banner").textContent = `Merge queued${queue.position === null ? "" : ` at position ${queue.position}`}. Waiting for GitHub.`; + } + if (queue?.state === mergePollState) backOffMergePoll(); + renderMerge(); + } catch (error) { + if (generation !== mergeGeneration || !data?.merge?.available) return; + $("banner").textContent = `Could not refresh merge-queue status. ${error.message}`; + backOffMergePoll(); + scheduleMergePoll(); + } +} async function act(command) { if (busy || !data) return false; busy = true; reviewGeneration++; + mergeGeneration++; + mergePollState = null; + mergePollDelay = 2000; renderAttachment(); try { rememberDraft(); @@ -145,16 +227,7 @@ function render() { `#${data.plan.issue} ${data.plan.summary} · r${data.plan.revision}`; $("progress").textContent = `${data.approved} of ${data.items.length} approved`; - const merge = data.merge; - $("merge").hidden = !merge?.available; - $("merge-details").hidden = !merge?.available || merge.ready; - if (merge?.available) { - $("merge").disabled = !merge.ready; - $("merge").textContent = merge.ready - ? "Merge PR" - : `${merge.blockers.length} blocker${merge.blockers.length === 1 ? "" : "s"}`; - $("merge-details").textContent = "Review blockers"; - } + renderMerge(); $("banner").textContent = data.demo ? "Demo repository · real Git changes and local SQLite storage. No tests or AI review have been run for this demo." : ""; @@ -379,22 +452,28 @@ $("approve").onclick = () => { $("reload").onclick = refresh; $("merge-details").onclick = () => { if (!data?.merge?.available) return; + const queue = data.merge.queue; showDialog( - `

Merge blockers

    ${data.merge.blockers.map((blocker) => `
  • ${esc(blocker.message)}
  • `).join("")}
`, + `

${queue ? "Merge status" : "Merge blockers"}

${queue ? `

${esc(queue.state)} · reviewed head ${esc(queue.reviewedHead.slice(0, 12))}

${queue.phase ? `

GitHub phase: ${esc(queue.phase)}${queue.position === null ? "" : ` · position ${queue.position}`}

` : ""}${queue.reason ? `

${esc(queue.reason)}

` : ""}${queue.observationError ? `

${esc(queue.observationError)}

` : ""}` : ""}
    ${data.merge.blockers.map((blocker) => `
  • ${esc(blocker.message)}
  • `).join("")}
`, ); }; $("merge").onclick = async () => { - if (busy || !data?.merge?.ready || !window.confirm("Merge this reviewed pull request?")) return; + if (busy || !data?.merge?.ready || !window.confirm(data.merge.action === "retry" ? "Retry merging this exact reviewed head?" : "Merge this reviewed pull request?")) return; busy = true; + mergeGeneration++; + mergePollState = null; + mergePollDelay = 2000; $("merge").disabled = true; try { + rememberDraft(); const updated = await api("/api/action", { action: "merge", token: data.token }); - const blocker = { code: "merge-submitted", message: "Merge was submitted. Refresh to confirm GitHub state." }; + const queued = updated.mergeQueue?.state === "queued" || updated.mergeQueue?.state === "submitting"; + const blocker = { code: queued ? "queue-active" : "merge-submitted", message: queued ? "The reviewed head is queued. Waiting for GitHub to confirm the outcome." : "Merge was submitted. Refresh to confirm GitHub state." }; data = updated.mergeRefreshRequired - ? { ...data, merge: { ...data.merge, ready: false, blockers: [blocker] } } - : { ...updated, merge: { ...updated.merge, ready: false, blockers: [blocker] } }; + ? { ...data, merge: { ...data.merge, ready: false, action: null, queue: updated.mergeQueue ?? null, blockers: [blocker] } } + : { ...updated, merge: { ...updated.merge, ready: false, action: null, blockers: [blocker] } }; render(); - $("banner").textContent = `Merge submitted. ${updated.mergeResult.url}`; + $("banner").textContent = `${queued ? "Merge queued" : "Merge submitted"}. ${updated.mergeResult.url}`; } catch (error) { data = { ...data, merge: { ...data.merge, ready: false, blockers: [{ code: "stale-merge", message: `${error.message} Refresh before trying again.` }] } }; render(); diff --git a/web/server.ts b/web/server.ts index 7c3807b..348d53b 100644 --- a/web/server.ts +++ b/web/server.ts @@ -1,4 +1,4 @@ -import { createServer } from 'node:http'; +import { createServer, type IncomingMessage } from 'node:http'; import { readFileSync } from 'node:fs'; import { fileURLToPath } from 'node:url'; import { randomBytes, timingSafeEqual } from 'node:crypto'; @@ -7,7 +7,8 @@ import { Questions, type QuestionAgent } from '../runner/questions.ts'; import { GhMergeGateway, type MergeGateway } from '../github/merge.ts'; import { MergeCoordinator } from '../runner/merge.ts'; const publicRoot = new URL('./public/', import.meta.url); -export async function startServer(config: ReviewConfig, port = 4318, questionAgent?: QuestionAgent, mergeGateway?: MergeGateway) { +export async function startServer(config: ReviewConfig, port = 4318, questionAgent?: QuestionAgent, mergeGateway?: MergeGateway, shutdownDrainMs = 14_500) { + if (!Number.isSafeInteger(shutdownDrainMs) || shutdownDrainMs < 1 || shutdownDrainMs > 14_500) throw new Error('Invalid shutdown drain deadline.'); const service = new ReviewService(config), token = randomBytes(32).toString('hex'); let questions: Questions, merges: MergeCoordinator | null; try { @@ -15,12 +16,15 @@ export async function startServer(config: ReviewConfig, port = 4318, questionAge questions=new Questions(service,questionAgent); merges = !config.demo && (mergeGateway || config.github) ? new MergeCoordinator(service, mergeGateway ?? new GhMergeGateway(config.github!)) : null; } catch (error) { service.close(); throw error; } - const load=async()=>{const view=service.load();return {...view,notes:view.notes.map(note=>({...note,answerActive:questions.isRunning(note.id)})),merge:merges?await merges.displayStatus(view):{available:false}};}; + const loadReview=()=>{const view=service.load();return {...view,notes:view.notes.map(note=>({...note,answerActive:questions.isRunning(note.id)}))};}; + const load=async(signal?:AbortSignal)=>{const view=loadReview();return {...view,merge:merges?await merges.displayStatus(view,signal):{available:false}};}; const answerStatuses=()=>service.store.getReviewNotes(config.identity) .filter(note=>note.kind==='question') .map(note=>({id:note.id,answer:note.answer,answerActive:questions.isRunning(note.id)})); let stopping = false; + const activeRequests=new Set<{abort:AbortController;request:IncomingMessage;readingBody:boolean}>(); const server = createServer(async (req, res) => { + const requestAbort=new AbortController(),activeRequest={abort:requestAbort,request:req,readingBody:false};activeRequests.add(activeRequest); const address = server.address(); const actualPort = address && typeof address !== 'string' ? address.port : port; const origin = `http://127.0.0.1:${actualPort}`; res.setHeader('Cache-Control', 'no-store'); res.setHeader('X-Content-Type-Options', 'nosniff'); @@ -32,26 +36,29 @@ export async function startServer(config: ReviewConfig, port = 4318, questionAge if (path.startsWith('/api/')) { const supplied = req.headers['x-codeboost-token']; if (typeof supplied !== 'string' || !/^[a-f0-9]{64}$/.test(supplied) || !timingSafeEqual(Buffer.from(supplied), Buffer.from(token))) { json(403, { error: 'Open the private local URL printed by the CLI.' }); return; } - if (stopping && req.method === 'POST') { json(503, { error: 'The review server is shutting down.' }); return; } + if (stopping) { json(503, { error: 'The review server is shutting down.' }); return; } if (req.method === 'GET' && path === '/api/settings') { json(200,{questionProvider:service.store.questionProvider()});return; } if (req.method === 'GET' && path === '/api/questions') { json(200,{notes:answerStatuses()});return; } - if (req.method === 'GET' && path === '/api/review') { json(200, await load()); return; } + if (req.method === 'GET' && path === '/api/merge') { if(!merges)throw new Error('Merging is not configured for this review.');json(200,{queue:await merges.pollQueue()});return; } + if (req.method === 'GET' && path === '/api/review') { json(200, await load(requestAbort.signal)); return; } if (req.method !== 'POST' || !['/api/action','/api/settings'].includes(path) || req.headers['content-type'] !== 'application/json') { json(405, { error: 'Unsupported request.' }); return; } const chunks: Buffer[] = []; let size = 0; - for await (const chunk of req) { size += chunk.length; if (size > 16384) { json(413, { error: 'Request too large.' }); return; } chunks.push(chunk); } + activeRequest.readingBody=true; + try { for await (const chunk of req) { size += chunk.length; if (size > 16384) { json(413, { error: 'Request too large.' }); return; } chunks.push(chunk); } } + finally { activeRequest.readingBody=false; } const body = new TextDecoder('utf-8', { fatal: true }).decode(Buffer.concat(chunks)); const input=JSON.parse(body); if (stopping && input.action === 'merge') { json(503, { error: 'The review server is shutting down.' }); return; } if(path==='/api/settings') {service.store.setQuestionProvider(input.questionProvider);json(200,{questionProvider:service.store.questionProvider()});return;} if(input.action==='retry-question') { const view=service.load();if(input.token!==view.token)throw new Error('Stale review state. Refresh and retry.'); - questions.start(input.id,view);json(200,await load());return; + questions.start(input.id,view);json(200,await load(requestAbort.signal));return; } if(input.action==='merge') { if(!merges)throw new Error('Merging is not configured for this review.'); const merged=await merges.merge(input.token); - try { json(200,{...(await load()),mergeResult:merged.result,mergeRefreshRequired:false}); } - catch { json(200,{mergeResult:merged.result,mergeRefreshRequired:true}); } + try { const mergeQueue=merges.queueSnapshot();json(200,{...loadReview(),merge:{...merged.status,queue:mergeQueue},mergeResult:merged.result,mergeQueue,mergeRefreshRequired:false}); } + catch { json(200,{mergeResult:merged.result,mergeQueue:null,mergeRefreshRequired:true}); } return; } const view=service.act(input); @@ -60,7 +67,7 @@ export async function startServer(config: ReviewConfig, port = 4318, questionAge // The saved question remains visible and retryable when capacity is reached. } } - json(200,await load());return; + json(200,await load(requestAbort.signal));return; } if (req.method !== 'GET') { json(405, { error: 'Method not allowed.' }); return; } const files: Record = { '/': ['index.html', 'text/html'], '/app.js': ['app.js', 'text/javascript'], '/style.css': ['style.css', 'text/css'] }; @@ -73,6 +80,7 @@ export async function startServer(config: ReviewConfig, port = 4318, questionAge } json(404, { error: 'Not found.' }); } catch (error) { json(409, { error: error instanceof Error ? error.message : 'Review failed.' }); } + finally {activeRequests.delete(activeRequest);} }); server.requestTimeout = 15000; await new Promise((resolve, reject) => { server.once('error', reject); server.listen(port, '127.0.0.1', () => { server.removeListener('error', reject); resolve(); }); }).catch(error => { service.close(); throw error; }); @@ -80,6 +88,14 @@ export async function startServer(config: ReviewConfig, port = 4318, questionAge return { server, service, token, url: `http://127.0.0.1:${address.port}/#${token}`, close: async () => { stopping = true; const closing = new Promise((resolve, reject) => server.close(error => error ? reject(error) : resolve())); + let timer: ReturnType | undefined; + await Promise.race([closing, new Promise(resolve => { timer=setTimeout(resolve,shutdownDrainMs); })]); + if(timer)clearTimeout(timer); + for(const active of activeRequests) { + const reason=new Error('Request cancelled during shutdown.'); + active.abort.abort(reason); + if(active.readingBody)active.request.destroy(reason); + } await merges?.close(); await closing; await questions.close();