Skip to content

Commit a68eb7a

Browse files
fix(playground,test): withdraw a failed catalog publication before releasing its staging link; bind custom observeProgress
- #persistSnapshot rolls the sidecar back while the staging link still exists when the post-link directory fsync fails, so a concurrent reader keeps seeing an in-progress publication until the path is withdrawn instead of adopting a briefly singly linked file - contractProgressObserver invokes a client's observeProgress method with the client as receiver
1 parent e14d36f commit a68eb7a

4 files changed

Lines changed: 100 additions & 2 deletions

File tree

‎packages/agent-bundle/src/dev/playground/native-playground-service.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1149,10 +1149,19 @@ export class NativePlaygroundService {
11491149
try { await handle.close(); }
11501150
catch (error) { cleanupFailures.push(error); }
11511151
}
1152+
// A failed publication withdraws its sidecar while the staging link still
1153+
// exists: concurrent readers keep seeing an in-progress (doubly linked)
1154+
// publication until the path is gone, never a settled singly linked file
1155+
// that is about to be rolled back.
1156+
if (primary !== undefined && created && publicationIdentity !== undefined) {
1157+
try { await this.#publicationReceipt(path, publicationIdentity, true, true).rollback(); }
1158+
catch (error) { cleanupFailures.push(error); }
1159+
created = false;
1160+
}
11521161
try { await this.#catalogStorage.remove(temporary, { force: true }); }
11531162
catch (error) { cleanupFailures.push(error); }
11541163
}
1155-
if ((primary !== undefined || cleanupFailures.length > 0) && created && publicationIdentity !== undefined) {
1164+
if (primary === undefined && cleanupFailures.length > 0 && created && publicationIdentity !== undefined) {
11561165
try { await this.#publicationReceipt(path, publicationIdentity, true, true).rollback(); }
11571166
catch (error) { cleanupFailures.push(error); }
11581167
}

‎packages/agent-bundle/src/test/contract.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1149,7 +1149,8 @@ const sdkProgressObserver = (
11491149
export const contractProgressObserver = (
11501150
client: ContractMatrixClient,
11511151
): ContractMatrixProgressSource['observeProgress'] => {
1152-
if (typeof client.observeProgress === 'function') return client.observeProgress;
1152+
const { observeProgress } = client;
1153+
if (typeof observeProgress === 'function') return (listener) => observeProgress.call(client, listener);
11531154
const handlers = (client as { readonly _notificationHandlers?: unknown })._notificationHandlers;
11541155
if (handlers instanceof Map) return sdkProgressObserver(handlers as Map<string, ClientNotificationHandler>);
11551156
throw new Error(

‎packages/agent-bundle/tests/dev-contract-runner.test.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -118,6 +118,24 @@ it('exposes live progress notifications through the session trace for lifecycle
118118
expect(seen).toEqual(['token-1']);
119119
});
120120

121+
it('invokes a custom observeProgress method with the client as its receiver', () => {
122+
class InstanceClient {
123+
readonly listeners = new Set<(notification: { readonly params?: { readonly progressToken?: string | number } }) => void>();
124+
observeProgress(listener: (notification: { readonly params?: { readonly progressToken?: string | number } }) => void): () => void {
125+
this.listeners.add(listener);
126+
return () => this.listeners.delete(listener);
127+
}
128+
}
129+
const client = new InstanceClient();
130+
const seen: unknown[] = [];
131+
const stop = contractProgressObserver(client as unknown as ContractMatrixClient)((notification) =>
132+
seen.push(notification.params?.progressToken));
133+
for (const listener of client.listeners) listener({ params: { progressToken: 'bound' } });
134+
stop();
135+
expect(seen).toEqual(['bound']);
136+
expect(client.listeners.size).toBe(0);
137+
});
138+
121139
it('names the missing notification path for a client that is neither an SDK Client nor exposes observeProgress', () => {
122140
const bare = {} as ContractMatrixClient;
123141
expect(() => contractProgressObserver(bare))

‎packages/agent-bundle/tests/native-playground-service.test.ts‎

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1599,6 +1599,76 @@ it('never adopts a staged sidecar that its publisher rolls back, and republishes
15991599
}
16001600
});
16011601

1602+
it('withdraws a sidecar whose directory fsync fails before releasing its staging link, so a waiting reader never adopts it', async () => {
1603+
const root = await mkdtemp(join(tmpdir(), 'agent-bundle-native-playground-fsync-withdrawn-'));
1604+
const catalogDirectory = join(root, 'catalog');
1605+
const reference = epoch('epoch-fsync-withdrawn', join(root, 'artifact'));
1606+
const sidecar = join(catalogDirectory, `${reference.epoch.id}.json`);
1607+
const directoryFailure = new Error('directory sync failed');
1608+
let signalLinked!: () => void;
1609+
const linked = new Promise<void>((resolvePromise) => { signalLinked = resolvePromise; });
1610+
let releaseDirectorySync!: () => void;
1611+
const directorySync = new Promise<void>((resolvePromise) => { releaseDirectorySync = resolvePromise; });
1612+
const winnerStorage: NativePlaygroundCatalogStorage = {
1613+
link: async (source, destination) => {
1614+
await link(source, destination);
1615+
signalLinked();
1616+
},
1617+
mkdir,
1618+
open: async (path, flags, mode) => {
1619+
const handle = await open(path, flags, mode);
1620+
return new Proxy(handle, {
1621+
get(target, property) {
1622+
if (property === 'sync' && String(path) === catalogDirectory) {
1623+
return async () => {
1624+
// The publisher stalls on its directory fsync after the link, then fails it.
1625+
await directorySync;
1626+
throw directoryFailure;
1627+
};
1628+
}
1629+
const value = Reflect.get(target, property, target);
1630+
return typeof value === 'function' ? value.bind(target) : value;
1631+
},
1632+
});
1633+
},
1634+
remove: rm,
1635+
} as NativePlaygroundCatalogStorage;
1636+
const serviceFor = (storage?: NativePlaygroundCatalogStorage): NativePlaygroundService => new NativePlaygroundService({
1637+
catalogDirectory,
1638+
...(storage === undefined ? {} : { catalogStorage: storage }),
1639+
discover: async () => suite(),
1640+
inspectArtifact: async (candidate) => Object.freeze({
1641+
binding: Object.freeze({ manifestPath: 'agent-bundle.manifest.json', source: 'explicit' as const, targetDigests: candidate.epoch.targetDigests }),
1642+
root: candidate.root,
1643+
}),
1644+
planFixture: async () => fixturePlan,
1645+
projectRoot: '/project',
1646+
});
1647+
const winner = serviceFor(winnerStorage);
1648+
const loser = serviceFor();
1649+
try {
1650+
const winning = winner.catalog(reference);
1651+
await linked;
1652+
expect((await stat(sidecar)).nlink).toBe(2);
1653+
const losing = loser.catalog(reference);
1654+
await expect(Promise.race([
1655+
losing.then(() => 'adopted' as const),
1656+
new Promise<'pending'>((resolvePromise) => { setTimeout(() => resolvePromise('pending'), 150); }),
1657+
])).resolves.toBe('pending');
1658+
releaseDirectorySync();
1659+
// The rollback's own directory fsync fails the same way, so both are retained.
1660+
await expect(winning).rejects.toMatchObject({ errors: [directoryFailure, directoryFailure] });
1661+
const adopted = await losing;
1662+
expect((await stat(sidecar)).nlink).toBe(1);
1663+
expect((await readdir(catalogDirectory)).filter((name) => name.includes('.stage-'))).toEqual([]);
1664+
expect(await loser.catalog(reference)).toEqual(adopted);
1665+
await Promise.all([winner.close(), loser.close()]);
1666+
} finally {
1667+
releaseDirectorySync();
1668+
await rm(root, { force: true, recursive: true });
1669+
}
1670+
});
1671+
16021672
it('still rejects a persisted catalog aliased by a hard link that is not an epoch staging file', async () => {
16031673
const root = await mkdtemp(join(tmpdir(), 'agent-bundle-native-playground-aliased-catalog-'));
16041674
const catalogDirectory = join(root, 'catalog');

0 commit comments

Comments
 (0)