- 1
//! Structured sandbox execution events for live observation. - 2
//! - 3
//! Execution tools (such as bash) emit these events through a - 4
//! sink threaded into `ToolContext`. The events flow through `AgentEvent` → - 5
//! `EventBus` → SSE → the Workbench panel, giving the user full visibility - 6
//! into what the sandbox/command execution is doing without having to trust the chat summary. - 7
- 8
use serde::Serialize; - 9
use tokio::sync::mpsc; - 10
- 11
/// A single sandbox execution event. Serialized through `AgentEvent::Sandbox` - 12
/// and delivered to the client via SSE. - 13
#[derive(Debug, Clone, Serialize, serde::Deserialize)] - 14
#[serde(tag = "kind")] - 15
pub enum SandboxEvent { - 16
/// A sandbox execution is starting. - 17
ExecutionStarted { - 18
execution_id: String, - 19
/// Session that owns this execution. This is the stable rehydration - 20
/// key when a parent delegates work to a child session. - 21
owner_session_id: Option<String>, - 22
tool: String, - 23
/// The code or command about to execute (first ~2000 chars). - 24
code_preview: String, - 25
/// Language hint for syntax highlighting. - 26
language: String, - 27
/// The scratch directory being used. - 28
scratch_dir: String, - 29
}, - 30
- 31
/// A chunk of stdout arrived from the running process. - 32
Stdout { - 33
execution_id: String, - 34
chunk: String, - 35
}, - 36
- 37
/// A chunk of stderr arrived from the running process. - 38
Stderr { - 39
execution_id: String, - 40
chunk: String, - 41
}, - 42
OutputTruncated { - 43
execution_id: String, - 44
}, - 45
- 46
/// A package was installed in the sandbox environment. - 47
PackageInstalled { - 48
execution_id: String, - 49
packages: Vec<String>, - 50
}, - 51
- 52
/// A file artifact was generated in the scratch directory. - 53
ArtifactGenerated { - 54
execution_id: String, - 55
path: String, - 56
/// e.g. "image/png", "text/csv", "text/html" - 57
mime_type: String, - 58
size_bytes: u64, - 59
}, - 60
- 61
/// Periodic resource usage telemetry for long-running executions. - 62
ProcessTelemetry { - 63
execution_id: String, - 64
elapsed_ms: u64, - 65
cpu_percent: f32, - 66
memory_bytes: u64, - 67
}, - 68
- 69
/// The execution finished. - 70
ExecutionFinished { - 71
execution_id: String, - 72
exit_code: i32, - 73
duration_ms: u64, - 74
/// Paths of all artifacts generated during this execution. - 75
artifacts: Vec<String>, - 76
}, - 77
} - 78
- 79
/// Channel-based sink that sandbox tools write events into. The agent loop - 80
/// reads the receiving end and forwards events as `AgentEvent::Sandbox`. - 81
#[derive(Clone)] - 82
pub struct SandboxEventSink { - 83
tx: mpsc::UnboundedSender<SandboxEvent>, - 84
execution_id: String, - 85
owner_session_id: Option<String>, - 86
} - 87
- 88
impl SandboxEventSink { - 89
pub fn new() -> (Self, mpsc::UnboundedReceiver<SandboxEvent>) { - 90
Self::new_with_id("unidentified".to_string()) - 91
} - 92
- 93
pub fn new_with_id(execution_id: String) -> (Self, mpsc::UnboundedReceiver<SandboxEvent>) { - 94
let (tx, rx) = mpsc::unbounded_channel(); - 95
( - 96
SandboxEventSink { - 97
tx, - 98
execution_id, - 99
owner_session_id: None, - 100
}, - 101
rx, - 102
) - 103
} - 104
- 105
pub fn with_owner_session(mut self, session_id: impl Into<String>) -> Self { - 106
self.owner_session_id = Some(session_id.into()); - 107
self - 108
} - 109
- 110
/// Emit an event. Best-effort: dropped if the receiver is gone. - 111
pub fn emit(&self, event: SandboxEvent) { - 112
let event = match event { - 113
SandboxEvent::ExecutionStarted { - 114
execution_id, - 115
owner_session_id, - 116
tool, - 117
code_preview, - 118
language, - 119
scratch_dir, - 120
} => SandboxEvent::ExecutionStarted { - 121
execution_id, - 122
owner_session_id: owner_session_id.or_else(|| self.owner_session_id.clone()), - 123
tool, - 124
code_preview, - 125
language, - 126
scratch_dir, - 127
}, - 128
other => other, - 129
}; - 130
let _ = self.tx.send(event); - 131
} - 132
- 133
pub fn execution_id(&self) -> &str { - 134
&self.execution_id - 135
} - 136
- 137
pub fn emit_execution_started( - 138
&self, - 139
tool: &str, - 140
code: &str, - 141
language: &str, - 142
scratch_dir: &str, - 143
) { - 144
let preview = if code.chars().count() > 2000 { - 145
let truncated: String = code.chars().take(2000).collect(); - 146
format!("{truncated}…") - 147
} else { - 148
code.to_string() - 149
}; - 150
self.emit(SandboxEvent::ExecutionStarted { - 151
execution_id: self.execution_id.clone(), - 152
owner_session_id: self.owner_session_id.clone(), - 153
tool: tool.to_string(), - 154
code_preview: preview, - 155
language: language.to_string(), - 156
scratch_dir: scratch_dir.to_string(), - 157
}); - 158
} - 159
- 160
pub fn emit_stdout(&self, chunk: &str) { - 161
if !chunk.is_empty() { - 162
self.emit(SandboxEvent::Stdout { - 163
execution_id: self.execution_id.clone(), - 164
chunk: chunk.to_string(), - 165
}); - 166
} - 167
} - 168
- 169
pub fn emit_stderr(&self, chunk: &str) { - 170
if !chunk.is_empty() { - 171
self.emit(SandboxEvent::Stderr { - 172
execution_id: self.execution_id.clone(), - 173
chunk: chunk.to_string(), - 174
}); - 175
} - 176
} - 177
- 178
pub fn emit_output_truncated(&self) { - 179
self.emit(SandboxEvent::OutputTruncated { - 180
execution_id: self.execution_id.clone(), - 181
}); - 182
} - 183
- 184
pub fn emit_packages_installed(&self, packages: &[String]) { - 185
if !packages.is_empty() { - 186
self.emit(SandboxEvent::PackageInstalled { - 187
execution_id: self.execution_id.clone(), - 188
packages: packages.to_vec(), - 189
}); - 190
} - 191
} - 192
- 193
pub fn emit_artifact(&self, path: &str, mime_type: &str, size_bytes: u64) { - 194
self.emit(SandboxEvent::ArtifactGenerated { - 195
execution_id: self.execution_id.clone(), - 196
path: path.to_string(), - 197
mime_type: mime_type.to_string(), - 198
size_bytes, - 199
}); - 200
} - 201
- 202
pub fn emit_telemetry(&self, elapsed_ms: u64, cpu_percent: f32, memory_bytes: u64) { - 203
self.emit(SandboxEvent::ProcessTelemetry { - 204
execution_id: self.execution_id.clone(), - 205
elapsed_ms, - 206
cpu_percent, - 207
memory_bytes, - 208
}); - 209
} - 210
- 211
pub fn emit_finished(&self, exit_code: i32, duration_ms: u64, artifacts: Vec<String>) { - 212
self.emit(SandboxEvent::ExecutionFinished { - 213
execution_id: self.execution_id.clone(), - 214
exit_code, - 215
duration_ms, - 216
artifacts, - 217
}); - 218
} - 219
} - 220
- 221
/// Folds carriage returns (`\r`) in terminal output so progress bars and - 222
/// in-place line rewrites don't produce duplicate or chaotic lines. - 223
pub fn fold_carriage_returns(text: &str) -> String { - 224
if !text.contains('\r') { - 225
return text.to_string(); - 226
} - 227
let has_trailing_newline = text.ends_with('\n'); - 228
let lines: Vec<&str> = text.lines().collect(); - 229
let mut out = String::with_capacity(text.len()); - 230
for (i, line) in lines.iter().enumerate() { - 231
if line.contains('\r') { - 232
if let Some(last) = line.split('\r').rfind(|s| !s.is_empty()) { - 233
out.push_str(last); - 234
} - 235
} else { - 236
out.push_str(line); - 237
} - 238
if i + 1 < lines.len() || has_trailing_newline { - 239
out.push('\n'); - 240
} - 241
} - 242
out - 243
} - 244
- 245
#[cfg(test)] - 246
mod tests { - 247
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 248
use super::*; - 249
- 250
#[test] - 251
fn sink_emits_and_receiver_collects() { - 252
let (sink, mut rx) = SandboxEventSink::new(); - 253
sink.emit_execution_started("bash", "echo 'hello'", "bash", "/workspace"); - 254
sink.emit_stdout("hello\n"); - 255
sink.emit_finished(0, 42, vec![]); - 256
- 257
let ev1 = rx.try_recv().unwrap(); - 258
assert!(matches!(ev1, SandboxEvent::ExecutionStarted { .. })); - 259
let ev2 = rx.try_recv().unwrap(); - 260
assert!(matches!(ev2, SandboxEvent::Stdout { .. })); - 261
let ev3 = rx.try_recv().unwrap(); - 262
assert!(matches!( - 263
ev3, - 264
SandboxEvent::ExecutionFinished { exit_code: 0, .. } - 265
)); - 266
} - 267
- 268
#[test] - 269
fn code_preview_truncation() { - 270
let (sink, mut rx) = SandboxEventSink::new(); - 271
let long_code = "x".repeat(3000); - 272
sink.emit_execution_started("bash", &long_code, "bash", "/workspace"); - 273
- 274
let ev = rx.try_recv().unwrap(); - 275
assert!(matches!(ev, SandboxEvent::ExecutionStarted { .. })); - 276
if let SandboxEvent::ExecutionStarted { code_preview, .. } = ev { - 277
assert!(code_preview.len() < 2100); - 278
assert!(code_preview.ends_with('…')); - 279
} - 280
} - 281
- 282
#[test] - 283
fn code_preview_truncation_is_utf8_safe() { - 284
let (sink, mut rx) = SandboxEventSink::new(); - 285
let code = format!("{}€", "x".repeat(1999)); - 286
sink.emit_execution_started("bash", &code, "bash", "."); - 287
let ev = rx.try_recv().unwrap(); - 288
assert!(matches!(ev, SandboxEvent::ExecutionStarted { .. })); - 289
} - 290
- 291
#[test] - 292
fn execution_start_records_owner_session_for_tree_rehydration() { - 293
let (sink, mut rx) = SandboxEventSink::new_with_id("exec-child".into()); - 294
let sink = sink.with_owner_session("session-child"); - 295
sink.emit_execution_started("bash", "echo hi", "bash", "/scratch"); - 296
let event = rx.try_recv().unwrap(); - 297
match event { - 298
SandboxEvent::ExecutionStarted { - 299
owner_session_id, .. - 300
} => { - 301
assert_eq!(owner_session_id.as_deref(), Some("session-child")); - 302
} - 303
_ => panic!("expected execution start"), - 304
} - 305
} - 306
- 307
#[test] - 308
fn empty_stdout_not_emitted() { - 309
let (sink, mut rx) = SandboxEventSink::new(); - 310
sink.emit_stdout(""); - 311
assert!(rx.try_recv().is_err()); - 312
} - 313
- 314
#[test] - 315
fn telemetry_emitted_and_collected() { - 316
let (sink, mut rx) = SandboxEventSink::new(); - 317
sink.emit_telemetry(1500, 12.5, 45_000_000); - 318
let ev = rx.try_recv().unwrap(); - 319
assert!(matches!( - 320
ev, - 321
SandboxEvent::ProcessTelemetry { - 322
elapsed_ms: 1500, - 323
memory_bytes: 45_000_000, - 324
.. - 325
} - 326
)); - 327
} - 328
- 329
#[test] - 330
fn test_fold_carriage_returns() { - 331
let raw = "Downloading: 10%\rDownloading: 50%\rDownloading: 100%\nDone!\n"; - 332
let folded = fold_carriage_returns(raw); - 333
assert_eq!(folded, "Downloading: 100%\nDone!\n"); - 334
- 335
let no_cr = "hello\nworld\n"; - 336
assert_eq!(fold_carriage_returns(no_cr), no_cr); - 337
} - 338
} - 339
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.