Files
memind/server/portal-gateway-services-bootstrap.mjs
john 286069449b
Memind CI / Test, build, and release guards (push) Failing after 2m14s
feat: add guarded portal canary release
2026-07-26 19:51:44 +08:00

452 lines
13 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 { isDirectChatSessionId } from '../direct-chat-service.mjs';
import { scanWorkspaceFilesForProhibitedBrowserStorage } from '../mindspace-browser-storage-policy.mjs';
import { evaluatePageDataHtmlContent } from '../mindspace-page-data-finish-guard.mjs';
import { normalizeWorkspaceRelativePath } from '../mindspace-pages.mjs';
import { resolveMindSpaceUserPublishDir } from '../mindspace-runtime-config.mjs';
import { policyAllowsAction } from '../page-access-policy.mjs';
import { detectPageDataDatasetUsageFromHtml } from '../page-data-html-detect.mjs';
import { readPageAccessPolicy } from '../page-data-policy-store.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({
h5Root,
resolveMindSpaceUserPublishDirFn =
resolveMindSpaceUserPublishDir,
normalizeWorkspaceRelativePathFn =
normalizeWorkspaceRelativePath,
evaluatePageDataHtmlContentFn =
evaluatePageDataHtmlContent,
readPageAccessPolicyFn = readPageAccessPolicy,
detectPageDataDatasetUsageFromHtmlFn =
detectPageDataDatasetUsageFromHtml,
policyAllowsActionFn = policyAllowsAction,
scanWorkspaceFilesForProhibitedBrowserStorageFn =
scanWorkspaceFilesForProhibitedBrowserStorage,
resolvePathFn = path.resolve,
pathSeparator = path.sep,
existsSyncFn = fs.existsSync,
readFileSyncFn = fs.readFileSync,
} = {}) {
if (!h5Root) {
throw new Error(
'createPortalRunDeliverablesValidator requires h5Root',
);
}
return async function validateRunDeliverables({
userId,
deliverables,
}) {
const publishDir =
resolveMindSpaceUserPublishDirFn(h5Root, {
id: userId,
});
const resolvedPublishDir =
resolvePathFn(publishDir);
const pageDataErrors = [];
for (const page of deliverables?.pages ?? []) {
const relativePath =
normalizeWorkspaceRelativePathFn(
page.workspaceRelativePath,
);
if (!relativePath?.startsWith('public/')) continue;
const filePath = resolvePathFn(
publishDir,
relativePath,
);
if (
!filePath.startsWith(
`${resolvedPublishDir}${pathSeparator}`,
) ||
!existsSyncFn(filePath)
) {
continue;
}
const html = readFileSyncFn(filePath, 'utf8');
const evaluation =
evaluatePageDataHtmlContentFn(html, {
relativePath,
});
if (!evaluation.usesPageDataApi) continue;
for (const issue of evaluation.issues) {
pageDataErrors.push({
code: issue,
message: `${relativePath} Page Data HTML 不可交付:${issue}`,
});
}
const policy = page.pageId
? readPageAccessPolicyFn(
publishDir,
page.pageId,
)
: null;
for (const [dataset, actions] of
detectPageDataDatasetUsageFromHtmlFn(html)) {
for (const action of ['read', 'insert']) {
if (
actions?.[action] &&
!policyAllowsActionFn(
policy,
dataset,
action,
)
) {
pageDataErrors.push({
code: 'page_data_policy_action_missing',
message: `${relativePath}${dataset}.${action} 未获最终 policy 授权或 dataset 已关闭`,
});
}
}
}
}
const violations =
scanWorkspaceFilesForProhibitedBrowserStorageFn({
publishDir,
relativePaths: (deliverables?.pages ?? [])
.map((page) => page.workspaceRelativePath)
.filter(Boolean),
});
return {
errors: [
...pageDataErrors,
...violations.map((violation) => ({
code: 'browser_storage_forbidden',
message: `${violation.relativePath} 使用 ${violation.apis.join(', ')}`,
})),
],
};
};
}
export function bootstrapPortalGatewayServices({
pool,
h5Root,
env = process.env,
apiTarget,
apiTargets,
apiSecret,
userAuth,
sessionAccess,
sessionStreamStore,
llmProviderService,
subscriptionService,
sessionSnapshotService,
conversationMemoryService,
memoryV2,
systemDisclosurePolicyService,
mindSpaceAssets,
directChatService,
chatIntentRouter,
syncUserGeneratedPages,
isSessionPageDeliveryActive,
createTkmindProxyFn = createTkmindProxy,
createToolGatewayFn = createToolGateway,
createAgentRunGatewayFn = createAgentRunGateway,
createOrchestratorAdminConfigServiceFn =
createOrchestratorAdminConfigService,
createWorkflowShadowObserverFn =
createWorkflowShadowObserver,
createRunDeliverablesValidatorFn =
createPortalRunDeliverablesValidator,
startAgentRunRecoveryLoopFn =
startPortalAgentRunRecoveryLoop,
readAssetFileFn = fs.promises.readFile,
} = {}) {
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,
sessionSnapshotService,
conversationMemoryService,
memoryV2,
systemDisclosurePolicyService,
localFetchAsset: mindSpaceAssets
? async (userId, assetId) => {
const { asset, path: assetPath } =
await mindSpaceAssets.readAsset(
userId,
assetId,
);
const buffer =
await readAssetFileFn(assetPath);
return {
buffer,
mimeType: asset.mimeType,
};
}
: null,
});
const toolGateway = createToolGatewayFn({
llmProviderService,
});
const validateRunDeliverables =
createRunDeliverablesValidatorFn({ h5Root });
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 agentRunGateway = createAgentRunGatewayFn({
pool,
userAuth,
sessionAccess,
tkmindProxy,
toolGateway,
directChatService,
systemDisclosurePolicyService,
chatIntentRouter,
sessionSnapshotService,
conversationMemoryService,
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,
agentRunRecoveryTimer,
validateRunDeliverables,
};
}