Files
memind/session-broker.mjs
john 14a00774d9 feat(h5): web 联网能力、实时查询路由与 session Finish 对齐
- 新增 web 能力并挂载 platform/web(web_search/fetch_url)
- 实时查询强制 web skill,router fallback 与 await session Finish
- Session Broker 覆盖率/指标、stream replay 与相关单测/E2E 脚本

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-06 16:06:26 +08:00

212 lines
6.7 KiB
JavaScript

/**
* Session Broker v0 — thin facade over userAuth session ownership + goosed target mapping.
*
* Does NOT own message history, tool state, memory, or SSE replay.
* See docs/h5-session-architecture-20260706.md Patch 1.
*/
import { createSessionBrokerMetrics, isSessionBrokerMetricsEnabled } from './session-broker-metrics.mjs';
function envFlag(value, fallback = false) {
const raw = String(value ?? '').trim().toLowerCase();
if (!raw) return fallback;
return ['1', 'true', 'yes', 'on'].includes(raw);
}
export function isSessionBrokerEnabled(env = process.env) {
return envFlag(env.MEMIND_SESSION_BROKER_ENABLED, false);
}
function requireUserAuthFn(userAuth, name) {
if (typeof userAuth?.[name] !== 'function') {
throw new Error(`userAuth.${name} is required for Session Broker`);
}
}
export function createSessionBroker({
userAuth,
logger = console,
metricsEnabled = isSessionBrokerMetricsEnabled(),
} = {}) {
requireUserAuthFn(userAuth, 'registerAgentSession');
requireUserAuthFn(userAuth, 'getSessionTarget');
requireUserAuthFn(userAuth, 'ownsSession');
requireUserAuthFn(userAuth, 'unregisterAgentSession');
const metrics = createSessionBrokerMetrics({ logger, enabled: metricsEnabled });
async function validateOwnership(userId, sessionId) {
if (!userId || !sessionId) return false;
const owned = await userAuth.ownsSession(userId, sessionId);
if (!owned) {
metrics.ownershipDenied({ userId, sessionId });
}
return owned;
}
async function registerSession({ userId, sessionId, target = 0, origin = 'h5' } = {}) {
if (!userId || !sessionId) {
throw new Error('registerSession requires userId and sessionId');
}
if (await userAuth.ownsSession(userId, sessionId)) {
metrics.registerDuplicate({ userId, sessionId, target, origin });
}
await userAuth.registerAgentSession(userId, sessionId, target);
if (typeof userAuth.setSessionOrigin === 'function' && (origin === 'h5' || origin === 'wechat')) {
await userAuth.setSessionOrigin(sessionId, origin);
}
}
async function unregisterSession({ userId, sessionId } = {}) {
if (!userId || !sessionId) return;
await userAuth.unregisterAgentSession(userId, sessionId);
}
async function resolveSessionTarget(sessionId) {
if (!sessionId) return { target: null, node: 0 };
const resolved = await userAuth.getSessionTarget(sessionId);
if (!resolved?.target && Number(resolved?.node ?? 0) > 0) {
metrics.resolveTargetMiss({
sessionId,
node: Number(resolved.node),
reason: 'legacy_node_fallback',
});
}
return resolved;
}
async function resolveSession(sessionId) {
if (!sessionId) return null;
const { target, node } = await userAuth.getSessionTarget(sessionId);
let origin = 'h5';
if (typeof userAuth.getSessionOrigins === 'function') {
const origins = await userAuth.getSessionOrigins([sessionId]);
origin = origins.get(sessionId) ?? 'h5';
}
return { sessionId, target, node, origin };
}
return {
validateOwnership,
registerSession,
unregisterSession,
resolveSession,
resolveSessionTarget,
};
}
/**
* Unified session ownership/target access with optional broker routing.
* When disabled, delegates to userAuth with identical behavior (Patch 2 fallback).
*/
export function createSessionAccess({
userAuth,
enabled = isSessionBrokerEnabled(),
broker = null,
logger = console,
metricsEnabled = isSessionBrokerMetricsEnabled(),
} = {}) {
if (!userAuth) {
throw new Error('createSessionAccess requires userAuth');
}
const useBroker = Boolean(enabled);
const sessionBroker = useBroker
? (broker ?? createSessionBroker({ userAuth, logger, metricsEnabled }))
: null;
async function validateOwnership(userId, sessionId) {
if (useBroker) return sessionBroker.validateOwnership(userId, sessionId);
if (!userId || !sessionId) return false;
return userAuth.ownsSession(userId, sessionId);
}
async function ownsSession(userId, sessionId) {
return validateOwnership(userId, sessionId);
}
async function registerSession({ userId, sessionId, target = 0, origin = 'h5' } = {}) {
if (useBroker) return sessionBroker.registerSession({ userId, sessionId, target, origin });
await userAuth.registerAgentSession(userId, sessionId, target);
if (
typeof userAuth.setSessionOrigin === 'function' &&
(origin === 'h5' || origin === 'wechat')
) {
await userAuth.setSessionOrigin(sessionId, origin);
}
}
async function registerAgentSession(userId, sessionId, target = 0) {
return registerSession({ userId, sessionId, target, origin: 'h5' });
}
async function unregisterSession({ userId, sessionId } = {}) {
if (useBroker) return sessionBroker.unregisterSession({ userId, sessionId });
if (!userId || !sessionId) return;
await userAuth.unregisterAgentSession(userId, sessionId);
}
async function resolveSessionTarget(sessionId) {
if (useBroker) return sessionBroker.resolveSessionTarget(sessionId);
if (!sessionId) return { target: null, node: 0 };
return userAuth.getSessionTarget(sessionId);
}
async function getSessionTarget(sessionId) {
return resolveSessionTarget(sessionId);
}
async function resolveSession(sessionId) {
if (useBroker) return sessionBroker.resolveSession(sessionId);
if (!sessionId) return null;
const { target, node } = await userAuth.getSessionTarget(sessionId);
let origin = 'h5';
if (typeof userAuth.getSessionOrigins === 'function') {
const origins = await userAuth.getSessionOrigins([sessionId]);
origin = origins.get(sessionId) ?? 'h5';
}
return { sessionId, target, node, origin };
}
return {
enabled: useBroker,
broker: sessionBroker,
validateOwnership,
ownsSession,
registerSession,
registerAgentSession,
unregisterSession,
resolveSession,
resolveSessionTarget,
getSessionTarget,
};
}
const noopSessionAccess = {
enabled: false,
broker: null,
async validateOwnership() {
return true;
},
async ownsSession() {
return true;
},
async registerSession() {},
async registerAgentSession() {},
async unregisterSession() {},
async resolveSession() {
return null;
},
async resolveSessionTarget() {
return { target: null, node: 0 };
},
async getSessionTarget() {
return { target: null, node: 0 };
},
};
/** Resolve shared session access for server wiring and tests. */
export function resolveSessionAccess({ userAuth, sessionAccess = null, enabled = isSessionBrokerEnabled() } = {}) {
if (sessionAccess) return sessionAccess;
if (userAuth) return createSessionAccess({ userAuth, enabled });
return noopSessionAccess;
}