370 lines
9.7 KiB
JavaScript
370 lines
9.7 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 {
|
|
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,
|
|
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,
|
|
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 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,
|
|
};
|
|
}
|