Skip to content

Commit 653b552

Browse files
committed
wip
1 parent e38b91f commit 653b552

6 files changed

Lines changed: 333 additions & 8 deletions

File tree

TODO.md

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -602,7 +602,29 @@ storage, state, file-history) и умирает БЕЗ dispose (process.exit); r
602602
- Примечание: clientId "shared" неявно зарезервирован — shared storage/state
603603
персистятся под этим ключом в тех же бакетах схем.
604604

605-
## Найденные и исправленные баги (26 итого)
605+
## 27-й заход: последние гарды — MCP listTools и фоновые таймеры (+6 тестов, +1 баг)
606+
607+
### Баг №27: два незакрытых пути после гарантийного свипа
608+
Файлы: `src/client/ClientMCP.ts`, `src/client/ClientAgent.ts`, `src/classes/Chat.ts`
609+
- MCP listTools с throw: await в _resolveTools → реджект queued EXECUTE_FN →
610+
вечное зависание + unhandled; вдобавок memoize fetchTools кэшировал
611+
РЕДЖЕКТНУТЫЙ промис — клиент оставался сломанным навсегда. Исправление:
612+
ClientMCP не кэширует реджект (fetchTools.clear при ошибке), _resolveTools
613+
продолжает без MCP-тулов + errorSubject (session.complete реджектится).
614+
- Бросающий chat-коллбек (onDispose/onCheckActivity) в фоновом cleanup-интервале
615+
ChatUtils → unhandled rejection, роняющий процесс. Исправление: try/catch
616+
вокруг тела итерации handleCleanup.
617+
618+
### Проверено и чисто (зафиксировано тестами)
619+
- queued-диспатчи Storage/State НЕ клинятся после броска (createIndex/dispatchFn):
620+
реджект доставлен вызывающему, следующая операция работает, unhandled нет
621+
(примечание: упавший upsert оставляет item в _itemMap — полузапись);
622+
- calculateSimilarity с throw в take (execpool-путь) — реджект доставлен, чисто;
623+
- makeAutoDispose с бросающим onDestroy — чисто.
624+
625+
### Новые тесты (6), сьют вырос до 271/271 — test/spec/finalguard.test.mjs
626+
627+
## Найденные и исправленные баги (27 итого)
606628

607629
### 1. Дедлок waitForOutput при functools-kit v4 (причина 39 упавших тестов)
608630
Файл: `src/client/ClientSwarm.ts`

src/classes/Chat.ts

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { SwarmName } from "../interfaces/Swarm.interface";
22
import { session } from "../functions/target/session";
3-
import { makeExtendable, singleshot, Subject } from "functools-kit";
3+
import { getErrorMessage, makeExtendable, singleshot, Subject } from "functools-kit";
44
import { SessionId } from "../interfaces/Session.interface";
55
import { GLOBAL_CONFIG } from "../config/params";
66
import swarm from "../lib";
@@ -261,10 +261,18 @@ export class ChatUtils implements IChatControl {
261261
const handleCleanup = async () => {
262262
const now = Date.now();
263263
for (const chat of this._chats.values()) {
264-
if (await chat.checkLastActivity(now)) {
265-
continue;
264+
try {
265+
if (await chat.checkLastActivity(now)) {
266+
continue;
267+
}
268+
await chat.dispose();
269+
} catch (error) {
270+
// The cleanup interval is fire-and-forget: a throwing user callback
271+
// (onDispose, onCheckActivity) must not become an unhandled rejection.
272+
console.error(
273+
`agent-swarm chat cleanup error error=${getErrorMessage(error)}`
274+
);
266275
}
267-
await chat.dispose();
268276
}
269277
};
270278
setInterval(handleCleanup, GLOBAL_CONFIG.CC_CHAT_INACTIVITY_CHECK);

src/client/ClientAgent.ts

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -871,7 +871,20 @@ export class ClientAgent implements IAgent {
871871
})
872872
);
873873
}
874-
const mcpToolList = await this.params.mcp.listTools(this.params.clientId);
874+
let mcpToolList: Awaited<ReturnType<typeof this.params.mcp.listTools>> = [];
875+
try {
876+
mcpToolList = await this.params.mcp.listTools(this.params.clientId);
877+
} catch (error) {
878+
// A throwing MCP listTools would reject the queued EXECUTE_FN and hang
879+
// the pending waitForOutput: continue without MCP tools and surface the
880+
// error to the caller through errorSubject.
881+
console.error(
882+
`agent-swarm mcp listTools error agentName=${
883+
this.params.agentName
884+
} clientId=${this.params.clientId} error=${getErrorMessage(error)}`
885+
);
886+
await errorSubject.next([this.params.clientId, error as Error]);
887+
}
875888
if (mcpToolList.length) {
876889
let commitActionFound = false;
877890
let navigationFound = false;

src/client/ClientMCP.ts

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,15 @@ export class ClientMCP implements IMCP {
6060
if (this.params.callbacks?.onList) {
6161
this.params.callbacks.onList(clientId);
6262
}
63-
const toolMap = await this.fetchTools(clientId);
63+
let toolMap: Awaited<ReturnType<typeof this.fetchTools>>;
64+
try {
65+
toolMap = await this.fetchTools(clientId);
66+
} catch (error) {
67+
// Never cache a rejected fetch: the memoized promise would keep this
68+
// client permanently broken even after the MCP recovers.
69+
this.fetchTools.clear(clientId);
70+
throw error;
71+
}
6472
return Array.from(toolMap.values());
6573
}
6674

@@ -76,7 +84,13 @@ export class ClientMCP implements IMCP {
7684
clientId,
7785
}
7886
);
79-
const toolMap = await this.fetchTools(clientId);
87+
let toolMap: Awaited<ReturnType<typeof this.fetchTools>>;
88+
try {
89+
toolMap = await this.fetchTools(clientId);
90+
} catch (error) {
91+
this.fetchTools.clear(clientId);
92+
throw error;
93+
}
8094
return toolMap.has(toolName);
8195
}
8296

test/index.mjs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ import "./spec/modelguard.test.mjs";
3939
import "./spec/hookguard.test.mjs";
4040
import "./spec/sweepguard.test.mjs";
4141
import "./spec/crashrecovery.test.mjs";
42+
import "./spec/finalguard.test.mjs";
4243

4344
run(import.meta.url, () => {
4445
console.log("All tests are finished");

test/spec/finalguard.test.mjs

Lines changed: 267 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,267 @@
1+
import { test } from "worker-testbed";
2+
3+
import {
4+
addAgent,
5+
addCompletion,
6+
addEmbedding,
7+
addMCP,
8+
addState,
9+
addStorage,
10+
addSwarm,
11+
makeAutoDispose,
12+
session,
13+
setConfig,
14+
Chat,
15+
State,
16+
Storage,
17+
} from "../../build/index.mjs";
18+
import { randomString, sleep } from "functools-kit";
19+
20+
const HANG = Symbol("hang");
21+
const raceHang = (p, ms = 8000) => Promise.race([p, sleep(ms).then(() => HANG)]);
22+
23+
const trackUnhandled = () => {
24+
const unhandled = [];
25+
process.on("unhandledRejection", (reason) => {
26+
unhandled.push(String(reason?.message ?? reason));
27+
});
28+
return unhandled;
29+
};
30+
31+
const addEcho = (name) =>
32+
addCompletion({
33+
completionName: name,
34+
getCompletion: async ({ agentName, messages }) => {
35+
const [last] = messages.slice(-1);
36+
return { agentName, content: `echo:${last.content}`, role: "assistant" };
37+
},
38+
});
39+
40+
test("Will reject without hang when MCP listTools throws", async ({ pass, fail }) => {
41+
setConfig({ CC_PERSIST_ENABLED_BY_DEFAULT: false });
42+
const unhandled = trackUnhandled();
43+
44+
const MOCK_COMPLETION = addEcho("mock-completion");
45+
const TEST_MCP = addMCP({
46+
mcpName: "test-mcp",
47+
listTools: async () => {
48+
throw new Error("listTools exploded");
49+
},
50+
callTool: async () => {},
51+
});
52+
const TEST_AGENT = addAgent({ agentName: "test-agent", completion: MOCK_COMPLETION, prompt: "", mcp: [TEST_MCP] });
53+
const TEST_SWARM = addSwarm({ swarmName: "test-swarm", agentList: [TEST_AGENT], defaultAgent: TEST_AGENT });
54+
55+
const chatSession = session(randomString(), TEST_SWARM);
56+
let out;
57+
try {
58+
out = await raceHang(chatSession.complete("hi"));
59+
} catch (error) {
60+
out = `THREW:${error.message}`;
61+
}
62+
await raceHang(chatSession.dispose());
63+
await sleep(100);
64+
65+
if (String(out).startsWith("THREW:listTools exploded") && unhandled.length === 0) {
66+
pass();
67+
return;
68+
}
69+
fail(`out=${String(out)} unhandled=${JSON.stringify(unhandled)}`);
70+
});
71+
72+
test("Will keep storage queue alive after throwing createIndex", async ({ pass, fail }) => {
73+
setConfig({ CC_PERSIST_ENABLED_BY_DEFAULT: false });
74+
const unhandled = trackUnhandled();
75+
76+
addEcho("mock-completion");
77+
addEmbedding({
78+
embeddingName: "test-embedding",
79+
createEmbedding: async (t) => [t.length],
80+
calculateSimilarity: async (a, b) => (a[0] === b[0] ? 1 : 0),
81+
});
82+
let boom = true;
83+
addStorage({
84+
storageName: "test-storage",
85+
embedding: "test-embedding",
86+
createIndex: (i) => {
87+
if (boom) throw new Error("createIndex exploded");
88+
return i.text;
89+
},
90+
});
91+
const TEST_AGENT = addAgent({
92+
agentName: "test-agent",
93+
completion: "mock-completion",
94+
prompt: "",
95+
storages: ["test-storage"],
96+
});
97+
const TEST_SWARM = addSwarm({ swarmName: "test-swarm", agentList: [TEST_AGENT], defaultAgent: TEST_AGENT });
98+
99+
const CLIENT_ID = randomString();
100+
const chatSession = session(CLIENT_ID, TEST_SWARM);
101+
const base = { clientId: CLIENT_ID, agentName: TEST_AGENT, storageName: "test-storage" };
102+
let first = "";
103+
try {
104+
await raceHang(Storage.upsert({ ...base, item: { id: 1, text: "x" } }));
105+
} catch (error) {
106+
first = error.message;
107+
}
108+
boom = false;
109+
const second = await raceHang(Storage.upsert({ ...base, item: { id: 2, text: "y" } }));
110+
const items = await raceHang(Storage.list(base));
111+
await chatSession.dispose();
112+
await sleep(100);
113+
114+
const ok =
115+
first.includes("createIndex exploded") &&
116+
second !== HANG &&
117+
items !== HANG &&
118+
items.some((i) => i.id === 2) &&
119+
unhandled.length === 0;
120+
if (ok) {
121+
pass();
122+
return;
123+
}
124+
fail(`first=${first} items=${JSON.stringify(items)} unhandled=${JSON.stringify(unhandled)}`);
125+
});
126+
127+
test("Will keep state queue alive after throwing dispatch", async ({ pass, fail }) => {
128+
setConfig({ CC_PERSIST_ENABLED_BY_DEFAULT: false });
129+
const unhandled = trackUnhandled();
130+
131+
addEcho("mock-completion");
132+
addState({ stateName: "test-state", getDefaultState: () => ({ v: 0 }) });
133+
const TEST_AGENT = addAgent({
134+
agentName: "test-agent",
135+
completion: "mock-completion",
136+
prompt: "",
137+
states: ["test-state"],
138+
});
139+
const TEST_SWARM = addSwarm({ swarmName: "test-swarm", agentList: [TEST_AGENT], defaultAgent: TEST_AGENT });
140+
141+
const CLIENT_ID = randomString();
142+
const chatSession = session(CLIENT_ID, TEST_SWARM);
143+
const context = { clientId: CLIENT_ID, agentName: TEST_AGENT, stateName: "test-state" };
144+
let first = "";
145+
try {
146+
await raceHang(
147+
State.setState(() => {
148+
throw new Error("dispatch exploded");
149+
}, context)
150+
);
151+
} catch (error) {
152+
first = error.message;
153+
}
154+
const second = await raceHang(State.setState(() => ({ v: 7 }), context));
155+
const got = await raceHang(State.getState(context));
156+
await chatSession.dispose();
157+
await sleep(100);
158+
159+
const ok =
160+
first.includes("dispatch exploded") &&
161+
second !== HANG &&
162+
got !== HANG &&
163+
got.v === 7 &&
164+
unhandled.length === 0;
165+
if (ok) {
166+
pass();
167+
return;
168+
}
169+
fail(`first=${first} got=${JSON.stringify(got)} unhandled=${JSON.stringify(unhandled)}`);
170+
});
171+
172+
test("Will deliver rejection from throwing calculateSimilarity in take", async ({ pass, fail }) => {
173+
setConfig({ CC_PERSIST_ENABLED_BY_DEFAULT: false });
174+
const unhandled = trackUnhandled();
175+
176+
addEcho("mock-completion");
177+
addEmbedding({
178+
embeddingName: "test-embedding",
179+
createEmbedding: async (t) => [t.length],
180+
calculateSimilarity: async () => {
181+
throw new Error("similarity exploded");
182+
},
183+
});
184+
addStorage({ storageName: "test-storage", embedding: "test-embedding", createIndex: (i) => i.text });
185+
const TEST_AGENT = addAgent({
186+
agentName: "test-agent",
187+
completion: "mock-completion",
188+
prompt: "",
189+
storages: ["test-storage"],
190+
});
191+
const TEST_SWARM = addSwarm({ swarmName: "test-swarm", agentList: [TEST_AGENT], defaultAgent: TEST_AGENT });
192+
193+
const CLIENT_ID = randomString();
194+
const chatSession = session(CLIENT_ID, TEST_SWARM);
195+
const base = { clientId: CLIENT_ID, agentName: TEST_AGENT, storageName: "test-storage" };
196+
await Storage.upsert({ ...base, item: { id: 1, text: "aa" } });
197+
let takeError = "";
198+
try {
199+
await raceHang(Storage.take({ ...base, search: "aa", total: 5, score: 0.1 }));
200+
} catch (error) {
201+
takeError = error.message;
202+
}
203+
await chatSession.dispose();
204+
await sleep(100);
205+
206+
if (takeError.includes("similarity exploded") && unhandled.length === 0) {
207+
pass();
208+
return;
209+
}
210+
fail(`takeError=${takeError} unhandled=${JSON.stringify(unhandled)}`);
211+
});
212+
213+
test("Will survive throwing onDestroy in makeAutoDispose timer", async ({ pass, fail }) => {
214+
setConfig({ CC_PERSIST_ENABLED_BY_DEFAULT: false });
215+
const unhandled = trackUnhandled();
216+
217+
const MOCK_COMPLETION = addEcho("mock-completion");
218+
const TEST_AGENT = addAgent({ agentName: "test-agent", completion: MOCK_COMPLETION, prompt: "" });
219+
const TEST_SWARM = addSwarm({ swarmName: "test-swarm", agentList: [TEST_AGENT], defaultAgent: TEST_AGENT });
220+
221+
const CLIENT_ID = randomString();
222+
session(CLIENT_ID, TEST_SWARM);
223+
const { tick } = makeAutoDispose(CLIENT_ID, TEST_SWARM, {
224+
timeoutSeconds: 1,
225+
onDestroy: () => {
226+
throw new Error("onDestroy exploded");
227+
},
228+
});
229+
tick();
230+
await sleep(2_500);
231+
232+
if (unhandled.length === 0) {
233+
pass();
234+
return;
235+
}
236+
fail(`unhandled=${JSON.stringify(unhandled)}`);
237+
});
238+
239+
test("Will survive throwing chat onDispose in cleanup interval", async ({ pass, fail }) => {
240+
setConfig({
241+
CC_PERSIST_ENABLED_BY_DEFAULT: false,
242+
CC_CHAT_INACTIVITY_CHECK: 100,
243+
CC_CHAT_INACTIVITY_TIMEOUT: 300,
244+
});
245+
const unhandled = trackUnhandled();
246+
247+
Chat.useChatCallbacks({
248+
onDispose: () => {
249+
throw new Error("chat onDispose exploded");
250+
},
251+
});
252+
253+
const MOCK_COMPLETION = addEcho("mock-completion");
254+
const TEST_AGENT = addAgent({ agentName: "test-agent", completion: MOCK_COMPLETION, prompt: "" });
255+
const TEST_SWARM = addSwarm({ swarmName: "test-swarm", agentList: [TEST_AGENT], defaultAgent: TEST_AGENT });
256+
257+
const CLIENT_ID = randomString();
258+
await Chat.beginChat(CLIENT_ID, TEST_SWARM);
259+
const answer = await Chat.sendMessage(CLIENT_ID, "ping", TEST_SWARM);
260+
await sleep(1_000);
261+
262+
if (answer === "echo:ping" && unhandled.length === 0) {
263+
pass();
264+
return;
265+
}
266+
fail(`answer=${String(answer)} unhandled=${JSON.stringify(unhandled)}`);
267+
});

0 commit comments

Comments
 (0)