452 lines
13 KiB
JavaScript
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,
|
|
};
|
|
}
|