- 347
.arg(effective) - 348
.current_dir(&ctx.cwd) - 349
.stdin(Stdio::piped()) - 350
.stdout(Stdio::piped()) - 351
.stderr(Stdio::piped()); - 352
crate::bash::scrub_environment(&mut cmd); - 353
cmd.env(WORKER_ENV, "1"); - 354
crate::bash::isolate_process_group(&mut cmd); - 355
- 356
let mut child = match cmd.spawn() { - 357
Ok(child) => child, - 358
Err(error) => return ToolOutput::error(format!("tool broker spawn failed: {error}")), - 359
}; - 360
let request = WorkerRequest { - 361
version: PROTOCOL_VERSION, - 362
task: WorkerTask::Tool { - 363
tool: tool.to_string(), - 364
args: request_args, - 365
execution_id: ctx - 366
.sandbox_sink - 367
.as_ref() - 368
.map(|sink| sink.execution_id().to_string()) - 369
.unwrap_or_else(|| "unidentified".into()), - 370
agent_id: ctx.agent_id.clone(), - 371
new_documents: new_documents.to_vec(), - 372
}, - 373
}; - 374
let payload = match serde_json::to_vec(&request) { - 375
Ok(payload) => payload, - 376
Err(error) => return ToolOutput::error(format!("tool broker encode failed: {error}")), - 377
}; - 378
let Some(mut stdin) = child.stdin.take() else { - 379
crate::bash::kill_process_group(&child.id()); - 380
return ToolOutput::error("tool broker has no stdin"); - 381
}; - 382
if let Err(error) = stdin.write_all(&payload).await { - 383
crate::bash::kill_process_group(&child.id()); - 384
let _ = child.wait().await; - 385
return ToolOutput::error(format!("tool broker request failed: {error}")); - 386
} - 387
drop(stdin); - 388
- 389
let Some(mut stdout) = child.stdout.take() else { - 390
return ToolOutput::error("tool broker has no stdout"); - 391
}; - 392
let Some(mut stderr) = child.stderr.take() else { - 393
return ToolOutput::error("tool broker has no stderr"); - 394
}; - 395
let sink = ctx.sandbox_sink.clone(); - 396
let partial = Arc::new(std::sync::Mutex::new(String::new())); - 397
let partial_reader = partial.clone(); - 398
let stderr_reader = tokio::spawn(async move { - 399
let mut bytes = Vec::new(); - 400
let mut line = Vec::new(); - 401
while let Ok(n) = stderr.read_buf(&mut line).await { - 402
if n == 0 { - 403
break; - 404
} - 405
while let Some(pos) = line.iter().position(|b| *b == b'\n') { - 406
let frame: Vec<u8> = line.drain(..=pos).collect(); - 407
if let Some(payload) = frame.strip_prefix(b"VAK_EVENT:") - 408
&& let Ok(event) = serde_json::from_slice::<SandboxEvent>(payload) - 409
{ - 410
if let Ok(mut output) = partial_reader.lock() { - 411
match &event { - 412
SandboxEvent::Stdout { chunk, .. } => { - 413
output.push_str(chunk); - 414
} - 415
SandboxEvent::Stderr { chunk, .. } => { - 416
output.push_str("[stderr]\n"); - 417
output.push_str(chunk); - 418
} - 419
_ => {} - 420
} - 421
} - 422
if let Some(ref sink) = sink { - 423
sink.emit(event); - 424
} - 425
} else { - 426
bytes.extend(frame); - 427
} - 428
} - 429
} - 430
bytes.extend(line); - 431
bytes - 432
}); - 433
let stdout_reader = tokio::spawn(async move { - 434
let mut bytes = Vec::new(); - 435
let _ = stdout.read_to_end(&mut bytes).await; - 436
bytes - 437
}); - 438
let pid = child.id(); - 439
let (status, cancelled) = tokio::select! { - 440
_ = ctx.cancel.cancelled() => { - 441
crate::bash::kill_process_group(&pid); - 442
(child.wait().await, true) - 443
} - 444
output = child.wait() => (output, false), - 445
}; - 446
let status = match status { - 447
Ok(status) => status, - 448
Err(error) => return ToolOutput::error(format!("tool broker wait failed: {error}")), - 449
}; - 450
let stdout = stdout_reader.await.unwrap_or_default(); - 451
let stderr = stderr_reader.await.unwrap_or_default(); - 452
if cancelled { - 453
let mut content = String::from("tool broker cancelled"); - 454
if let Ok(output) = partial.lock() - 455
&& !output.is_empty() - 456
{ - 457
content.push_str("\n\n[partial output]\n"); - 458
content.push_str(&output); - 459
} - 460
let diagnostics = String::from_utf8_lossy(&stderr).trim().to_string(); - 461
if !diagnostics.is_empty() { - 462
content.push_str("\n\n[diagnostics]\n"); - 463
content.push_str(&diagnostics); - 464
} - 465
let _ = stdout; - 466
return ToolOutput::error(content); - 467
} - 468
if stdout.len() as u64 > MAX_PROTOCOL_BYTES || stderr.len() as u64 > MAX_PROTOCOL_BYTES { - 469
return ToolOutput::error("tool broker output exceeded protocol limit"); - 470
} - 471
if !status.success() { - 472
let detail = String::from_utf8_lossy(&stderr); - 473
return ToolOutput::error(format!( - 474
"tool broker exited with {}: {}", - 475
status.code().unwrap_or(-1), - 476
detail.trim() - 477
)); - 478
} - 479
let response: WorkerResponse = match serde_json::from_slice(&stdout) { - 480
Ok(response) => response, - 481
Err(error) => { - 482
let detail = String::from_utf8_lossy(&stderr); - 483
return ToolOutput::error(format!( - 484
"tool broker returned invalid protocol: {error}; stderr: {}", - 485
detail.trim() - 486
)); - 487
} - 488
}; - 489
if response.version != PROTOCOL_VERSION { - 490
return ToolOutput::error(format!( - 491
"tool broker protocol mismatch: expected {PROTOCOL_VERSION}, got {}", - 492
response.version - 493
)); - 494
} - 495
if response.is_error { - 496
ToolOutput::error(response.content) - 497
} else { - 498
ToolOutput::ok(response.content) - 499
} - 500
} - 501
- 502
pub async fn worker_main() -> i32 { - 503
let mut bytes = Vec::new(); - 504
if tokio::io::stdin() - 505
.take(MAX_PROTOCOL_BYTES + 1) - 506
.read_to_end(&mut bytes) - 507
.await - 508
.is_err() - 509
|| bytes.len() as u64 > MAX_PROTOCOL_BYTES - 510
{ - 511
return 125; - 512
} - 513
let request: WorkerRequest = match serde_json::from_slice::<WorkerRequest>(&bytes) { - 514
Ok(request) if request.version == PROTOCOL_VERSION => request, - 515
_ => return 125, - 516
}; - 517
let (tool_name, args, execution_id, agent_id, new_documents) = match request.task { - 518
WorkerTask::Tool { - 519
tool, - 520
args, - 521
execution_id, - 522
agent_id, - 523
new_documents, - 524
} => (tool, args, execution_id, agent_id, new_documents), - 525
WorkerTask::OfficeReview { - 526
before, - 527
after, - 528
lineage, - 529
} => { - 530
let (content, is_error) = - 531
match office_review_in_worker(before.as_deref(), &after, lineage.as_ref()) { - 532
Ok(content) => (content, false), - 533
Err(error) => (error, true), - 534
}; - 535
return write_response(WorkerResponse { - 536
version: PROTOCOL_VERSION, - 537
content, - 538
is_error, - 539
events: Vec::new(), - 540
}) - 541
.await; - 542
} - 543
WorkerTask::OfficeProject { path, view } => { - 544
let (content, is_error) = match office_project_in_worker(&path, view) { - 545
Ok(content) => (content, false), - 546
Err(error) => (error, true), - 547
}; - 548
return write_response(WorkerResponse { - 549
version: PROTOCOL_VERSION, - 550
content, - 551
is_error, - 552
events: Vec::new(), - 553
}) - 554
.await; - 555
} - 556
WorkerTask::OfficeApply { lineage, out } => { - 557
let (content, is_error) = match office_apply_in_worker(&lineage, &out) { - 558
Ok(content) => (content, false), - 559
Err(error) => (error, true), - 560
}; - 561
return write_response(WorkerResponse { - 562
version: PROTOCOL_VERSION, - 563
content, - 564
is_error, - 565
events: Vec::new(), - 566
}) - 567
.await; - 568
} - 569
WorkerTask::OfficeNarrow { - 570
lineage, - 571
draft, - 572
keep, - 573
out, - 574
} => { - 575
let (content, is_error) = match office_narrow_in_worker(&lineage, &draft, &keep, &out) { - 576
Ok(content) => (content, false), - 577
Err(error) => (error, true), - 578
}; - 579
return write_response(WorkerResponse { - 580
version: PROTOCOL_VERSION, - 581
content, - 582
is_error, - 583
events: Vec::new(), - 584
}) - 585
.await; - 586
} - 587
WorkerTask::VerifyTargets { root, checks } => { - 588
let results = vak_sandbox::default_target_verifiers().verify(&root, &checks); - 589
let content = serde_json::to_string(&results).unwrap_or_default(); - 590
return write_response(WorkerResponse { - 591
version: PROTOCOL_VERSION, - 592
content, - 593
is_error: false, - 594
events: Vec::new(), - 595
}) - 596
.await; - 597
} - 598
}; - 599
let tool = crate::default_tools() - 600
.into_iter() - 601
.find(|candidate| candidate.name() == tool_name); - 602
let (output, events) = match tool { - 603
Some(tool) => { - 604
let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")); - 605
let (sink, mut rx) = crate::sandbox_events::SandboxEventSink::new_with_id(execution_id); - 606
let mut ctx = ToolContext::new(cwd) - 607
.with_sandbox_sink(sink) - 608
.with_new_documents(new_documents); - 609
if let Some(agent_id) = agent_id { - 610
ctx = ctx.with_agent_id(agent_id); - 611
} - 612
let event_forwarder = tokio::spawn(async move { - 613
let mut stderr = tokio::io::stderr(); - 614
while let Some(event) = rx.recv().await { - 615
let Ok(mut frame) = serde_json::to_vec(&event) else { - 616
continue; - 617
}; - 618
let mut prefixed = b"VAK_EVENT:".to_vec(); - 619
prefixed.append(&mut frame); - 620
prefixed.push(b'\n'); - 621
if stderr.write_all(&prefixed).await.is_err() { - 622
break; - 623
} - 624
let _ = stderr.flush().await; - 625
} - 626
}); - 627
let output = tool.execute(&args, &ctx).await; - 628
drop(ctx.sandbox_sink); - 629
let _ = event_forwarder.await; - 630
(output, Vec::new()) - 631
} - 632
None => ( - 633
ToolOutput::error(format!("worker does not expose tool '{tool_name}'")), - 634
Vec::new(), - 635
), - 636
}; - 637
write_response(WorkerResponse { - 638
version: PROTOCOL_VERSION, - 639
content: output.content, - 640
is_error: output.is_error, - 641
events, - 642
}) - 643
.await - 644
} - 645
- 646
async fn write_response(response: WorkerResponse) -> i32 { - 647
let payload = match serde_json::to_vec(&response) { - 648
Ok(payload) if payload.len() as u64 <= MAX_PROTOCOL_BYTES => payload, - 649
Ok(payload) => { - 650
let refusal = WorkerResponse { - 651
version: PROTOCOL_VERSION, - 652
content: format!( - 653
"the result is {} bytes, over the {} byte worker protocol limit; request a smaller range", - 654
payload.len(), - 655
MAX_PROTOCOL_BYTES - 656
), - 657
is_error: true, - 658
events: Vec::new(), - 659
}; - 660
match serde_json::to_vec(&refusal) { - 661
Ok(payload) => payload, - 662
Err(_) => return 125, - 663
} - 664
} - 665
Err(_) => return 125, - 666
}; - 667
let mut stdout = tokio::io::stdout(); - 668
if stdout.write_all(&payload).await.is_err() || stdout.flush().await.is_err() { - 669
return 125; - 670
} - 671
0 - 672
} - 673
- 674
/// Runs the registered target verifiers over `checks` under `root` in a - 675
/// worker process (docs/design/72-openxml-documents.md, F4). The worker runs - 676
/// under a read-only, network-denied sandbox rooted at `root` whatever the - 677
/// session's permission mode, and within [`VERIFY_DEADLINE`]. Anything that - 678
/// stops the worker from answering fails every planned check; a check never - 679
/// passes because verification could not run. - 680
pub async fn verify_targets( - 681
worker_exe: &Path, - 682
root: &Path, - 683
checks: &[vak_sandbox::TargetCheckPlan], - 684
) -> Vec<vak_sandbox::TargetCheckResult> { - 685
if checks.is_empty() { - 686
return Vec::new(); - 687
} - 688
let failed = |reason: String| { - 689
checks - 690
.iter() - 691
.map(|check| vak_sandbox::TargetCheckResult { - 692
verifier: check.verifier.clone(), - 693
path: check.path.clone(), - 694
status: "failed".into(), - 695
evidence: format!("verification did not run: {reason}"), - 696
}) - 697
.collect::<Vec<_>>() - 698
}; - 699
let task = WorkerTask::VerifyTargets { - 700
root: root.to_path_buf(), - 701
checks: checks.to_vec(), - 702
}; - 703
let content = match run_task(worker_exe, root, &[root], false, task).await { - 704
Ok(content) => content, - 705
Err(reason) => return failed(reason), - 706
}; - 707
match serde_json::from_str::<Vec<vak_sandbox::TargetCheckResult>>(&content) { - 708
Ok(results) if results.len() == checks.len() => results, - 709
Ok(_) => failed("worker answered a different number of checks".into()), - 710
Err(error) => failed(format!("worker returned invalid results: {error}")), - 711
} - 712
} - 713
- 714
/// What an Office draft (`after`) changes compared with the file it would - 715
/// replace (`before`, absent for a new file), computed in a worker under the - 716
/// same read-only sandbox and deadline as verification, because every file - 717
/// involved is hostile input (invariant 39). With the draft's `lineage`, the - 718
/// answer also carries the choices a person can keep or leave out - 719
/// (`vak_ooxml::review`), or the reason the draft can only be taken whole. - 720
pub async fn office_review( - 721
worker_exe: &Path, - 722
before: Option<&Path>, - 723
after: &Path, - 724
lineage: Option<&OfficeLineage>, - 725
) -> Result<Value, String> { - 726
let Some(after_dir) = after.parent() else { - 727
return Err("the draft has no directory".into()); - 728
}; - 729
let mut roots = vec![after_dir]; - 730
roots.extend(before.and_then(Path::parent)); - 731
roots.extend(lineage.and_then(|lineage| lineage.origin.directory())); - 732
let task = WorkerTask::OfficeReview { - 733
before: before.map(Path::to_path_buf), - 734
after: after.to_path_buf(), - 735
lineage: lineage.cloned(), - 736
}; - 737
let content = run_task(worker_exe, after_dir, &roots, false, task).await?; - 738
serde_json::from_str(&content) - 739
.map_err(|error| format!("worker returned an invalid review: {error}")) - 740
} - 741
- 742
/// Replays the choices in `keep` from `lineage` and writes the narrower - 743
/// version of `draft` to `out`, refusing a draft the lineage does not - 744
/// reproduce. The worker can write only `out`'s directory, which must - 745
/// exist, and read only the lineage's source and the draft. - 746
pub async fn office_narrow( - 747
worker_exe: &Path, - 748
lineage: &OfficeLineage, - 749
draft: &Path, - 750
keep: &[String], - 751
out: &Path, - 752
) -> Result<Value, String> { - 753
let Some(out_dir) = out.parent() else { - 754
return Err("the output has no directory".into()); - 755
}; - 756
let mut roots = vec![out_dir]; - 757
roots.extend(lineage.origin.directory()); - 758
roots.extend(draft.parent()); - 759
let task = WorkerTask::OfficeNarrow { - 760
lineage: lineage.clone(), - 761
draft: draft.to_path_buf(), - 762
keep: keep.to_vec(), - 763
out: out.to_path_buf(), - 764
}; - 765
let content = run_task(worker_exe, out_dir, &roots, true, task).await?; - 766
serde_json::from_str(&content) - 767
.map_err(|error| format!("worker returned an invalid answer: {error}")) - 768
} - 769
- 770
/// Applies `lineage`'s ops to its source and writes the result to `out`, a - 771
/// new file (`vak office apply`, docs/design/72 P5). The same checked apply - 772
/// as the `office_apply` tool; the worker can write only `out`'s directory - 773
/// and read only the source's. The answer lists each op's result and - 774
/// postcondition, the engine's notices, the semantic change list against - 775
/// the source, and what the edit does to signatures and labels. - 776
pub async fn office_apply_to( - 777
worker_exe: &Path, - 778
lineage: &OfficeLineage, - 779
out: &Path, - 780
) -> Result<Value, String> { - 781
let Some(out_dir) = out.parent() else { - 782
return Err("the output has no directory".into()); - 783
}; - 784
let mut roots = vec![out_dir]; - 785
roots.extend(lineage.origin.directory()); - 786
let task = WorkerTask::OfficeApply { - 787
lineage: lineage.clone(), - 788
out: out.to_path_buf(), - 789
}; - 790
let content = run_task(worker_exe, out_dir, &roots, true, task).await?; - 791
serde_json::from_str(&content) - 792
.map_err(|error| format!("worker returned an invalid answer: {error}")) - 793
} - 794
- 795
fn office_apply_in_worker(lineage: &OfficeLineage, out: &Path) -> Result<String, String> { - 796
if std::fs::symlink_metadata(out).is_ok() { - 797
return Err(format!( - 798
"{} already exists; name a new file for the result", - 799
out.display() - 800
)); - 801
} - 802
let target = out - 803
.extension() - 804
.and_then(|extension| extension.to_str()) - 805
.and_then(vak_ooxml::Format::from_extension) - 806
.ok_or_else(|| { - 807
format!( - 808
"{} is not named as a Word, Excel or PowerPoint file", - 809
out.display() - 810
) - 811
})?; - 812
let context = lineage_context(lineage); - 813
let (source, applied) = - 814
crate::office_apply::apply_checked(&lineage.origin, &lineage.ops, &context, target)?; - 815
let before = match &lineage.origin { - 816
OfficeOrigin::File { path, .. } => Some( - 817
vak_ooxml::read::read(std::io::Cursor::new(source), vak_ooxml::Limits::default()) - 818
.map_err(|error| format!("{} could not be read: {error}", path.display()))?, - 819
), - 820
OfficeOrigin::Blank => None, - 821
}; - 822
crate::office_apply::write_atomically(out, &applied.bytes)?; - 823
serde_json::to_string(&serde_json::json!({ - 824
"path": out, - 825
"sha256": crate::office_apply::sha256_hex(&applied.bytes), - 826
"results": applied.results, - 827
"notices": applied.notices, - 828
"changes": vak_ooxml::diff::diff(before.as_ref(), &applied.document), - 829
"impact": vak_ooxml::diff::impact(before.as_ref(), &applied.document), - 830
})) - 831
.map_err(|error| error.to_string()) - 832
} - 833
- 834
/// What a view draws of the Office file at `path`: a page of its content - 835
/// or its structure, with the file's digest so an anchor chosen in the - 836
/// view is bound to these exact bytes. Parsed in a worker under the - 837
/// read-only sandbox, like every Office read (invariant 39). - 838
pub async fn office_project( - 839
worker_exe: &Path, - 840
path: &Path, - 841
view: OfficeView, - 842
) -> Result<Value, String> { - 843
let Some(dir) = path.parent() else { - 844
return Err("the file has no directory".into()); - 845
}; - 846
let task = WorkerTask::OfficeProject { - 847
path: path.to_path_buf(), - 848
view, - 849
}; - 850
let content = run_task(worker_exe, dir, &[dir], false, task).await?; - 851
serde_json::from_str(&content) - 852
.map_err(|error| format!("worker returned an invalid projection: {error}")) - 853
} - 854
- 855
fn office_project_in_worker(path: &Path, view: OfficeView) -> Result<String, String> { - 856
let limits = vak_ooxml::Limits::default(); - 857
let bytes = crate::office_apply::read_bounded(path, &limits)?; - 858
let sha256 = crate::office_apply::sha256_hex(&bytes); - 859
let mut body = match view { - 860
OfficeView::Content { .. } | OfficeView::At { .. } | OfficeView::Facts => { - 861
let document = vak_ooxml::read::read(std::io::Cursor::new(bytes), limits) - 862
.map_err(|error| format!("{} could not be read: {error}", path.display()))?; - 863
let focus = match &view { - 864
OfficeView::At { anchor } => vak_ooxml::projection::locate(&document, anchor), - 865
_ => None, - 866
}; - 867
let from = match &view { - 868
OfficeView::Content { from } => *from, - 869
OfficeView::At { .. } => focus - 870
.map(|index| { - 871
vak_ooxml::projection::page_start( - 872
&document, - 873
index, - 874
vak_ooxml::projection::PAGE_BYTES, - 875
) - 876
}) - 877
.unwrap_or(0), - 878
_ => document.units.len(), - 879
}; - 880
let mut page = serde_json::to_value(vak_ooxml::projection::project( - 881
&document, - 882
from, - 883
vak_ooxml::projection::PAGE_BYTES, - 884
)); - 885
if let (Ok(page), Some(index)) = (&mut page, focus) { - 886
page["focus"] = Value::String(document.units[index].anchor.clone()); - 887
} - 888
page - 889
} - 890
OfficeView::Structure => { - 891
let structure = - 892
vak_ooxml::projection::structure(std::io::Cursor::new(bytes), limits) - 893
.map_err(|error| format!("{} could not be read: {error}", path.display()))?; - 894
serde_json::to_value(structure) - 895
} - 896
} - 897
.map_err(|error| error.to_string())?; - 898
body["sha256"] = Value::String(sha256); - 899
serde_json::to_string(&body).map_err(|error| error.to_string()) - 900
} - 901
- 902
fn read_document(path: &Path) -> Result<vak_ooxml::read::Document, String> { - 903
let file = std::fs::File::open(path) - 904
.map_err(|error| format!("cannot open {}: {error}", path.display()))?; - 905
vak_ooxml::read::read(file, vak_ooxml::Limits::default()) - 906
.map_err(|error| format!("{} could not be read: {error}", path.display())) - 907
} - 908
- 909
/// The bytes the lineage starts from: its file, refused unless it still has - 910
/// the digest the draft was made against, or the built-in blank for - 911
/// `target`, the format the draft is named as. - 912
fn lineage_source( - 913
lineage: &OfficeLineage, - 914
target: Option<vak_ooxml::Format>, - 915
) -> Result<Vec<u8>, String> { - 916
let (path, base_digest) = match &lineage.origin { - 917
OfficeOrigin::File { path, base_digest } => (path, base_digest), - 918
OfficeOrigin::Blank => { - 919
let Some(target) = target else { - 920
return Err("the draft is not named as a Word, Excel or PowerPoint file".into()); - 921
}; - 922
return vak_ooxml::blank::blank(target); - 923
} - 924
}; - 925
let limits = vak_ooxml::Limits::default(); - 926
let bytes = crate::office_apply::read_bounded(path, &limits)?; - 927
let digest = crate::office_apply::sha256_hex(&bytes); - 928
let expected = base_digest - 929
.trim() - 930
.trim_end_matches('…') - 931
.to_ascii_lowercase(); - 932
if expected.len() < 16 || !digest.starts_with(&expected) { - 933
return Err(format!( - 934
"{} changed after the draft was made from it", - 935
path.display() - 936
)); - 937
} - 938
Ok(bytes) - 939
} - 940
- 941
fn lineage_context(lineage: &OfficeLineage) -> vak_ooxml::edit::EditContext { - 942
vak_ooxml::edit::EditContext { - 943
author: lineage.author.clone(), - 944
date: chrono::Utc::now().format("%Y-%m-%dT%H:%M:%SZ").to_string(), - 945
tracked: !lineage.new_file && lineage.origin != OfficeOrigin::Blank, - 946
} - 947
} - 948
- 949
fn target_of(path: &Path) -> Option<vak_ooxml::Format> { - 950
path.extension() - 951
.and_then(|extension| extension.to_str()) - 952
.and_then(vak_ooxml::Format::from_extension) - 953
} - 954
- 955
fn office_review_in_worker( - 956
before: Option<&Path>, - 957
after: &Path, - 958
lineage: Option<&OfficeLineage>, - 959
) -> Result<String, String> { - 960
let current = before.map(read_document).transpose()?; - 961
let draft = read_document(after)?; - 962
let diff = vak_ooxml::diff::diff(current.as_ref(), &draft); - 963
let mut body = serde_json::json!({ - 964
"summary": diff.summary, - 965
"changes": diff.changes, - 966
"flags": draft.inspection.flags(), - 967
"impact": vak_ooxml::diff::impact(current.as_ref(), &draft), - 968
}); - 969
let choices = match lineage { - 970
None => Err("the draft was not made by office_apply in this conversation".to_string()), - 971
Some(lineage) => lineage_source(lineage, target_of(after)).and_then(|source| { - 972
vak_ooxml::review::choices( - 973
&source, - 974
&lineage.ops, - 975
&lineage_context(lineage), - 976
vak_ooxml::Limits::default(), - 977
target_of(after), - 978
&draft, - 979
) - 980
}), - 981
}; - 982
match choices { - 983
Ok(choices) => body["choices"] = serde_json::json!(choices), - 984
Err(reason) => body["choices_unavailable"] = Value::String(reason), - 985
} - 986
serde_json::to_string(&body).map_err(|error| error.to_string()) - 987
} - 988
- 989
fn office_narrow_in_worker( - 990
lineage: &OfficeLineage, - 991
draft: &Path, - 992
keep: &[String], - 993
out: &Path, - 994
) -> Result<String, String> { - 995
let source = lineage_source(lineage, target_of(out))?; - 996
let draft = read_document(draft)?; - 997
let applied = vak_ooxml::review::narrow( - 998
&source, - 999
&lineage.ops, - 1000
keep, - 1001
&lineage_context(lineage), - 1002
vak_ooxml::Limits::default(), - 1003
target_of(out), - 1004
&draft, - 1005
) - 1006
.map_err(|error| format!("nothing was written: {error}"))?; - 1007
crate::office_apply::write_atomically(out, &applied.bytes)?; - 1008
serde_json::to_string(&serde_json::json!({ - 1009
"results": applied.results, - 1010
"sha256": crate::office_apply::sha256_hex(&applied.bytes), - 1011
})) - 1012
.map_err(|error| error.to_string()) - 1013
} - 1014
- 1015
/// Spawns one worker for `task` under a network-denied sandbox that can - 1016
/// read `roots` and write nothing, or only `roots[0]` when `writes_first`, - 1017
/// and returns its answer's content. Every failure, including the deadline, - 1018
/// is an error; nothing is guessed. - 1019
async fn run_task( - 1020
worker_exe: &Path, - 1021
cwd: &Path, - 1022
roots: &[&Path], - 1023
writes_first: bool, - 1024
task: WorkerTask, - 1025
) -> Result<String, String> { - 1026
if !worker_exe.is_file() { - 1027
return Err(format!( - 1028
"worker executable not found: {}", - 1029
worker_exe.display() - 1030
)); - 1031
} - 1032
let request = WorkerRequest { - 1033
version: PROTOCOL_VERSION, - 1034
task, - 1035
}; - 1036
let payload = - 1037
serde_json::to_vec(&request).map_err(|error| format!("request encode failed: {error}"))?; - 1038
let worker_command = format!( - 1039
"{} {}", - 1040
shell_quote(&worker_exe.display().to_string()), - 1041
WORKER_SUBCOMMAND - 1042
); - 1043
let effective = match verification_sandbox(roots, writes_first, worker_exe) { - 1044
Some(sandbox) => sandbox.wrap(&worker_command), - 1045
None => worker_command, - 1046
}; - 1047
let mut command = tokio::process::Command::new(crate::bash::POSIX_SHELL); - 1048
command - 1049
.arg("-c") - 1050
.arg(effective) - 1051
.current_dir(cwd) - 1052
.stdin(Stdio::piped()) - 1053
.stdout(Stdio::piped()) - 1054
.stderr(Stdio::piped()) - 1055
.kill_on_drop(true); - 1056
crate::bash::scrub_environment(&mut command); - 1057
command.env(WORKER_ENV, "1"); - 1058
crate::bash::isolate_process_group(&mut command); - 1059
let mut child = command - 1060
.spawn() - 1061
.map_err(|error| format!("worker spawn failed: {error}"))?; - 1062
let pid = child.id(); - 1063
let Some(mut stdin) = child.stdin.take() else { - 1064
crate::bash::kill_process_group(&pid); - 1065
return Err("worker has no stdin".into()); - 1066
}; - 1067
if let Err(error) = stdin.write_all(&payload).await { - 1068
crate::bash::kill_process_group(&pid); - 1069
return Err(format!("worker request failed: {error}")); - 1070
} - 1071
drop(stdin); - 1072
let (Some(stdout), Some(stderr)) = (child.stdout.take(), child.stderr.take()) else { - 1073
crate::bash::kill_process_group(&pid); - 1074
return Err("worker has no output pipes".into()); - 1075
}; - 1076
let stdout_reader = tokio::spawn(async move { - 1077
let mut bytes = Vec::new(); - 1078
let _ = stdout - 1079
.take(MAX_PROTOCOL_BYTES + 1) - 1080
.read_to_end(&mut bytes) - 1081
.await; - 1082
bytes - 1083
}); - 1084
let stderr_reader = tokio::spawn(async move { - 1085
let mut bytes = Vec::new(); - 1086
let _ = stderr.take(64 * 1024).read_to_end(&mut bytes).await; - 1087
bytes - 1088
}); - 1089
let status = match tokio::time::timeout(VERIFY_DEADLINE, child.wait()).await { - 1090
Ok(Ok(status)) => status, - 1091
Ok(Err(error)) => return Err(format!("worker wait failed: {error}")), - 1092
Err(_) => { - 1093
crate::bash::kill_process_group(&pid); - 1094
let _ = child.wait().await; - 1095
return Err(format!( - 1096
"worker exceeded the {}s deadline", - 1097
VERIFY_DEADLINE.as_secs() - 1098
)); - 1099
} - 1100
}; - 1101
let stdout = stdout_reader.await.unwrap_or_default(); - 1102
let stderr = stderr_reader.await.unwrap_or_default(); - 1103
if !status.success() { - 1104
return Err(format!( - 1105
"worker exited with {}: {}", - 1106
status.code().unwrap_or(-1), - 1107
String::from_utf8_lossy(&stderr).trim() - 1108
)); - 1109
} - 1110
if stdout.len() as u64 > MAX_PROTOCOL_BYTES { - 1111
return Err("worker output exceeded the protocol limit".into()); - 1112
} - 1113
let response: WorkerResponse = serde_json::from_slice(&stdout) - 1114
.map_err(|error| format!("worker returned invalid protocol: {error}"))?; - 1115
if response.version != PROTOCOL_VERSION || response.is_error { - 1116
return Err(format!("worker refused the request: {}", response.content)); - 1117
} - 1118
Ok(response.content) - 1119
} - 1120
- 1121
/// Network-denied, able to read `roots` and the worker executable, and to - 1122
/// write `roots[0]` only when `writes_first`. `None` where no OS backend - 1123
/// exists; the in-code bounds of each reader then stand alone. - 1124
fn verification_sandbox( - 1125
roots: &[&Path], - 1126
writes_first: bool, - 1127
worker_exe: &Path, - 1128
) -> Option<Arc<dyn crate::sandbox::Sandbox>> { - 1129
let (first, rest) = roots.split_first()?; - 1130
let mut extra: Vec<PathBuf> = rest - 1131
.iter() - 1132
.filter_map(|root| root.canonicalize().ok()) - 1133
.collect(); - 1134
extra.extend(worker_exe.parent().map(Path::to_path_buf)); - 1135
let mode = if writes_first { - 1136
crate::sandbox::SandboxMode::WorkspaceWrite - 1137
} else { - 1138
crate::sandbox::SandboxMode::ReadOnly - 1139
}; - 1140
#[cfg(target_os = "macos")] - 1141
{ - 1142
let mut sandbox = crate::sandbox::Seatbelt::task_copy(mode, first); - 1143
sandbox.read_paths.extend(extra); - 1144
Some(Arc::new(sandbox)) - 1145
} - 1146
#[cfg(target_os = "linux")] - 1147
{ - 1148
let mut sandbox = crate::landlock::Landlock::task_copy(mode, first); - 1149
sandbox.read_paths.extend(extra); - 1150
Some(Arc::new(sandbox)) - 1151
} - 1152
#[cfg(not(any(target_os = "macos", target_os = "linux")))] - 1153
{ - 1154
let _ = (first, extra, mode); - 1155
None - 1156
} - 1157
} - 1158
- 1159
fn shell_quote(value: &str) -> String { - 1160
let mut out = String::with_capacity(value.len() + 2); - 1161
out.push('\''); - 1162
for character in value.chars() { - 1163
if character == '\'' { - 1164
out.push_str("'\\''"); - 1165
} else { - 1166
out.push(character); - 1167
} - 1168
} - 1169
out.push('\''); - 1170
out - 1171
} - 1172
- 1173
/// Decides how a sandbox wrapper applies to a brokered tool invocation. - 1174
/// - 1175
/// Two strategies: - 1176
/// - **`WorkerProcess` target** (Seatbelt/Landlock): the entire worker - 1177
/// process is wrapped, so the sandbox binary itself is sandboxed at exec. - 1178
/// - **`ToolCommand` target** (Docker): the worker runs on the host and only - 1179
/// the tool's own command is wrapped, so the model-controlled shell lands - 1180
/// inside the container. - 1181
/// - 1182
/// Returns the effective shell command string and mutates `request_args` - 1183
/// in place when the command-level wrapping path is taken. - 1184
fn resolve_worker_command( - 1185
sandbox: Option<&dyn crate::sandbox::Sandbox>, - 1186
tool: &str, - 1187
request_args: &mut serde_json::Value, - 1188
worker_command: &str, - 1189
) -> String { - 1190
match sandbox { - 1191
Some(sb) if sb.target() == SandboxTarget::WorkerProcess => sb.wrap(worker_command), - 1192
Some(sb) => { - 1193
if tool == "bash" - 1194
&& let Some(command) = request_args - 1195
.get("command") - 1196
.and_then(serde_json::Value::as_str) - 1197
{ - 1198
request_args["command"] = serde_json::Value::String(sb.wrap(command)); - 1199
} - 1200
worker_command.to_string() - 1201
} - 1202
None => worker_command.to_string(), - 1203
} - 1204
} - 1205
- 1206
#[cfg(test)] - 1207
mod tests { - 1208
#![allow(clippy::unwrap_used, clippy::expect_used)] - 1209
- 1210
use super::*; - 1211
use crate::sandbox::{Sandbox, SandboxMode, Seatbelt}; - 1212
use serde_json::json; - 1213
- 1214
#[test] - 1215
fn no_sandbox_passes_worker_command_through() { - 1216
let worker_cmd = "/path/to/vak __tool_worker"; - 1217
let mut args = json!({"command": "echo hi"}); - 1218
let effective = resolve_worker_command(None, "bash", &mut args, worker_cmd); - 1219
assert_eq!(effective, worker_cmd); - 1220
// args untouched - 1221
assert_eq!(args["command"], json!("echo hi")); - 1222
} - 1223
- 1224
#[test] - 1225
fn worker_process_target_wraps_the_worker_executable() { - 1226
let dir = std::env::temp_dir(); - 1227
let sb: Arc<dyn Sandbox> = Arc::new(Seatbelt::new(SandboxMode::ReadOnly, &dir)); - 1228
assert_eq!(sb.target(), SandboxTarget::WorkerProcess); - 1229
let worker_cmd = "/vak __tool_worker"; - 1230
let mut args = json!({"command": "echo hi"}); - 1231
let effective = resolve_worker_command(Some(&*sb), "bash", &mut args, worker_cmd); - 1232
// The worker executable itself — not the inner bash command — is - 1233
// wrapped, so the whole broker runs under sandbox-exec. - 1234
assert!( - 1235
effective.starts_with("sandbox-exec -p "), - 1236
"expected sandbox-exec wrapper, got: {effective}" - 1237
); - 1238
assert!(effective.contains("__tool_worker")); - 1239
// The inner bash command must NOT be wrapped — it travels through - 1240
// the worker protocol, not the shell wrapper. - 1241
assert!(!effective.contains("echo hi")); - 1242
} - 1243
- 1244
#[test] - 1245
fn tool_command_target_wraps_bash_command_in_args() { - 1246
// A Docker-style sandbox (ToolCommand target) must wrap the bash - 1247
// command inside the worker request args — the host worker stays - 1248
// a protocol adapter and only the model-controlled shell is sandboxed. - 1249
let sb: Arc<dyn Sandbox> = Arc::new(DockerStub { - 1250
target: SandboxTarget::ToolCommand, - 1251
wrap_fn: |cmd| format!("docker-run-wrapped({})", cmd), - 1252
}); - 1253
assert_eq!(sb.target(), SandboxTarget::ToolCommand); - 1254
- 1255
let worker_cmd = "/vak __tool_worker"; - 1256
- 1257
// bash tool: command gets wrapped in the args - 1258
let mut args = json!({"command": "echo from-bash"}); - 1259
let effective = resolve_worker_command(Some(&*sb), "bash", &mut args, worker_cmd); - 1260
assert_eq!( - 1261
effective, worker_cmd, - 1262
"worker command passes through unwrapped" - 1263
); - 1264
assert_eq!( - 1265
args["command"], - 1266
json!("docker-run-wrapped(echo from-bash)"), - 1267
"bash command must be wrapped in the request args" - 1268
); - 1269
- 1270
// non-bash tool: args untouched, worker command passes through - 1271
let mut args2 = json!({"path": "src/main.rs"}); - 1272
let effective2 = resolve_worker_command(Some(&*sb), "read", &mut args2, worker_cmd); - 1273
assert_eq!(effective2, worker_cmd); - 1274
assert_eq!(args2["path"], json!("src/main.rs")); - 1275
} - 1276
- 1277
#[test] - 1278
fn tool_command_target_does_not_wrap_bash_without_command_arg() { - 1279
let sb: Arc<dyn Sandbox> = Arc::new(DockerStub { - 1280
target: SandboxTarget::ToolCommand, - 1281
wrap_fn: |cmd| format!("wrapped({})", cmd), - 1282
}); - 1283
let worker_cmd = "/vak __tool_worker"; - 1284
let mut args = json!({}); - 1285
let effective = resolve_worker_command(Some(&*sb), "bash", &mut args, worker_cmd); - 1286
assert_eq!(effective, worker_cmd); - 1287
assert_eq!( - 1288
args["command"], - 1289
json!(Value::Null), - 1290
"no command key to wrap" - 1291
); - 1292
} - 1293
- 1294
/// Minimal Sandbox stub that lets us test ToolCommand-target wrapping - 1295
/// without a Docker daemon. - 1296
struct DockerStub { - 1297
target: SandboxTarget, - 1298
wrap_fn: fn(&str) -> String, - 1299
} - 1300
- 1301
impl Sandbox for DockerStub { - 1302
fn name(&self) -> &str { - 1303
"docker-stub" - 1304
} - 1305
fn wrap(&self, command: &str) -> String { - 1306
(self.wrap_fn)(command) - 1307
} - 1308
fn target(&self) -> SandboxTarget { - 1309
self.target - 1310
} - 1311
} - 1312
} - 1313
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.