Files
Genarrative/src/services/creative-agent/creativeAgentSse.ts
T

134 lines
3.4 KiB
TypeScript

import type {
CreativeAgentSessionSnapshot,
CreativeAgentSseEvent,
CreativeDraftEditResult,
} from '../../../packages/shared/src/contracts/creativeAgent';
import type { TextStreamOptions } from '../aiTypes';
import { readSseJsonStream } from '../sseStream';
type CreativeAgentSseOptions = TextStreamOptions & {
fallbackMessage: string;
incompleteMessage: string;
onEvent?: (event: CreativeAgentSseEvent) => void;
};
type CreativeAgentSseResult = {
session: CreativeAgentSessionSnapshot | null;
draftEditResult: CreativeDraftEditResult | null;
};
function normalizeCreativeAgentSseEvent(
eventName: string,
data: Record<string, unknown>,
): CreativeAgentSseEvent | null {
switch (eventName) {
case 'stage':
case 'agent_message_delta':
case 'thought_summary_delta':
case 'puzzle_template_catalog':
case 'puzzle_template_selection':
case 'puzzle_cost_range':
case 'puzzle_level_plan':
case 'tool_started':
case 'tool_completed':
case 'reflection':
case 'target_session':
case 'session':
case 'error':
case 'done':
return {
event: eventName,
data: data as never,
} as CreativeAgentSseEvent;
default:
return null;
}
}
function handleParsedCreativeAgentEvent(
eventName: string,
parsed: Record<string, unknown>,
options: CreativeAgentSseOptions,
): Partial<CreativeAgentSseResult> | null {
const normalizedEvent = normalizeCreativeAgentSseEvent(eventName, parsed);
if (normalizedEvent) {
options.onEvent?.(normalizedEvent);
}
if (eventName === 'agent_message_delta') {
const textDelta = parsed.textDelta;
if (typeof textDelta === 'string') {
options.onUpdate?.(textDelta);
}
return null;
}
if (eventName === 'session' && parsed.session) {
return {
session: parsed.session as CreativeAgentSessionSnapshot,
draftEditResult:
'puzzleSession' in parsed
? (parsed as unknown as CreativeDraftEditResult)
: null,
};
}
if (eventName === 'draft_edit_result' && parsed.session) {
return {
session: parsed.session as CreativeAgentSessionSnapshot,
draftEditResult: parsed as unknown as CreativeDraftEditResult,
};
}
if (eventName === 'error') {
const message =
typeof parsed.message === 'string' && parsed.message.trim()
? parsed.message.trim()
: options.fallbackMessage;
throw new Error(message);
}
return null;
}
export async function readCreativeAgentSessionFromSse(
response: Response,
options: CreativeAgentSseOptions,
) {
const result = await readCreativeAgentResultFromSse(response, options);
if (!result.session) {
throw new Error(options.incompleteMessage);
}
return result.session;
}
export async function readCreativeAgentResultFromSse(
response: Response,
options: CreativeAgentSseOptions,
): Promise<CreativeAgentSseResult> {
const result: CreativeAgentSseResult = {
session: null,
draftEditResult: null,
};
await readSseJsonStream(response, ({ eventName, parsed }) => {
const nextResult = handleParsedCreativeAgentEvent(
eventName,
parsed,
options,
);
if (nextResult?.session) {
result.session = nextResult.session;
}
if (nextResult?.draftEditResult) {
result.draftEditResult = nextResult.draftEditResult;
}
});
if (!result.session) {
throw new Error(options.incompleteMessage);
}
return result;
}