From 6322069f1a02f16afe3423c567660637e5ad7f86 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Pierzcha=C5=82a?= Date: Sat, 19 Sep 2026 14:18:05 +0200 Subject: [PATCH] refactor: share the bounded app-log poller in capture-kit (#2616) --- packages/capture-kit/package.json | 4 + .../src/app-log-polling.fixtures.ts | 87 +++++ .../capture-kit/src/app-log-polling.test.ts | 301 ++++++++++++++++++ packages/capture-kit/src/app-log-polling.ts | 209 ++++++++++++ .../src/app-log-poller.test.ts | 201 ++++-------- .../provider-limrun/src/app-log-poller.ts | 182 +---------- .../src/request-cancellation.ts | 2 +- scripts/layering/package-boundaries.test.ts | 1 + 8 files changed, 682 insertions(+), 305 deletions(-) create mode 100644 packages/capture-kit/src/app-log-polling.fixtures.ts create mode 100644 packages/capture-kit/src/app-log-polling.test.ts create mode 100644 packages/capture-kit/src/app-log-polling.ts diff --git a/packages/capture-kit/package.json b/packages/capture-kit/package.json index de728bbb4b..8d61bf2d0c 100644 --- a/packages/capture-kit/package.json +++ b/packages/capture-kit/package.json @@ -18,6 +18,10 @@ "types": "./src/snapshot/android-replacement-surface-occlusion.ts", "default": "./src/snapshot/android-replacement-surface-occlusion.ts" }, + "./app-log-polling": { + "types": "./src/app-log-polling.ts", + "default": "./src/app-log-polling.ts" + }, "./audio-probe-admission-ledger": { "types": "./src/capture-admission/audio-probe-admission-ledger.ts", "default": "./src/capture-admission/audio-probe-admission-ledger.ts" diff --git a/packages/capture-kit/src/app-log-polling.fixtures.ts b/packages/capture-kit/src/app-log-polling.fixtures.ts new file mode 100644 index 0000000000..7f439c962f --- /dev/null +++ b/packages/capture-kit/src/app-log-polling.fixtures.ts @@ -0,0 +1,87 @@ +import type { AppLogRuntimeHost } from '@agent-device/contracts/app-log-runtime'; + +export type DeferredSleeps = Readonly<{ + wait(milliseconds: number): Promise; + resolveNext(milliseconds: number): void; + resolveAny(): void; + hasPending(milliseconds: number): boolean; +}>; + +export function createDeferredSleeps(): DeferredSleeps { + const pending: Array<{ milliseconds: number; resolve: () => void }> = []; + return { + wait: async (milliseconds: number) => + await new Promise((resolve) => pending.push({ milliseconds, resolve })), + resolveNext: (milliseconds: number) => { + const index = pending.findIndex((entry) => entry.milliseconds === milliseconds); + if (index < 0) throw new Error(`No ${milliseconds}ms sleep is pending`); + pending.splice(index, 1)[0]?.resolve(); + }, + resolveAny: () => pending.shift()?.resolve(), + hasPending: (milliseconds: number) => + pending.some((entry) => entry.milliseconds === milliseconds), + }; +} + +export function createPollerHost(options: { + existingTail: string; + writes: string[]; + sleeps: DeferredSleeps; + onOutputDispose?: () => void; + outputDisposeError?: boolean; + failure?: 'readTail' | 'openAppend'; +}): AppLogRuntimeHost { + return { + appleTools: { + isXcrunAvailable: async () => false, + run: async () => { + throw new Error('unused'); + }, + }, + toolchains: { prepare: async () => undefined }, + artifacts: { + resolveSession: () => ({ + outputPath: '/sessions/one/app.log', + pidPath: '/sessions/one/app-log.pid', + }), + }, + commands: { + which: async () => undefined, + run: async () => ({ stdout: '', stderr: '', exitCode: 0 }), + }, + outputs: { + readTail: async () => { + if (options.failure === 'readTail') throw new Error('tail failed'); + return options.existingTail; + }, + openAppend: async () => { + if (options.failure === 'openAppend') throw new Error('open failed'); + return { + write: async (chunk) => { + options.writes.push(String(chunk)); + }, + [Symbol.asyncDispose]: async () => { + if (options.outputDisposeError) throw new Error('output cleanup failed'); + options.onOutputDispose?.(); + }, + }; + }, + }, + processTransports: { + resolve: async () => ({ mode: 'local' }), + }, + processes: { + start: async () => { + throw new Error('unused'); + }, + readMarker: async () => ({ status: 'missing' }), + clearMarker: async () => {}, + inspect: async () => 'missing', + terminate: async () => 'already-missing', + }, + clock: { + now: () => 100, + sleep: async (milliseconds) => await options.sleeps.wait(milliseconds), + }, + }; +} diff --git a/packages/capture-kit/src/app-log-polling.test.ts b/packages/capture-kit/src/app-log-polling.test.ts new file mode 100644 index 0000000000..e796b96e7b --- /dev/null +++ b/packages/capture-kit/src/app-log-polling.test.ts @@ -0,0 +1,301 @@ +import type { AppLogLiveHandle, AppLogRuntimeHost } from '@agent-device/contracts/app-log-runtime'; +import { describe, expect, test, vi } from 'vitest'; +import { + startAppLogPoller, + type AppLogPollerInput, + type AppLogPollerReader, +} from './app-log-polling.ts'; +import { createDeferredSleeps, createPollerHost } from './app-log-polling.fixtures.ts'; + +const CLEANUP_MESSAGE = 'capture cleanup did not settle every owned resource'; + +function reader(overrides: Partial = {}): AppLogPollerReader { + return { + readLogs: async () => '', + [Symbol.asyncDispose]: async () => {}, + ...overrides, + }; +} + +async function start(options: { + host: AppLogRuntimeHost; + reader: AppLogPollerReader; +}): Promise { + return await startAppLogPoller(pollerInput({ host: options.host, reader: options.reader })); +} + +describe('app-log poller', () => { + test('deduplicates the persisted tail, echoes backend identity, and stops before disposing', async () => { + const sleeps = createDeferredSleeps(); + const writes: string[] = []; + const disposals: string[] = []; + const readLogs = vi.fn( + async (_appBundleId: string, _lineLimit: number) => 'old line\nshared\nnew line\n', + ); + const logReader = reader({ + readLogs, + [Symbol.asyncDispose]: async () => { + disposals.push('reader'); + }, + }); + const handle = await start({ + host: createPollerHost({ + existingTail: 'old line\nshared\n[agent-device][mark][time] checkpoint\n', + writes, + sleeps, + onOutputDispose: () => { + disposals.push('output'); + }, + }), + reader: logReader, + }); + await vi.waitFor(() => expect(writes).toEqual(['new line\n'])); + expect(handle.inspect().backend).toBe('ios-simulator'); + expect(handle.inspect().state).toBe('active'); + + const finishing = handle.finish(); + expect(disposals).toEqual([]); + sleeps.resolveNext(1_000); + await finishing; + expect(readLogs).toHaveBeenCalledTimes(1); + expect(disposals).toEqual(['reader', 'output']); + expect(handle.inspect().state).toBe('ended'); + }); + + test('streams only the appended delta across reads and normalizes the trailing newline', async () => { + const sleeps = createDeferredSleeps(); + const writes: string[] = []; + const tail = ['partial', 'partial done\n', 'partial done\n']; + let read = 0; + const handle = await start({ + host: createPollerHost({ existingTail: '', writes, sleeps }), + reader: reader({ readLogs: async () => tail[read++] ?? '' }), + }); + await vi.waitFor(() => expect(writes).toEqual(['partial\n'])); + sleeps.resolveNext(1_000); + await vi.waitFor(() => expect(writes).toEqual(['partial\n', ' done\n'])); + sleeps.resolveNext(1_000); + await vi.waitFor(() => expect(sleeps.hasPending(1_000)).toBe(true)); + expect(writes).toEqual(['partial\n', ' done\n']); + expect(handle.inspect().state).toBe('active'); + const finishing = handle.finish(); + sleeps.resolveNext(1_000); + await finishing; + expect(handle.inspect().state).toBe('ended'); + }); + + test.each(['readTail', 'openAppend'] as const)( + 'rolls back the reader when %s fails during acquisition', + async (failure) => { + const dispose = vi.fn(async () => {}); + await expect( + startAppLogPoller( + pollerInput({ + host: createPollerHost({ + existingTail: '', + writes: [], + sleeps: createDeferredSleeps(), + failure, + }), + reader: reader({ [Symbol.asyncDispose]: dispose }), + }), + ), + ).rejects.toThrow(`${failure === 'readTail' ? 'tail' : 'open'} failed`); + expect(dispose).toHaveBeenCalledOnce(); + }, + ); + + test('uses linear overlap matching for a near-limit tail without overlap', async () => { + const sleeps = createDeferredSleeps(); + const writes: string[] = []; + const handle = await start({ + host: createPollerHost({ existingTail: `${'a'.repeat(240_000)}\n`, writes, sleeps }), + reader: reader({ readLogs: async () => `${'b'.repeat(240_000)}\n` }), + }); + await vi.waitFor(() => expect(writes[0]?.length).toBe(240_001)); + const finishing = handle.finish(); + sleeps.resolveNext(1_000); + await finishing; + }); + + test('reports the input cleanup message when reader disposal rejects, still disposing the output', async () => { + const sleeps = createDeferredSleeps(); + let outputDisposed = false; + const handle = await start({ + host: createPollerHost({ + existingTail: '', + writes: [], + sleeps, + onOutputDispose: () => { + outputDisposed = true; + }, + }), + reader: reader({ + [Symbol.asyncDispose]: async () => { + throw new Error('reader cleanup failed'); + }, + }), + }); + await vi.waitFor(() => expect(sleeps.hasPending(1_000)).toBe(true)); + const finishing = handle.finish(); + sleeps.resolveNext(1_000); + await expect(finishing).resolves.toMatchObject({ + status: 'cleanup-pending', + message: CLEANUP_MESSAGE, + }); + expect(outputDisposed).toBe(true); + }); + + test('reports the cleanup message when output disposal rejects', async () => { + const sleeps = createDeferredSleeps(); + let readerDisposed = false; + const handle = await start({ + host: createPollerHost({ + existingTail: '', + writes: [], + sleeps, + outputDisposeError: true, + }), + reader: reader({ + [Symbol.asyncDispose]: async () => { + readerDisposed = true; + }, + }), + }); + await vi.waitFor(() => expect(sleeps.hasPending(1_000)).toBe(true)); + const finishing = handle.finish(); + sleeps.resolveNext(1_000); + await expect(finishing).resolves.toMatchObject({ + status: 'cleanup-pending', + message: CLEANUP_MESSAGE, + }); + expect(readerDisposed).toBe(true); + }); + + test('fails a bounded read that never settles, then disposes', async () => { + const sleeps = createDeferredSleeps(); + let readerDisposed = false; + const readLogs = vi.fn(async (_appId: string, _limit: number) => { + return await new Promise(() => {}); + }); + const handle = await start({ + host: createPollerHost({ existingTail: '', writes: [], sleeps }), + reader: reader({ + readLogs, + [Symbol.asyncDispose]: async () => { + readerDisposed = true; + }, + }), + }); + await vi.waitFor(() => expect(sleeps.hasPending(5_000)).toBe(true)); + sleeps.resolveNext(5_000); + await vi.waitFor(() => expect(handle.inspect().state).toBe('failed')); + expect(readLogs).toHaveBeenCalledTimes(1); + expect(readerDisposed).toBe(false); + await expect(handle.finish()).resolves.toMatchObject({ status: 'completed' }); + expect(readerDisposed).toBe(true); + }); + + test('recovers after a read rejection and resumes writing on the next read', async () => { + const sleeps = createDeferredSleeps(); + const writes: string[] = []; + let call = 0; + const handle = await start({ + host: createPollerHost({ existingTail: '', writes, sleeps }), + reader: reader({ + readLogs: async () => { + call += 1; + if (call === 1) throw new Error('transient read failure'); + return 'recovered line\n'; + }, + }), + }); + await vi.waitFor(() => expect(sleeps.hasPending(1_000)).toBe(true)); + expect(handle.inspect().state).toBe('recovering'); + expect(writes).toEqual([]); + sleeps.resolveNext(1_000); + await vi.waitFor(() => expect(writes).toEqual(['recovered line\n'])); + expect(handle.inspect().state).toBe('active'); + const finishing = handle.finish(); + sleeps.resolveNext(1_000); + await finishing; + }); + + test('stops before writing when finish begins during an in-flight read', async () => { + const sleeps = createDeferredSleeps(); + const writes: string[] = []; + let readerDisposed = false; + let settleRead!: (text: string) => void; + const readLogs = vi.fn( + () => + new Promise((resolve) => { + settleRead = resolve; + }), + ); + const handle = await start({ + host: createPollerHost({ existingTail: '', writes, sleeps }), + reader: reader({ + readLogs, + [Symbol.asyncDispose]: async () => { + readerDisposed = true; + }, + }), + }); + await vi.waitFor(() => expect(readLogs).toHaveBeenCalledTimes(1)); + const finishing = handle.finish(); + settleRead('late line\n'); + for (let index = 0; index < 4; index += 1) { + await Promise.resolve(); + sleeps.resolveAny(); + } + await finishing; + expect(writes).toEqual([]); + expect(readerDisposed).toBe(true); + }); + + test('finishes idempotently and disposes each owned resource once', async () => { + const sleeps = createDeferredSleeps(); + const writes: string[] = []; + let readerDisposes = 0; + let outputDisposes = 0; + const handle = await start({ + host: createPollerHost({ + existingTail: '', + writes, + sleeps, + onOutputDispose: () => { + outputDisposes += 1; + }, + }), + reader: reader({ + readLogs: async () => 'one line\n', + [Symbol.asyncDispose]: async () => { + readerDisposes += 1; + }, + }), + }); + await vi.waitFor(() => expect(writes).toEqual(['one line\n'])); + const first = handle.finish(); + const second = handle.finish(); + sleeps.resolveNext(1_000); + await expect(Promise.all([first, second])).resolves.toEqual([ + expect.objectContaining({ status: 'completed' }), + expect.objectContaining({ status: 'completed' }), + ]); + expect(readerDisposes).toBe(1); + expect(outputDisposes).toBe(1); + }); +}); + +function pollerInput( + input: Readonly<{ host: AppLogRuntimeHost; reader: AppLogPollerReader }>, +): AppLogPollerInput { + return { + host: input.host, + reader: input.reader, + backend: 'ios-simulator', + appBundleId: 'com.example.app', + outputPath: '/sessions/one/app.log', + cleanupFailureMessage: CLEANUP_MESSAGE, + }; +} diff --git a/packages/capture-kit/src/app-log-polling.ts b/packages/capture-kit/src/app-log-polling.ts new file mode 100644 index 0000000000..4427e4404d --- /dev/null +++ b/packages/capture-kit/src/app-log-polling.ts @@ -0,0 +1,209 @@ +import type { + AppLogLiveHandle, + AppLogLiveSnapshot, + AppLogOutputSink, + AppLogRuntimeHost, +} from '@agent-device/contracts/app-log-runtime'; +import type { FinishOutcome } from '@agent-device/contracts/durable-resource'; +import type { LogBackend } from '@agent-device/contracts/observability'; +import { AsyncCleanupStack } from '@agent-device/contracts/async-lifecycle'; +import { createAppLogLiveHandleFromFinish } from './app-log-live-handle.ts'; + +const READ_TIMEOUT_MS = 5_000; +const POLL_INTERVAL_MS = 1_000; +const READ_LINE_LIMIT = 1_000; +const TAIL_READ_BYTES = 256 * 1024; +const MARK_PREFIX = '[agent-device][mark]'; +const READ_ABORT_MESSAGE = 'App-log provider read aborted'; + +/** + * Provider-supplied app-log source. A read cannot be cancelled, so the poller settles a bounded + * read on its own timeout and never proves that the provider request ended. + */ +export type AppLogPollerReader = AsyncDisposable & + Readonly<{ + readLogs(appBundleId: string, lineLimit: number): Promise; + }>; + +export type AppLogPollerInput = Readonly<{ + host: AppLogRuntimeHost; + reader: AppLogPollerReader; + backend: LogBackend; + appBundleId: string; + outputPath: string; + cleanupFailureMessage: string; +}>; + +export async function startAppLogPoller(input: AppLogPollerInput): Promise { + const rollback = new AsyncCleanupStack(); + let adopted = false; + rollback.defer(async () => { + if (!adopted) await input.reader[Symbol.asyncDispose](); + }); + try { + const existingTail = await input.host.outputs.readTail(input.outputPath, TAIL_READ_BYTES); + const output = await input.host.outputs.openAppend(input.outputPath); + rollback.defer(async () => { + if (!adopted) await output[Symbol.asyncDispose](); + }); + const handle = createPollerHandle(input, output, remoteOnlyTail(existingTail)); + adopted = true; + return handle; + } finally { + await rollback[Symbol.asyncDispose](); + } +} + +function createPollerHandle( + input: AppLogPollerInput, + output: AppLogOutputSink, + existingTail: string, +): AppLogLiveHandle { + const { backend } = input; + const startedAt = input.host.clock.now(); + let state: AppLogLiveSnapshot['state'] = 'active'; + let stopped = false; + let previous = existingTail; + const polling = (async () => { + while (!stopped) { + try { + const read = await boundedRead(input); + if (read.status === 'timeout') { + state = 'failed'; + return; + } + if (stopped) return; + const delta = appendedTail(previous, read.text); + previous = read.text; + if (delta) await output.write(delta.endsWith('\n') ? delta : `${delta}\n`); + state = 'active'; + } catch { + if (stopped) return; + state = 'recovering'; + } + await input.host.clock.sleep(POLL_INTERVAL_MS); + } + })(); + let finishPromise: + | Promise> + | undefined; + const finish = async () => + (finishPromise ??= (async () => { + stopped = true; + await polling; + const failures = await disposeAll([input.reader, output]); + if (failures.length > 0) { + state = 'failed'; + return { + status: 'cleanup-pending', + reason: 'transport-failed', + message: input.cleanupFailureMessage, + } as const; + } + state = 'ended'; + return { + status: 'completed', + result: { + backend, + outputPath: input.outputPath, + completedAt: input.host.clock.now(), + }, + } as const; + })()); + return createAppLogLiveHandleFromFinish({ + inspect: () => ({ backend, state, startedAt }), + finish, + }); +} + +async function boundedRead( + input: AppLogPollerInput, +): Promise | Readonly<{ status: 'timeout' }>> { + const controller = new AbortController(); + const read = settleOnAbort( + input.reader.readLogs(input.appBundleId, READ_LINE_LIMIT), + controller.signal, + ).then((text) => ({ status: 'read' as const, text })); + try { + const result = await Promise.race([ + read, + input.host.clock + .sleep(READ_TIMEOUT_MS, controller.signal) + .then(() => ({ status: 'timeout' as const })), + ]); + if (result.status === 'timeout') { + controller.abort(); + await read.catch(() => undefined); + } + return result; + } finally { + controller.abort(); + } +} + +/** A reader that cannot cancel its request must still let the bounded read settle on abort. */ +async function settleOnAbort(source: Promise, signal: AbortSignal): Promise { + return await new Promise((resolve, reject) => { + const aborted = () => reject(signal.reason ?? new Error(READ_ABORT_MESSAGE)); + if (signal.aborted) { + void source.catch(() => undefined); + aborted(); + return; + } + signal.addEventListener('abort', aborted, { once: true }); + void source.then( + (value) => { + signal.removeEventListener('abort', aborted); + resolve(value); + }, + (error: unknown) => { + signal.removeEventListener('abort', aborted); + reject(error); + }, + ); + }); +} + +/** The new tail minus its overlap with the previous one (longest suffix/prefix match). */ +function appendedTail(previous: string, current: string): string { + if (!previous || !current) return current; + const prefix = buildPrefixTable(current); + const maximum = Math.min(previous.length, current.length); + const suffix = previous.slice(previous.length - maximum); + return current.slice(suffixPrefixOverlap(suffix, current, prefix)); +} + +function buildPrefixTable(text: string): Uint32Array { + const prefix = new Uint32Array(text.length); + let matched = 0; + for (let index = 1; index < text.length; index += 1) { + while (matched > 0 && text[index] !== text[matched]) matched = prefix[matched - 1]!; + if (text[index] === text[matched]) matched += 1; + prefix[index] = matched; + } + return prefix; +} + +function suffixPrefixOverlap(suffix: string, current: string, prefix: Uint32Array): number { + let matched = 0; + for (let index = 0; index < suffix.length; index += 1) { + while (matched > 0 && suffix[index] !== current[matched]) matched = prefix[matched - 1]!; + if (suffix[index] === current[matched]) matched += 1; + if (matched === current.length && index < suffix.length - 1) matched = prefix[matched - 1]!; + } + return matched; +} + +function remoteOnlyTail(tail: string): string { + return tail + .split('\n') + .filter((line) => !line.startsWith(MARK_PREFIX)) + .join('\n'); +} + +async function disposeAll(resources: readonly AsyncDisposable[]): Promise { + const results = await Promise.allSettled( + resources.map(async (resource) => await resource[Symbol.asyncDispose]()), + ); + return results.flatMap((result) => (result.status === 'rejected' ? [result.reason] : [])); +} diff --git a/packages/provider-limrun/src/app-log-poller.test.ts b/packages/provider-limrun/src/app-log-poller.test.ts index 5e87015eff..3b213de6c0 100644 --- a/packages/provider-limrun/src/app-log-poller.test.ts +++ b/packages/provider-limrun/src/app-log-poller.test.ts @@ -2,152 +2,89 @@ import type { AppLogRuntimeHost } from '@agent-device/contracts/app-log-runtime' import { describe, expect, test, vi } from 'vitest'; import { startLimrunAppLogPoller, type LimrunAppLogReader } from './app-log-poller.ts'; -describe('Limrun app-log poller', () => { - test('deduplicates the persisted tail and stops before disposing its resources', async () => { - const sleeps = deferredSleeps(); - const writes: string[] = []; - let outputDisposed = false; - let readerDisposed = false; - const reader: LimrunAppLogReader = { - platform: 'ios', - leaseId: 'lease-1', - instanceId: 'instance-1', - readLogs: vi.fn(async () => 'old line\nshared\nnew line\n'), - [Symbol.asyncDispose]: async () => { - readerDisposed = true; - }, - }; - const handle = await startLimrunAppLogPoller({ - host: pollerHost({ - existingTail: 'old line\nshared\n[agent-device][mark][time] checkpoint\n', - writes, - sleeps, - onOutputDispose: () => { - outputDisposed = true; - }, - }), - reader, - appBundleId: 'com.example.app', - outputPath: '/sessions/one/app.log', - }); - await vi.waitFor(() => expect(writes).toEqual(['new line\n'])); - - const finishing = handle.finish(); - expect(readerDisposed).toBe(false); - expect(outputDisposed).toBe(false); - sleeps.resolveNext(1_000); - await finishing; - expect(reader.readLogs).toHaveBeenCalledTimes(1); - expect(readerDisposed).toBe(true); - expect(outputDisposed).toBe(true); - }); +function reader( + platform: LimrunAppLogReader['platform'], + overrides: Partial = {}, +): LimrunAppLogReader { + return { + platform, + leaseId: 'lease-1', + instanceId: 'instance-1', + readLogs: async () => 'one line\n', + [Symbol.asyncDispose]: async () => {}, + ...overrides, + }; +} - test.each(['readTail', 'openAppend'] as const)( - 'rolls back the reader when %s fails during acquisition', - async (failure) => { - const dispose = vi.fn(async () => {}); - const reader: LimrunAppLogReader = { - platform: 'ios', - leaseId: 'lease-1', - instanceId: 'instance-1', - readLogs: async () => '', - [Symbol.asyncDispose]: dispose, - }; - const host = pollerHost({ - existingTail: '', - writes: [], - sleeps: deferredSleeps(), - failure, +describe('Limrun app-log poller', () => { + test.each([ + ['ios', 'ios-simulator'], + ['android', 'android'], + ] as const)( + 'derives the %s backend identity independently of the shared poller', + async (platform, backend) => { + const sleeps = deferredSleeps(); + const writes: string[] = []; + const handle = await startLimrunAppLogPoller({ + host: pollerHost({ existingTail: '', writes, sleeps }), + reader: reader(platform), + appBundleId: 'com.example.app', + outputPath: '/sessions/one/app.log', + }); + await vi.waitFor(() => expect(writes).toEqual(['one line\n'])); + expect(handle.inspect().backend).toBe(backend); + const finishing = handle.finish(); + sleeps.resolveNext(1_000); + await expect(finishing).resolves.toMatchObject({ + status: 'completed', + result: { backend }, }); - await expect( - startLimrunAppLogPoller({ - host, - reader, - appBundleId: 'com.example.app', - outputPath: '/sessions/one/app.log', - }), - ).rejects.toThrow(`${failure === 'readTail' ? 'tail' : 'open'} failed`); - expect(dispose).toHaveBeenCalledOnce(); }, ); - test('uses linear overlap matching for a near-limit tail without overlap', async () => { + test('keeps the Limrun cleanup wording when disposal fails', async () => { const sleeps = deferredSleeps(); - const writes: string[] = []; - const reader: LimrunAppLogReader = { - platform: 'ios', - leaseId: 'lease-1', - instanceId: 'instance-1', - readLogs: async () => `${'b'.repeat(240_000)}\n`, - [Symbol.asyncDispose]: async () => {}, - }; const handle = await startLimrunAppLogPoller({ - host: pollerHost({ existingTail: `${'a'.repeat(240_000)}\n`, writes, sleeps }), - reader, - appBundleId: 'com.example.app', - outputPath: '/sessions/one/app.log', - }); - await vi.waitFor(() => expect(writes[0]?.length).toBe(240_001)); - const finishing = handle.finish(); - sleeps.resolveNext(1_000); - await finishing; - }); - - test('still disposes the output when reader cleanup rejects', async () => { - const sleeps = deferredSleeps(); - let outputDisposed = false; - const reader: LimrunAppLogReader = { - platform: 'ios', - leaseId: 'lease-1', - instanceId: 'instance-1', - readLogs: async () => '', - [Symbol.asyncDispose]: async () => { - throw new Error('reader cleanup failed'); - }, - }; - const handle = await startLimrunAppLogPoller({ - host: pollerHost({ - existingTail: '', - writes: [], - sleeps, - onOutputDispose: () => { - outputDisposed = true; + host: pollerHost({ existingTail: '', writes: [], sleeps }), + reader: reader('ios', { + [Symbol.asyncDispose]: async () => { + throw new Error('reader cleanup failed'); }, }), - reader, appBundleId: 'com.example.app', outputPath: '/sessions/one/app.log', }); await vi.waitFor(() => expect(sleeps.hasPending(1_000)).toBe(true)); const finishing = handle.finish(); sleeps.resolveNext(1_000); - await expect(finishing).resolves.toMatchObject({ status: 'cleanup-pending' }); - expect(outputDisposed).toBe(true); + await expect(finishing).resolves.toMatchObject({ + status: 'cleanup-pending', + message: 'Limrun app-log cleanup did not settle every owned resource', + }); }); - test('allows one bounded read only and settles it before disposal', async () => { + test('settles an uncancelable read whose reader ignores the poller signal', async () => { const sleeps = deferredSleeps(); let readerDisposed = false; - const reader: LimrunAppLogReader = { - platform: 'android', - leaseId: 'lease-1', - instanceId: 'instance-1', - readLogs: vi.fn(async () => await new Promise(() => {})), - [Symbol.asyncDispose]: async () => { - readerDisposed = true; - }, - }; + const readLogs = vi.fn(async (_appBundleId: string, _lineLimit: number) => { + return await new Promise(() => {}); + }); const handle = await startLimrunAppLogPoller({ host: pollerHost({ existingTail: '', writes: [], sleeps }), - reader, + reader: reader('android', { + readLogs, + [Symbol.asyncDispose]: async () => { + readerDisposed = true; + }, + }), appBundleId: 'com.example.app', outputPath: '/sessions/one/app.log', }); - const finishing = handle.finish(); - expect(readerDisposed).toBe(false); + await vi.waitFor(() => expect(sleeps.hasPending(5_000)).toBe(true)); sleeps.resolveNext(5_000); - await finishing; - expect(reader.readLogs).toHaveBeenCalledTimes(1); + await vi.waitFor(() => expect(handle.inspect().state).toBe('failed')); + expect(readLogs).toHaveBeenCalledTimes(1); + await expect(handle.finish()).resolves.toMatchObject({ status: 'completed' }); expect(readerDisposed).toBe(true); }); }); @@ -156,8 +93,6 @@ function pollerHost(options: { existingTail: string; writes: string[]; sleeps: ReturnType; - onOutputDispose?: () => void; - failure?: 'readTail' | 'openAppend'; }): AppLogRuntimeHost { return { appleTools: { @@ -178,19 +113,13 @@ function pollerHost(options: { run: async () => ({ stdout: '', stderr: '', exitCode: 0 }), }, outputs: { - readTail: async () => { - if (options.failure === 'readTail') throw new Error('tail failed'); - return options.existingTail; - }, - openAppend: async () => { - if (options.failure === 'openAppend') throw new Error('open failed'); - return { - write: async (chunk) => { - options.writes.push(String(chunk)); - }, - [Symbol.asyncDispose]: async () => options.onOutputDispose?.(), - }; - }, + readTail: async () => options.existingTail, + openAppend: async () => ({ + write: async (chunk) => { + options.writes.push(String(chunk)); + }, + [Symbol.asyncDispose]: async () => {}, + }), }, processTransports: { resolve: async () => ({ mode: 'local' }), diff --git a/packages/provider-limrun/src/app-log-poller.ts b/packages/provider-limrun/src/app-log-poller.ts index 5a4f9a9257..e05b67b37c 100644 --- a/packages/provider-limrun/src/app-log-poller.ts +++ b/packages/provider-limrun/src/app-log-poller.ts @@ -1,14 +1,8 @@ -import type { - AppLogLiveHandle, - AppLogLiveSnapshot, - AppLogOutputSink, - AppLogRuntimeHost, -} from '@agent-device/contracts/app-log-runtime'; -import type { FinishOutcome } from '@agent-device/contracts/durable-resource'; -import { AsyncCleanupStack } from '@agent-device/contracts/async-lifecycle'; -import { createAppLogLiveHandleFromFinish } from '@agent-device/capture-kit'; -import type { LogBackend } from '@agent-device/contracts/observability'; -import { awaitLimrunOperation } from './request-cancellation.ts'; +import type { AppLogLiveHandle, AppLogRuntimeHost } from '@agent-device/contracts/app-log-runtime'; +import { + startAppLogPoller, + type AppLogPollerReader, +} from '@agent-device/capture-kit/app-log-polling'; export type LimrunAppLogReader = AsyncDisposable & Readonly<{ @@ -18,170 +12,22 @@ export type LimrunAppLogReader = AsyncDisposable & readLogs(appBundleId: string, lineLimit: number): Promise; }>; +const LIMRUN_CLEANUP_FAILURE_MESSAGE = 'Limrun app-log cleanup did not settle every owned resource'; + export async function startLimrunAppLogPoller(options: { host: AppLogRuntimeHost; reader: LimrunAppLogReader; appBundleId: string; outputPath: string; }): Promise { - const rollback = new AsyncCleanupStack(); - let adopted = false; - rollback.defer(async () => { - if (!adopted) await options.reader[Symbol.asyncDispose](); + return await startAppLogPoller({ + host: options.host, + reader: options.reader satisfies AppLogPollerReader, + backend: backendForReader(options.reader), + appBundleId: options.appBundleId, + outputPath: options.outputPath, + cleanupFailureMessage: LIMRUN_CLEANUP_FAILURE_MESSAGE, }); - try { - const existingTail = await options.host.outputs.readTail(options.outputPath, 256 * 1024); - const output = await options.host.outputs.openAppend(options.outputPath); - rollback.defer(async () => { - if (!adopted) await output[Symbol.asyncDispose](); - }); - const handle = createPollerHandle(options, output, remoteOnlyTail(existingTail)); - adopted = true; - return handle; - } finally { - await rollback[Symbol.asyncDispose](); - } -} - -function createPollerHandle( - options: { - host: AppLogRuntimeHost; - reader: LimrunAppLogReader; - appBundleId: string; - outputPath: string; - }, - output: AppLogOutputSink, - existingTail: string, -): AppLogLiveHandle { - const backend = backendForReader(options.reader); - const startedAt = options.host.clock.now(); - let state: AppLogLiveSnapshot['state'] = 'active'; - let stopped = false; - let previous = existingTail; - const polling = (async () => { - while (!stopped) { - try { - const read = await boundedRead(options); - if (read.status === 'timeout') { - state = 'failed'; - return; - } - if (stopped) return; - const delta = appendedTail(previous, read.text); - previous = read.text; - if (delta) await output.write(delta.endsWith('\n') ? delta : `${delta}\n`); - state = 'active'; - } catch { - if (stopped) return; - state = 'recovering'; - } - await options.host.clock.sleep(1_000); - } - })(); - let finishPromise: - | Promise> - | undefined; - const finish = async () => - (finishPromise ??= (async () => { - stopped = true; - await polling; - const failures = await disposeAll([options.reader, output]); - if (failures.length > 0) { - state = 'failed'; - return { - status: 'cleanup-pending', - reason: 'transport-failed', - message: 'Limrun app-log cleanup did not settle every owned resource', - } as const; - } - state = 'ended'; - return { - status: 'completed', - result: { - backend, - outputPath: options.outputPath, - completedAt: options.host.clock.now(), - }, - } as const; - })()); - return createAppLogLiveHandleFromFinish({ - inspect: () => ({ backend, state, startedAt }), - finish, - }); -} - -async function boundedRead(options: { - host: AppLogRuntimeHost; - reader: LimrunAppLogReader; - appBundleId: string; -}): Promise | Readonly<{ status: 'timeout' }>> { - const controller = new AbortController(); - const read = awaitLimrunOperation( - options.reader.readLogs(options.appBundleId, 1_000), - controller.signal, - 'App-log provider read aborted', - ).then((text) => ({ status: 'read' as const, text })); - try { - const result = await Promise.race([ - read, - options.host.clock - .sleep(5_000, controller.signal) - .then(() => ({ status: 'timeout' as const })), - ]); - if (result.status === 'timeout') { - controller.abort(); - await read.catch(() => undefined); - } - return result; - } finally { - controller.abort(); - } -} - -function appendedTail(previous: string, current: string): string { - if (!previous || !current) return current; - const prefix = buildPrefixTable(current); - const maximum = Math.min(previous.length, current.length); - const suffix = previous.slice(previous.length - maximum); - return current.slice(suffixPrefixOverlap(suffix, current, prefix)); -} - -function buildPrefixTable(text: string): Uint32Array { - const prefix = new Uint32Array(text.length); - let matched = 0; - for (let index = 1; index < text.length; index += 1) { - while (matched > 0 && text[index] !== text[matched]) matched = prefix[matched - 1]!; - if (text[index] === text[matched]) matched += 1; - prefix[index] = matched; - } - return prefix; -} - -function suffixPrefixOverlap(suffix: string, current: string, prefix: Uint32Array): number { - let matched = 0; - for (let index = 0; index < suffix.length; index += 1) { - while (matched > 0 && suffix[index] !== current[matched]) matched = prefix[matched - 1]!; - if (suffix[index] === current[matched]) matched += 1; - if (matched === current.length && index < suffix.length - 1) matched = prefix[matched - 1]!; - } - return matched; -} - -function remoteOnlyTail(tail: string): string { - return tail - .split('\n') - .filter((line) => !line.startsWith('[agent-device][mark]')) - .join('\n'); -} - -async function disposeAll(resources: readonly AsyncDisposable[]): Promise { - const results = await Promise.allSettled( - resources.map(async (resource) => await resource[Symbol.asyncDispose]()), - ); - const failures = results.flatMap((result) => - result.status === 'rejected' ? [result.reason] : [], - ); - return failures; } function backendForReader(reader: LimrunAppLogReader): 'ios-simulator' | 'android' { diff --git a/packages/provider-limrun/src/request-cancellation.ts b/packages/provider-limrun/src/request-cancellation.ts index d377dcf3c9..5b5a666469 100644 --- a/packages/provider-limrun/src/request-cancellation.ts +++ b/packages/provider-limrun/src/request-cancellation.ts @@ -2,7 +2,7 @@ * The Limrun WebSocket client does not expose an AbortSignal parameter. Keep the request-bound * caller cancellable while retaining a handled continuation for the provider operation it started. */ -export async function awaitLimrunOperation( +async function awaitLimrunOperation( source: Promise, signal?: AbortSignal, abortMessage = 'Limrun provider operation aborted', diff --git a/scripts/layering/package-boundaries.test.ts b/scripts/layering/package-boundaries.test.ts index c2d9218270..8e62aa1fa4 100644 --- a/scripts/layering/package-boundaries.test.ts +++ b/scripts/layering/package-boundaries.test.ts @@ -393,6 +393,7 @@ test('the real tree parses, declares, and passes R11', () => { assert.deepEqual([...captureKitPackage.exportTargets.keys()].sort(), [ '@agent-device/capture-kit', '@agent-device/capture-kit/android-replacement-surface-occlusion', + '@agent-device/capture-kit/app-log-polling', '@agent-device/capture-kit/audio-probe-admission-ledger', '@agent-device/capture-kit/audio-probe-recovery', '@agent-device/capture-kit/audio-probe-resource-store',