Skip to content

Commit ae13e57

Browse files
authored
Merge pull request #152 from Reloaded-Project/event-abstractions
Stream framework-owned RunEvents from HookedAgent::run_stream
2 parents 9d4726f + fdb8e3b commit ae13e57

11 files changed

Lines changed: 1710 additions & 150 deletions

File tree

‎src/Cargo.lock‎

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎src/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,7 @@ serdes-ai-models = { version = "0.2.6", default-features = false }
7676
serdes-ai-streaming = "0.2"
7777

7878
# Internal crates
79-
reloaded-code-core = { version = "0.2.0", path = "reloaded-code-core", default-features = false }
79+
reloaded-code-core = { version = "0.2.2", path = "reloaded-code-core", default-features = false }
8080
reloaded-code-bubblewrap = { version = "0.1.0", path = "reloaded-code-bubblewrap" }
8181
reloaded-code-agents = { version = "0.1.0", path = "reloaded-code-agents" }
8282
reloaded-code-models-dev = { version = "0.1.0", path = "reloaded-code-models-dev" }

‎src/reloaded-code-core/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[package]
22
name = "reloaded-code-core"
3-
version = "0.2.0"
3+
version = "0.2.2"
44
edition = "2021"
55
description = "Lightweight, high-performance core types and utilities for coding tools - framework agnostic"
66
repository = "https://github.com/Reloaded-Project/ReloadedCode"

‎src/reloaded-code-core/src/hooks/mod.rs‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,13 @@
2020
//! - [`HookRunContext`] - Context given to hook run lifecycle events
2121
//! - [`EndReason`] - Why a run ended
2222
//!
23+
//! Run event types:
24+
//! - [`RunEvent`] - Framework-owned event yielded by run streams
25+
//! - [`RunMessage`] - Distilled transcript message in a completed run
26+
//! - [`RunMessageRole`] - Author role of a transcript message
27+
//! - [`RunToolCallSummary`] - Distilled tool call summary
28+
//! - [`RunToolResultSummary`] - Distilled tool result summary
29+
//!
2330
//! Observers are plain hooks: code before `original` is "start", code
2431
//! after is "end". They participate in the same hook chain.
2532
//!
@@ -37,11 +44,13 @@
3744
3845
pub use self::builder::HookSetBuilder;
3946
pub use self::hook_set::HookSet;
47+
pub use self::run_event::*;
4048
pub use self::run_hook::*;
4149
pub use self::tool_hook::*;
4250

4351
mod builder;
4452
mod hook_set;
53+
mod run_event;
4554
mod run_hook;
4655
mod tool_hook;
4756

Lines changed: 333 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,333 @@
1+
//! Run event types: the framework-owned streaming item type.
2+
//!
3+
//! [`RunEvent`] is the item type a run stream yields. Adapters
4+
//! translate their vendor-specific stream events into it, so consumers
5+
//! match one stable framework-owned enum instead of vendor types.
6+
//!
7+
//! # Transcript distillation
8+
//!
9+
//! [`RunEvent::RunComplete`] carries a distilled transcript
10+
//! ([`RunMessage`]): each message records its role, text, and tool
11+
//! call/result summaries. The record serves display and audit;
12+
//! consumers needing model-replay detail must use the underlying
13+
//! agent directly.
14+
//!
15+
//! # Extensibility
16+
//!
17+
//! [`RunEvent`] is `#[non_exhaustive]`: variants may be appended
18+
//! without a breaking release. Consumers match it with a wildcard arm.
19+
20+
use serde::{Deserialize, Serialize};
21+
22+
/// Framework-owned event yielded by a run stream.
23+
///
24+
/// One variant per observable streaming milestone: run start, step
25+
/// boundaries, context telemetry, text and thinking deltas, tool
26+
/// activity (call start, argument deltas, call complete, executed),
27+
/// output-ready, run complete, error, and cancellation.
28+
///
29+
/// The enum is `#[non_exhaustive]`: variants may be appended in a
30+
/// future release without a breaking change, so matches outside this
31+
/// crate need a wildcard arm.
32+
#[non_exhaustive]
33+
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
34+
pub enum RunEvent {
35+
/// The run started.
36+
RunStart {
37+
/// Identifier of the started run.
38+
run_id: String,
39+
},
40+
/// A model-request step started; the step's tool events may follow
41+
/// the matching [`RunEvent::StepEnd`] (see that variant).
42+
///
43+
/// Optional: emitted only by backends that report step boundaries;
44+
/// absence is normal.
45+
StepStart {
46+
/// Index of the started step; numbering is backend-defined.
47+
step: u32,
48+
},
49+
/// Context-size telemetry measured while building a model request.
50+
///
51+
/// Optional: emitted only by backends that report context metrics;
52+
/// absence is normal.
53+
ContextInfo {
54+
/// Estimated token count of the request.
55+
estimated_tokens: usize,
56+
/// Serialized request size in bytes (messages plus tools).
57+
request_bytes: usize,
58+
/// Model's context window limit, when known.
59+
context_limit: Option<u64>,
60+
},
61+
/// Context was compressed to fit within limits.
62+
///
63+
/// Optional: emitted only by backends that compress context and
64+
/// report it; absence is normal.
65+
ContextCompressed {
66+
/// Token count before compression.
67+
original_tokens: usize,
68+
/// Token count after compression.
69+
compressed_tokens: usize,
70+
/// Strategy used, e.g. "truncate" or "summarize".
71+
strategy: String,
72+
/// Number of messages before compression.
73+
messages_before: usize,
74+
/// Number of messages after compression.
75+
messages_after: usize,
76+
},
77+
/// Incremental assistant text arrived.
78+
TextDelta {
79+
/// Text fragment appended since the previous delta.
80+
text: String,
81+
},
82+
/// Incremental reasoning text arrived (reasoning models).
83+
ThinkingDelta {
84+
/// Thinking fragment appended since the previous delta.
85+
text: String,
86+
},
87+
/// A tool call started; its arguments may still be streaming.
88+
ToolCallStart {
89+
/// Name of the tool being called.
90+
tool_name: String,
91+
/// Call id correlating this call with its completion and result.
92+
tool_call_id: Option<String>,
93+
},
94+
/// Incremental tool-call argument fragment arrived.
95+
///
96+
/// Optional: a backend may omit it; absence is normal. Complete
97+
/// arguments remain available in the [`RunEvent::RunComplete`]
98+
/// transcript as [`RunToolCallSummary::arguments_json`].
99+
ToolCallDelta {
100+
/// Call id correlating this fragment with its call.
101+
tool_call_id: Option<String>,
102+
/// Argument fragment appended since the previous delta.
103+
delta: String,
104+
},
105+
/// A tool call's arguments finished streaming.
106+
ToolCallComplete {
107+
/// Name of the tool being called.
108+
tool_name: String,
109+
/// Call id correlating this call with its start and result.
110+
tool_call_id: Option<String>,
111+
},
112+
/// A tool finished executing.
113+
ToolExecuted {
114+
/// Name of the tool that ran.
115+
tool_name: String,
116+
/// Call id of the executed call.
117+
tool_call_id: Option<String>,
118+
/// Whether the tool reported success.
119+
success: bool,
120+
/// Error text when the tool failed.
121+
error: Option<String>,
122+
},
123+
/// A model-request step's response finished.
124+
///
125+
/// Backends that execute tools inside a step emit this before the
126+
/// step's tool calls run, so tool events may arrive after the
127+
/// matching [`RunEvent::StepEnd`].
128+
///
129+
/// Optional: emitted only by backends that report step boundaries;
130+
/// absence is normal.
131+
StepEnd {
132+
/// Index of the finished step, matching its
133+
/// [`RunEvent::StepStart`].
134+
step: u32,
135+
},
136+
/// The run's final output is ready to consume.
137+
OutputReady,
138+
/// The run completed.
139+
RunComplete {
140+
/// Identifier of the completed run.
141+
run_id: String,
142+
/// Distilled transcript of the run.
143+
messages: Vec<RunMessage>,
144+
},
145+
/// The run failed.
146+
Error {
147+
/// Human-readable description of the failure.
148+
message: String,
149+
},
150+
/// The run was cancelled.
151+
Cancelled {
152+
/// Partial text accumulated before cancellation.
153+
partial_text: Option<String>,
154+
/// Partial thinking content accumulated before cancellation.
155+
partial_thinking: Option<String>,
156+
/// Tool names whose calls were still in progress when
157+
/// cancelled.
158+
pending_tools: Vec<String>,
159+
},
160+
}
161+
162+
/// Distilled transcript message carried by [`RunEvent::RunComplete`].
163+
///
164+
/// One participant turn: what it said ([`Self::text`]), which tools it
165+
/// requested ([`Self::tool_calls`]), or which tool result it returns
166+
/// ([`Self::tool_result`]). Missing fields mean the message carries no
167+
/// content of that kind.
168+
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
169+
pub struct RunMessage {
170+
/// Author of the message.
171+
pub role: RunMessageRole,
172+
/// Text content of the message, if any.
173+
pub text: Option<String>,
174+
/// Tool calls the message requests, in order.
175+
pub tool_calls: Vec<RunToolCallSummary>,
176+
/// Tool result the message returns, if any.
177+
pub tool_result: Option<RunToolResultSummary>,
178+
}
179+
180+
/// Author role of a [`RunMessage`].
181+
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
182+
pub enum RunMessageRole {
183+
/// Framework-level instruction.
184+
System,
185+
/// Human input.
186+
User,
187+
/// Model output.
188+
Assistant,
189+
/// Tool output answering an assistant tool call.
190+
Tool,
191+
}
192+
193+
/// Distilled summary of one tool call requested during a run.
194+
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
195+
pub struct RunToolCallSummary {
196+
/// Name of the requested tool.
197+
pub tool_name: String,
198+
/// Call id correlating the call with its result.
199+
pub tool_call_id: Option<String>,
200+
/// JSON-serialized arguments, when the call carries any.
201+
pub arguments_json: Option<String>,
202+
}
203+
204+
/// Distilled summary of one tool result returned during a run.
205+
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
206+
pub struct RunToolResultSummary {
207+
/// Call id of the tool call this result answers.
208+
pub tool_call_id: Option<String>,
209+
/// Result payload rendered as text for display and audit.
210+
pub output: String,
211+
}
212+
213+
#[cfg(test)]
214+
mod tests {
215+
use super::*;
216+
217+
/// Full-shape transcript: every role, a tool call, and a tool result.
218+
fn complete_transcript() -> Vec<RunMessage> {
219+
vec![
220+
RunMessage {
221+
role: RunMessageRole::System,
222+
text: Some("sys".into()),
223+
tool_calls: Vec::new(),
224+
tool_result: None,
225+
},
226+
RunMessage {
227+
role: RunMessageRole::User,
228+
text: Some("read a.txt".into()),
229+
tool_calls: Vec::new(),
230+
tool_result: None,
231+
},
232+
RunMessage {
233+
role: RunMessageRole::Assistant,
234+
text: Some("checking".into()),
235+
tool_calls: vec![RunToolCallSummary {
236+
tool_name: "read_file".into(),
237+
tool_call_id: Some("call_1".into()),
238+
arguments_json: Some(r#"{"path":"a.txt"}"#.into()),
239+
}],
240+
tool_result: None,
241+
},
242+
RunMessage {
243+
role: RunMessageRole::Tool,
244+
text: None,
245+
tool_calls: Vec::new(),
246+
tool_result: Some(RunToolResultSummary {
247+
tool_call_id: Some("call_1".into()),
248+
output: "contents".into(),
249+
}),
250+
},
251+
]
252+
}
253+
254+
#[test]
255+
fn run_complete_serde_roundtrip_preserves_transcript() {
256+
let event = RunEvent::RunComplete {
257+
run_id: "run-42".into(),
258+
messages: complete_transcript(),
259+
};
260+
let json = serde_json::to_string(&event).unwrap();
261+
// Pin the wire shape: variants are externally tagged, so the
262+
// variant name is the JSON object key consumers see.
263+
assert!(json.starts_with("{\"RunComplete\":"));
264+
let restored: RunEvent = serde_json::from_str(&json).unwrap();
265+
assert_eq!(restored, event);
266+
}
267+
268+
#[test]
269+
fn run_event_variants_serde_roundtrip() {
270+
let events = vec![
271+
RunEvent::RunStart {
272+
run_id: "run-42".into(),
273+
},
274+
RunEvent::StepStart { step: 0 },
275+
RunEvent::ContextInfo {
276+
estimated_tokens: 128,
277+
request_bytes: 512,
278+
context_limit: Some(8192),
279+
},
280+
RunEvent::ContextInfo {
281+
estimated_tokens: 128,
282+
request_bytes: 512,
283+
context_limit: None,
284+
},
285+
RunEvent::ContextCompressed {
286+
original_tokens: 9000,
287+
compressed_tokens: 4000,
288+
strategy: "truncate".into(),
289+
messages_before: 20,
290+
messages_after: 6,
291+
},
292+
RunEvent::TextDelta {
293+
text: "chunk".into(),
294+
},
295+
RunEvent::ThinkingDelta {
296+
text: "thought".into(),
297+
},
298+
RunEvent::ToolCallStart {
299+
tool_name: "read_file".into(),
300+
tool_call_id: Some("call_1".into()),
301+
},
302+
RunEvent::ToolCallDelta {
303+
tool_call_id: Some("call_1".into()),
304+
delta: "{\"path\":".into(),
305+
},
306+
RunEvent::ToolCallComplete {
307+
tool_name: "read_file".into(),
308+
tool_call_id: Some("call_1".into()),
309+
},
310+
RunEvent::ToolExecuted {
311+
tool_name: "read_file".into(),
312+
tool_call_id: Some("call_1".into()),
313+
success: false,
314+
error: Some("missing".into()),
315+
},
316+
RunEvent::StepEnd { step: 0 },
317+
RunEvent::OutputReady,
318+
RunEvent::Error {
319+
message: "boom".into(),
320+
},
321+
RunEvent::Cancelled {
322+
partial_text: Some("partial".into()),
323+
partial_thinking: None,
324+
pending_tools: vec!["read_file".into()],
325+
},
326+
];
327+
for event in events {
328+
let json = serde_json::to_string(&event).unwrap();
329+
let restored: RunEvent = serde_json::from_str(&json).unwrap();
330+
assert_eq!(restored, event);
331+
}
332+
}
333+
}

0 commit comments

Comments
 (0)