|
| 1 | +""" |
| 2 | +Agent Builder — AG-UI Event Conversion |
| 3 | +
|
| 4 | +Calculations: pure functions that convert ReAct execution data |
| 5 | +into AG-UI protocol events. No I/O, no database, no side effects. |
| 6 | +Each function takes data in, returns AG-UI event(s) out. |
| 7 | +""" |
| 8 | + |
| 9 | +import json |
| 10 | +import uuid |
| 11 | +from typing import Any, Optional |
| 12 | + |
| 13 | +from ag_ui.core import ( |
| 14 | + CustomEvent, |
| 15 | + EventType, |
| 16 | + RunErrorEvent, |
| 17 | + RunFinishedEvent, |
| 18 | + RunStartedEvent, |
| 19 | + StepFinishedEvent, |
| 20 | + StepStartedEvent, |
| 21 | + TextMessageContentEvent, |
| 22 | + TextMessageEndEvent, |
| 23 | + TextMessageStartEvent, |
| 24 | + ToolCallArgsEvent, |
| 25 | + ToolCallEndEvent, |
| 26 | + ToolCallResultEvent, |
| 27 | + ToolCallStartEvent, |
| 28 | + StateSnapshotEvent, |
| 29 | +) |
| 30 | +from ag_ui.encoder import EventEncoder |
| 31 | + |
| 32 | + |
| 33 | +# Shared encoder — stateless, safe to reuse |
| 34 | +_encoder = EventEncoder() |
| 35 | + |
| 36 | + |
| 37 | +# --------------------------------------------------------------------------- |
| 38 | +# Pure conversion functions |
| 39 | +# --------------------------------------------------------------------------- |
| 40 | + |
| 41 | +def run_started_event(thread_id: str, run_id: str) -> RunStartedEvent: |
| 42 | + return RunStartedEvent( |
| 43 | + type=EventType.RUN_STARTED, |
| 44 | + thread_id=thread_id, |
| 45 | + run_id=run_id, |
| 46 | + ) |
| 47 | + |
| 48 | + |
| 49 | +def run_finished_event(thread_id: str, run_id: str) -> RunFinishedEvent: |
| 50 | + return RunFinishedEvent( |
| 51 | + type=EventType.RUN_FINISHED, |
| 52 | + thread_id=thread_id, |
| 53 | + run_id=run_id, |
| 54 | + ) |
| 55 | + |
| 56 | + |
| 57 | +def run_error_event(message: str) -> RunErrorEvent: |
| 58 | + return RunErrorEvent( |
| 59 | + type=EventType.RUN_ERROR, |
| 60 | + message=message, |
| 61 | + ) |
| 62 | + |
| 63 | + |
| 64 | +def text_message_start(message_id: str, role: str = "assistant") -> TextMessageStartEvent: |
| 65 | + return TextMessageStartEvent( |
| 66 | + type=EventType.TEXT_MESSAGE_START, |
| 67 | + message_id=message_id, |
| 68 | + role=role, |
| 69 | + ) |
| 70 | + |
| 71 | + |
| 72 | +def text_message_content(message_id: str, delta: str) -> TextMessageContentEvent: |
| 73 | + return TextMessageContentEvent( |
| 74 | + type=EventType.TEXT_MESSAGE_CONTENT, |
| 75 | + message_id=message_id, |
| 76 | + delta=delta, |
| 77 | + ) |
| 78 | + |
| 79 | + |
| 80 | +def text_message_end(message_id: str) -> TextMessageEndEvent: |
| 81 | + return TextMessageEndEvent( |
| 82 | + type=EventType.TEXT_MESSAGE_END, |
| 83 | + message_id=message_id, |
| 84 | + ) |
| 85 | + |
| 86 | + |
| 87 | +def text_message_events(content: str, message_id: Optional[str] = None) -> list: |
| 88 | + """Convert a complete text into the three-event TEXT_MESSAGE sequence.""" |
| 89 | + msg_id = message_id or str(uuid.uuid4()) |
| 90 | + return [ |
| 91 | + text_message_start(msg_id, role="assistant"), |
| 92 | + text_message_content(msg_id, content), |
| 93 | + text_message_end(msg_id), |
| 94 | + ] |
| 95 | + |
| 96 | + |
| 97 | +def tool_call_start(tool_call_id: str, tool_name: str, parent_message_id: Optional[str] = None) -> ToolCallStartEvent: |
| 98 | + return ToolCallStartEvent( |
| 99 | + type=EventType.TOOL_CALL_START, |
| 100 | + tool_call_id=tool_call_id, |
| 101 | + tool_call_name=tool_name, |
| 102 | + parent_message_id=parent_message_id, |
| 103 | + ) |
| 104 | + |
| 105 | + |
| 106 | +def tool_call_args(tool_call_id: str, args: Any) -> ToolCallArgsEvent: |
| 107 | + args_str = json.dumps(args) if not isinstance(args, str) else args |
| 108 | + return ToolCallArgsEvent( |
| 109 | + type=EventType.TOOL_CALL_ARGS, |
| 110 | + tool_call_id=tool_call_id, |
| 111 | + delta=args_str, |
| 112 | + ) |
| 113 | + |
| 114 | + |
| 115 | +def tool_call_end(tool_call_id: str) -> ToolCallEndEvent: |
| 116 | + return ToolCallEndEvent( |
| 117 | + type=EventType.TOOL_CALL_END, |
| 118 | + tool_call_id=tool_call_id, |
| 119 | + ) |
| 120 | + |
| 121 | + |
| 122 | +def tool_call_result(tool_call_id: str, content: str, message_id: Optional[str] = None) -> ToolCallResultEvent: |
| 123 | + msg_id = message_id or str(uuid.uuid4()) |
| 124 | + return ToolCallResultEvent( |
| 125 | + type=EventType.TOOL_CALL_RESULT, |
| 126 | + message_id=msg_id, |
| 127 | + tool_call_id=tool_call_id, |
| 128 | + content=content, |
| 129 | + role="tool", |
| 130 | + ) |
| 131 | + |
| 132 | + |
| 133 | +def tool_call_events( |
| 134 | + tool_name: str, |
| 135 | + args: Any, |
| 136 | + result: str, |
| 137 | + tool_call_id: Optional[str] = None, |
| 138 | + parent_message_id: Optional[str] = None, |
| 139 | +) -> list: |
| 140 | + """Convert a complete tool invocation into the four-event TOOL_CALL sequence.""" |
| 141 | + tc_id = tool_call_id or str(uuid.uuid4()) |
| 142 | + return [ |
| 143 | + tool_call_start(tc_id, tool_name, parent_message_id), |
| 144 | + tool_call_args(tc_id, args), |
| 145 | + tool_call_end(tc_id), |
| 146 | + tool_call_result(tc_id, result), |
| 147 | + ] |
| 148 | + |
| 149 | + |
| 150 | +def step_started_event(step_name: str) -> StepStartedEvent: |
| 151 | + return StepStartedEvent( |
| 152 | + type=EventType.STEP_STARTED, |
| 153 | + step_name=step_name, |
| 154 | + ) |
| 155 | + |
| 156 | + |
| 157 | +def step_finished_event(step_name: str) -> StepFinishedEvent: |
| 158 | + return StepFinishedEvent( |
| 159 | + type=EventType.STEP_FINISHED, |
| 160 | + step_name=step_name, |
| 161 | + ) |
| 162 | + |
| 163 | + |
| 164 | +def state_snapshot_event(snapshot: dict[str, Any]) -> StateSnapshotEvent: |
| 165 | + return StateSnapshotEvent( |
| 166 | + type=EventType.STATE_SNAPSHOT, |
| 167 | + snapshot=snapshot, |
| 168 | + ) |
| 169 | + |
| 170 | + |
| 171 | +def structured_output_event( |
| 172 | + output: Any, |
| 173 | + schema: Optional[dict[str, Any]] = None, |
| 174 | +) -> CustomEvent: |
| 175 | + """Emit structured agent output as a CUSTOM event with widget metadata. |
| 176 | +
|
| 177 | + The frontend uses the schema's x-ui annotations to pick the right |
| 178 | + CopilotKit tool render for each field. |
| 179 | + """ |
| 180 | + return CustomEvent( |
| 181 | + type=EventType.CUSTOM, |
| 182 | + name="structured_output", |
| 183 | + value={ |
| 184 | + "output": output, |
| 185 | + "schema": schema or {}, |
| 186 | + }, |
| 187 | + ) |
| 188 | + |
| 189 | + |
| 190 | +# --------------------------------------------------------------------------- |
| 191 | +# SSE encoding — thin wrapper, still a pure calculation |
| 192 | +# --------------------------------------------------------------------------- |
| 193 | + |
| 194 | +def encode_event(event: Any) -> str: |
| 195 | + """Serialize an AG-UI event to an SSE data line.""" |
| 196 | + return _encoder.encode(event) |
0 commit comments