Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 26 additions & 7 deletions packages/cli/src/cli/commands/fleet.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1257,8 +1257,20 @@ describe('fleet command support', () => {
workerCwd: '/srv/agent-workforce/cloud/packages/web',
};
const events: string[] = [];
const materializeCloudRelayfileRepository = vi.fn(async () => {
const materializeCloudRelayfileRepository = vi.fn(async (_input, options) => {
events.push('materialize');
options?.onProgress?.({
status: 'queued',
waitedMs: 65_000,
wait: {
reason: 'waiting_on_output_lock',
ownerJobId: 'older-clone-job',
ownerStatus: 'running',
ownerStartedAt: '2026-10-02T10:00:00.000Z',
ownerUpdatedAt: '2026-10-02T10:01:00.000Z',
leaseExpiresAt: '2026-10-02T10:16:00.000Z',
},
});
return {
cloudWorkspaceId: 'cloud-workspace',
repository: 'AgentWorkforce/cloud',
Expand Down Expand Up @@ -1298,6 +1310,7 @@ describe('fleet command support', () => {
release: vi.fn(async () => ({ released: true, deleted: true })),
},
}));
const warn = vi.fn();
const program = new Command();
program.exitOverride();
registerFleetCommands(program, {
Expand All @@ -1322,7 +1335,7 @@ describe('fleet command support', () => {
deleteCloudFleetSandbox: vi.fn(async () => undefined),
createFleetWorkspaceClient: vi.fn() as never,
log: () => undefined,
warn: () => undefined,
warn,
error: () => undefined,
});

Expand All @@ -1347,11 +1360,17 @@ describe('fleet command support', () => {
);

expect(events).toEqual(['materialize', 'ensure', 'spawn']);
expect(materializeCloudRelayfileRepository).toHaveBeenCalledWith({
workspaceId: 'rw_abc',
repository: 'AgentWorkforce/cloud',
revision,
});
expect(materializeCloudRelayfileRepository).toHaveBeenCalledWith(
{
workspaceId: 'rw_abc',
repository: 'AgentWorkforce/cloud',
revision,
},
{ onProgress: expect.any(Function) }
);
expect(warn).toHaveBeenCalledWith(
'Materializing AgentWorkforce/cloud into Relayfile: waiting on output lock held by clone job older-clone-job (running; owner updated 2026-10-02T10:01:00.000Z; lease expires 2026-10-02T10:16:00.000Z; 65s elapsed).'
);
expect(ensureCloudFleetSandbox).toHaveBeenCalledWith(
expect.objectContaining({
mountRelayfile: true,
Expand Down
28 changes: 23 additions & 5 deletions packages/cli/src/cli/commands/fleet.ts
Original file line number Diff line number Diff line change
Expand Up @@ -685,11 +685,29 @@ export function registerFleetCommands(
throw new Error('--workspace-id does not match the captured workspace identity.');
}
if (!checkoutRepository && mountSandboxRelayfile && sandboxRepository) {
liveRepository = await deps.materializeCloudRelayfileRepository({
workspaceId: relayWorkspaceId,
repository: sandboxRepository.repository,
revision: sandboxRepository.revision,
});
const repository = sandboxRepository.repository;
liveRepository = await deps.materializeCloudRelayfileRepository(
{
workspaceId: relayWorkspaceId,
repository,
revision: sandboxRepository.revision,
},
{
onProgress: (progress) => {
const elapsedSeconds = Math.floor(progress.waitedMs / 1_000);
if (progress.wait) {
const ownerState = progress.wait.ownerStatus ?? 'unknown';
deps.warn(
`Materializing ${repository} into Relayfile: waiting on output lock held by clone job ${progress.wait.ownerJobId} (${ownerState}; owner updated ${progress.wait.ownerUpdatedAt ?? 'unknown'}; lease expires ${progress.wait.leaseExpiresAt}; ${elapsedSeconds}s elapsed).`
);
return;
}
deps.warn(
`Materializing ${repository} into Relayfile (${progress.status}; ${elapsedSeconds}s elapsed).`
);
},
}
);
}
const sandboxId = sandboxIdOption ?? (sandboxName === undefined ? `sbx_${randomUUID()}` : undefined);
const deterministicSandboxName =
Expand Down
48 changes: 46 additions & 2 deletions packages/cloud/src/fleet-sandbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -114,11 +114,39 @@ describe('Cloud fleet sandbox client', () => {
.mockResolvedValueOnce({
response: Response.json({
ok: true,
wait: {
reason: 'waiting_on_output_lock',
ownerJobId: 'older-clone-job',
ownerStatus: 'running',
ownerStartedAt: '2026-10-02T10:00:00.000Z',
ownerUpdatedAt: '2026-10-02T10:01:00.000Z',
leaseExpiresAt: '2026-10-02T10:16:00.000Z',
},
job: {
owner: 'AgentWorkforce',
repo: 'cloud',
ref: revision,
status: 'running',
status: 'queued',
},
}),
auth: refreshedAuth,
})
.mockResolvedValueOnce({
response: Response.json({
ok: true,
wait: {
reason: 'waiting_on_output_lock',
ownerJobId: 'older-clone-job',
ownerStatus: 'running',
ownerStartedAt: '2026-10-02T10:00:00.000Z',
ownerUpdatedAt: '2026-10-02T10:01:02.000Z',
leaseExpiresAt: '2026-10-02T10:16:02.000Z',
},
job: {
owner: 'AgentWorkforce',
repo: 'cloud',
ref: revision,
status: 'queued',
},
}),
auth: refreshedAuth,
Expand Down Expand Up @@ -148,14 +176,17 @@ describe('Cloud fleet sandbox client', () => {
auth: refreshedAuth,
});

const onProgress = vi.fn(() => {
throw new Error('progress observer failed');
});
await expect(
materializeCloudRelayfileRepository(
{
workspaceId: 'rw_abc',
repository: 'AgentWorkforce/cloud',
revision,
},
{ pollIntervalMs: 1 }
{ pollIntervalMs: 1, onProgress }
)
).resolves.toEqual({
cloudWorkspaceId: CLOUD_WORKSPACE_ID,
Expand All @@ -178,6 +209,19 @@ describe('Cloud fleet sandbox client', () => {
sourceProfile: 'complete-v1',
});
expect(JSON.stringify(requestCall)).not.toContain('githubToken');
expect(onProgress).toHaveBeenCalledTimes(1);
expect(onProgress).toHaveBeenCalledWith({
status: 'queued',
waitedMs: expect.any(Number),
wait: {
reason: 'waiting_on_output_lock',
ownerJobId: 'older-clone-job',
ownerStatus: 'running',
ownerStartedAt: '2026-10-02T10:00:00.000Z',
ownerUpdatedAt: '2026-10-02T10:01:00.000Z',
leaseExpiresAt: '2026-10-02T10:16:00.000Z',
},
});
});

it('rejects a completed clone that did not produce the requested live Relayfile revision', async () => {
Expand Down
90 changes: 90 additions & 0 deletions packages/cloud/src/fleet-sandbox.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,8 +55,23 @@ export type CloudRelayfileRepositoryMaterialization = {
sentinelPath: string;
};

export type CloudRelayfileRepositoryProgress = {
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
status: 'queued' | 'running' | 'retrying';
waitedMs: number;
wait?: {
reason: 'waiting_on_output_lock';
ownerJobId: string;
ownerStatus: 'queued' | 'running' | 'retrying' | 'completed' | 'failed' | null;
ownerStartedAt: string | null;
ownerUpdatedAt: string | null;
leaseExpiresAt: string;
};
};

export type CloudRelayfileRepositoryMaterializeOptions = CloudFleetSandboxRequestOptions & {
pollIntervalMs?: number;
/** Receives bounded, credential-free status while the Cloud clone remains non-terminal. */
onProgress?: (progress: CloudRelayfileRepositoryProgress) => void;
};

export type CloudFleetSandboxProviderId =
Expand Down Expand Up @@ -468,6 +483,53 @@ function expectedRelayfileRepositoryPaths(
};
}

const CLOUD_CLONE_PROGRESS_STATUSES = new Set(['queued', 'running', 'retrying', 'completed', 'failed']);
const CLOUD_CLONE_PROGRESS_REPEAT_MS = 30_000;

function normalizeProgressTimestamp(value: unknown): string | null | undefined {
if (value === null) return null;
if (typeof value !== 'string' || value.length > 64) return undefined;
const timestamp = Date.parse(value);
return Number.isFinite(timestamp) ? new Date(timestamp).toISOString() : undefined;
}

function readCloudRelayfileRepositoryWait(
payload: JsonRecord
): CloudRelayfileRepositoryProgress['wait'] | undefined {
if (!isObject(payload.wait) || payload.wait.reason !== 'waiting_on_output_lock') return undefined;
const wait = payload.wait;
const ownerJobId = readString(wait, 'ownerJobId');
const ownerStatusValue = wait.ownerStatus;
const ownerStatus =
ownerStatusValue === null
? null
: typeof ownerStatusValue === 'string' && CLOUD_CLONE_PROGRESS_STATUSES.has(ownerStatusValue)
? (ownerStatusValue as NonNullable<CloudRelayfileRepositoryProgress['wait']>['ownerStatus'])
: undefined;
const ownerStartedAt = normalizeProgressTimestamp(wait.ownerStartedAt);
const ownerUpdatedAt = normalizeProgressTimestamp(wait.ownerUpdatedAt);
const leaseExpiresAt = normalizeProgressTimestamp(wait.leaseExpiresAt);
if (
!ownerJobId ||
ownerJobId.length > 128 ||
!/^[A-Za-z0-9_-]+$/.test(ownerJobId) ||
ownerStatus === undefined ||
ownerStartedAt === undefined ||
ownerUpdatedAt === undefined ||
!leaseExpiresAt
) {
return undefined;
}
return {
reason: 'waiting_on_output_lock',
ownerJobId,
ownerStatus,
ownerStartedAt,
ownerUpdatedAt,
leaseExpiresAt,
};
}

function waitForDelay(ms: number, signal: AbortSignal): Promise<void> {
if (signal.aborted) return Promise.reject(signal.reason);
return new Promise((resolve, reject) => {
Expand Down Expand Up @@ -816,6 +878,9 @@ export async function materializeCloudRelayfileRepository(
}
const jobId = requiredString(requestPayload, 'jobId', 'Cloud Relayfile repository materializer');
const expectedPaths = expectedRelayfileRepositoryPaths(owner, repo);
const materializationStartedAt = Date.now();
let lastProgressSignature: string | undefined;
let lastProgressAt = 0;

for (;;) {
const statusResult = await authorizedApiFetch(
Expand Down Expand Up @@ -889,6 +954,31 @@ export async function materializeCloudRelayfileRepository(
if (status !== 'queued' && status !== 'running' && status !== 'retrying') {
throw new Error('Cloud Relayfile repository materialization reported an unknown status.');
}
const now = Date.now();
const wait = readCloudRelayfileRepositoryWait(statusPayload);
const progress: CloudRelayfileRepositoryProgress = {
status,
waitedMs: Math.max(0, now - materializationStartedAt),
...(wait ? { wait } : {}),
};
const progressSignature = JSON.stringify([
progress.status,
progress.wait?.reason,
progress.wait?.ownerJobId,
progress.wait?.ownerStatus,
]);
if (
options.onProgress &&
(progressSignature !== lastProgressSignature || now - lastProgressAt >= CLOUD_CLONE_PROGRESS_REPEAT_MS)
) {
lastProgressSignature = progressSignature;
lastProgressAt = now;
try {
options.onProgress(progress);
} catch {
// Progress is advisory; an observer must not abort a healthy materialization.
}
}
await waitForDelay(pollIntervalMs, signal);
}
}
Expand Down
1 change: 1 addition & 0 deletions packages/cloud/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ export {
type MaterializeCloudRelayfileRepositoryInput,
type CloudRelayfileRepositoryMaterialization,
type CloudRelayfileRepositoryMaterializeOptions,
type CloudRelayfileRepositoryProgress,
type CloudFleetSandboxReady,
type CloudFleetSandboxReused,
type CloudFleetSandboxProvisioningTimeout,
Expand Down
Loading