| |
| |
| |
| |
| |
|
|
|
|
| import type { A2ATask } from "./taskManager";
|
|
|
| export interface SSEChunkEvent {
|
| jsonrpc: "2.0";
|
| method: "message/stream";
|
| params: {
|
| task: { id: string; state: string };
|
| chunk?: { type: string; content: string };
|
| metadata?: Record<string, unknown>;
|
| };
|
| }
|
|
|
| |
| |
|
|
| export function formatSSE(event: SSEChunkEvent): string {
|
| return `data: ${JSON.stringify(event)}\n\n`;
|
| }
|
|
|
| |
| |
|
|
| export function createChunkEvent(taskId: string, content: string): string {
|
| return formatSSE({
|
| jsonrpc: "2.0",
|
| method: "message/stream",
|
| params: {
|
| task: { id: taskId, state: "working" },
|
| chunk: { type: "text", content },
|
| },
|
| });
|
| }
|
|
|
| |
| |
|
|
| export function createCompletionEvent(taskId: string, metadata: Record<string, unknown>): string {
|
| return formatSSE({
|
| jsonrpc: "2.0",
|
| method: "message/stream",
|
| params: {
|
| task: { id: taskId, state: "completed" },
|
| metadata,
|
| },
|
| });
|
| }
|
|
|
| |
| |
|
|
| export function createHeartbeat(taskId: string): string {
|
| return `: heartbeat ${new Date().toISOString()}\n\n`;
|
| }
|
|
|
| |
| |
|
|
| export function createFailureEvent(taskId: string, error: string): string {
|
| return formatSSE({
|
| jsonrpc: "2.0",
|
| method: "message/stream",
|
| params: {
|
| task: { id: taskId, state: "failed" },
|
| metadata: { error },
|
| },
|
| });
|
| }
|
|
|
| |
| |
|
|
| export const SSE_HEADERS = {
|
| "Content-Type": "text/event-stream",
|
| "Cache-Control": "no-cache, no-transform",
|
| Connection: "keep-alive",
|
| "X-Accel-Buffering": "no",
|
| } as const;
|
|
|
| |
| |
| |
|
|
| export function createA2AStream(
|
| task: A2ATask,
|
| executeSkill: (
|
| task: A2ATask
|
| ) => Promise<{ artifacts: Array<{ content: string }>; metadata: Record<string, unknown> }>,
|
| abortSignal?: AbortSignal,
|
| lifecycle?: {
|
| onStart?: () => void;
|
| onEnd?: () => void;
|
| }
|
| ): ReadableStream<Uint8Array> {
|
| const encoder = new TextEncoder();
|
|
|
| return new ReadableStream({
|
| async start(controller) {
|
| lifecycle?.onStart?.();
|
|
|
| const heartbeatInterval = setInterval(() => {
|
| try {
|
| controller.enqueue(encoder.encode(createHeartbeat(task.id)));
|
| } catch {
|
|
|
| }
|
| }, 15_000);
|
|
|
| try {
|
|
|
| if (abortSignal?.aborted) {
|
| controller.enqueue(encoder.encode(createFailureEvent(task.id, "Cancelled")));
|
| controller.close();
|
| return;
|
| }
|
|
|
|
|
| const result = await executeSkill(task);
|
|
|
|
|
| for (const artifact of result.artifacts) {
|
| if (abortSignal?.aborted) break;
|
| controller.enqueue(encoder.encode(createChunkEvent(task.id, artifact.content)));
|
| }
|
|
|
| if (abortSignal?.aborted) {
|
| controller.enqueue(encoder.encode(createFailureEvent(task.id, "Cancelled")));
|
| return;
|
| }
|
|
|
|
|
| controller.enqueue(encoder.encode(createCompletionEvent(task.id, result.metadata)));
|
| } catch (err) {
|
| const msg = err instanceof Error ? err.message : String(err);
|
| controller.enqueue(encoder.encode(createFailureEvent(task.id, msg)));
|
| } finally {
|
| clearInterval(heartbeatInterval);
|
| lifecycle?.onEnd?.();
|
| controller.close();
|
| }
|
| },
|
| });
|
| }
|
|
|