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; }