- 1
//! Real API client for the vak server, replacing the previous hardcoded - 2
//! mock surface. Every piece of data the terminal renders is fetched live - 3
//! from the running server over HTTP/SSE, authenticated with a bearer token. - 4
//! - 5
//! (docs/design/55-rich-terminal-surface.md §6 — the real-time data plane.) - 6
- 7
use std::collections::HashMap; - 8
- 9
use futures::StreamExt; - 10
use reqwest::header::{AUTHORIZATION, CONTENT_TYPE}; - 11
use serde::de::DeserializeOwned; - 12
use serde::{Deserialize, Serialize}; - 13
use thiserror::Error; - 14
use tokio::sync::mpsc; - 15
- 16
// ---- Error type ------------------------------------------------------------- - 17
- 18
#[derive(Debug, Error)] - 19
pub enum ApiError { - 20
#[error("network error: {0}")] - 21
Network(String), - 22
- 23
#[error("server error: {status} {body}")] - 24
Server { status: u16, body: String }, - 25
- 26
#[error("parse error: {0}")] - 27
Parse(String), - 28
- 29
#[error("connection refused: {0}")] - 30
Unreachable(String), - 31
} - 32
- 33
// ---- Terminal events emitted by the SSE watcher ---- - 34
- 35
/// Events the terminal's main loop receives from background SSE tasks. - 36
#[derive(Debug, Clone)] - 37
pub enum TerminalEvent { - 38
/// A real agent event (TurnStart, ToolCall, ApprovalRequested, etc.). - 39
Agent { event: serde_json::Value, seq: u64 }, - 40
/// A real presentation timeline frame. - 41
Presentation { frame: serde_json::Value, seq: u64 }, - 42
/// Periodic health refresh. - 43
Health(HealthReport), - 44
/// Connection established to the server. - 45
Connected, - 46
/// An error occurred (stream dropped, parse failure, etc.). - 47
Error(String), - 48
} - 49
- 50
// ---- API client ------------------------------------------------------------- - 51
- 52
/// A bearer-token-authenticated HTTP+SSE client to the vak server. - 53
pub struct ApiClient { - 54
client: reqwest::Client, - 55
base_url: String, - 56
token: String, - 57
} - 58
- 59
impl ApiClient { - 60
pub fn new(base_url: impl Into<String>, token: impl Into<String>) -> Self { - 61
let client = reqwest::Client::builder() - 62
.timeout(std::time::Duration::from_secs(10)) - 63
.build() - 64
.unwrap_or_default(); - 65
Self { - 66
client, - 67
base_url: base_url.into(), - 68
token: token.into(), - 69
} - 70
} - 71
- 72
pub fn base_url(&self) -> &str { - 73
&self.base_url - 74
} - 75
- 76
fn auth(&self) -> String { - 77
format!("Bearer {}", self.token) - 78
} - 79
- 80
fn url(&self, path: &str) -> String { - 81
let base = self.base_url.trim_end_matches('/'); - 82
let path = path.trim_start_matches('/'); - 83
format!("{base}/{path}") - 84
} - 85
- 86
async fn get_json<T: DeserializeOwned>(&self, path: &str) -> Result<T, ApiError> { - 87
let url = self.url(path); - 88
let res = self - 89
.client - 90
.get(&url) - 91
.header(AUTHORIZATION, self.auth()) - 92
.send() - 93
.await - 94
.map_err(|e| ApiError::Network(e.to_string()))?; - 95
if !res.status().is_success() { - 96
return Err(ApiError::Server { - 97
status: res.status().as_u16(), - 98
body: res.text().await.unwrap_or_default(), - 99
}); - 100
} - 101
let json = res - 102
.json::<T>() - 103
.await - 104
.map_err(|e| ApiError::Parse(e.to_string()))?; - 105
Ok(json) - 106
} - 107
- 108
async fn post_json<T: Serialize>( - 109
&self, - 110
path: &str, - 111
body: &T, - 112
) -> Result<serde_json::Value, ApiError> { - 113
let url = self.url(path); - 114
let res = self - 115
.client - 116
.post(&url) - 117
.header(AUTHORIZATION, self.auth()) - 118
.header(CONTENT_TYPE, "application/json") - 119
.json(body) - 120
.send() - 121
.await - 122
.map_err(|e| ApiError::Network(e.to_string()))?; - 123
if !res.status().is_success() { - 124
return Err(ApiError::Server { - 125
status: res.status().as_u16(), - 126
body: res.text().await.unwrap_or_default(), - 127
}); - 128
} - 129
let json = res - 130
.json::<serde_json::Value>() - 131
.await - 132
.map_err(|e| ApiError::Parse(e.to_string()))?; - 133
Ok(json) - 134
} - 135
- 136
pub async fn health(&self) -> Result<HealthReport, ApiError> { - 137
let json = self.get_json::<serde_json::Value>("/health").await?; - 138
Ok(HealthReport::from_json(json)) - 139
} - 140
- 141
pub async fn list_sessions(&self) -> Result<Vec<SessionInfo>, ApiError> { - 142
let json = self.get_json::<serde_json::Value>("/sessions").await?; - 143
let arr = json - 144
.get("sessions") - 145
.and_then(|v| v.as_array()) - 146
.ok_or_else(|| ApiError::Parse("missing sessions array".into()))?; - 147
arr.iter() - 148
.map(|v| serde_json::from_value(v.clone()).map_err(|e| ApiError::Parse(e.to_string()))) - 149
.collect() - 150
} - 151
- 152
pub async fn create_session(&self) -> Result<String, ApiError> { - 153
let json = self.post_json("/sessions", &serde_json::json!({})).await?; - 154
json.get("session_id") - 155
.and_then(|v| v.as_str()) - 156
.map(String::from) - 157
.ok_or_else(|| ApiError::Parse("missing session_id".into())) - 158
} - 159
- 160
pub async fn attach_session(&self, id: &str) -> Result<AttachResponse, ApiError> { - 161
let json = self - 162
.post_json( - 163
&format!("/sessions/{}/attach", id), - 164
&serde_json::json!({ "session_id": id }), - 165
) - 166
.await?; - 167
serde_json::from_value(json).map_err(|e| ApiError::Parse(e.to_string())) - 168
} - 169
- 170
pub async fn list_providers(&self) -> Result<ProviderList, ApiError> { - 171
let json = self.get_json::<serde_json::Value>("/providers").await?; - 172
serde_json::from_value(json).map_err(|e| ApiError::Parse(e.to_string())) - 173
} - 174
- 175
pub async fn discover_models(&self, provider: &str) -> Result<Vec<String>, ApiError> { - 176
let json = self - 177
.get_json::<serde_json::Value>(&format!("/providers/{}/models", provider)) - 178
.await?; - 179
json.get("models") - 180
.and_then(|v| v.as_array()) - 181
.map(|arr| { - 182
arr.iter() - 183
.filter_map(|v| v.as_str().map(String::from)) - 184
.collect() - 185
}) - 186
.ok_or_else(|| ApiError::Parse("missing models array".into())) - 187
} - 188
- 189
pub async fn get_mcp_servers(&self) -> Result<HashMap<String, McpServerDef>, ApiError> { - 190
let json = self.get_json::<serde_json::Value>("/config/mcp").await?; - 191
let servers = json - 192
.get("servers") - 193
.and_then(|v| v.as_object()) - 194
.ok_or_else(|| ApiError::Parse("missing servers map".into()))?; - 195
let mut result = HashMap::new(); - 196
for (name, def) in servers { - 197
if let Ok(mcp) = serde_json::from_value::<McpServerDef>(def.clone()) { - 198
result.insert(name.clone(), mcp); - 199
} - 200
} - 201
Ok(result) - 202
} - 203
- 204
pub async fn get_gateway_approvals(&self) -> Result<GatewayApprovals, ApiError> { - 205
let json = self - 206
.get_json::<serde_json::Value>("/gateway/approvals") - 207
.await?; - 208
serde_json::from_value(json).map_err(|e| ApiError::Parse(e.to_string())) - 209
} - 210
- 211
pub async fn list_bots(&self) -> Result<Vec<BotInfo>, ApiError> { - 212
let json = self.get_json::<serde_json::Value>("/gateway/bots").await?; - 213
json.get("bots") - 214
.and_then(|v| v.as_array()) - 215
.map(|arr| { - 216
arr.iter() - 217
.filter_map(|v| serde_json::from_value(v.clone()).ok()) - 218
.collect() - 219
}) - 220
.ok_or_else(|| ApiError::Parse("missing bots array".into())) - 221
} - 222
- 223
pub async fn get_config(&self) -> Result<ConfigSnapshot, ApiError> { - 224
let json = self.get_json::<serde_json::Value>("/config").await?; - 225
serde_json::from_value(json).map_err(|e| ApiError::Parse(e.to_string())) - 226
} - 227
- 228
pub async fn get_control_state(&self, session_id: &str) -> Result<ControlState, ApiError> { - 229
let json = self - 230
.get_json::<serde_json::Value>(&format!("/sessions/{}/control-state", session_id)) - 231
.await?; - 232
Ok(ControlState { - 233
running: json - 234
.get("running") - 235
.and_then(|v| v.as_bool()) - 236
.unwrap_or(false), - 237
paused: json - 238
.get("paused") - 239
.and_then(|v| v.as_bool()) - 240
.unwrap_or(false), - 241
revision: json.get("revision").and_then(|v| v.as_i64()).unwrap_or(0) as u64, - 242
}) - 243
} - 244
- 245
pub async fn get_launch_servers( - 246
&self, - 247
session_id: &str, - 248
) -> Result<Vec<LaunchServer>, ApiError> { - 249
let json = self - 250
.get_json::<serde_json::Value>(&format!("/sessions/{}/launch", session_id)) - 251
.await?; - 252
let arr = json - 253
.get("servers") - 254
.and_then(|v| v.as_array()) - 255
.ok_or_else(|| ApiError::Parse("missing servers array".into()))?; - 256
arr.iter() - 257
.map(|v| serde_json::from_value(v.clone()).map_err(|e| ApiError::Parse(e.to_string()))) - 258
.collect() - 259
} - 260
- 261
pub async fn answer_approval( - 262
&self, - 263
session_id: &str, - 264
req_id: &str, - 265
approve: bool, - 266
remember: bool, - 267
) -> Result<ApprovalAnswer, ApiError> { - 268
let json = self - 269
.post_json( - 270
&format!("/sessions/{}/approvals/{}", session_id, req_id), - 271
&serde_json::json!({ "approve": approve, "remember": remember }), - 272
) - 273
.await?; - 274
serde_json::from_value(json).map_err(|e| ApiError::Parse(e.to_string())) - 275
} - 276
- 277
pub async fn run_prompt(&self, session_id: &str, prompt: &str) -> Result<(), ApiError> { - 278
let _ = self - 279
.post_json( - 280
&format!("/sessions/{}/run", session_id), - 281
&serde_json::json!({ "prompt": prompt }), - 282
) - 283
.await?; - 284
Ok(()) - 285
} - 286
- 287
pub async fn send_steering(&self, session_id: &str, text: &str) -> Result<(), ApiError> { - 288
let _ = self - 289
.post_json( - 290
&format!("/sessions/{}/steering", session_id), - 291
&serde_json::json!({ "text": text }), - 292
) - 293
.await?; - 294
Ok(()) - 295
} - 296
- 297
pub async fn set_permission_mode(&self, mode: &str) -> Result<(), ApiError> { - 298
let _ = self - 299
.post_json("/config/mode", &serde_json::json!({ "mode": mode })) - 300
.await?; - 301
Ok(()) - 302
} - 303
- 304
pub async fn patch_config(&self, patch: &ConfigPatch) -> Result<(), ApiError> { - 305
let body = serde_json::to_value(patch).unwrap_or_default(); - 306
let _ = self.post_json("/config", &body).await?; - 307
Ok(()) - 308
} - 309
- 310
pub async fn start_launch(&self, session_id: &str, name: &str) -> Result<bool, ApiError> { - 311
let json = self - 312
.post_json( - 313
&format!("/sessions/{}/launch/start", session_id), - 314
&serde_json::json!({ "name": name }), - 315
) - 316
.await?; - 317
Ok(json - 318
.get("started") - 319
.and_then(|v| v.as_bool()) - 320
.unwrap_or(false)) - 321
} - 322
- 323
// ---- SSE streaming ---- - 324
- 325
/// Connect to an SSE endpoint and drain ALL events from it through the - 326
/// channel, converting each parsed SSE frame into the appropriate - 327
/// `TerminalEvent` variant. Returns `Ok(())` when the remote closes - 328
/// the stream; the caller is expected to reconnect. - 329
async fn drain_sse( - 330
client: &reqwest::Client, - 331
url: &str, - 332
token: &str, - 333
tx: &mpsc::UnboundedSender<TerminalEvent>, - 334
kind: SseKind, - 335
) -> Result<(), ApiError> { - 336
let res = client - 337
.get(url) - 338
.header(AUTHORIZATION, token) - 339
.send() - 340
.await - 341
.map_err(|e| ApiError::Network(e.to_string()))?; - 342
- 343
if !res.status().is_success() { - 344
return Err(ApiError::Server { - 345
status: res.status().as_u16(), - 346
body: res.text().await.unwrap_or_default(), - 347
}); - 348
} - 349
- 350
let mut buf: Vec<u8> = Vec::new(); - 351
let mut stream = res.bytes_stream(); - 352
while let Some(chunk) = stream.next().await { - 353
match chunk { - 354
Ok(bytes) => { - 355
buf.extend_from_slice(&bytes); - 356
for ev in parse_sse_buffer(&mut buf) { - 357
let seq = ev - 358
.id - 359
.as_deref() - 360
.and_then(|s| s.parse::<u64>().ok()) - 361
.unwrap_or(0); - 362
let terminal_ev = match kind { - 363
SseKind::Agent => TerminalEvent::Agent { - 364
event: ev.event, - 365
seq, - 366
}, - 367
SseKind::Presentation => TerminalEvent::Presentation { - 368
frame: ev.event, - 369
seq, - 370
}, - 371
}; - 372
if tx.send(terminal_ev).is_err() { - 373
return Ok(()); // receiver dropped - 374
} - 375
} - 376
} - 377
Err(e) => return Err(ApiError::Network(e.to_string())), - 378
} - 379
} - 380
Ok(()) - 381
} - 382
- 383
/// Spawn a background task that keeps the agent-event SSE stream - 384
/// open (reconnecting on disconnect) and forwards `AgentEvent` JSON - 385
/// objects through the channel. - 386
pub fn spawn_agent_event_watcher( - 387
&self, - 388
session_id: &str, - 389
tx: mpsc::UnboundedSender<TerminalEvent>, - 390
) -> tokio::task::JoinHandle<()> { - 391
let url = self.url(&format!("/sessions/{}/events", session_id)); - 392
let token = self.auth(); - 393
let client = self.client.clone(); - 394
- 395
tokio::spawn(async move { - 396
let _ = tx.send(TerminalEvent::Connected); - 397
loop { - 398
match Self::drain_sse(&client, &url, &token, &tx, SseKind::Agent).await { - 399
Ok(()) => { - 400
// Stream closed normally — wait briefly and reconnect. - 401
tokio::time::sleep(std::time::Duration::from_secs(2)).await; - 402
} - 403
Err(e) => { - 404
let _ = tx.send(TerminalEvent::Error(e.to_string())); - 405
tokio::time::sleep(std::time::Duration::from_secs(2)).await; - 406
} - 407
} - 408
} - 409
}) - 410
} - 411
- 412
/// Spawn a background task that keeps the presentation SSE stream - 413
/// open and forwards timeline frames through the channel. - 414
pub fn spawn_presentation_watcher( - 415
&self, - 416
session_id: &str, - 417
tx: mpsc::UnboundedSender<TerminalEvent>, - 418
) -> tokio::task::JoinHandle<()> { - 419
let url = self.url(&format!("/sessions/{}/presentation/events", session_id)); - 420
let token = self.auth(); - 421
let client = self.client.clone(); - 422
- 423
tokio::spawn(async move { - 424
loop { - 425
match Self::drain_sse(&client, &url, &token, &tx, SseKind::Presentation).await { - 426
Ok(()) => { - 427
tokio::time::sleep(std::time::Duration::from_secs(2)).await; - 428
} - 429
Err(e) => { - 430
let _ = tx.send(TerminalEvent::Error(e.to_string())); - 431
tokio::time::sleep(std::time::Duration::from_secs(2)).await; - 432
} - 433
} - 434
} - 435
}) - 436
} - 437
- 438
/// Spawn a background task that polls `/health` periodically and - 439
/// forwards `HealthReport` through the channel. - 440
pub fn spawn_health_watcher( - 441
self: std::sync::Arc<Self>, - 442
tx: mpsc::UnboundedSender<TerminalEvent>, - 443
) -> tokio::task::JoinHandle<()> { - 444
tokio::spawn(async move { - 445
let mut interval = - 446
tokio::time::interval(std::time::Duration::from_secs(HEALTH_POLL_INTERVAL_SECS)); - 447
loop { - 448
interval.tick().await; - 449
match self.health().await { - 450
Ok(report) => { - 451
let _ = tx.send(TerminalEvent::Health(report)); - 452
} - 453
Err(e) => { - 454
let _ = tx.send(TerminalEvent::Error(e.to_string())); - 455
} - 456
} - 457
} - 458
}) - 459
} - 460
} - 461
- 462
#[derive(Debug, Clone, Copy)] - 463
enum SseKind { - 464
Agent, - 465
Presentation, - 466
} - 467
- 468
const HEALTH_POLL_INTERVAL_SECS: u64 = 3; - 469
- 470
// ---- Typed response structures ---------------------------------------------- - 471
- 472
/// Decoded from `GET /health`. - 473
#[derive(Debug, Clone, Default)] - 474
pub struct HealthReport { - 475
pub provider: String, - 476
pub model: String, - 477
pub permission_mode: String, - 478
pub sandbox: String, - 479
pub context_window: u64, - 480
pub cwd: String, - 481
pub healthy: bool, - 482
pub circuit_breaker_healthy: bool, - 483
pub failures: u32, - 484
pub warnings: Vec<String>, - 485
} - 486
- 487
impl HealthReport { - 488
pub fn from_json(json: serde_json::Value) -> Self { - 489
Self { - 490
provider: json - 491
.get("provider") - 492
.and_then(|v| v.as_str()) - 493
.unwrap_or("unknown") - 494
.to_string(), - 495
model: json - 496
.get("model") - 497
.and_then(|v| v.as_str()) - 498
.unwrap_or("unknown") - 499
.to_string(), - 500
permission_mode: json - 501
.get("permission_mode") - 502
.and_then(|v| v.as_str()) - 503
.unwrap_or("unknown") - 504
.to_string(), - 505
sandbox: json - 506
.get("sandbox") - 507
.and_then(|v| v.as_str()) - 508
.unwrap_or("unknown") - 509
.to_string(), - 510
context_window: json - 511
.get("context_window") - 512
.and_then(|v| v.as_u64()) - 513
.unwrap_or(0), - 514
cwd: json - 515
.get("cwd") - 516
.and_then(|v| v.as_str()) - 517
.unwrap_or("") - 518
.to_string(), - 519
healthy: json.get("posture").and_then(|v| v.as_str()) == Some("healthy"), - 520
circuit_breaker_healthy: json.get("failures").and_then(|v| v.as_u64()).unwrap_or(0) - 521
== 0, - 522
failures: json.get("failures").and_then(|v| v.as_u64()).unwrap_or(0) as u32, - 523
warnings: json - 524
.get("warnings") - 525
.and_then(|v| v.as_array()) - 526
.map(|arr| { - 527
arr.iter() - 528
.filter_map(|v| v.as_str().map(String::from)) - 529
.collect() - 530
}) - 531
.unwrap_or_default(), - 532
} - 533
} - 534
} - 535
- 536
/// Decoded from `GET /sessions`. - 537
#[derive(Debug, Clone, Deserialize)] - 538
pub struct SessionInfo { - 539
pub session_id: String, - 540
pub cwd: Option<String>, - 541
pub created_at: Option<String>, - 542
pub updated_at: Option<String>, - 543
pub entries: Option<u64>, - 544
pub title: Option<String>, - 545
pub running: Option<bool>, - 546
pub archived: Option<bool>, - 547
} - 548
- 549
impl SessionInfo { - 550
/// A human-readable status string derived from the real server state. - 551
pub fn status(&self) -> &str { - 552
if self.running == Some(true) { - 553
"active" - 554
} else if self.archived == Some(true) { - 555
"archived" - 556
} else { - 557
"idle" - 558
} - 559
} - 560
} - 561
- 562
#[derive(Debug, Clone, Deserialize)] - 563
pub struct AttachResponse { - 564
pub session_id: String, - 565
} - 566
- 567
/// Decoded from `GET /providers`. - 568
#[derive(Debug, Clone, Deserialize, Default)] - 569
pub struct ProviderList { - 570
pub current: String, - 571
pub current_model: String, - 572
pub current_configured: bool, - 573
pub providers: Vec<ProviderInfo>, - 574
} - 575
- 576
#[derive(Debug, Clone, Deserialize)] - 577
pub struct ProviderInfo { - 578
pub name: String, - 579
pub env_var: Option<String>, - 580
pub pool_env_var: Option<String>, - 581
pub pool_size: Option<usize>, - 582
pub credential_ids: Vec<String>, - 583
pub requires_key: bool, - 584
pub configured: bool, - 585
} - 586
- 587
/// Decoded from `GET /config`. - 588
#[derive(Debug, Clone, Serialize, Deserialize, Default)] - 589
pub struct ConfigSnapshot { - 590
#[serde(default)] - 591
pub provider: Option<String>, - 592
#[serde(default)] - 593
pub model: Option<String>, - 594
#[serde(default)] - 595
pub max_turns: Option<u32>, - 596
#[serde(default)] - 597
pub permission_mode: Option<String>, - 598
#[serde(default)] - 599
pub approval_mode: Option<String>, - 600
#[serde(default)] - 601
pub theme: Option<String>, - 602
} - 603
- 604
/// Patch body for `PATCH /config`. - 605
#[derive(Debug, Clone, Serialize, Default)] - 606
pub struct ConfigPatch { - 607
#[serde(skip_serializing_if = "Option::is_none")] - 608
pub provider: Option<String>, - 609
#[serde(skip_serializing_if = "Option::is_none")] - 610
pub model: Option<String>, - 611
#[serde(skip_serializing_if = "Option::is_none")] - 612
pub permission_mode: Option<String>, - 613
#[serde(skip_serializing_if = "Option::is_none")] - 614
pub approval_mode: Option<String>, - 615
#[serde(skip_serializing_if = "Option::is_none")] - 616
pub max_turns: Option<u32>, - 617
} - 618
- 619
/// MCP server definition from `GET /config/mcp`. - 620
#[derive(Debug, Clone, Serialize, Deserialize)] - 621
pub struct McpServerDef { - 622
pub command: String, - 623
#[serde(default)] - 624
pub args: Vec<String>, - 625
#[serde(default)] - 626
pub env: HashMap<String, String>, - 627
#[serde(default)] - 628
pub network: bool, - 629
} - 630
- 631
/// Decoded from `GET /gateway/approvals`. - 632
#[derive(Debug, Clone, Deserialize, Default)] - 633
pub struct GatewayApprovals { - 634
pub mode: String, - 635
pub approver: Option<String>, - 636
pub timeout_secs: u64, - 637
pub enabled: bool, - 638
pub forwarding: bool, - 639
pub candidates: Vec<String>, - 640
} - 641
- 642
/// Bot info from `GET /gateway/bots`. - 643
#[derive(Debug, Clone, Deserialize)] - 644
pub struct BotInfo { - 645
pub id: String, - 646
pub surface: String, - 647
pub name: String, - 648
pub configured: bool, - 649
#[serde(default)] - 650
pub policy_inherited: bool, - 651
} - 652
- 653
/// Decoded from `POST /sessions/{id}/approvals/{req_id}`. - 654
#[derive(Debug, Clone, Deserialize)] - 655
pub struct ApprovalAnswer { - 656
pub approved: bool, - 657
pub learned_rule: Option<String>, - 658
pub learn_error: Option<String>, - 659
} - 660
- 661
/// Decoded from `GET /sessions/{id}/control-state`. - 662
#[derive(Debug, Clone, Default)] - 663
pub struct ControlState { - 664
pub running: bool, - 665
pub paused: bool, - 666
pub revision: u64, - 667
} - 668
- 669
/// Decoded from `GET /sessions/{id}/launch`. - 670
#[derive(Debug, Clone, Deserialize)] - 671
pub struct LaunchServer { - 672
pub name: String, - 673
pub cmd: String, - 674
#[serde(default)] - 675
pub args: Vec<String>, - 676
pub port: Option<u16>, - 677
pub running: bool, - 678
} - 679
- 680
// ---- SSE frame parsing ------------------------------------------------------ - 681
- 682
/// Parsed SSE frame from the server's event streams. - 683
#[derive(Debug, Clone)] - 684
pub struct SseEvent { - 685
pub event: serde_json::Value, - 686
pub id: Option<String>, - 687
} - 688
- 689
/// Parsed SSE line: a field name and its value. - 690
#[derive(Debug, Clone)] - 691
struct SseLine { - 692
field: String, - 693
value: String, - 694
} - 695
- 696
impl SseLine { - 697
/// Parse a single SSE line (without the trailing newline). - 698
/// Returns `None` for comment lines (starting with `:`). - 699
fn parse(line: &str) -> Option<Self> { - 700
if line.is_empty() { - 701
return Some(SseLine { - 702
field: String::new(), - 703
value: String::new(), - 704
}); - 705
} - 706
if line.starts_with(':') { - 707
return None; - 708
} - 709
if let Some(colon_pos) = line.find(':') { - 710
let field = &line[..colon_pos]; - 711
let value = if line.as_bytes().get(colon_pos + 1) == Some(&b' ') { - 712
&line[colon_pos + 2..] - 713
} else { - 714
&line[colon_pos + 1..] - 715
}; - 716
Some(SseLine { - 717
field: field.to_string(), - 718
value: value.to_string(), - 719
}) - 720
} else { - 721
Some(SseLine { - 722
field: line.to_string(), - 723
value: String::new(), - 724
}) - 725
} - 726
} - 727
} - 728
- 729
/// Consume complete SSE frames from `buf` (a byte buffer of raw event data), - 730
/// returning decoded `SseEvent` objects. Bytes that don't form a complete - 731
/// frame (no trailing blank line) remain in `buf` for the next call. - 732
pub fn parse_sse_buffer(buf: &mut Vec<u8>) -> Vec<SseEvent> { - 733
let mut events = Vec::new(); - 734
let mut data_lines: Vec<String> = Vec::new(); - 735
let mut event_type = String::new(); - 736
let mut id = None::<String>; - 737
- 738
while let Some(nl) = buf.iter().position(|&b| b == b'\n') { - 739
let line_bytes: Vec<u8> = buf.drain(..=nl).collect(); - 740
let line_str = String::from_utf8_lossy(&line_bytes); - 741
let trimmed = line_str.trim_end_matches(['\r', '\n']); - 742
- 743
match SseLine::parse(trimmed) { - 744
None => {} // comment line - 745
Some(SseLine { - 746
field: f, - 747
value: _v, - 748
}) if f.is_empty() => { - 749
// blank line — dispatch accumulated event - 750
if !data_lines.is_empty() || !event_type.is_empty() || id.is_some() { - 751
let data = data_lines.join("\n"); - 752
let event = serde_json::from_str::<serde_json::Value>(&data).unwrap_or_else( - 753
|_| serde_json::json!({ "_raw": data, "_type": event_type }), - 754
); - 755
events.push(SseEvent { - 756
event, - 757
id: id.clone(), - 758
}); - 759
data_lines.clear(); - 760
event_type.clear(); - 761
id = None; - 762
} - 763
} - 764
Some(SseLine { field: f, value: v }) if f == "data" => { - 765
data_lines.push(v); - 766
} - 767
Some(SseLine { field: f, value: v }) if f == "event" => { - 768
event_type = v; - 769
} - 770
Some(SseLine { field: f, value: v }) if f == "id" => { - 771
id = Some(v); - 772
} - 773
_ => {} - 774
} - 775
} - 776
- 777
events - 778
} - 779
- 780
#[cfg(test)] - 781
mod tests { - 782
use super::*; - 783
- 784
#[test] - 785
fn test_parse_sse_single_event() { - 786
let mut buf = b"id:42\ndata:{\"test\":true}\n\n".to_vec(); - 787
let events = parse_sse_buffer(&mut buf); - 788
assert_eq!(events.len(), 1); - 789
assert_eq!(events[0].id, Some("42".to_string())); - 790
assert_eq!( - 791
events[0].event.get("test").and_then(|v| v.as_bool()), - 792
Some(true) - 793
); - 794
assert!(buf.is_empty()); - 795
} - 796
- 797
#[test] - 798
fn test_parse_sse_multi_line_data() { - 799
let mut buf = b"data:{\"a\":1,\ndata:\"b\":2}\n\n".to_vec(); - 800
let events = parse_sse_buffer(&mut buf); - 801
assert_eq!(events.len(), 1); - 802
assert_eq!(events[0].event.get("a").and_then(|v| v.as_i64()), Some(1)); - 803
assert_eq!(events[0].event.get("b").and_then(|v| v.as_i64()), Some(2)); - 804
} - 805
- 806
#[test] - 807
fn test_parse_sse_incomplete_buffer() { - 808
let mut buf = b"id:42\ndata:{\"test\":".to_vec(); - 809
let events = parse_sse_buffer(&mut buf); - 810
assert!(events.is_empty()); - 811
// The incomplete data should still be in the buffer - 812
assert!(!buf.is_empty()); - 813
} - 814
- 815
#[test] - 816
fn test_parse_sse_comment_lines() { - 817
let mut buf = b":keepalive\ndata:{\"ok\":true}\n\n".to_vec(); - 818
let events = parse_sse_buffer(&mut buf); - 819
assert_eq!(events.len(), 1); - 820
assert_eq!( - 821
events[0].event.get("ok").and_then(|v| v.as_bool()), - 822
Some(true) - 823
); - 824
} - 825
- 826
#[test] - 827
fn test_parse_sse_two_events() { - 828
let mut buf = b"id:1\ndata:{\"n\":1}\n\nid:2\ndata:{\"n\":2}\n\n".to_vec(); - 829
let events = parse_sse_buffer(&mut buf); - 830
assert_eq!(events.len(), 2); - 831
assert_eq!(events[0].event.get("n").and_then(|v| v.as_i64()), Some(1)); - 832
assert_eq!(events[1].event.get("n").and_then(|v| v.as_i64()), Some(2)); - 833
} - 834
- 835
#[test] - 836
fn test_url_construction() { - 837
let client = ApiClient::new("http://localhost:8901", "mytoken"); - 838
assert_eq!(client.url("/health"), "http://localhost:8901/health"); - 839
assert_eq!( - 840
client.url("/sessions/abc/events"), - 841
"http://localhost:8901/sessions/abc/events" - 842
); - 843
} - 844
} - 845
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.