Skip to content

Commit 5eb99a9

Browse files
authored
Infra: unify plugin split runtime state (#50725)
Merged via squash. Prepared head SHA: 570b7b9 Co-authored-by: huntharo <[email protected]> Co-authored-by: huntharo <[email protected]> Reviewed-by: @huntharo
1 parent 1643d15 commit 5eb99a9

13 files changed

Lines changed: 406 additions & 56 deletions

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -202,6 +202,7 @@ Docs: https://docs.openclaw.ai
202202
- Gateway/probe: honor caller `--timeout` for active local loopback probes in `gateway status`, keep inactive remote-mode loopback probes fast, and clamp probe timers to JS-safe bounds so slow local/container gateways stop reporting false timeouts. (#47533) Thanks @MonkeyLeeT.
203203
- Config/startup: keep bundled web-search allowlist compatibility on a lightweight manifest path so config validation no longer pulls bundled web-search registry imports into startup, while still avoiding accidental auto-allow of config-loaded override plugins. (#51574) Thanks @RichardCao.
204204
- Gateway/chat.send: persist uploaded image references across reloads and compaction without delaying first-turn dispatch or double-submitting the same image to vision models. (#51324) Thanks @fuller-stack-dev.
205+
- Plugins/runtime state: share plugin-facing infra singleton state across duplicate module graphs and keep session-binding adapter ownership stable until the active owner unregisters. (#50725) thanks @huntharo.
205206

206207
### Breaking
207208

extensions/discord/src/monitor/thread-bindings.manager.ts

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import {
55
resolveThreadBindingConversationIdFromBindingId,
66
unregisterSessionBindingAdapter,
77
type BindingTargetKind,
8+
type SessionBindingAdapter,
89
type SessionBindingRecord,
910
} from "openclaw/plugin-sdk/conversation-runtime";
1011
import { normalizeAccountId, resolveAgentIdFromSessionKey } from "openclaw/plugin-sdk/routing";
@@ -556,6 +557,7 @@ export function createThreadBindingManager(
556557
unregisterSessionBindingAdapter({
557558
channel: "discord",
558559
accountId,
560+
adapter: sessionBindingAdapter,
559561
});
560562
forgetThreadBindingToken(accountId);
561563
},
@@ -572,7 +574,7 @@ export function createThreadBindingManager(
572574
}
573575
}
574576

575-
registerSessionBindingAdapter({
577+
const sessionBindingAdapter: SessionBindingAdapter = {
576578
channel: "discord",
577579
accountId,
578580
capabilities: {
@@ -682,7 +684,9 @@ export function createThreadBindingManager(
682684
});
683685
return removed ? [toSessionBindingRecord(removed, { idleTimeoutMs, maxAgeMs })] : [];
684686
},
685-
});
687+
};
688+
689+
registerSessionBindingAdapter(sessionBindingAdapter);
686690

687691
registerManager(manager);
688692
return manager;

extensions/feishu/src/thread-bindings.ts

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import {
66
resolveThreadBindingConversationIdFromBindingId,
77
unregisterSessionBindingAdapter,
88
type BindingTargetKind,
9+
type SessionBindingAdapter,
910
type SessionBindingRecord,
1011
} from "openclaw/plugin-sdk/conversation-runtime";
1112
import { normalizeAccountId, resolveAgentIdFromSessionKey } from "openclaw/plugin-sdk/routing";
@@ -233,11 +234,15 @@ export function createFeishuThreadBindingManager(params: {
233234
}
234235
}
235236
getState().managersByAccountId.delete(accountId);
236-
unregisterSessionBindingAdapter({ channel: "feishu", accountId });
237+
unregisterSessionBindingAdapter({
238+
channel: "feishu",
239+
accountId,
240+
adapter: sessionBindingAdapter,
241+
});
237242
},
238243
};
239244

240-
registerSessionBindingAdapter({
245+
const sessionBindingAdapter: SessionBindingAdapter = {
241246
channel: "feishu",
242247
accountId,
243248
capabilities: {
@@ -292,7 +297,9 @@ export function createFeishuThreadBindingManager(params: {
292297
const removed = manager.unbindConversation(conversationId);
293298
return removed ? [toSessionBindingRecord(removed, { idleTimeoutMs, maxAgeMs })] : [];
294299
},
295-
});
300+
};
301+
302+
registerSessionBindingAdapter(sessionBindingAdapter);
296303

297304
getState().managersByAccountId.set(accountId, manager);
298305
return manager;

extensions/matrix/src/matrix/thread-bindings.ts

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import path from "node:path";
2+
import type { SessionBindingAdapter } from "openclaw/plugin-sdk/conversation-runtime";
23
import {
34
readJsonFileWithFallback,
45
registerSessionBindingAdapter,
@@ -367,6 +368,7 @@ export async function createMatrixThreadBindingManager(params: {
367368
unregisterSessionBindingAdapter({
368369
channel: "matrix",
369370
accountId: params.accountId,
371+
adapter: sessionBindingAdapter,
370372
});
371373
if (getMatrixThreadBindingManagerEntry(params.accountId)?.manager === manager) {
372374
deleteMatrixThreadBindingManagerEntry(params.accountId);
@@ -413,7 +415,7 @@ export async function createMatrixThreadBindingManager(params: {
413415
return removed.map((record) => toSessionBindingRecord(record, defaults));
414416
};
415417

416-
registerSessionBindingAdapter({
418+
const sessionBindingAdapter: SessionBindingAdapter = {
417419
channel: "matrix",
418420
accountId: params.accountId,
419421
capabilities: { placements: ["current", "child"], bindSupported: true, unbindSupported: true },
@@ -512,7 +514,9 @@ export async function createMatrixThreadBindingManager(params: {
512514
);
513515
return removed;
514516
},
515-
});
517+
};
518+
519+
registerSessionBindingAdapter(sessionBindingAdapter);
516520

517521
if (params.enableSweeper !== false) {
518522
sweepTimer = setInterval(() => {

extensions/telegram/src/thread-bindings.ts

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import {
88
resolveThreadBindingEffectiveExpiresAt,
99
unregisterSessionBindingAdapter,
1010
type BindingTargetKind,
11+
type SessionBindingAdapter,
1112
type SessionBindingRecord,
1213
} from "openclaw/plugin-sdk/conversation-runtime";
1314
import { writeJsonAtomic } from "openclaw/plugin-sdk/infra-runtime";
@@ -542,15 +543,19 @@ export function createTelegramThreadBindingManager(
542543
clearInterval(sweepTimer);
543544
sweepTimer = null;
544545
}
545-
unregisterSessionBindingAdapter({ channel: "telegram", accountId });
546+
unregisterSessionBindingAdapter({
547+
channel: "telegram",
548+
accountId,
549+
adapter: sessionBindingAdapter,
550+
});
546551
const existingManager = getThreadBindingsState().managersByAccountId.get(accountId);
547552
if (existingManager === manager) {
548553
getThreadBindingsState().managersByAccountId.delete(accountId);
549554
}
550555
},
551556
};
552557

553-
registerSessionBindingAdapter({
558+
const sessionBindingAdapter: SessionBindingAdapter = {
554559
channel: "telegram",
555560
accountId,
556561
capabilities: {
@@ -687,7 +692,9 @@ export function createTelegramThreadBindingManager(
687692
]
688693
: [];
689694
},
690-
});
695+
};
696+
697+
registerSessionBindingAdapter(sessionBindingAdapter);
691698

692699
const sweeperEnabled = params.enableSweeper !== false;
693700
if (sweeperEnabled) {

src/infra/agent-events.test.ts

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,14 @@ import {
88
resetAgentRunContextForTest,
99
} from "./agent-events.js";
1010

11+
type AgentEventsModule = typeof import("./agent-events.js");
12+
13+
const agentEventsModuleUrl = new URL("./agent-events.ts", import.meta.url).href;
14+
15+
async function importAgentEventsModule(cacheBust: string): Promise<AgentEventsModule> {
16+
return (await import(`${agentEventsModuleUrl}?t=${cacheBust}`)) as AgentEventsModule;
17+
}
18+
1119
describe("agent-events sequencing", () => {
1220
test("stores and clears run context", async () => {
1321
resetAgentRunContextForTest();
@@ -144,4 +152,42 @@ describe("agent-events sequencing", () => {
144152

145153
expect(seen).toEqual(["run-safe"]);
146154
});
155+
156+
test("shares run context, listeners, and sequence state across duplicate module instances", async () => {
157+
const first = await importAgentEventsModule(`first-${Date.now()}`);
158+
const second = await importAgentEventsModule(`second-${Date.now()}`);
159+
160+
first.resetAgentEventsForTest();
161+
first.registerAgentRunContext("run-dup", { sessionKey: "session-dup" });
162+
163+
const seen: Array<{ seq: number; sessionKey?: string }> = [];
164+
const stop = first.onAgentEvent((evt) => {
165+
if (evt.runId === "run-dup") {
166+
seen.push({ seq: evt.seq, sessionKey: evt.sessionKey });
167+
}
168+
});
169+
170+
second.emitAgentEvent({
171+
runId: "run-dup",
172+
stream: "assistant",
173+
data: { text: "from second" },
174+
sessionKey: " ",
175+
});
176+
first.emitAgentEvent({
177+
runId: "run-dup",
178+
stream: "assistant",
179+
data: { text: "from first" },
180+
sessionKey: " ",
181+
});
182+
183+
stop();
184+
185+
expect(second.getAgentRunContext("run-dup")).toEqual({ sessionKey: "session-dup" });
186+
expect(seen).toEqual([
187+
{ seq: 1, sessionKey: "session-dup" },
188+
{ seq: 2, sessionKey: "session-dup" },
189+
]);
190+
191+
first.resetAgentEventsForTest();
192+
});
147193
});

src/infra/agent-events.ts

Lines changed: 31 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import type { VerboseLevel } from "../auto-reply/thinking.js";
2+
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
23

34
export type AgentEventStream = "lifecycle" | "tool" | "assistant" | "error" | (string & {});
45

@@ -19,18 +20,27 @@ export type AgentRunContext = {
1920
isControlUiVisible?: boolean;
2021
};
2122

22-
// Keep per-run counters so streams stay strictly monotonic per runId.
23-
const seqByRun = new Map<string, number>();
24-
const listeners = new Set<(evt: AgentEventPayload) => void>();
25-
const runContextById = new Map<string, AgentRunContext>();
23+
type AgentEventState = {
24+
seqByRun: Map<string, number>;
25+
listeners: Set<(evt: AgentEventPayload) => void>;
26+
runContextById: Map<string, AgentRunContext>;
27+
};
28+
29+
const AGENT_EVENT_STATE_KEY = Symbol.for("openclaw.agentEvents.state");
30+
31+
const state = resolveGlobalSingleton<AgentEventState>(AGENT_EVENT_STATE_KEY, () => ({
32+
seqByRun: new Map<string, number>(),
33+
listeners: new Set<(evt: AgentEventPayload) => void>(),
34+
runContextById: new Map<string, AgentRunContext>(),
35+
}));
2636

2737
export function registerAgentRunContext(runId: string, context: AgentRunContext) {
2838
if (!runId) {
2939
return;
3040
}
31-
const existing = runContextById.get(runId);
41+
const existing = state.runContextById.get(runId);
3242
if (!existing) {
33-
runContextById.set(runId, { ...context });
43+
state.runContextById.set(runId, { ...context });
3444
return;
3545
}
3646
if (context.sessionKey && existing.sessionKey !== context.sessionKey) {
@@ -48,21 +58,21 @@ export function registerAgentRunContext(runId: string, context: AgentRunContext)
4858
}
4959

5060
export function getAgentRunContext(runId: string) {
51-
return runContextById.get(runId);
61+
return state.runContextById.get(runId);
5262
}
5363

5464
export function clearAgentRunContext(runId: string) {
55-
runContextById.delete(runId);
65+
state.runContextById.delete(runId);
5666
}
5767

5868
export function resetAgentRunContextForTest() {
59-
runContextById.clear();
69+
state.runContextById.clear();
6070
}
6171

6272
export function emitAgentEvent(event: Omit<AgentEventPayload, "seq" | "ts">) {
63-
const nextSeq = (seqByRun.get(event.runId) ?? 0) + 1;
64-
seqByRun.set(event.runId, nextSeq);
65-
const context = runContextById.get(event.runId);
73+
const nextSeq = (state.seqByRun.get(event.runId) ?? 0) + 1;
74+
state.seqByRun.set(event.runId, nextSeq);
75+
const context = state.runContextById.get(event.runId);
6676
const isControlUiVisible = context?.isControlUiVisible ?? true;
6777
const eventSessionKey =
6878
typeof event.sessionKey === "string" && event.sessionKey.trim() ? event.sessionKey : undefined;
@@ -73,7 +83,7 @@ export function emitAgentEvent(event: Omit<AgentEventPayload, "seq" | "ts">) {
7383
seq: nextSeq,
7484
ts: Date.now(),
7585
};
76-
for (const listener of listeners) {
86+
for (const listener of state.listeners) {
7787
try {
7888
listener(enriched);
7989
} catch {
@@ -83,6 +93,12 @@ export function emitAgentEvent(event: Omit<AgentEventPayload, "seq" | "ts">) {
8393
}
8494

8595
export function onAgentEvent(listener: (evt: AgentEventPayload) => void) {
86-
listeners.add(listener);
87-
return () => listeners.delete(listener);
96+
state.listeners.add(listener);
97+
return () => state.listeners.delete(listener);
98+
}
99+
100+
export function resetAgentEventsForTest() {
101+
state.seqByRun.clear();
102+
state.listeners.clear();
103+
state.runContextById.clear();
88104
}

src/infra/heartbeat-events.test.ts

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,18 @@ import {
33
emitHeartbeatEvent,
44
getLastHeartbeatEvent,
55
onHeartbeatEvent,
6+
resetHeartbeatEventsForTest,
67
resolveIndicatorType,
78
} from "./heartbeat-events.js";
89

10+
type HeartbeatEventsModule = typeof import("./heartbeat-events.js");
11+
12+
const heartbeatEventsModuleUrl = new URL("./heartbeat-events.ts", import.meta.url).href;
13+
14+
async function importHeartbeatEventsModule(cacheBust: string): Promise<HeartbeatEventsModule> {
15+
return (await import(`${heartbeatEventsModuleUrl}?t=${cacheBust}`)) as HeartbeatEventsModule;
16+
}
17+
918
describe("resolveIndicatorType", () => {
1019
it("maps heartbeat statuses to indicator types", () => {
1120
expect(resolveIndicatorType("ok-empty")).toBe("ok");
@@ -23,6 +32,7 @@ describe("heartbeat events", () => {
2332
});
2433

2534
afterEach(() => {
35+
resetHeartbeatEventsForTest();
2636
vi.useRealTimers();
2737
});
2838

@@ -56,4 +66,28 @@ describe("heartbeat events", () => {
5666

5767
expect(seen).toEqual(["first:ok-empty", "third:ok-empty"]);
5868
});
69+
70+
it("shares heartbeat state across duplicate module instances", async () => {
71+
const first = await importHeartbeatEventsModule(`first-${Date.now()}`);
72+
const second = await importHeartbeatEventsModule(`second-${Date.now()}`);
73+
74+
first.resetHeartbeatEventsForTest();
75+
76+
const seen: string[] = [];
77+
const stop = first.onHeartbeatEvent((evt) => {
78+
seen.push(evt.status);
79+
});
80+
81+
second.emitHeartbeatEvent({ status: "ok-token", preview: "pong" });
82+
83+
expect(first.getLastHeartbeatEvent()).toEqual({
84+
ts: 1767960000000,
85+
status: "ok-token",
86+
preview: "pong",
87+
});
88+
expect(seen).toEqual(["ok-token"]);
89+
90+
stop();
91+
first.resetHeartbeatEventsForTest();
92+
});
5993
});

0 commit comments

Comments
 (0)