Files
memind/server/portal-gateway-services-bootstrap.mjs
john 391ba0b705
Memind CI / Test, build, and release guards (pull_request) Failing after 1m29s
feat(billing): make metering formula admin-configurable
Allow margin/FX/cost-mode knobs to be overridden via h5_billing_admin_config so memind_adm can adjust DeepSeek billing without editing env.

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

376 lines
9.9 KiB
JavaScript

import crypto from 'node:crypto';
import fs from 'node:fs';
import path from 'node:path';
import { createAgentRunGateway } from '../agent-run-gateway.mjs';
import { createGoalRunService } from '../goal-run-service.mjs';
import { isDirectChatSessionId } from '../direct-chat-service.mjs';
import {
cancelSessionActiveRequest,
quiesceSessionStdioExtensions,
} from '../session-runtime-lifecycle.mjs';
import { createOrchestratorAdminConfigService } from '../services/orchestrator/admin-config.mjs';
import { createWorkflowShadowObserver } from '../services/orchestrator/shadow-observer.mjs';
import { createTkmindProxy } from '../tkmind-proxy.mjs';
import { createToolGateway } from '../tool-gateway.mjs';
import { isPassiveCanaryRuntime } from './portal-runtime-role.mjs';
function isEnabledFlag(value, fallback = '') {
return ['1', 'true', 'yes', 'on'].includes(
String(value ?? fallback).trim().toLowerCase(),
);
}
export function startPortalAgentRunRecoveryLoop(
agentRunGateway,
{
env = process.env,
logger = console,
setIntervalFn = setInterval,
} = {},
) {
if (
typeof agentRunGateway?.dispatchQueuedRuns !== 'function' ||
typeof agentRunGateway?.recoverStaleRunningRuns !== 'function'
) {
return null;
}
const configuredInterval = Number(
env.MEMIND_AGENT_RUN_RECOVERY_INTERVAL_MS ?? 30_000,
);
const intervalMs = Number.isFinite(configuredInterval)
? Math.max(5_000, configuredInterval)
: 30_000;
const configuredHeartbeatStaleMs = Number(
env.MEMIND_AGENT_RUN_HEARTBEAT_STALE_MS ?? 90_000,
);
const heartbeatStaleMs =
Number.isFinite(configuredHeartbeatStaleMs)
? Math.max(intervalMs * 2, configuredHeartbeatStaleMs)
: Math.max(intervalMs * 2, 90_000);
let sweepActive = false;
const sweep = async () => {
if (sweepActive) return;
sweepActive = true;
try {
const heartbeatRecovery =
await agentRunGateway.recoverStaleRunningRuns({
staleMs: heartbeatStaleMs,
limit: 20,
dryRun: false,
reason: 'worker_heartbeat_stale_recovery',
});
const result =
await agentRunGateway.dispatchQueuedRuns();
const recovered =
Number(heartbeatRecovery?.recovered ?? 0) +
Number(result?.staleRecovery?.recovered ?? 0);
if (recovered > 0) {
logger.warn(
`[AgentRun] recovered ${recovered} stale run(s)`,
);
}
} catch (error) {
logger.warn(
'[AgentRun] background recovery sweep failed:',
error instanceof Error ? error.message : error,
);
} finally {
sweepActive = false;
}
};
void sweep();
const timer = setIntervalFn(sweep, intervalMs);
timer?.unref?.();
return timer;
}
export function createPortalRunDeliverablesValidator({
getWorkspacePublicationDelivery = () =>
null,
} = {}) {
if (
typeof getWorkspacePublicationDelivery !==
'function'
) {
throw new Error(
'createPortalRunDeliverablesValidator requires MindSpace delivery access',
);
}
return async function validateRunDeliverables(
input,
) {
const service =
getWorkspacePublicationDelivery();
if (
!service ||
typeof service.validateRunDeliverables !==
'function'
) {
throw new Error(
'MindSpace 交付验证服务未启用',
);
}
return service.validateRunDeliverables(
input,
);
};
}
export function bootstrapPortalGatewayServices({
pool,
h5Root,
env = process.env,
apiTarget,
apiTargets,
apiSecret,
userAuth,
sessionAccess,
sessionStreamStore,
llmProviderService,
subscriptionService,
billingConfigService = null,
sessionSnapshotService,
conversationMemoryService,
memoryV2,
systemDisclosurePolicyService,
mindSpaceAssets,
getWorkspacePublicationDelivery = () =>
null,
directChatService,
chatIntentRouter,
syncUserGeneratedPages,
isSessionPageDeliveryActive,
createTkmindProxyFn = createTkmindProxy,
createToolGatewayFn = createToolGateway,
createAgentRunGatewayFn = createAgentRunGateway,
createOrchestratorAdminConfigServiceFn =
createOrchestratorAdminConfigService,
createWorkflowShadowObserverFn =
createWorkflowShadowObserver,
createRunDeliverablesValidatorFn =
createPortalRunDeliverablesValidator,
startAgentRunRecoveryLoopFn =
startPortalAgentRunRecoveryLoop,
} = {}) {
if (
!pool ||
!h5Root ||
!userAuth ||
!sessionAccess ||
!llmProviderService ||
typeof syncUserGeneratedPages !== 'function' ||
typeof isSessionPageDeliveryActive !== 'function'
) {
throw new Error(
'bootstrapPortalGatewayServices requires gateway dependencies',
);
}
// GOOSED PROXY BOUNDARY: H5 chat → goosed unique entry.
const tkmindProxy = createTkmindProxyFn({
apiTarget,
apiTargets,
apiSecret,
userAuth,
sessionAccess,
sessionStreamStore,
llmProviderService,
subscriptionService,
billingConfigService,
sessionSnapshotService,
conversationMemoryService,
memoryV2,
systemDisclosurePolicyService,
localFetchAsset: mindSpaceAssets
? async (userId, assetId) => {
const { asset, bodyBase64 } =
await mindSpaceAssets
.readAssetContent(
userId,
assetId,
);
return {
buffer: Buffer.from(
bodyBase64,
'base64',
),
mimeType: asset.mimeType,
};
}
: null,
});
const toolGateway = createToolGatewayFn({
llmProviderService,
});
const validateRunDeliverables =
createRunDeliverablesValidatorFn({
getWorkspacePublicationDelivery,
});
const workflowShadowEnabled = isEnabledFlag(
env.MEMIND_ORCHESTRATOR_SHADOW_OBSERVATION_ENABLED,
);
const workflowShadowObserver = workflowShadowEnabled
? createWorkflowShadowObserverFn({
configService:
createOrchestratorAdminConfigServiceFn(pool),
serviceToken:
env.MEMIND_ORCHESTRATOR_SERVICE_TOKEN,
logger: console,
})
: null;
let runtimeRoot;
try {
runtimeRoot = fs.realpathSync(h5Root);
} catch {
runtimeRoot = path.resolve(h5Root);
}
const workerIdentity = {
workerId:
String(
env.MEMIND_AGENT_RUN_WORKER_ID ?? '',
).trim() || `${process.pid}:${crypto.randomUUID()}`,
runtimeRoot,
buildId:
String(
env.MEMIND_RUNTIME_BUILD_ID ??
env.MEMIND_RELEASE_ID ??
env.GIT_COMMIT ??
'dev',
).trim() || 'dev',
};
const goalRunService = createGoalRunService({ pool });
const agentRunGateway = createAgentRunGatewayFn({
pool,
userAuth,
sessionAccess,
tkmindProxy,
toolGateway,
directChatService,
systemDisclosurePolicyService,
chatIntentRouter,
sessionSnapshotService,
conversationMemoryService,
goalRunService,
observeWorkflowRun: workflowShadowObserver,
observeWorkflowValidation:
workflowShadowObserver?.observeValidation ?? null,
enforcePageDataWorkflowValidation:
workflowShadowEnabled
&& isEnabledFlag(
env.MEMIND_ORCHESTRATOR_PAGE_DATA_VALIDATION_GATE_ENABLED,
),
observePersonalMemoryOnSuccess: async ({
userId,
sessionId,
userMessage,
}) => {
if (!memoryV2?.observePersonalMemory) return;
await memoryV2.observePersonalMemory({
userId,
sessionId,
messages: [userMessage],
});
},
syncUserPagesOnSuccess: async ({
userId,
sessionId,
runStartedAtMs,
}) =>
syncUserGeneratedPages(userId, {
sessionId,
sinceMs: runStartedAtMs,
}),
isSessionExternallyBusy: ({ sessionId }) =>
isSessionPageDeliveryActive(sessionId),
validateRunDeliverables,
cancelSessionOnRetry: async ({
sessionId,
requestId,
}) => {
if (isDirectChatSessionId(sessionId)) {
return { cancelled: false, skipped: true };
}
const target =
await tkmindProxy.resolveTarget(sessionId);
const sessionApiFetch = (pathname, init) =>
tkmindProxy.apiFetchTo(
target,
pathname,
init,
);
return cancelSessionActiveRequest(
sessionApiFetch,
sessionId,
requestId,
);
},
quiesceSessionOnTerminal: async ({
sessionId,
requestId,
status,
}) => {
if (isDirectChatSessionId(sessionId)) {
return { removed: [], skipped: true };
}
const target =
await tkmindProxy.resolveTarget(sessionId);
const sessionApiFetch = (pathname, init) =>
tkmindProxy.apiFetchTo(
target,
pathname,
init,
);
const cancellation =
status === 'failed'
? await cancelSessionActiveRequest(
sessionApiFetch,
sessionId,
requestId,
)
: { cancelled: false, skipped: true };
const quiesced =
await quiesceSessionStdioExtensions(
sessionApiFetch,
sessionId,
);
return { ...quiesced, cancellation };
},
autoDispatch: isEnabledFlag(
env.MEMIND_AGENT_RUN_AUTODISPATCH,
'1',
),
maxConcurrentRuns: Number(
env.MEMIND_AGENT_RUN_QUEUE_CONCURRENCY ?? 1,
),
runTimeoutMs: Number(
env.MEMIND_AGENT_RUN_TIMEOUT_MS ??
15 * 60 * 1000,
),
maxConcurrentShadowObservations: Number(
env.MEMIND_ORCHESTRATOR_SHADOW_MAX_CONCURRENCY ??
2,
),
maxQueuedShadowObservations: Number(
env.MEMIND_ORCHESTRATOR_SHADOW_MAX_QUEUE ??
100,
),
workerIdentity,
});
const agentRunRecoveryTimer = isPassiveCanaryRuntime(env)
? null
: startAgentRunRecoveryLoopFn(
agentRunGateway,
{ env, logger: console },
);
return {
tkmindProxy,
toolGateway,
agentRunGateway,
goalRunService,
agentRunRecoveryTimer,
validateRunDeliverables,
};
}