diff --git a/package-lock.json b/package-lock.json index c53677b..1747447 100644 --- a/package-lock.json +++ b/package-lock.json @@ -17,7 +17,7 @@ "@agentworkforce/delivery": "^4.1.23", "@agentworkforce/runtime": "^4.1.23", "@relayfile/relay-helpers": "^0.4.6", - "@relayfile/sdk": "^0.10.32", + "@relayfile/sdk": "0.10.34", "@relayflows/core": "^1.0.3", "agent-relay": "^10.6.4", "proper-lockfile": "^4.1.2", @@ -2735,9 +2735,9 @@ } }, "node_modules/@relayfile/core": { - "version": "0.10.32", - "resolved": "https://registry.npmjs.org/@relayfile/core/-/core-0.10.32.tgz", - "integrity": "sha512-iCOvD3Zfkv4zXj/Fw6adwQfsBFaXrJNqgK8WxlvxjFZ1lItMsJKV94qHKIytzNSEpXnl9ZI1cNfE2Jiiqxs51A==", + "version": "0.10.34", + "resolved": "https://registry.npmjs.org/@relayfile/core/-/core-0.10.34.tgz", + "integrity": "sha512-aqeRhruuacoiXQ7khMkWLyFJaq5SEa0w457rg39dBXbpJZ3uIYRHb5m75AteLgvMAaffu7SwVBcuoV1AoMM8Nw==", "license": "Apache-2.0", "engines": { "node": ">=18" @@ -2757,9 +2757,9 @@ } }, "node_modules/@relayfile/mount-darwin-arm64": { - "version": "0.10.32", - "resolved": "https://registry.npmjs.org/@relayfile/mount-darwin-arm64/-/mount-darwin-arm64-0.10.32.tgz", - "integrity": "sha512-I2TOeymcAw/U1oL1APQ15QkMmaRxJWCrWv0GkMqgcDWPgpmhd3Y/grg7PYrORLRZAVQC5fcQwdxZMn3s+nZJOA==", + "version": "0.10.34", + "resolved": "https://registry.npmjs.org/@relayfile/mount-darwin-arm64/-/mount-darwin-arm64-0.10.34.tgz", + "integrity": "sha512-SqJcKMumdz+euMOgbCGIfNo243MnPYcPtOEQb6UKeUuE5e18BQ08Lqy6VMe44tuYnoCE7gxwyQu5Ae03eOwAjw==", "cpu": [ "arm64" ], @@ -2770,9 +2770,9 @@ ] }, "node_modules/@relayfile/mount-darwin-x64": { - "version": "0.10.32", - "resolved": "https://registry.npmjs.org/@relayfile/mount-darwin-x64/-/mount-darwin-x64-0.10.32.tgz", - "integrity": "sha512-ek2yR23sRmEfDMw1DjQ18NvStBIX9C6QaXzCX//PGv3ISoeaHTC8yEG2VNBqYYEAAGG4Om1yY/POH81lG3bG7g==", + "version": "0.10.34", + "resolved": "https://registry.npmjs.org/@relayfile/mount-darwin-x64/-/mount-darwin-x64-0.10.34.tgz", + "integrity": "sha512-j5A799CrsEeM1kLglVU5aZ0VMAOgzwhBn5x3srwhJe3Dt9PQZ0Nf6iU2J2vgGbp209Nrpsj1PNuYiDtX1tKAdA==", "cpu": [ "x64" ], @@ -2783,9 +2783,9 @@ ] }, "node_modules/@relayfile/mount-linux-arm64": { - "version": "0.10.32", - "resolved": "https://registry.npmjs.org/@relayfile/mount-linux-arm64/-/mount-linux-arm64-0.10.32.tgz", - "integrity": "sha512-m0cpVFQZTYEKDkgLvss7InY+ZzD767ApC7Q+r0aVmC+1FDG9xOqPQ877pX6v9oJja/IpLikQACy55aR2V3AdJg==", + "version": "0.10.34", + "resolved": "https://registry.npmjs.org/@relayfile/mount-linux-arm64/-/mount-linux-arm64-0.10.34.tgz", + "integrity": "sha512-haxDMza884Y+aQ32URcsCYQDARy+x9hex4DVfTUBTcMsTfzhjLCqrmcj1ZSluM947M2XLjvvc8YQlkBdSjAOXw==", "cpu": [ "arm64" ], @@ -2796,9 +2796,9 @@ ] }, "node_modules/@relayfile/mount-linux-x64": { - "version": "0.10.32", - "resolved": "https://registry.npmjs.org/@relayfile/mount-linux-x64/-/mount-linux-x64-0.10.32.tgz", - "integrity": "sha512-FwiPOYrNaDoOrK5idfPWguQwweV4bHKtzyhRB5qRoBrmbMRBjsMJiezO0swoWIXu1sBbGGzzj/6k0yt8Yht+uA==", + "version": "0.10.34", + "resolved": "https://registry.npmjs.org/@relayfile/mount-linux-x64/-/mount-linux-x64-0.10.34.tgz", + "integrity": "sha512-pi/lUzUjs7eOGEJrJHhy1yy7NXeeHrf2GAwf2HQTklqPZ1jD1mjDPjFO59jryXoqM3YLVOsd3LC6HuS5bvPgiQ==", "cpu": [ "x64" ], @@ -2835,12 +2835,12 @@ } }, "node_modules/@relayfile/sdk": { - "version": "0.10.32", - "resolved": "https://registry.npmjs.org/@relayfile/sdk/-/sdk-0.10.32.tgz", - "integrity": "sha512-+Djgg7+8Oz/iuoTl05VjuoFi3S/LeFaHiHxel5b2QAOkeSxkJgNwy7ujnhF6yfHqrltc5sAmNUP3zKQqVUIuSg==", + "version": "0.10.34", + "resolved": "https://registry.npmjs.org/@relayfile/sdk/-/sdk-0.10.34.tgz", + "integrity": "sha512-VU52piPz+12Ooogi5dxpf36rUPAdAkIYmpHMEuNwQcvq6S7fXxCmOw4B71k81C+KuwYxZQTIftpxIgZdvLzJQw==", "license": "Apache-2.0", "dependencies": { - "@relayfile/core": "0.10.32", + "@relayfile/core": "0.10.34", "ignore": "^7.0.5", "tar": "^7.5.10" }, @@ -2848,10 +2848,10 @@ "node": ">=18" }, "optionalDependencies": { - "@relayfile/mount-darwin-arm64": "0.10.32", - "@relayfile/mount-darwin-x64": "0.10.32", - "@relayfile/mount-linux-arm64": "0.10.32", - "@relayfile/mount-linux-x64": "0.10.32" + "@relayfile/mount-darwin-arm64": "0.10.34", + "@relayfile/mount-darwin-x64": "0.10.34", + "@relayfile/mount-linux-arm64": "0.10.34", + "@relayfile/mount-linux-x64": "0.10.34" } }, "node_modules/@relayflows/browser-primitive": { diff --git a/package.json b/package.json index 01b9f40..2a4a184 100644 --- a/package.json +++ b/package.json @@ -84,7 +84,7 @@ "@agent-relay/integration-prompts": "^10.6.4", "@agent-relay/sdk": "^10.6.4", "@relayfile/relay-helpers": "^0.4.6", - "@relayfile/sdk": "^0.10.32", + "@relayfile/sdk": "0.10.34", "@relayflows/core": "^1.0.3", "agent-relay": "^10.6.4", "proper-lockfile": "^4.1.2", diff --git a/src/index.ts b/src/index.ts index 71064f1..8661da8 100644 --- a/src/index.ts +++ b/src/index.ts @@ -219,6 +219,7 @@ export type { } from './safety/factory-scope' export { canonicalMountPaths, + createResourceSubscriptionsSdkClient, createWorkspaceScopedEventClient, deliveryTargetsFor, eventPathGlobsForIntegration, @@ -236,6 +237,8 @@ export { linearScopePredicates, normalizeChangePath, relayfileSdkPathFiltersFor, + ResourceSubscriptionsUnavailableError, + isResourceSubscriptionsUnavailable, parseSlackThreadReply, slackThreadReplyGlob, slackListenDms, @@ -262,6 +265,13 @@ export type { WorkspaceScopedEventClientOptions, WorkspaceScopedSubscribeOptions, ChangeEvent as SubscriptionChangeEvent, + AcceptedResourceDelivery, + ResourceDeliveryClaim, + ResourceSubscription, + ResourceSubscriptionInput, + ResourceSubscriptionsClient, + ResourceSubscriptionsSdk, + ResourceSubscriptionsSdkClientOptions, } from './subscriptions' export type { Capability, diff --git a/src/mount/relayfile-cloud-mount-client.test.ts b/src/mount/relayfile-cloud-mount-client.test.ts index 7b17338..1bc59b1 100644 --- a/src/mount/relayfile-cloud-mount-client.test.ts +++ b/src/mount/relayfile-cloud-mount-client.test.ts @@ -1,6 +1,12 @@ import { describe, expect, it, vi } from 'vitest' import { CloudAuthError, type CloudSession, type StoredAuth } from '@agent-relay/cloud' -import type { ChangeEvent, OperationStatusResponse } from '@relayfile/sdk' +import type { + AcceptDurableSubscriptionDeliveryInput, + ChangeEvent, + ClaimDurableSubscriptionDeliveriesInput, + CreateOrRenewDurableResourceSubscriptionInput, + OperationStatusResponse, +} from '@relayfile/sdk' import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' @@ -56,6 +62,14 @@ class FakeRelayFileClient implements RelayFileClientLike { readonly getEventsCalls: Array<{ workspaceId: string; opts?: { cursor?: string; limit?: number; provider?: string; last?: number } }> = [] readonly listLastNChangesCalls: Array<{ limit: number; context?: { workspaceId: string } }> = [] readonly getOpCalls: Array<{ workspaceId: string; opId: string }> = [] + readonly createSubscriptionCalls: CreateOrRenewDurableResourceSubscriptionInput[] = [] + readonly claimDeliveryCalls: ClaimDurableSubscriptionDeliveriesInput[] = [] + readonly acceptDeliveryCalls: AcceptDurableSubscriptionDeliveryInput[] = [] + readonly cancelSubscriptionCalls: Array<{ + workspaceId: string + subscriptionId: string + options?: { signal?: AbortSignal } + }> = [] getSyncStatus?: RelayFileClientLike['getSyncStatus'] treePageSize?: number @@ -181,6 +195,80 @@ class FakeRelayFileClient implements RelayFileClientLike { } } + async createOrRenewDurableResourceSubscription(input: CreateOrRenewDurableResourceSubscriptionInput) { + this.createSubscriptionCalls.push(input) + return { + id: 'sub-1', + ownerId: 'configured-factory-agent', + subscriberId: input.subscriberId, + provider: input.provider, + resourceRef: input.resourceRef, + eventTypes: input.eventTypes, + terminalEventTypes: input.terminalEventTypes ?? [], + intent: input.intent ?? null, + status: 'active' as const, + createdAt: '2026-07-21T00:00:00.000Z', + updatedAt: '2026-07-21T00:00:00.000Z', + expiresAt: '2026-12-31T00:00:00.000Z', + retiredAt: null, + } + } + + async claimDurableSubscriptionDeliveries(input: ClaimDurableSubscriptionDeliveriesInput) { + this.claimDeliveryCalls.push(input) + return { + deliveries: [{ + id: 'delivery-1', + claimToken: 'claim-token-1', + subscriptionId: 'sub-1', + ownerId: 'configured-factory-agent', + subscriberId: 'factory-babysitter:uuid-1', + provider: 'github', + resourceRef: '/github/repos/AgentWorkforce__pear/pulls/by-id/1.json', + event: { + id: 'event-1', + type: 'pull_request.closed', + path: '/github/repos/AgentWorkforce__pear/pulls/by-id/1.json', + revision: '2', + origin: 'github', + provider: 'github', + correlationId: 'corr-1', + timestamp: '2026-07-21T00:00:00.000Z', + }, + terminal: true, + status: 'claimed' as const, + createdAt: '2026-07-21T00:00:00.000Z', + claimedAt: '2026-07-21T00:00:01.000Z', + claimLeaseExpiresAt: '2026-07-21T00:01:01.000Z', + acceptedAt: null, + }], + } + } + + async acceptDurableSubscriptionDelivery(input: AcceptDurableSubscriptionDeliveryInput) { + this.acceptDeliveryCalls.push(input) + if (input.claimToken !== 'claim-token-1') { + throw Object.assign(new Error('delivery claim mismatch'), { status: 409 }) + } + const claimed = (await this.claimDurableSubscriptionDeliveries({ workspaceId: input.workspaceId })).deliveries[0]! + return { + delivery: { + ...claimed, + claimToken: null, + status: 'accepted' as const, + acceptedAt: '2026-07-21T00:00:02.000Z', + }, + } + } + + async cancelDurableResourceSubscription( + workspaceId: string, + subscriptionId: string, + options?: { signal?: AbortSignal }, + ) { + this.cancelSubscriptionCalls.push({ workspaceId, subscriptionId, options }) + } + async getToken() { return 'relayfile-token' } @@ -680,6 +768,80 @@ describe('RelayfileCloudMountClient', () => { expect(cloudSessionProvider).toHaveBeenCalledTimes(2) }) + it('adapts durable resource subscriptions through the canonical Relayfile SDK methods', async () => { + const fake = new FakeRelayFileClient() + const mount = new RelayfileCloudMountClient({ workspaceId: 'rw_test', client: fake }) + + const client = mount.resourceSubscriptions! + await expect(client.createOrRenew('rw_test', { + provider: 'github', + resourceRef: '/github/repos/AgentWorkforce__pear/pulls/by-id/1.json', + eventTypes: ['pull_request_review_comment.created'], + terminalEventTypes: ['pull_request.closed'], + subscriberId: 'factory-babysitter:uuid-1', + ttlSeconds: 3600, + })).resolves.toMatchObject({ subscriptionId: 'sub-1', ownerId: 'configured-factory-agent' }) + await expect(client.claimDeliveryClaims('rw_test')).resolves.toEqual([expect.objectContaining({ + deliveryId: 'delivery-1', + claimToken: 'claim-token-1', + terminal: true, + })]) + await expect(client.acceptDelivery('rw_test', { deliveryId: 'delivery-1', claimToken: 'wrong-token' })) + .rejects.toMatchObject({ status: 409 }) + await expect(client.acceptDelivery('rw_test', { deliveryId: 'delivery-1', claimToken: 'claim-token-1' })) + .resolves.toEqual({ deliveryId: 'delivery-1', subscriptionId: 'sub-1', terminal: true }) + await client.cancel('rw_test', { subscriptionId: 'sub-1' }) + + expect(fake.createSubscriptionCalls).toEqual([expect.objectContaining({ + workspaceId: 'rw_test', + provider: 'github', + resourceRef: '/github/repos/AgentWorkforce__pear/pulls/by-id/1.json', + eventTypes: ['pull_request_review_comment.created'], + terminalEventTypes: ['pull_request.closed'], + subscriberId: 'factory-babysitter:uuid-1', + ttlSeconds: 3600, + })]) + expect(fake.claimDeliveryCalls[0]).toMatchObject({ workspaceId: 'rw_test' }) + expect(fake.acceptDeliveryCalls).toEqual([ + expect.objectContaining({ workspaceId: 'rw_test', deliveryId: 'delivery-1', claimToken: 'wrong-token' }), + expect.objectContaining({ workspaceId: 'rw_test', deliveryId: 'delivery-1', claimToken: 'claim-token-1' }), + ]) + expect(fake.cancelSubscriptionCalls).toEqual([ + expect.objectContaining({ workspaceId: 'rw_test', subscriptionId: 'sub-1' }), + ]) + }) + + it('fails closed for non-claimed SDK deliveries and forwards lifecycle cancellation', async () => { + const fake = new FakeRelayFileClient() + fake.claimDurableSubscriptionDeliveries = vi.fn(async (input) => ({ + deliveries: [{ + ...(await new FakeRelayFileClient().claimDurableSubscriptionDeliveries(input)).deliveries[0]!, + claimToken: null, + status: 'pending' as const, + }], + })) + const malformed = new RelayfileCloudMountClient({ + workspaceId: 'rw_test', + client: fake, + }) + await expect(malformed.resourceSubscriptions!.claimDeliveryClaims('rw_test')) + .rejects.toThrow(/without a live claim/u) + + const controller = new AbortController() + const cancelledSdk = new FakeRelayFileClient() + const claim = vi.spyOn(cancelledSdk, 'claimDurableSubscriptionDeliveries') + const cancelled = new RelayfileCloudMountClient({ + workspaceId: 'rw_test', + client: cancelledSdk, + resourceSubscriptionSignal: controller.signal, + }) + await cancelled.resourceSubscriptions!.claimDeliveryClaims('rw_test') + expect(claim).toHaveBeenCalledWith(expect.objectContaining({ + workspaceId: 'rw_test', + signal: controller.signal, + })) + }) + it('coalesces concurrent shared session resolutions for relayfile token refresh', async () => { const setup = { joinWorkspace: vi.fn(async () => ({ diff --git a/src/mount/relayfile-cloud-mount-client.ts b/src/mount/relayfile-cloud-mount-client.ts index 0174370..b9f3841 100644 --- a/src/mount/relayfile-cloud-mount-client.ts +++ b/src/mount/relayfile-cloud-mount-client.ts @@ -36,11 +36,14 @@ import type { SubscribeOptions, } from '../ports' import { + createResourceSubscriptionsSdkClient, createWorkspaceScopedEventClient, type RelayfileEventClient, + type ResourceSubscriptionsSdk, type TokenProvider, type WorkspaceEventClientSource, } from '../subscriptions' +import type { ResourceSubscriptionsClient } from '../subscriptions' import { RelayfileGithubConnectionWrite } from './relayfile-github-connection-write' import { ensureLocalMount as runLocalMountPreflight, @@ -186,6 +189,10 @@ export interface RelayfileCloudMountClientConfig { tokenProvider?: TokenProvider baseUrl?: string eventClient?: RelayfileEventClient + /** Override the standard Relayfile SDK durable-subscription adapter. */ + resourceSubscriptions?: ResourceSubscriptionsClient + /** Optional lifecycle cancellation forwarded through the Relayfile SDK. */ + resourceSubscriptionSignal?: AbortSignal logger?: Logger onLocalMountHealth?: (event: LocalMountHealthEvent) => Promise | void /** Internal health cadence override for tests. */ @@ -196,8 +203,7 @@ export interface RelayfileCloudMountClientConfig { isAllowedDelete?: (path: string, currentContent: unknown) => boolean | Promise } -export type RelayFileClientLike = - { +export type RelayFileClientLike = { readFile(workspaceId: string, path: string): Promise writeFile(input: WriteFileInput): Promise deleteFile(input: DeleteFileInput): Promise @@ -209,7 +215,15 @@ export type RelayFileClientLike = getSyncStatus?(workspaceId: string, options?: { provider?: string }): Promise getToken?(): Promise | string getBaseUrl?(): string - } + } & Partial + +const hasResourceSubscriptionsSdk = ( + client: RelayFileClientLike, +): client is RelayFileClientLike & ResourceSubscriptionsSdk => + typeof client.createOrRenewDurableResourceSubscription === 'function' && + typeof client.claimDurableSubscriptionDeliveries === 'function' && + typeof client.acceptDurableSubscriptionDelivery === 'function' && + typeof client.cancelDurableResourceSubscription === 'function' export function relayfileWorkspaceTokenProvider( client: RelayFileClientLike, @@ -225,6 +239,7 @@ export class RelayfileCloudMountClient implements MountClient { readonly workspaceId: string readonly writebackTransport = 'relayfile-cloud' readonly githubWrite: GithubConnectionWrite + readonly resourceSubscriptions?: ResourceSubscriptionsClient readonly integrationConnections?: FactoryIntegrationConnections readonly #client: RelayFileClientLike @@ -269,6 +284,13 @@ export class RelayfileCloudMountClient implements MountClient { this.#tokenProvider = config.tokenProvider ?? (() => this.#client.getToken?.()) this.#baseUrl = config.baseUrl ?? this.#client.getBaseUrl?.() this.#eventClient = config.eventClient + this.resourceSubscriptions = config.resourceSubscriptions ?? ( + hasResourceSubscriptionsSdk(this.#client) + ? createResourceSubscriptionsSdkClient(this.#client, { + signal: config.resourceSubscriptionSignal, + }) + : undefined + ) this.#logger = config.logger this.#onLocalMountHealth = config.onLocalMountHealth this.#localMountHealthIntervalMs = Math.max( diff --git a/src/orchestrator/factory.test.ts b/src/orchestrator/factory.test.ts index 4c64c20..cbd163e 100644 --- a/src/orchestrator/factory.test.ts +++ b/src/orchestrator/factory.test.ts @@ -35,6 +35,14 @@ import { InMemoryStateStore } from '../state/in-memory-state-store' import { FileStateStore } from '../state/file-state-store' import { githubIssuePathParts, githubRepoSubscriptionGlobs, keyFromPath } from './factory' import { globMatchesPath } from '../subscriptions/globs' +import { + ResourceSubscriptionsUnavailableError, + type AcceptedResourceDelivery, + type ResourceDeliveryClaim, + type ResourceSubscription, + type ResourceSubscriptionInput, + type ResourceSubscriptionsClient, +} from '../subscriptions' import { InternalFleetClient, type HarnessDriverClientLike } from '../fleet/internal-fleet-client' import type { ConversationMessage, ConversationSessionState, DispatchLifecycle } from '../ports/state' @@ -16127,6 +16135,99 @@ describe('FactoryLoop PR babysitter', () => { }) } + class FakeResourceSubscriptions implements ResourceSubscriptionsClient { + readonly createCalls: Array<{ workspaceId: string; input: ResourceSubscriptionInput }> = [] + readonly claimCalls: Array<{ workspaceId: string; limit?: number }> = [] + readonly accepted: Array<{ workspaceId: string; deliveryId: string; claimToken: string }> = [] + readonly cancelled: Array<{ workspaceId: string; subscriptionId: string }> = [] + readonly records: ResourceSubscription[] = [] + claims: ResourceDeliveryClaim[] = [] + readonly leasedClaims = new Map() + readonly acceptedClaims = new Map() + ownerId = 'configured-factory-agent' + unavailable = false + createFailure?: Error + claimFailure?: Error + onAccept?: (deliveryId: string) => Promise | void + + async createOrRenew(workspaceId: string, input: ResourceSubscriptionInput): Promise { + if (this.unavailable) throw new ResourceSubscriptionsUnavailableError() + if (this.createFailure) throw this.createFailure + this.createCalls.push({ workspaceId, input: structuredClone(input) }) + const identity = JSON.stringify([input.subscriberId, input.resourceRef, [...input.eventTypes].sort()]) + let record = this.records.find((candidate) => JSON.stringify([ + candidate.subscriberId, + candidate.resourceRef, + [...candidate.eventTypes].sort(), + ]) === identity) + if (!record) { + record = { + ...structuredClone(input), + ownerId: this.ownerId, + eventTypes: [...input.eventTypes].sort(), + subscriptionId: `sub-${this.records.length + 1}`, + expiresAt: '2026-12-31T00:00:00.000Z', + } + this.records.push(record) + } + return structuredClone(record) + } + + claimFor( + subscription: ResourceSubscription, + eventType: string, + deliveryId: string, + ): ResourceDeliveryClaim { + return { + deliveryId, + claimToken: `claim-token-${deliveryId}`, + subscriptionId: subscription.subscriptionId, + provider: subscription.provider, + resourceRef: subscription.resourceRef, + eventType, + ownerId: subscription.ownerId, + subscriberId: subscription.subscriberId, + terminal: subscription.terminalEventTypes?.includes(eventType) === true, + } + } + + async claimDeliveryClaims(workspaceId: string, input?: { limit?: number }): Promise { + if (this.unavailable) throw new ResourceSubscriptionsUnavailableError() + if (this.claimFailure) throw this.claimFailure + this.claimCalls.push({ workspaceId, ...input }) + // Model Relayfile's atomic claim lease: overlapping runtime polls can + // observe a delivery at most once, independent of client-side locking. + const claimed = this.claims.splice(0, input?.limit ?? this.claims.length) + for (const claim of claimed) this.leasedClaims.set(claim.deliveryId, claim) + return claimed.map((claim) => structuredClone(claim)) + } + + async acceptDelivery( + workspaceId: string, + input: { deliveryId: string; claimToken: string }, + ): Promise { + if (this.unavailable) throw new ResourceSubscriptionsUnavailableError() + const accepted = this.acceptedClaims.get(input.deliveryId) + if (accepted) { + if (accepted.claimToken !== input.claimToken) throw new Error(`claim ${input.deliveryId} has an invalid token`) + this.accepted.push({ workspaceId, ...input }) + return structuredClone(accepted.receipt) + } + const claim = this.leasedClaims.get(input.deliveryId) + if (!claim) throw new Error(`claim ${input.deliveryId} is unavailable`) + if (claim.claimToken !== input.claimToken) throw new Error(`claim ${input.deliveryId} has an invalid token`) + this.accepted.push({ workspaceId, ...input }) + const receipt = { deliveryId: claim.deliveryId, subscriptionId: claim.subscriptionId, terminal: claim.terminal } + this.acceptedClaims.set(input.deliveryId, { claimToken: input.claimToken, receipt }) + await this.onAccept?.(input.deliveryId) + this.leasedClaims.delete(input.deliveryId) + return structuredClone(receipt) + } + + async cancel(workspaceId: string, input: { subscriptionId: string }): Promise { + this.cancelled.push({ workspaceId, ...input }) + } + } it('releases a weak-match babysitter when exact branch reconciliation proves a different PR', async () => { const issue = realIssueFile(495, ready, { title: 'Real exact PR ownership' }) const mount = new FakeMountClient({ [issuePath(495)]: issue }, { @@ -16821,6 +16922,402 @@ describe('FactoryLoop PR babysitter', () => { expect(fleet.spawns.filter((s) => s.name === 'ar-403-babysit')).toHaveLength(1) }) + it('uses Relayfile by-id delivery claims for renamed PR activity, renews idempotently, and retires terminal claims after the local queue is durable', async () => { + const number = 604 + const issue = realIssueFile(number, ready, { title: 'Real durable resource babysitter' }) + const prPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/metadata.json` + const renamedCommentPath = `/github/repos/AgentWorkforce/pear/pulls/${number}__renamed-after-review/comments/6041.json` + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + const subscriptions = new FakeResourceSubscriptions() + mount.resourceSubscriptions = subscriptions + const fleet = new FakeFleetClient() + const state = new InMemoryStateStore({ batchSize: 2 }) + const factory = createFactory(babysitterConfig(), { mount, fleet, triage: new StaticTriage(), stateStore: state }) + + try { + await factory.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + await factory.dispatch(await factory.triageIssue(parseLinearIssue(issuePath(number), issue))) + mount.files.set(prPath, { content: { number, state: 'open', draft: false, head_ref: `ar-${number}-fix` } }) + mount.emit(changeEvent(prPath, 'durable-pr-open')) + await flush() + await vi.waitFor(() => expect(fleet.spawns.map((spawn) => spawn.name)).toContain(`ar-${number}-babysit`)) + await vi.waitFor(() => expect(subscriptions.records).toHaveLength(1)) + + const [[sessionKey, initialSession]] = await state.listBabysitterSessions('factory-test') + expect(initialSession.resourceSubscription).toMatchObject({ + resourceRef: `/github/repos/AgentWorkforce__pear/pulls/by-id/${number}.json`, + ownerId: 'configured-factory-agent', + subscriberId: `factory-babysitter:uuid-${number}`, + }) + const subscription = initialSession.resourceSubscription! + expect(subscriptions.createCalls[0]?.input).toMatchObject({ + eventTypes: expect.arrayContaining(['pull_request_review_comment.created']), + terminalEventTypes: ['pull_request.closed'], + }) + expect(subscriptions.createCalls[0]?.input.eventTypes).not.toContain('pull_request.closed') + subscriptions.claims = [ + { + deliveryId: 'delivery-unowned', + claimToken: 'claim-token-unowned', + subscriptionId: 'sub-unowned', + resourceRef: '/github/repos/AgentWorkforce__pear/pulls/by-id/999.json', + eventType: 'pull_request_review_comment.created', + ownerId: 'configured-factory-agent', + subscriberId: 'factory-babysitter:unowned', + provider: 'github', + terminal: false, + }, + subscriptions.claimFor(subscription, 'pull_request_review_comment.created', 'delivery-renamed'), + ] + let sessionAtAcceptance: Awaited> | undefined + subscriptions.onAccept = async (deliveryId) => { + if (deliveryId === 'delivery-renamed') sessionAtAcceptance = await state.listBabysitterSessions('factory-test') + } + + // The raw path carries a new title slug. Factory does not transform it or + // scan its local owners: the service's stable by-id delivery claim picks + // the one session that must wake. + // Two raw events may race in the runtime; the service lease gives the + // babysitter one delivery and therefore one injected wake. + mount.emit(changeEvent(renamedCommentPath, 'renamed-pr-comment-a')) + mount.emit(changeEvent(renamedCommentPath, 'renamed-pr-comment-b')) + await vi.waitFor(() => expect( + fleet.messages.filter((message) => message.text.startsWith(' message.text.startsWith(' entry.deliveryId)).toEqual(['delivery-renamed']) + expect(sessionAtAcceptance?.find(([key]) => key === sessionKey)?.[1]).toMatchObject({ + pendingKinds: ['pull-request-state'], + pendingDeliveryClaims: [{ deliveryId: 'delivery-renamed', claimToken: 'claim-token-delivery-renamed' }], + }) + expect(factory.status().counters.babysitterResourceDeliveriesIgnoredUnowned).toBe(1) + + // A duplicate PR observation renews the same server identity rather than + // making another subscription record. There is no claim, so no wake. + const renamedPrPath = `/github/repos/AgentWorkforce/pear/pulls/${number}__renamed-after-review/metadata.json` + mount.files.set(renamedPrPath, { content: { number, state: 'open', draft: false, head_ref: `ar-${number}-fix` } }) + mount.emit(changeEvent(renamedPrPath, 'durable-pr-repeat')) + await vi.waitFor(() => expect(subscriptions.createCalls).toHaveLength(2)) + expect(subscriptions.records).toHaveLength(1) + expect(fleet.messages.filter((message) => message.text.startsWith(' key === sessionKey)![1] + subscriptions.claims = [subscriptions.claimFor( + subscriptions.records[0]!, + 'pull_request.closed', + 'delivery-terminal', + )] + let terminalSessionAtAcceptance: Awaited> | undefined + subscriptions.onAccept = async (deliveryId) => { + if (deliveryId === 'delivery-terminal') terminalSessionAtAcceptance = await state.listBabysitterSessions('factory-test') + } + mount.emit(changeEvent(renamedCommentPath, 'renamed-pr-terminal')) + await vi.waitFor(() => expect(subscriptions.accepted.map((entry) => entry.deliveryId)) + .toEqual(['delivery-renamed', 'delivery-terminal'])) + expect(terminalSessionAtAcceptance?.find(([key]) => key === sessionKey)?.[1].resourceSubscription).toMatchObject({ + terminal: true, + }) + expect((await state.listBabysitterSessions('factory-test')).find(([key]) => key === sessionKey)?.[1].resourceSubscription) + .toMatchObject({ terminal: true }) + expect(factory.status().counters.babysitterResourceSubscriptionsRetiredTerminal).toBe(1) + await vi.waitFor(() => expect( + fleet.messages.filter((message) => message.text.startsWith(' message.text.startsWith(' setTimeout(resolve, 900)) + expect(subscriptions.createCalls).toHaveLength(2) + expect(fleet.messages.filter((message) => message.text.startsWith(' setTimeout(resolve, 900)) + expect(fleet.messages.filter((message) => message.text.startsWith(' { + const number = 605 + const issue = realIssueFile(number, ready, { title: 'Real durable subscription fallback' }) + const prPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/metadata.json` + const commentPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/comments/6051.json` + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + const subscriptions = new FakeResourceSubscriptions() + subscriptions.unavailable = true + mount.resourceSubscriptions = subscriptions + const fleet = new FakeFleetClient() + const factory = createFactory(babysitterConfig(), { mount, fleet, triage: new StaticTriage() }) + + try { + await factory.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + await factory.dispatch(await factory.triageIssue(parseLinearIssue(issuePath(number), issue))) + mount.files.set(prPath, { content: { number, state: 'open', draft: false, head_ref: `ar-${number}-fix` } }) + mount.emit(changeEvent(prPath, 'fallback-pr-open')) + await vi.waitFor(() => expect(fleet.spawns.map((spawn) => spawn.name)).toContain(`ar-${number}-babysit`)) + + mount.emit(changeEvent(commentPath, 'fallback-comment')) + await vi.waitFor(() => expect( + fleet.messages.filter((message) => message.text.startsWith(' message.to), + ).toEqual([`ar-${number}-babysit`])) + expect(factory.status().counters.babysitterResourceSubscriptionUnavailable).toBeGreaterThan(0) + } finally { + await factory.stop() + } + }) + + it('keeps local routing for an unregistered owner, then renews its durable subscription without double delivery', async () => { + vi.useFakeTimers() + const number = 610 + const issue = realIssueFile(number, ready, { title: 'Real durable renewal' }) + const commentPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/comments/6101.json` + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + seedPrMeta(mount, 'AgentWorkforce/pear', number, { state: 'open', draft: false }) + const subscriptions = new FakeResourceSubscriptions() + subscriptions.createFailure = Object.assign(new Error('transient create failure'), { status: 503 }) + mount.resourceSubscriptions = subscriptions + const fleet = new FakeFleetClient() + const factory = createFactory(babysitterConfig(), { + mount, + fleet, + triage: new StaticTriage(), + probePrResolver: async () => ({ repo: 'AgentWorkforce/pear', prNumber: number }), + }) + + try { + await factory.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + await factory.dispatch(await factory.triageIssue(parseLinearIssue(issuePath(number), issue))) + fleet.emitAgentExit(`ar-${number}-impl-pear`, 'worker_exited') + await vi.advanceTimersByTimeAsync(0) + await vi.waitFor(() => expect( + factory.status().counters.babysitterResourceSubscriptionRenewFailures, + ).toBe(1)) + expect(subscriptions.records).toEqual([]) + + mount.emit(changeEvent(commentPath, 'unregistered-owner-comment')) + await vi.advanceTimersByTimeAsync(800) + expect( + fleet.messages.filter((message) => message.text.startsWith(' message.to), + ).toEqual([`ar-${number}-babysit`]) + + subscriptions.createFailure = undefined + await vi.advanceTimersByTimeAsync(5_000) + expect(subscriptions.records).toHaveLength(1) + + mount.emit(changeEvent(commentPath, 'registered-owner-comment')) + await vi.advanceTimersByTimeAsync(800) + expect(fleet.messages.filter((message) => message.text.startsWith(' { + const number = 607 + const issue = realIssueFile(number, ready, { title: 'Real durable terminal close' }) + const prPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/metadata.json` + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + const subscriptions = new FakeResourceSubscriptions() + mount.resourceSubscriptions = subscriptions + const fleet = new FakeFleetClient() + const state = new InMemoryStateStore({ batchSize: 2 }) + seedPrMeta(mount, 'AgentWorkforce/pear', number, { state: 'open', draft: false }) + const factory = createFactory(babysitterConfig(), { + mount, + fleet, + triage: new StaticTriage(), + stateStore: state, + probePrResolver: async () => ({ repo: 'AgentWorkforce/pear', prNumber: number }), + }) + + try { + await factory.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + await factory.dispatch(await factory.triageIssue(parseLinearIssue(issuePath(number), issue))) + fleet.emitAgentExit(`ar-${number}-impl-pear`, 'worker_exited') + await vi.waitFor(() => expect(subscriptions.records).toHaveLength(1)) + await vi.waitFor(() => expect(fleet.spawns.map((spawn) => spawn.name)).toContain(`ar-${number}-babysit`)) + + const subscription = subscriptions.records[0]! + subscriptions.claims = [subscriptions.claimFor(subscription, 'pull_request.closed', 'delivery-close')] + mount.files.set(prPath, { content: { number, state: 'closed', draft: false, head_ref: `ar-${number}-fix` } }) + mount.emit(changeEvent(prPath, 'terminal-close')) + + await vi.waitFor(() => expect(subscriptions.accepted.map((entry) => entry.deliveryId)).toEqual(['delivery-close'])) + expect(subscriptions.cancelled).toEqual([{ + workspaceId: 'factory-test', + subscriptionId: subscription.subscriptionId, + }]) + await expect(state.listBabysitterSessions('factory-test')).resolves.toEqual([]) + } finally { + await factory.stop() + } + }) + + it('reuses the persisted claim token after a crash following terminal acceptance', async () => { + class FailingTerminalAcceptancePersistStore extends InMemoryStateStore { + failFinalTerminalPersist = false + + override async setBabysitterSession(...args: Parameters): Promise { + const [, , session] = args + if ( + this.failFinalTerminalPersist && + session.resourceSubscription?.terminal && + !session.pendingDeliveryClaims?.length + ) { + this.failFinalTerminalPersist = false + throw new Error('process stopped after Relayfile accepted the terminal delivery') + } + await super.setBabysitterSession(...args) + } + } + + const number = 609 + const issue = realIssueFile(number, ready, { title: 'Real durable acceptance restart' }) + const prPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/metadata.json` + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + const subscriptions = new FakeResourceSubscriptions() + mount.resourceSubscriptions = subscriptions + const state = new FailingTerminalAcceptancePersistStore({ batchSize: 2 }) + const first = createFactory(babysitterConfig(), { + mount, + fleet: new FakeFleetClient(), + triage: new StaticTriage(), + stateStore: state, + }) + let restarted: ReturnType | undefined + + try { + await first.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + await first.dispatch(await first.triageIssue(parseLinearIssue(issuePath(number), issue))) + mount.files.set(prPath, { content: { number, state: 'open', draft: false, head_ref: `ar-${number}-fix` } }) + mount.emit(changeEvent(prPath, 'acceptance-restart-open')) + await flush() + await vi.waitFor(() => expect(subscriptions.records).toHaveLength(1)) + + const subscription = subscriptions.records[0]! + subscriptions.claims = [subscriptions.claimFor(subscription, 'pull_request.closed', 'delivery-accepted-before-crash')] + state.failFinalTerminalPersist = true + mount.emit(changeEvent(`/github/repos/AgentWorkforce/pear/pulls/${number}/comments/6091.json`, 'acceptance-restart-terminal')) + + await vi.waitFor(() => expect(subscriptions.accepted.map((entry) => entry.deliveryId)) + .toEqual(['delivery-accepted-before-crash'])) + await vi.waitFor(async () => expect(await state.listBabysitterSessions('factory-test')).toEqual([ + [expect.any(String), expect.objectContaining({ + resourceSubscription: expect.objectContaining({ terminal: true }), + pendingDeliveryClaims: [{ + deliveryId: 'delivery-accepted-before-crash', + claimToken: 'claim-token-delivery-accepted-before-crash', + }], + })], + ])) + await first.stop() + + restarted = createFactory(babysitterConfig(), { + mount, + fleet: new FakeFleetClient(), + triage: new StaticTriage(), + stateStore: state, + }) + await restarted.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + + await vi.waitFor(() => expect(subscriptions.accepted.map((entry) => entry.deliveryId)) + .toEqual(['delivery-accepted-before-crash', 'delivery-accepted-before-crash'])) + await vi.waitFor(async () => { + const [[, restored]] = await state.listBabysitterSessions('factory-test') + expect(restored?.resourceSubscription).toMatchObject({ terminal: true }) + expect(restored?.pendingDeliveryClaims).toBeUndefined() + }) + } finally { + await first.stop() + await restarted?.stop() + } + }) + + it('cancels the durable subscription before completion clears the babysitter owner', async () => { + const number = 608 + const issue = realIssueFile(number, ready, { title: 'Real durable completion cancellation' }) + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + seedPrMeta(mount, 'AgentWorkforce/pear', number, { state: 'open', draft: false }) + const subscriptions = new FakeResourceSubscriptions() + mount.resourceSubscriptions = subscriptions + const fleet = new FakeFleetClient() + const factory = createFactory(babysitterConfig(), { + mount, + fleet, + triage: new StaticTriage(), + probePrResolver: async () => ({ repo: 'AgentWorkforce/pear', prNumber: number }), + }) + + try { + await factory.dispatch(await factory.triageIssue(parseLinearIssue(issuePath(number), issue))) + fleet.emitAgentExit(`ar-${number}-impl-pear`, 'worker_exited') + await vi.waitFor(() => expect(subscriptions.records).toHaveLength(1)) + const subscription = subscriptions.records[0]! + + fleet.emitAgentMessage({ from: `ar-${number}-babysit`, target: 'factory', body: `[factory-pr-ready] AR-${number}` }) + await vi.waitFor(() => expect(subscriptions.cancelled).toEqual([{ + workspaceId: 'factory-test', + subscriptionId: subscription.subscriptionId, + }])) + } finally { + await factory.stop() + } + }) + + it('does not double-route locally when an established durable delivery lookup is transiently unavailable', async () => { + const number = 606 + const issue = realIssueFile(number, ready, { title: 'Real durable transient claim retry' }) + const prPath = `/github/repos/AgentWorkforce/pear/pulls/${number}/metadata.json` + const commentPath = `/github/repos/AgentWorkforce/pear/pulls/${number}__renamed/comments/6061.json` + const mount = new FakeMountClient({ [issuePath(number)]: issue }) + const subscriptions = new FakeResourceSubscriptions() + mount.resourceSubscriptions = subscriptions + const fleet = new FakeFleetClient() + const factory = createFactory(babysitterConfig(), { mount, fleet, triage: new StaticTriage() }) + + try { + await factory.start({ mode: 'live', liveSubscription: { transport: 'subscribe' } }) + await factory.dispatch(await factory.triageIssue(parseLinearIssue(issuePath(number), issue))) + mount.files.set(prPath, { content: { number, state: 'open', draft: false, head_ref: `ar-${number}-fix` } }) + mount.emit(changeEvent(prPath, 'transient-pr-open')) + await vi.waitFor(() => expect(subscriptions.records).toHaveLength(1)) + const subscription = subscriptions.records[0]! + + subscriptions.claimFailure = Object.assign(new Error('temporarily unavailable'), { status: 503 }) + mount.emit(changeEvent(commentPath, 'transient-claim-event')) + await new Promise((resolve) => setTimeout(resolve, 900)) + expect(fleet.messages.filter((message) => message.text.startsWith(' expect( + fleet.messages.filter((message) => message.text.startsWith(' message.to), + ).toEqual([`ar-${number}-babysit`])) + expect(subscriptions.accepted.map((entry) => entry.deliveryId)).toEqual(['delivery-after-recovery']) + } finally { + await factory.stop() + } + }) + it('publishes a ready PR and reaches Human Review through lifecycle actions without control identities', async () => { class LifecycleActionFleet extends FakeFleetClient { override readonly lifecycleActionName = 'factory.lifecycle' diff --git a/src/orchestrator/factory.ts b/src/orchestrator/factory.ts index 8be3fdd..1dc7eab 100644 --- a/src/orchestrator/factory.ts +++ b/src/orchestrator/factory.ts @@ -60,6 +60,7 @@ import { import { resolveTestGuidance } from '../dispatch/test-guidance' import { HeuristicTriage, TieredTriage, babysitterSpec, isShapeLabel, scopeFromLabels } from '../triage' import { agentNameForRole, sanitizeAgentSlug } from '../triage/agent-names' +import { isResourceSubscriptionsUnavailable, type ResourceSubscription } from '../subscriptions' import type { DispatchResult, Factory, @@ -144,7 +145,19 @@ type BabysitterWakeKind = | 'checks-failed' | 'merge-conflict' | 'base-diverged' -type BabysitterPrRef = { repo: string; prNumber: number; path?: string; agentName: string } +type BabysitterResourceSubscription = Pick< + ResourceSubscription, + 'subscriptionId' | 'provider' | 'resourceRef' | 'subscriberId' | 'ownerId' | 'expiresAt' +> & { terminal?: boolean } +type BabysitterPendingDeliveryClaim = { deliveryId: string; claimToken: string } +type BabysitterPrRef = { + repo: string + prNumber: number + path?: string + agentName: string + resourceSubscription?: BabysitterResourceSubscription + pendingDeliveryClaims?: BabysitterPendingDeliveryClaim[] +} type BabysitterWakeState = { issue: IssueRef repo: string @@ -262,6 +275,24 @@ const INJECTION_RETRY_ATTEMPT_TIMEOUT_MS = 15_000 const INJECTION_MAX_ATTEMPTS = 6 const BABYSITTER_EVENT_COALESCE_MS = 750 const BABYSITTER_EVENT_RETRY_MS = 1_000 +const BABYSITTER_SUBSCRIPTION_TTL_SECONDS = 60 * 60 +// Relayfile receives provider-native GitHub events, not the materialized file +// changes that the legacy local router consumed. `closed` is separately +// indexed as a terminal event below, so terminal delivery does not need a +// broad normal-event subscription. +const BABYSITTER_SUBSCRIPTION_EVENT_TYPES = [ + 'pull_request.opened', + 'pull_request.reopened', + 'pull_request.synchronize', + 'pull_request.ready_for_review', + 'pull_request_review.submitted', + 'pull_request_review_comment.created', + 'issue_comment.created', + 'check_run.completed', +] +const BABYSITTER_SUBSCRIPTION_TERMINAL_EVENT_TYPES = ['pull_request.closed'] +const BABYSITTER_RESOURCE_DELIVERY_RETRY_MS = 5_000 +const BABYSITTER_RESOURCE_SUBSCRIPTION_RENEW_MS = (BABYSITTER_SUBSCRIPTION_TTL_SECONDS * 1_000) / 2 // A babysitter wake fails with a "registration lag" error (agent_not_found / // recipient unavailable) both when a freshly spawned agent has not finished // enrolling AND when an agent is up but its relay identity never becomes @@ -460,6 +491,14 @@ export class FactoryLoop implements Factory { // webhook-fed mount path so readiness can re-read PR meta without a gh call. readonly #babysitterPr = new Map() readonly #babysitterIssueRefs = new Map() + // Relayfile matches subscription IDs server-side. This direct index means a + // delivery claim never requires the legacy local repo/PR scan to find an + // owning babysitter. + readonly #babysitterSubscriptionOwners = new Map() + #babysitterResourceSubscriptionFault = false + #babysitterResourceSubscriptionUnavailable = false + #babysitterResourceDeliveryRetryTimer?: ReturnType + #babysitterResourceSubscriptionRenewTimer?: ReturnType readonly #babysitterReady = new Set() readonly #babysitterWakeStates = new Map() // A babysitter announces this fence before invoking destructive git tooling @@ -811,6 +850,10 @@ export class FactoryLoop implements Factory { async stop(): Promise { this.#started = false this.#stopping = true + if (this.#babysitterResourceDeliveryRetryTimer) clearTimeout(this.#babysitterResourceDeliveryRetryTimer) + this.#babysitterResourceDeliveryRetryTimer = undefined + if (this.#babysitterResourceSubscriptionRenewTimer) clearTimeout(this.#babysitterResourceSubscriptionRenewTimer) + this.#babysitterResourceSubscriptionRenewTimer = undefined if (this.#dispatchLifecycleRenewTimer) clearInterval(this.#dispatchLifecycleRenewTimer) this.#dispatchLifecycleRenewTimer = undefined for (const timer of this.#dispatchLifecycleRetryTimers.values()) clearTimeout(timer) @@ -861,6 +904,7 @@ export class FactoryLoop implements Factory { this.#babysitterSpawned.clear() this.#babysitterPr.clear() this.#babysitterIssueRefs.clear() + this.#babysitterSubscriptionOwners.clear() this.#babysitterReady.clear() this.#babysitterCriticalAgents.clear() const subscription = this.#subscription @@ -8564,8 +8608,13 @@ export class FactoryLoop implements Factory { prNumber: session.prNumber, path: session.path, agentName: session.agentName, + resourceSubscription: session.resourceSubscription, + pendingDeliveryClaims: session.pendingDeliveryClaims, } this.#babysitterPr.set(ownershipKey, ref) + if (ref.resourceSubscription) { + this.#babysitterSubscriptionOwners.set(ref.resourceSubscription.subscriptionId, ownershipKey) + } this.#babysitterIssueRefs.set(ownershipKey, { ...session.issue }) this.#babysitterSpawned.add(ownershipKey) if (session.critical) this.#babysitterCriticalAgents.add(session.agentName) @@ -8576,10 +8625,21 @@ export class FactoryLoop implements Factory { this.#increment('babysitterOwnershipRestored') const pendingKinds = session.pendingKinds.filter(isBabysitterWakeKind) if (pendingKinds.length > 0) { - await this.#queueBabysitterWake(session.issue, ref, pendingKinds, tracked) + // A terminal marker can coexist with the one wake persisted before + // acknowledgement. Rehydrate that durable hand-off once, while later + // raw events remain quarantined by #queueBabysitterWake. + await this.#queueBabysitterWake(session.issue, ref, pendingKinds, tracked, { allowTerminal: true }) this.#increment('babysitterPendingWakesRestored') } + await this.#ensureBabysitterResourceSubscription(session.issue, ref, tracked) } + // A crash after the local queue write but before (or just after) the + // remote acceptance leaves an ID in state. Accept is idempotent once the + // lease was accepted, so settle those durable hand-offs before claiming + // new work; lease expiry is retried below when a server has not released + // the original claim yet. + await this.#retryPendingBabysitterDeliveryAcceptances() + await this.#routeDurableBabysitterDeliveries() } async #reconcileRestoredBabysitterReceipts(onlyRecord?: InFlightIssue): Promise { @@ -8680,6 +8740,28 @@ export class FactoryLoop implements Factory { this.#babysitterWakeStates.delete(key) this.#babysitterCriticalAgents.delete(state.agentName) } + if (mayClearDurable && ref?.resourceSubscription) { + this.#babysitterSubscriptionOwners.delete(ref.resourceSubscription.subscriptionId) + const subscriptions = this.#mount.resourceSubscriptions + if (subscriptions) { + try { + await subscriptions.cancel(this.#workspaceId, { + subscriptionId: ref.resourceSubscription.subscriptionId, + }) + this.#increment('babysitterResourceSubscriptionsCancelled') + } catch (error) { + // Cancellation is deliberately idempotent. A terminal acceptance may + // have retired this record already; an outage leaves the bounded TTL + // as the leak backstop and must not prevent local session cleanup. + this.#increment('babysitterResourceSubscriptionCancelFailures') + this.#logger.warn?.('[factory] could not cancel durable babysitter resource subscription', { + issue: issue?.key, + subscriptionId: ref.resourceSubscription.subscriptionId, + error: describeError(error).errorMessage, + }) + } + } + } this.#babysitterPr.delete(ownershipKey) this.#babysitterIssueRefs.delete(ownershipKey) this.#babysitterSpawned.delete(ownershipKey) @@ -8695,7 +8777,331 @@ export class FactoryLoop implements Factory { await Promise.all(keys.map(async (key) => this.#cancelBabysitterWake(key))) } + async #ensureBabysitterResourceSubscription( + issue: IssueRef, + ref: BabysitterPrRef, + tracked?: TrackedAgent, + ): Promise { + const subscriptions = this.#mount.resourceSubscriptions + if (!subscriptions || !ref.agentName || !await this.#assertIssueDispatchLifecycleOwner(issue)) return + // A terminal claim is persisted before Relayfile acceptance so a crash in + // that gap can never renew a retired record into a fresh generation. + if (ref.resourceSubscription?.terminal) return + + const resourceRef = babysitterResourceRef(ref.repo, ref.prNumber) + const subscriberId = babysitterSubscriberId(issue) + try { + const subscription = await subscriptions.createOrRenew(this.#workspaceId, { + provider: 'github', + resourceRef, + eventTypes: [...BABYSITTER_SUBSCRIPTION_EVENT_TYPES], + terminalEventTypes: [...BABYSITTER_SUBSCRIPTION_TERMINAL_EVENT_TYPES], + subscriberId, + ttlSeconds: BABYSITTER_SUBSCRIPTION_TTL_SECONDS, + }) + if ( + !subscription.subscriptionId || + subscription.provider !== 'github' || + subscription.resourceRef !== resourceRef || + subscription.subscriberId !== subscriberId || + !subscription.ownerId || + !subscription.expiresAt || + !subscription.terminalEventTypes?.includes('pull_request.closed') + ) { + throw new Error('Relayfile returned an invalid durable resource subscription') + } + if (ref.resourceSubscription?.subscriptionId && ref.resourceSubscription.subscriptionId !== subscription.subscriptionId) { + this.#babysitterSubscriptionOwners.delete(ref.resourceSubscription.subscriptionId) + } + ref.resourceSubscription = { + subscriptionId: subscription.subscriptionId, + provider: subscription.provider, + resourceRef: subscription.resourceRef, + subscriberId: subscription.subscriberId, + ownerId: subscription.ownerId, + expiresAt: subscription.expiresAt, + } + this.#babysitterSubscriptionOwners.set( + subscription.subscriptionId, + babysitterOwnershipKey(issue, ref), + ) + await this.#persistBabysitterSession(issue, ref, tracked) + this.#babysitterResourceSubscriptionFault = false + this.#babysitterResourceSubscriptionUnavailable = false + this.#scheduleBabysitterResourceSubscriptionRenewal() + this.#increment('babysitterResourceSubscriptionsRenewed') + } catch (error) { + if (isResourceSubscriptionsUnavailable(error)) { + this.#babysitterResourceSubscriptionFault = false + this.#babysitterResourceSubscriptionUnavailable = true + this.#increment('babysitterResourceSubscriptionUnavailable') + return + } + this.#babysitterResourceSubscriptionFault = true + this.#scheduleDurableBabysitterDeliveryRetry() + this.#increment('babysitterResourceSubscriptionRenewFailures') + this.#logger.warn?.('[factory] could not create or renew durable babysitter resource subscription', { + issue: issue.key, + repo: ref.repo, + prNumber: ref.prNumber, + error: describeError(error).errorMessage, + }) + } + } + + async #routeDurableBabysitterDeliveries(): Promise { + const subscriptions = this.#mount.resourceSubscriptions + if (!subscriptions || !this.#config.babysitter.enabled || this.#stopping) return false + // Do not bypass the proven local router until every active babysitter has + // completed its own create/renew. This closes the rollout and transient + // provisioning gap without making a successful API response for some + // other subscription suppress an unregistered PR's wake. + if ([...this.#babysitterPr.values()].some((ref) => ref.agentName && !ref.resourceSubscription)) { + return false + } + + let claims: Awaited> + try { + claims = await subscriptions.claimDeliveryClaims(this.#workspaceId) + } catch (error) { + if (isResourceSubscriptionsUnavailable(error)) { + this.#babysitterResourceSubscriptionFault = false + this.#babysitterResourceSubscriptionUnavailable = true + this.#increment('babysitterResourceSubscriptionUnavailable') + } else { + this.#babysitterResourceSubscriptionFault = true + this.#scheduleDurableBabysitterDeliveryRetry() + this.#increment('babysitterResourceDeliveryLookupFailures') + this.#logger.warn?.('[factory] durable babysitter delivery-claim lookup failed; retaining durable delivery retry', { + error: describeError(error).errorMessage, + }) + } + return !isResourceSubscriptionsUnavailable(error) + } + this.#babysitterResourceSubscriptionFault = false + this.#babysitterResourceSubscriptionUnavailable = false + + for (const claim of claims) { + const issueIdentity = this.#babysitterSubscriptionOwners.get(claim.subscriptionId) + const issue = issueIdentity ? this.#babysitterIssueRefs.get(issueIdentity) : undefined + const ref = issueIdentity ? this.#babysitterPr.get(issueIdentity) : undefined + const subscription = ref?.resourceSubscription + if ( + !issue || + !ref || + !subscription || + subscription.subscriptionId !== claim.subscriptionId || + subscription.provider !== claim.provider || + subscription.resourceRef !== claim.resourceRef || + subscription.subscriberId !== claim.subscriberId || + subscription.ownerId !== claim.ownerId + ) { + // The service is owner-isolated, but Factory may have just retired a + // local owner. Never route a stale or other-session claim by resource. + this.#increment('babysitterResourceDeliveriesIgnoredUnowned') + continue + } + if (!await this.#assertIssueDispatchLifecycleOwner(issue)) { + this.#increment('babysitterResourceDeliveriesIgnoredNonOwner') + continue + } + + // A terminal delivery may be reclaimed after a process crash before its + // remote acceptance. Only that already-persisted delivery may finish; + // no later claim is allowed to wake or re-open the retired session. + if (subscription.terminal && !ref.pendingDeliveryClaims?.some((pending) => pending.deliveryId === claim.deliveryId)) { + this.#increment('babysitterResourceDeliveriesIgnoredTerminal') + continue + } + + const batch = await this.#batch() + const tracked = batch.getIssue(issue)?.agents.get(ref.agentName) + ?? [...(batch.getIssue(issue)?.agents.values() ?? [])].find((agent) => agent.spec.role === 'babysitter') + ?? durableBabysitterTrackedAgent({ + issue, + repo: ref.repo, + prNumber: ref.prNumber, + path: ref.path, + agentName: ref.agentName, + critical: false, + pendingKinds: [], + resourceSubscription: subscription, + pendingDeliveryClaims: ref.pendingDeliveryClaims, + }) + + const pendingClaim = ref.pendingDeliveryClaims?.find((pending) => pending.deliveryId === claim.deliveryId) + const alreadyQueued = Boolean(pendingClaim) + if (!alreadyQueued) { + const queued = await this.#queueBabysitterWake(issue, ref, ['pull-request-state'], tracked) + if (!queued) continue + } + if (!pendingClaim || pendingClaim.claimToken !== claim.claimToken) { + ref.pendingDeliveryClaims = [ + ...(ref.pendingDeliveryClaims ?? []).filter((pending) => pending.deliveryId !== claim.deliveryId), + { deliveryId: claim.deliveryId, claimToken: claim.claimToken }, + ] + // The claim lease joins Factory's durable pending-wake state before + // the external acceptance. A crash after this point can retry the + // exact hand-off without delivering the same wake a second time. + await this.#persistBabysitterSession(issue, ref, tracked) + } + + try { + if (claim.terminal && !subscription.terminal) { + subscription.terminal = true + await this.#persistBabysitterSession(issue, ref, tracked) + } + const accepted = await subscriptions.acceptDelivery(this.#workspaceId, { + deliveryId: claim.deliveryId, + claimToken: claim.claimToken, + }) + if (accepted.deliveryId !== claim.deliveryId || accepted.subscriptionId !== claim.subscriptionId) { + throw new Error('Relayfile accepted a different durable delivery claim') + } + if (accepted.terminal || claim.terminal) { + // Keep the terminal marker and subscription identity locally until + // normal PR/session teardown. That quarantines the babysitter from + // both legacy fallback and a restart-time create-or-renew. + subscription.terminal = true + ref.pendingDeliveryClaims = (ref.pendingDeliveryClaims ?? []).filter((pending) => pending.deliveryId !== claim.deliveryId) + await this.#persistBabysitterSession(issue, ref, tracked) + this.#increment('babysitterResourceSubscriptionsRetiredTerminal') + } else { + ref.pendingDeliveryClaims = (ref.pendingDeliveryClaims ?? []).filter((pending) => pending.deliveryId !== claim.deliveryId) + await this.#persistBabysitterSession(issue, ref, tracked) + } + this.#increment('babysitterResourceDeliveriesAccepted') + } catch (error) { + this.#increment('babysitterResourceDeliveryAcceptFailures') + this.#logger.warn?.('[factory] durable babysitter delivery claim remains pending after wake queue', { + issue: issue.key, + subscriptionId: claim.subscriptionId, + deliveryId: claim.deliveryId, + error: describeError(error).errorMessage, + }) + } + } + if ([...this.#babysitterPr.values()].some((ref) => ref.pendingDeliveryClaims?.length)) { + this.#scheduleDurableBabysitterDeliveryRetry() + } + return true + } + + async #retryPendingBabysitterDeliveryAcceptances(): Promise { + const subscriptions = this.#mount.resourceSubscriptions + if (!subscriptions) return + const retrySubscriptionRenewal = this.#babysitterResourceSubscriptionFault + + for (const [issueIdentity, ref] of this.#babysitterPr) { + const issue = this.#babysitterIssueRefs.get(issueIdentity) + if (issue && ref.agentName && !ref.resourceSubscription?.terminal && (!ref.resourceSubscription || retrySubscriptionRenewal)) { + const batch = await this.#batch() + const tracked = batch.getIssue(issue)?.agents.get(ref.agentName) + ?? [...(batch.getIssue(issue)?.agents.values() ?? [])].find((agent) => agent.spec.role === 'babysitter') + await this.#ensureBabysitterResourceSubscription(issue, ref, tracked) + } + const subscription = ref.resourceSubscription + const pendingDeliveryClaims = [...(ref.pendingDeliveryClaims ?? [])] + if (!subscription || !issue || pendingDeliveryClaims.length === 0) continue + const batch = await this.#batch() + const tracked = batch.getIssue(issue)?.agents.get(ref.agentName) + ?? [...(batch.getIssue(issue)?.agents.values() ?? [])].find((agent) => agent.spec.role === 'babysitter') + ?? durableBabysitterTrackedAgent({ + issue, + repo: ref.repo, + prNumber: ref.prNumber, + path: ref.path, + agentName: ref.agentName, + critical: false, + pendingKinds: [], + resourceSubscription: subscription, + pendingDeliveryClaims, + }) + for (const { deliveryId, claimToken } of pendingDeliveryClaims) { + try { + const accepted = await subscriptions.acceptDelivery(this.#workspaceId, { deliveryId, claimToken }) + if (accepted.deliveryId !== deliveryId || accepted.subscriptionId !== subscription.subscriptionId) { + throw new Error('Relayfile accepted a different durable delivery claim') + } + if (accepted.terminal) subscription.terminal = true + ref.pendingDeliveryClaims = (ref.pendingDeliveryClaims ?? []).filter((pending) => pending.deliveryId !== deliveryId) + await this.#persistBabysitterSession(issue, ref, tracked) + this.#increment('babysitterResourceDeliveriesAcceptedAfterRestore') + } catch (error) { + if (isResourceSubscriptionsUnavailable(error)) { + this.#babysitterResourceSubscriptionFault = false + this.#babysitterResourceSubscriptionUnavailable = true + this.#increment('babysitterResourceSubscriptionUnavailable') + } else { + this.#babysitterResourceSubscriptionFault = true + this.#increment('babysitterResourceDeliveryAcceptFailures') + this.#logger.warn?.('[factory] durable babysitter delivery acceptance remains pending after restore', { + issue: issue.key, + subscriptionId: subscription.subscriptionId, + deliveryId, + error: describeError(error).errorMessage, + }) + } + } + } + } + if ([...this.#babysitterPr.values()].some((ref) => ref.pendingDeliveryClaims?.length)) { + this.#scheduleDurableBabysitterDeliveryRetry() + } + } + + #scheduleBabysitterResourceSubscriptionRenewal(): void { + if ( + this.#babysitterResourceSubscriptionRenewTimer || + this.#stopping || + !this.#mount.resourceSubscriptions || + ![...this.#babysitterPr.values()].some((ref) => ref.resourceSubscription && !ref.resourceSubscription.terminal) + ) return + this.#babysitterResourceSubscriptionRenewTimer = setTimeout(() => { + this.#babysitterResourceSubscriptionRenewTimer = undefined + void (async () => { + const batch = await this.#batch() + for (const [issueIdentity, ref] of this.#babysitterPr) { + const issue = this.#babysitterIssueRefs.get(issueIdentity) + if (!issue || !ref.resourceSubscription || ref.resourceSubscription.terminal) continue + const tracked = batch.getIssue(issue)?.agents.get(ref.agentName) + ?? [...(batch.getIssue(issue)?.agents.values() ?? [])].find((agent) => agent.spec.role === 'babysitter') + await this.#ensureBabysitterResourceSubscription(issue, ref, tracked) + } + })().catch((error) => { + this.#logger.warn?.('[factory] durable babysitter subscription renewal rejected', { + error: describeError(error).errorMessage, + }) + }).finally(() => { + this.#scheduleBabysitterResourceSubscriptionRenewal() + }) + }, BABYSITTER_RESOURCE_SUBSCRIPTION_RENEW_MS) + this.#babysitterResourceSubscriptionRenewTimer.unref?.() + } + + #scheduleDurableBabysitterDeliveryRetry(): void { + if (this.#babysitterResourceDeliveryRetryTimer || this.#stopping || !this.#mount.resourceSubscriptions) return + this.#babysitterResourceDeliveryRetryTimer = setTimeout(() => { + this.#babysitterResourceDeliveryRetryTimer = undefined + void (async () => { + await this.#retryPendingBabysitterDeliveryAcceptances() + await this.#routeDurableBabysitterDeliveries() + })().catch((error) => { + this.#logger.warn?.('[factory] durable babysitter delivery retry rejected', { + error: describeError(error).errorMessage, + }) + }) + }, BABYSITTER_RESOURCE_DELIVERY_RETRY_MS) + this.#babysitterResourceDeliveryRetryTimer.unref?.() + } + async #routeBabysitterEvent(path: string, extraKinds: Iterable = []): Promise { + // A successful Relayfile claim lookup is the new exact demux. While some + // owners are still unregistered, the legacy path router remains available + // only to those owners. Registered owners retain and retry service claims, + // so a mixed rollout or transient create failure neither double-delivers a + // registered PR nor drops an event for an unregistered PR. + if (await this.#routeDurableBabysitterDeliveries()) return const event = githubBabysitterEventPathParts(path) if (!event || !this.#config.babysitter.enabled || this.#stopping) return let targets: Array<{ prNumber: number; kinds: BabysitterWakeKind[] }> @@ -8726,6 +9132,14 @@ export class FactoryLoop implements Factory { this.#logger.debug?.('[factory] ignored unowned PR event for babysitter routing', { ...event, prNumber: target.prNumber }) continue } + if ( + this.#mount.resourceSubscriptions && + !this.#babysitterResourceSubscriptionUnavailable && + owner.ref.resourceSubscription + ) { + this.#increment('babysitterEventsDeferredToDurableSubscription') + continue + } if (!await this.#assertIssueDispatchLifecycleOwner(owner.issue)) { this.#increment('babysitterEventsIgnoredNonOwner') continue @@ -8790,10 +9204,15 @@ export class FactoryLoop implements Factory { ref: BabysitterPrRef, kinds: Iterable, tracked: TrackedAgent, - ): Promise { + options: { allowTerminal?: boolean } = {}, + ): Promise { + if (ref.resourceSubscription?.terminal && !options.allowTerminal) { + this.#increment('babysitterEventsIgnoredTerminal') + return false + } if (!await this.#assertIssueDispatchLifecycleOwner(issue)) { this.#increment('babysitterEventsIgnoredNonOwner') - return + return false } // Owner lookup and queueing straddle async mount/state reads. Revalidate // the exact composite owner so a concurrent close/merge cancellation can @@ -8806,7 +9225,7 @@ export class FactoryLoop implements Factory { githubPrIdentity(current.repo, current.prNumber) !== githubPrIdentity(ref.repo, ref.prNumber) ) { this.#increment('babysitterEventsIgnoredStaleOwner') - return + return false } // Any new event invalidates a prior readiness assertion for this exact PR. this.#babysitterReady.delete(ownershipKey) @@ -8836,9 +9255,10 @@ export class FactoryLoop implements Factory { if (state.deferredSubmitTargets || state.inFlight || this.#babysitterCriticalAgents.has(state.agentName)) { this.#increment('babysitterEventWakesDeferred') - return + return true } this.#scheduleBabysitterWake(state, BABYSITTER_EVENT_COALESCE_MS) + return true } async #recordPendingBabysitterWake(state: BabysitterWakeState): Promise { @@ -8884,6 +9304,8 @@ export class FactoryLoop implements Factory { path: ref.path, critical: this.#babysitterCriticalAgents.has(ref.agentName), pendingKinds: pending?.kinds.filter(isBabysitterWakeKind).sort(compareBabysitterWakeKinds) ?? [], + ...(ref.resourceSubscription ? { resourceSubscription: { ...ref.resourceSubscription } } : {}), + ...(ref.pendingDeliveryClaims?.length ? { pendingDeliveryClaims: structuredClone(ref.pendingDeliveryClaims) } : {}), }) } @@ -9162,10 +9584,19 @@ export class FactoryLoop implements Factory { } if (!this.#config.babysitter.enabled) return if (snapshot.state && snapshot.state.trim().toUpperCase() !== 'OPEN') { + // A provider close produces a separately indexed terminal claim. Keep + // the durable owner until it has been accepted (or a transient retry + // has claimed it); do not let closed-state cleanup erase that hand-off. + await this.#routeDurableBabysitterDeliveries() + if (this.#babysitterResourceSubscriptionFault) return await this.#cancelBabysitterWake(owned.key) return } if (snapshot.draft) this.#increment('babysitterDraftPrSkipped') + // PR meta events are also the normal renewal heartbeat for the durable + // record. The store's identity makes this a create-or-renew, never a + // second subscription for the same babysitter. + await this.#ensureBabysitterResourceSubscription(owned.issue, owned.ref, owned.tracked) await this.#routeBabysitterEvent(path, babysitterWakeKindsFromSnapshot(snapshot)) return } @@ -9220,6 +9651,12 @@ export class FactoryLoop implements Factory { } if (snapshot.state && snapshot.state.trim().toUpperCase() !== 'OPEN') { + // `pull_request.closed` is a separately indexed Relayfile terminal + // event. Claim and accept its durable hand-off before the local closed + // PR cleanup drops the subscription owner. On a transient service fault, + // retain the owner so the retry loop can claim it without local fallback. + await this.#routeDurableBabysitterDeliveries() + if (this.#babysitterResourceSubscriptionFault) return if (babysitterKey && existing) await this.#cancelBabysitterWake(babysitterKey) return } @@ -9428,6 +9865,11 @@ export class FactoryLoop implements Factory { await this.#babysitterSpawnInFlight.get(babysitterKey) const settled = this.#babysitterPr.get(babysitterKey) if (settled && prRef.path) settled.path = prRef.path + if (settled) { + const tracked = record.agents.get(settled.agentName) + ?? [...record.agents.values()].find((agent) => agent.spec.role === 'babysitter') + await this.#ensureBabysitterResourceSubscription(record.issue, settled, tracked) + } return } const wantedPr = githubPrIdentity(prRef.repo, prRef.prNumber) @@ -9449,7 +9891,9 @@ export class FactoryLoop implements Factory { agentName: tracked.result?.name ?? trackedName, }) this.#babysitterSpawned.add(babysitterKey) - await this.#persistBabysitterSession(record.issue, this.#babysitterPr.get(babysitterKey)!, tracked) + const ref = this.#babysitterPr.get(babysitterKey)! + await this.#persistBabysitterSession(record.issue, ref, tracked) + await this.#ensureBabysitterResourceSubscription(record.issue, ref, tracked) await this.#retargetSlackConversationToBabysitter(record) return } @@ -9550,7 +9994,9 @@ export class FactoryLoop implements Factory { path: prRef.path, agentName: tracked?.result?.name ?? spawned.name, }) - await this.#persistBabysitterSession(record.issue, this.#babysitterPr.get(babysitterKey)!, tracked) + const ref = this.#babysitterPr.get(babysitterKey)! + await this.#persistBabysitterSession(record.issue, ref, tracked) + await this.#ensureBabysitterResourceSubscription(record.issue, ref, tracked) await this.#retargetSlackConversationToBabysitter(record) await this.#writeInFlightRegistry() if (!await this.#saveDispatchLifecycle(record, 'running')) return @@ -9975,6 +10421,8 @@ export class FactoryLoop implements Factory { const stateKey = issueStateKey(record.issue) this.#probePrGhBackoffUntilMs.delete(stateKey) this.#probePrResolvedCache.delete(stateKey) + // Cancellation must see the subscription identity so it can issue the + // idempotent Relayfile DELETE before clearing the local owner maps. await this.#cancelBabysittersForIssue(record.issue) const durable = await this.#state.getDispatchLifecycle(this.#workspaceId, issueKey(record.issue)).catch(() => undefined) if (!this.#usesDurableDispatchLifecycle() || (durable && isTerminalDispatchLifecycle(durable))) { @@ -13969,6 +14417,19 @@ const validPrNumber = (value: number): boolean => Number.isInteger(value) && val const githubPrIdentity = (repo: string, prNumber: number): string | undefined => validGithubRepo(repo) && validPrNumber(prNumber) ? `${repo.toLowerCase()}#${prNumber}` : undefined +// This is the public Relayfile stable identity for a GitHub pull request. It +// deliberately comes from the PR's repo/number ownership record, never by +// transforming an incoming canonical path (whose title slug can be renamed). +const babysitterResourceRef = (repo: string, prNumber: number): string => { + if (!validGithubRepo(repo) || !validPrNumber(prNumber)) { + throw new Error('Cannot create a durable babysitter subscription for an invalid GitHub PR identity') + } + const [owner, name] = repo.split('/') + return `/github/repos/${owner}__${name}/pulls/by-id/${prNumber}.json` +} + +const babysitterSubscriberId = (issue: IssueRef): string => `factory-babysitter:${issue.uuid}` + const babysitterOwnershipKey = ( issue: IssueRef, ref: Pick, diff --git a/src/ports/mount.ts b/src/ports/mount.ts index 506a830..c8c4765 100644 --- a/src/ports/mount.ts +++ b/src/ports/mount.ts @@ -2,6 +2,7 @@ import type { ChangeEvent as RelayFileChangeEvent, Subscription as RelayFileSubscription, } from '@relayfile/sdk' +import type { ResourceSubscriptionsClient } from '../subscriptions/resource-subscriptions' export type ChangeEvent = RelayFileChangeEvent export type Subscription = RelayFileSubscription @@ -98,6 +99,12 @@ export interface GithubConnectionWrite { export interface MountClient { readonly writebackTransport?: 'relayfile-cloud' | 'test' readonly githubWrite?: GithubConnectionWrite + /** + * Optional durable Relayfile resource-subscription API. Its absence means + * this mount targets a pre-subscription service and consumers retain their + * established local routing behaviour. + */ + readonly resourceSubscriptions?: ResourceSubscriptionsClient readonly integrationConnections?: FactoryIntegrationConnections /** Ensure the SDK-authenticated Relayfile mirror exists below a checkout. */ ensureLocalMount?(startDir: string, options?: LocalMountOptions): Promise diff --git a/src/ports/state.ts b/src/ports/state.ts index ec90740..c85d967 100644 --- a/src/ports/state.ts +++ b/src/ports/state.ts @@ -75,6 +75,19 @@ export type BabysitterSessionState = { path?: string critical: boolean pendingKinds: string[] + /** Durable Relayfile subscription identity, if the workspace supports it. */ + resourceSubscription?: { + subscriptionId: string + provider: string + resourceRef: string + subscriberId: string + ownerId: string + expiresAt: string + /** Terminal delivery accepted or awaiting acceptance; never renew this record. */ + terminal?: boolean + } + /** Claims already queued locally but not yet accepted by Relayfile. */ + pendingDeliveryClaims?: Array<{ deliveryId: string; claimToken: string }> } export type ConversationMessage = { diff --git a/src/state/file-state-store.test.ts b/src/state/file-state-store.test.ts index d9f594d..ee55ff0 100644 --- a/src/state/file-state-store.test.ts +++ b/src/state/file-state-store.test.ts @@ -107,6 +107,16 @@ describe('FileStateStore', () => { path: '/github/repos/AgentWorkforce/factory/pulls/87/metadata.json', critical: true, pendingKinds: ['checks-failed', 'review-comment'], + resourceSubscription: { + subscriptionId: 'sub-87', + provider: 'github', + resourceRef: '/github/repos/AgentWorkforce__factory/pulls/by-id/87.json', + subscriberId: 'factory-babysitter:uuid-87', + ownerId: 'factory-runtime', + expiresAt: '2026-12-31T00:00:00.000Z', + terminal: true, + }, + pendingDeliveryClaims: [{ deliveryId: 'delivery-87-terminal', claimToken: 'claim-token-87' }], } const first = new FileStateStore({ batchSize: 2, watchStatePath }) await first.setBabysitterSession('workspace-1', 'AR-87:uuid-87:/linear/issues/AR-87__uuid-87.json', session) diff --git a/src/state/file-state-store.ts b/src/state/file-state-store.ts index 750bbad..f770521 100644 --- a/src/state/file-state-store.ts +++ b/src/state/file-state-store.ts @@ -1049,7 +1049,23 @@ const parseBabysitterSessions = (value: Record): Record typeof kind === 'string') + !candidate.pendingKinds.every((kind) => typeof kind === 'string') || + (candidate.pendingDeliveryClaims !== undefined && ( + !Array.isArray(candidate.pendingDeliveryClaims) || + !candidate.pendingDeliveryClaims.every((claim) => + isRecord(claim) && typeof claim.deliveryId === 'string' && typeof claim.claimToken === 'string' + ) + )) || + (candidate.resourceSubscription !== undefined && ( + !isRecord(candidate.resourceSubscription) || + typeof candidate.resourceSubscription.subscriptionId !== 'string' || + typeof candidate.resourceSubscription.provider !== 'string' || + typeof candidate.resourceSubscription.resourceRef !== 'string' || + typeof candidate.resourceSubscription.subscriberId !== 'string' || + typeof candidate.resourceSubscription.ownerId !== 'string' || + typeof candidate.resourceSubscription.expiresAt !== 'string' || + (candidate.resourceSubscription.terminal !== undefined && typeof candidate.resourceSubscription.terminal !== 'boolean') + )) ) { throw new Error('Factory GitHub watch state file is invalid') } @@ -1061,6 +1077,23 @@ const parseBabysitterSessions = (value: Record): Record>).map((claim) => ({ + deliveryId: claim.deliveryId as string, + claimToken: claim.claimToken as string, + })), + }), } } return sessions diff --git a/src/subscriptions/index.ts b/src/subscriptions/index.ts index c5c69a6..806527f 100644 --- a/src/subscriptions/index.ts +++ b/src/subscriptions/index.ts @@ -46,6 +46,20 @@ export { filesystemEventToChangeEvent, integrationRelayFileSyncOptions, } from './event-client' +export { + createResourceSubscriptionsSdkClient, + isResourceSubscriptionsUnavailable, + ResourceSubscriptionsUnavailableError, +} from './resource-subscriptions' +export type { + AcceptedResourceDelivery, + ResourceDeliveryClaim, + ResourceSubscription, + ResourceSubscriptionInput, + ResourceSubscriptionsClient, + ResourceSubscriptionsSdk, + ResourceSubscriptionsSdkClientOptions, +} from './resource-subscriptions' export type { ChangeEvent, FilesystemEventLike, diff --git a/src/subscriptions/resource-subscriptions.ts b/src/subscriptions/resource-subscriptions.ts new file mode 100644 index 0000000..225ab9b --- /dev/null +++ b/src/subscriptions/resource-subscriptions.ts @@ -0,0 +1,191 @@ +import type { + AcceptDurableSubscriptionDeliveryInput, + ClaimDurableSubscriptionDeliveriesInput, + CreateOrRenewDurableResourceSubscriptionInput, + DurableResourceSubscription, + DurableSubscriptionDelivery, + DurableSubscriptionDeliveryListResponse, + DurableSubscriptionDeliveryResponse, +} from '@relayfile/sdk' + +/** + * Consumer-facing contract for Relayfile's durable resource-subscription API. + * + * Relayfile owns matching and delivery-claim persistence. A runtime owns the + * meaning of its opaque IDs and must only accept a claim after it has made the + * corresponding wake durable on its side. + */ +export type ResourceSubscriptionInput = Omit< + CreateOrRenewDurableResourceSubscriptionInput, + 'workspaceId' | 'correlationId' | 'signal' +> + +export type ResourceSubscription = { + subscriptionId: string + /** Derived by Relayfile from the authenticated workspace bearer token. */ + ownerId: string + provider: string + resourceRef: string + eventTypes: string[] + terminalEventTypes?: string[] + subscriberId: string + intent?: string + expiresAt: string +} + +export type ResourceDeliveryClaim = { + deliveryId: string + /** Opaque, short-lived lease credential required to accept this claim. */ + claimToken: string + subscriptionId: string + resourceRef: string + eventType: string + subscriberId: string + ownerId: string + provider: string + terminal: boolean +} + +export type AcceptedResourceDelivery = { + deliveryId: string + subscriptionId: string + terminal: boolean +} + +export interface ResourceSubscriptionsClient { + createOrRenew(workspaceId: string, input: ResourceSubscriptionInput): Promise + /** Atomically lease pending claims for the authenticated owner. */ + claimDeliveryClaims(workspaceId: string, input?: { limit?: number }): Promise + acceptDelivery(workspaceId: string, input: { deliveryId: string; claimToken: string }): Promise + cancel(workspaceId: string, input: { subscriptionId: string }): Promise +} + +/** The canonical Relayfile SDK surface used by Factory. */ +export interface ResourceSubscriptionsSdk { + createOrRenewDurableResourceSubscription( + input: CreateOrRenewDurableResourceSubscriptionInput, + ): Promise + claimDurableSubscriptionDeliveries( + input: ClaimDurableSubscriptionDeliveriesInput, + ): Promise + acceptDurableSubscriptionDelivery( + input: AcceptDurableSubscriptionDeliveryInput, + ): Promise + cancelDurableResourceSubscription( + workspaceId: string, + subscriptionId: string, + options?: { signal?: AbortSignal }, + ): Promise +} + +export type ResourceSubscriptionsSdkClientOptions = { + /** Optional lifecycle cancellation forwarded through the SDK on every call. */ + signal?: AbortSignal +} + +/** The service is absent or too old; callers may retain their legacy route. */ +export class ResourceSubscriptionsUnavailableError extends Error { + constructor(message = 'Relayfile durable resource subscriptions are unavailable') { + super(message) + this.name = 'ResourceSubscriptionsUnavailableError' + } +} + +export const isResourceSubscriptionsUnavailable = (error: unknown): boolean => { + if (error instanceof ResourceSubscriptionsUnavailableError) return true + if (!error || typeof error !== 'object') return false + const record = error as Record + const status = record.status ?? record.statusCode ?? record.httpStatus + return status === 404 +} + +/** + * Adapts Factory's narrow orchestration port to Relayfile's public SDK. + * Authentication, request timeouts, retries, URL construction, and response + * transport remain owned by RelayFileClient. + */ +export function createResourceSubscriptionsSdkClient( + sdk: ResourceSubscriptionsSdk, + options: ResourceSubscriptionsSdkClientOptions = {}, +): ResourceSubscriptionsClient { + return { + async createOrRenew(workspaceId, input) { + const subscription = await sdk.createOrRenewDurableResourceSubscription({ + workspaceId, + ...input, + signal: options.signal, + }) + return toResourceSubscription(subscription) + }, + + async claimDeliveryClaims(workspaceId, input) { + const { deliveries } = await sdk.claimDurableSubscriptionDeliveries({ + workspaceId, + ...(input?.limit === undefined ? {} : { limit: input.limit }), + signal: options.signal, + }) + if (!Array.isArray(deliveries)) { + throw new Error('Relayfile SDK returned an invalid durable-delivery envelope') + } + return deliveries.map(toResourceDeliveryClaim) + }, + + async acceptDelivery(workspaceId, input) { + const { delivery } = await sdk.acceptDurableSubscriptionDelivery({ + workspaceId, + deliveryId: input.deliveryId, + claimToken: input.claimToken, + signal: options.signal, + }) + if (delivery.status !== 'accepted') { + throw new Error(`Relayfile SDK returned delivery ${delivery.id} with non-accepted status ${delivery.status}`) + } + return { + deliveryId: delivery.id, + subscriptionId: delivery.subscriptionId, + terminal: delivery.terminal, + } + }, + + async cancel(workspaceId, input) { + await sdk.cancelDurableResourceSubscription( + workspaceId, + input.subscriptionId, + { signal: options.signal }, + ) + }, + } +} + +const toResourceSubscription = ( + subscription: DurableResourceSubscription, +): ResourceSubscription => ({ + subscriptionId: subscription.id, + ownerId: subscription.ownerId, + subscriberId: subscription.subscriberId, + provider: subscription.provider, + resourceRef: subscription.resourceRef, + eventTypes: subscription.eventTypes, + terminalEventTypes: subscription.terminalEventTypes, + ...(subscription.intent ? { intent: subscription.intent } : {}), + expiresAt: subscription.expiresAt, +}) + +const toResourceDeliveryClaim = ( + delivery: DurableSubscriptionDelivery, +): ResourceDeliveryClaim => { + if (delivery.status !== 'claimed' || !delivery.claimToken) { + throw new Error(`Relayfile SDK returned delivery ${delivery.id} without a live claim`) + } + return { + deliveryId: delivery.id, + claimToken: delivery.claimToken, + subscriptionId: delivery.subscriptionId, + ownerId: delivery.ownerId, + subscriberId: delivery.subscriberId, + provider: delivery.provider, + resourceRef: delivery.resourceRef, + eventType: delivery.event.type, + terminal: delivery.terminal, + } +} diff --git a/src/testing/fakes.ts b/src/testing/fakes.ts index 9a7db15..64eaa06 100644 --- a/src/testing/fakes.ts +++ b/src/testing/fakes.ts @@ -18,6 +18,7 @@ import type { PreviewSweepInput, PreviewSweepResult, } from '../ports' +import type { ResourceSubscriptionsClient } from '../subscriptions' type ExitListener = (name: string, reason?: string) => void type DeliveryFailedListener = (info: { to: string; msgId?: string; reason?: string }) => void @@ -27,6 +28,7 @@ type AgentLifecycleSignalListener = (signal: AgentLifecycleSignal) => void | Pro export class FakeMountClient implements MountClient { readonly writebackTransport = 'test' githubWrite?: GithubConnectionWrite + resourceSubscriptions?: ResourceSubscriptionsClient readonly files = new Map() readonly writes: Array<{ path: string; content: unknown }> = [] readonly deletes: string[] = []