Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 35 additions & 12 deletions src/messaging/process-message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ import { isDebugMode } from "./debug-mode.js";
import { sendWeixinErrorNotice } from "./error-notice.js";
import { applyWeixinMessageSendingHook, emitWeixinMessageSent } from "./outbound-hooks.js";
import {
setContextToken,
weixinMessageToMsgContext,
getContextTokenFromMsgContext,
isMediaItem,
Expand All @@ -38,6 +37,16 @@ import { handleSlashCommand } from "./slash-commands.js";

const MEDIA_OUTBOUND_TEMP_DIR = path.join(resolvePreferredOpenClawTmpDir(), "weixin/media/outbound-temp");

type DispatchReplyOptions = NonNullable<
Parameters<PluginRuntime["channel"]["reply"]["dispatchReplyFromConfig"]>[0]["replyOptions"]
> & {
queuedFollowupLifecycle?: {
onEnqueued?: () => void;
onComplete?: () => void;
};
onTurnAdopted?: () => void | Promise<void>;
};

/** Dependencies for processOneMessage, injected by the monitor loop. */
export type ProcessMessageDeps = {
accountId: string;
Expand All @@ -49,6 +58,7 @@ export type ProcessMessageDeps = {
typingTicket?: string;
log: (msg: string) => void;
errLog: (m: string) => void;
onReplyAdmitted?: () => void;
};

/** Extract text body from item_list (for slash command detection). */
Expand Down Expand Up @@ -166,7 +176,7 @@ export async function processOneMessage(
const ctx = weixinMessageToMsgContext(full, deps.accountId, mediaOpts);

// --- Framework command authorization ---
const rawBody = ctx.Body?.trim() ?? "";
const rawBody = textBody.trim() || (ctx.Body?.trim() ?? "");
ctx.CommandBody = rawBody;

const senderId = full.from_user_id ?? "";
Expand Down Expand Up @@ -245,7 +255,10 @@ export async function processOneMessage(
agentId: route.agentId,
});
const finalized = deps.channelRuntime.reply.finalizeInboundContext(
ctx as Parameters<typeof deps.channelRuntime.reply.finalizeInboundContext>[0],
{
...ctx,
BodyForAgent: ctx.Body,
} as Parameters<typeof deps.channelRuntime.reply.finalizeInboundContext>[0],
);

logger.info(
Expand All @@ -270,9 +283,6 @@ export async function processOneMessage(
);

const contextToken = getContextTokenFromMsgContext(ctx);
if (contextToken) {
setContextToken(deps.accountId, full.from_user_id ?? "", contextToken);
}
const runId = randomUUID();
const replyProgressSender = resolveReplyProgressMessagesEnabled(deps.config)
? new WeixinReplyProgressSender({
Expand Down Expand Up @@ -446,6 +456,23 @@ export async function processOneMessage(
},
});

let queuedFollowup = false;
const dispatchReplyOptions: DispatchReplyOptions = {
...replyOptions,
...(replyProgressSender?.replyOptions ?? {}),
// Newer hosts use this marker for active-run admission; older hosts ignore it.
queuedFollowupLifecycle: {
onEnqueued: () => {
queuedFollowup = true;
deps.onReplyAdmitted?.();
},
onComplete: () => void replyProgressSender?.finalize(),
},
onAgentRunStart: () => deps.onReplyAdmitted?.(),
onTurnAdopted: deps.onReplyAdmitted,
disableBlockStreaming: true,
};

logger.debug(`dispatchReplyFromConfig: starting agentId=${route.agentId ?? "(none)"}`);
try {
await deps.channelRuntime.reply.withReplyDispatcher({
Expand All @@ -455,11 +482,7 @@ export async function processOneMessage(
ctx: finalized,
cfg: deps.config,
dispatcher,
replyOptions: {
...replyOptions,
...(replyProgressSender?.replyOptions ?? {}),
disableBlockStreaming: true,
},
replyOptions: dispatchReplyOptions,
}),
});
logger.debug(`dispatchReplyFromConfig: done agentId=${route.agentId ?? "(none)"}`);
Expand All @@ -470,7 +493,7 @@ export async function processOneMessage(
throw err;
} finally {
markDispatchIdle();
await replyProgressSender?.finalize();
if (!queuedFollowup) await replyProgressSender?.finalize();

logger.info(
`debug-check: accountId=${deps.accountId} debug=${String(debug)} hasContextToken=${Boolean(contextToken)}`,
Expand Down
172 changes: 172 additions & 0 deletions src/monitor/monitor.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
import { describe, expect, it, vi } from "vitest";

import { MessageItemType, type GetUpdatesResp, type WeixinMessage } from "../api/types.js";
import type { ProcessMessageDeps } from "../messaging/process-message.js";

const getUpdatesMock = vi.fn<(opts: { abortSignal?: AbortSignal }) => Promise<GetUpdatesResp>>();
const getForUserMock = vi.fn<
(userId: string, contextToken?: string) => Promise<{ typingTicket: string }>
>();
const processOneMessageMock =
vi.fn<(message: WeixinMessage, deps: ProcessMessageDeps) => Promise<void>>();
const saveGetUpdatesBufMock = vi.fn<(filePath: string, value: string) => void>();
const setContextTokenMock = vi.fn<(accountId: string, userId: string, token: string) => void>();

vi.mock("../api/api.js", () => ({
getUpdates: (opts: { abortSignal?: AbortSignal }) => getUpdatesMock(opts),
classifyFetchError: (err: unknown) => ({
type: "mock",
description: String(err),
code: undefined,
}),
}));

vi.mock("../api/config-cache.js", () => ({
WeixinConfigManager: class {
async getForUser(userId: string, contextToken?: string): Promise<{ typingTicket: string }> {
return getForUserMock(userId, contextToken);
}
},
}));

vi.mock("../messaging/process-message.js", () => ({
processOneMessage: (message: WeixinMessage, deps: ProcessMessageDeps) =>
processOneMessageMock(message, deps),
}));

vi.mock("../messaging/inbound.js", () => ({
setContextToken: (accountId: string, userId: string, token: string) =>
setContextTokenMock(accountId, userId, token),
}));

vi.mock("../storage/sync-buf.js", () => ({
getSyncBufFilePath: () => "sync-buf",
loadGetUpdatesBuf: () => undefined,
saveGetUpdatesBuf: (filePath: string, value: string) =>
saveGetUpdatesBufMock(filePath, value),
}));

vi.mock("../util/logger.js", () => ({
logger: {
withAccount: () => ({
info: vi.fn(),
debug: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
}),
},
}));

describe("monitorWeixinProvider", () => {
it("orders ordinary admission while approvals bypass an active ordinary turn", async () => {
vi.resetModules();
const { monitorWeixinProvider } = await import("./monitor.js");
const abortController = new AbortController();
const firstPreprocessing = createDeferred();
const firstRun = createDeferred();
const approvalStarted = createDeferred();
const secondStarted = createDeferred();
const started: string[] = [];
const responses: GetUpdatesResp[] = [
{
ret: 0,
msgs: [makeMessage("first", { message_id: 101, context_token: "token-1" })],
get_updates_buf: "cursor-1",
},
{
ret: 0,
msgs: [
makeMessage("second", { message_id: 102, context_token: "token-2" }),
makeMessage("third", { message_id: 104, context_token: "token-4" }),
makeMessage("/approve plugin:test approve", {
message_id: 103,
context_token: "token-3",
}),
],
get_updates_buf: "cursor-2",
},
];

getUpdatesMock.mockImplementation(async ({ abortSignal }) => {
const next = responses.shift();
if (next) return next;
return await new Promise<GetUpdatesResp>((_, reject) => {
if (abortSignal?.aborted) {
reject(new Error("aborted"));
return;
}
abortSignal?.addEventListener("abort", () => reject(new Error("aborted")), { once: true });
});
});
getForUserMock.mockResolvedValue({ typingTicket: "ticket" });
processOneMessageMock.mockImplementation(async (message, deps) => {
const text = getText(message);
started.push(text);
if (text === "first") {
await firstPreprocessing.promise;
deps.onReplyAdmitted?.();
await firstRun.promise;
return;
}
if (text.startsWith("/approve plugin:")) {
approvalStarted.resolve();
return;
}
if (text === "second") {
secondStarted.resolve();
abortController.abort();
}
});

const monitor = monitorWeixinProvider({
baseUrl: "https://example.test",
cdnBaseUrl: "https://cdn.example.test",
accountId: "acc-monitor",
config: {} as never,
channelRuntime: {} as never,
abortSignal: abortController.signal,
runtime: { log: vi.fn(), error: vi.fn() },
});

try {
await approvalStarted.promise;
expect(started).toEqual(["first", "/approve plugin:test approve"]);
firstPreprocessing.resolve();
await secondStarted.promise;
await monitor;
await new Promise((resolve) => setTimeout(resolve, 0));
expect(started).toEqual(["first", "/approve plugin:test approve", "second"]);
expect(saveGetUpdatesBufMock).toHaveBeenLastCalledWith("sync-buf", "cursor-2");
expect(setContextTokenMock.mock.calls).toEqual([
["acc-monitor", "user-a", "token-1"],
["acc-monitor", "user-a", "token-2"],
["acc-monitor", "user-a", "token-4"],
["acc-monitor", "user-a", "token-3"],
]);
} finally {
firstPreprocessing.resolve();
firstRun.resolve();
await monitor;
}
});
});

function makeMessage(text: string, overrides: Partial<WeixinMessage> = {}): WeixinMessage {
return {
from_user_id: "user-a",
item_list: [{ type: MessageItemType.TEXT, text_item: { text } }],
...overrides,
};
}

function getText(message: WeixinMessage): string {
return message.item_list?.[0]?.text_item?.text ?? "";
}

function createDeferred(): { promise: Promise<void>; resolve: () => void } {
let resolve!: () => void;
const promise = new Promise<void>((innerResolve) => {
resolve = innerResolve;
});
return { promise, resolve };
}
Loading