Skip to content

Commit e99ca64

Browse files
committed
fix: address review feedback on utf8 line chunking, shutdown isolation, and test contracts
1 parent a8b6143 commit e99ca64

10 files changed

Lines changed: 136 additions & 40 deletions

File tree

‎packages/agent-core-v2/src/os/backends/node-local/hostFsService.ts‎

Lines changed: 36 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import { ScopeActivation, registerScopedService } from '#/_base/di/scope';
1414
import { decodeTextWithErrors, type TextDecodeErrors } from '#/_base/execEnv/decodeText';
1515

1616
import { type HostDirEntry, type HostFileStat, IHostFileSystem } from '#/os/interface/hostFileSystem';
17-
import { toHostFsError } from '#/os/interface/hostFsErrors';
17+
import { toHostFsError, HostFsError, OsFsErrors } from '#/os/interface/hostFsErrors';
1818
import { atomicWrite } from '#/_base/utils/fs';
1919

2020
const NEWLINE = Buffer.from([0x0a]);
@@ -38,6 +38,33 @@ function* splitLinesKeepingTerminator(text: string): Generator<string> {
3838
}
3939
}
4040

41+
function trimToValidUtf8Boundary(buf: Buffer): Buffer {
42+
let i = buf.length - 1;
43+
let continuations = 0;
44+
while (i >= 0 && continuations < 3) {
45+
const byte = buf[i];
46+
if (byte === undefined || (byte & 0xc0) !== 0x80) break;
47+
continuations += 1;
48+
i -= 1;
49+
}
50+
if (i < 0) return buf;
51+
const lead = buf[i];
52+
if (lead === undefined) return buf;
53+
if ((lead & 0x80) === 0) {
54+
return continuations === 0 ? buf : buf.subarray(0, i + 1);
55+
}
56+
let needed = 0;
57+
if ((lead & 0xe0) === 0xc0) needed = 1;
58+
else if ((lead & 0xf0) === 0xe0) needed = 2;
59+
else if ((lead & 0xf8) === 0xf0) needed = 3;
60+
else return buf.subarray(0, i);
61+
62+
if (continuations < needed) {
63+
return buf.subarray(0, i);
64+
}
65+
return buf;
66+
}
67+
4168
export class HostFileSystem implements IHostFileSystem {
4269
declare readonly _serviceBrand: undefined;
4370

@@ -121,6 +148,9 @@ export class HostFileSystem implements IHostFileSystem {
121148
const errors = options?.errors ?? 'strict';
122149

123150
if (!isUtf8Encoding(encoding)) {
151+
if (options?.maxLineBytes !== undefined) {
152+
throw new HostFsError(OsFsErrors.codes.OS_FS_UNKNOWN, 'maxLineBytes is only supported for UTF-8 encoding');
153+
}
124154
const content = decodeTextWithErrors(await readFile(path), encoding, errors);
125155
yield* splitLinesKeepingTerminator(content);
126156
return;
@@ -159,8 +189,9 @@ export class HostFileSystem implements IHostFileSystem {
159189
const takeLine = (piece: Buffer): Buffer => {
160190
if (pending.length === 0 && piece.length <= maxLineBytes) return piece;
161191
retain(piece);
162-
const parts = pendingTruncated && piece.at(-1) === 0x0a ? [...pending, NEWLINE] : pending;
163-
const line = Buffer.concat(parts);
192+
const kept = Buffer.concat(pending);
193+
const safe = pendingTruncated ? trimToValidUtf8Boundary(kept) : kept;
194+
const line = pendingTruncated && piece.at(-1) === 0x0a ? Buffer.concat([safe, NEWLINE]) : safe;
164195
pending = [];
165196
pendingBytes = 0;
166197
pendingTruncated = false;
@@ -190,7 +221,8 @@ export class HostFileSystem implements IHostFileSystem {
190221
}
191222

192223
if (pending.length > 0) {
193-
const line = Buffer.concat(pending);
224+
const kept = Buffer.concat(pending);
225+
const line = pendingTruncated ? trimToValidUtf8Boundary(kept) : kept;
194226
yield decodeTextWithErrors(line, 'utf-8', errors, pendingOffset !== 0);
195227
}
196228
} finally {

‎packages/agent-core-v2/test/agent/prompt/promptService.test.ts‎

Lines changed: 40 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,12 @@ import { Readable } from 'node:stream';
44

55
import { DisposableStore } from '#/_base/di/lifecycle';
66
import { createServices } from '#/_base/di/test';
7-
import { Event } from '#/_base/event';
7+
import { Emitter, Event } from '#/_base/event';
88
import { IAgentBlobService } from '#/agent/blob/agentBlobService';
99
import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory';
1010
import type { ContextMessage } from '#/agent/contextMemory/types';
1111
import type { ContentPart } from '#/kosong/contract/message';
12-
import { IAgentFullCompactionService } from '#/agent/fullCompaction/fullCompaction';
12+
import { IAgentFullCompactionService, type FullCompactionTask } from '#/agent/fullCompaction/fullCompaction';
1313
import { IAgentLoopService } from '#/agent/loop/loop';
1414
import { TurnSteer } from '#/agent/loop/turnOps';
1515
import { IAgentPromptService } from '#/agent/prompt/prompt';
@@ -74,13 +74,19 @@ function harness(loopOptions: StubLoopOptions = { pendingTurnResult: true }) {
7474
},
7575
});
7676
const loop = stubLoopWithHooks(loopOptions);
77-
const fullCompaction = {
77+
let activeTask: FullCompactionTask | null = null;
78+
const finishCompactionEmitter = new Emitter<FullCompactionTask>();
79+
disposables.add(finishCompactionEmitter);
80+
const fullCompaction: IAgentFullCompactionService = {
7881
_serviceBrand: undefined,
79-
compacting: null,
82+
get compacting() {
83+
return activeTask;
84+
},
8085
begin: () => false,
86+
cancel: () => {},
8187
hooks: createHooks(['onWillCompact']),
82-
onDidFinishCompaction: Event.None,
83-
} as unknown as IAgentFullCompactionService;
88+
onDidFinishCompaction: finishCompactionEmitter.event,
89+
};
8490
const intake = {
8591
get: vi.fn(async () => ({
8692
meta: {
@@ -124,7 +130,21 @@ function harness(loopOptions: StubLoopOptions = { pendingTurnResult: true }) {
124130
(ix.get(IEventBus) as ISessionEventBus).activateAgent(
125131
ix.get(IAgentScopeContext).agentContext,
126132
);
127-
return { prompt: ix.get(IAgentPromptService), loop, context, fullCompaction, eventBus: ix.get(IEventBus), intake };
133+
return {
134+
prompt: ix.get(IAgentPromptService),
135+
loop,
136+
context,
137+
fullCompaction,
138+
eventBus: ix.get(IEventBus),
139+
intake,
140+
setCompacting: (task: FullCompactionTask | null) => {
141+
activeTask = task;
142+
},
143+
finishCompaction: (task: FullCompactionTask) => {
144+
activeTask = null;
145+
finishCompactionEmitter.fire(task);
146+
},
147+
};
128148
}
129149

130150
describe('AgentPromptService', () => {
@@ -305,14 +325,23 @@ describe('AgentPromptService', () => {
305325
});
306326

307327
it('parks a queued prompt while compaction runs instead of recursing', async () => {
308-
const { prompt, fullCompaction } = harness();
309-
(fullCompaction as unknown as { compacting: unknown }).compacting = { promise: new Promise(() => {}), abortController: new AbortController() };
328+
const { prompt, setCompacting, finishCompaction } = harness();
329+
const task: FullCompactionTask = {
330+
promise: new Promise(() => {}),
331+
abortController: new AbortController(),
332+
trigger: 'manual',
333+
tokenCount: 100,
334+
};
335+
setCompacting(task);
310336
const handle = await prompt.enqueue({ id: 'parked', message: message('later') });
311337
expect(handle.state).toBe('pending');
312-
await (prompt as unknown as { startNext(): Promise<void> }).startNext();
313-
await (prompt as unknown as { startNext(): Promise<void> }).startNext();
314338
expect(prompt.list().pending.map((item) => item.id)).toEqual(['parked']);
315339
expect(prompt.list().active).toBeUndefined();
340+
341+
finishCompaction(task);
342+
await expect(handle.launched).resolves.toBeDefined();
343+
expect(prompt.list().pending).toEqual([]);
344+
expect(prompt.list().active?.id).toBe('parked');
316345
});
317346

318347
it('aborts pending prompts and settles completion', async () => {

‎packages/agent-core-v2/test/agent/toolExecutor/toolScheduler.test.ts‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -256,7 +256,6 @@ describe('ToolScheduler leases and budgets', () => {
256256
abandoned.resolve();
257257
await waitOneMacrotask();
258258

259-
expect(drained).toEqual([]);
260259
expect(started).toEqual(['abandoned']);
261260
settleEffects();
262261
await waitOneMacrotask();

‎packages/agent-core-v2/test/kosong/provider/dsml-tool-parser-conformance.test.ts‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { readFileSync } from 'node:fs';
2-
import { resolve } from 'node:path';
2+
import { dirname, resolve } from 'node:path';
3+
import { fileURLToPath } from 'node:url';
34

45
import { describe, expect, it } from 'vitest';
56

@@ -320,8 +321,9 @@ describe('DsmlStreamParser conformance', () => {
320321

321322
describe('implementation parity with @pymodel/kosong', () => {
322323
it('keeps both parser copies byte-identical except for the message import', () => {
323-
const here = resolve(__dirname, '../../../src/kosong/provider/bases/openai/dsml-tool-parser.ts');
324-
const legacy = resolve(__dirname, '../../../../kosong/src/providers/dsml-tool-parser.ts');
324+
const dir = import.meta.dirname;
325+
const here = resolve(dir, '../../../src/kosong/provider/bases/openai/dsml-tool-parser.ts');
326+
const legacy = resolve(dir, '../../../../kosong/src/providers/dsml-tool-parser.ts');
325327
const v2 = readFileSync(here, 'utf8').replace(
326328
"from '#/kosong/contract/message';",
327329
"from '#/message';",

‎packages/agent-core-v2/test/os/backends/node-local/hostFsService.test.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,24 @@ describe('HostFileSystem readLines budget', () => {
8787

8888
expect(lines).toEqual([`${long}\n`]);
8989
});
90+
91+
it('preserves complete UTF-8 code points when line budget cuts into multibyte characters', async () => {
92+
const path = join(dir, 'multibyte.txt');
93+
await writeFile(path, `${'€'.repeat(100)}\n`, 'utf-8');
94+
95+
const lines: string[] = [];
96+
for await (const line of fs.readLines(path, { maxLineBytes: 16, errors: 'strict' })) lines.push(line);
97+
98+
expect(lines).toEqual([`${'€'.repeat(5)}\n`]);
99+
});
100+
101+
it('rejects maxLineBytes when encoding is not UTF-8', async () => {
102+
const path = join(dir, 'other.txt');
103+
await writeFile(path, 'content', 'utf-8');
104+
105+
const iter = fs.readLines(path, { encoding: 'latin1', maxLineBytes: 4 });
106+
await expect(iter.next()).rejects.toThrow('maxLineBytes is only supported for UTF-8 encoding');
107+
});
90108
});
91109

92110
describe('HostFileSystem stat / lstat', () => {

‎packages/agent-core-v2/test/session/sessionMetadata/sessionMetadata.test.ts‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -234,8 +234,13 @@ describe('SessionMetadata', () => {
234234
await meta.ready;
235235
await meta.registerAgent('sub-1', { homedir: '/tmp/sub-1', type: 'sub', parentAgentId: 'main' });
236236
const before = (await meta.read()).updatedAt;
237-
await meta.unregisterAgent('sub-1');
238-
await meta.unregisterAgent('missing');
237+
const nowSpy = vi.spyOn(Date, 'now').mockReturnValue(before + 10_000);
238+
try {
239+
await meta.unregisterAgent('sub-1');
240+
await meta.unregisterAgent('missing');
241+
} finally {
242+
nowSpy.mockRestore();
243+
}
239244
const after = await meta.read();
240245
expect(after.agents).toEqual({});
241246
expect(after.updatedAt).toBe(before);

‎packages/agent-core-v2/test/session/subagent/runAgentTurn.test.ts‎

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -21,36 +21,37 @@ function fakeHandle(totals: TokenUsage[]): IAgentScopeHandle {
2121
id: 1,
2222
signal: new AbortController().signal,
2323
ready: Promise.resolve(),
24-
result: Promise.resolve({ type: 'completed', steps: 1, truncated: false } as never),
24+
result: Promise.resolve({ type: 'completed', steps: 1, truncated: false }),
2525
cancel: () => true,
2626
};
2727
return {
2828
id: 'sub',
2929
kind: LifecycleScope.Agent,
3030
accessor: {
31-
get: (serviceId: unknown) => {
31+
get: <T>(serviceId: unknown): T => {
3232
if (serviceId === ISessionUsageService) {
3333
return {
3434
status: () => ({ total: totals[Math.min(cursor, totals.length - 1)] }),
35-
};
35+
} as T;
3636
}
3737
if (serviceId === IAgentPromptService) {
3838
return {
3939
enqueue: async () => {
4040
cursor += 1;
4141
return { launched: Promise.resolve(turn) };
4242
},
43-
};
43+
} as T;
4444
}
45-
if (serviceId === IAgentLoopService) return { cancel: () => true };
45+
if (serviceId === IAgentLoopService) return { cancel: () => true } as T;
4646
if (serviceId === IAgentContextMemoryService) {
47-
return { get: () => [{ role: 'assistant', content: [{ type: 'text', text: 'done' }] }] };
47+
return { get: () => [{ role: 'assistant', content: [{ type: 'text', text: 'done' }] }] } as T;
4848
}
49-
if (serviceId === IAgentScopeContext) return { agentContext };
49+
if (serviceId === IAgentScopeContext) return { agentContext } as T;
5050
throw new Error(`unexpected service ${String(serviceId)}`);
5151
},
5252
},
53-
} as unknown as IAgentScopeHandle;
53+
dispose: () => {},
54+
};
5455
}
5556

5657
describe('runAgentTurn usage attribution', () => {

‎packages/agent-gateway/src/start.ts‎

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -318,11 +318,22 @@ export async function startServer(opts: ServerStartOptions): Promise<RunningServ
318318
await phase('config-publisher', () => configChangedPublisher.close());
319319
await phase('http', () => app.close());
320320
await phase('subscriptions', () => {
321-
configWarningSubscription.dispose();
322-
pluginChangeSubscription.dispose();
323-
capabilityInstallSubscription.dispose();
324-
authFailureLimiter?.dispose();
325-
modelCatalogRefreshScheduler.dispose();
321+
for (const sub of [
322+
configWarningSubscription,
323+
pluginChangeSubscription,
324+
capabilityInstallSubscription,
325+
authFailureLimiter,
326+
modelCatalogRefreshScheduler,
327+
]) {
328+
try {
329+
sub?.dispose();
330+
} catch (error) {
331+
logger.warn(
332+
{ err: error instanceof Error ? error.message : String(error) },
333+
'subscription disposal failed; continuing',
334+
);
335+
}
336+
}
326337
});
327338
await phase('telemetry', () => shutdownServerTelemetry(telemetry), false);
328339
await phase('session-metadata', () => drainSessionMetadataWrites());

‎packages/agent-gateway/test/wsConnectionV1.test.ts‎

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -71,11 +71,8 @@ function makeBroadcaster(): SessionEventBroadcaster {
7171
} as unknown as SessionEventBroadcaster;
7272
}
7373

74-
function withBroadcaster(overrides: Record<string, unknown>): SessionEventBroadcaster {
75-
return Object.assign(
76-
makeBroadcaster() as unknown as Record<string, unknown>,
77-
overrides,
78-
) as unknown as SessionEventBroadcaster;
74+
function withBroadcaster(overrides: Partial<SessionEventBroadcaster>): SessionEventBroadcaster {
75+
return Object.assign(makeBroadcaster(), overrides);
7976
}
8077

8178
function makeRegistry(): IConnectionRegistry {

‎packages/kosong/test/openai-legacy.test.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1539,7 +1539,9 @@ describe('OpenAILegacyChatProvider', () => {
15391539
};
15401540
}
15411541

1542-
(provider as any)._client.chat.completions.create = vi
1542+
(
1543+
provider as unknown as { _client: { chat: { completions: { create: unknown } } } }
1544+
)._client.chat.completions.create = vi
15431545
.fn()
15441546
.mockResolvedValue(mockedStream());
15451547

0 commit comments

Comments
 (0)