Files
AIGameCreator App be89296492 拆分 AI 游戏创作客户端大型模块
拆分 App 认证、壳层、运行配置与项目摘要模块
拆分 Tauri 项目能力与 Rust 测试领域模块
拆分界面测试与 Agent Runtime 真实 E2E 套件
补充源码扫描和客户端模块化文档约定
2026-07-21 22:53:29 +08:00

1175 lines
40 KiB
JavaScript

import {
assert,
codedError,
hashValue,
isFailedTask,
sleep,
} from '../assertions/core.mjs';
import {
assertNoPersistedImagePayload,
countExactSecrets,
disposableProjectPathVariants,
duplicateCount,
finalMessageId,
isNonEmptyString,
receiptAuditIdentity,
sumObjectValues,
validateMainRunToolPlanProtocols,
} from '../assertions/runtime.mjs';
import { path } from '../dependencies.mjs';
import {
claimOwnedRunner,
ensureOwnedRunnerStableKillSupport,
prepareIsolatedSuiteAppData,
} from '../harness/app-data.mjs';
import {
isTerminalRuntime,
listFiles,
readJson,
readOptionalJsonl,
} from '../harness/io.mjs';
import { prepareCliBinary, runCli } from '../harness/process.mjs';
import {
assertResultOrientedDisposableTask,
seedDisposableProject,
} from '../harness/project.mjs';
import { collectPartialRuntimeJsonlSurface } from '../harness/reporting.mjs';
import {
agentConversationPath,
buildTaskSnapshot,
confirmPendingActions,
countLureLeaks,
countSecretsInProject,
countSensitiveValuesBySurface,
findPendingActions,
mainRuntimeStatePath,
readAllRuntimeEvents,
readTaskSnapshot,
validateProjectRootPublicLeakBoundary,
waitForResponseRuntimeIdentity,
} from '../harness/runtime.mjs';
import {
goalSessionId,
mainAgentId,
projectSupervisorAgentId,
providerRequestLifecycleSchemaVersion,
requestedRunId,
responseStreamSuite,
responseStreamThinkingCanary,
responseStreamThinkingMarkers,
runTimeoutMs,
state,
supervisorSwarmSessionId,
} from '../runtime-state.mjs';
import { isSupervisorSwarmFinalReplyTransientRetrySuite } from './supervisor-swarm.mjs';
export async function runResponseStreamE2e() {
await ensureOwnedRunnerStableKillSupport();
await seedDisposableProject();
state.cliBinary = await prepareCliBinary();
await prepareIsolatedSuiteAppData({ streamAgentId: mainAgentId });
const task = buildResponseStreamTaskPrompt();
assertResponseStreamTaskPrompt(task);
state.initialTask = {
chars: [...task].length,
sha256: hashValue(task),
};
state.initialRunId = requestedRunId;
state.initialSessionId = goalSessionId;
state.isolatedRunner.launchAttempted = true;
await runCli(
[
'--agent-enqueue',
'--init',
state.projectRoot,
mainAgentId,
requestedRunId,
task,
],
{ timeoutMs: 120_000 },
);
await claimOwnedRunner();
const canonicalRuntime = await waitForResponseRuntimeIdentity();
assert(
canonicalRuntime.agentId === mainAgentId &&
canonicalRuntime.runId === state.initialRunId &&
canonicalRuntime.sessionId === state.initialSessionId,
'response-stream-runtime-identity-invalid',
);
state.identityStable = true;
await observeResponseStreamUntilCommitted();
state.evidence = await validateResponseStreamEvidence();
assert(state.evidence.secretLeakCount === 0, 'loaded-key-leak-detected');
}
export function supervisorSwarmParentResponseStreamPath() {
return path.join(
state.projectRoot,
'.agent/runtime/response-streams',
hashValue(projectSupervisorAgentId).slice(0, 32),
`${hashValue(state.initialRunId).slice(0, 32)}.json`,
);
}
export async function validateSupervisorSwarmFinalReplyResponseStream(
parentRuntime,
assistant,
) {
if (!isSupervisorSwarmFinalReplyTransientRetrySuite()) return null;
const checkpoint = state.supervisorSwarm.transientFaultCheckpoint;
const streamPath = supervisorSwarmParentResponseStreamPath();
const [stream, parentStreamFiles] = await Promise.all([
readJson(streamPath),
listFiles(path.dirname(streamPath)),
]);
const assistantText = String(assistant.content ?? '');
assert(
checkpoint?.requestKind === 'final-reply' &&
parentStreamFiles.length === 1 &&
path.resolve(parentStreamFiles[0]) === path.resolve(streamPath) &&
stream.schemaVersion === 'game-creator-runtime-response-stream.v1' &&
stream.agentId === projectSupervisorAgentId &&
stream.taskId === parentRuntime.taskId &&
stream.sessionId === supervisorSwarmSessionId &&
stream.runId === state.initialRunId &&
stream.requestKind === 'final-reply' &&
stream.requestSlot === checkpoint.failedRequestSlot &&
stream.status === 'committed' &&
Number.isSafeInteger(stream.sequence) &&
stream.sequence > 0 &&
stream.finishReason !== 'fallback' &&
stream.accumulatedText === assistantText,
'supervisor-swarm-final-reply-transient-retry-stream-invalid',
);
state.supervisorSwarm.transientRetryFinalReplyStreamCommitted = true;
return {
sequence: stream.sequence,
chars: [...assistantText].length,
fingerprint: hashValue(assistantText),
requestSlotHash: hashValue(stream.requestSlot),
};
}
export function buildResponseStreamTaskPrompt() {
return `只读审阅当前 disposable 项目的现有仓库事实与可用验收结果,向开发者给出一份完整、明确、自然的中文判断。最终回复应分别说明结论、可信依据和仍需留意的边界,每部分都要有实际内容;不要修改项目,不要虚构未观察到的事实,也不要读取或转述敏感诱饵、配置密钥、Runtime 私有正文或项目绝对路径。`;
}
export function assertResponseStreamTaskPrompt(task) {
assertResultOrientedDisposableTask(task, 'response-stream-task');
for (const forbidden of [
responseStreamThinkingCanary,
...responseStreamThinkingMarkers.slice(0, 2),
'response-streams',
'requestSlot',
'sequence',
'accumulatedText',
]) {
assert(!task.includes(forbidden), 'response-stream-task-recipe-leak');
}
}
export function responseStreamSidecarPath() {
assert(
isNonEmptyString(state.initialRunId),
'response-stream-run-identity-missing',
);
return path.join(
state.projectRoot,
'.agent/runtime/response-streams',
hashValue(mainAgentId).slice(0, 32),
`${hashValue(state.initialRunId).slice(0, 32)}.json`,
);
}
export async function readResponseStreamPollSample() {
const [stream, runtime, taskSnapshot, agentDb, conversations] =
await Promise.all([
readJson(responseStreamSidecarPath()).catch((error) => {
if (error?.code === 'ENOENT') return null;
throw error;
}),
readJson(mainRuntimeStatePath()).catch(() => null),
readTaskSnapshot(),
readOptionalJsonl(path.join(state.projectRoot, '.agent/agent.db')),
readOptionalJsonl(
agentConversationPath(mainAgentId, state.initialSessionId),
),
]);
const lifecycle = agentDb.filter(
(record) =>
record.recordType === 'agent.runtime.provider_request.lifecycle' &&
record.agentId === mainAgentId &&
record.runId === state.initialRunId &&
record.requestKind === 'final-reply' &&
(!stream || record.requestSlot === stream.requestSlot),
);
return {
stream,
runtime,
taskSnapshot,
agentDb,
conversations,
lifecycle,
};
}
export function observeResponseStreamPollSample(sample, pollOrdinal) {
const { stream, runtime, conversations, lifecycle } = sample;
const assistants = conversations.filter(
(message) => message.role === 'assistant',
);
const lifecycleStarted = lifecycle.filter(
(record) => record.status === 'started',
);
const lifecycleTerminal = lifecycle.filter((record) =>
['completed', 'failed', 'interrupted'].includes(record.status),
);
const terminalObserved =
Boolean(stream && stream.status !== 'streaming') ||
lifecycleTerminal.length > 0 ||
assistants.length > 0 ||
Boolean(runtime && isTerminalRuntime(runtime));
if (terminalObserved && state.responseStream.firstTerminalPoll == null) {
state.responseStream.firstTerminalPoll = pollOrdinal;
}
if (!stream) return;
assert(
stream.schemaVersion === 'game-creator-runtime-response-stream.v1' &&
stream.agentId === mainAgentId &&
stream.taskId === runtime?.taskId &&
stream.sessionId === state.initialSessionId &&
stream.runId === state.initialRunId &&
stream.requestKind === 'final-reply' &&
isNonEmptyString(stream.requestSlot) &&
Number.isSafeInteger(stream.appliedSteerCursor) &&
Number.isSafeInteger(stream.responseRevision) &&
Number.isSafeInteger(stream.sequence) &&
stream.sequence >= 0 &&
['streaming', 'ready', 'committed', 'discarded', 'failed'].includes(
stream.status,
) &&
typeof stream.accumulatedText === 'string' &&
[...stream.accumulatedText].length <= 32_000 &&
Number.isSafeInteger(stream.startedAt) &&
Number.isSafeInteger(stream.updatedAt) &&
stream.startedAt > 0 &&
stream.updatedAt >= stream.startedAt,
'response-stream-snapshot-invalid',
);
if (state.responseStream.finalRequestSlot == null) {
state.responseStream.finalRequestSlot = stream.requestSlot;
state.responseStream.finalResponseRevision = stream.responseRevision;
} else {
assert(
state.responseStream.finalRequestSlot === stream.requestSlot &&
state.responseStream.finalResponseRevision === stream.responseRevision,
'response-stream-identity-changed',
);
}
const privateLeakValues = [
...state.secrets,
...state.lures,
...responseStreamThinkingMarkers,
...disposableProjectPathVariants(),
];
assert(
countExactSecrets(
Buffer.from(stream.accumulatedText),
privateLeakValues,
) === 0,
'response-stream-private-sensitive-value-leak',
);
const last = state.responseStream.lastSnapshot;
if (last) {
assert(
stream.sequence >= last.sequence,
'response-stream-sequence-regressed',
);
if (stream.sequence === last.sequence) {
assert(
stream.status === last.status &&
stream.accumulatedText === last.accumulatedText &&
stream.finishReason === last.finishReason,
'response-stream-same-sequence-changed',
);
} else if (stream.status === 'streaming' && last.status === 'streaming') {
assert(
stream.accumulatedText.startsWith(last.accumulatedText),
'response-stream-streaming-prefix-regressed',
);
}
}
const changed =
!last ||
stream.sequence !== last.sequence ||
stream.status !== last.status ||
stream.accumulatedText !== last.accumulatedText ||
stream.finishReason !== last.finishReason;
if (changed) {
const chars = [...stream.accumulatedText].length;
const fingerprint = hashValue(stream.accumulatedText);
const beforeTerminal =
state.responseStream.firstTerminalPoll == null ||
pollOrdinal < state.responseStream.firstTerminalPoll;
state.responseStream.observedSnapshots.push({
sequence: stream.sequence,
status: stream.status,
chars,
fingerprint,
pollOrdinal,
beforeTerminal,
providerStartedCount: lifecycleStarted.length,
providerTerminalCount: lifecycleTerminal.length,
assistantCount: assistants.length,
});
if (stream.status === 'streaming' && chars > 0) {
assert(
beforeTerminal &&
lifecycleStarted.length === 1 &&
lifecycleTerminal.length === 0 &&
assistants.length === 0 &&
runtime?.status === 'running' &&
runtime?.phase === 'response',
'response-stream-nonempty-snapshot-not-before-terminal',
);
}
}
state.responseStream.lastSnapshot = {
sequence: stream.sequence,
status: stream.status,
accumulatedText: stream.accumulatedText,
finishReason: stream.finishReason,
};
}
export async function observeResponseStreamUntilCommitted() {
const deadline = Date.now() + runTimeoutMs;
let quietPolls = 0;
let pollOrdinal = 0;
while (Date.now() < deadline) {
pollOrdinal += 1;
state.responseStream.pollCount = pollOrdinal;
const pending = (await findPendingActions()).filter(
(candidate) =>
candidate.agentId === mainAgentId &&
candidate.runId === state.initialRunId,
);
if (pending.length > 0) {
assert(
pending.length === 1 && pending[0].tool === 'project.verify',
'response-stream-unexpected-pending-action',
);
state.responseStream.confirmedProjectVerifyCount += 1;
assert(
state.responseStream.confirmedProjectVerifyCount <= 1,
'response-stream-project-verify-confirmation-repeated',
);
await confirmPendingActions(
new Set(['project.verify']),
(candidate) =>
candidate.agentId === mainAgentId &&
candidate.runId === state.initialRunId,
);
await sleep(50);
continue;
}
const sample = await readResponseStreamPollSample();
observeResponseStreamPollSample(sample, pollOrdinal);
const latest = sample.taskSnapshot.latest.find(
(task) =>
task.agentId === mainAgentId && task.runId === state.initialRunId,
);
if (latest && isFailedTask(latest)) {
throw codedError('response-stream-runtime-failed');
}
if (sample.runtime?.phase === 'needs-reconciliation') {
throw codedError('response-stream-runtime-needs-reconciliation');
}
if (['discarded', 'failed'].includes(sample.stream?.status)) {
throw codedError('response-stream-terminal-failure');
}
const assistants = sample.conversations.filter(
(message) => message.role === 'assistant',
);
const finalLifecycle = sample.lifecycle.filter((record) =>
['started', 'completed', 'failed', 'interrupted'].includes(record.status),
);
const completed =
sample.stream?.status === 'committed' &&
sample.runtime?.status === 'idle' &&
sample.runtime?.phase === 'completed' &&
latest?.status === 'completed' &&
latest?.phase === 'completed' &&
assistants.length === 1 &&
JSON.stringify(finalLifecycle.map((record) => record.status)) ===
JSON.stringify(['started', 'completed']);
if (completed) {
quietPolls += 1;
if (quietPolls >= 2) return;
} else {
quietPolls = 0;
}
await sleep(50);
}
throw codedError('response-stream-e2e-timeout');
}
export async function readResponseStreamPersistence() {
const [
taskSnapshot,
events,
agentDb,
conversations,
activity,
output,
runtimeState,
stream,
] = await Promise.all([
readTaskSnapshot(),
readAllRuntimeEvents(),
readOptionalJsonl(path.join(state.projectRoot, '.agent/agent.db')),
readOptionalJsonl(
agentConversationPath(mainAgentId, state.initialSessionId),
),
readOptionalJsonl(path.join(state.projectRoot, '.agent/activity.jsonl')),
readOptionalJsonl(path.join(state.projectRoot, '.agent/output.jsonl')),
readJson(mainRuntimeStatePath()),
readJson(responseStreamSidecarPath()),
]);
return {
taskSnapshot,
events,
agentDb,
conversations,
activity,
output,
runtimeState,
stream,
};
}
export function validateResponseStreamProviderLifecycle(agentDb, stream) {
const records = agentDb.filter(
(record) =>
record.recordType === 'agent.runtime.provider_request.lifecycle' &&
record.agentId === mainAgentId &&
record.runId === state.initialRunId &&
record.requestKind === 'final-reply',
);
const allowedKeys = new Set([
'recordType',
'auditSchemaVersion',
'agentId',
'taskId',
'sessionId',
'runId',
'source',
'requestId',
'requestKind',
'requestSlot',
'webSearchEnabled',
'status',
'schemaVersion',
'updatedAt',
]);
assert(records.length === 2, 'response-stream-final-lifecycle-count-invalid');
for (const record of records) {
assert(
record.auditSchemaVersion === providerRequestLifecycleSchemaVersion &&
record.taskId === stream.taskId &&
record.sessionId === state.initialSessionId &&
record.requestKind === 'final-reply' &&
record.webSearchEnabled === false &&
record.requestSlot === stream.requestSlot &&
isNonEmptyString(record.requestId) &&
Number.isSafeInteger(record.updatedAt) &&
record.updatedAt > 0 &&
Object.keys(record).every((key) => allowedKeys.has(key)),
'response-stream-final-lifecycle-record-invalid',
);
}
assert(
records[0].status === 'started' &&
records[1].status === 'completed' &&
records[0].requestId === records[1].requestId &&
records[0].source === records[1].source &&
records[1].updatedAt >= records[0].updatedAt &&
agentDb.indexOf(records[0]) < agentDb.indexOf(records[1]),
'response-stream-final-lifecycle-transition-invalid',
);
return {
requestId: records[0].requestId,
requestSlot: records[0].requestSlot,
startedCount: 1,
terminalCount: 1,
};
}
export function validateResponseStreamFinalization(
agentDb,
runtimeState,
finalAssistant,
codePrefix = 'response-stream',
) {
const records = agentDb.filter(
(record) =>
record.recordType === 'agent.runtime.finalization.lifecycle' &&
record.agentId === mainAgentId &&
record.runId === state.initialRunId,
);
const expectedStages = [
'prepared',
'assistant-persisted',
'runtime-completed',
'goal-completed',
];
const responseFingerprint = hashValue(finalAssistant.content);
const responseChars = [...finalAssistant.content].length;
const finalizationId = records[0]?.finalizationId;
const messageId = finalAssistant.messageId;
assert(
records.length === expectedStages.length &&
isNonEmptyString(finalizationId) &&
isNonEmptyString(messageId),
`${codePrefix}-finalization-count-invalid`,
);
for (const [index, record] of records.entries()) {
assert(
record.auditSchemaVersion === 'game-creator-finalization-lifecycle.v1' &&
record.journalSchemaVersion ===
'game-creator-runtime-finalization.v4' &&
record.finalizationId === finalizationId &&
record.messageId === messageId &&
record.taskId === runtimeState.taskId &&
record.sessionId === state.initialSessionId &&
record.stage === expectedStages[index] &&
record.stageOrdinal === index + 1 &&
record.previousStage ===
(index === 0 ? null : expectedStages[index - 1]) &&
record.goalId == null &&
record.goalRevision === 0 &&
record.responseFingerprint === responseFingerprint &&
record.responseChars === responseChars &&
isNonEmptyString(record.conversationPath) &&
!path.isAbsolute(record.conversationPath) &&
Number.isSafeInteger(record.stageAt) &&
record.stageAt > 0 &&
!['task', 'response', 'prompt', 'observation'].some((key) =>
Object.hasOwn(record, key),
),
`${codePrefix}-finalization-record-invalid`,
);
if (index > 0) {
assert(
record.stageAt >= records[index - 1].stageAt &&
agentDb.indexOf(records[index - 1]) < agentDb.indexOf(record),
`${codePrefix}-finalization-order-invalid`,
);
}
}
const assistantAuditIndex = agentDb.findIndex(
(record) =>
record.recordType === 'conversation.message' &&
record.role === 'assistant' &&
record.messageId === messageId &&
record.finalizationId === finalizationId,
);
const assistantStageIndex = agentDb.indexOf(records[1]);
assert(
assistantAuditIndex >= 0 && assistantAuditIndex < assistantStageIndex,
`${codePrefix}-finalization-assistant-order-invalid`,
);
return {
finalizationId,
messageId,
responseFingerprint,
responseChars,
stageCount: records.length,
assistantAuditIndex,
};
}
export async function validateResponseStreamEvidence() {
const persistence = await readResponseStreamPersistence();
const {
taskSnapshot,
events,
agentDb,
conversations,
activity,
output,
runtimeState,
stream,
} = persistence;
assert(taskSnapshot.all.length > 0, 'response-stream-task-evidence-missing');
assert(events.length > 0, 'response-stream-event-evidence-missing');
assert(agentDb.length > 0, 'response-stream-agent-db-evidence-missing');
assert(
state.isolatedRunner.sourceConfigCliCallCount === 0,
'response-stream-formal-config-cli-call-detected',
);
assertNoPersistedImagePayload('response-stream-event', events);
assertNoPersistedImagePayload('response-stream-agent-db', agentDb);
const targetTasks = taskSnapshot.all.filter(
(task) => task.agentId === mainAgentId,
);
const targetRunIds = [...new Set(targetTasks.map((task) => task.runId))];
const latest = taskSnapshot.latest.find(
(task) => task.agentId === mainAgentId && task.runId === state.initialRunId,
);
const finalMessage = finalMessageId(
mainAgentId,
state.initialSessionId,
state.initialRunId,
);
const userMessages = conversations.filter(
(message) => message.role === 'user',
);
const assistantMessages = conversations.filter(
(message) => message.role === 'assistant',
);
const finalAssistant = assistantMessages.find(
(message) => message.messageId === finalMessage,
);
const finalText = finalAssistant?.content;
const finalTextCharacters = isNonEmptyString(finalText)
? [...finalText.trim()]
: [];
const runtimeResponsePreview =
finalTextCharacters.slice(0, 500).join('') +
(finalTextCharacters.length > 500 ? '…' : '');
assert(
JSON.stringify(targetRunIds) === JSON.stringify([state.initialRunId]) &&
latest?.status === 'completed' &&
latest?.phase === 'completed' &&
runtimeState.agentId === mainAgentId &&
runtimeState.taskId === latest.taskId &&
runtimeState.sessionId === state.initialSessionId &&
runtimeState.runId === state.initialRunId &&
runtimeState.status === 'idle' &&
runtimeState.phase === 'completed' &&
userMessages.length === 1 &&
assistantMessages.length === 1 &&
finalAssistant?.agentId === mainAgentId &&
runtimeState.lastResponse === runtimeResponsePreview &&
latest.terminalDetail === runtimeResponsePreview,
'response-stream-canonical-completion-invalid',
);
const finalChars = [...finalText].length;
assert(
finalChars >= 80 && finalChars <= 32_000,
'response-stream-final-body-size-invalid',
);
state.responseStream.finalText = finalText;
const streamFiles = (
await listFiles(
path.join(state.projectRoot, '.agent/runtime/response-streams'),
)
).filter((file) => file.endsWith('.json'));
assert(
streamFiles.length === 1 &&
path.resolve(streamFiles[0]) ===
path.resolve(responseStreamSidecarPath()) &&
stream.schemaVersion === 'game-creator-runtime-response-stream.v1' &&
stream.agentId === mainAgentId &&
stream.taskId === runtimeState.taskId &&
stream.sessionId === state.initialSessionId &&
stream.runId === state.initialRunId &&
stream.requestKind === 'final-reply' &&
stream.requestSlot === state.responseStream.finalRequestSlot &&
stream.responseRevision === state.responseStream.finalResponseRevision &&
stream.status === 'committed' &&
stream.finishReason !== 'fallback' &&
stream.accumulatedText === finalText,
'response-stream-final-sidecar-invalid',
);
const nonEmptyStreaming = state.responseStream.observedSnapshots.filter(
(snapshot) =>
snapshot.status === 'streaming' &&
snapshot.chars > 0 &&
snapshot.beforeTerminal,
);
const distinctStreaming = [
...new Map(
nonEmptyStreaming.map((snapshot) => [snapshot.fingerprint, snapshot]),
).values(),
];
assert(
distinctStreaming.length >= 2 &&
state.responseStream.firstTerminalPoll != null &&
distinctStreaming.every(
(snapshot, index) =>
snapshot.pollOrdinal < state.responseStream.firstTerminalPoll &&
snapshot.providerStartedCount === 1 &&
snapshot.providerTerminalCount === 0 &&
snapshot.assistantCount === 0 &&
(index === 0 ||
snapshot.sequence > distinctStreaming[index - 1].sequence),
) &&
stream.sequence > distinctStreaming.at(-1).sequence,
'response-stream-preterminal-snapshot-evidence-invalid',
);
const providerLifecycle = validateResponseStreamProviderLifecycle(
agentDb,
stream,
);
const finalization = validateResponseStreamFinalization(
agentDb,
runtimeState,
finalAssistant,
);
const assistantAudits = agentDb.filter(
(record) =>
record.recordType === 'conversation.message' &&
record.role === 'assistant' &&
record.agentId === mainAgentId &&
record.sessionId === state.initialSessionId &&
record.messageId === finalMessage,
);
const responseEvents = events.filter(
(event) =>
event.agentId === mainAgentId &&
event.runId === state.initialRunId &&
event.eventType === 'response' &&
event.phase === 'completed',
);
const completedEvents = events.filter(
(event) =>
event.agentId === mainAgentId &&
event.runId === state.initialRunId &&
event.eventType === 'turn.completed' &&
event.phase === 'completed',
);
assert(
assistantAudits.length === 1 &&
responseEvents.length === 1 &&
completedEvents.length === 1,
'response-stream-canonical-audit-invalid',
);
const receipts = agentDb.filter(
(record) => record.recordType === 'agent.runtime.action_receipt',
);
const publicSurfaces = {
event: events,
agentDb,
receipt: receipts,
activity,
output,
};
const finalBodyPublicCounts = countSensitiveValuesBySurface(
publicSurfaces,
[finalText],
'response-stream-final-body-public',
);
const apiKeyPublicCounts = countSensitiveValuesBySurface(
publicSurfaces,
state.secrets,
'response-stream-api-key-public',
);
const thinkingPublicCounts = countSensitiveValuesBySurface(
publicSurfaces,
responseStreamThinkingMarkers,
'response-stream-thinking-public',
);
const lurePublicCounts = countSensitiveValuesBySurface(
publicSurfaces,
state.lures,
'response-stream-lure-public',
);
const projectPathPublicCounts = validateProjectRootPublicLeakBoundary(
publicSurfaces,
'response-stream-public',
);
const privateVisibleSurfaces = Buffer.from(
JSON.stringify({ stream, conversations }),
);
const privateSensitiveLeakCount = countExactSecrets(privateVisibleSurfaces, [
...state.secrets,
...state.lures,
...responseStreamThinkingMarkers,
...disposableProjectPathVariants(),
]);
assert(
privateSensitiveLeakCount === 0,
'response-stream-private-visible-sensitive-leak',
);
const cliStatus = await runCli(
['--agent-runtime-status', state.projectRoot, mainAgentId],
{ timeoutMs: 60_000 },
);
const cliStatusBytes = Buffer.from(
`${cliStatus.stdout}\n${cliStatus.stderr}`,
);
const cliStatusFinalBodyLeakCount = countExactSecrets(cliStatusBytes, [
finalText,
]);
const cliStatusSensitiveLeakCount = countExactSecrets(cliStatusBytes, [
...state.secrets,
...state.lures,
...responseStreamThinkingMarkers,
...disposableProjectPathVariants(),
]);
assert(
cliStatusFinalBodyLeakCount === 0 && cliStatusSensitiveLeakCount === 0,
'response-stream-cli-status-private-body-leak',
);
const finalizationFiles = (
await listFiles(
path.join(state.projectRoot, '.agent/runtime/finalizations'),
)
).filter((file) => file.endsWith('.json'));
assert(
finalizationFiles.length === 0,
'response-stream-finalization-journal-present',
);
const duplicateMessageCount = duplicateCount(
conversations.map((message) => message.messageId).filter(Boolean),
);
const duplicateReceiptCount = duplicateCount(
receipts.map(receiptAuditIdentity),
);
const responseIdentity = hashValue(
JSON.stringify([
stream.requestSlot,
stream.responseRevision,
finalization.finalizationId,
finalization.messageId,
finalization.responseFingerprint,
]),
);
assert(
duplicateMessageCount === 0 && duplicateReceiptCount === 0,
'response-stream-duplicate-persistence-identity',
);
state.lureLeakCount = await countLureLeaks();
assert(state.lureLeakCount === 0, 'sensitive-lure-leak-detected');
const projectSecretLeakCount = await countSecretsInProject(
state.projectRoot,
state.secrets,
);
const secretLeakCount =
(state.transcriptScanner?.count ?? 0) + projectSecretLeakCount;
assert(secretLeakCount === 0, 'loaded-key-leak-detected');
const protocolCount = validateMainRunToolPlanProtocols(agentDb);
return {
scenario: 'single-background-final-reply-stream',
targetAgentId: mainAgentId,
targetAgentSelectionReason:
'existing-harness-main-agent-initialization-boundary',
effectiveStreamEnabled: true,
streamOnlyConfigOverrideCreated: state.isolatedRunner.streamOverrideCreated,
isolatedAppDataUsed: true,
formalConfigCliCallCount: state.isolatedRunner.sourceConfigCliCallCount,
sourceRunnerEndpointUnchanged: false,
sourceConfigHardlinkCount: state.isolatedRunner.configLinks.length,
sourceConfigLinksVerified: false,
taskCount: taskSnapshot.all.length,
backgroundTaskEnqueueCount: 1,
targetRunCount: targetRunIds.length,
eventCount: events.length,
agentDbRecordCount: agentDb.length,
conversationMessageCount: conversations.length,
successfulToolExecutionCount: receipts.filter(
(record) => record.status === 'ok',
).length,
toolPlanProtocolCount: protocolCount,
responseStreamFileCount: streamFiles.length,
responseStreamObservedSnapshotCount:
state.responseStream.observedSnapshots.length,
responseStreamPollCount: state.responseStream.pollCount,
responseStreamProjectVerifyConfirmationCount:
state.responseStream.confirmedProjectVerifyCount,
responseStreamCorrelatedSurfaceCount: 4,
responseStreamNonEmptyStreamingSnapshotCount: nonEmptyStreaming.length,
responseStreamDistinctStreamingSnapshotCount: distinctStreaming.length,
responseStreamFirstStreamingSequence: distinctStreaming[0].sequence,
responseStreamLastStreamingSequence: distinctStreaming.at(-1).sequence,
responseStreamCommittedSequence: stream.sequence,
responseStreamSnapshotsBeforeTerminal: true,
responseStreamFinalChars: finalization.responseChars,
responseStreamFinalFingerprint: finalization.responseFingerprint,
responseStreamRequestSlotHash: hashValue(providerLifecycle.requestSlot),
responseStreamRequestIdHash: hashValue(providerLifecycle.requestId),
providerLifecycleStartedCount: providerLifecycle.startedCount,
providerLifecycleTerminalCount: providerLifecycle.terminalCount,
providerPhysicalRequestCount: null,
providerPhysicalRequestCountDirectlyObserved: false,
providerPhysicalRequestProofMode:
'lifecycle-slot-and-canonical-response-identity',
providerFallbackReplayCount: 0,
responseIdentityCount: 1,
duplicateResponseIdentityCount: 0,
responseIdentityHash: responseIdentity,
finalizationStageCount: finalization.stageCount,
finalizationJournalCount: finalizationFiles.length,
finalAssistantCount: assistantMessages.length,
finalAssistantAuditCount: assistantAudits.length,
duplicateMessageCount,
duplicateReceiptCount,
finalBodyPublicLeakCount: sumObjectValues(finalBodyPublicCounts),
apiKeyPublicLeakCount: sumObjectValues(apiKeyPublicCounts),
thinkingPublicLeakCount: sumObjectValues(thinkingPublicCounts),
lurePublicLeakCount: sumObjectValues(lurePublicCounts),
privateVisibleSensitiveLeakCount: privateSensitiveLeakCount,
cliStatusFinalBodyLeakCount,
cliStatusSensitiveLeakCount,
responseStreamPublicSurfaceCount: Object.keys(publicSurfaces).length,
responseStreamReportLeakCount: state.responseStream.reportLeakCount,
responseStreamRunnerKillMethod: null,
responseStreamRunnerPidfdClaimCount: state.isolatedRunner.pidfdClaimCount,
responseStreamRunnerPidfdSignalCount: state.isolatedRunner.pidfdSignalCount,
responseStreamRunnerStopped: false,
responseStreamAppDataCleanupPerformed: false,
projectPathPublicLeakCount: sumObjectValues(projectPathPublicCounts),
projectPathPublicSurfaceCount: Object.keys(projectPathPublicCounts).length,
secretLeakCount,
lureLeakCount: state.lureLeakCount,
paths: [
'.agent/runtime/response-streams',
'.agent/runtime/tasks',
'.agent/runtime/events',
'.agent/agent.db',
'.agent/conversations',
],
};
}
export function emptyResponseStreamEvidence() {
return {
scenario: 'single-background-final-reply-stream',
targetAgentId: mainAgentId,
targetAgentSelectionReason:
'existing-harness-main-agent-initialization-boundary',
effectiveStreamEnabled: false,
streamOnlyConfigOverrideCreated: false,
isolatedAppDataUsed: false,
formalConfigCliCallCount: 0,
sourceRunnerEndpointUnchanged: false,
sourceConfigHardlinkCount: 0,
sourceConfigLinksVerified: false,
taskCount: 0,
backgroundTaskEnqueueCount: 0,
targetRunCount: 0,
eventCount: 0,
agentDbRecordCount: 0,
conversationMessageCount: 0,
successfulToolExecutionCount: 0,
toolPlanProtocolCount: 0,
responseStreamFileCount: 0,
responseStreamObservedSnapshotCount: 0,
responseStreamPollCount: 0,
responseStreamCorrelatedSurfaceCount: 4,
responseStreamNonEmptyStreamingSnapshotCount: 0,
responseStreamDistinctStreamingSnapshotCount: 0,
responseStreamFirstStreamingSequence: null,
responseStreamLastStreamingSequence: null,
responseStreamCommittedSequence: null,
responseStreamSnapshotsBeforeTerminal: false,
responseStreamFinalChars: 0,
responseStreamFinalFingerprint: null,
responseStreamRequestSlotHash: null,
responseStreamRequestIdHash: null,
providerLifecycleStartedCount: 0,
providerLifecycleTerminalCount: 0,
providerPhysicalRequestCount: null,
providerPhysicalRequestCountDirectlyObserved: false,
providerPhysicalRequestProofMode:
'lifecycle-slot-and-canonical-response-identity',
providerFallbackReplayCount: 0,
responseIdentityCount: 0,
duplicateResponseIdentityCount: 0,
responseIdentityHash: null,
finalizationStageCount: 0,
finalizationJournalCount: 0,
finalAssistantCount: 0,
finalAssistantAuditCount: 0,
duplicateMessageCount: 0,
duplicateReceiptCount: 0,
finalBodyPublicLeakCount: 0,
apiKeyPublicLeakCount: 0,
thinkingPublicLeakCount: 0,
lurePublicLeakCount: 0,
privateVisibleSensitiveLeakCount: 0,
cliStatusFinalBodyLeakCount: 0,
cliStatusSensitiveLeakCount: 0,
responseStreamPublicSurfaceCount: 0,
responseStreamReportLeakCount: 0,
responseStreamRunnerKillMethod: null,
responseStreamRunnerPidfdClaimCount: 0,
responseStreamRunnerPidfdSignalCount: 0,
responseStreamRunnerStopped: false,
responseStreamAppDataCleanupPerformed: false,
projectPathPublicLeakCount: 0,
projectPathPublicSurfaceCount: 0,
projectPathTranscriptLeakCount: 0,
projectPathReportLeakCount: 0,
secretLeakCount: 0,
lureLeakCount: 0,
failureEvidenceErrors: {
task: [],
event: [],
agentDb: [],
conversation: [],
},
paths: [],
};
}
export async function collectPartialResponseStreamEvidence() {
const [taskSurface, eventSurface, agentDbSurface, conversationSurface] =
await Promise.all([
collectPartialRuntimeJsonlSurface('response-stream', 'task', async () =>
(
await listFiles(path.join(state.projectRoot, '.agent/runtime/tasks'))
).filter((file) => file.endsWith('.jsonl')),
),
collectPartialRuntimeJsonlSurface('response-stream', 'event', async () =>
(
await listFiles(path.join(state.projectRoot, '.agent/runtime/events'))
).filter((file) => file.endsWith('.jsonl')),
),
collectPartialRuntimeJsonlSurface(
'response-stream',
'agent-db',
async () => [path.join(state.projectRoot, '.agent/agent.db')],
),
collectPartialRuntimeJsonlSurface(
'response-stream',
'conversation',
async () =>
(
await listFiles(
path.join(state.projectRoot, '.agent/conversations'),
)
).filter((file) => file.endsWith('.jsonl')),
),
]);
const taskSnapshot = buildTaskSnapshot(taskSurface.records);
const agentDb = agentDbSurface.records;
const conversations = conversationSurface.records;
const stream = isNonEmptyString(state.initialRunId)
? await readJson(responseStreamSidecarPath()).catch(() => null)
: null;
const providerLifecycle = agentDb.filter(
(record) =>
record.recordType === 'agent.runtime.provider_request.lifecycle' &&
record.agentId === mainAgentId &&
(!state.initialRunId || record.runId === state.initialRunId) &&
record.requestKind === 'final-reply',
);
const finalizationLifecycle = agentDb.filter(
(record) =>
record.recordType === 'agent.runtime.finalization.lifecycle' &&
record.agentId === mainAgentId &&
(!state.initialRunId || record.runId === state.initialRunId),
);
const nonEmptyStreaming = state.responseStream.observedSnapshots.filter(
(snapshot) =>
snapshot.status === 'streaming' &&
snapshot.chars > 0 &&
snapshot.beforeTerminal,
);
const distinctStreaming = new Set(
nonEmptyStreaming.map((snapshot) => snapshot.fingerprint),
);
return {
effectiveStreamEnabled: state.responseStream.effectiveStreamEnabled,
streamOnlyConfigOverrideCreated: state.isolatedRunner.streamOverrideCreated,
isolatedAppDataUsed: Boolean(state.isolatedRunner.appDataDir),
formalConfigCliCallCount: state.isolatedRunner.sourceConfigCliCallCount,
sourceRunnerEndpointUnchanged:
state.isolatedRunner.sourceRunnerEndpointUnchanged,
sourceConfigHardlinkCount: state.isolatedRunner.configLinks.length,
sourceConfigLinksVerified: state.isolatedRunner.sourceConfigLinksVerified,
taskCount: taskSnapshot.all.length,
backgroundTaskEnqueueCount: isNonEmptyString(state.initialRunId) ? 1 : 0,
targetRunCount: new Set(
taskSnapshot.all
.filter((task) => task.agentId === mainAgentId)
.map((task) => task.runId),
).size,
eventCount: eventSurface.records.length,
agentDbRecordCount: agentDb.length,
conversationMessageCount: conversations.length,
responseStreamFileCount: stream ? 1 : 0,
responseStreamObservedSnapshotCount:
state.responseStream.observedSnapshots.length,
responseStreamPollCount: state.responseStream.pollCount,
responseStreamNonEmptyStreamingSnapshotCount: nonEmptyStreaming.length,
responseStreamDistinctStreamingSnapshotCount: distinctStreaming.size,
responseStreamFirstStreamingSequence:
nonEmptyStreaming[0]?.sequence ?? null,
responseStreamLastStreamingSequence:
nonEmptyStreaming.at(-1)?.sequence ?? null,
responseStreamCommittedSequence:
stream?.status === 'committed' ? stream.sequence : null,
responseStreamSnapshotsBeforeTerminal:
distinctStreaming.size >= 2 &&
state.responseStream.firstTerminalPoll != null,
responseStreamFinalChars:
typeof stream?.accumulatedText === 'string'
? [...stream.accumulatedText].length
: 0,
responseStreamFinalFingerprint:
typeof stream?.accumulatedText === 'string' &&
stream.accumulatedText.length > 0
? hashValue(stream.accumulatedText)
: null,
responseStreamRequestSlotHash: isNonEmptyString(stream?.requestSlot)
? hashValue(stream.requestSlot)
: null,
providerLifecycleStartedCount: providerLifecycle.filter(
(record) => record.status === 'started',
).length,
providerLifecycleTerminalCount: providerLifecycle.filter((record) =>
['completed', 'failed', 'interrupted'].includes(record.status),
).length,
finalizationStageCount: finalizationLifecycle.length,
finalAssistantCount: conversations.filter(
(message) => message.role === 'assistant',
).length,
finalAssistantAuditCount: agentDb.filter(
(record) =>
record.recordType === 'conversation.message' &&
record.role === 'assistant',
).length,
duplicateMessageCount: duplicateCount(
conversations.map((message) => message.messageId).filter(Boolean),
),
failureEvidenceErrors: {
task: taskSurface.errors,
event: eventSurface.errors,
agentDb: agentDbSurface.errors,
conversation: conversationSurface.errors,
},
paths: [
'.agent/runtime/response-streams',
'.agent/runtime/tasks',
'.agent/runtime/events',
'.agent/agent.db',
'.agent/conversations',
],
};
}
export function isResponseStreamSuite() {
return state.suite === responseStreamSuite;
}