- 1001
// the live flow file BEFORE touching anything. Drift fails closed - 1002
// unless explicitly accepted; the frozen snapshot always wins. - 1003
if resume { - 1004
let snapshot_body = std::fs::read_to_string(&state_path).unwrap_or_default(); - 1005
match vak_flow::adopt::recovery_audit(&snapshot_body, Some(&toml_str)) { - 1006
Ok((snapshot, action)) => match action { - 1007
"resume" => println!("[recovery-audit] snapshot={snapshot} action=resume"), - 1008
"repair" if !accept_drift => { - 1009
eprintln!( - 1010
"[recovery-audit] snapshot={snapshot} — live flow file drifted from the frozen definition\n \ - 1011
resume executes the FROZEN copy; pass --accept-drift to acknowledge." - 1012
); - 1013
return 2; - 1014
} - 1015
_ => println!("[recovery-audit] snapshot={snapshot} action={action}"), - 1016
}, - 1017
Err(e) => { - 1018
eprintln!("error: cannot read run ledger for audit: {e}"); - 1019
return 2; - 1020
} - 1021
} - 1022
} - 1023
- 1024
let mut state = load_state(&state_path, &name, &toml_str); - 1025
- 1026
let prepared = core.prepare_turn().await; - 1027
let mut outcome = - 1028
vak_intent::OutcomeSpec::from_reading(name.clone(), &vak_intent::Reading::default(), 0); - 1029
outcome.max_turns = Some(core.effective_max_turns()); - 1030
let deps = vak_flow::ExecutorDeps { - 1031
prompt_layers: Vec::new(), - 1032
provider, - 1033
system_prompt: prepared.system_prompt, - 1034
model: core.effective_model(), - 1035
tools: prepared.tools, - 1036
read_only_tools: prepared.read_only_tools, - 1037
max_turns: core.effective_max_turns(), - 1038
outcome: Some(outcome), - 1039
max_retries: 0, - 1040
retry_base_backoff_ms: 100, - 1041
request_timeout: Some(std::time::Duration::from_secs(600)), - 1042
circuit_breaker: None, - 1043
run_retry_attempts: 0, - 1044
run_retry_base_backoff_ms: 1000, - 1045
dispatch_ceiling: 1, - 1046
spend_gate: None, - 1047
permission: Some(std::sync::Arc::new(engine)), - 1048
mode: match core.effective_permission_mode() { - 1049
vak_config::PermissionMode::ReadOnly => vak_permission::Mode::ReadOnly, - 1050
vak_config::PermissionMode::WorkspaceWrite => vak_permission::Mode::WorkspaceWrite, - 1051
vak_config::PermissionMode::FullAccess => vak_permission::Mode::FullAccess, - 1052
}, - 1053
approval_mode: match core.effective_approval_mode() { - 1054
vak_config::ApprovalMode::Ask => vak_agent::ApprovalMode::Ask, - 1055
vak_config::ApprovalMode::ApproveSafe => vak_agent::ApprovalMode::ApproveSafe, - 1056
vak_config::ApprovalMode::AutoApprove => vak_agent::ApprovalMode::AutoApprove, - 1057
}, - 1058
approver, - 1059
sandbox: core.agent_sandbox(), - 1060
cwd: core.cwd().clone(), - 1061
sessions_home: core.sessions_home().clone(), - 1062
parent_session_id, - 1063
state_path: state_path.clone(), - 1064
agent_identity: core.agent_identity().cloned(), - 1065
conversation_context: core.conversation_context().cloned(), - 1066
work: None, - 1067
}; - 1068
let executor = vak_flow::Executor::new(deps); - 1069
- 1070
let cancel = CancellationToken::new(); - 1071
{ - 1072
let cancel = cancel.clone(); - 1073
tokio::spawn(async move { - 1074
if tokio::signal::ctrl_c().await.is_ok() { - 1075
eprintln!("\n[cancelling…]"); - 1076
cancel.cancel(); - 1077
} - 1078
}); - 1079
} - 1080
- 1081
let (tx, mut rx) = tokio::sync::mpsc::channel::<String>(256); - 1082
// Layer-aware progress strip (docs/design/10-flows.md): prefix ✓/✗/⊘ events - 1083
// with their topological layer, other lines verbatim. - 1084
let layer_of: std::collections::HashMap<String, usize> = match vak_flow::parse::layers(&flow) { - 1085
Ok(layers) => layers - 1086
.iter() - 1087
.enumerate() - 1088
.flat_map(|(li, l)| l.iter().map(move |id| (id.clone(), li + 1))) - 1089
.collect(), - 1090
Err(_) => Default::default(), - 1091
}; - 1092
let total_layers = layer_of.values().copied().max().unwrap_or(0); - 1093
- 1094
let runner = tokio::spawn(async move { executor.run(&flow, &mut state, cancel, tx).await }); - 1095
while let Some(line) = rx.recv().await { - 1096
let marker = line - 1097
.strip_prefix('✓') - 1098
.or_else(|| line.strip_prefix('✗')) - 1099
.or_else(|| line.strip_prefix('⊘')); - 1100
if let Some(rest) = marker { - 1101
let id = rest.trim(); - 1102
let layer = layer_of.get(id).copied().unwrap_or(0); - 1103
eprintln!("[L{}/{}] {}", layer, total_layers, line); - 1104
} else { - 1105
eprintln!("{line}"); - 1106
} - 1107
} - 1108
- 1109
match runner.await { - 1110
Ok(outcome) => match outcome { - 1111
vak_flow::FlowOutcome::Completed { outputs } => { - 1112
// Snapshot from the persisted ledger (state moved into the runner). - 1113
if let Ok(body) = std::fs::read_to_string(&state_path) - 1114
&& let Ok(st) = serde_json::from_str::<vak_flow::FlowState>(&body) - 1115
{ - 1116
let snap = vak_flow::graph::graph_snapshot(&st); - 1117
eprintln!( - 1118
"── snapshot: {} completed / {} failed / {} skipped · {} layer(s)", - 1119
snap.completed, snap.failed, snap.skipped, snap.layers_total - 1120
); - 1121
} - 1122
eprintln!("── flow completed · state {}", state_path.display()); - 1123
for (id, out) in outputs { - 1124
println!("[{id}]\n{out}\n"); - 1125
} - 1126
0 - 1127
} - 1128
vak_flow::FlowOutcome::Failed { node, reason, .. } => { - 1129
eprintln!("── flow failed at '{node}': {reason}"); - 1130
eprintln!(" resume with: vak flow run {name} --resume"); - 1131
1 - 1132
} - 1133
vak_flow::FlowOutcome::Aborted => { - 1134
eprintln!("── flow aborted · resume with: vak flow run {name} --resume"); - 1135
1 - 1136
} - 1137
}, - 1138
Err(e) => { - 1139
eprintln!("error: flow runner crashed: {e}"); - 1140
2 - 1141
} - 1142
} - 1143
} - 1144
- 1145
fn load_state(state_path: &PathBuf, flow_name: &str, definition_toml: &str) -> vak_flow::FlowState { - 1146
if let Ok(text) = std::fs::read_to_string(state_path) - 1147
&& let Ok(state) = serde_json::from_str::<vak_flow::FlowState>(&text) - 1148
{ - 1149
return state; - 1150
} - 1151
vak_flow::FlowState { - 1152
run_id: state_path - 1153
.file_stem() - 1154
.map(|s| s.to_string_lossy().into_owned()) - 1155
.unwrap_or_default(), - 1156
flow_name: flow_name.to_string(), - 1157
definition_toml: definition_toml.to_string(), - 1158
started_at: chrono::Utc::now(), - 1159
outcome: None, - 1160
nodes: Default::default(), - 1161
} - 1162
} - 1163
- 1164
#[allow(clippy::too_many_arguments)] - 1165
async fn run_exec( - 1166
cwd: PathBuf, - 1167
prompt: String, - 1168
agent: Option<String>, - 1169
model: Option<String>, - 1170
provider: Option<String>, - 1171
max_turns: usize, - 1172
_json: bool, - 1173
yes: bool, - 1174
permission_mode: Option<String>, - 1175
write_paths: Vec<PathBuf>, - 1176
worktree: bool, - 1177
resume_session: Option<String>, - 1178
accept_drift: bool, - 1179
managed: bool, - 1180
goal: Option<String>, - 1181
criteria: Vec<String>, - 1182
trusted: bool, - 1183
) -> i32 { - 1184
let mut effective_cwd = cwd.clone(); - 1185
let mut created_worktree: Option<vak_core::worktree::Worktree> = None; - 1186
if worktree { - 1187
match vak_core::worktree::create(&cwd, &format!("exec-{}", timestamp_id())) { - 1188
Ok(wt) => { - 1189
eprintln!("▸ isolated worktree: {} ({})", wt.path.display(), wt.branch); - 1190
effective_cwd = wt.path.clone(); - 1191
created_worktree = Some(wt); - 1192
} - 1193
Err(e) => { - 1194
eprintln!("error: worktree isolation failed: {e}"); - 1195
return 2; - 1196
} - 1197
} - 1198
} - 1199
let mut core = match Core::new_with_trust(effective_cwd.clone(), trusted) - 1200
.map(|c| with_cli_surface(c, yes)) - 1201
{ - 1202
Ok(c) => c, - 1203
Err(e) => { - 1204
eprintln!("error: {e}"); - 1205
return 2; - 1206
} - 1207
}; - 1208
if let Some(agent_id) = &agent - 1209
&& agent_id != "vak" - 1210
{ - 1211
let profiles = match vak_server::agents::effective(&core) { - 1212
Ok(p) => p, - 1213
Err(e) => { - 1214
eprintln!("error: failed to load agents: {e}"); - 1215
return 2; - 1216
} - 1217
}; - 1218
let Some(profile) = profiles - 1219
.into_iter() - 1220
.find(|p| p.id == *agent_id || p.name.eq_ignore_ascii_case(agent_id)) - 1221
else { - 1222
eprintln!( - 1223
"error: agent '{agent_id}' not found. Run 'vak agents list' to view configured agents." - 1224
); - 1225
return 2; - 1226
}; - 1227
if !profile.is_admissible() { - 1228
eprintln!( - 1229
"error: agent '{agent_id}' is {:?} and cannot execute runs", - 1230
profile.lifecycle - 1231
); - 1232
return 2; - 1233
} - 1234
eprintln!("▸ agent: {} ({})", profile.name, profile.id); - 1235
core = core.with_agent_identity(Some(profile.identity())); - 1236
} - 1237
print_config_warnings(&core); - 1238
update_check::maybe_check_update(core.config()); - 1239
if provider.is_some() || model.is_some() { - 1240
let route = core.effective_route(); - 1241
core.set_route( - 1242
provider.unwrap_or(route.provider), - 1243
model.unwrap_or(route.model), - 1244
); - 1245
} - 1246
core.set_max_turns(max_turns); - 1247
if let Some(pm) = permission_mode { - 1248
match vak_config::PermissionMode::deserialize_str(&pm) { - 1249
Some(m) => core.set_permission_mode(m), - 1250
None => { - 1251
eprintln!( - 1252
"error: unknown --permission-mode '{pm}' (read-only | workspace-write | full-access)" - 1253
); - 1254
return 2; - 1255
} - 1256
} - 1257
} - 1258
- 1259
let session = match resume_session { - 1260
Some(sid) => match core.open_session(&sid).await { - 1261
Ok(s) => { - 1262
// A named session resumes its FROZEN prompt. Unlike a chat - 1263
// binding — which the operator never named, so the gateway - 1264
// rotates it — the user asked for this session by id, so - 1265
// silently running a different prompt would be the wrong - 1266
// surprise. Fail closed and make them acknowledge, the same - 1267
// shape `flow run --resume` already uses for a drifted - 1268
// definition (docs/design/45-prompt-layers.md). - 1269
if let Some(drift) = s.header().and_then(|h| core.prompt_drift(&h.contract)) - 1270
&& !accept_drift - 1271
{ - 1272
eprintln!("error: prompt layers changed since session '{sid}' was created:"); - 1273
for line in drift.lines() { - 1274
eprintln!(" {line}"); - 1275
} - 1276
eprintln!( - 1277
" resume executes the FROZEN prompt; pass --accept-drift to \ - 1278
acknowledge, or start a new session to pick up the change." - 1279
); - 1280
return 2; - 1281
} - 1282
if let Some(agent_id) = &agent - 1283
&& let Some(h_agent) = s.header().and_then(|h| h.agent.as_ref()) - 1284
&& h_agent.id != *agent_id - 1285
&& *agent_id != "vak" - 1286
{ - 1287
eprintln!( - 1288
"error: session '{sid}' was created for agent '{}', but you requested agent '{agent_id}'", - 1289
h_agent.id - 1290
); - 1291
return 2; - 1292
} - 1293
eprintln!("▸ resuming session {sid}"); - 1294
s - 1295
} - 1296
Err(e) => { - 1297
eprintln!("error: cannot open session '{sid}': {e}"); - 1298
return 2; - 1299
} - 1300
}, - 1301
None => match core.start_session().await { - 1302
Ok(s) => s, - 1303
Err(e) => { - 1304
eprintln!("error: {e}"); - 1305
return 2; - 1306
} - 1307
}, - 1308
}; - 1309
let session_path = session.path().to_path_buf(); - 1310
- 1311
let (tx, mut rx) = tokio::sync::mpsc::channel::<AgentEvent>(1024); - 1312
let cancel = CancellationToken::new(); - 1313
let cancel_for_signal = cancel.clone(); - 1314
tokio::spawn(async move { - 1315
if tokio::signal::ctrl_c().await.is_ok() { - 1316
eprintln!("\n[cancelling…]"); - 1317
cancel_for_signal.cancel(); - 1318
} - 1319
}); - 1320
- 1321
let approver: Option<std::sync::Arc<dyn vak_agent::Approver>> = Some(if yes { - 1322
std::sync::Arc::new(vak_agent::AutoApprove) - 1323
} else { - 1324
std::sync::Arc::new(vak_agent::AutoDeny) - 1325
}); - 1326
let permission = if write_paths.is_empty() { - 1327
None - 1328
} else { - 1329
match core.build_permission_engine(&core.extra_allow_snapshot()) { - 1330
Ok(engine) => Some(std::sync::Arc::new( - 1331
engine.restrict_write_paths(core.cwd(), &write_paths), - 1332
)), - 1333
Err(error) => { - 1334
eprintln!("error: {error}"); - 1335
return 2; - 1336
} - 1337
} - 1338
}; - 1339
- 1340
// The runner consumes `core`; keep a handle for the post-turn - 1341
// reflection seam. - 1342
let reflection_core = core.clone(); - 1343
let runner = tokio::spawn(async move { - 1344
let fut = async { - 1345
if let Some(objective) = goal.as_deref() { - 1346
core.run_goal_turn_with( - 1347
session, - 1348
&prompt, - 1349
objective, - 1350
criteria.clone(), - 1351
cancel, - 1352
approver, - 1353
permission, - 1354
None, - 1355
tx, - 1356
) - 1357
.await - 1358
} else if managed { - 1359
core.run_managed_turn_with(session, &prompt, cancel, approver, permission, None, tx) - 1360
.await - 1361
} else { - 1362
core.run_turn_with(session, &prompt, cancel, approver, permission, None, tx) - 1363
.await - 1364
} - 1365
}; - 1366
fut.await - 1367
}); - 1368
- 1369
let mut total_in: u64 = 0; - 1370
let mut total_out: u64 = 0; - 1371
let mut tool_args: std::collections::HashMap<String, (String, String)> = - 1372
std::collections::HashMap::new(); - 1373
let stdout = std::io::stdout(); - 1374
let mut out = stdout.lock(); - 1375
- 1376
while let Some(ev) = rx.recv().await { - 1377
match ev { - 1378
AgentEvent::TurnStart { turn } => { - 1379
if turn > 0 { - 1380
writeln!(out).ok(); - 1381
} - 1382
} - 1383
AgentEvent::Stream(StreamEvent::TextDelta { delta, .. }) => { - 1384
write!(out, "{delta}").ok(); - 1385
out.flush().ok(); - 1386
} - 1387
AgentEvent::Stream(StreamEvent::ThinkingDelta { .. }) => {} - 1388
AgentEvent::ToolCallStart { - 1389
id, - 1390
name, - 1391
args_json, - 1392
} => { - 1393
tool_args.insert(id, (name.clone(), args_json.clone())); - 1394
eprintln!("▸ {} {}", name, format::summarize_args(&name, &args_json)); - 1395
} - 1396
AgentEvent::ToolCallEnd { - 1397
id, - 1398
name, - 1399
is_error, - 1400
result_preview, - 1401
} => { - 1402
eprintln!("{} {name}", if is_error { "✗" } else { "✓" }); - 1403
let stored = tool_args.remove(&id); - 1404
if name == "edit" - 1405
&& !is_error - 1406
&& let Some((_, args)) = &stored - 1407
&& let Some(diff) = format::edit_diff_text(args, 10) - 1408
{ - 1409
for line in format::strip_ansi(&diff).lines() { - 1410
eprintln!(" {line}"); - 1411
} - 1412
} else if is_error - 1413
&& let Some(prev) = result_preview - 1414
&& !prev.trim().is_empty() - 1415
{ - 1416
let tail: Vec<&str> = prev - 1417
.lines() - 1418
.filter(|l| !l.trim().is_empty()) - 1419
.rev() - 1420
.take(6) - 1421
.collect(); - 1422
for line in tail.iter().rev() { - 1423
eprintln!(" │ {line}"); - 1424
} - 1425
} - 1426
} - 1427
AgentEvent::RetryScheduled { - 1428
attempt, - 1429
delay_ms, - 1430
reason, - 1431
} => { - 1432
eprintln!("⟳ [{attempt}] backing off {delay_ms}ms — {reason}"); - 1433
} - 1434
AgentEvent::RouteFallback { - 1435
to_provider, - 1436
to_model, - 1437
} => { - 1438
eprintln!("⤵ route fallback → {to_provider}/{to_model} (frozen ladder leg)"); - 1439
} - 1440
AgentEvent::ContextCompacting { estimated_tokens } => { - 1441
eprintln!("◌ compacting context (~{estimated_tokens} tokens)"); - 1442
} - 1443
AgentEvent::ContextCompacted { - 1444
before_tokens, - 1445
after_tokens, - 1446
.. - 1447
} => { - 1448
eprintln!("◌ compacted ~{before_tokens} → ~{after_tokens} tokens"); - 1449
} - 1450
AgentEvent::TurnEnd { usage } => { - 1451
total_in += usage.prompt_tokens(); - 1452
total_out += usage.output_tokens; - 1453
} - 1454
AgentEvent::Sandbox(sb_ev) => { - 1455
use vak_tools::sandbox_events::{SandboxEvent, fold_carriage_returns}; - 1456
match sb_ev { - 1457
SandboxEvent::ExecutionStarted { - 1458
tool, - 1459
code_preview, - 1460
scratch_dir, - 1461
.. - 1462
} => { - 1463
let first_line = code_preview.lines().next().unwrap_or("").trim(); - 1464
let display_cmd = if first_line.len() > 60 { - 1465
format!("{}…", &first_line[..60]) - 1466
} else { - 1467
first_line.to_string() - 1468
}; - 1469
eprintln!(" ┌─ [sandbox:{tool}] {display_cmd}"); - 1470
if !scratch_dir.is_empty() { - 1471
eprintln!(" │ scratch: {scratch_dir}"); - 1472
} - 1473
} - 1474
SandboxEvent::Stdout { chunk, .. } => { - 1475
let folded = fold_carriage_returns(&chunk); - 1476
for line in folded.lines() { - 1477
if !line.trim().is_empty() { - 1478
eprintln!(" │ {line}"); - 1479
} - 1480
} - 1481
} - 1482
SandboxEvent::Stderr { chunk, .. } => { - 1483
let folded = fold_carriage_returns(&chunk); - 1484
for line in folded.lines() { - 1485
if !line.trim().is_empty() { - 1486
eprintln!(" │! {line}"); - 1487
} - 1488
} - 1489
} - 1490
SandboxEvent::OutputTruncated { .. } => { - 1491
eprintln!(" │! output truncated after the configured capture limit"); - 1492
} - 1493
SandboxEvent::PackageInstalled { packages, .. } => { - 1494
eprintln!(" │ 📦 packages: {}", packages.join(", ")); - 1495
} - 1496
SandboxEvent::ArtifactGenerated { - 1497
path, - 1498
mime_type, - 1499
size_bytes, - 1500
.. - 1501
} => { - 1502
eprintln!(" │ 📄 artifact: {path} ({size_bytes} B, {mime_type})"); - 1503
} - 1504
SandboxEvent::ProcessTelemetry { - 1505
elapsed_ms, - 1506
memory_bytes, - 1507
.. - 1508
} => { - 1509
if memory_bytes > 0 { - 1510
let mb = memory_bytes as f64 / (1024.0 * 1024.0); - 1511
eprintln!(" │ ⏱ {}ms | RSS: {:.1}MB", elapsed_ms, mb); - 1512
} - 1513
} - 1514
SandboxEvent::ExecutionFinished { - 1515
exit_code, - 1516
duration_ms, - 1517
artifacts, - 1518
.. - 1519
} => { - 1520
let status_sym = if exit_code == 0 { "✓" } else { "✗" }; - 1521
let mut summary = format!( - 1522
" └─ {status_sym} finished in {duration_ms}ms (exit: {exit_code})" - 1523
); - 1524
if !artifacts.is_empty() { - 1525
summary.push_str(&format!(" [{} artifact(s)]", artifacts.len())); - 1526
} - 1527
eprintln!("{summary}"); - 1528
} - 1529
} - 1530
} - 1531
_ => {} - 1532
} - 1533
} - 1534
- 1535
let (outcome, session) = match runner.await { - 1536
Ok(Ok((o, session))) => (o, Some(session)), - 1537
Ok(Err(e)) => { - 1538
eprintln!("error: {e}"); - 1539
return 2; - 1540
} - 1541
Err(e) => { - 1542
eprintln!("error: runner crashed: {e}"); - 1543
return 2; - 1544
} - 1545
}; - 1546
- 1547
writeln!(out).ok(); - 1548
eprintln!( - 1549
"\n── {} · Σ tokens in {total_in} / out {total_out} · session {}", - 1550
match &outcome { - 1551
TurnOutcome::Completed { .. } => "completed", - 1552
TurnOutcome::Aborted { .. } => "aborted", - 1553
TurnOutcome::Failed { .. } => "failed", - 1554
TurnOutcome::MaxTurnsReached => "max turns reached", - 1555
}, - 1556
session_path.display() - 1557
); - 1558
if let TurnOutcome::Failed { error } = outcome { - 1559
eprintln!("error: {error}"); - 1560
if let Some(wt) = created_worktree { - 1561
let _ = vak_core::worktree::remove(&cwd, &wt); - 1562
} - 1563
return 1; - 1564
} - 1565
// Post-turn reflection seam (docs/design/29 P1): the final assistant - 1566
// text is already committed to the ledger, so no extra tail is passed. - 1567
// Bounded so a stuck auxiliary stream cannot hang the exit. - 1568
if let Some(session) = session.as_ref() { - 1569
let pass = tokio::time::timeout( - 1570
REFLECTION_CALL_TIMEOUT, - 1571
reflection_core.reflect_after_turn(session, ""), - 1572
) - 1573
.await; - 1574
if let Ok(outcome) = pass - 1575
&& let Some(line) = exec_reflection_line(&outcome) - 1576
{ - 1577
eprintln!("{line}"); - 1578
} - 1579
} - 1580
if let Some(wt) = created_worktree { - 1581
eprintln!( - 1582
"▸ worktree kept for inspection: {} (branch {}) — remove with git worktree remove", - 1583
wt.path.display(), - 1584
wt.branch - 1585
); - 1586
} - 1587
0 - 1588
} - 1589
- 1590
fn timestamp_id() -> String { - 1591
std::time::SystemTime::now() - 1592
.duration_since(std::time::UNIX_EPOCH) - 1593
.unwrap_or_default() - 1594
.as_nanos() - 1595
.to_string() - 1596
} - 1597
- 1598
/// Upper bound on one background reflection pass so a stuck auxiliary - 1599
/// stream cannot hang process exit. - 1600
const REFLECTION_CALL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(120); - 1601
- 1602
/// Footnote line for a post-turn reflection outcome: one dim line when the - 1603
/// pass persisted something or was skipped for a reason worth acting on - 1604
/// (budget), silent otherwise. - 1605
fn exec_reflection_line(outcome: &vak_core::reflection::ReflectionOutcome) -> Option<String> { - 1606
match outcome { - 1607
vak_core::reflection::ReflectionOutcome::Reflected { notes_added, .. } => { - 1608
Some(format!("· reflected: {notes_added} note(s)")) - 1609
} - 1610
vak_core::reflection::ReflectionOutcome::Skipped { reason: "budget" } => { - 1611
Some("· reflection skipped: budget cap reached".to_string()) - 1612
} - 1613
_ => None, - 1614
} - 1615
} - 1616
- 1617
/// `vak config permissions` — the whole permission answer in one place. - 1618
/// - 1619
/// The CLI could previously only dump the merged config, which prints the - 1620
/// mode but not what the engine actually evaluates. An operator debugging a - 1621
/// refusal had no terminal command that showed the rules, the approval mode, - 1622
/// and which capabilities the composed policy will refuse on this surface. - 1623
fn run_config_permissions(cwd: PathBuf) -> i32 { - 1624
// Trust matters here: an untrusted project's `allow` rules are stripped - 1625
// by the loader, so reading with the wrong trust would print a rule set - 1626
// no run would ever use. - 1627
let trusted = vak_core::trust::is_trusted(&cwd); - 1628
let core = match Core::new_with_trust(cwd, trusted) { - 1629
Ok(core) => core, - 1630
Err(e) => { - 1631
eprintln!("error: {e}"); - 1632
return 2; - 1633
} - 1634
}; - 1635
let (allow, ask, deny) = core.effective_permission_rules(); - 1636
println!("permission_mode = {:?}", core.effective_permission_mode()); - 1637
println!( - 1638
"approval_mode = {}", - 1639
core.effective_approval_mode().as_str() - 1640
); - 1641
println!("sandbox = {}", core.effective_sandbox_name()); - 1642
println!("project_trusted = {trusted}"); - 1643
for (label, list) in [("deny", &deny), ("ask", &ask), ("allow", &allow)] { - 1644
if list.is_empty() { - 1645
println!("{label:<16} = (none)"); - 1646
} else { - 1647
println!("{label:<16} = {}", list.join(", ")); - 1648
} - 1649
} - 1650
let learned = core.extra_allow_snapshot(); - 1651
if !learned.is_empty() { - 1652
println!("learned = {}", learned.join(", ")); - 1653
} - 1654
// The composed answer, not the configured one: this is the same - 1655
// computation the system prompt and the tool registry read. - 1656
let standings = core.capability_standings(); - 1657
if standings.is_empty() { - 1658
println!("\ncapabilities = (none configured)"); - 1659
return 0; - 1660
} - 1661
println!("\ncapabilities"); - 1662
for standing in &standings { - 1663
let mark = match standing.reach { - 1664
vak_core::reach::Reach::Open => "open", - 1665
vak_core::reach::Reach::Gated => "gated", - 1666
vak_core::reach::Reach::Blocked => "BLOCKED", - 1667
}; - 1668
println!(" {mark:<8} {}", standing.label); - 1669
if !standing.reason.is_empty() { - 1670
println!(" {}", standing.reason); - 1671
} - 1672
if !standing.remedy.is_empty() { - 1673
println!(" fix: {}", standing.remedy); - 1674
} - 1675
} - 1676
0 - 1677
} - 1678
- 1679
/// Where a `vak config set-*` write lands. Mirrors the two scopes every - 1680
/// other layered setting already uses. - 1681
fn scope_config_path(scope: cli::PromptScope, cwd: &std::path::Path) -> Option<PathBuf> { - 1682
match scope { - 1683
cli::PromptScope::User => vak_config::global_path(), - 1684
cli::PromptScope::Project => Some(vak_config::project_path(cwd)), - 1685
} - 1686
} - 1687
- 1688
fn run_config_set_mode(cwd: PathBuf, mode: &str, scope: cli::PromptScope) -> i32 { - 1689
let Some(parsed) = vak_config::PermissionMode::deserialize_str(mode) else { - 1690
eprintln!("error: unknown mode '{mode}' (read-only | workspace-write | full-access)"); - 1691
return 2; - 1692
}; - 1693
let Some(path) = scope_config_path(scope, &cwd) else { - 1694
eprintln!("error: user home unavailable"); - 1695
return 2; - 1696
}; - 1697
if let Err(e) = - 1698
vak_config::persist_preferences_to(path, None, None, None, Some(parsed), None, None) - 1699
{ - 1700
eprintln!("error: {e}"); - 1701
return 2; - 1702
} - 1703
println!("permission_mode = {mode} ({} layer)", scope_label(scope)); - 1704
// A running server keeps its own copy; say so rather than implying the - 1705
// change reached every surface already. - 1706
println!("A running `vak serve` picks this up on its next config refresh."); - 1707
0 - 1708
} - 1709
- 1710
fn run_config_set_approval(cwd: PathBuf, mode: &str, scope: cli::PromptScope) -> i32 { - 1711
let Some(parsed) = vak_config::ApprovalMode::parse(mode) else { - 1712
eprintln!("error: unknown mode '{mode}' (ask | approve-safe | auto-approve)"); - 1713
return 2; - 1714
}; - 1715
let Some(path) = scope_config_path(scope, &cwd) else { - 1716
eprintln!("error: user home unavailable"); - 1717
return 2; - 1718
}; - 1719
if let Err(e) = - 1720
vak_config::persist_preferences_to(path, None, None, None, None, Some(parsed), None) - 1721
{ - 1722
eprintln!("error: {e}"); - 1723
return 2; - 1724
} - 1725
println!("approval_mode = {mode} ({} layer)", scope_label(scope)); - 1726
println!("A running `vak serve` picks this up on its next config refresh."); - 1727
0 - 1728
} - 1729
- 1730
fn scope_label(scope: cli::PromptScope) -> &'static str { - 1731
match scope { - 1732
cli::PromptScope::User => "Shared", - 1733
cli::PromptScope::Project => "project", - 1734
} - 1735
} - 1736
- 1737
fn run_config_dump(cwd: PathBuf) { - 1738
match Core::new(cwd.clone()) { - 1739
Ok(core) => { - 1740
println!("# Vakyartha effective config"); - 1741
println!("version = {}", vak_core::APP_VERSION); - 1742
println!("cwd = {}", core.cwd().display()); - 1743
println!("provider = {}", core.effective_provider()); - 1744
println!("model = {}", core.effective_model()); - 1745
println!("max_tokens = {}", core.config().max_tokens); - 1746
println!("max_turns = {}", core.effective_max_turns()); - 1747
println!("permission_mode = {:?}", core.effective_permission_mode()); - 1748
println!( - 1749
"approval_mode = {}", - 1750
core.effective_approval_mode().as_str() - 1751
); - 1752
println!("sandbox = {}", core.effective_sandbox_name()); - 1753
println!("sessions_home = {}", core.sessions_home().display()); - 1754
println!( - 1755
"anthropic_base = {}", - 1756
core.config() - 1757
.anthropic_base_url - 1758
.clone() - 1759
.unwrap_or_else(|| "https://api.anthropic.com".into()) - 1760
); - 1761
println!("tools = {}", core.tool_names().join(", ")); - 1762
let f = &core.config().finops; - 1763
println!( - 1764
"finops = run_cap {} · day_cap {} · overrides {}", - 1765
f.max_run_usd - 1766
.map(|v| format!("${v:.2}")) - 1767
.unwrap_or_else(|| "none".into()), - 1768
f.max_day_usd - 1769
.map(|v| format!("${v:.2}")) - 1770
.unwrap_or_else(|| "none".into()), - 1771
f.price_overrides.len(), - 1772
); - 1773
println!( - 1774
"goal = handoff_reset {} · max_audit_blocks {}", - 1775
core.config().goal.handoff_reset, - 1776
core.config().goal.max_audit_blocks, - 1777
); - 1778
let r = &core.config().route; - 1779
println!( - 1780
"route = objective {} · fallback_models [{}] · max_fallbacks {} · quality_hints [{}]", - 1781
r.objective, - 1782
r.fallback_models.join(", "), - 1783
r.max_fallbacks, - 1784
r.quality_hints.join(", "), - 1785
); - 1786
for w in &core.config().warnings { - 1787
println!("warning = {w}"); - 1788
} - 1789
} - 1790
Err(e) => { - 1791
eprintln!("error: {e}"); - 1792
} - 1793
} - 1794
} - 1795
- 1796
fn run_sessions_list(cwd: PathBuf) { - 1797
let Ok(core) = Core::new(cwd) else { - 1798
return; - 1799
}; - 1800
let mut session_dirs = Vec::new(); - 1801
let direct = vak_session::SessionPath::sessions_dir(&core.sessions_home(), core.cwd()); - 1802
if direct.exists() { - 1803
session_dirs.push(direct.clone()); - 1804
} - 1805
let shared = core.shared_data_home(); - 1806
if let Ok(agents) = std::fs::read_dir(shared.join("agents")) { - 1807
for agent in agents.flatten() { - 1808
let s = vak_session::SessionPath::sessions_dir(&agent.path(), core.cwd()); - 1809
if s.exists() && !session_dirs.contains(&s) { - 1810
session_dirs.push(s); - 1811
} - 1812
} - 1813
} - 1814
let trashed = vak_core::trash::trashed(&shared); - 1815
let mut rows: Vec<(std::time::SystemTime, u64, String)> = Vec::new(); - 1816
for dir in &session_dirs { - 1817
if let Ok(entries) = std::fs::read_dir(dir) { - 1818
for entry in entries.flatten() { - 1819
let name = entry.file_name().to_string_lossy().into_owned(); - 1820
let maybe_row = (name.ends_with(".jsonl") - 1821
&& !trashed.contains(name.trim_end_matches(".jsonl")) - 1822
&& !rows.iter().any(|(_, _, n)| n == &name)) - 1823
.then(|| { - 1824
entry - 1825
.metadata() - 1826
.ok() - 1827
.and_then(|meta| meta.modified().ok().map(|m| (m, meta.len()))) - 1828
}) - 1829
.flatten(); - 1830
if let Some((mtime, len)) = maybe_row { - 1831
rows.push((mtime, len, name)); - 1832
} - 1833
} - 1834
} - 1835
} - 1836
rows.sort_by_key(|(mtime, _, _)| std::cmp::Reverse(*mtime)); - 1837
if rows.is_empty() { - 1838
println!("no sessions yet ({})", direct.display()); - 1839
return; - 1840
} - 1841
for (_mtime, size, name) in rows { - 1842
println!("{name} {size:>10} bytes"); - 1843
} - 1844
} - 1845
- 1846
#[allow(clippy::too_many_arguments)] - 1847
async fn run_plan( - 1848
cwd: PathBuf, - 1849
task: String, - 1850
yes: bool, - 1851
permission_mode: Option<String>, - 1852
write_paths: Vec<PathBuf>, - 1853
worktree: bool, - 1854
trusted: bool, - 1855
) -> i32 { - 1856
let mut effective_cwd = cwd.clone(); - 1857
if worktree { - 1858
match vak_core::worktree::create(&cwd, &format!("plan-{}", timestamp_id())) { - 1859
Ok(wt) => { - 1860
eprintln!("▸ isolated worktree: {} ({})", wt.path.display(), wt.branch); - 1861
effective_cwd = wt.path.clone(); - 1862
} - 1863
Err(e) => { - 1864
eprintln!("error: worktree isolation failed: {e}"); - 1865
return 2; - 1866
} - 1867
} - 1868
} - 1869
let core = match Core::new_with_trust(effective_cwd.clone(), trusted) - 1870
.map(|c| with_cli_surface(c, yes)) - 1871
{ - 1872
Ok(c) => c, - 1873
Err(e) => { - 1874
eprintln!("error: {e}"); - 1875
return 2; - 1876
} - 1877
}; - 1878
// Applied before the session is started, so the frozen prompt describes - 1879
// the boundary this run will actually have. - 1880
if let Some(pm) = permission_mode { - 1881
match vak_config::PermissionMode::deserialize_str(&pm) { - 1882
Some(m) => core.set_permission_mode(m), - 1883
None => { - 1884
eprintln!( - 1885
"error: unknown --permission-mode '{pm}' (read-only | workspace-write | full-access)" - 1886
); - 1887
return 2; - 1888
} - 1889
} - 1890
} - 1891
let provider = match core.provider() { - 1892
Ok(p) => p, - 1893
Err(e) => { - 1894
eprintln!("error: {e}"); - 1895
return 2; - 1896
} - 1897
}; - 1898
let session = match core.start_session().await { - 1899
Ok(s) => s, - 1900
Err(e) => { - 1901
eprintln!("error: {e}"); - 1902
return 2; - 1903
} - 1904
}; - 1905
let parent_session_id = session - 1906
.header() - 1907
.map(|h| h.session_id.clone()) - 1908
.unwrap_or_default(); - 1909
- 1910
let approver: Option<std::sync::Arc<dyn vak_agent::Approver>> = Some(if yes { - 1911
std::sync::Arc::new(vak_agent::AutoApprove) - 1912
} else { - 1913
std::sync::Arc::new(vak_agent::AutoDeny) - 1914
}); - 1915
let engine = match core.build_permission_engine(&core.extra_allow_snapshot()) { - 1916
// Same execution contract `exec --write-path` carries: direct file - 1917
// mutations outside the declared paths are refused before rules or - 1918
// permission mode are consulted. A planner generates its own write - 1919
// targets, which is precisely why it needs this. - 1920
Ok(e) if !write_paths.is_empty() => e.restrict_write_paths(core.cwd(), &write_paths), - 1921
Ok(e) => e, - 1922
Err(e) => { - 1923
eprintln!("error: {e}"); - 1924
return 2; - 1925
} - 1926
}; - 1927
- 1928
let prepared = core.prepare_turn().await; - 1929
let mut outcome = - 1930
vak_intent::OutcomeSpec::from_reading(task.clone(), &vak_intent::Reading::default(), 0); - 1931
outcome.max_turns = Some(core.effective_max_turns()); - 1932
let deps = vak_flow::ExecutorDeps { - 1933
prompt_layers: Vec::new(), - 1934
provider, - 1935
system_prompt: prepared.system_prompt, - 1936
model: core.effective_model(), - 1937
tools: prepared.tools, - 1938
read_only_tools: prepared.read_only_tools, - 1939
max_turns: core.effective_max_turns(), - 1940
outcome: Some(outcome), - 1941
max_retries: 0, - 1942
retry_base_backoff_ms: 100, - 1943
request_timeout: Some(std::time::Duration::from_secs(600)), - 1944
circuit_breaker: None, - 1945
run_retry_attempts: 0, - 1946
run_retry_base_backoff_ms: 1000, - 1947
dispatch_ceiling: 1, - 1948
spend_gate: None, - 1949
permission: Some(std::sync::Arc::new(engine)), - 1950
mode: match core.effective_permission_mode() { - 1951
vak_config::PermissionMode::ReadOnly => vak_permission::Mode::ReadOnly, - 1952
vak_config::PermissionMode::WorkspaceWrite => vak_permission::Mode::WorkspaceWrite, - 1953
vak_config::PermissionMode::FullAccess => vak_permission::Mode::FullAccess, - 1954
}, - 1955
approval_mode: match core.effective_approval_mode() { - 1956
vak_config::ApprovalMode::Ask => vak_agent::ApprovalMode::Ask, - 1957
vak_config::ApprovalMode::ApproveSafe => vak_agent::ApprovalMode::ApproveSafe, - 1958
vak_config::ApprovalMode::AutoApprove => vak_agent::ApprovalMode::AutoApprove, - 1959
}, - 1960
approver, - 1961
sandbox: core.agent_sandbox(), - 1962
cwd: core.cwd().clone(), - 1963
sessions_home: core.sessions_home().clone(), - 1964
parent_session_id, - 1965
state_path: core.sessions_home().join("flow-runs/plan"), - 1966
agent_identity: core.agent_identity().cloned(), - 1967
conversation_context: core.conversation_context().cloned(), - 1968
work: None, - 1969
}; - 1970
- 1971
let cancel = CancellationToken::new(); - 1972
{ - 1973
let cancel = cancel.clone(); - 1974
tokio::spawn(async move { - 1975
if tokio::signal::ctrl_c().await.is_ok() { - 1976
eprintln!("\n[cancelling…]"); - 1977
cancel.cancel(); - 1978
} - 1979
}); - 1980
} - 1981
- 1982
let (tx, mut rx) = tokio::sync::mpsc::channel::<String>(256); - 1983
let runner = tokio::spawn(async move { - 1984
vak_flow::plan_and_run(std::sync::Arc::new(deps), &task, cancel, tx).await - 1985
}); - 1986
- 1987
while let Some(line) = rx.recv().await { - 1988
eprintln!("{line}"); - 1989
} - 1990
- 1991
match runner.await { - 1992
Ok(vak_flow::PlanOutcome::Completed { outputs, attempts }) => { - 1993
eprintln!("── plan completed after {attempts} attempt(s)"); - 1994
for (id, out) in outputs { - 1995
println!("[{id}]\n{out}\n"); - 1996
} - 1997
0 - 1998
} - 1999
Ok(vak_flow::PlanOutcome::PlanningFailed { reason }) => { - 2000
eprintln!("── planning_failed (fail-closed): {reason}");
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.