feat: add group chat activation mode

This commit is contained in:
Peter Steinberger
2025-12-22 19:32:12 +01:00
parent a0dd504991
commit 15e468f5dd
9 changed files with 232 additions and 14 deletions

View File

@@ -19,6 +19,7 @@ import * as commandQueue from "../process/command-queue.js";
import {
HEARTBEAT_PROMPT,
HEARTBEAT_TOKEN,
SILENT_REPLY_TOKEN,
monitorWebProvider,
resolveHeartbeatRecipients,
resolveReplyHeartbeatMinutes,
@@ -1431,6 +1432,81 @@ describe("web auto-reply", () => {
expect(payload.Body).toContain("[from: Bob (+222)]");
});
it("supports always-on group activation with silent token and preserves history", async () => {
const sendMedia = vi.fn();
const reply = vi.fn().mockResolvedValue(undefined);
const sendComposing = vi.fn();
const resolver = vi
.fn()
.mockResolvedValueOnce({ text: SILENT_REPLY_TOKEN })
.mockResolvedValueOnce({ text: "ok" });
setLoadConfigMock(() => ({
inbound: {
groupChat: { activation: "always", mentionPatterns: ["@clawd"] },
},
}));
let capturedOnMessage:
| ((msg: import("./inbound.js").WebInboundMessage) => Promise<void>)
| undefined;
const listenerFactory = async (opts: {
onMessage: (
msg: import("./inbound.js").WebInboundMessage,
) => Promise<void>;
}) => {
capturedOnMessage = opts.onMessage;
return { close: vi.fn() };
};
await monitorWebProvider(false, listenerFactory, false, resolver);
expect(capturedOnMessage).toBeDefined();
await capturedOnMessage?.({
body: "first",
from: "123@g.us",
conversationId: "123@g.us",
chatId: "123@g.us",
chatType: "group",
to: "+2",
id: "g-always-1",
senderE164: "+111",
senderName: "Alice",
selfE164: "+999",
sendComposing,
reply,
sendMedia,
});
expect(resolver).toHaveBeenCalledTimes(1);
expect(reply).not.toHaveBeenCalled();
await capturedOnMessage?.({
body: "second",
from: "123@g.us",
conversationId: "123@g.us",
chatId: "123@g.us",
chatType: "group",
to: "+2",
id: "g-always-2",
senderE164: "+222",
senderName: "Bob",
selfE164: "+999",
sendComposing,
reply,
sendMedia,
});
expect(resolver).toHaveBeenCalledTimes(2);
const payload = resolver.mock.calls[1][0];
expect(payload.Body).toContain("Chat messages since your last reply");
expect(payload.Body).toContain("Alice: first");
expect(payload.Body).toContain("Bob: second");
expect(reply).toHaveBeenCalledTimes(1);
resetLoadConfigMock();
});
it("ignores JID mentions in self-chat mode (group chats)", async () => {
const sendMedia = vi.fn();
const reply = vi.fn().mockResolvedValue(undefined);

View File

@@ -2,8 +2,13 @@ import { chunkText } from "../auto-reply/chunk.js";
import { formatAgentEnvelope } from "../auto-reply/envelope.js";
import { getReplyFromConfig } from "../auto-reply/reply.js";
import type { ReplyPayload } from "../auto-reply/types.js";
import {
HEARTBEAT_TOKEN,
SILENT_REPLY_TOKEN,
} from "../auto-reply/tokens.js";
import { waitForever } from "../cli/wait.js";
import { loadConfig } from "../config/config.js";
import { resolveGroupChatActivation } from "../config/group-chat.js";
import {
DEFAULT_IDLE_MINUTES,
loadSessionStore,
@@ -77,8 +82,8 @@ const formatDuration = (ms: number) =>
ms >= 1000 ? `${(ms / 1000).toFixed(2)}s` : `${ms}ms`;
const DEFAULT_REPLY_HEARTBEAT_MINUTES = 30;
export const HEARTBEAT_TOKEN = "HEARTBEAT_OK";
export const HEARTBEAT_PROMPT = "HEARTBEAT";
export { HEARTBEAT_TOKEN, SILENT_REPLY_TOKEN };
export type WebProviderStatus = {
running: boolean;
@@ -110,7 +115,9 @@ type MentionConfig = {
function buildMentionConfig(cfg: ReturnType<typeof loadConfig>): MentionConfig {
const gc = cfg.inbound?.groupChat;
const requireMention = gc?.requireMention !== false; // default true
const activation = resolveGroupChatActivation(cfg);
const requireMention =
activation === "always" ? false : gc?.requireMention !== false; // default true
const mentionRegexes =
gc?.mentionPatterns
?.map((p) => {
@@ -207,6 +214,14 @@ export function stripHeartbeatToken(raw?: string) {
};
}
function isSilentReply(payload?: ReplyPayload): boolean {
if (!payload) return false;
const text = payload.text?.trim();
if (!text || text !== SILENT_REPLY_TOKEN) return false;
if (payload.mediaUrl || payload.mediaUrls?.length) return false;
return true;
}
export async function runWebHeartbeatOnce(opts: {
cfg?: ReturnType<typeof loadConfig>;
to: string;
@@ -760,6 +775,7 @@ export async function monitorWebProvider(
string,
Array<{ sender: string; body: string; timestamp?: number }>
>();
const groupMemberNames = new Map<string, Map<string, string>>();
const sleep =
tuning.sleep ??
((ms: number, signal?: AbortSignal) =>
@@ -773,6 +789,60 @@ export async function monitorWebProvider(
}),
);
const noteGroupMember = (
conversationId: string,
e164?: string,
name?: string,
) => {
if (!e164 || !name) return;
const normalized = normalizeE164(e164);
const key = normalized ?? e164;
if (!key) return;
let roster = groupMemberNames.get(conversationId);
if (!roster) {
roster = new Map();
groupMemberNames.set(conversationId, roster);
}
roster.set(key, name);
};
const formatGroupMembers = (
participants: string[] | undefined,
roster: Map<string, string> | undefined,
fallbackE164?: string,
) => {
const seen = new Set<string>();
const ordered: string[] = [];
if (participants?.length) {
for (const entry of participants) {
if (!entry) continue;
const normalized = normalizeE164(entry) ?? entry;
if (!normalized || seen.has(normalized)) continue;
seen.add(normalized);
ordered.push(normalized);
}
}
if (roster) {
for (const entry of roster.keys()) {
const normalized = normalizeE164(entry) ?? entry;
if (!normalized || seen.has(normalized)) continue;
seen.add(normalized);
ordered.push(normalized);
}
}
if (ordered.length === 0 && fallbackE164) {
const normalized = normalizeE164(fallbackE164) ?? fallbackE164;
if (normalized) ordered.push(normalized);
}
if (ordered.length === 0) return undefined;
return ordered
.map((entry) => {
const name = roster?.get(entry);
return name ? `${name} (${entry})` : entry;
})
.join(", ");
};
// Avoid noisy MaxListenersExceeded warnings in test environments where
// multiple gateway instances may be constructed.
const currentMaxListeners = process.getMaxListeners?.() ?? 10;
@@ -843,6 +913,7 @@ export async function monitorWebProvider(
emitStatus();
const conversationId = msg.conversationId ?? msg.from;
let combinedBody = buildLine(msg);
let shouldClearGroupHistory = false;
if (msg.chatType === "group") {
const history = groupHistories.get(conversationId) ?? [];
@@ -867,8 +938,7 @@ export async function monitorWebProvider(
? `${msg.senderName} (${msg.senderE164})`
: (msg.senderName ?? msg.senderE164 ?? "Unknown");
combinedBody = `${combinedBody}\\n[from: ${senderLabel}]`;
// Clear stored history after using it
groupHistories.set(conversationId, []);
shouldClearGroupHistory = true;
}
// Echo detection uses combined body so we don't respond twice.
@@ -933,6 +1003,7 @@ export async function monitorWebProvider(
}
const responsePrefix = cfg.inbound?.responsePrefix;
let didSendReply = false;
let toolSendChain: Promise<void> = Promise.resolve();
const sendToolResult = (payload: ReplyPayload) => {
if (
@@ -942,6 +1013,7 @@ export async function monitorWebProvider(
) {
return;
}
if (isSilentReply(payload)) return;
const toolPayload: ReplyPayload = { ...payload };
if (
responsePrefix &&
@@ -961,6 +1033,7 @@ export async function monitorWebProvider(
connectionId,
skipLog: true,
});
didSendReply = true;
if (toolPayload.text) {
recentlySent.add(toolPayload.text);
if (recentlySent.size > MAX_RECENT_MESSAGES) {
@@ -987,7 +1060,11 @@ export async function monitorWebProvider(
MediaType: msg.mediaType,
ChatType: msg.chatType,
GroupSubject: msg.groupSubject,
GroupMembers: msg.groupParticipants?.join(", "),
GroupMembers: formatGroupMembers(
msg.groupParticipants,
groupMemberNames.get(conversationId),
msg.senderE164,
),
SenderName: msg.senderName,
SenderE164: msg.senderE164,
Surface: "whatsapp",
@@ -1004,14 +1081,24 @@ export async function monitorWebProvider(
: [replyResult]
: [];
if (replyList.length === 0) {
logVerbose("Skipping auto-reply: no text/media returned from resolver");
const sendableReplies = replyList.filter(
(payload) => !isSilentReply(payload),
);
if (sendableReplies.length === 0) {
await toolSendChain;
if (shouldClearGroupHistory && didSendReply) {
groupHistories.set(conversationId, []);
}
logVerbose(
"Skipping auto-reply: silent token or no text/media returned from resolver",
);
return;
}
await toolSendChain;
for (const replyPayload of replyList) {
for (const replyPayload of sendableReplies) {
if (
responsePrefix &&
replyPayload.text &&
@@ -1029,6 +1116,7 @@ export async function monitorWebProvider(
replyLogger,
connectionId,
});
didSendReply = true;
if (replyPayload.text) {
recentlySent.add(replyPayload.text);
@@ -1065,6 +1153,10 @@ export async function monitorWebProvider(
);
}
}
if (shouldClearGroupHistory && didSendReply) {
groupHistories.set(conversationId, []);
}
};
const listener = await (listenerFactory ?? monitorWebInbox)({
@@ -1096,6 +1188,7 @@ export async function monitorWebProvider(
}
if (msg.chatType === "group") {
noteGroupMember(conversationId, msg.senderE164, msg.senderName);
const history =
groupHistories.get(conversationId) ??
([] as Array<{ sender: string; body: string; timestamp?: number }>);