惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
酷 壳 – CoolShell
酷 壳 – CoolShell
博客园_首页
Engineering at Meta
Engineering at Meta
量子位
A
About on SuperTechFans
阮一峰的网络日志
阮一峰的网络日志
Recent Announcements
Recent Announcements
博客园 - 司徒正美
V
Visual Studio Blog
H
Hackread – Cybersecurity News, Data Breaches, AI and More
The GitHub Blog
The GitHub Blog
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
F
Fortinet All Blogs
Martin Fowler
Martin Fowler
腾讯CDC
Jina AI
Jina AI
C
Check Point Blog
H
Help Net Security
罗磊的独立博客
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
V
V2EX
爱范儿
爱范儿
I
InfoQ

Recent Commits to openclaw:main

test: merge chat side-result checks · openclaw/openclaw@ddd2c2a test: merge cron history checks · openclaw/openclaw@f7eb746 test: merge responsive navigation shell checks · openclaw/openclaw@c2e4b47 docs(changelog): add codex oauth fixes · openclaw/openclaw@628e6cd test: merge navigation routing cases · openclaw/openclaw@5d8cecb Tests: mock channel registry bundled fallback · openclaw/openclaw@2b08233 Secrets: avoid broad web search discovery for single plugin config · openclaw/openclaw@a464f59 test: merge config view browser checks · openclaw/openclaw@20cf511 fix(status): align oauth health with runtime · openclaw/openclaw@eed7116 feat: add macOS screen snapshots for monitor preview (#67954) thanks … · openclaw/openclaw@f377db1 fix: report shared auth scopes in hello-ok (#67810) thanks @BunsDev · openclaw/openclaw@0b6c39b Auto-reply: avoid eager bundled route fallback · openclaw/openclaw@3ea1bf4 Tests: narrow session binding contract setup · openclaw/openclaw@54e4e16 fix(macOS): enable undo/redo in webchat composer text input (#34962) · openclaw/openclaw@00951dc Tests: speed up channel setup promotion · openclaw/openclaw@82b529a Docs: refresh agent instructions · openclaw/openclaw@5775fe2 fix(auth): serialize OAuth refresh across agents to fix #26322 (#67876) · openclaw/openclaw@8e79080 test: allow ollama public surface boundary test · openclaw/openclaw@7d4f1a6 Docs: add test performance guardrails · openclaw/openclaw@89706d3 Tests: restore context-engine usage proof · openclaw/openclaw@e4c4f95 Tests: slim context engine runtime coverage · openclaw/openclaw@74c198f ci: retry failed custom checkouts · openclaw/openclaw@0ee5baf test: trim duplicate provider auth onboarding cases · openclaw/openclaw@1ffc02e matrix: fix sessions_spawn --thread subagent session spawning (#67643) · openclaw/openclaw@1ce2596 test: reduce auth choice fixture churn · openclaw/openclaw@857b9cd test: mock health status config boundaries · openclaw/openclaw@9d5ab4a test: mock onboard config io boundary · openclaw/openclaw@299694d test: mock legacy state plugin boundaries · openclaw/openclaw@2713089 test: mock channel install boundaries · openclaw/openclaw@b945248 test: mock doctor preview channel boundaries · openclaw/openclaw@b1a3ad4
refactor: move channel message sdk compat into core · ope...
steipete · 2026-05-27 · via Recent Commits to openclaw:main

@@ -0,0 +1,332 @@

1+

/**

2+

* Shared inbound reply dispatch helpers for channel message adapters and

3+

* deprecated SDK compatibility facades.

4+

*/

5+6+

import { withReplyDispatcher } from "../../auto-reply/dispatch.js";

7+

import type { GetReplyOptions } from "../../auto-reply/get-reply-options.types.js";

8+

import {

9+

dispatchReplyFromConfig,

10+

type DispatchFromConfigResult,

11+

} from "../../auto-reply/reply/dispatch-from-config.js";

12+

import type { DispatchReplyWithBufferedBlockDispatcher } from "../../auto-reply/reply/provider-dispatcher.types.js";

13+

import type { ReplyDispatcher } from "../../auto-reply/reply/reply-dispatcher.types.js";

14+

import type { FinalizedMsgContext } from "../../auto-reply/templating.js";

15+

import type { OpenClawConfig } from "../../config/types.openclaw.js";

16+

import {

17+

normalizeOutboundReplyPayload,

18+

type OutboundReplyPayload,

19+

} from "../../infra/outbound/reply-payload-normalize.js";

20+

import {

21+

hasFinalChannelTurnDispatch,

22+

hasVisibleChannelTurnDispatch,

23+

deliverInboundReplyWithMessageSendContext,

24+

dispatchChannelInboundReply as dispatchChannelInboundReplyCore,

25+

isDurableInboundReplyDeliveryHandled,

26+

resolveChannelTurnDispatchCounts,

27+

recordDroppedChannelInboundHistory,

28+

runChannelInboundEvent as runChannelInboundEventCore,

29+

runPreparedInboundReply as runPreparedInboundReplyCore,

30+

throwIfDurableInboundReplyDeliveryFailed,

31+

} from "../turn/kernel.js";

32+

import type {

33+

ChannelTurnResult,

34+

DispatchedChannelTurnResult,

35+

DurableInboundReplyDeliveryOptions,

36+

} from "../turn/kernel.js";

37+

import type {

38+

AssembledChannelTurn,

39+

PreparedChannelTurn,

40+

RunChannelTurnParams,

41+

} from "../turn/types.js";

42+43+

export type {

44+

ChannelTurnDroppedHistoryOptions,

45+

ChannelTurnDroppedHistoryOptions as ChannelInboundDroppedHistoryOptions,

46+

ChannelTurnRecordOptions,

47+

ChannelTurnRecordOptions as InboundReplyRecordOptions,

48+

} from "../turn/types.js";

49+

export type { DurableInboundReplyDeliveryParams } from "../turn/kernel.js";

50+

export type { ChannelBotLoopProtectionFacts } from "../turn/kernel.js";

51+

export { recordChannelBotPairLoopAndCheckSuppression } from "../turn/kernel.js";

52+53+

type ReplyOptionsWithoutModelSelected = Omit<

54+

Omit<GetReplyOptions, "onBlockReply">,

55+

"onModelSelected"

56+

>;

57+

type RecordInboundSessionFn = typeof import("../session.js").recordInboundSession;

58+59+

type ReplyDispatchFromConfigOptions = Omit<GetReplyOptions, "onBlockReply">;

60+

export type ChannelInboundEventRunnerParams<

61+

TRaw,

62+

TDispatchResult = DispatchFromConfigResult,

63+

> = RunChannelTurnParams<TRaw, TDispatchResult>;

64+

export type PreparedInboundReply<TDispatchResult> = PreparedChannelTurn<TDispatchResult>;

65+

export type AssembledInboundReply = AssembledChannelTurn;

66+

export type InboundReplyDispatchResult<TDispatchResult> = ChannelTurnResult<TDispatchResult>;

67+68+

/** Run an already prepared inbound reply through shared session-record + dispatch ordering. */

69+

type PreparedInboundReplyTurnWithBotLoopProtection<TDispatchResult> =

70+

PreparedChannelTurn<TDispatchResult> & {

71+

botLoopProtection: NonNullable<PreparedChannelTurn<TDispatchResult>["botLoopProtection"]>;

72+

};

73+74+

type PreparedInboundReplyTurnWithoutBotLoopProtection<TDispatchResult> = Omit<

75+

PreparedChannelTurn<TDispatchResult>,

76+

"botLoopProtection"

77+

> & {

78+

botLoopProtection?: undefined;

79+

};

80+81+

export function runPreparedInboundReply<TDispatchResult>(

82+

params: PreparedInboundReplyTurnWithBotLoopProtection<TDispatchResult>,

83+

): Promise<ChannelTurnResult<TDispatchResult>>;

84+

export function runPreparedInboundReply<TDispatchResult>(

85+

params: PreparedInboundReplyTurnWithoutBotLoopProtection<TDispatchResult>,

86+

): Promise<DispatchedChannelTurnResult<TDispatchResult>>;

87+

export function runPreparedInboundReply<TDispatchResult>(

88+

params: PreparedChannelTurn<TDispatchResult>,

89+

): Promise<ChannelTurnResult<TDispatchResult>>;

90+

export async function runPreparedInboundReply<TDispatchResult>(

91+

params: PreparedChannelTurn<TDispatchResult>,

92+

): Promise<ChannelTurnResult<TDispatchResult>> {

93+

return await runPreparedInboundReplyCore(params);

94+

}

95+96+

/** @deprecated Use `runPreparedInboundReply`. */

97+

export function runPreparedInboundReplyTurn<TDispatchResult>(

98+

params: PreparedInboundReplyTurnWithBotLoopProtection<TDispatchResult>,

99+

): Promise<ChannelTurnResult<TDispatchResult>>;

100+

export function runPreparedInboundReplyTurn<TDispatchResult>(

101+

params: PreparedInboundReplyTurnWithoutBotLoopProtection<TDispatchResult>,

102+

): Promise<DispatchedChannelTurnResult<TDispatchResult>>;

103+

export function runPreparedInboundReplyTurn<TDispatchResult>(

104+

params: PreparedChannelTurn<TDispatchResult>,

105+

): Promise<ChannelTurnResult<TDispatchResult>>;

106+

export async function runPreparedInboundReplyTurn<TDispatchResult>(

107+

params: PreparedChannelTurn<TDispatchResult>,

108+

): Promise<ChannelTurnResult<TDispatchResult>> {

109+

return await runPreparedInboundReply(params);

110+

}

111+112+

export async function runChannelInboundEvent<TRaw, TDispatchResult = DispatchFromConfigResult>(

113+

params: ChannelInboundEventRunnerParams<TRaw, TDispatchResult>,

114+

) {

115+

return await runChannelInboundEventCore(params);

116+

}

117+118+

/** @deprecated Use `runChannelInboundEvent`. */

119+

export async function runInboundReplyTurn<TRaw, TDispatchResult = DispatchFromConfigResult>(

120+

params: ChannelInboundEventRunnerParams<TRaw, TDispatchResult>,

121+

) {

122+

return await runChannelInboundEvent(params);

123+

}

124+125+

export async function dispatchChannelInboundReply(params: AssembledInboundReply) {

126+

return await dispatchChannelInboundReplyCore(params);

127+

}

128+129+

export {

130+

hasFinalChannelTurnDispatch as hasFinalInboundReplyDispatch,

131+

hasVisibleChannelTurnDispatch as hasVisibleInboundReplyDispatch,

132+

deliverInboundReplyWithMessageSendContext as deliverDurableInboundReplyPayload,

133+

deliverInboundReplyWithMessageSendContext,

134+

recordDroppedChannelInboundHistory as recordDroppedChannelTurnHistory,

135+

recordDroppedChannelInboundHistory,

136+

resolveChannelTurnDispatchCounts as resolveInboundReplyDispatchCounts,

137+

};

138+139+

/** Run `dispatchReplyFromConfig` with a dispatcher that always gets its settled callback. */

140+

export async function dispatchReplyFromConfigWithSettledDispatcher(params: {

141+

cfg: OpenClawConfig;

142+

ctxPayload: FinalizedMsgContext;

143+

dispatcher: ReplyDispatcher;

144+

onSettled: () => void | Promise<void>;

145+

replyOptions?: ReplyDispatchFromConfigOptions;

146+

configOverride?: OpenClawConfig;

147+

}): Promise<DispatchFromConfigResult> {

148+

return await withReplyDispatcher({

149+

dispatcher: params.dispatcher,

150+

onSettled: params.onSettled,

151+

run: () =>

152+

dispatchReplyFromConfig({

153+

ctx: params.ctxPayload,

154+

cfg: params.cfg,

155+

dispatcher: params.dispatcher,

156+

replyOptions: params.replyOptions,

157+

configOverride: params.configOverride,

158+

}),

159+

});

160+

}

161+162+

/** Assemble the common inbound reply dispatch dependencies for a resolved route. */

163+

export function buildInboundReplyDispatchBase(params: {

164+

cfg: OpenClawConfig;

165+

channel: string;

166+

accountId?: string;

167+

route: {

168+

agentId: string;

169+

sessionKey: string;

170+

};

171+

storePath: string;

172+

ctxPayload: FinalizedMsgContext;

173+

core: {

174+

channel: {

175+

session: {

176+

recordInboundSession: RecordInboundSessionFn;

177+

};

178+

reply: {

179+

dispatchReplyWithBufferedBlockDispatcher: DispatchReplyWithBufferedBlockDispatcher;

180+

};

181+

};

182+

};

183+

}) {

184+

return {

185+

cfg: params.cfg,

186+

channel: params.channel,

187+

accountId: params.accountId,

188+

agentId: params.route.agentId,

189+

routeSessionKey: params.route.sessionKey,

190+

storePath: params.storePath,

191+

ctxPayload: params.ctxPayload,

192+

recordInboundSession: params.core.channel.session.recordInboundSession,

193+

dispatchReplyWithBufferedBlockDispatcher:

194+

params.core.channel.reply.dispatchReplyWithBufferedBlockDispatcher,

195+

};

196+

}

197+198+

type BuildInboundReplyDispatchBaseParams = Parameters<typeof buildInboundReplyDispatchBase>[0];

199+

type RecordChannelMessageReplyDispatchParams = {

200+

cfg: OpenClawConfig;

201+

channel: string;

202+

accountId?: string;

203+

agentId: string;

204+

routeSessionKey: string;

205+

storePath: string;

206+

ctxPayload: FinalizedMsgContext;

207+

recordInboundSession: RecordInboundSessionFn;

208+

dispatchReplyWithBufferedBlockDispatcher: DispatchReplyWithBufferedBlockDispatcher;

209+

deliver: (payload: OutboundReplyPayload) => Promise<void>;

210+

durable?: false | DurableInboundReplyDeliveryOptions;

211+

onRecordError: (err: unknown) => void;

212+

onDispatchError: (err: unknown, info: { kind: string }) => void;

213+

replyOptions?: ReplyOptionsWithoutModelSelected;

214+

};

215+216+

/**

217+

* Resolve the shared dispatch base and immediately record + dispatch one inbound reply turn.

218+

*

219+

* @deprecated Compatibility reply-dispatch bridge. New channel plugins should

220+

* expose a `message` adapter via `defineChannelMessageAdapter(...)` and route

221+

* sends through `deliverInboundReplyWithMessageSendContext(...)` or

222+

* `sendDurableMessageBatch(...)`.

223+

*/

224+

export async function dispatchChannelMessageReplyWithBase(

225+

params: BuildInboundReplyDispatchBaseParams &

226+

Pick<

227+

RecordChannelMessageReplyDispatchParams,

228+

"deliver" | "durable" | "onRecordError" | "onDispatchError" | "replyOptions"

229+

>,

230+

): Promise<void> {

231+

const dispatchBase = buildInboundReplyDispatchBase(params);

232+

await recordChannelMessageReplyDispatch({

233+

...dispatchBase,

234+

deliver: params.deliver,

235+

durable: params.durable,

236+

onRecordError: params.onRecordError,

237+

onDispatchError: params.onDispatchError,

238+

replyOptions: params.replyOptions,

239+

});

240+

}

241+242+

/**

243+

* Resolve the shared dispatch base and immediately record + dispatch one inbound reply turn.

244+

*

245+

* @deprecated Legacy inbound reply helper. New channel plugins should expose a

246+

* `message` adapter via `defineChannelMessageAdapter(...)` and use

247+

* `dispatchChannelMessageReplyWithBase` only for compatibility dispatchers that

248+

* have not moved to the message lifecycle yet.

249+

*/

250+

export async function dispatchInboundReplyWithBase(

251+

params: Parameters<typeof dispatchChannelMessageReplyWithBase>[0],

252+

): Promise<void> {

253+

await dispatchChannelMessageReplyWithBase(params);

254+

}

255+256+

/**

257+

* Record the inbound session first, then dispatch the reply using normalized outbound delivery.

258+

*

259+

* @deprecated Compatibility reply-dispatch bridge. New channel plugins should

260+

* expose a `message` adapter via `defineChannelMessageAdapter(...)` and route

261+

* sends through `deliverInboundReplyWithMessageSendContext(...)` or

262+

* `sendDurableMessageBatch(...)`.

263+

*/

264+

export async function recordChannelMessageReplyDispatch(

265+

params: RecordChannelMessageReplyDispatchParams,

266+

): Promise<void> {

267+

await dispatchChannelInboundReplyCore({

268+

cfg: params.cfg,

269+

channel: params.channel,

270+

accountId: params.accountId,

271+

agentId: params.agentId,

272+

routeSessionKey: params.routeSessionKey,

273+

storePath: params.storePath,

274+

ctxPayload: params.ctxPayload,

275+

recordInboundSession: params.recordInboundSession,

276+

dispatchReplyWithBufferedBlockDispatcher: params.dispatchReplyWithBufferedBlockDispatcher,

277+

delivery: {

278+

preparePayload: (payload) =>

279+

(payload && typeof payload === "object"

280+

? normalizeOutboundReplyPayload(payload as Record<string, unknown>)

281+

: {}) as OutboundReplyPayload,

282+

deliver: async (payload, info) => {

283+

if (params.durable) {

284+

const durable = await deliverInboundReplyWithMessageSendContext({

285+

cfg: params.cfg,

286+

channel: params.channel,

287+

accountId: params.accountId,

288+

agentId: params.agentId,

289+

ctxPayload: params.ctxPayload,

290+

payload,

291+

info,

292+

...params.durable,

293+

});

294+

throwIfDurableInboundReplyDeliveryFailed(durable);

295+

if (isDurableInboundReplyDeliveryHandled(durable)) {

296+

return durable.delivery;

297+

}

298+

}

299+

return await params.deliver(payload as OutboundReplyPayload);

300+

},

301+

onError: params.onDispatchError,

302+

},

303+

replyPipeline: {},

304+

replyOptions: params.replyOptions,

305+

record: {

306+

onRecordError: params.onRecordError,

307+

},

308+

});

309+

}

310+311+

/**

312+

* Record the inbound session first, then dispatch the reply using normalized outbound delivery.

313+

*

314+

* @deprecated Legacy inbound reply helper. New channel plugins should expose a

315+

* `message` adapter via `defineChannelMessageAdapter(...)` and use

316+

* `recordChannelMessageReplyDispatch` only for compatibility dispatchers that

317+

* have not moved to the message lifecycle yet.

318+

*/

319+

export async function recordInboundSessionAndDispatchReply(

320+

params: RecordChannelMessageReplyDispatchParams,

321+

): Promise<void> {

322+

await recordChannelMessageReplyDispatch(params);

323+

}

324+325+

/** @deprecated Compatibility helper for legacy reply dispatch bridges. */

326+

export const buildChannelMessageReplyDispatchBase = buildInboundReplyDispatchBase;

327+

/** @deprecated Compatibility helper for legacy reply dispatch results. */

328+

export const hasFinalChannelMessageReplyDispatch = hasFinalChannelTurnDispatch;

329+

/** @deprecated Compatibility helper for legacy reply dispatch results. */

330+

export const hasVisibleChannelMessageReplyDispatch = hasVisibleChannelTurnDispatch;

331+

/** @deprecated Compatibility helper for legacy reply dispatch results. */

332+

export const resolveChannelMessageReplyDispatchCounts = resolveChannelTurnDispatchCounts;