mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-10 10:30:47 +08:00
feat: add message source context enhancement
This commit is contained in:
parent
5928b33057
commit
0cb79f74e9
69 changed files with 6217 additions and 597 deletions
|
|
@ -28,6 +28,7 @@ import {
|
|||
} from '../shared/preset-command.mjs';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
import { captureContextEnhancement, enhanceContextContent } from '../shared/context-enhancement.mjs';
|
||||
import {
|
||||
BatchInputManager,
|
||||
batchInputBusyMessage,
|
||||
|
|
@ -372,6 +373,7 @@ export class DingtalkHarnessBridge {
|
|||
#clientSecret;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#status;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
|
|
@ -383,7 +385,8 @@ export class DingtalkHarnessBridge {
|
|||
#interactionKeys = new Map();
|
||||
#interactionTasks = new Set();
|
||||
#commandTasks = new Set();
|
||||
#acceptedMessageIds = new Set();
|
||||
// Keep the accepted configuration through the existing queue/reply lifecycle.
|
||||
#acceptedMessageIds = new Map();
|
||||
#approvals;
|
||||
#batchInputs = new BatchInputManager();
|
||||
|
||||
|
|
@ -393,6 +396,7 @@ export class DingtalkHarnessBridge {
|
|||
clientSecret,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
status = createDingtalkBridgeStatus(),
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
|
|
@ -410,6 +414,7 @@ export class DingtalkHarnessBridge {
|
|||
this.#clientSecret = clientSecret.trim();
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#status = status;
|
||||
this.#logger = logger;
|
||||
this.#approvals = new HarnessApprovalQueue({ label: 'DingTalk', logger });
|
||||
|
|
@ -428,13 +433,17 @@ export class DingtalkHarnessBridge {
|
|||
return structuredClone(this.#status);
|
||||
}
|
||||
|
||||
accept(message) {
|
||||
accept(message, { contextSnapshot } = {}) {
|
||||
if (this.#signal?.aborted) return Promise.resolve();
|
||||
const messageId = nonEmptyString(message?.msgId);
|
||||
const sender = senderStaffId(message);
|
||||
if (!messageId || !sender || this.#state.hasSeen(messageId)
|
||||
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
|
||||
this.#acceptedMessageIds.add(messageId);
|
||||
this.#acceptedMessageIds.set(messageId, contextSnapshot === undefined ? captureContextEnhancement(
|
||||
this.#contextEnhancement,
|
||||
message.conversationType === '1' || message.conversationType === 1 ? 'direct'
|
||||
: message.conversationType === '2' || message.conversationType === 2 ? 'group' : null,
|
||||
) : contextSnapshot);
|
||||
|
||||
let key;
|
||||
try {
|
||||
|
|
@ -931,9 +940,17 @@ export class DingtalkHarnessBridge {
|
|||
return;
|
||||
}
|
||||
|
||||
const content = hasImages
|
||||
let content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
const snapshot = this.#acceptedMessageIds.get(messageId);
|
||||
if (snapshot) {
|
||||
content = enhanceContextContent(content ?? text, snapshot, () => ({
|
||||
channel: 'dingtalk',
|
||||
senderId: sender,
|
||||
senderName: message.senderNick,
|
||||
}));
|
||||
}
|
||||
if (typeof this.#api.createAiCard === 'function'
|
||||
&& typeof this.#api.updateAiCard === 'function'
|
||||
&& typeof this.#api.finishAiCard === 'function') {
|
||||
|
|
@ -951,7 +968,7 @@ export class DingtalkHarnessBridge {
|
|||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
...(hasImages ? { content } : { text }),
|
||||
...(content !== undefined ? { content } : { text }),
|
||||
createOptions: { signal: this.#signal },
|
||||
existsOptions: { signal: this.#signal },
|
||||
askOptions: {
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ import {
|
|||
} from './dingtalk-bridge.mjs';
|
||||
import { sendRememberedConnectionTest } from '../shared/connection-test.mjs';
|
||||
import { t } from '../shared/i18n.mjs';
|
||||
import { captureContextEnhancement } from '../shared/context-enhancement.mjs';
|
||||
|
||||
function nonEmptyString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
|
|
@ -128,6 +129,7 @@ export class DingtalkRuntime {
|
|||
#clientSecret;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#maxMessageChars;
|
||||
|
|
@ -149,6 +151,7 @@ export class DingtalkRuntime {
|
|||
clientSecret,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
maxMessageChars = 4_000,
|
||||
|
|
@ -166,6 +169,7 @@ export class DingtalkRuntime {
|
|||
this.#clientSecret = clientSecret.trim();
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#maxMessageChars = maxMessageChars;
|
||||
|
|
@ -233,6 +237,7 @@ export class DingtalkRuntime {
|
|||
approvedSenders: this.#config.approvedSenders,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
@ -268,6 +273,15 @@ export class DingtalkRuntime {
|
|||
}
|
||||
}
|
||||
|
||||
// Retain the committed, read-only settings at receipt without moving
|
||||
// JSON parsing out of the existing asynchronous callback path.
|
||||
const contextEnhancement = this.#contextEnhancement;
|
||||
let receivedSettings;
|
||||
try {
|
||||
receivedSettings = contextEnhancement?.getSettings?.();
|
||||
} catch {
|
||||
// Optional settings failures leave the original message path active.
|
||||
}
|
||||
const task = Promise.resolve().then(async () => {
|
||||
if (this.#bridge !== bridge) return;
|
||||
let message;
|
||||
|
|
@ -282,7 +296,12 @@ export class DingtalkRuntime {
|
|||
}
|
||||
if (!message || typeof message !== 'object') return;
|
||||
this.#status.lastCallbackAt = Date.now();
|
||||
await bridge.accept(message);
|
||||
const contextSnapshot = captureContextEnhancement({
|
||||
getSettings: () => receivedSettings,
|
||||
get botId() { return contextEnhancement?.botId; },
|
||||
}, message.conversationType === '1' || message.conversationType === 1 ? 'direct'
|
||||
: message.conversationType === '2' || message.conversationType === 2 ? 'group' : null);
|
||||
await bridge.accept(message, { contextSnapshot });
|
||||
}).catch(() => {
|
||||
if (signal.aborted || this.#bridge !== bridge) return;
|
||||
this.#status.lastError = t('钉钉消息处理失败。');
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ import { createEditableMessageStream, splitMessageText } from '../shared/editabl
|
|||
import { fetchFileStream } from '../shared/file-download.mjs';
|
||||
import { fetchImageBuffer } from '../shared/image-prompt.mjs';
|
||||
import { t } from '../shared/i18n.mjs';
|
||||
import { captureContextEnhancement } from '../shared/context-enhancement.mjs';
|
||||
import { DiscordApi } from './discord-api.mjs';
|
||||
import { createDiscordBridgeStatus, DiscordHarnessBridge } from './discord-bridge.mjs';
|
||||
|
||||
|
|
@ -213,6 +214,10 @@ export function normalizeDiscordMessage(message, botId, { fetchImpl = fetch } =
|
|||
return {
|
||||
messageId: String(message.id),
|
||||
senderId: String(message.author.id),
|
||||
contextSource: () => ({
|
||||
senderName: [message.member?.nick, message.author.global_name, message.author.username]
|
||||
.find((value) => typeof value === 'string' && value.trim()),
|
||||
}),
|
||||
senderIsBot: message.author.bot === true,
|
||||
kind: direct ? 'direct' : 'group',
|
||||
conversationId: String(message.channel_id),
|
||||
|
|
@ -427,6 +432,7 @@ export class DiscordRuntime {
|
|||
#token;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
|
|
@ -457,6 +463,7 @@ export class DiscordRuntime {
|
|||
token,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 20_000,
|
||||
|
|
@ -472,6 +479,7 @@ export class DiscordRuntime {
|
|||
this.#token = token;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
|
|
@ -534,6 +542,7 @@ export class DiscordRuntime {
|
|||
bot: client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
@ -707,20 +716,25 @@ export class DiscordRuntime {
|
|||
if (!messageId || this.#state.hasSeen(messageId)) return;
|
||||
let route = this.#routing.get(messageId);
|
||||
if (!route) {
|
||||
route = resolveDiscordMessageRoute(message, this.#config.platformId, {
|
||||
const contextSnapshot = captureContextEnhancement(
|
||||
this.#contextEnhancement,
|
||||
message.guild_id ? 'group' : 'direct',
|
||||
);
|
||||
const pendingRoute = resolveDiscordMessageRoute(message, this.#config.platformId, {
|
||||
api: this.#api,
|
||||
channel: this.#channels.get(String(message.channel_id)),
|
||||
signal: this.#abortController?.signal,
|
||||
onChannel: (resolved) => this.#rememberChannel(resolved),
|
||||
});
|
||||
route = { pendingRoute, contextSnapshot };
|
||||
this.#routing.set(messageId, route);
|
||||
void route.finally(() => {
|
||||
void pendingRoute.finally(() => {
|
||||
if (this.#routing.get(messageId) === route) this.#routing.delete(messageId);
|
||||
}).catch(() => undefined);
|
||||
}
|
||||
try {
|
||||
const normalized = await route;
|
||||
if (normalized) await bridge.accept(normalized);
|
||||
const normalized = await route.pendingRoute;
|
||||
if (normalized) await bridge.accept(normalized, { contextSnapshot: route.contextSnapshot });
|
||||
} catch (error) {
|
||||
if (error?.code === 'discord-thread-create-uncertain') {
|
||||
await this.#state.markSeen(messageId);
|
||||
|
|
|
|||
|
|
@ -46,6 +46,7 @@ import {
|
|||
} from '../shared/preset-command.mjs';
|
||||
import { runWorkspaceCommand, resolveSessionListWorkspace, workspacePathSnapshot } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
import { captureContextEnhancement, enhanceContextContent } from '../shared/context-enhancement.mjs';
|
||||
import { deliverOutboundArtifacts } from '../shared/semantic/artifact-delivery.mjs';
|
||||
import {
|
||||
createDeliveryReceipt,
|
||||
|
|
@ -383,12 +384,14 @@ export class FeishuHarnessBridge {
|
|||
#channel;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#queues = new Map();
|
||||
#batchInputs = new BatchInputManager();
|
||||
#pendingInteractions = new Map();
|
||||
#interactionKeys = new Map();
|
||||
#resolvedQuestionReplies = new Map();
|
||||
#acceptedMessageIds = new Set();
|
||||
// Keep the accepted configuration through the existing queue/reply lifecycle.
|
||||
#acceptedMessageIds = new Map();
|
||||
#interactionTasks = new Set();
|
||||
#commandTasks = new Set();
|
||||
/** All accepted card work, including tasks waiting behind an earlier click. */
|
||||
|
|
@ -447,6 +450,7 @@ export class FeishuHarnessBridge {
|
|||
channel,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
status,
|
||||
allowedSenderOpenIds = new Set(),
|
||||
botId,
|
||||
|
|
@ -483,6 +487,7 @@ export class FeishuHarnessBridge {
|
|||
this.#channel = channel;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#status = status;
|
||||
this.#allowedSenderOpenIds = allowedSenderOpenIds;
|
||||
this.#botId = nonEmptyString(botId);
|
||||
|
|
@ -551,7 +556,10 @@ export class FeishuHarnessBridge {
|
|||
if (chatId) rememberConnectionTestTarget(this.#state, { chatId });
|
||||
}
|
||||
|
||||
this.#acceptedMessageIds.add(messageId);
|
||||
this.#acceptedMessageIds.set(messageId, captureContextEnhancement(
|
||||
this.#contextEnhancement,
|
||||
event.message.chat_type === 'p2p' ? 'direct' : event.message.chat_type === 'group' ? 'group' : null,
|
||||
));
|
||||
const processingReaction = this.#beginReaction(messageId);
|
||||
const commandMessage = extractInboundMessage(event, this.#client);
|
||||
const commandText = nonEmptyString(commandMessage.content) ?? '';
|
||||
|
|
@ -3143,9 +3151,16 @@ export class FeishuHarnessBridge {
|
|||
askCompleted = true;
|
||||
onAskComplete?.();
|
||||
};
|
||||
const content = hasInboundImages(message)
|
||||
let content = hasInboundImages(message)
|
||||
? await promptContentForMessage(message, { signal: this.#signal })
|
||||
: undefined;
|
||||
const snapshot = this.#acceptedMessageIds.get(messageId);
|
||||
if (snapshot) {
|
||||
content = enhanceContextContent(content ?? text, snapshot, () => ({
|
||||
channel: 'feishu',
|
||||
senderId: senderOpenId(event),
|
||||
}));
|
||||
}
|
||||
if (!this.#channel?.stream) {
|
||||
const { answer, artifacts = [] } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
|
|
|
|||
|
|
@ -104,6 +104,7 @@ export class FeishuRuntime {
|
|||
#ownerOpenIds;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
#requestTimeoutMs;
|
||||
|
|
@ -131,6 +132,7 @@ export class FeishuRuntime {
|
|||
ownerOpenIds,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
repair,
|
||||
replyTimeoutMs = 600000,
|
||||
connectTimeoutMs = 15000,
|
||||
|
|
@ -162,6 +164,7 @@ export class FeishuRuntime {
|
|||
this.#ownerOpenIds = normalizedOwners;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#repair = repair ?? null;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
|
|
@ -260,6 +263,7 @@ export class FeishuRuntime {
|
|||
channel,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
allowedSenderOpenIds: new Set(this.#ownerOpenIds),
|
||||
botId: this.#botId,
|
||||
|
|
|
|||
|
|
@ -21,6 +21,7 @@ import {
|
|||
runPresetCommand,
|
||||
} from '../shared/preset-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
import { captureContextEnhancement, enhanceContextContent } from '../shared/context-enhancement.mjs';
|
||||
import {
|
||||
BatchInputManager,
|
||||
batchInputBusyMessage,
|
||||
|
|
@ -365,6 +366,7 @@ export class QqHarnessBridge {
|
|||
#ownerUserOpenid;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#status;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
|
|
@ -374,7 +376,8 @@ export class QqHarnessBridge {
|
|||
#queues = new Map();
|
||||
#pendingInteractions = new Map();
|
||||
#interactionKeys = new Map();
|
||||
#acceptedMessageIds = new Set();
|
||||
// Keep the accepted configuration through the existing queue/reply lifecycle.
|
||||
#acceptedMessageIds = new Map();
|
||||
#approvalTasks = new Set();
|
||||
#commandTasks = new Set();
|
||||
#approvals;
|
||||
|
|
@ -385,6 +388,7 @@ export class QqHarnessBridge {
|
|||
ownerUserOpenid,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
status = createQqBridgeStatus(),
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
|
|
@ -403,6 +407,7 @@ export class QqHarnessBridge {
|
|||
this.#ownerUserOpenid = ownerUserOpenid;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#status = status;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
|
|
@ -425,7 +430,10 @@ export class QqHarnessBridge {
|
|||
|| this.#state.hasSeen(messageId)
|
||||
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
|
||||
const key = conversationKey(message);
|
||||
this.#acceptedMessageIds.add(messageId);
|
||||
this.#acceptedMessageIds.set(messageId, captureContextEnhancement(
|
||||
this.#contextEnhancement,
|
||||
message.kind === 'c2c' ? 'direct' : 'group',
|
||||
));
|
||||
if (message.kind === 'c2c'
|
||||
&& (this.#ownerUserOpenid === '*' || sender === this.#ownerUserOpenid)
|
||||
&& message.replyTarget?.scope === 'c2c'
|
||||
|
|
@ -777,9 +785,17 @@ export class QqHarnessBridge {
|
|||
return;
|
||||
}
|
||||
|
||||
const content = hasImages
|
||||
let content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
const snapshot = this.#acceptedMessageIds.get(messageId);
|
||||
if (snapshot) {
|
||||
content = enhanceContextContent(content ?? text, snapshot, () => ({
|
||||
channel: 'qq',
|
||||
senderId: sender,
|
||||
senderName: message.kind === 'group' ? message.senderName : undefined,
|
||||
}));
|
||||
}
|
||||
// QQ stream_messages can acknowledge a final frame without rendering it in
|
||||
// some C2C clients. Standard Markdown delivery is the reliable reply path.
|
||||
const toolErrors = [];
|
||||
|
|
@ -793,7 +809,7 @@ export class QqHarnessBridge {
|
|||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
...(hasImages ? { content } : { text }),
|
||||
...(content !== undefined ? { content } : { text }),
|
||||
createOptions: { signal: this.#signal },
|
||||
existsOptions: { signal: this.#signal },
|
||||
askOptions: {
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ export class QqRuntime {
|
|||
#appSecret;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
|
|
@ -48,6 +49,7 @@ export class QqRuntime {
|
|||
appSecret,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 20_000,
|
||||
|
|
@ -61,6 +63,7 @@ export class QqRuntime {
|
|||
this.#appSecret = appSecret;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
|
|
@ -139,6 +142,7 @@ export class QqRuntime {
|
|||
ownerUserOpenid: this.#config.ownerUserOpenid,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
|
|
@ -14,6 +14,11 @@ import {
|
|||
validateAgentPresetId,
|
||||
} from './agent-preset.mjs';
|
||||
import { CONNECTION_TEST_STATE_IDENTITY } from './connection-test.mjs';
|
||||
import {
|
||||
DEFAULT_CONTEXT_ENHANCEMENT_CONFIG,
|
||||
normalizeContextEnhancementConfig,
|
||||
validateContextEnhancementConfig,
|
||||
} from './context-enhancement.mjs';
|
||||
import { WORKSPACE_SESSION_STALE } from './workspace-session.mjs';
|
||||
|
||||
const EMPTY_DOCUMENT = Object.freeze({ version: 1, workspaces: Object.freeze({}) });
|
||||
|
|
@ -68,7 +73,17 @@ function normalizeDocument(value) {
|
|||
}
|
||||
}
|
||||
}
|
||||
return { version: 1, workspaces, agentPresets };
|
||||
const contextEnhancement = Object.create(null);
|
||||
// Enhancement damage is isolated from the existing workspace/preset document.
|
||||
if (value.contextEnhancement && typeof value.contextEnhancement === 'object'
|
||||
&& !Array.isArray(value.contextEnhancement)) {
|
||||
for (const [botId, config] of Object.entries(value.contextEnhancement)) {
|
||||
if (/^[A-Za-z0-9_-]{1,128}$/.test(botId)) {
|
||||
contextEnhancement[botId] = normalizeContextEnhancementConfig(config);
|
||||
}
|
||||
}
|
||||
}
|
||||
return { version: 1, workspaces, agentPresets, contextEnhancement };
|
||||
}
|
||||
|
||||
export async function validateWorkspacePath(value) {
|
||||
|
|
@ -99,6 +114,7 @@ export class BotWorkspaceStore {
|
|||
#defaultWorkspace;
|
||||
#workspaces = {};
|
||||
#agentPresets = {};
|
||||
#contextEnhancement = {};
|
||||
#generations = new Map();
|
||||
#nextGeneration = 1;
|
||||
#incarnations = new Map();
|
||||
|
|
@ -121,10 +137,12 @@ export class BotWorkspaceStore {
|
|||
if (!normalized) throw new Error('dsh-im workspace config is invalid');
|
||||
this.#workspaces = normalized.workspaces;
|
||||
this.#agentPresets = normalized.agentPresets;
|
||||
this.#contextEnhancement = normalized.contextEnhancement;
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#workspaces = {};
|
||||
this.#agentPresets = {};
|
||||
this.#contextEnhancement = {};
|
||||
}
|
||||
this.#generations.clear();
|
||||
this.#nextGeneration = 1;
|
||||
|
|
@ -156,6 +174,13 @@ export class BotWorkspaceStore {
|
|||
return this.#agentPresets[botIdOf(botId)] ?? null;
|
||||
}
|
||||
|
||||
contextEnhancementFor(botId) {
|
||||
const id = botIdOf(botId);
|
||||
return this.has(id) && Object.hasOwn(this.#contextEnhancement, id)
|
||||
? this.#contextEnhancement[id]
|
||||
: DEFAULT_CONTEXT_ENHANCEMENT_CONFIG;
|
||||
}
|
||||
|
||||
generationFor(botId) {
|
||||
return this.#generations.get(botIdOf(botId)) ?? null;
|
||||
}
|
||||
|
|
@ -270,6 +295,24 @@ export class BotWorkspaceStore {
|
|||
});
|
||||
}
|
||||
|
||||
async setContextEnhancement(botId, value, { incarnation } = {}) {
|
||||
const id = botIdOf(botId);
|
||||
const expectedIncarnation = incarnation === undefined ? this.incarnationFor(id) : incarnation;
|
||||
const config = validateContextEnhancementConfig(value);
|
||||
return this.#enqueue(id, async () => {
|
||||
if (!this.has(id) || expectedIncarnation !== this.incarnationFor(id)) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
throw error;
|
||||
}
|
||||
const next = { ...this.#contextEnhancement, [id]: config };
|
||||
// Messages keep the previous committed snapshot until rename succeeds.
|
||||
await this.#persist(next);
|
||||
this.#contextEnhancement = next;
|
||||
return config;
|
||||
});
|
||||
}
|
||||
|
||||
async bindWorkspaceSession(botId, value, {
|
||||
conversationKey,
|
||||
sessionId,
|
||||
|
|
@ -408,6 +451,14 @@ export class BotWorkspaceStore {
|
|||
});
|
||||
}
|
||||
|
||||
/** A failed retirement must reach disk before the same config ID can be rebound. */
|
||||
flushPendingRemoval(botId) {
|
||||
if (!this.#dirtyRemovals.has(botId)) return undefined;
|
||||
return this.#enqueue(botId, async () => {
|
||||
if (this.#dirtyRemovals.has(botId)) await this.#persistCurrentDocument();
|
||||
});
|
||||
}
|
||||
|
||||
async remove(botId) {
|
||||
const result = await this.retireAfterConfigCommit(botId);
|
||||
if (result.error) throw result.error;
|
||||
|
|
@ -419,6 +470,7 @@ export class BotWorkspaceStore {
|
|||
const candidates = new Set([
|
||||
...Object.keys(this.#workspaces),
|
||||
...Object.keys(this.#agentPresets),
|
||||
...Object.keys(this.#contextEnhancement),
|
||||
...this.#dirtyRemovals,
|
||||
]);
|
||||
for (const botId of candidates) {
|
||||
|
|
@ -435,6 +487,7 @@ export class BotWorkspaceStore {
|
|||
...bot,
|
||||
workspace: this.workspaceFor(bot.botId),
|
||||
agentPreset: this.agentPresetFor(bot.botId),
|
||||
contextEnhancement: this.contextEnhancementFor(bot.botId),
|
||||
}
|
||||
: bot),
|
||||
};
|
||||
|
|
@ -464,9 +517,11 @@ export class BotWorkspaceStore {
|
|||
async #retireCurrentIncarnation(id) {
|
||||
const hadWorkspace = Object.hasOwn(this.#workspaces, id);
|
||||
const hadPreset = Object.hasOwn(this.#agentPresets, id);
|
||||
const needsCleanup = hadWorkspace || hadPreset || this.#dirtyRemovals.has(id);
|
||||
const hadContextEnhancement = Object.hasOwn(this.#contextEnhancement, id);
|
||||
const needsCleanup = hadWorkspace || hadPreset || hadContextEnhancement || this.#dirtyRemovals.has(id);
|
||||
delete this.#workspaces[id];
|
||||
delete this.#agentPresets[id];
|
||||
delete this.#contextEnhancement[id];
|
||||
this.#generations.delete(id);
|
||||
this.#incarnations.delete(id);
|
||||
if (!needsCleanup) return {
|
||||
|
|
@ -496,11 +551,14 @@ export class BotWorkspaceStore {
|
|||
return queued;
|
||||
}
|
||||
|
||||
async #persist() {
|
||||
async #persist(contextEnhancement = this.#contextEnhancement) {
|
||||
const document = { version: 1, workspaces: this.#workspaces };
|
||||
if (Object.keys(this.#agentPresets).length > 0) {
|
||||
document.agentPresets = this.#agentPresets;
|
||||
}
|
||||
if (Object.keys(contextEnhancement).length > 0) {
|
||||
document.contextEnhancement = contextEnhancement;
|
||||
}
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, `${JSON.stringify(document, null, 2)}\n`, {
|
||||
|
|
@ -513,7 +571,8 @@ export class BotWorkspaceStore {
|
|||
|
||||
async #persistCurrentDocument() {
|
||||
if (Object.keys(this.#workspaces).length > 0
|
||||
|| Object.keys(this.#agentPresets).length > 0) {
|
||||
|| Object.keys(this.#agentPresets).length > 0
|
||||
|| Object.keys(this.#contextEnhancement).length > 0) {
|
||||
await this.#persist();
|
||||
return;
|
||||
}
|
||||
|
|
@ -569,10 +628,16 @@ function targetStatus(controller) {
|
|||
return Promise.resolve(controller.status());
|
||||
}
|
||||
|
||||
/** Observe the config store's durable removal commit without changing its API. */
|
||||
/** Observe durable removals and finish failed cleanup before a same-ID config save. */
|
||||
export function observeBotWorkspaceRemovals(
|
||||
configStore,
|
||||
{ workspaces, method = 'remove', botIdFromRemoved = (removed) => removed?.botId },
|
||||
{
|
||||
workspaces,
|
||||
method = 'remove',
|
||||
botIdFromRemoved = (removed) => removed?.botId,
|
||||
saveMethod = 'save',
|
||||
botIdFromSave = (config) => config?.botId,
|
||||
},
|
||||
) {
|
||||
if (!configStore || !workspaces || typeof configStore[method] !== 'function') {
|
||||
throw new TypeError('configStore removal observer dependencies are required');
|
||||
|
|
@ -588,6 +653,12 @@ export function observeBotWorkspaceRemovals(
|
|||
return removed;
|
||||
};
|
||||
}
|
||||
if (property === saveMethod && typeof value === 'function') {
|
||||
return (...args) => {
|
||||
const cleanup = workspaces.flushPendingRemoval(botIdFromSave(args[0], args));
|
||||
return cleanup ? cleanup.then(() => value.apply(target, args)) : value.apply(target, args);
|
||||
};
|
||||
}
|
||||
return typeof value === 'function' ? value.bind(target) : value;
|
||||
},
|
||||
});
|
||||
|
|
@ -998,6 +1069,31 @@ export function createWorkspaceAwareController(controller, { workspaces, stateFo
|
|||
);
|
||||
});
|
||||
};
|
||||
const updateContextEnhancement = (botId, value, projectStatus) => {
|
||||
const incarnation = workspaces.incarnationFor(botId);
|
||||
const config = validateContextEnhancementConfig(value);
|
||||
return withBotTransition(botId, async () => {
|
||||
const snapshot = await controller.status();
|
||||
if (!snapshot?.bots?.some((bot) => bot?.botId === botId)) {
|
||||
const error = new Error('找不到要修改的机器人。');
|
||||
error.code = 'workspace-bot-not-found';
|
||||
throw error;
|
||||
}
|
||||
const catalog = await resolveAgentPresetCatalog(agentPresetCatalog);
|
||||
const decorated = workspaces.decorateStatus(snapshot);
|
||||
const updated = {
|
||||
...decorated,
|
||||
bots: decorated.bots.map((bot) => bot?.botId === botId
|
||||
? { ...bot, contextEnhancement: config } : bot),
|
||||
...(catalog ? { agentPresetCatalog: catalog } : {}),
|
||||
};
|
||||
// QR/status projection can fail too. Prepare the complete response before
|
||||
// commit so a failed save never publishes new running settings.
|
||||
const result = projectStatus ? await projectStatus(updated) : updated;
|
||||
await workspaces.setContextEnhancement(botId, config, { incarnation });
|
||||
return result;
|
||||
});
|
||||
};
|
||||
const deleteWithWorkspace = (botId, invokeDelete) => withBotTransition(botId, async () => {
|
||||
// Fence the old runtime without changing the durable mapping. A crash
|
||||
// before the controller removes its config therefore keeps the bot's
|
||||
|
|
@ -1037,6 +1133,7 @@ export function createWorkspaceAwareController(controller, { workspaces, stateFo
|
|||
get(target, property) {
|
||||
if (property === 'updateWorkspace') return updateWorkspace;
|
||||
if (property === 'updateAgentPreset') return updateAgentPreset;
|
||||
if (property === 'updateContextEnhancement') return updateContextEnhancement;
|
||||
const value = Reflect.get(target, property, target);
|
||||
if (typeof value !== 'function') return value;
|
||||
if (property === 'deleteBot') {
|
||||
|
|
|
|||
135
src/channels/shared/context-enhancement.mjs
Normal file
135
src/channels/shared/context-enhancement.mjs
Normal file
|
|
@ -0,0 +1,135 @@
|
|||
// Shared by the Host and settings UI; keep this module browser-compatible.
|
||||
export const CONTEXT_ENHANCEMENT_FIELDS = Object.freeze([
|
||||
'channel', 'conversationType', 'senderId', 'senderName', 'botId',
|
||||
]);
|
||||
|
||||
export const CONTEXT_ENHANCEMENT_GUIDANCE_MAX_LENGTH = 8_000;
|
||||
export const CONTEXT_GUIDANCE_EXAMPLE = `仅依据当前消息的 <dsh_im_source> 中实际提供的字段理解来源;没有提供的字段不要猜测或补全。
|
||||
conversationType是群聊时回复严肃一点,conversationType是私聊时回复一定要幽默搞笑,像周星驰的电影一样搞笑`;
|
||||
|
||||
// Kept as an alias for integrations that imported the original template name.
|
||||
export const DEFAULT_CONTEXT_GUIDANCE = CONTEXT_GUIDANCE_EXAMPLE;
|
||||
|
||||
export const DEFAULT_CONTEXT_ENHANCEMENT_CONFIG = Object.freeze({
|
||||
groupEnabled: false,
|
||||
directEnabled: false,
|
||||
fields: Object.freeze(['senderId']),
|
||||
guidance: '',
|
||||
});
|
||||
|
||||
const CONFIG_KEYS = ['groupEnabled', 'directEnabled', 'fields', 'guidance'];
|
||||
const CHANNELS = new Set([
|
||||
'wecom', 'weixin', 'feishu', 'dingtalk', 'qq',
|
||||
'slack', 'telegram', 'discord', 'whatsapp',
|
||||
]);
|
||||
const SOURCE_LIMITS = { channel: 16, conversationType: 6, senderId: 256, senderName: 256, botId: 128 };
|
||||
const CONTROL_CHARACTERS = /[\u0000-\u001f\u007f-\u009f\u202a-\u202e\u2066-\u2069]/g;
|
||||
|
||||
function invalidConfig(message) {
|
||||
const error = new TypeError(message);
|
||||
error.code = 'context-enhancement-invalid';
|
||||
return error;
|
||||
}
|
||||
|
||||
/** Validate the complete atomic save, preserving explicit empty selections/text. */
|
||||
export function validateContextEnhancementConfig(input) {
|
||||
if (!input || typeof input !== 'object' || Array.isArray(input)
|
||||
|| ![Object.prototype, null].includes(Object.getPrototypeOf(input))
|
||||
|| Reflect.ownKeys(input).length !== CONFIG_KEYS.length
|
||||
|| !CONFIG_KEYS.every((key) => Object.hasOwn(input, key))) {
|
||||
throw invalidConfig('请提交完整的上下文增强设置。');
|
||||
}
|
||||
const { groupEnabled, directEnabled, fields, guidance } = input;
|
||||
if (typeof groupEnabled !== 'boolean' || typeof directEnabled !== 'boolean') {
|
||||
throw invalidConfig('群聊和私聊开关必须是布尔值。');
|
||||
}
|
||||
if (!Array.isArray(fields) || ![...fields].every((field) => CONTEXT_ENHANCEMENT_FIELDS.includes(field))) {
|
||||
throw invalidConfig('来源字段只能选择已定义的五个字段。');
|
||||
}
|
||||
if (typeof guidance !== 'string' || guidance.length > CONTEXT_ENHANCEMENT_GUIDANCE_MAX_LENGTH) {
|
||||
throw invalidConfig(`增强提示词不得超过 ${CONTEXT_ENHANCEMENT_GUIDANCE_MAX_LENGTH} 个字符。`);
|
||||
}
|
||||
return Object.freeze({
|
||||
groupEnabled,
|
||||
directEnabled,
|
||||
fields: Object.freeze(CONTEXT_ENHANCEMENT_FIELDS.filter((field) => fields.includes(field))),
|
||||
guidance: guidance.trim() ? guidance : '',
|
||||
});
|
||||
}
|
||||
|
||||
/** Missing or damaged enhancement settings must never break an existing bot. */
|
||||
export function normalizeContextEnhancementConfig(input) {
|
||||
try {
|
||||
return validateContextEnhancementConfig(input);
|
||||
} catch {
|
||||
return DEFAULT_CONTEXT_ENHANCEMENT_CONFIG;
|
||||
}
|
||||
}
|
||||
|
||||
/** Capture before queueing. The off path reads only the applicable switch. */
|
||||
export function captureContextEnhancement(provider, conversationType) {
|
||||
if (conversationType !== 'group' && conversationType !== 'direct') return null;
|
||||
try {
|
||||
const settings = provider?.getSettings?.();
|
||||
const enabledKey = conversationType === 'group' ? 'groupEnabled' : 'directEnabled';
|
||||
if (settings?.[enabledKey] !== true) return null;
|
||||
const config = normalizeContextEnhancementConfig(settings);
|
||||
if (config[enabledKey] !== true) return null;
|
||||
return Object.freeze({ config, botId: provider.botId, conversationType });
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function sourceString(value, field) {
|
||||
if (field === 'senderId' && (typeof value === 'bigint' || Number.isFinite(value))) {
|
||||
value = String(value);
|
||||
}
|
||||
if (typeof value !== 'string') return undefined;
|
||||
const normalized = value.replace(CONTROL_CHARACTERS, '').trim().slice(0, SOURCE_LIMITS[field]);
|
||||
if (!normalized || (field === 'channel' && !CHANNELS.has(normalized))) return undefined;
|
||||
return normalized;
|
||||
}
|
||||
|
||||
function sourceBlock(snapshot, sourceFactory) {
|
||||
const { fields } = snapshot.config;
|
||||
const needsSource = fields.some((field) => ['channel', 'senderId', 'senderName'].includes(field));
|
||||
const source = needsSource ? sourceFactory?.() : null;
|
||||
const projected = {};
|
||||
for (const field of fields) {
|
||||
const value = field === 'botId' || field === 'conversationType'
|
||||
? snapshot[field] : source?.[field];
|
||||
const normalized = sourceString(value, field);
|
||||
if (normalized !== undefined) projected[field] = normalized;
|
||||
}
|
||||
if (Object.keys(projected).length === 0) return '';
|
||||
const json = JSON.stringify(projected).replace(/[<>&]/g, (character) => ({
|
||||
'<': '\\u003c', '>': '\\u003e', '&': '\\u0026',
|
||||
})[character]);
|
||||
return `<dsh_im_source>${json}</dsh_im_source>`;
|
||||
}
|
||||
|
||||
function guidanceBlock(guidance) {
|
||||
if (!guidance.trim()) return '';
|
||||
const body = guidance.replace(/<\/?dsh_im_source_guidance\b[^>]*(?:>|$)/gi, (tag) => (
|
||||
tag.replace(/</g, '<').replace(/>/g, '>')
|
||||
));
|
||||
return `<dsh_im_source_guidance>\n${body}\n</dsh_im_source_guidance>`;
|
||||
}
|
||||
|
||||
/** Add one text prefix; never inspect sources, format or copy content when off. */
|
||||
export function enhanceContextContent(content, snapshot, sourceFactory) {
|
||||
if (!snapshot) return content;
|
||||
try {
|
||||
const blocks = [sourceBlock(snapshot, sourceFactory), guidanceBlock(snapshot.config.guidance)]
|
||||
.filter(Boolean);
|
||||
if (blocks.length === 0) return content;
|
||||
const prefix = blocks.join('\n\n');
|
||||
if (typeof content === 'string') return `${prefix}\n\n${content}`;
|
||||
if (Array.isArray(content)) return [{ type: 'text', text: prefix }, ...content];
|
||||
return content;
|
||||
} catch {
|
||||
// Only enhancement errors are isolated; the caller's original flow proceeds.
|
||||
return content;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,4 +1,5 @@
|
|||
import { t } from './i18n.mjs';
|
||||
import { captureContextEnhancement, enhanceContextContent } from './context-enhancement.mjs';
|
||||
import { runWorkspaceCommand } from './workspace-command.mjs';
|
||||
import { runCompactCommand } from './compact-command.mjs';
|
||||
import { isHistoryCommand, runHistoryCommand } from './history-command.mjs';
|
||||
|
|
@ -127,6 +128,7 @@ export class TextHarnessBridge {
|
|||
#bot;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#status;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
|
|
@ -134,7 +136,8 @@ export class TextHarnessBridge {
|
|||
#queues = new Map();
|
||||
#pendingInteractions = new Map();
|
||||
#interactionKeys = new Map();
|
||||
#acceptedMessageIds = new Set();
|
||||
// Keep the accepted configuration through the existing queue/reply lifecycle.
|
||||
#acceptedMessageIds = new Map();
|
||||
#approvalTasks = new Set();
|
||||
#commandTasks = new Set();
|
||||
#approvals;
|
||||
|
|
@ -145,6 +148,7 @@ export class TextHarnessBridge {
|
|||
bot,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
status = createTextBridgeStatus(),
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
|
|
@ -157,6 +161,7 @@ export class TextHarnessBridge {
|
|||
this.#bot = bot;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#status = status;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
|
|
@ -171,7 +176,7 @@ export class TextHarnessBridge {
|
|||
return structuredClone(this.#status);
|
||||
}
|
||||
|
||||
accept(message) {
|
||||
accept(message, { contextSnapshot } = {}) {
|
||||
if (this.#signal?.aborted) return Promise.resolve();
|
||||
const conversationId = cleanText(message?.conversationId);
|
||||
const kind = message?.kind === 'group' ? 'group' : 'direct';
|
||||
|
|
@ -182,7 +187,9 @@ export class TextHarnessBridge {
|
|||
|| this.#state.hasSeen(messageId) || this.#acceptedMessageIds.has(messageId)) {
|
||||
return Promise.resolve();
|
||||
}
|
||||
this.#acceptedMessageIds.add(messageId);
|
||||
this.#acceptedMessageIds.set(messageId, contextSnapshot === undefined
|
||||
? captureContextEnhancement(this.#contextEnhancement, message?.kind)
|
||||
: contextSnapshot);
|
||||
const statusReaction = beginStatusReaction({
|
||||
adapter: this.#bot,
|
||||
target: normalized.kind === 'direct' || normalized.addressed === true
|
||||
|
|
@ -614,9 +621,17 @@ export class TextHarnessBridge {
|
|||
);
|
||||
}
|
||||
}
|
||||
const content = hasImages
|
||||
let content = hasImages
|
||||
? await promptContentForMessage(message, { signal: this.#signal })
|
||||
: undefined;
|
||||
const snapshot = this.#acceptedMessageIds.get(messageId);
|
||||
if (snapshot) {
|
||||
content = enhanceContextContent(content ?? text, snapshot, () => ({
|
||||
channel: this.#descriptor.key,
|
||||
senderId,
|
||||
senderName: message.contextSource?.()?.senderName,
|
||||
}));
|
||||
}
|
||||
const { answer, artifacts = [] } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
|
|
|
|||
|
|
@ -346,6 +346,7 @@ export class SlackRuntime {
|
|||
#appToken;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
|
|
@ -369,6 +370,7 @@ export class SlackRuntime {
|
|||
appToken,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 20_000,
|
||||
|
|
@ -384,6 +386,7 @@ export class SlackRuntime {
|
|||
this.#appToken = appToken;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
|
|
@ -436,6 +439,7 @@ export class SlackRuntime {
|
|||
bot: client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ import { randomInt } from 'node:crypto';
|
|||
import { createEditableMessageStream, splitMessageText } from '../shared/editable-message-stream.mjs';
|
||||
import { createTextDeliveryBlock } from '../shared/semantic/delivery.mjs';
|
||||
import { t } from '../shared/i18n.mjs';
|
||||
import { captureContextEnhancement } from '../shared/context-enhancement.mjs';
|
||||
import { COMMANDS_MENU_BUTTON, TelegramApi } from './telegram-api.mjs';
|
||||
import { createTelegramBridgeStatus, TelegramHarnessBridge } from './telegram-bridge.mjs';
|
||||
import {
|
||||
|
|
@ -159,6 +160,11 @@ export function normalizeTelegramUpdate(update, {
|
|||
return {
|
||||
messageId: String(update.update_id),
|
||||
senderId: String(senderId),
|
||||
contextSource: () => ({
|
||||
senderName: [message.from?.first_name, message.from?.last_name]
|
||||
.filter((value) => typeof value === 'string' && value.trim())
|
||||
.map((value) => value.trim()).join(' ') || message.from?.username,
|
||||
}),
|
||||
senderIsBot: message.from?.is_bot === true,
|
||||
kind: direct ? 'direct' : 'group',
|
||||
conversationId: messageThreadId === undefined
|
||||
|
|
@ -659,6 +665,7 @@ export class TelegramRuntime {
|
|||
#token;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#createApi;
|
||||
|
|
@ -676,6 +683,7 @@ export class TelegramRuntime {
|
|||
token,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
createApi = (options) => new TelegramApi(options),
|
||||
|
|
@ -687,6 +695,7 @@ export class TelegramRuntime {
|
|||
this.#token = token;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#createApi = createApi;
|
||||
|
|
@ -758,6 +767,7 @@ export class TelegramRuntime {
|
|||
bot: client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
@ -798,7 +808,19 @@ export class TelegramRuntime {
|
|||
while (!signal.aborted) {
|
||||
const updates = await this.#api.getUpdates({ offset: cursor, timeout: 25, signal });
|
||||
this.#status.lastCheckedAt = Date.now();
|
||||
for (const update of updates) {
|
||||
if (signal.aborted) return;
|
||||
// All updates have arrived together; cursor persistence must not move the
|
||||
// settings boundary for the later messages in this received batch.
|
||||
const received = updates.map((update) => {
|
||||
const chatType = update?.message?.chat?.type;
|
||||
return {
|
||||
update,
|
||||
contextSnapshot: captureContextEnhancement(this.#contextEnhancement,
|
||||
chatType === 'private' ? 'direct'
|
||||
: chatType === 'group' || chatType === 'supergroup' ? 'group' : null),
|
||||
};
|
||||
});
|
||||
for (const { update, contextSnapshot } of received) {
|
||||
if (signal.aborted) return;
|
||||
const message = normalizeTelegramUpdate(update, {
|
||||
botId: this.#config.platformId,
|
||||
|
|
@ -810,7 +832,7 @@ export class TelegramRuntime {
|
|||
accessMode: this.#accessMode,
|
||||
allowedPrivateUserIds: this.#allowedPrivateUserIds,
|
||||
})) {
|
||||
void this.#bridge.accept(message).catch((error) => {
|
||||
void this.#bridge.accept(message, { contextSnapshot }).catch((error) => {
|
||||
if (signal.aborted) return;
|
||||
this.#logger.error?.(
|
||||
`[dsh-im:telegram] bot ${this.#config.botId} message handling failed:`,
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import {
|
|||
} from '../shared/preset-command.mjs';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
import { captureContextEnhancement, enhanceContextContent } from '../shared/context-enhancement.mjs';
|
||||
import {
|
||||
hasInboundImages,
|
||||
ImagePromptError,
|
||||
|
|
@ -488,6 +489,7 @@ export class WecomHarnessBridge {
|
|||
#client;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#status;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
|
|
@ -497,7 +499,8 @@ export class WecomHarnessBridge {
|
|||
#queues = new Map();
|
||||
#pendingInteractions = new Map();
|
||||
#interactionKeys = new Map();
|
||||
#acceptedMessageIds = new Set();
|
||||
// Keep the accepted configuration through the existing queue/reply lifecycle.
|
||||
#acceptedMessageIds = new Map();
|
||||
#approvalTasks = new Set();
|
||||
#commandTasks = new Set();
|
||||
#approvals;
|
||||
|
|
@ -508,6 +511,7 @@ export class WecomHarnessBridge {
|
|||
client,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
status = createWecomBridgeStatus(),
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
|
|
@ -525,6 +529,7 @@ export class WecomHarnessBridge {
|
|||
this.#client = client;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#status = status;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
|
|
@ -552,7 +557,10 @@ export class WecomHarnessBridge {
|
|||
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
|
||||
|
||||
const key = conversationKey(frame);
|
||||
this.#acceptedMessageIds.add(messageId);
|
||||
this.#acceptedMessageIds.set(messageId, captureContextEnhancement(
|
||||
this.#contextEnhancement,
|
||||
body.chattype === 'single' ? 'direct' : 'group',
|
||||
));
|
||||
if (body.chattype === 'single') {
|
||||
rememberConnectionTestTarget(this.#state, { chatId });
|
||||
}
|
||||
|
|
@ -948,9 +956,16 @@ export class WecomHarnessBridge {
|
|||
this.#logger.warn?.('[dsh-im:wecom] unable to start a stream; using an active reply:', error);
|
||||
}
|
||||
|
||||
const content = hasImages
|
||||
let content = hasImages
|
||||
? await promptContentForMessage(message, { signal: this.#signal })
|
||||
: undefined;
|
||||
const snapshot = this.#acceptedMessageIds.get(messageId);
|
||||
if (snapshot) {
|
||||
content = enhanceContextContent(content ?? text, snapshot, () => ({
|
||||
channel: 'wecom',
|
||||
senderId,
|
||||
}));
|
||||
}
|
||||
await this.#state.markSeen(messageId);
|
||||
promptRecorded = true;
|
||||
const { answer, artifacts = [] } = await askInWorkspaceSession({
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@ export class WecomRuntime {
|
|||
#secret;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
|
|
@ -45,6 +46,7 @@ export class WecomRuntime {
|
|||
secret,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 20_000,
|
||||
|
|
@ -58,6 +60,7 @@ export class WecomRuntime {
|
|||
this.#secret = secret;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
|
|
@ -107,6 +110,7 @@ export class WecomRuntime {
|
|||
client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
|
|
@ -33,6 +33,7 @@ import {
|
|||
} from '../shared/preset-command.mjs';
|
||||
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
|
||||
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
|
||||
import { captureContextEnhancement, enhanceContextContent } from '../shared/context-enhancement.mjs';
|
||||
import {
|
||||
hasInboundImages,
|
||||
imagePromptDiagnostic,
|
||||
|
|
@ -258,6 +259,7 @@ export class WeixinHarnessBridge {
|
|||
#ownerUserId;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#status;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
|
|
@ -267,7 +269,8 @@ export class WeixinHarnessBridge {
|
|||
#queues = new Map();
|
||||
#pendingInteractions = new Map();
|
||||
#interactionKeys = new Map();
|
||||
#acceptedMessageIds = new Set();
|
||||
// Keep the accepted configuration through the existing queue/reply lifecycle.
|
||||
#acceptedMessageIds = new Map();
|
||||
#approvalTasks = new Set();
|
||||
#commandTasks = new Set();
|
||||
#approvals;
|
||||
|
|
@ -288,6 +291,7 @@ export class WeixinHarnessBridge {
|
|||
ownerUserId,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
status = createWeixinBridgeStatus(),
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
|
|
@ -307,6 +311,7 @@ export class WeixinHarnessBridge {
|
|||
this.#ownerUserId = ownerUserId;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#status = status;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
|
|
@ -327,7 +332,10 @@ export class WeixinHarnessBridge {
|
|||
const sender = nonEmptyString(message?.from_user_id);
|
||||
if (!messageId || !sender || this.#state.hasSeen(messageId)
|
||||
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
|
||||
this.#acceptedMessageIds.add(messageId);
|
||||
this.#acceptedMessageIds.set(messageId, captureContextEnhancement(
|
||||
this.#contextEnhancement,
|
||||
'direct',
|
||||
));
|
||||
if (sender === this.#ownerUserId) {
|
||||
rememberConnectionTestTarget(this.#state, { toUserId: sender });
|
||||
}
|
||||
|
|
@ -650,16 +658,23 @@ export class WeixinHarnessBridge {
|
|||
let artifacts = [];
|
||||
await this.#startTyping(sender, contextToken);
|
||||
try {
|
||||
const content = hasImages
|
||||
let content = hasImages
|
||||
? await promptContentForMessage(promptMessage, { signal: this.#signal })
|
||||
: undefined;
|
||||
const snapshot = this.#acceptedMessageIds.get(messageId);
|
||||
if (snapshot) {
|
||||
content = enhanceContextContent(content ?? text, snapshot, () => ({
|
||||
channel: 'weixin',
|
||||
senderId: sender,
|
||||
}));
|
||||
}
|
||||
await this.#state.markSeen(messageId);
|
||||
promptRecorded = true;
|
||||
({ answer, artifacts = [] } = await askInWorkspaceSession({
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
key,
|
||||
...(hasImages ? { content } : { text }),
|
||||
...(content !== undefined ? { content } : { text }),
|
||||
createOptions: { signal: this.#signal },
|
||||
existsOptions: { signal: this.#signal },
|
||||
askOptions: {
|
||||
|
|
|
|||
|
|
@ -108,6 +108,7 @@ export class WeixinRuntime {
|
|||
#token;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#maxMessageChars;
|
||||
|
|
@ -124,6 +125,7 @@ export class WeixinRuntime {
|
|||
token,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
maxMessageChars = DEFAULT_WEIXIN_MAX_MESSAGE_CHARS,
|
||||
|
|
@ -137,6 +139,7 @@ export class WeixinRuntime {
|
|||
this.#token = token;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#maxMessageChars = maxMessageChars;
|
||||
|
|
@ -178,6 +181,7 @@ export class WeixinRuntime {
|
|||
ownerUserId: this.#config.ownerUserId,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
|
|
@ -234,6 +234,7 @@ export function normalizeWhatsappMessage(message, accountJid, {
|
|||
messageId: `${remoteJid}:${messageId}`,
|
||||
providerMessageId: messageId,
|
||||
senderId: senderJid,
|
||||
contextSource: () => ({ senderName: message.pushName }),
|
||||
senderAlternateId: typeof senderAlternateJid === 'string' ? senderAlternateJid : '',
|
||||
senderIsBot: false,
|
||||
kind: group ? 'group' : 'direct',
|
||||
|
|
@ -536,6 +537,7 @@ export class WhatsappRuntime {
|
|||
#authDir;
|
||||
#harness;
|
||||
#state;
|
||||
#contextEnhancement;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
|
|
@ -555,6 +557,7 @@ export class WhatsappRuntime {
|
|||
authDir,
|
||||
harness,
|
||||
state,
|
||||
contextEnhancement,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 30_000,
|
||||
|
|
@ -568,6 +571,7 @@ export class WhatsappRuntime {
|
|||
this.#authDir = authDir;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#contextEnhancement = contextEnhancement;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
|
|
@ -673,6 +677,7 @@ export class WhatsappRuntime {
|
|||
bot: client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
contextEnhancement: this.#contextEnhancement,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue