- 1
//! Capability-checked communication between isolated agent workspaces. - 2
//! - 3
//! The optional Unix transport is the only supported transport. It keeps - 4
//! isolated Docker tasks off the network while preserving the same capability - 5
//! checks as in-process callers. - 6
- 7
use std::collections::{HashMap, HashSet, VecDeque}; - 8
use std::sync::{Arc, Mutex, MutexGuard}; - 9
- 10
use uuid::Uuid; - 11
- 12
const DEFAULT_MESSAGE_LIMIT: usize = 256 * 1024; - 13
const MAX_QUEUE_MESSAGES: usize = 256; - 14
const MAX_QUEUE_BYTES: usize = 16 * 1024 * 1024; - 15
const MAX_FRAME_BYTES: usize = 1024 * 1024; - 16
const CAPABILITY_TTL: std::time::Duration = std::time::Duration::from_secs(15 * 60); - 17
- 18
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] - 19
pub struct WorkspaceNetworkPolicy { - 20
pub enabled: bool, - 21
pub allowed_peers: HashSet<String>, - 22
pub max_message_bytes: usize, - 23
} - 24
- 25
impl Default for WorkspaceNetworkPolicy { - 26
fn default() -> Self { - 27
Self { - 28
enabled: false, - 29
allowed_peers: HashSet::new(), - 30
max_message_bytes: DEFAULT_MESSAGE_LIMIT, - 31
} - 32
} - 33
} - 34
- 35
#[derive(Debug, Clone, PartialEq, Eq)] - 36
pub struct AgentMessage { - 37
pub from_workspace: String, - 38
pub to_workspace: String, - 39
pub body: Vec<u8>, - 40
} - 41
- 42
#[derive(Debug, Clone, PartialEq, Eq)] - 43
pub struct BrokerCapability { - 44
pub workspace: String, - 45
token: String, - 46
} - 47
- 48
impl BrokerCapability { - 49
pub fn from_parts(workspace: impl Into<String>, token: impl Into<String>) -> Self { - 50
Self { - 51
workspace: workspace.into(), - 52
token: token.into(), - 53
} - 54
} - 55
- 56
pub fn token(&self) -> &str { - 57
&self.token - 58
} - 59
} - 60
- 61
#[derive(Clone)] - 62
pub struct AgentNetworkBroker { - 63
inner: Arc<Mutex<BrokerState>>, - 64
} - 65
- 66
pub struct AgentNetworkTool { - 67
broker: AgentNetworkBroker, - 68
workspace: String, - 69
} - 70
- 71
impl AgentNetworkTool { - 72
pub fn new(broker: AgentNetworkBroker, workspace: impl Into<String>) -> Self { - 73
Self { - 74
broker, - 75
workspace: workspace.into(), - 76
} - 77
} - 78
} - 79
- 80
#[async_trait::async_trait] - 81
impl vak_tools::Tool for AgentNetworkTool { - 82
fn name(&self) -> &str { - 83
"agent_network" - 84
} - 85
- 86
fn serves(&self) -> &'static [&'static str] { - 87
&["messaging"] - 88
} - 89
- 90
fn description(&self) -> &str { - 91
"Send a bounded message to, or receive a message from, an explicitly authorized agent workspace." - 92
} - 93
- 94
fn schema(&self) -> serde_json::Value { - 95
serde_json::json!({ - 96
"type": "object", - 97
"properties": { - 98
"action": {"type": "string", "enum": ["send", "receive"]}, - 99
"destination_workspace": {"type": "string"}, - 100
"message": {"type": "string", "maxLength": MAX_MESSAGE_SCHEMA_CHARS} - 101
}, - 102
"required": ["action"], - 103
"additionalProperties": false - 104
}) - 105
} - 106
- 107
async fn execute( - 108
&self, - 109
args: &serde_json::Value, - 110
_ctx: &vak_tools::ToolContext, - 111
) -> vak_tools::ToolOutput { - 112
let Some(capability) = self.broker.capability_for(&self.workspace) else { - 113
return vak_tools::ToolOutput::error("agent network is not enabled for this workspace"); - 114
}; - 115
match args.get("action").and_then(serde_json::Value::as_str) { - 116
Some("send") => { - 117
let Some(destination) = args - 118
.get("destination_workspace") - 119
.and_then(serde_json::Value::as_str) - 120
else { - 121
return vak_tools::ToolOutput::error( - 122
"agent_network send requires destination_workspace", - 123
); - 124
}; - 125
let Some(message) = args.get("message").and_then(serde_json::Value::as_str) else { - 126
return vak_tools::ToolOutput::error("agent_network send requires message"); - 127
}; - 128
match self - 129
.broker - 130
.send_to(&capability, destination, message.as_bytes().to_vec()) - 131
{ - 132
Ok(()) => vak_tools::ToolOutput::ok("message queued"), - 133
Err(error) => vak_tools::ToolOutput::error(error), - 134
} - 135
} - 136
Some("receive") => match self.broker.receive(&capability) { - 137
Ok(Some(message)) => vak_tools::ToolOutput::ok( - 138
serde_json::json!({ - 139
"from_workspace": message.from_workspace, - 140
"message": String::from_utf8_lossy(&message.body) - 141
}) - 142
.to_string(), - 143
), - 144
Ok(None) => vak_tools::ToolOutput::ok("no messages available"), - 145
Err(error) => vak_tools::ToolOutput::error(error), - 146
}, - 147
_ => vak_tools::ToolOutput::error("agent_network action must be send or receive"), - 148
} - 149
} - 150
- 151
fn claims(&self, _args: &serde_json::Value) -> vak_tools::ResourceClaims { - 152
vak_tools::ResourceClaims { - 153
exclusive: true, - 154
..Default::default() - 155
} - 156
} - 157
} - 158
- 159
const MAX_MESSAGE_SCHEMA_CHARS: usize = 262_144; - 160
- 161
impl Default for AgentNetworkBroker { - 162
fn default() -> Self { - 163
static BROKER: std::sync::OnceLock<Arc<Mutex<BrokerState>>> = std::sync::OnceLock::new(); - 164
Self { - 165
inner: BROKER - 166
.get_or_init(|| Arc::new(Mutex::new(BrokerState::default()))) - 167
.clone(), - 168
} - 169
} - 170
} - 171
- 172
#[derive(Default)] - 173
struct BrokerState { - 174
policies: HashMap<String, WorkspaceNetworkPolicy>, - 175
capabilities: HashMap<String, CapabilityRecord>, - 176
queues: HashMap<String, VecDeque<AgentMessage>>, - 177
} - 178
- 179
struct CapabilityRecord { - 180
token: String, - 181
expires_at: std::time::Instant, - 182
} - 183
- 184
impl AgentNetworkBroker { - 185
fn capability_for(&self, workspace: &str) -> Option<BrokerCapability> { - 186
let state = self.lock(); - 187
state - 188
.capabilities - 189
.get(workspace) - 190
.map(|record| BrokerCapability { - 191
workspace: workspace.to_string(), - 192
token: record.token.clone(), - 193
}) - 194
} - 195
fn lock(&self) -> MutexGuard<'_, BrokerState> { - 196
self.inner - 197
.lock() - 198
.unwrap_or_else(std::sync::PoisonError::into_inner) - 199
} - 200
- 201
pub fn register( - 202
&self, - 203
workspace: impl Into<String>, - 204
policy: WorkspaceNetworkPolicy, - 205
) -> BrokerCapability { - 206
let workspace = workspace.into(); - 207
let token = Uuid::now_v7().to_string(); - 208
let mut state = self.lock(); - 209
state.policies.insert(workspace.clone(), policy); - 210
state.capabilities.insert( - 211
workspace.clone(), - 212
CapabilityRecord { - 213
token: token.clone(), - 214
expires_at: std::time::Instant::now() + CAPABILITY_TTL, - 215
}, - 216
); - 217
state.queues.insert(workspace.clone(), VecDeque::new()); - 218
BrokerCapability { workspace, token } - 219
} - 220
- 221
/// Whether `workspace` holds a live registration whose policy enables - 222
/// networking with at least one peer — the condition under which its - 223
/// agent is offered the `agent_network` tool. Off by default: nothing is - 224
/// registered until an operator authorizes the workspace. - 225
pub fn is_enabled(&self, workspace: &str) -> bool { - 226
let state = self.lock(); - 227
state - 228
.policies - 229
.get(workspace) - 230
.is_some_and(|policy| policy.enabled && !policy.allowed_peers.is_empty()) - 231
&& state - 232
.capabilities - 233
.get(workspace) - 234
.is_some_and(|record| record.expires_at > std::time::Instant::now()) - 235
} - 236
- 237
pub fn revoke(&self, capability: &BrokerCapability) -> bool { - 238
let mut state = self.lock(); - 239
if state - 240
.capabilities - 241
.get(&capability.workspace) - 242
.is_some_and(|record| { - 243
record.token == capability.token && record.expires_at > std::time::Instant::now() - 244
}) - 245
{ - 246
state.capabilities.remove(&capability.workspace); - 247
state.policies.remove(&capability.workspace); - 248
state.queues.remove(&capability.workspace); - 249
true - 250
} else { - 251
false - 252
} - 253
} - 254
- 255
pub fn send( - 256
&self, - 257
sender: &BrokerCapability, - 258
destination: &BrokerCapability, - 259
body: Vec<u8>, - 260
) -> Result<(), String> { - 261
let mut state = self.lock(); - 262
validate_capability(&state, sender)?; - 263
validate_capability(&state, destination)?; - 264
if sender.workspace == destination.workspace { - 265
return Err("agent network requires a distinct destination workspace".into()); - 266
} - 267
let source = state - 268
.policies - 269
.get(&sender.workspace) - 270
.ok_or_else(|| "sender workspace is not registered".to_string())?; - 271
let target = state - 272
.policies - 273
.get(&destination.workspace) - 274
.ok_or_else(|| "destination workspace is not registered".to_string())?; - 275
if !source.enabled || !target.enabled { - 276
return Err("agent network is disabled for one workspace".into()); - 277
} - 278
if !source.allowed_peers.contains(&destination.workspace) - 279
|| !target.allowed_peers.contains(&sender.workspace) - 280
{ - 281
return Err("workspace-to-workspace edge is not authorized by both policies".into()); - 282
} - 283
let limit = source.max_message_bytes.min(target.max_message_bytes); - 284
if body.len() > limit { - 285
return Err(format!("agent message exceeds {limit} byte limit")); - 286
} - 287
let queue = state - 288
.queues - 289
.entry(destination.workspace.clone()) - 290
.or_default(); - 291
let queued_bytes: usize = queue.iter().map(|message| message.body.len()).sum(); - 292
if queue.len() >= MAX_QUEUE_MESSAGES - 293
|| queued_bytes.saturating_add(body.len()) > MAX_QUEUE_BYTES - 294
{ - 295
return Err("destination agent queue is full".into()); - 296
} - 297
queue.push_back(AgentMessage { - 298
from_workspace: sender.workspace.clone(), - 299
to_workspace: destination.workspace.clone(), - 300
body, - 301
}); - 302
Ok(()) - 303
} - 304
- 305
pub fn send_to( - 306
&self, - 307
sender: &BrokerCapability, - 308
destination_workspace: &str, - 309
body: Vec<u8>, - 310
) -> Result<(), String> { - 311
let destination = { - 312
let state = self.lock(); - 313
let token = state - 314
.capabilities - 315
.get(destination_workspace) - 316
.ok_or_else(|| "destination workspace is not registered".to_string())? - 317
.token - 318
.clone(); - 319
BrokerCapability::from_parts(destination_workspace, token) - 320
}; - 321
self.send(sender, &destination, body) - 322
} - 323
- 324
pub fn receive(&self, receiver: &BrokerCapability) -> Result<Option<AgentMessage>, String> { - 325
let mut state = self.lock(); - 326
validate_capability(&state, receiver)?; - 327
Ok(state - 328
.queues - 329
.get_mut(&receiver.workspace) - 330
.and_then(VecDeque::pop_front)) - 331
} - 332
- 333
pub fn save_policies(&self, path: &std::path::Path) -> Result<(), String> { - 334
let state = self.lock(); - 335
let policies = state.policies.clone(); - 336
let parent = path - 337
.parent() - 338
.ok_or_else(|| "agent network policy file has no parent".to_string())?; - 339
std::fs::create_dir_all(parent) - 340
.map_err(|error| format!("agent network policy directory: {error}"))?; - 341
let text = serde_json::to_vec_pretty(&policies) - 342
.map_err(|error| format!("serialize agent network policies: {error}"))?; - 343
let temporary = path.with_extension("json.tmp"); - 344
std::fs::write(&temporary, text) - 345
.map_err(|error| format!("write agent network policies: {error}"))?; - 346
std::fs::rename(&temporary, path) - 347
.map_err(|error| format!("replace agent network policies: {error}"))?; - 348
Ok(()) - 349
} - 350
- 351
pub fn load_policies(&self, path: &std::path::Path) -> Result<usize, String> { - 352
let text = match std::fs::read_to_string(path) { - 353
Ok(text) => text, - 354
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(0), - 355
Err(error) => return Err(format!("read agent network policies: {error}")), - 356
}; - 357
let policies = serde_json::from_str::<HashMap<String, WorkspaceNetworkPolicy>>(&text) - 358
.map_err(|error| format!("parse agent network policies: {error}"))?; - 359
let mut loaded = 0; - 360
for (workspace, policy) in policies { - 361
let Some(workspace) = canonical_workspace(&workspace) else { - 362
continue; - 363
}; - 364
if policy - 365
.allowed_peers - 366
.iter() - 367
.any(|peer| canonical_workspace(peer).is_none()) - 368
{ - 369
continue; - 370
} - 371
self.register(workspace, policy); - 372
loaded += 1; - 373
} - 374
Ok(loaded) - 375
} - 376
- 377
pub fn socket_path(sessions_home: &std::path::Path) -> std::path::PathBuf { - 378
sessions_home.join("agent-network").join("broker.sock") - 379
} - 380
- 381
#[cfg(unix)] - 382
pub async fn serve_unix(&self, socket: &std::path::Path) -> Result<(), String> { - 383
use std::os::unix::fs::FileTypeExt; - 384
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; - 385
use tokio::net::UnixListener; - 386
- 387
let parent = socket - 388
.parent() - 389
.ok_or_else(|| "agent network socket has no parent".to_string())?; - 390
tokio::fs::create_dir_all(parent) - 391
.await - 392
.map_err(|error| format!("agent network socket directory: {error}"))?; - 393
if socket.exists() { - 394
let metadata = tokio::fs::symlink_metadata(socket) - 395
.await - 396
.map_err(|error| format!("agent network socket metadata: {error}"))?; - 397
if !metadata.file_type().is_socket() { - 398
return Err("agent network socket path is not a socket".into()); - 399
} - 400
if tokio::net::UnixStream::connect(socket).await.is_ok() { - 401
return Err("agent network broker is already running".into()); - 402
} - 403
tokio::fs::remove_file(socket) - 404
.await - 405
.map_err(|error| format!("remove stale agent network socket: {error}"))?; - 406
} - 407
let listener = UnixListener::bind(socket) - 408
.map_err(|error| format!("bind agent network socket: {error}"))?; - 409
use std::os::unix::fs::PermissionsExt; - 410
tokio::fs::set_permissions(socket, std::fs::Permissions::from_mode(0o600)) - 411
.await - 412
.map_err(|error| format!("secure agent network socket: {error}"))?; - 413
- 414
loop { - 415
let (stream, _) = listener - 416
.accept() - 417
.await - 418
.map_err(|error| format!("accept agent network connection: {error}"))?; - 419
let broker = self.clone(); - 420
tokio::spawn(async move { - 421
let (read, mut write) = stream.into_split(); - 422
let mut reader = BufReader::new(read); - 423
let mut line = Vec::with_capacity(8192); - 424
loop { - 425
line.clear(); - 426
let mut complete = false; - 427
loop { - 428
let available = match reader.fill_buf().await { - 429
Ok(available) if !available.is_empty() => available, - 430
_ => break, - 431
}; - 432
let take = available - 433
.iter() - 434
.position(|byte| *byte == b'\n') - 435
.map_or(available.len(), |position| position + 1); - 436
if line.len().saturating_add(take) > MAX_FRAME_BYTES { - 437
return; - 438
} - 439
line.extend_from_slice(&available[..take]); - 440
reader.consume(take); - 441
if line.last() == Some(&b'\n') { - 442
line.pop(); - 443
complete = true; - 444
break; - 445
} - 446
} - 447
if !complete { - 448
break; - 449
} - 450
let Ok(line) = std::str::from_utf8(&line) else { - 451
break; - 452
}; - 453
let response = broker.handle_wire(line); - 454
let mut encoded = match serde_json::to_vec(&response) { - 455
Ok(value) => value, - 456
Err(_) => break, - 457
}; - 458
encoded.push(b'\n'); - 459
if write.write_all(&encoded).await.is_err() { - 460
break; - 461
} - 462
} - 463
}); - 464
} - 465
} - 466
- 467
fn handle_wire(&self, line: &str) -> serde_json::Value { - 468
let Ok(request) = serde_json::from_str::<serde_json::Value>(line) else { - 469
return serde_json::json!({"ok": false, "error": "invalid JSON"}); - 470
}; - 471
match request.get("op").and_then(serde_json::Value::as_str) { - 472
Some("send") => self.handle_send_wire(&request), - 473
Some("receive") => self.handle_receive_wire(&request), - 474
_ => { - 475
serde_json::json!({"ok": false, "error": "operation is not available on the agent socket"}) - 476
} - 477
} - 478
} - 479
- 480
fn handle_send_wire(&self, request: &serde_json::Value) -> serde_json::Value { - 481
use base64::Engine as _; - 482
let (Some(workspace), Some(token), Some(destination), Some(body)) = ( - 483
request.get("sender_workspace").and_then(|v| v.as_str()), - 484
request.get("capability").and_then(|v| v.as_str()), - 485
request - 486
.get("destination_workspace") - 487
.and_then(|v| v.as_str()), - 488
request.get("body").and_then(|v| v.as_str()), - 489
) else { - 490
return serde_json::json!({"ok": false, "error": "send fields are required"}); - 491
}; - 492
let Ok(body) = base64::engine::general_purpose::STANDARD.decode(body) else { - 493
return serde_json::json!({"ok": false, "error": "body must be standard base64"}); - 494
}; - 495
let sender = BrokerCapability::from_parts(workspace, token); - 496
match self.send_to(&sender, destination, body) { - 497
Ok(()) => serde_json::json!({"ok": true}), - 498
Err(error) => serde_json::json!({"ok": false, "error": error}), - 499
} - 500
} - 501
- 502
fn handle_receive_wire(&self, request: &serde_json::Value) -> serde_json::Value { - 503
use base64::Engine as _; - 504
let (Some(workspace), Some(token)) = ( - 505
request.get("workspace").and_then(|v| v.as_str()), - 506
request.get("capability").and_then(|v| v.as_str()), - 507
) else { - 508
return serde_json::json!({"ok": false, "error": "workspace and capability are required"}); - 509
}; - 510
let receiver = BrokerCapability::from_parts(workspace, token); - 511
match self.receive(&receiver) { - 512
Ok(Some(message)) => serde_json::json!({"ok": true, "message": { - 513
"from_workspace": message.from_workspace, - 514
"to_workspace": message.to_workspace, - 515
"body": base64::engine::general_purpose::STANDARD.encode(message.body) - 516
}}), - 517
Ok(None) => serde_json::json!({"ok": true, "message": null}), - 518
Err(error) => serde_json::json!({"ok": false, "error": error}), - 519
} - 520
} - 521
} - 522
- 523
fn canonical_workspace(raw: &str) -> Option<String> { - 524
let path = std::path::Path::new(raw).canonicalize().ok()?; - 525
path.is_dir().then(|| path.display().to_string()) - 526
} - 527
- 528
fn validate_capability(state: &BrokerState, capability: &BrokerCapability) -> Result<(), String> { - 529
if state - 530
.capabilities - 531
.get(&capability.workspace) - 532
.is_some_and(|record| { - 533
record.token == capability.token && record.expires_at > std::time::Instant::now() - 534
}) - 535
{ - 536
Ok(()) - 537
} else { - 538
Err("agent network capability is invalid or revoked".into()) - 539
} - 540
} - 541
- 542
#[cfg(test)] - 543
mod tests { - 544
use super::*; - 545
- 546
fn policy(peer: &str) -> WorkspaceNetworkPolicy { - 547
WorkspaceNetworkPolicy { - 548
enabled: true, - 549
allowed_peers: [peer.to_string()].into_iter().collect(), - 550
max_message_bytes: 8, - 551
} - 552
} - 553
- 554
#[test] - 555
fn requires_explicit_bidirectional_authorization() { - 556
let broker = AgentNetworkBroker::default(); - 557
let a_name = format!("a-{}", Uuid::now_v7()); - 558
let b_name = format!("b-{}", Uuid::now_v7()); - 559
let a = broker.register(a_name, policy(&b_name)); - 560
let b = broker.register( - 561
b_name, - 562
WorkspaceNetworkPolicy { - 563
enabled: true, - 564
..WorkspaceNetworkPolicy::default() - 565
}, - 566
); - 567
assert!(matches!( - 568
broker.send(&a, &b, b"hello".to_vec()), - 569
Err(error) if error.contains("both policies") - 570
)); - 571
} - 572
- 573
#[test] - 574
fn delivers_bounded_messages_and_revocation_closes_access() { - 575
let broker = AgentNetworkBroker::default(); - 576
let a_name = format!("a-{}", Uuid::now_v7()); - 577
let b_name = format!("b-{}", Uuid::now_v7()); - 578
let a = broker.register(a_name.clone(), policy(&b_name)); - 579
let b = broker.register(b_name, policy(&a.workspace)); - 580
assert!(broker.send(&a, &b, b"hello".to_vec()).is_ok()); - 581
let message = broker.receive(&b).ok().flatten(); - 582
assert_eq!(message.map(|message| message.body), Some(b"hello".to_vec())); - 583
assert!(broker.send(&a, &b, vec![0; 9]).is_err()); - 584
assert!(broker.revoke(&b)); - 585
assert!(broker.receive(&b).is_err()); - 586
assert!(broker.send(&a, &b, b"again".to_vec()).is_err()); - 587
} - 588
- 589
#[cfg(unix)] - 590
#[tokio::test] - 591
async fn unix_transport_preserves_broker_authorization() { - 592
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; - 593
use tokio::net::UnixStream; - 594
- 595
let root = tempfile::tempdir().ok(); - 596
let Some(root) = root else { return }; - 597
let a_dir = tempfile::tempdir().ok(); - 598
let b_dir = tempfile::tempdir().ok(); - 599
let (Some(a_dir), Some(b_dir)) = (a_dir, b_dir) else { - 600
return; - 601
}; - 602
let a = a_dir.path().display().to_string(); - 603
let b = b_dir.path().display().to_string(); - 604
let socket = root.path().join("broker.sock"); - 605
let broker = AgentNetworkBroker::default(); - 606
let a_cap = broker.register(a.clone(), policy(&b)); - 607
let b_cap = broker.register(b.clone(), policy(&a)); - 608
let server = broker.clone(); - 609
let server_socket = socket.clone(); - 610
let task = tokio::spawn(async move { - 611
let _ = server.serve_unix(&server_socket).await; - 612
}); - 613
for _ in 0..50 { - 614
if socket.exists() { - 615
break; - 616
} - 617
tokio::task::yield_now().await; - 618
} - 619
let stream = match UnixStream::connect(&socket).await { - 620
Ok(stream) => stream, - 621
Err(_) => { - 622
task.abort(); - 623
return; - 624
} - 625
}; - 626
let mut reader = BufReader::new(stream); - 627
use base64::Engine as _; - 628
let send = serde_json::json!({ - 629
"op": "send", "sender_workspace": a, - 630
"capability": a_cap.token(), "destination_workspace": b, - 631
"body": base64::engine::general_purpose::STANDARD.encode("hello") - 632
}); - 633
let mut line = serde_json::to_vec(&send).unwrap_or_default(); - 634
line.push(b'\n'); - 635
assert!(reader.get_mut().write_all(&line).await.is_ok()); - 636
let mut response = String::new(); - 637
assert!(reader.read_line(&mut response).await.is_ok()); - 638
assert!(response.contains("\"ok\":true")); - 639
- 640
let receive = serde_json::json!({ - 641
"op": "receive", "workspace": b_cap.workspace, - 642
"capability": b_cap.token() - 643
}); - 644
let mut line = serde_json::to_vec(&receive).unwrap_or_default(); - 645
line.push(b'\n'); - 646
assert!(reader.get_mut().write_all(&line).await.is_ok()); - 647
response.clear(); - 648
assert!(reader.read_line(&mut response).await.is_ok()); - 649
assert!(response.contains("from_workspace") && response.contains("aGVsbG8=")); - 650
task.abort(); - 651
} - 652
} - 653
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.