Files
simstudioai--sim/apps/sim/providers/streaming-execution.test.ts
WeHub Mirror 6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
WeHub snapshot of cb28d14c6f2c081de7a0d8729a8c816c9adef67a
2026-08-10 11:17:50 +08:00

289 lines
9.8 KiB
TypeScript

/**
* @vitest-environment node
*/
import { describe, expect, it, vi } from 'vitest'
import type { NormalizedBlockOutput } from '@/executor/types'
import { createAgentEventReadableStream } from '@/providers/stream-events'
import { createStreamingExecution } from '@/providers/streaming-execution'
/**
* Builds a fake stream factory mirroring the providers' `createReadableStreamFrom*`
* helpers: it returns a sentinel stream and synchronously invokes the drain
* callback so the test can assert the populated output without a real stream.
*/
function fakeStreamFactory(
drain: (handles: { output: NormalizedBlockOutput; finalizeTiming: () => void }) => void
) {
const stream = new ReadableStream()
return {
stream,
createStream: (handles: { output: NormalizedBlockOutput; finalizeTiming: () => void }) => {
drain(handles)
return stream
},
}
}
describe('createStreamingExecution', () => {
const providerStartTime = 1_000
const providerStartTimeISO = new Date(providerStartTime).toISOString()
it('assembles the simple (no-tools) shape and finalizes timing on drain', () => {
const drainTime = 5_000
vi.spyOn(Date, 'now').mockReturnValue(drainTime)
const { stream, createStream } = fakeStreamFactory(({ output, finalizeTiming }) => {
output.content = 'hello'
output.tokens = { input: 10, output: 20, total: 30 }
output.cost = { input: 0.1, output: 0.2, total: 0.3 }
finalizeTiming()
})
const result = createStreamingExecution({
model: 'test-model',
providerStartTime,
providerStartTimeISO,
timing: { kind: 'simple', segmentName: 'test-model' },
initialTokens: { input: 0, output: 0, total: 0 },
initialCost: { input: 0, output: 0, total: 0 },
isStreaming: true,
createStream,
})
expect(result.stream).toBe(stream)
const output = result.execution.output
expect(output.content).toBe('hello')
expect(output.model).toBe('test-model')
expect(output.tokens).toEqual({ input: 10, output: 20, total: 30 })
expect(output.cost).toEqual({ input: 0.1, output: 0.2, total: 0.3 })
expect(output.toolCalls).toBeUndefined()
const timing = output.providerTiming
expect(timing?.startTime).toBe(providerStartTimeISO)
expect(timing?.endTime).toBe(new Date(drainTime).toISOString())
expect(timing?.duration).toBe(drainTime - providerStartTime)
expect(timing?.modelTime).toBeUndefined()
const segment = timing?.timeSegments?.[0]
expect(segment).toMatchObject({
type: 'model',
name: 'test-model',
startTime: providerStartTime,
})
expect(segment?.endTime).toBe(drainTime)
expect(segment?.duration).toBe(drainTime - providerStartTime)
expect(result.execution.success).toBe(true)
expect(result.execution.logs).toEqual([])
expect(result.execution.isStreaming).toBe(true)
expect(result.execution.metadata?.startTime).toBe(providerStartTimeISO)
vi.restoreAllMocks()
})
it('assembles the accumulated (post-tools) shape with pre-built segments', () => {
const drainTime = 7_000
vi.spyOn(Date, 'now').mockReturnValue(drainTime)
const timeSegments = [
{ type: 'model' as const, name: 'iter 1', startTime: 1_000, endTime: 2_000, duration: 1_000 },
{ type: 'tool' as const, name: 'lookup', startTime: 2_000, endTime: 2_500, duration: 500 },
]
const { createStream } = fakeStreamFactory(({ output }) => {
output.content = 'final'
output.tokens = { input: 110, output: 220, total: 330 }
output.cost = { input: 1.1, output: 2.2, toolCost: 0.5, total: 3.8 }
})
const result = createStreamingExecution({
model: 'tool-model',
providerStartTime,
providerStartTimeISO,
timing: {
kind: 'accumulated',
modelTime: 1_500,
toolsTime: 500,
firstResponseTime: 800,
iterations: 2,
timeSegments,
},
initialTokens: { input: 100, output: 200, total: 300 },
initialCost: { input: 1, output: 2, toolCost: undefined, total: 3 },
toolCalls: { list: [{ name: 'lookup' }], count: 1 },
isStreaming: true,
createStream,
})
const output = result.execution.output
expect(output.content).toBe('final')
expect(output.tokens).toEqual({ input: 110, output: 220, total: 330 })
expect(output.cost).toEqual({ input: 1.1, output: 2.2, toolCost: 0.5, total: 3.8 })
expect(output.toolCalls).toEqual({ list: [{ name: 'lookup' }], count: 1 })
const timing = output.providerTiming
expect(timing?.modelTime).toBe(1_500)
expect(timing?.toolsTime).toBe(500)
expect(timing?.firstResponseTime).toBe(800)
expect(timing?.iterations).toBe(2)
expect(timing?.timeSegments).toBe(timeSegments)
expect(timing?.startTime).toBe(providerStartTimeISO)
expect(timing?.endTime).toBe(new Date(drainTime).toISOString())
expect(timing?.duration).toBe(drainTime - providerStartTime)
vi.restoreAllMocks()
})
it('finalizes timing when the provider stream closes', async () => {
const constructTime = 1_200
const drainTime = 4_500
const nowMock = vi.spyOn(Date, 'now').mockReturnValue(constructTime)
const result = createStreamingExecution({
model: 'no-finalize',
providerStartTime,
providerStartTimeISO,
timing: {
kind: 'accumulated',
modelTime: 0,
toolsTime: 0,
firstResponseTime: 0,
iterations: 1,
timeSegments: [],
},
initialTokens: { input: 0, output: 0, total: 0 },
initialCost: { input: 0, output: 0, total: 0 },
createStream: ({ output }) => {
output.content = 'no-timing-mutation'
return new ReadableStream({
start(controller) {
controller.close()
},
})
},
})
nowMock.mockReturnValue(drainTime)
await result.stream.getReader().read()
const timing = result.execution.output.providerTiming
expect(timing?.endTime).toBe(new Date(drainTime).toISOString())
expect(timing?.duration).toBe(drainTime - providerStartTime)
expect(result.execution.metadata?.endTime).toBe(new Date(drainTime).toISOString())
expect(result.execution.metadata?.duration).toBe(drainTime - providerStartTime)
expect(result.execution.isStreaming).toBeUndefined()
vi.restoreAllMocks()
})
it('finalizeTiming touches only top-level aggregate for accumulated timing', () => {
const constructTime = 1_000
const drainTime = 9_000
const nowMock = vi.spyOn(Date, 'now').mockReturnValue(constructTime)
const segment = { type: 'model' as const, name: 's', startTime: 1, endTime: 2, duration: 1 }
const result = createStreamingExecution({
model: 'm',
providerStartTime,
providerStartTimeISO,
timing: {
kind: 'accumulated',
modelTime: 0,
toolsTime: 0,
firstResponseTime: 0,
iterations: 1,
timeSegments: [segment],
},
initialTokens: { input: 0, output: 0, total: 0 },
initialCost: { input: 0, output: 0, total: 0 },
createStream: ({ finalizeTiming }) => {
nowMock.mockReturnValue(drainTime)
finalizeTiming()
return new ReadableStream()
},
})
const timing = result.execution.output.providerTiming
expect(timing?.endTime).toBe(new Date(drainTime).toISOString())
expect(timing?.duration).toBe(drainTime - providerStartTime)
expect(timing?.timeSegments?.[0]).toEqual({
type: 'model',
name: 's',
startTime: 1,
endTime: 2,
duration: 1,
})
vi.restoreAllMocks()
})
it('propagates cancellation and finalizes timing', async () => {
const constructTime = 1_100
const cancelTime = 3_200
const nowMock = vi.spyOn(Date, 'now').mockReturnValue(constructTime)
const onCancel = vi.fn()
const result = createStreamingExecution({
model: 'cancel-model',
providerStartTime,
providerStartTimeISO,
timing: { kind: 'simple', segmentName: 'cancel-model' },
initialTokens: { input: 0, output: 0, total: 0 },
initialCost: { input: 0, output: 0, total: 0 },
createStream: () =>
new ReadableStream({
cancel: onCancel,
}),
})
nowMock.mockReturnValue(cancelTime)
await result.stream.cancel('consumer disconnected')
expect(onCancel).toHaveBeenCalledWith('consumer disconnected')
expect(result.execution.output.providerTiming?.duration).toBe(cancelTime - providerStartTime)
expect(result.execution.metadata?.duration).toBe(cancelTime - providerStartTime)
vi.restoreAllMocks()
})
it('defaults streamFormat to text and can attach agent-events-v1 object streams', async () => {
const textResult = createStreamingExecution({
model: 'm',
providerStartTime,
providerStartTimeISO,
timing: { kind: 'simple', segmentName: 'm' },
initialTokens: { input: 0, output: 0, total: 0 },
initialCost: { input: 0, output: 0, total: 0 },
createStream: () => new ReadableStream(),
})
expect(textResult.streamFormat).toBe('text')
const events = [
{ type: 'thinking_delta' as const, text: 'reason' },
{ type: 'text_delta' as const, text: 'answer', turn: 'final' as const },
]
const eventResult = createStreamingExecution({
model: 'm',
providerStartTime,
providerStartTimeISO,
timing: { kind: 'simple', segmentName: 'm' },
initialTokens: { input: 0, output: 0, total: 0 },
initialCost: { input: 0, output: 0, total: 0 },
streamFormat: 'agent-events-v1',
createStream: () => createAgentEventReadableStream(events),
})
expect(eventResult.streamFormat).toBe('agent-events-v1')
const reader = eventResult.stream.getReader()
const received = []
while (true) {
const { done, value } = await reader.read()
if (done) break
received.push(value)
}
expect(received).toEqual(events)
})
})