refactor(msteams): isolate thread context routing

This commit is contained in:
Vincent Koc
2026-07-29 16:55:51 +08:00
parent 29e98dd96d
commit 12d3a2f2f6
2 changed files with 267 additions and 209 deletions
@@ -13,10 +13,7 @@ import {
} from "openclaw/plugin-sdk/channel-inbound";
import { fanInChannelIngressLifecycles } from "openclaw/plugin-sdk/channel-ingress-runtime";
import { bindIngressLifecycleToReplyOptions } from "openclaw/plugin-sdk/channel-outbound";
import {
filterSupplementalContextItems,
resolveChannelContextVisibilityMode,
} from "openclaw/plugin-sdk/context-visibility-runtime";
import { resolveChannelContextVisibilityMode } from "openclaw/plugin-sdk/context-visibility-runtime";
import {
DEFAULT_GROUP_HISTORY_LIMIT,
createChannelHistoryWindow,
@@ -24,31 +21,15 @@ import {
} from "openclaw/plugin-sdk/reply-history";
import { sliceUtf16Safe, truncateUtf16Safe } from "openclaw/plugin-sdk/text-utility-runtime";
import { formatUnknownError } from "../errors.js";
import {
fetchChannelMessage,
fetchChatMessageText,
fetchThreadReplies,
formatThreadContext,
type GraphThreadMessage,
} from "../graph-thread.js";
import { normalizeMSTeamsConversationId, parseMSTeamsActivityTimestamp } from "../inbound.js";
import type { MSTeamsMessageHandlerDeps } from "../monitor-handler.types.js";
import type { MSTeamsIngressLifecycle } from "../msteams-ingress.js";
import { resolveMSTeamsAllowlistMatch, resolveMSTeamsReplyPolicy } from "../policy.js";
import { extractMSTeamsPollVote } from "../polls.js";
import { createMSTeamsReplyDispatcher } from "../reply-dispatcher.js";
import { createMSTeamsInboundDeadline, withMSTeamsRequestDeadline } from "../request-timeout.js";
import { getMSTeamsRuntime } from "../runtime.js";
import type { MSTeamsTurnContext } from "../sdk-types.js";
import { recordMSTeamsSentMessage } from "../sent-message-cache.js";
import { resolveTeamGroupId } from "../team-identity.js";
import {
fetchParentMessageCached,
formatParentContextEvent,
markParentContextInjected,
shouldInjectParentContext,
summarizeParentMessage,
} from "../thread-parent-context.js";
import { admitMSTeamsMessage } from "./access.js";
import { prepareMSTeamsInboundContent } from "./inbound-content.js";
import {
@@ -56,7 +37,7 @@ import {
prepareMSTeamsDebounceEntry,
type MSTeamsDebounceEntry,
} from "./inbound-facts.js";
import { resolveMSTeamsRouteSessionKey } from "./thread-session.js";
import { prepareMSTeamsThreadRouting, resolveMSTeamsThreadContext } from "./thread-context.js";
export function createMSTeamsMessageHandler(deps: MSTeamsMessageHandlerDeps) {
const {
@@ -120,9 +101,6 @@ export function createMSTeamsMessageHandler(deps: MSTeamsMessageHandlerDeps) {
const historyBody = [text, formatMediaPlaceholderText(advertisedMedia)]
.filter(Boolean)
.join("\n");
let quoteSenderId: string | undefined;
let quoteSenderName: string | undefined;
log.info("received message", {
rawText: truncateUtf16Safe(rawText, 50),
text: truncateUtf16Safe(text, 50),
@@ -202,27 +180,18 @@ export function createMSTeamsMessageHandler(deps: MSTeamsMessageHandlerDeps) {
: `msteams:group:${conversationId}`;
const teamsTo = isDirectMessage ? `user:${senderId}` : `conversation:${conversationId}`;
const route = core.channel.routing.resolveAgentRoute({
const threadRouting = prepareMSTeamsThreadRouting({
cfg,
channel: "msteams",
teamId,
peer: {
kind: isDirectMessage ? "direct" : isChannel ? "channel" : "group",
id: isDirectMessage ? senderId : conversationId,
},
});
// Isolate channel thread sessions: each thread gets its own session key so
// context does not bleed across threads. Prefer conversationMessageId (the
// ;messageid= portion of conversation.id, i.e. the thread root) over
// activity.replyToId (which may point to a non-root parent in deep threads).
// DMs and group chats are unaffected — only channel thread replies fork.
route.sessionKey = resolveMSTeamsRouteSessionKey({
baseSessionKey: route.sessionKey,
context,
isDirectMessage,
isChannel,
conversationMessageId,
replyToId: activity.replyToId,
senderId,
conversationId,
conversationMessageId: conversationMessageId ?? undefined,
teamId,
log,
});
const { route, deadline: preprocessingDeadline } = threadRouting;
const preview = sliceUtf16Safe(rawBody.replace(/\s+/g, " "), 0, 160);
const inboundLabel = isDirectMessage
@@ -284,29 +253,6 @@ export function createMSTeamsMessageHandler(deps: MSTeamsMessageHandlerDeps) {
return;
}
}
const preprocessingDeadline = createMSTeamsInboundDeadline();
let teamAadGroupId = activity.channelData?.team?.aadGroupId?.trim() || undefined;
const conversationTeamId = isChannel ? teamId : undefined;
let teamGroupIdPromise: Promise<string | undefined> | undefined;
const resolveChannelTeamGroupId = async (): Promise<string | undefined> => {
if (!conversationTeamId) {
return undefined;
}
teamGroupIdPromise ??= resolveTeamGroupId({
conversationTeamId,
aadGroupId: teamAadGroupId,
getTeamDetails: context.getTeamDetails,
deadline: preprocessingDeadline,
}).catch((err: unknown) => {
log.debug?.("failed to resolve Teams AAD group ID", {
teamId: conversationTeamId,
error: formatUnknownError(err),
});
return undefined;
});
teamAadGroupId = await teamGroupIdPromise;
return teamAadGroupId;
};
const content = await prepareMSTeamsInboundContent({
entry: params,
rawBody,
@@ -315,8 +261,8 @@ export function createMSTeamsMessageHandler(deps: MSTeamsMessageHandlerDeps) {
conversationType,
conversationId,
conversationMessageId: conversationMessageId ?? undefined,
teamAadGroupId,
resolveTeamAadGroupId: resolveChannelTeamGroupId,
teamAadGroupId: threadRouting.getTeamAadGroupId(),
resolveTeamAadGroupId: threadRouting.resolveTeamAadGroupId,
mediaMaxBytes,
tokenProvider,
mediaAllowHosts: msteamsCfg?.mediaAllowHosts,
@@ -330,149 +276,22 @@ export function createMSTeamsMessageHandler(deps: MSTeamsMessageHandlerDeps) {
}
const { agentBody, inboundMedia } = content;
enqueuePrimaryMessageSystemEvent();
teamAadGroupId = await resolveChannelTeamGroupId();
// Media is the primary payload, so optional quote enrichment only gets the
// remaining preprocessing budget. DMs alone may fetch the full quote: group
// and channel quotes retain their visibility-filtered preview.
let quoteBodyFull: string | undefined;
const quoteMessageId = quoteInfo?.id;
if (quoteMessageId && isDirectMessage && conversationId.startsWith("19:")) {
try {
const graphToken = await withMSTeamsRequestDeadline({
deadline: preprocessingDeadline,
label: "MS Teams quote token",
work: () => tokenProvider.getAccessToken("https://graph.microsoft.com"),
});
quoteBodyFull = await withMSTeamsRequestDeadline({
deadline: preprocessingDeadline,
label: "MS Teams quote lookup",
work: () =>
fetchChatMessageText(graphToken, conversationId, quoteMessageId, preprocessingDeadline),
});
} catch (err) {
log.debug?.("failed to fetch full quoted message text", {
error: formatUnknownError(err),
});
}
}
// Fetch thread history when the message is a reply inside a Teams channel thread.
// This is a best-effort enhancement; errors are logged and do not block the reply.
//
// We also enqueue a compact `Replying to @sender: …` system event when the parent
// is resolvable. On brand-new thread sessions (see PR #62713), this gives the agent
// immediate parent context even before the fuller `[Thread history]` block is assembled.
// Parent fetches are cached (5 min LRU, 100 entries) and per-session deduped so
// consecutive replies in the same thread do not re-inject identical context.
let threadContext: string | undefined;
const threadParentId = activity.replyToId;
const channelGroupId = teamAadGroupId;
if (threadParentId && isChannel && channelGroupId) {
try {
const graphToken = await withMSTeamsRequestDeadline({
deadline: preprocessingDeadline,
label: "MS Teams thread token",
work: () => tokenProvider.getAccessToken("https://graph.microsoft.com"),
});
// Use allSettled so a failure in one fetch does not discard the other.
// For example, reply-fetch 403 should not throw away a successful parent fetch.
const [parentResult, repliesResult] = await withMSTeamsRequestDeadline({
deadline: preprocessingDeadline,
label: "MS Teams thread history",
work: () =>
Promise.allSettled([
fetchParentMessageCached(
graphToken,
channelGroupId,
conversationId,
threadParentId,
(token, groupId, requestedChannelId, messageId) =>
fetchChannelMessage(
token,
groupId,
requestedChannelId,
messageId,
preprocessingDeadline,
),
),
fetchThreadReplies(
graphToken,
channelGroupId,
conversationId,
threadParentId,
50,
preprocessingDeadline,
),
]),
});
const parentMsg = parentResult.status === "fulfilled" ? parentResult.value : undefined;
const replies = repliesResult.status === "fulfilled" ? repliesResult.value : [];
if (parentResult.status === "rejected") {
log.debug?.("failed to fetch parent message", {
error: formatUnknownError(parentResult.reason),
});
}
if (repliesResult.status === "rejected") {
log.debug?.("failed to fetch thread replies", {
error: formatUnknownError(repliesResult.reason),
});
}
const isThreadSenderAllowed = (msg: GraphThreadMessage) =>
resolveInboundSupplementalSenderAllowed({
isGroup: isChannel,
groupPolicy,
allowFrom: effectiveGroupAllowFrom,
isSenderAllowed: (allowFrom) =>
resolveMSTeamsAllowlistMatch({
allowFrom,
senderId: msg.from?.user?.id ?? "",
senderName: msg.from?.user?.displayName,
allowNameMatching,
}).allowed,
});
const parentSummary = summarizeParentMessage(parentMsg);
const visibleParentMessages = parentMsg
? filterSupplementalContextItems({
items: [parentMsg],
mode: contextVisibilityMode,
kind: "thread",
isSenderAllowed: isThreadSenderAllowed,
}).items
: [];
if (
parentSummary &&
visibleParentMessages.length > 0 &&
shouldInjectParentContext(route.sessionKey, threadParentId)
) {
core.system.enqueueSystemEvent(formatParentContextEvent(parentSummary), {
sessionKey: route.sessionKey,
contextKey: `msteams:thread-parent:${conversationId}:${threadParentId}`,
});
markParentContextInjected(route.sessionKey, threadParentId);
}
const allMessages = parentMsg ? [parentMsg, ...replies] : replies;
quoteSenderId = parentMsg?.from?.user?.id ?? parentMsg?.from?.application?.id ?? undefined;
quoteSenderName =
parentMsg?.from?.user?.displayName ??
parentMsg?.from?.application?.displayName ??
quoteInfo?.sender;
const { items: threadMessages } = filterSupplementalContextItems({
items: allMessages,
mode: contextVisibilityMode,
kind: "thread",
isSenderAllowed: isThreadSenderAllowed,
});
const formatted = formatThreadContext(threadMessages, activity.id);
if (formatted) {
threadContext = formatted;
}
} catch (err) {
log.debug?.("failed to fetch thread history", { error: formatUnknownError(err) });
// Graceful degradation: thread history is an optional enhancement.
}
}
quoteSenderName ??= quoteInfo?.sender;
const { teamAadGroupId, quoteBodyFull, quoteSenderId, quoteSenderName, threadContext } =
await resolveMSTeamsThreadContext({
routing: threadRouting,
context,
tokenProvider,
quoteInfo,
isDirectMessage,
isChannel,
conversationId,
contextVisibilityMode,
groupPolicy,
effectiveGroupAllowFrom,
allowNameMatching,
log,
});
const envelopeFrom = isDirectMessage ? senderName : conversationType;
const buildEnvelope = createChannelInboundEnvelopeBuilder({ cfg, route });
@@ -0,0 +1,239 @@
// Msteams plugin module owns thread routing and Graph parent context.
import { resolveInboundSupplementalSenderAllowed } from "openclaw/plugin-sdk/channel-inbound";
import { filterSupplementalContextItems } from "openclaw/plugin-sdk/context-visibility-runtime";
import type { OpenClawConfig } from "../../runtime-api.js";
import { formatUnknownError } from "../errors.js";
import {
fetchChannelMessage,
fetchChatMessageText,
fetchThreadReplies,
formatThreadContext,
type GraphThreadMessage,
} from "../graph-thread.js";
import type { extractMSTeamsQuoteInfo } from "../inbound.js";
import type { MSTeamsMessageHandlerDeps } from "../monitor-handler.types.js";
import { resolveMSTeamsAllowlistMatch } from "../policy.js";
import { createMSTeamsInboundDeadline, withMSTeamsRequestDeadline } from "../request-timeout.js";
import { getMSTeamsRuntime } from "../runtime.js";
import type { MSTeamsTurnContext } from "../sdk-types.js";
import { resolveTeamGroupId } from "../team-identity.js";
import {
fetchParentMessageCached,
formatParentContextEvent,
markParentContextInjected,
shouldInjectParentContext,
summarizeParentMessage,
} from "../thread-parent-context.js";
import { resolveMSTeamsRouteSessionKey } from "./thread-session.js";
export function prepareMSTeamsThreadRouting(params: {
cfg: OpenClawConfig;
context: MSTeamsTurnContext;
isDirectMessage: boolean;
isChannel: boolean;
senderId: string;
conversationId: string;
conversationMessageId?: string;
teamId?: string;
log: MSTeamsMessageHandlerDeps["log"];
}) {
const core = getMSTeamsRuntime();
const route = core.channel.routing.resolveAgentRoute({
cfg: params.cfg,
channel: "msteams",
teamId: params.teamId,
peer: {
kind: params.isDirectMessage ? "direct" : params.isChannel ? "channel" : "group",
id: params.isDirectMessage ? params.senderId : params.conversationId,
},
});
route.sessionKey = resolveMSTeamsRouteSessionKey({
baseSessionKey: route.sessionKey,
isChannel: params.isChannel,
conversationMessageId: params.conversationMessageId,
replyToId: params.context.activity.replyToId,
});
const deadline = createMSTeamsInboundDeadline();
let teamAadGroupId = params.context.activity.channelData?.team?.aadGroupId?.trim() || undefined;
const conversationTeamId = params.isChannel ? params.teamId : undefined;
let teamGroupIdPromise: Promise<string | undefined> | undefined;
const resolveTeamAadGroupId = async (): Promise<string | undefined> => {
if (!conversationTeamId) {
return undefined;
}
teamGroupIdPromise ??= resolveTeamGroupId({
conversationTeamId,
aadGroupId: teamAadGroupId,
getTeamDetails: params.context.getTeamDetails,
deadline,
}).catch((err: unknown) => {
params.log.debug?.("failed to resolve Teams AAD group ID", {
teamId: conversationTeamId,
error: formatUnknownError(err),
});
return undefined;
});
teamAadGroupId = await teamGroupIdPromise;
return teamAadGroupId;
};
return {
route,
deadline,
resolveTeamAadGroupId,
getTeamAadGroupId: () => teamAadGroupId,
};
}
export async function resolveMSTeamsThreadContext(params: {
routing: ReturnType<typeof prepareMSTeamsThreadRouting>;
context: MSTeamsTurnContext;
tokenProvider: MSTeamsMessageHandlerDeps["tokenProvider"];
quoteInfo: ReturnType<typeof extractMSTeamsQuoteInfo>;
isDirectMessage: boolean;
isChannel: boolean;
conversationId: string;
contextVisibilityMode: "all" | "allowlist" | "allowlist_quote";
groupPolicy: Parameters<typeof resolveInboundSupplementalSenderAllowed>[0]["groupPolicy"];
effectiveGroupAllowFrom: Parameters<
typeof resolveInboundSupplementalSenderAllowed
>[0]["allowFrom"];
allowNameMatching: boolean;
log: MSTeamsMessageHandlerDeps["log"];
}) {
const core = getMSTeamsRuntime();
const activity = params.context.activity;
const { route, deadline } = params.routing;
const teamAadGroupId = await params.routing.resolveTeamAadGroupId();
let quoteBodyFull: string | undefined;
let quoteSenderId: string | undefined;
let quoteSenderName: string | undefined;
const quoteMessageId = params.quoteInfo?.id;
if (quoteMessageId && params.isDirectMessage && params.conversationId.startsWith("19:")) {
try {
const graphToken = await withMSTeamsRequestDeadline({
deadline,
label: "MS Teams quote token",
work: () => params.tokenProvider.getAccessToken("https://graph.microsoft.com"),
});
quoteBodyFull = await withMSTeamsRequestDeadline({
deadline,
label: "MS Teams quote lookup",
work: () =>
fetchChatMessageText(graphToken, params.conversationId, quoteMessageId, deadline),
});
} catch (err) {
params.log.debug?.("failed to fetch full quoted message text", {
error: formatUnknownError(err),
});
}
}
let threadContext: string | undefined;
const threadParentId = activity.replyToId;
if (threadParentId && params.isChannel && teamAadGroupId) {
try {
const graphToken = await withMSTeamsRequestDeadline({
deadline,
label: "MS Teams thread token",
work: () => params.tokenProvider.getAccessToken("https://graph.microsoft.com"),
});
// Keep successful parent context when the replies endpoint fails, and vice versa.
const [parentResult, repliesResult] = await withMSTeamsRequestDeadline({
deadline,
label: "MS Teams thread history",
work: () =>
Promise.allSettled([
fetchParentMessageCached(
graphToken,
teamAadGroupId,
params.conversationId,
threadParentId,
(token, groupId, requestedChannelId, messageId) =>
fetchChannelMessage(token, groupId, requestedChannelId, messageId, deadline),
),
fetchThreadReplies(
graphToken,
teamAadGroupId,
params.conversationId,
threadParentId,
50,
deadline,
),
]),
});
const parentMsg = parentResult.status === "fulfilled" ? parentResult.value : undefined;
const replies = repliesResult.status === "fulfilled" ? repliesResult.value : [];
if (parentResult.status === "rejected") {
params.log.debug?.("failed to fetch parent message", {
error: formatUnknownError(parentResult.reason),
});
}
if (repliesResult.status === "rejected") {
params.log.debug?.("failed to fetch thread replies", {
error: formatUnknownError(repliesResult.reason),
});
}
const isThreadSenderAllowed = (message: GraphThreadMessage) =>
resolveInboundSupplementalSenderAllowed({
isGroup: params.isChannel,
groupPolicy: params.groupPolicy,
allowFrom: params.effectiveGroupAllowFrom,
isSenderAllowed: (allowFrom) =>
resolveMSTeamsAllowlistMatch({
allowFrom,
senderId: message.from?.user?.id ?? "",
senderName: message.from?.user?.displayName,
allowNameMatching: params.allowNameMatching,
}).allowed,
});
const parentSummary = summarizeParentMessage(parentMsg);
const visibleParentMessages = parentMsg
? filterSupplementalContextItems({
items: [parentMsg],
mode: params.contextVisibilityMode,
kind: "thread",
isSenderAllowed: isThreadSenderAllowed,
}).items
: [];
if (
parentSummary &&
visibleParentMessages.length > 0 &&
shouldInjectParentContext(route.sessionKey, threadParentId)
) {
core.system.enqueueSystemEvent(formatParentContextEvent(parentSummary), {
sessionKey: route.sessionKey,
contextKey: `msteams:thread-parent:${params.conversationId}:${threadParentId}`,
});
markParentContextInjected(route.sessionKey, threadParentId);
}
const allMessages = parentMsg ? [parentMsg, ...replies] : replies;
quoteSenderId = parentMsg?.from?.user?.id ?? parentMsg?.from?.application?.id ?? undefined;
quoteSenderName =
parentMsg?.from?.user?.displayName ??
parentMsg?.from?.application?.displayName ??
params.quoteInfo?.sender;
const { items: threadMessages } = filterSupplementalContextItems({
items: allMessages,
mode: params.contextVisibilityMode,
kind: "thread",
isSenderAllowed: isThreadSenderAllowed,
});
threadContext = formatThreadContext(threadMessages, activity.id) || undefined;
} catch (err) {
params.log.debug?.("failed to fetch thread history", {
error: formatUnknownError(err),
});
}
}
quoteSenderName ??= params.quoteInfo?.sender;
return {
teamAadGroupId,
quoteBodyFull,
quoteSenderId,
quoteSenderName,
threadContext,
};
}