feat(whatsapp): real group @mentions and cleaner inbound text

Wire mention JIDs/names through scheduled delivery and immediate replies (including attachment captions), sanitize group message text for the model, and surface channel-derived sessions in admin Chat.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-07-13 20:05:43 +08:00
parent cbe68a54cd
commit ad5d1db3c6
19 changed files with 1507 additions and 58 deletions

View file

@ -386,9 +386,10 @@ def _resolve_chat_session(store: SqliteStore, ctx: dict[str, Any], session_id: s
user_id = str(ctx.get("user_id") or "")
sid = str(session_id or "").strip()
if _is_administrator_chat_viewer(ctx):
sess = store.get_session_for_administrator_username(
sess = store.get_session_for_administrator_chat_view(
session_id=sid,
username=_chat_username(ctx),
tenant_id=tenant_id,
)
if sess is not None:
return sess
@ -891,9 +892,11 @@ def include_chat_routes(router: APIRouter, *, resolve_auth: Callable[[SqliteStor
user_id = str(ctx.get("user_id") or "")
if _is_administrator_chat_viewer(ctx):
uname = _chat_username(ctx)
meta = store.get_sessions_list_meta_for_administrator_username(username=uname)
rows = store.list_sessions_for_administrator_username(
username=uname, limit=limit, offset=offset
meta = store.get_sessions_list_meta_for_administrator_chat_view(
username=uname, tenant_id=tenant_id
)
rows = store.list_sessions_for_administrator_chat_view(
username=uname, tenant_id=tenant_id, limit=limit, offset=offset
)
else:
meta = store.get_sessions_list_meta_for_user(tenant_id=tenant_id, user_id=user_id)

View file

@ -149,6 +149,7 @@ def run_attempt(*, store: Any, data: AttemptRunnerInput) -> AttemptRunnerOutput:
tools=data.tools,
user_text=data.msg.text,
attachments=data.msg.attachments,
inbound_metadata=data.msg.metadata if isinstance(data.msg.metadata, dict) else None,
trace_id=data.trace_id,
parent_span_id=data.parent_span_id,
run_id=data.run_id,

View file

@ -0,0 +1,57 @@
from __future__ import annotations
from typing import Any
def resolve_channel_account_ui_owner(
store: Any,
*,
channel: str,
account_id: str,
tenant_id: str,
fallback_user_id: str = "",
) -> tuple[str, str]:
"""UI list owner for channel-derived sessions: the oclaw user who owns the channel account."""
acct = store.find_user_by_channel_account(channel=str(channel or ""), account_id=str(account_id or "")) or {}
owner_tid = str(acct.get("tenant_id") or tenant_id or "").strip()
owner_uid = str(acct.get("user_id") or "").strip()
if owner_uid:
return owner_tid, owner_uid
admin = store.get_user_by_username(tenant_id=owner_tid or tenant_id, username="administrator")
if isinstance(admin, dict):
owner_uid = str(admin.get("id") or "").strip()
owner_tid = str(admin.get("tenant_id") or owner_tid or tenant_id).strip()
if owner_uid:
return owner_tid, owner_uid
return str(tenant_id or "").strip(), str(fallback_user_id or "").strip()
def assign_channel_session_to_account_owner(
store: Any,
*,
session_id: str,
channel: str,
account_id: str,
tenant_id: str,
fallback_user_id: str = "",
) -> None:
owner_tid, owner_uid = resolve_channel_account_ui_owner(
store,
channel=channel,
account_id=account_id,
tenant_id=tenant_id,
fallback_user_id=fallback_user_id,
)
if not owner_tid or not owner_uid:
return
setter = getattr(store, "set_ui_session_owner", None)
if callable(setter):
setter(session_id=str(session_id), tenant_id=owner_tid, user_id=owner_uid)
return
store.ensure_ui_session_owner(session_id=str(session_id), tenant_id=owner_tid, user_id=owner_uid)
__all__ = [
"assign_channel_session_to_account_owner",
"resolve_channel_account_ui_owner",
]

View file

@ -825,14 +825,13 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
account = store.find_user_by_channel_account(channel=inbound.channel, account_id=account_id) or {}
from runtime.orchestration.group_ingest import (
build_group_focus_instruction,
build_group_quoted_context_block,
build_group_sender_context,
enrich_alert_group_question,
extract_group_quoted_message,
extract_quoted_ume_alert_text,
mentions_include_bot,
metadata_mentions_bot,
prepare_group_user_text_for_model,
resolve_group_policy,
session_user_key,
should_inject_quoted_context,
@ -965,7 +964,18 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
),
)
channel_session_id = str(session_id)
store.ensure_ui_session_owner(session_id=session_id, tenant_id=tenant_id, user_id=user_id)
from runtime.application.gateway.channel_session_owner import (
assign_channel_session_to_account_owner,
)
assign_channel_session_to_account_owner(
store,
session_id=session_id,
channel=str(inbound.channel or ""),
account_id=account_id,
tenant_id=tenant_id,
fallback_user_id=user_id,
)
if str(inbound.channel or "").strip().lower() in {"wechat", "weixin"}:
from runtime.scheduler.channel_delivery import (
extract_context_token_from_inbound_metadata,
@ -1021,26 +1031,90 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
if cmd_reply is not None:
reply = cmd_reply
elif not reply:
user_text = (inbound.text or "").strip()
raw_user_text = (inbound.text or "").strip()
user_text = raw_user_text
meta_for_inbound = inbound.metadata if isinstance(inbound.metadata, dict) else {}
mention_jids_for_ctx: list[str] = []
seen_mention_jids: set[str] = set()
for src in (list(inbound.mentions or []), meta_for_inbound.get("mentioned_jids") or []):
if not isinstance(src, list):
continue
for item in src:
jid = str(item or "").strip()
if jid and jid not in seen_mention_jids:
seen_mention_jids.add(jid)
mention_jids_for_ctx.append(jid)
bot_jid_for_filter = str(meta_for_inbound.get("bot_jid") or bot_jid or "").strip().lower()
bot_lid_for_filter = str(meta_for_inbound.get("bot_lid") or "").strip().lower()
filtered_mention_jids: list[str] = []
for jid in mention_jids_for_ctx:
low = jid.lower()
if (bot_jid_for_filter and low == bot_jid_for_filter) or (
bot_lid_for_filter and low == bot_lid_for_filter
):
continue
filtered_mention_jids.append(jid)
mention_names_for_ctx: list[str] = []
if str(inbound.channel or "").strip().lower() == "whatsapp" and filtered_mention_jids:
from runtime.orchestration.group_ingest import _bot_identity_jids, jid_base_local, jid_phone
from runtime.scheduler.whatsapp_mentions import (
_lookup_push_name,
extract_whatsapp_mention_names,
)
raw_names = extract_whatsapp_mention_names(raw_user_text)
bot_name_hints = {
str(meta_for_inbound.get("bot_push_name") or "").strip().lower(),
"oliver",
}
for bid in _bot_identity_jids(
bot_jid=bot_jid, metadata=meta_for_inbound
):
local = jid_base_local(bid)
if local:
bot_name_hints.add(local.lower())
phone = jid_phone(bid)
if phone:
bot_name_hints.add(phone)
text_names = [n for n in raw_names if n.lower() not in bot_name_hints and n]
text_name_idx = 0
upsert_contact = getattr(store, "upsert_whatsapp_contact", None)
wa_acct = str(account_id or "").strip() or "wa-default"
for jid in filtered_mention_jids:
push_name = _lookup_push_name(
store,
tenant_id=str(tenant_id or ""),
account_id=wa_acct,
jid=jid,
)
if not push_name and text_name_idx < len(text_names):
push_name = text_names[text_name_idx]
text_name_idx += 1
mention_names_for_ctx.append(str(push_name or "").strip())
if callable(upsert_contact):
try:
upsert_contact(
tenant_id=str(tenant_id or ""),
account_id=wa_acct,
external_user_id=jid,
push_name=str(push_name or "").strip(),
)
except Exception:
pass
if inbound.is_group:
meta_for_group = inbound.metadata if isinstance(inbound.metadata, dict) else {}
meta_for_group = meta_for_inbound
bot_reached = metadata_mentions_bot(meta_for_group) or mentions_include_bot(
mentions=list(inbound.mentions or []),
bot_jid=bot_jid,
metadata=meta_for_group,
) or text_mentions_bot(text=user_text, bot_jid=bot_jid)
) or text_mentions_bot(text=raw_user_text, bot_jid=bot_jid)
if bot_reached:
quoted_alert = extract_quoted_ume_alert_text(metadata=meta_for_group)
if quoted_alert:
user_text = enrich_alert_group_question(
user_text=user_text,
raw_user_text = enrich_alert_group_question(
user_text=raw_user_text,
quoted_alert=quoted_alert,
)
sender_ctx = build_group_sender_context(
metadata=meta_for_group,
external_user_id=inbound.external_user_id,
)
group_rule = build_group_focus_instruction()
quoted_ctx = ""
quoted_info = extract_group_quoted_message(metadata=meta_for_group)
quoted_text = str(quoted_info.get("quoted_text") or "").strip()
@ -1051,9 +1125,21 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
recent_messages=recent_messages,
):
quoted_ctx = build_group_quoted_context_block(metadata=meta_for_group)
prefix_parts = [sender_ctx, group_rule, quoted_ctx]
prefix = "\n".join(part for part in prefix_parts if str(part).strip())
user_text = f"{prefix}\n{user_text}" if user_text else prefix
wa_acct = str(account_id or "").strip() or "wa-default"
user_text = prepare_group_user_text_for_model(
text=raw_user_text,
metadata=meta_for_group,
mentions=list(inbound.mentions or []),
bot_jid=bot_jid,
session_scope=group_policy.session_scope,
external_user_id=inbound.external_user_id,
filtered_mention_jids=filtered_mention_jids,
mention_names=mention_names_for_ctx,
store=store,
tenant_id=str(tenant_id or ""),
account_id=wa_acct,
quoted_ctx=quoted_ctx,
)
gw_attachments = _channel_attachments_for_gateway(
list(inbound.attachments or [])
)
@ -1094,6 +1180,11 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]:
"external_user_id": inbound.external_user_id,
"external_chat_id": inbound.external_chat_id,
"group_sender_id": inbound.external_user_id,
"mentioned_jids": mention_jids_for_ctx,
"mention_names": mention_names_for_ctx,
"raw_inbound_text": raw_user_text,
"bot_jid": str(meta_for_inbound.get("bot_jid") or bot_jid or ""),
"bot_lid": str(meta_for_inbound.get("bot_lid") or ""),
},
)
manager = _build_admin_gateway_executor(

View file

@ -479,6 +479,7 @@ class ToolExecutionContext:
session_id: str
lang: str = "zh"
user_text: str = ""
inbound_metadata: dict[str, Any] | None = None
specialist: str = ""
task_kind: str = ""
policy_engine: Any | None = None
@ -529,11 +530,59 @@ class ToolExecutor:
try:
tname = str(tc.name or "")
if tname == "schedule_create":
md = ctx.inbound_metadata if isinstance(ctx.inbound_metadata, dict) else {}
cur_spec = str(ctx.specialist or "").strip().lower()
if cur_spec and not str(tool_args.get("specialist") or "").strip() and not str(
tool_args.get("selected_specialist") or ""
).strip():
tool_args["selected_specialist"] = cur_spec
raw_mentions = md.get("mentioned_jids") or md.get("mentionedJids") or md.get("mentions") or []
mention_list = raw_mentions if isinstance(raw_mentions, list) else []
bot_jid = str(md.get("bot_jid") or "").strip().lower()
bot_lid = str(md.get("bot_lid") or "").strip().lower()
cleaned: list[str] = []
seen: set[str] = set()
for m in mention_list:
jid = str(m or "").strip()
if not jid:
continue
low = jid.lower()
if (bot_jid and low == bot_jid) or (bot_lid and low == bot_lid):
continue
if low in seen:
continue
seen.add(low)
cleaned.append(jid)
if cleaned:
tool_args["whatsapp_mention_jids"] = cleaned
elif tool_args.get("whatsapp_mention_jids") is not None:
from runtime.scheduler.whatsapp_mentions import normalize_whatsapp_mention_jids
tool_args["whatsapp_mention_jids"] = normalize_whatsapp_mention_jids(
tool_args.get("whatsapp_mention_jids")
)
if tool_args.get("whatsapp_mention_names") is None:
from runtime.scheduler.whatsapp_mentions import extract_whatsapp_mention_names
user_text = str(ctx.user_text or "")
names = extract_whatsapp_mention_names(user_text)
bot_names = {
str(x or "").strip().lower()
for x in (
md.get("bot_push_name"),
"oliver",
)
if str(x or "").strip()
}
filtered_names = [n for n in names if n.lower() not in bot_names]
if filtered_names:
tool_args["whatsapp_mention_names"] = filtered_names
elif md.get("mention_names"):
names_md = md.get("mention_names")
if isinstance(names_md, list):
tool_args["whatsapp_mention_names"] = [
str(x or "").strip() for x in names_md if str(x or "").strip()
]
except Exception:
pass
tool_args = filter_arguments_to_schema(tool.parameters, tool_args)

View file

@ -864,6 +864,16 @@ def _dsml_mismatch_user_message(*, lang: str) -> str:
)
def _empty_final_user_message(*, lang: str, had_tools: bool) -> str:
if had_tools:
if str(lang or "").strip().lower().startswith("zh"):
return "工具执行未返回可展示结果,请稍后重试;若持续失败,请更换查询关键词。"
return "Tool execution did not produce a user-visible result. Please retry with a narrower query."
if str(lang or "").strip().lower().startswith("zh"):
return "抱歉,这次没有生成可展示内容,请重试。"
return "Sorry, no user-visible content was produced this turn. Please try again."
def _promote_dsml_tool_calls(
*,
allow: bool,
@ -1131,6 +1141,7 @@ def _execute_tool_step(
run_id: str | None = None,
attempt_no: int | None = None,
turn_uuid: str | None = None,
inbound_metadata: dict[str, Any] | None = None,
) -> tuple[int, dict[str, tuple[dict[str, Any], int]]]:
t0 = time.perf_counter()
_tool_messages, results_by_id = skill_exec.execute_skill_uses(
@ -1140,6 +1151,7 @@ def _execute_tool_step(
session_id=session_id,
lang=lang,
user_text=user_text,
inbound_metadata=inbound_metadata,
specialist="oclaw",
trace_id=trace_id,
parent_span_id=parent_span_id,
@ -1390,6 +1402,7 @@ def run_oclaw_direct_loop(
wire_policy_role: str | None = None,
prompt_build_context: dict[str, Any] | None = None,
turn_uuid: str | None = None,
inbound_metadata: dict[str, Any] | None = None,
) -> TurnRunOutcome:
"""A minimal oclaw-style loop: model -> tool_uses -> execute -> tool_results -> continue."""
_check_stop(should_stop)
@ -1599,6 +1612,7 @@ def run_oclaw_direct_loop(
run_id=run_id,
attempt_no=attempt_no,
turn_uuid=turn_uuid,
inbound_metadata=inbound_metadata,
)
for tc in step.llm_tool_calls:
@ -1662,6 +1676,18 @@ def run_oclaw_direct_loop(
)
final_text = step.assistant_text
if not str(final_text or "").strip():
fallback = _empty_final_user_message(lang=lang, had_tools=bool(tool_traces))
store.add_message(
session_id=session_id,
role="assistant",
content=fallback,
turn_uuid=turn_uuid,
event_type="assistant_text",
event_payload={"runtime_fallback": "empty_final_text"},
)
final_text = fallback
return TurnRunOutcome(
final_text=str(final_text or ""),
tool_traces=tuple(tool_traces),

View file

@ -1251,6 +1251,21 @@ class OclawGateway:
)
else:
reply = specialist_reply
if not str(reply or "").strip():
rs = getattr(core_out, "run_state", None)
if rs is not None and str(getattr(rs, "status", "") or "") == "failed":
reason = ""
attempts = getattr(rs, "attempts", ()) or ()
if attempts:
reason = str(getattr(attempts[-1], "reason", "") or "")
if not reason:
reason = str(getattr(rs, "last_error_code", "") or "unknown_error")
base = render_prompt(
"fallback/runtime_error.en.md" if str(lang or "").startswith("en") else "fallback/runtime_error.zh.md",
strict=True,
)
detail = str(reason or "").strip().replace("\n", " ")[:400]
reply = f"{base}\n(detail: {detail})" if detail else base
except Exception as exc:
base = render_prompt(
"fallback/runtime_error.en.md" if str(lang or "").startswith("en") else "fallback/runtime_error.zh.md",

View file

@ -527,14 +527,151 @@ function readMentionJids(meta: Json): string[] {
return raw.map((j) => 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 readMentionNames(meta: Json): string[] {
const raw = (meta as any).mention_names ?? (meta as any).mentionNames;
if (!Array.isArray(raw)) return [];
return raw.map((n) => String(n || "").trim()).filter(Boolean);
}
function escapeRegExp(text: string): string {
return String(text || "").replace(/[.*+?^${}()|[\]\\]/g, "\\$&");
}
function mentionTokenForJid(jid: string): string {
const local = jidBaseLocal(jid);
return local ? `@${local}` : "";
}
function resolveOutboundMentionJids(
sock: ReturnType<typeof makeWASocket> | null,
chatId: string,
mentionJids: string[],
): string[] {
const isGroup = String(chatId || "").toLowerCase().endsWith("@g.us");
const lidMapping = (sock as any)?.signalRepository?.lidMapping;
const out: string[] = [];
const seen = new Set<string>();
const push = (jid: string) => {
const j = String(jid || "").trim();
if (!j) return;
const key = jidBaseLocal(j);
if (!key || seen.has(key)) return;
seen.add(key);
out.push(j);
};
for (let raw of mentionJids) {
let jid = String(raw || "").trim();
if (!jid) continue;
if (!jid.includes("@") && /^\d+$/.test(jid)) {
jid = `${jid}@lid`;
}
const low = jid.toLowerCase();
if (low.endsWith("@lid")) {
push(jid);
continue;
}
if (isGroup && lidMapping && typeof lidMapping.getLIDForPN === "function") {
try {
const pn = low.includes("@s.whatsapp") ? jidNormalizedUser(jid) : jid;
const lid = lidMapping.getLIDForPN(pn);
if (lid) {
push(String(lid));
continue;
}
} catch {
// ignore
}
}
push(jid.includes("@") ? jid : jidNormalizedUser(jid));
}
return out;
}
function alignMentionTextWithJids(text: string, mentionJids: string[], mentionNames: string[]): string {
let body = String(text || "").trim();
const tokens = mentionJids.map((j) => mentionTokenForJid(j)).filter(Boolean);
if (!tokens.length) return body;
for (let i = 0; i < mentionJids.length; i++) {
const token = tokens[i] || "";
if (!token) continue;
const name = String(mentionNames[i] || "").trim();
if (name) {
body = body.replace(new RegExp(`@${escapeRegExp(name)}(?=\\s|$|[,。!?!?,.])`, "gu"), token);
}
}
const missing = tokens.filter((token) => !body.includes(token));
if (missing.length) {
body = body.replace(/^(@\S+\s*)+/, "").trim();
body = `${missing.join(" ")} ${body}`.trim();
}
return body;
}
function decodeOutboundSource(raw: string): Json {
const text = String(raw || "").trim();
if (!text) return {};
if (text.startsWith("{")) {
try {
const data = JSON.parse(text);
return data && typeof data === "object" ? (data as Json) : {};
} catch {
return {};
}
}
return { kind: text };
}
function isOutboundMentionTextReady(meta: Json): boolean {
return Boolean((meta as any).mention_text_ready ?? (meta as any).mentionTextReady);
}
function buildMentionedOutboundText(
body: string,
meta: Json,
sock?: ReturnType<typeof makeWASocket> | null,
chatId?: string,
): { text: string; mentions?: string[] } {
const trimmed = String(body || "").trim();
if (!trimmed) return { text: trimmed };
let mentionJids = readMentionJids(meta);
if (!mentionJids.length) return { text: trimmed };
const chat = String(chatId || "").trim();
if (sock && chat) {
mentionJids = resolveOutboundMentionJids(sock, chat, mentionJids);
}
const mentionNames = readMentionNames(meta);
if (isOutboundMentionTextReady(meta)) {
const aligned = alignMentionTextWithJids(trimmed, mentionJids, mentionNames);
return { text: aligned, mentions: mentionJids };
}
const cleaned = trimmed
.replace(/@\+?\d+(?:\s+\d+)*\s*/g, "")
.replace(/@\S+\s*/g, "")
.trim();
const prefix = mentionJids.map((j) => mentionTokenForJid(j)).filter(Boolean).join(" ");
const outText = prefix ? `${prefix} ${cleaned}`.trim() : trimmed;
return { text: outText, mentions: mentionJids };
}
function buildOutboundSendContent(
text: string,
sourceRaw: string,
sock?: ReturnType<typeof makeWASocket> | null,
chatId?: string,
): { text: string; mentions?: string[] } {
const meta = decodeOutboundSource(sourceRaw);
return buildMentionedOutboundText(text, meta, sock, chatId);
}
function buildQuotedMessage(params: {
@ -565,13 +702,10 @@ function buildTextSendOptions(params: {
deliverTo: string;
text: string;
reply: Json;
sock?: ReturnType<typeof makeWASocket> | null;
}): { 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 content = buildMentionedOutboundText(params.text, meta, params.sock, params.deliverTo);
const quoted = buildQuotedMessage({
chatId: String((meta as any).quote_remote_jid || (meta as any).quoteRemoteJid || params.deliverTo).trim(),
@ -628,25 +762,28 @@ function buildOutboundMediaMessage(params: {
mime: string;
fileName: string;
caption?: string;
mentions?: string[];
}): Json {
const mime = String(params.mime || "application/octet-stream").trim() || "application/octet-stream";
const m = mime.toLowerCase();
const caption = String(params.caption || "").trim();
const captionOpt = caption ? { caption } : {};
const mentionOpt = params.mentions?.length ? { mentions: params.mentions } : {};
if (m.startsWith("image/")) {
return { image: params.data as any, mimetype: mime, ...captionOpt };
return { image: params.data as any, mimetype: mime, ...captionOpt, ...mentionOpt };
}
if (m.startsWith("video/")) {
return { video: params.data as any, mimetype: mime, ...captionOpt };
return { video: params.data as any, mimetype: mime, ...captionOpt, ...mentionOpt };
}
if (m.startsWith("audio/")) {
return { audio: params.data as any, mimetype: mime, ...captionOpt };
return { audio: params.data as any, mimetype: mime, ...captionOpt, ...mentionOpt };
}
return {
document: params.data as any,
mimetype: mime,
fileName: sanitizeFileName(params.fileName || "attachment.bin"),
...captionOpt,
...mentionOpt,
};
}
@ -663,7 +800,12 @@ 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 textOpts = buildTextSendOptions({
deliverTo: params.deliverTo,
text: outText,
reply,
sock: params.sock,
});
const sendOpts = textOpts.quoted ? { quoted: textOpts.quoted } : undefined;
const sendMediaBuffer = async (data: Buffer, mime: string, fileName: string): Promise<boolean> => {
@ -673,6 +815,7 @@ async function sendReplyWithAttachments(params: {
mime,
fileName,
caption: textOpts.content.text,
mentions: textOpts.content.mentions,
});
await s.sendMessage(params.deliverTo, msg as any, sendOpts);
return true;
@ -753,13 +896,17 @@ async function pollOutboundQueue(sock: ReturnType<typeof makeWASocket>): Promise
const id = String((item as any).id || "").trim();
const chatId = String((item as any).chat_id || "").trim();
const text = String((item as any).text || "").trim();
const source = String((item as any).source || "").trim();
if (!id || !chatId || !text) continue;
let ok = true;
let err = "";
try {
const sent = await sock.sendMessage(chatId, { text });
const sendContent = buildOutboundSendContent(text, source, sock, chatId);
const sent = await sock.sendMessage(chatId, sendContent);
const stanzaId = String((sent as any)?.key?.id || "").trim();
log(`outbound sent id=${id} chat=${chatId} stanza=${stanzaId || "?"}`);
log(
`outbound sent id=${id} chat=${chatId} stanza=${stanzaId || "?"} mentions=${(sendContent.mentions || []).length} mention0=${(sendContent.mentions || [])[0] || ""} text=${sendContent.text.slice(0, 80)}`,
);
try {
await fetch(`${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/ack`, {
method: "POST",

View file

@ -356,9 +356,178 @@ def build_group_sender_context(*, metadata: dict[str, Any] | None, external_user
push_name = str(raw.get("pushName") or meta.get("push_name") or meta.get("display_name") or "").strip()
sender = str(external_user_id or "").strip()
label = push_name or sender or "unknown"
if sender and push_name and sender not in push_name:
return f"[群成员: {label} ({sender})]"
return f"[群成员: {label}]"
return f"[发言: {label}]"
def _mention_tokens_for_jids(jids: list[str]) -> list[str]:
tokens: list[str] = []
seen: set[str] = set()
for jid in jids or []:
local = jid_base_local(jid)
if local:
tok = f"@{local}"
if tok.lower() not in seen:
seen.add(tok.lower())
tokens.append(tok)
phone = jid_phone(jid)
if phone and len(phone) >= 6:
tok = f"@{phone}"
if tok.lower() not in seen:
seen.add(tok.lower())
tokens.append(tok)
return tokens
def strip_bot_mentions_from_text(
*,
text: str,
bot_jid: str | None,
metadata: dict[str, Any] | None = None,
) -> str:
body = str(text or "")
identities = _bot_identity_jids(bot_jid=bot_jid, metadata=metadata)
tokens = _mention_tokens_for_jids(identities)
extra_names: set[str] = set()
if isinstance(metadata, dict):
for name in (metadata.get("bot_push_name"), "oliver"):
n = str(name or "").strip()
if n:
extra_names.add(n.lower())
for token in sorted(tokens, key=len, reverse=True):
body = re.sub(
re.escape(token) + r"(?=\s|$|[,。!?!?,.])",
"",
body,
flags=re.IGNORECASE,
)
for name in sorted(extra_names, key=len, reverse=True):
body = re.sub(
rf"@{re.escape(name)}(?=\s|$|[,。!?!?,.])",
"",
body,
flags=re.IGNORECASE,
)
body = re.sub(r"\s{2,}", " ", body).strip()
return body
def _filter_non_bot_mention_jids(
mention_jids: list[str],
*,
bot_jid: str | None,
metadata: dict[str, Any] | None = None,
) -> list[str]:
identities = _bot_identity_jids(bot_jid=bot_jid, metadata=metadata)
out: list[str] = []
seen: set[str] = set()
for raw in mention_jids or []:
jid = str(raw or "").strip()
if not jid:
continue
if any(jids_same_user(jid, identity) for identity in identities):
continue
key = jid_base_local(jid) or jid.lower()
if key in seen:
continue
seen.add(key)
out.append(jid)
return out
def normalize_mentioned_users_in_text(
*,
text: str,
mention_jids: list[str],
mention_names: list[str] | None = None,
store: Any = None,
tenant_id: str = "",
account_id: str = "",
) -> str:
body = str(text or "")
names = [str(x or "").strip() for x in (mention_names or [])]
for idx, jid in enumerate(mention_jids or []):
jid_s = str(jid or "").strip()
if not jid_s:
continue
name = names[idx] if idx < len(names) and names[idx] else ""
if not name and store is not None:
from runtime.scheduler.whatsapp_mentions import _lookup_push_name
name = _lookup_push_name(
store,
tenant_id=str(tenant_id or ""),
account_id=str(account_id or ""),
jid=jid_s,
)
if not name:
local = jid_base_local(jid_s)
if local:
body = re.sub(
rf"@{re.escape(local)}(?=\s|$|[,。!?!?,.])",
"",
body,
)
continue
local = jid_base_local(jid_s)
if local:
body = re.sub(
rf"@{re.escape(local)}(?=\s|$|[,。!?!?,.])",
f"@{name}",
body,
)
phone = jid_phone(jid_s)
if phone and len(phone) >= 6:
body = re.sub(
rf"@{re.escape(phone)}(?=\s|$|[,。!?!?,.])",
f"@{name}",
body,
)
body = re.sub(r"\s{2,}", " ", body).strip()
return body
def prepare_group_user_text_for_model(
*,
text: str,
metadata: dict[str, Any] | None,
mentions: list[str],
bot_jid: str | None,
session_scope: str,
external_user_id: str,
filtered_mention_jids: list[str] | None = None,
mention_names: list[str] | None = None,
store: Any = None,
tenant_id: str = "",
account_id: str = "",
quoted_ctx: str = "",
) -> str:
body = str(text or "").strip()
body = strip_bot_mentions_from_text(text=body, bot_jid=bot_jid, metadata=metadata)
target_jids = (
filtered_mention_jids
if filtered_mention_jids is not None
else _filter_non_bot_mention_jids(list(mentions or []), bot_jid=bot_jid, metadata=metadata)
)
body = normalize_mentioned_users_in_text(
text=body,
mention_jids=target_jids,
mention_names=mention_names,
store=store,
tenant_id=tenant_id,
account_id=account_id,
)
prefix_parts: list[str] = []
if normalize_group_session_scope(session_scope) == "chat":
prefix_parts.append(
build_group_sender_context(metadata=metadata, external_user_id=external_user_id)
)
quote = str(quoted_ctx or "").strip()
if quote:
prefix_parts.append(quote)
prefix = "\n".join(part for part in prefix_parts if str(part).strip())
if prefix and body:
return f"{prefix}\n{body}"
return prefix or body
def build_group_focus_instruction(*, lang: str = "zh") -> str:
@ -408,6 +577,7 @@ def build_whatsapp_group_reply_metadata(
}
if push_name:
out["quote_push_name"] = push_name
out["mention_names"] = [push_name]
return out
@ -424,9 +594,12 @@ __all__ = [
"extract_quoted_ume_alert_text",
"mentions_include_bot",
"metadata_mentions_bot",
"normalize_mentioned_users_in_text",
"normalize_jid",
"normalize_group_session_scope",
"normalize_jids",
"prepare_group_user_text_for_model",
"strip_bot_mentions_from_text",
"infer_is_group_from_chat_id",
"is_nonsend_channel_reply_text",
"resolve_is_group",

View file

@ -6,6 +6,12 @@ from typing import Any
from runtime.orchestration.group_ingest import is_nonsend_channel_reply_text, should_send_channel_reply_text
from runtime.scheduler.session_resolver import parse_delivery_json
from runtime.scheduler.whatsapp_mentions import (
encode_whatsapp_outbound_source,
format_whatsapp_mention_text,
infer_whatsapp_mention_jids_from_text,
normalize_whatsapp_mention_jids,
)
def _encode_weixin_outbound_source(
@ -216,15 +222,56 @@ def deliver_scheduled_reply(
wa.get("account_id") or resolved_account_id or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default"
).strip()
if wa_enabled and chat_id and should_send_channel_reply_text(text):
mention_jids = normalize_whatsapp_mention_jids(wa.get("mention_jids"))
if not mention_jids:
mention_jids = infer_whatsapp_mention_jids_from_text(
text,
store=store,
tenant_id=tenant_id,
account_id=account_id,
)
mention_names = wa.get("mention_names") if isinstance(wa.get("mention_names"), list) else None
if mention_names is None and mention_jids:
from runtime.scheduler.whatsapp_mentions import _lookup_push_name
derived_names: list[str] = []
for jid in mention_jids:
derived_names.append(
_lookup_push_name(
store,
tenant_id=tenant_id,
account_id=account_id,
jid=str(jid or ""),
)
)
if any(derived_names):
mention_names = derived_names
out_text = format_whatsapp_mention_text(
text,
mention_jids,
store=store,
tenant_id=tenant_id,
account_id=account_id,
mention_names=mention_names,
)
msg_id = store.enqueue_channel_outbound_message(
channel="whatsapp",
chat_id=chat_id,
text=text,
text=out_text,
tenant_id=tenant_id,
account_id=account_id,
source="scheduled_job",
source=encode_whatsapp_outbound_source(
mention_jids=mention_jids,
mention_names=mention_names,
mention_text_ready=bool(mention_jids),
),
)
results["whatsapp"] = {"ok": True, "message_id": msg_id, "chat_id": chat_id}
results["whatsapp"] = {
"ok": True,
"message_id": msg_id,
"chat_id": chat_id,
"mention_jids": mention_jids,
}
wx = delivery.get("weixin") if isinstance(delivery.get("weixin"), dict) else {}
wx_enabled = bool(wx.get("enabled", True))

View file

@ -189,8 +189,18 @@ def resolve_scheduled_session(
external_user_id=external_user_id,
session_title=f"Scheduled · {job_name}",
)
if user_id:
store.ensure_ui_session_owner(session_id=session_id, tenant_id=tenant_id, user_id=user_id)
from runtime.application.gateway.channel_session_owner import (
assign_channel_session_to_account_owner,
)
assign_channel_session_to_account_owner(
store,
session_id=session_id,
channel="whatsapp",
account_id=account_id,
tenant_id=tenant_id,
fallback_user_id=user_id,
)
return ResolvedSession(
session_id=session_id,
tenant_id=tenant_id,
@ -219,8 +229,18 @@ def resolve_scheduled_session(
external_user_id=external_user_id,
session_title=f"Scheduled · {job_name}",
)
if user_id:
store.ensure_ui_session_owner(session_id=session_id, tenant_id=tenant_id, user_id=user_id)
from runtime.application.gateway.channel_session_owner import (
assign_channel_session_to_account_owner,
)
assign_channel_session_to_account_owner(
store,
session_id=session_id,
channel=channel,
account_id=account_id,
tenant_id=tenant_id,
fallback_user_id=user_id,
)
return ResolvedSession(
session_id=session_id,
tenant_id=tenant_id,

View file

@ -0,0 +1,313 @@
from __future__ import annotations
import json
import re
from typing import Any
from runtime.extensions.whatsapp.access_control import contact_phone_key
from runtime.extensions.whatsapp.api import normalize_whatsapp_target
def encode_whatsapp_outbound_source(
*,
kind: str = "scheduled_job",
mention_jids: list[str] | None = None,
mention_names: list[str] | None = None,
mention_text_ready: bool = False,
) -> str:
payload: dict[str, Any] = {"kind": str(kind or "scheduled_job")}
jids = normalize_whatsapp_mention_jids(mention_jids)
if jids:
payload["mention_jids"] = jids
names = [str(x or "").strip() for x in (mention_names or []) if str(x or "").strip()]
if names:
payload["mention_names"] = names
if mention_text_ready:
payload["mention_text_ready"] = True
return json.dumps(payload, ensure_ascii=False)
def decode_whatsapp_outbound_source(raw: str) -> dict[str, Any]:
text = str(raw or "").strip()
if not text:
return {}
if text.startswith("{"):
try:
data = json.loads(text)
return data if isinstance(data, dict) else {}
except Exception:
return {}
return {"kind": text}
def normalize_whatsapp_mention_jids(raw: Any) -> list[str]:
items: list[str] = []
if isinstance(raw, str):
items = re.split(r"[\s,;]+", raw)
elif isinstance(raw, list):
for entry in raw:
if isinstance(entry, str):
items.extend(re.split(r"[\s,;]+", entry))
seen: set[str] = set()
out: list[str] = []
for item in items:
jid_raw = str(item or "").strip()
jid = jid_raw
low = jid_raw.lower()
if jid and low.endswith("@lid"):
jid = jid_raw
elif jid and ("@" not in jid_raw):
# Bare digits in group mentions are usually LID local parts, not phone numbers.
if re.fullmatch(r"\d{10,20}", jid_raw):
jid = f"{jid_raw}@lid"
else:
try:
jid = normalize_whatsapp_target(jid_raw)
except Exception:
jid = jid_raw
if not jid or jid in seen:
continue
seen.add(jid)
out.append(jid)
return out
def merge_whatsapp_mention_jids(delivery: dict[str, Any] | None, mention_jids: Any) -> dict[str, Any]:
merged = dict(delivery or {})
jids = normalize_whatsapp_mention_jids(mention_jids)
if not jids:
return merged
wa = merged.get("whatsapp") if isinstance(merged.get("whatsapp"), dict) else {}
wa = dict(wa)
wa["mention_jids"] = jids
merged["whatsapp"] = wa
return merged
def merge_whatsapp_mention_names(delivery: dict[str, Any] | None, mention_names: Any) -> dict[str, Any]:
merged = dict(delivery or {})
names = [str(x or "").strip() for x in (mention_names or []) if str(x or "").strip()]
if not names:
return merged
wa = merged.get("whatsapp") if isinstance(merged.get("whatsapp"), dict) else {}
wa = dict(wa)
wa["mention_names"] = names
merged["whatsapp"] = wa
return merged
def _normalize_digits(value: str) -> str:
return re.sub(r"\D", "", str(value or ""))
def _extract_phoneish_mentions(text: str) -> list[str]:
out: list[str] = []
for m in re.finditer(r"@([+\d][\d\s().-]{5,})", str(text or "")):
digits = _normalize_digits(m.group(1))
if len(digits) >= 6:
out.append(digits)
return out
def extract_whatsapp_mention_names(text: str) -> list[str]:
out: list[str] = []
# Keep it conservative: only plain @name tokens, stop at whitespace/punct.
for m in re.finditer(r"@([^\s@]{1,64})", str(text or "")):
token = str(m.group(1) or "").strip()
if not token:
continue
# Skip phone-ish patterns; handled separately.
if re.fullmatch(r"\+?\d[\d\s().-]{5,}", token):
continue
out.append(token)
# de-dupe while preserving order
seen: set[str] = set()
deduped: list[str] = []
for name in out:
if name not in seen:
seen.add(name)
deduped.append(name)
return deduped
def infer_whatsapp_mention_jids_from_names(
mention_names: list[str] | None,
*,
store: Any = None,
tenant_id: str = "",
account_id: str = "",
) -> list[str]:
names = [str(x or "").strip() for x in (mention_names or []) if str(x or "").strip()]
if not names or store is None:
return []
lister = getattr(store, "list_whatsapp_contacts", None)
if not callable(lister):
return []
try:
contacts = lister(tenant_id=str(tenant_id or ""), account_id=str(account_id or ""), limit=500)
except Exception:
return []
by_push: dict[str, list[str]] = {}
for row in contacts or []:
if not isinstance(row, dict):
continue
jid = str(row.get("external_user_id") or "").strip()
push = str(row.get("push_name") or "").strip()
if jid and push:
by_push.setdefault(push, []).append(jid)
out: list[str] = []
for name in names:
hits = by_push.get(name) or []
# Only accept unique match to avoid pinging wrong person.
if len(hits) == 1:
out.append(hits[0])
return out
def _phone_digits_from_jid(jid: str) -> str:
local = str(jid or "").split("@", 1)[0].strip()
return _normalize_digits(local)
def infer_whatsapp_mention_jids_from_text(
text: str,
*,
store: Any = None,
tenant_id: str = "",
account_id: str = "",
) -> list[str]:
if store is None:
return []
lister = getattr(store, "list_whatsapp_contacts", None)
if not callable(lister):
return []
phones = _extract_phoneish_mentions(text)
if not phones:
return []
try:
contacts = lister(tenant_id=str(tenant_id or ""), account_id=str(account_id or ""), limit=500)
except Exception:
return []
out: list[str] = []
seen: set[str] = set()
for phone in phones:
for row in contacts or []:
if not isinstance(row, dict):
continue
jid = str(row.get("external_user_id") or "").strip()
if not jid or jid in seen:
continue
key = _normalize_digits(contact_phone_key(row))
local = _phone_digits_from_jid(jid)
# Accept exact/suffix matches to handle country code and formatting noise.
if key and (key == phone or key.endswith(phone) or phone.endswith(key)):
seen.add(jid)
out.append(jid)
break
if local and (local == phone or local.endswith(phone) or phone.endswith(local)):
seen.add(jid)
out.append(jid)
break
return out
def _lookup_push_name(store: Any, *, tenant_id: str, account_id: str, jid: str) -> str:
getter = getattr(store, "get_whatsapp_contact", None)
if not callable(getter):
return ""
candidates = [str(jid or "").strip()]
local = str(jid or "").split("@", 1)[0].strip()
if local and f"{local}@lid" not in candidates:
candidates.append(f"{local}@lid")
if local and f"{local}@s.whatsapp.net" not in candidates:
candidates.append(f"{local}@s.whatsapp.net")
for candidate in candidates:
if not candidate:
continue
try:
row = getter(
tenant_id=str(tenant_id or ""),
account_id=str(account_id or ""),
external_user_id=candidate,
)
except Exception:
row = None
if isinstance(row, dict):
push = str(row.get("push_name") or "").strip()
if push:
return push
lister = getattr(store, "list_whatsapp_contacts", None)
if not callable(lister) or not local:
return ""
try:
contacts = lister(tenant_id=str(tenant_id or ""), account_id=str(account_id or ""), limit=500)
except Exception:
return ""
for row in contacts or []:
if not isinstance(row, dict):
continue
eid = str(row.get("external_user_id") or "").strip()
if not eid:
continue
elocal = eid.split("@", 1)[0].strip()
if elocal == local:
push = str(row.get("push_name") or "").strip()
if push:
return push
return ""
def mention_tag_for_jid(jid: str, *, push_name: str = "") -> str:
name = str(push_name or "").strip()
if name:
return f"@{name}"
local = str(jid or "").split("@", 1)[0].strip()
return f"@{local}" if local else ""
def format_whatsapp_mention_text(
text: str,
mention_jids: list[str],
*,
store: Any = None,
tenant_id: str = "",
account_id: str = "",
mention_names: list[str] | None = None,
) -> str:
"""Align outbound text with explicit mention JIDs; visible @ uses display name when known."""
body = str(text or "").strip()
jids = normalize_whatsapp_mention_jids(mention_jids)
if not body or not jids:
return body
names = [str(x or "").strip() for x in (mention_names or [])]
tags: list[str] = []
for idx, jid in enumerate(jids):
push_name = names[idx] if idx < len(names) and names[idx] else ""
if not push_name and store is not None:
push_name = _lookup_push_name(store, tenant_id=tenant_id, account_id=account_id, jid=jid)
tag = mention_tag_for_jid(jid, push_name=push_name)
if tag:
tags.append(tag)
if not tags:
return body
if any(tag in body for tag in tags):
return body
cleaned = re.sub(r"@\+?\d+(?:\s+\d+)*\s*", "", body)
cleaned = re.sub(r"@\S+\s*", "", cleaned)
cleaned = re.sub(r"\s{2,}", " ", cleaned).strip()
prefix = " ".join(tags) + " "
return f"{prefix}{cleaned}" if cleaned else prefix.strip()
__all__ = [
"decode_whatsapp_outbound_source",
"encode_whatsapp_outbound_source",
"extract_whatsapp_mention_names",
"format_whatsapp_mention_text",
"infer_whatsapp_mention_jids_from_names",
"infer_whatsapp_mention_jids_from_text",
"mention_tag_for_jid",
"merge_whatsapp_mention_jids",
"merge_whatsapp_mention_names",
"normalize_whatsapp_mention_jids",
]

View file

@ -22,6 +22,7 @@ class SkillExecutionContext:
session_id: str
lang: str = "zh"
user_text: str = ""
inbound_metadata: dict[str, Any] | None = None
specialist: str = "oclaw"
trace_id: str | None = None
parent_span_id: str | None = None
@ -160,6 +161,7 @@ class SkillExecutor:
session_id=ctx.session_id,
lang=ctx.lang,
user_text=ctx.user_text,
inbound_metadata=ctx.inbound_metadata,
specialist=ctx.specialist,
task_kind="turn",
policy_engine=None,

View file

@ -4,6 +4,7 @@ import json
from typing import Any
from runtime.scheduler.cron_service import build_delivery_for_session
from runtime.scheduler.whatsapp_mentions import merge_whatsapp_mention_jids, merge_whatsapp_mention_names
from runtime.scheduler.expressions import normalize_schedule_kind
from runtime.scheduler.system_timezone import default_system_timezone
from runtime.scheduler.service import run_scheduled_job_now
@ -61,6 +62,14 @@ def schedule_create_tool() -> ToolSpec:
session_id=str(args.get("session_id") or ""),
whatsapp_chat_id=str(args.get("whatsapp_chat_id") or ""),
)
delivery = merge_whatsapp_mention_jids(
delivery,
args.get("whatsapp_mention_jids") if args.get("whatsapp_mention_jids") is not None else None,
)
delivery = merge_whatsapp_mention_names(
delivery,
args.get("whatsapp_mention_names") if args.get("whatsapp_mention_names") is not None else None,
)
interaction_mode = normalize_interaction_mode(
str(args.get("interaction_mode") or "expert")
)
@ -89,7 +98,7 @@ def schedule_create_tool() -> ToolSpec:
return ToolSpec(
name="schedule_create",
description="Create a scheduled job (cron, once, or interval). Delivery follows the current chat channel (WhatsApp vs WeChat) unless delivery is set explicitly.",
description="Create a scheduled job (cron, once, or interval). Delivery follows the current chat channel (WhatsApp vs WeChat) unless delivery is set explicitly. For WhatsApp group @mentions, set whatsapp_mention_jids with explicit JIDs.",
parameters={
"type": "object",
"properties": {
@ -106,6 +115,11 @@ def schedule_create_tool() -> ToolSpec:
"selected_specialist": {"type": "string"},
"lang": {"type": "string"},
"whatsapp_chat_id": {"type": "string"},
"whatsapp_mention_jids": {
"type": "array",
"items": {"type": "string"},
"description": "WhatsApp JIDs to @mention on delivery (e.g. 628...@s.whatsapp.net). Stored in delivery.whatsapp.mention_jids.",
},
"delivery": {"type": "object"},
"description": {"type": "string"},
},

View file

@ -7,7 +7,7 @@ from typing import Any, Mapping
from sqlalchemy import and_, delete, distinct, func, insert, literal, or_, select, update
from sqlalchemy.engine import Engine
from svc.persistence.db.tables import app_user, chat_message, chat_session, ui_session_owner
from svc.persistence.db.tables import app_user, channel_session_v2, chat_message, chat_session, ui_session_owner
from svc.persistence.sqlite_store import ChatSession, SessionsListMeta
@ -197,6 +197,108 @@ class ChatSessionsSaRepository:
uname = str(username or "").strip().lower()
return func.lower(app_user.c.username) == uname
def _administrator_chat_session_ids_subquery(self, *, username: str, tenant_id: str) -> Any:
pred = self._administrator_username_predicate(username)
tid = str(tenant_id or "").strip()
admin_ids = (
select(chat_session.c.id.label("sid"))
.select_from(
chat_session.join(
ui_session_owner,
ui_session_owner.c.session_id == chat_session.c.id,
).join(
app_user,
and_(
app_user.c.id == ui_session_owner.c.user_id,
app_user.c.tenant_id == ui_session_owner.c.tenant_id,
),
)
)
.where(pred)
)
channel_ids = (
select(channel_session_v2.c.session_id.label("sid"))
.where(
channel_session_v2.c.tenant_id == tid,
channel_session_v2.c.session_id.isnot(None),
channel_session_v2.c.session_id != "",
)
)
return admin_ids.union(channel_ids).subquery()
def list_chat_sessions_for_administrator_chat_view(
self,
*,
username: str,
tenant_id: str,
limit: int | None,
offset: int,
) -> list[ChatSession]:
ids = self._administrator_chat_session_ids_subquery(
username=username, tenant_id=tenant_id
)
stmt = (
select(
chat_session.c.id,
chat_session.c.title,
chat_session.c.created_at,
chat_session.c.last_message_at,
)
.where(chat_session.c.id.in_(select(ids.c.sid)))
.order_by(*_activity_order())
)
if limit is not None:
stmt = stmt.limit(int(limit)).offset(int(offset))
with self._engine.connect() as conn:
rows = conn.execute(stmt).mappings().all()
return [_session_from_row(r) for r in rows]
def sessions_list_meta_for_administrator_chat_view(
self, *, username: str, tenant_id: str
) -> SessionsListMeta:
ids = self._administrator_chat_session_ids_subquery(
username=username, tenant_id=tenant_id
)
with self._engine.connect() as conn:
row = conn.execute(
select(
func.count().label("c"),
func.max(
func.coalesce(chat_session.c.last_message_at, chat_session.c.created_at)
).label("latest_activity_at"),
)
.select_from(chat_session)
.where(chat_session.c.id.in_(select(ids.c.sid)))
).mappings().first()
return SessionsListMeta(
session_count=int(row["c"] or 0) if row else 0,
latest_activity_at=str(row["latest_activity_at"])
if row and row.get("latest_activity_at") is not None
else None,
)
def fetch_chat_session_for_administrator_chat_view(
self, *, session_id: str, username: str, tenant_id: str
) -> ChatSession | None:
sid = str(session_id or "").strip()
if not sid:
return None
ids = self._administrator_chat_session_ids_subquery(
username=username, tenant_id=tenant_id
)
with self._engine.connect() as conn:
row = conn.execute(
select(
chat_session.c.id,
chat_session.c.title,
chat_session.c.created_at,
chat_session.c.last_message_at,
)
.where(chat_session.c.id == sid, chat_session.c.id.in_(select(ids.c.sid)))
.limit(1)
).mappings().first()
return _session_from_row(row) if row else None
def list_chat_sessions_for_administrator_username(
self,
*,

View file

@ -1546,6 +1546,14 @@ class SqliteStore(ScheduledJobStoreMixin):
created_at=utc_now_iso(),
)
def set_ui_session_owner(self, *, session_id: str, tenant_id: str, user_id: str) -> None:
self._ui_session_owner_repo().upsert_replace(
session_id=str(session_id),
tenant_id=str(tenant_id),
user_id=str(user_id),
created_at=utc_now_iso(),
)
def get_session(self, session_id: str) -> Optional[ChatSession]:
return self._chat_sessions_repo().fetch_chat_session_by_id(session_id=session_id)
@ -1629,6 +1637,35 @@ class SqliteStore(ScheduledJobStoreMixin):
username=username
)
def list_sessions_for_administrator_chat_view(
self,
*,
username: str,
tenant_id: str,
limit: int | None = None,
offset: int = 0,
) -> list[ChatSession]:
return self._chat_sessions_repo().list_chat_sessions_for_administrator_chat_view(
username=username,
tenant_id=tenant_id,
limit=limit,
offset=int(offset),
)
def get_sessions_list_meta_for_administrator_chat_view(
self, *, username: str, tenant_id: str
) -> SessionsListMeta:
return self._chat_sessions_repo().sessions_list_meta_for_administrator_chat_view(
username=username, tenant_id=tenant_id
)
def get_session_for_administrator_chat_view(
self, *, session_id: str, username: str, tenant_id: str
) -> Optional[ChatSession]:
return self._chat_sessions_repo().fetch_chat_session_for_administrator_chat_view(
session_id=session_id, username=username, tenant_id=tenant_id
)
def get_session_for_administrator_username(
self, *, session_id: str, username: str
) -> Optional[ChatSession]:

View file

@ -13,10 +13,13 @@ from runtime.orchestration.group_ingest import (
is_nonsend_channel_reply_text,
normalize_jid,
normalize_group_session_scope,
normalize_mentioned_users_in_text,
prepare_group_user_text_for_model,
resolve_group_policy,
session_user_key,
should_inject_quoted_context,
should_process_group_inbound,
strip_bot_mentions_from_text,
text_mentions_bot,
)
from interfaces.channels.base import InboundMessage
@ -356,8 +359,56 @@ def test_build_group_sender_context() -> None:
metadata={"raw": {"pushName": "Alice"}},
external_user_id="111@s.whatsapp.net",
)
assert "Alice" in ctx
assert "111@s.whatsapp.net" in ctx
assert ctx == "[发言: Alice]"
assert "111@s.whatsapp.net" not in ctx
def test_strip_bot_mentions_from_text() -> None:
meta = {"bot_lid": "162788605444170@lid", "bot_push_name": "oliver"}
out = strip_bot_mentions_from_text(
text="@162788605444170 生成小鹿照片",
bot_jid="8618142387786@s.whatsapp.net",
metadata=meta,
)
assert out == "生成小鹿照片"
def test_normalize_mentioned_users_in_text_uses_nickname() -> None:
out = normalize_mentioned_users_in_text(
text="每三分钟提醒@200846277140511 喝水",
mention_jids=["200846277140511@lid"],
mention_names=["吴华"],
)
assert out == "每三分钟提醒@吴华 喝水"
def test_prepare_group_user_text_for_model_user_in_chat() -> None:
out = prepare_group_user_text_for_model(
text="@162788605444170 每三分钟提醒@200846277140511 喝水",
metadata={"bot_lid": "162788605444170@lid", "bot_push_name": "oliver"},
mentions=["162788605444170@lid", "200846277140511@lid"],
bot_jid="8618142387786@s.whatsapp.net",
session_scope="user_in_chat",
external_user_id="8618142387786@s.whatsapp.net",
filtered_mention_jids=["200846277140511@lid"],
mention_names=["吴华"],
)
assert out == "每三分钟提醒@吴华 喝水"
assert "[发言:" not in out
def test_prepare_group_user_text_for_model_shared_chat_prefix() -> None:
out = prepare_group_user_text_for_model(
text="@999 帮忙",
metadata={"raw": {"pushName": "Alice"}},
mentions=["999@s.whatsapp.net"],
bot_jid="999@s.whatsapp.net",
session_scope="chat",
external_user_id="111@s.whatsapp.net",
)
assert out.startswith("[发言: Alice]")
assert out.endswith("帮忙")
assert "@999" not in out
def test_build_group_focus_instruction() -> None:
@ -421,6 +472,7 @@ def test_build_whatsapp_group_reply_metadata() -> None:
assert meta["mention_jids"] == ["111:12@s.whatsapp.net"]
assert meta["quote_text"] == "明天几点?"
assert meta["quote_participant"] == "111:12@s.whatsapp.net"
assert meta["mention_names"] == ["Alice"]
def test_default_group_session_is_per_user(fresh_sqlite_store: SqliteStore) -> None:
@ -608,7 +660,7 @@ def test_inbound_group_mention_uses_per_user_session_and_sender_prefix(
{
**base,
"user_id": "111@s.whatsapp.net",
"text": "@bot hi",
"text": "@999 hi",
"mentions": ["999@s.whatsapp.net"],
"metadata": {**base["metadata"], "raw": {"pushName": "Alice"}},
}
@ -617,7 +669,7 @@ def test_inbound_group_mention_uses_per_user_session_and_sender_prefix(
{
**base,
"user_id": "222@s.whatsapp.net",
"text": "@bot again",
"text": "@999 again",
"mentions": ["999@s.whatsapp.net"],
"metadata": {**base["metadata"], "raw": {"pushName": "Bob"}},
}
@ -625,9 +677,10 @@ def test_inbound_group_mention_uses_per_user_session_and_sender_prefix(
assert len(session_ids) == 2
assert session_ids[0] != session_ids[1]
assert "[群成员:" in captured["text"]
assert "群聊规则" in captured["text"]
assert "Bob" in captured["text"]
assert "[群成员:" not in captured["text"]
assert "[发言:" not in captured["text"]
assert captured["text"] == "again"
assert "群聊规则" not in captured["text"]
sid = store.get_or_create_channel_session_v2(
tenant_id=tenant_id,
@ -882,3 +935,112 @@ def test_inbound_group_reply_includes_quote_and_mention_metadata(
assert meta.get("quote_stanza_id") == "ABC123"
assert meta.get("mention_jids") == ["111@s.whatsapp.net"]
assert meta.get("quote_text") == "@bot 明天几点?"
def test_inbound_group_schedule_mention_metadata_preserved_after_text_clean(
monkeypatch: pytest.MonkeyPatch, fresh_sqlite_store: SqliteStore
) -> None:
store = fresh_sqlite_store
tenant_id, _ = _setup_whatsapp_identity(store)
store.upsert_whatsapp_contact(
tenant_id=tenant_id,
account_id="wa-default",
external_user_id="200846277140511@lid",
push_name="WuHua",
list_type="whitelist",
)
monkeypatch.setattr("svc.persistence.assistant_store.get_assistant_store", lambda: store)
captured: dict[str, object] = {}
class _Turn:
turn_uuid = "turn-sched"
reply_text = "ok"
class _Gw:
def __init__(self, *, store: object) -> None:
_ = store
def handle_turn(self, **kwargs: object) -> _Turn:
msg = kwargs.get("msg")
captured["text"] = str(getattr(msg, "text", "") or "")
captured["metadata"] = getattr(msg, "metadata", None)
return _Turn()
monkeypatch.setattr("runtime.gateway.OclawGateway", _Gw)
process_inbound_payload(
{
"channel": "whatsapp",
"account_id": "wa-default",
"user_id": "111@s.whatsapp.net",
"chat_id": "120363012345678@g.us",
"text": "@999 remind @200846277140511 water",
"is_group": True,
"mentions": ["999@s.whatsapp.net", "200846277140511@lid"],
"metadata": {
"bot_jid": "999@s.whatsapp.net",
"bot_lid": "162788605444170@lid",
"raw": {"pushName": "Alice"},
},
}
)
meta = captured.get("metadata") if isinstance(captured.get("metadata"), dict) else {}
assert "200846277140511@lid" in list(meta.get("mentioned_jids") or [])
assert meta.get("raw_inbound_text")
text = str(captured.get("text") or "")
assert "200846277140511" not in text
assert "@WuHua" in text
assert "@999" not in text
def test_channel_group_session_visible_in_administrator_chat_list(
monkeypatch: pytest.MonkeyPatch, fresh_sqlite_store: SqliteStore
) -> None:
store = fresh_sqlite_store
tenant_id, owner_user_id = _setup_whatsapp_identity(store, extra_user_ids=["222@s.whatsapp.net"])
monkeypatch.setattr("svc.persistence.assistant_store.get_assistant_store", lambda: store)
session_ids: list[str] = []
class _Turn:
turn_uuid = "turn-guest"
reply_text = "ok"
class _Gw:
def __init__(self, *, store: object) -> None:
_ = store
def handle_turn(self, **kwargs: object) -> _Turn:
msg = kwargs.get("msg")
session_ids.append(str(getattr(msg, "session_id", "") or ""))
return _Turn()
monkeypatch.setattr("runtime.gateway.OclawGateway", _Gw)
process_inbound_payload(
{
"channel": "whatsapp",
"account_id": "wa-default",
"user_id": "222@s.whatsapp.net",
"chat_id": "120363012345678@g.us",
"text": "@999 hello from bob",
"is_group": True,
"mentions": ["999@s.whatsapp.net"],
"metadata": {"bot_jid": "999@s.whatsapp.net", "raw": {"pushName": "Bob"}},
}
)
assert session_ids
owner = store.get_ui_session_owner(session_id=session_ids[0]) or {}
assert owner.get("user_id") == owner_user_id
listed = store.list_sessions_for_administrator_chat_view(
username="administrator",
tenant_id=tenant_id,
limit=20,
offset=0,
)
listed_ids = {s.id for s in listed}
assert session_ids[0] in listed_ids

View file

@ -0,0 +1,74 @@
from __future__ import annotations
import os
import tempfile
import unittest
from pathlib import Path
from runtime.scheduler.whatsapp_mentions import format_whatsapp_mention_text, merge_whatsapp_mention_jids, normalize_whatsapp_mention_jids
from svc.persistence.assistant_store import reset_assistant_store_singleton
from svc.persistence.sqlite_store import SqliteStore
class WhatsappMentionJidFieldTests(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory(ignore_cleanup_errors=True)
self.db = Path(self._tmp.name) / "wa_names.sqlite"
os.environ["OPS_ASSISTANT_DB_PATH"] = str(self.db)
os.environ["AIA_ASSISTANT_DB_BACKEND"] = "sqlite"
reset_assistant_store_singleton()
self.store = SqliteStore(str(self.db))
tenant = self.store.create_tenant("Team")
self.tenant_id = str(tenant["id"])
def tearDown(self) -> None:
reset_assistant_store_singleton()
self._tmp.cleanup()
def test_normalize_whatsapp_mention_jids(self) -> None:
out = normalize_whatsapp_mention_jids(
["6281@s.whatsapp.net", "6282@s.whatsapp.net, 6283@s.whatsapp.net"]
)
self.assertEqual(
out,
["6281@s.whatsapp.net", "6282@s.whatsapp.net", "6283@s.whatsapp.net"],
)
def test_merge_whatsapp_mention_jids_into_delivery(self) -> None:
delivery = merge_whatsapp_mention_jids(
{"whatsapp": {"enabled": True, "chat_id": "120@g.us"}},
["6281@s.whatsapp.net"],
)
self.assertEqual(delivery["whatsapp"]["mention_jids"], ["6281@s.whatsapp.net"])
def test_format_whatsapp_mention_text_uses_push_name(self) -> None:
out = format_whatsapp_mention_text(
"💧 @+1 33346537541839 该喝水啦",
["333465375410398@lid"],
store=self.store,
tenant_id=self.tenant_id,
account_id="wa-default",
mention_names=["Egista Hadi Putranto"],
)
self.assertEqual(out, "@Egista Hadi Putranto 💧 该喝水啦")
def test_format_whatsapp_mention_text_looks_up_push_name_by_jid(self) -> None:
self.store.upsert_whatsapp_contact(
tenant_id=self.tenant_id,
account_id="wa-default",
external_user_id="333465375410398@lid",
push_name="Egista Hadi Putranto",
phone="33346537541839",
list_type="whitelist",
)
out = format_whatsapp_mention_text(
"该喝水啦",
["333465375410398@lid"],
store=self.store,
tenant_id=self.tenant_id,
account_id="wa-default",
)
self.assertEqual(out, "@Egista Hadi Putranto 该喝水啦")
if __name__ == "__main__":
unittest.main()

View file

@ -0,0 +1,116 @@
from __future__ import annotations
import os
import tempfile
import unittest
from pathlib import Path
from runtime.scheduler.channel_delivery import deliver_scheduled_reply
from svc.persistence.assistant_store import reset_assistant_store_singleton
from svc.persistence.sqlite_store import SqliteStore
class WhatsappScheduledMentionTests(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory(ignore_cleanup_errors=True)
self.db = Path(self._tmp.name) / "wa_mentions.sqlite"
os.environ["OPS_ASSISTANT_DB_PATH"] = str(self.db)
os.environ["AIA_ASSISTANT_DB_BACKEND"] = "sqlite"
reset_assistant_store_singleton()
self.store = SqliteStore(str(self.db))
tenant = self.store.create_tenant("Team")
self.tenant_id = str(tenant["id"])
self.store.upsert_whatsapp_contact(
tenant_id=self.tenant_id,
account_id="wa-default",
external_user_id="333465375410398@lid",
push_name="Egista Hadi Putranto",
phone="33346537541839",
list_type="whitelist",
)
def tearDown(self) -> None:
reset_assistant_store_singleton()
self._tmp.cleanup()
def test_deliver_scheduled_reply_encodes_explicit_mention_jids(self) -> None:
delivery = {
"whatsapp": {
"enabled": True,
"target_type": "group",
"chat_id": "120363012345678@g.us",
"account_id": "wa-default",
"mention_jids": ["333465375410398@lid"],
},
"weixin": {"enabled": False},
}
reply_text = "💧 @+1 33346537541839 该喝水啦!记得保持水分哦~"
result = deliver_scheduled_reply(
self.store,
tenant_id=self.tenant_id,
reply_text=reply_text,
delivery_json=__import__("json").dumps(delivery),
)
self.assertTrue(result.get("ok"), result)
pending = self.store.list_pending_channel_outbound_messages(
channel="whatsapp",
account_id="wa-default",
limit=5,
)
self.assertEqual(len(pending), 1)
source = __import__("json").loads(str(pending[0].get("source") or "{}"))
self.assertEqual(source.get("mention_jids"), ["333465375410398@lid"])
out_text = str(pending[0].get("text") or "")
self.assertTrue(out_text.startswith("@Egista Hadi Putranto"))
self.assertIn("该喝水啦", out_text)
self.assertNotIn("@+1", out_text)
self.assertTrue(source.get("mention_text_ready"))
self.assertEqual(source.get("mention_names"), ["Egista Hadi Putranto"])
def test_deliver_scheduled_reply_infers_mentions_from_phoneish_text(self) -> None:
self.store.upsert_whatsapp_contact(
tenant_id=self.tenant_id,
account_id="wa-default",
external_user_id="200846277140511@lid",
push_name="吴华",
phone="0846277140511",
list_type="whitelist",
)
delivery = {
"whatsapp": {
"enabled": True,
"target_type": "group",
"chat_id": "120363012345678@g.us",
"account_id": "wa-default",
},
"weixin": {"enabled": False},
}
reply_text = "💧 @+20 0846277140511 该喝水啦!快去补充一下水分吧~ 💧"
result = deliver_scheduled_reply(
self.store,
tenant_id=self.tenant_id,
reply_text=reply_text,
delivery_json=__import__("json").dumps(delivery),
)
self.assertTrue(result.get("ok"), result)
pending = self.store.list_pending_channel_outbound_messages(
channel="whatsapp",
account_id="wa-default",
limit=5,
)
self.assertEqual(len(pending), 1)
source = __import__("json").loads(str(pending[0].get("source") or "{}"))
self.assertEqual(source.get("mention_jids"), ["200846277140511@lid"])
out_text = str(pending[0].get("text") or "")
self.assertTrue(out_text.startswith("@吴华"))
self.assertIn("该喝水啦", out_text)
self.assertNotIn("@+20", out_text)
def test_normalize_bare_digits_to_lid(self) -> None:
from runtime.scheduler.whatsapp_mentions import normalize_whatsapp_mention_jids
self.assertEqual(normalize_whatsapp_mention_jids(["200846277140511"]), ["200846277140511@lid"])
if __name__ == "__main__":
unittest.main()