mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-09 00:40:45 +08:00
fix(whatsapp): suppress silent group replies and quote sender on reply
Infer group chats from @g.us JIDs, drop (silent)/NO_REPLY outbound text, and @ plus quote the asker when replying in groups. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
900e3d0d3e
commit
a218139da6
4 changed files with 325 additions and 8 deletions
|
|
@ -336,7 +336,9 @@ def _parse_generic_inbound(channel_name: str, payload: dict[str, Any]) -> Inboun
|
|||
raise ValueError("missing user_id")
|
||||
if not chat_id:
|
||||
chat_id = user_id
|
||||
is_group = bool(payload.get("is_group"))
|
||||
from runtime.orchestration.group_ingest import resolve_is_group
|
||||
|
||||
is_group = resolve_is_group(payload_is_group=bool(payload.get("is_group")), chat_id=chat_id)
|
||||
mentions = payload.get("mentions") if isinstance(payload.get("mentions"), list) else []
|
||||
attachments = payload.get("attachments") if isinstance(payload.get("attachments"), list) else []
|
||||
return InboundMessage(
|
||||
|
|
@ -906,6 +908,11 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
|
|||
elif _should_suppress_channel_reply(channel=inbound.channel, text=reply):
|
||||
return {"ok": True, "replies": []}
|
||||
|
||||
from runtime.orchestration.group_ingest import is_nonsend_channel_reply_text
|
||||
|
||||
if is_nonsend_channel_reply_text(reply):
|
||||
return {"ok": True, "replies": []}
|
||||
|
||||
if channel_session_id and reply:
|
||||
_persist_channel_assistant_if_turn_missing(
|
||||
store=store,
|
||||
|
|
@ -920,13 +927,18 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
|
|||
if not reply_attachments:
|
||||
# If assistant didn't persist attachments, fall back to recent tool-produced media.
|
||||
reply_attachments = _collect_recent_tool_attachments(store=store, session_id=str(session_id))
|
||||
reply_metadata: dict[str, Any] = {}
|
||||
if str(inbound.channel or "").strip().lower() == "whatsapp" and inbound.is_group:
|
||||
from runtime.orchestration.group_ingest import build_whatsapp_group_reply_metadata
|
||||
|
||||
reply_metadata = build_whatsapp_group_reply_metadata(inbound=inbound)
|
||||
replies = [
|
||||
{
|
||||
"channel": inbound.channel,
|
||||
"chat_id": inbound.external_chat_id,
|
||||
"text": reply,
|
||||
"attachments": list(reply_attachments or []),
|
||||
"metadata": {},
|
||||
"metadata": reply_metadata,
|
||||
}
|
||||
]
|
||||
# For wechat/weixin sidecar delivery, add a local media_path when reply attachments refer to attachment_id.
|
||||
|
|
|
|||
|
|
@ -146,6 +146,86 @@ function decodeBase64Payload(raw: string): { mime: string; data: Buffer } | null
|
|||
}
|
||||
}
|
||||
|
||||
function readReplyMetadata(reply: Json): Json {
|
||||
const meta = (reply as any).metadata;
|
||||
return meta && typeof meta === "object" ? (meta as Json) : {};
|
||||
}
|
||||
|
||||
function readMentionJids(meta: Json): string[] {
|
||||
const raw = (meta as any).mention_jids ?? (meta as any).mentionJids;
|
||||
if (!Array.isArray(raw)) {
|
||||
const single = String((meta as any).reply_to_user_id || (meta as any).replyToUserId || "").trim();
|
||||
return single ? [jidNormalizedUser(single)] : [];
|
||||
}
|
||||
return raw.map((j) => jidNormalizedUser(String(j || "").trim())).filter(Boolean);
|
||||
}
|
||||
|
||||
function buildMentionPrefix(mentionJids: string[]): string {
|
||||
const tags = mentionJids
|
||||
.map((j) => {
|
||||
const user = String(j || "").split("@")[0]?.trim();
|
||||
return user ? `@${user}` : "";
|
||||
})
|
||||
.filter(Boolean);
|
||||
return tags.length ? `${tags.join(" ")} ` : "";
|
||||
}
|
||||
|
||||
function buildQuotedMessage(params: {
|
||||
chatId: string;
|
||||
stanzaId: string;
|
||||
participant: string;
|
||||
quoteText: string;
|
||||
}): proto.IWebMessageInfo | undefined {
|
||||
const stanzaId = String(params.stanzaId || "").trim();
|
||||
const chatId = String(params.chatId || "").trim();
|
||||
if (!stanzaId || !chatId) return undefined;
|
||||
const participant = jidNormalizedUser(String(params.participant || "").trim());
|
||||
const quoteText = String(params.quoteText || "").trim() || "...";
|
||||
return {
|
||||
key: {
|
||||
remoteJid: chatId,
|
||||
fromMe: false,
|
||||
id: stanzaId,
|
||||
...(participant ? { participant } : {}),
|
||||
},
|
||||
message: {
|
||||
conversation: quoteText,
|
||||
},
|
||||
} as proto.IWebMessageInfo;
|
||||
}
|
||||
|
||||
function buildTextSendOptions(params: {
|
||||
deliverTo: string;
|
||||
text: string;
|
||||
reply: Json;
|
||||
}): { content: { text: string; mentions?: string[] }; quoted?: proto.IWebMessageInfo } {
|
||||
const meta = readReplyMetadata(params.reply);
|
||||
const mentionJids = readMentionJids(meta);
|
||||
const body = String(params.text || "").trim();
|
||||
const prefix = mentionJids.length ? buildMentionPrefix(mentionJids) : "";
|
||||
const content: { text: string; mentions?: string[] } = { text: prefix ? `${prefix}${body}` : body };
|
||||
if (mentionJids.length) content.mentions = mentionJids;
|
||||
|
||||
const quoted = buildQuotedMessage({
|
||||
chatId: String((meta as any).quote_remote_jid || (meta as any).quoteRemoteJid || params.deliverTo).trim(),
|
||||
stanzaId: String((meta as any).quote_stanza_id || (meta as any).quoteStanzaId || "").trim(),
|
||||
participant: String((meta as any).quote_participant || (meta as any).quoteParticipant || "").trim(),
|
||||
quoteText: String((meta as any).quote_text || (meta as any).quoteText || "").trim(),
|
||||
});
|
||||
|
||||
return quoted ? { content, quoted } : { content };
|
||||
}
|
||||
|
||||
function shouldSendOutboundText(text: string): boolean {
|
||||
const t = String(text || "").trim();
|
||||
if (!t) return false;
|
||||
const normalized = t.toLowerCase().replace(/(/g, "(").replace(/)/g, ")");
|
||||
if (normalized === "(silent)" || normalized === "[silent]" || normalized === "no_reply" || normalized === "no reply" || normalized === "静默") {
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
async function sendReplyWithAttachments(params: {
|
||||
sock: ReturnType<typeof makeWASocket> | null;
|
||||
deliverTo: string;
|
||||
|
|
@ -159,12 +239,13 @@ async function sendReplyWithAttachments(params: {
|
|||
const mediaPath = String((reply as any).media_path || (reply as any).mediaPath || "").trim();
|
||||
const mediaUrl = String((reply as any).media_url || (reply as any).mediaUrl || "").trim();
|
||||
const attachments = Array.isArray((reply as any).attachments) ? ((reply as any).attachments as Json[]) : [];
|
||||
const textOpts = buildTextSendOptions({ deliverTo: params.deliverTo, text: outText, reply });
|
||||
|
||||
const sendMediaRef = async (source: string): Promise<boolean> => {
|
||||
if (!source) return false;
|
||||
const msg: Json = { document: source as any };
|
||||
if (outText) (msg as any).caption = outText;
|
||||
await s.sendMessage(params.deliverTo, msg as any);
|
||||
if (outText) (msg as any).caption = textOpts.content.text;
|
||||
await s.sendMessage(params.deliverTo, msg as any, textOpts.quoted ? { quoted: textOpts.quoted } : undefined);
|
||||
return true;
|
||||
};
|
||||
|
||||
|
|
@ -183,13 +264,17 @@ async function sendReplyWithAttachments(params: {
|
|||
mimetype: decoded.mime,
|
||||
fileName: String((att as any).name || (att as any).filename || "attachment.bin"),
|
||||
};
|
||||
if (outText) (msg as any).caption = outText;
|
||||
await s.sendMessage(params.deliverTo, msg as any);
|
||||
if (outText) (msg as any).caption = textOpts.content.text;
|
||||
await s.sendMessage(params.deliverTo, msg as any, textOpts.quoted ? { quoted: textOpts.quoted } : undefined);
|
||||
return;
|
||||
}
|
||||
|
||||
if (outText) {
|
||||
await s.sendMessage(params.deliverTo, { text: outText });
|
||||
await s.sendMessage(
|
||||
params.deliverTo,
|
||||
textOpts.content as any,
|
||||
textOpts.quoted ? { quoted: textOpts.quoted } : undefined,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -412,6 +497,10 @@ async function main(): Promise<void> {
|
|||
if (VERBOSE) log(`inbound ok replies=${replies.length}`);
|
||||
for (const r of replies) {
|
||||
const outText = String((r as any).text || "").trim();
|
||||
if (!shouldSendOutboundText(outText) && !(Array.isArray((r as any).attachments) && (r as any).attachments.length)) {
|
||||
if (VERBOSE) log(`skip outbound reply chat=${chatId} (nonsend text)`);
|
||||
continue;
|
||||
}
|
||||
const deliverTo = String((r as any).chat_id || chatId).trim() || chatId;
|
||||
await sendReplyWithAttachments({ sock, deliverTo, text: outText, reply: r });
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,10 +1,20 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import re
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
GROUP_SESSION_USER_SENTINEL = "__group__"
|
||||
_NONSEND_REPLY_TEXTS = frozenset(
|
||||
{
|
||||
"(silent)",
|
||||
"[silent]",
|
||||
"no_reply",
|
||||
"no reply",
|
||||
"静默",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def normalize_jid(jid: str) -> str:
|
||||
|
|
@ -31,6 +41,31 @@ def session_user_key(*, is_group: bool, external_user_id: str) -> str:
|
|||
return GROUP_SESSION_USER_SENTINEL if is_group else str(external_user_id or "").strip()
|
||||
|
||||
|
||||
def infer_is_group_from_chat_id(chat_id: str) -> bool:
|
||||
c = str(chat_id or "").strip().lower()
|
||||
return c.endswith("@g.us")
|
||||
|
||||
|
||||
def resolve_is_group(*, payload_is_group: bool, chat_id: str) -> bool:
|
||||
if payload_is_group:
|
||||
return True
|
||||
return infer_is_group_from_chat_id(chat_id)
|
||||
|
||||
|
||||
def is_nonsend_channel_reply_text(text: str) -> bool:
|
||||
t = str(text or "").strip()
|
||||
if not t:
|
||||
return True
|
||||
normalized = t.lower().replace("(", "(").replace(")", ")")
|
||||
if normalized in _NONSEND_REPLY_TEXTS:
|
||||
return True
|
||||
return bool(re.fullmatch(r"no_reply", normalized, flags=re.IGNORECASE))
|
||||
|
||||
|
||||
def should_send_channel_reply_text(text: str) -> bool:
|
||||
return not is_nonsend_channel_reply_text(text)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GroupPolicyConfig:
|
||||
require_mention: bool = True
|
||||
|
|
@ -123,13 +158,45 @@ def build_group_sender_context(*, metadata: dict[str, Any] | None, external_user
|
|||
return f"[群成员: {label}]"
|
||||
|
||||
|
||||
def build_whatsapp_group_reply_metadata(
|
||||
*,
|
||||
inbound: Any,
|
||||
) -> dict[str, Any]:
|
||||
"""Outbound hints for WhatsApp sidecar: @ sender + quote original message."""
|
||||
meta = inbound.metadata if isinstance(getattr(inbound, "metadata", None), dict) else {}
|
||||
raw = meta.get("raw") if isinstance(meta.get("raw"), dict) else {}
|
||||
sender_jid = normalize_jid(str(getattr(inbound, "external_user_id", "") or ""))
|
||||
chat_id = str(getattr(inbound, "external_chat_id", "") or "").strip()
|
||||
stanza_id = str(raw.get("id") or meta.get("message_id") or "").strip()
|
||||
participant = normalize_jid(str(raw.get("participant") or sender_jid or ""))
|
||||
quote_text = str(getattr(inbound, "text", "") or "").strip()
|
||||
push_name = str(raw.get("pushName") or meta.get("push_name") or "").strip()
|
||||
out: dict[str, Any] = {
|
||||
"is_group": True,
|
||||
"reply_to_user_id": sender_jid,
|
||||
"mention_jids": [sender_jid] if sender_jid else [],
|
||||
"quote_remote_jid": chat_id,
|
||||
"quote_stanza_id": stanza_id,
|
||||
"quote_participant": participant,
|
||||
"quote_text": quote_text,
|
||||
}
|
||||
if push_name:
|
||||
out["quote_push_name"] = push_name
|
||||
return out
|
||||
|
||||
|
||||
__all__ = [
|
||||
"GROUP_SESSION_USER_SENTINEL",
|
||||
"GroupPolicyConfig",
|
||||
"build_group_sender_context",
|
||||
"build_whatsapp_group_reply_metadata",
|
||||
"normalize_jid",
|
||||
"normalize_jids",
|
||||
"infer_is_group_from_chat_id",
|
||||
"is_nonsend_channel_reply_text",
|
||||
"resolve_is_group",
|
||||
"resolve_group_policy",
|
||||
"session_user_key",
|
||||
"should_process_group_inbound",
|
||||
"should_send_channel_reply_text",
|
||||
]
|
||||
|
|
|
|||
|
|
@ -5,12 +5,16 @@ import pytest
|
|||
from runtime.orchestration.group_ingest import (
|
||||
GROUP_SESSION_USER_SENTINEL,
|
||||
build_group_sender_context,
|
||||
build_whatsapp_group_reply_metadata,
|
||||
infer_is_group_from_chat_id,
|
||||
is_nonsend_channel_reply_text,
|
||||
normalize_jid,
|
||||
resolve_group_policy,
|
||||
session_user_key,
|
||||
should_process_group_inbound,
|
||||
)
|
||||
from runtime.application.gateway.inbound_service import process_inbound_payload
|
||||
from interfaces.channels.base import InboundMessage
|
||||
from runtime.application.gateway.inbound_service import _parse_generic_inbound, process_inbound_payload
|
||||
from svc.persistence.db.engine import clear_assistant_engine_cache
|
||||
from svc.persistence.sqlite_store import SqliteStore
|
||||
|
||||
|
|
@ -98,6 +102,85 @@ def test_resolve_group_policy_from_account_config() -> None:
|
|||
assert policy.triggers == ("!ask",)
|
||||
|
||||
|
||||
def test_infer_is_group_from_chat_id() -> None:
|
||||
assert infer_is_group_from_chat_id("120363012345678@g.us") is True
|
||||
assert infer_is_group_from_chat_id("111@s.whatsapp.net") is False
|
||||
|
||||
|
||||
def test_parse_generic_inbound_infers_whatsapp_group() -> None:
|
||||
inbound = _parse_generic_inbound(
|
||||
"whatsapp",
|
||||
{
|
||||
"user_id": "111@s.whatsapp.net",
|
||||
"chat_id": "120363012345678@g.us",
|
||||
"text": "hello",
|
||||
},
|
||||
)
|
||||
assert inbound.is_group is True
|
||||
|
||||
|
||||
def test_is_nonsend_channel_reply_text() -> None:
|
||||
assert is_nonsend_channel_reply_text("") is True
|
||||
assert is_nonsend_channel_reply_text("(silent)") is True
|
||||
assert is_nonsend_channel_reply_text("(silent)") is True
|
||||
assert is_nonsend_channel_reply_text("NO_REPLY") is True
|
||||
assert is_nonsend_channel_reply_text("你好") is False
|
||||
|
||||
|
||||
def test_inbound_group_without_is_group_flag_is_still_silent(
|
||||
monkeypatch: pytest.MonkeyPatch, fresh_sqlite_store: SqliteStore
|
||||
) -> None:
|
||||
store = fresh_sqlite_store
|
||||
_setup_whatsapp_identity(store)
|
||||
monkeypatch.setattr("svc.persistence.assistant_store.get_assistant_store", lambda: store)
|
||||
|
||||
out = process_inbound_payload(
|
||||
{
|
||||
"channel": "whatsapp",
|
||||
"account_id": "wa-default",
|
||||
"user_id": "111@s.whatsapp.net",
|
||||
"chat_id": "120363012345678@g.us",
|
||||
"text": "Test alarm",
|
||||
"metadata": {"bot_jid": "999@s.whatsapp.net", "source": "test"},
|
||||
}
|
||||
)
|
||||
assert out.get("replies") == []
|
||||
|
||||
|
||||
def test_inbound_suppresses_silent_llm_reply(
|
||||
monkeypatch: pytest.MonkeyPatch, fresh_sqlite_store: SqliteStore
|
||||
) -> None:
|
||||
store = fresh_sqlite_store
|
||||
_setup_whatsapp_identity(store)
|
||||
monkeypatch.setattr("svc.persistence.assistant_store.get_assistant_store", lambda: store)
|
||||
|
||||
class _Turn:
|
||||
turn_uuid = "turn-s"
|
||||
reply_text = "(silent)"
|
||||
|
||||
class _Gw:
|
||||
def __init__(self, *, store: object) -> None:
|
||||
_ = store
|
||||
|
||||
def handle_turn(self, **kwargs: object) -> _Turn:
|
||||
return _Turn()
|
||||
|
||||
monkeypatch.setattr("runtime.gateway.OclawGateway", _Gw)
|
||||
|
||||
out = process_inbound_payload(
|
||||
{
|
||||
"channel": "whatsapp",
|
||||
"account_id": "wa-default",
|
||||
"user_id": "111@s.whatsapp.net",
|
||||
"chat_id": "111@s.whatsapp.net",
|
||||
"text": "ping",
|
||||
"is_group": False,
|
||||
"metadata": {"bot_jid": "999@s.whatsapp.net"},
|
||||
}
|
||||
)
|
||||
assert out.get("replies") == []
|
||||
|
||||
|
||||
def test_build_group_sender_context() -> None:
|
||||
ctx = build_group_sender_context(
|
||||
metadata={"raw": {"pushName": "Alice"}},
|
||||
|
|
@ -107,6 +190,28 @@ def test_build_group_sender_context() -> None:
|
|||
assert "111@s.whatsapp.net" in ctx
|
||||
|
||||
|
||||
def test_build_whatsapp_group_reply_metadata() -> None:
|
||||
inbound = InboundMessage(
|
||||
channel="whatsapp",
|
||||
external_user_id="111:12@s.whatsapp.net",
|
||||
external_chat_id="120363012345678@g.us",
|
||||
text="明天几点?",
|
||||
is_group=True,
|
||||
metadata={
|
||||
"raw": {
|
||||
"id": "MSG123",
|
||||
"participant": "111:12@s.whatsapp.net",
|
||||
"pushName": "Alice",
|
||||
}
|
||||
},
|
||||
)
|
||||
meta = build_whatsapp_group_reply_metadata(inbound=inbound)
|
||||
assert meta["quote_stanza_id"] == "MSG123"
|
||||
assert meta["mention_jids"] == ["111@s.whatsapp.net"]
|
||||
assert meta["quote_text"] == "明天几点?"
|
||||
assert meta["quote_participant"] == "111@s.whatsapp.net"
|
||||
|
||||
|
||||
def test_shared_group_session_for_multiple_senders(fresh_sqlite_store: SqliteStore) -> None:
|
||||
store = fresh_sqlite_store
|
||||
tenant = store.create_tenant("WA")
|
||||
|
|
@ -285,3 +390,47 @@ def test_inbound_group_mention_uses_shared_session_and_sender_prefix(
|
|||
assert sid == session_ids[0]
|
||||
|
||||
_ = tenant_id, user_id
|
||||
|
||||
|
||||
def test_inbound_group_reply_includes_quote_and_mention_metadata(
|
||||
monkeypatch: pytest.MonkeyPatch, fresh_sqlite_store: SqliteStore
|
||||
) -> None:
|
||||
store = fresh_sqlite_store
|
||||
_setup_whatsapp_identity(store)
|
||||
monkeypatch.setattr("svc.persistence.assistant_store.get_assistant_store", lambda: store)
|
||||
|
||||
class _Turn:
|
||||
turn_uuid = "turn-g"
|
||||
reply_text = "下午三点"
|
||||
|
||||
class _Gw:
|
||||
def __init__(self, *, store: object) -> None:
|
||||
_ = store
|
||||
|
||||
def handle_turn(self, **kwargs: object) -> _Turn:
|
||||
return _Turn()
|
||||
|
||||
monkeypatch.setattr("runtime.gateway.OclawGateway", _Gw)
|
||||
|
||||
out = process_inbound_payload(
|
||||
{
|
||||
"channel": "whatsapp",
|
||||
"account_id": "wa-default",
|
||||
"user_id": "111@s.whatsapp.net",
|
||||
"chat_id": "120363012345678@g.us",
|
||||
"text": "@bot 明天几点?",
|
||||
"is_group": True,
|
||||
"mentions": ["999@s.whatsapp.net"],
|
||||
"metadata": {
|
||||
"bot_jid": "999@s.whatsapp.net",
|
||||
"raw": {"id": "ABC123", "participant": "111@s.whatsapp.net", "pushName": "Alice"},
|
||||
},
|
||||
}
|
||||
)
|
||||
replies = out.get("replies") if isinstance(out.get("replies"), list) else []
|
||||
assert replies
|
||||
meta = replies[0].get("metadata") if isinstance(replies[0], dict) else {}
|
||||
assert isinstance(meta, dict)
|
||||
assert meta.get("quote_stanza_id") == "ABC123"
|
||||
assert meta.get("mention_jids") == ["111@s.whatsapp.net"]
|
||||
assert meta.get("quote_text") == "@bot 明天几点?"
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue