Files
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

130 lines
3.7 KiB
TypeScript

import { createLogger } from '@sim/logger'
import { toError } from '@sim/utils/errors'
const logger = createLogger('CopilotSseParser')
export class FatalSseEventError extends Error {}
function createParseFailure(message: string, preview: string): FatalSseEventError {
logger.error(message, { preview })
return new FatalSseEventError(message)
}
function normalizeSseLine(line: string): string {
return line.endsWith('\r') ? line.slice(0, -1) : line
}
/**
* Processes an SSE stream by calling onEvent for each parsed event.
*
* @param onEvent Called per parsed event. Return true to stop processing.
*/
export async function processSSEStream(
reader: ReadableStreamDefaultReader<Uint8Array>,
decoder: TextDecoder,
abortSignal: AbortSignal | undefined,
onEvent: (event: unknown) => boolean | undefined | Promise<boolean | undefined>
): Promise<void> {
let buffer = ''
try {
try {
while (true) {
if (abortSignal?.aborted) {
logger.info('SSE stream aborted by signal')
break
}
const { done, value } = await reader.read()
if (done) break
buffer += decoder.decode(value, { stream: true })
const lines = buffer.split('\n')
buffer = lines.pop() || ''
let stopped = false
for (const line of lines) {
const normalizedLine = normalizeSseLine(line)
if (abortSignal?.aborted) {
logger.info('SSE stream aborted mid-chunk (between events)')
return
}
if (!normalizedLine.trim()) continue
if (!normalizedLine.startsWith('data: ')) continue
const jsonStr = normalizedLine.slice(6)
if (jsonStr === '[DONE]') continue
let parsed: unknown
try {
parsed = JSON.parse(jsonStr)
} catch (error) {
const preview = jsonStr.slice(0, 200)
const detail = toError(error).message
throw createParseFailure(`Failed to parse SSE event JSON: ${detail}`, preview)
}
try {
if (await onEvent(parsed)) {
stopped = true
break
}
} catch (error) {
if (error instanceof FatalSseEventError) {
throw error
}
logger.warn('Failed to handle SSE event', {
preview: jsonStr.slice(0, 200),
error: toError(error).message,
})
}
}
if (stopped) break
}
} catch (error) {
const aborted =
abortSignal?.aborted || (error instanceof DOMException && error.name === 'AbortError')
if (aborted) {
logger.info('SSE stream read aborted')
return
}
throw error
}
const normalizedBuffer = normalizeSseLine(buffer)
if (normalizedBuffer.trim() && normalizedBuffer.startsWith('data: ')) {
const jsonStr = normalizedBuffer.slice(6)
if (jsonStr === '[DONE]') {
return
}
let parsed: unknown
try {
parsed = JSON.parse(jsonStr)
} catch (error) {
const preview = normalizedBuffer.slice(0, 200)
const detail = toError(error).message
throw createParseFailure(`Failed to parse final SSE buffer JSON: ${detail}`, preview)
}
try {
await onEvent(parsed)
} catch (error) {
if (error instanceof FatalSseEventError) {
throw error
}
logger.warn('Failed to handle final SSE event', {
preview: normalizedBuffer.slice(0, 200),
error: toError(error).message,
})
}
}
} finally {
try {
reader.releaseLock()
} catch {
logger.warn('Failed to release SSE reader lock')
}
}
}