diff --git a/interfaces/admin/chat_api.py b/interfaces/admin/chat_api.py index b0703780..60b9a2ca 100644 --- a/interfaces/admin/chat_api.py +++ b/interfaces/admin/chat_api.py @@ -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) diff --git a/runtime/agent_core_attempt.py b/runtime/agent_core_attempt.py index eac6e2b1..d8880909 100644 --- a/runtime/agent_core_attempt.py +++ b/runtime/agent_core_attempt.py @@ -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, diff --git a/runtime/application/gateway/channel_session_owner.py b/runtime/application/gateway/channel_session_owner.py new file mode 100644 index 00000000..cd39e018 --- /dev/null +++ b/runtime/application/gateway/channel_session_owner.py @@ -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", +] diff --git a/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index f46a9f62..21a9768e 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -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( diff --git a/runtime/chat/tool_runtime.py b/runtime/chat/tool_runtime.py index bbecb947..5df1c4f8 100644 --- a/runtime/chat/tool_runtime.py +++ b/runtime/chat/tool_runtime.py @@ -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) diff --git a/runtime/direct_loop.py b/runtime/direct_loop.py index e9df7bec..5233eb46 100644 --- a/runtime/direct_loop.py +++ b/runtime/direct_loop.py @@ -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), diff --git a/runtime/gateway.py b/runtime/gateway.py index d31ebefd..13ba5ed0 100644 --- a/runtime/gateway.py +++ b/runtime/gateway.py @@ -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", diff --git a/runtime/operations/whatsapp_bridge/baileys_runner.ts b/runtime/operations/whatsapp_bridge/baileys_runner.ts index ff4eced5..a375bd68 100644 --- a/runtime/operations/whatsapp_bridge/baileys_runner.ts +++ b/runtime/operations/whatsapp_bridge/baileys_runner.ts @@ -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 | 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(); + + 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 | 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 | 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 | 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 => { @@ -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): 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", diff --git a/runtime/orchestration/group_ingest.py b/runtime/orchestration/group_ingest.py index b1c85e39..6d138fe0 100644 --- a/runtime/orchestration/group_ingest.py +++ b/runtime/orchestration/group_ingest.py @@ -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", diff --git a/runtime/scheduler/channel_delivery.py b/runtime/scheduler/channel_delivery.py index 389106d4..1b3fbd37 100644 --- a/runtime/scheduler/channel_delivery.py +++ b/runtime/scheduler/channel_delivery.py @@ -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)) diff --git a/runtime/scheduler/session_resolver.py b/runtime/scheduler/session_resolver.py index f327e687..f25aee57 100644 --- a/runtime/scheduler/session_resolver.py +++ b/runtime/scheduler/session_resolver.py @@ -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, diff --git a/runtime/scheduler/whatsapp_mentions.py b/runtime/scheduler/whatsapp_mentions.py new file mode 100644 index 00000000..414ff05e --- /dev/null +++ b/runtime/scheduler/whatsapp_mentions.py @@ -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", +] diff --git a/runtime/skill_executor.py b/runtime/skill_executor.py index 5303975b..7d61563a 100644 --- a/runtime/skill_executor.py +++ b/runtime/skill_executor.py @@ -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, diff --git a/runtime/tools/experts/productivity/schedule_tools.py b/runtime/tools/experts/productivity/schedule_tools.py index 4d66e7c7..d1b67e93 100644 --- a/runtime/tools/experts/productivity/schedule_tools.py +++ b/runtime/tools/experts/productivity/schedule_tools.py @@ -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"}, }, diff --git a/svc/persistence/sa_repos/chat_sessions.py b/svc/persistence/sa_repos/chat_sessions.py index dc9c2595..31d4fbf2 100644 --- a/svc/persistence/sa_repos/chat_sessions.py +++ b/svc/persistence/sa_repos/chat_sessions.py @@ -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, *, diff --git a/svc/persistence/sqlite_store.py b/svc/persistence/sqlite_store.py index 20332eae..32bd61c7 100644 --- a/svc/persistence/sqlite_store.py +++ b/svc/persistence/sqlite_store.py @@ -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]: diff --git a/tests/test_group_ingest.py b/tests/test_group_ingest.py index fafe0a8b..9c929b75 100644 --- a/tests/test_group_ingest.py +++ b/tests/test_group_ingest.py @@ -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 diff --git a/tests/test_whatsapp_mention_jid_field.py b/tests/test_whatsapp_mention_jid_field.py new file mode 100644 index 00000000..99fc4704 --- /dev/null +++ b/tests/test_whatsapp_mention_jid_field.py @@ -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() diff --git a/tests/test_whatsapp_scheduled_mentions.py b/tests/test_whatsapp_scheduled_mentions.py new file mode 100644 index 00000000..fcb4112d --- /dev/null +++ b/tests/test_whatsapp_scheduled_mentions.py @@ -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()