1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
1501
1502
1503
1504
1505
1506
1507
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
#![deny(warnings)]
use anyhow::Result;
use chrono::Local;
use serde_json::json;
use std::env;
use std::collections::VecDeque;
use std::io::{self, IsTerminal, Write};
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
mod agent;
mod acp;
mod bash;
mod commands;
mod config;
mod context;
mod files;
mod frame;
mod goal;
mod history;
mod hooks;
mod instructions;
mod plan;
mod lineedit;
mod queue;
mod recents;
mod llm;
mod mode;
mod output;
mod permissions;
mod providers;
mod question;
mod reminders;
mod sandbox;
mod settings;
mod session;
mod shell;
mod skills;
mod status;
mod tools;
mod ui;
use agent::Agent;
use config::ConfigManager;
use hooks::HookEvent;
use tools::ToolDefinition;
fn register_builtin_tools(agent: &mut Agent) {
// get_time tool
let time_def = ToolDefinition::new(
"get_time",
"Get the current date and time",
json!({ "type": "object", "properties": {} }),
);
agent.tools().register(time_def, Box::new(|_| {
Ok(json!({ "time": Local::now().to_string() }))
}));
// echo tool
let echo_def = ToolDefinition::new(
"echo",
"Echo back the input text",
json!({
"type": "object",
"properties": {
"text": { "type": "string" }
},
"required": ["text"]
}),
);
agent.tools().register(echo_def, Box::new(|args| {
let text = args["text"].as_str().unwrap_or("");
Ok(json!({ "echo": text }))
}));
files::register(agent.tools());
// bash tool
let bash_config = bash::BashConfig {
default_timeout: std::time::Duration::from_secs(agent.config().bash_timeout_secs.max(1)),
cancel: Some(agent.control().cancel_flag()),
sandbox: agent.config().sandbox.clone(),
shared_output_dir: Some(agent.spill_dir_handle()),
..Default::default()
};
agent.tools().register(bash::definition(), Box::new(move |args| {
Ok(json!(bash::run(&bash_config, &args)))
}));
// question tool: blocks until the turn loop answers (see question.rs).
// Headless (ACP) sessions have no one to answer, so it errors instead.
let broker = agent.questions();
agent.tools().register(question::definition(), Box::new(move |args| {
if !broker.is_interactive() {
anyhow::bail!("question tool needs an interactive terminal; end your turn with the question, or report_outcome(needs_input) with the question in the summary, instead");
}
let questions = question::parse(&args)?;
let answer = broker.ask_blocking(questions.clone());
Ok(json!(question::result_text(&questions, &answer)))
}));
}
fn register_hooks(agent: &mut Agent) {
// Log all lifecycle events
for event in [
HookEvent::BeforeContextLoad,
HookEvent::AfterContextLoad,
HookEvent::BeforeLLMSend,
HookEvent::AfterLLMResponse,
HookEvent::BeforeToolCall,
HookEvent::AfterToolCall,
] {
let event_name = event.to_string();
agent.hooks().register(event.clone(), Box::new(move |ctx| {
if ui::verbosity() < ui::Verbosity::Debug {
return;
}
match ctx.event {
HookEvent::BeforeContextLoad => {
let input = ctx.data.get("user_input").and_then(|v| v.as_str()).unwrap_or("");
ui::log(&format!("[hook] {} - user: {}", event_name, input));
}
HookEvent::AfterContextLoad => {
let count = ctx.data.get("message_count").and_then(|v| v.as_i64()).unwrap_or(0);
ui::log(&format!("[hook] {} - messages: {}", event_name, count));
}
HookEvent::BeforeLLMSend => {
let iter = ctx.data.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
ui::log(&format!("[hook] {} - iteration: {}", event_name, iter));
}
HookEvent::AfterLLMResponse => {
let has_tools = ctx.data.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
ui::log(&format!("[hook] {} - tool_calls: {}", event_name, has_tools));
}
HookEvent::BeforeToolCall => {
let name = ctx.data.get("tool_name").and_then(|v| v.as_str()).unwrap_or("");
ui::log(&format!("[hook] {} - tool: {}", event_name, name));
}
HookEvent::AfterToolCall => {
let name = ctx.data.get("tool_name").and_then(|v| v.as_str()).unwrap_or("");
ui::log(&format!("[hook] {} - tool: {}", event_name, name));
}
}
}));
}
}
/// Terminal input, read on demand so `/settings` prompts can use stdin directly.
enum TermInput {
Line(String),
Eof,
Interrupt,
/// Ctrl-O: expand or collapse thinking.
ToggleThinking,
/// A lone Esc press.
Escape,
/// Shift+Tab: cycle the agent mode (normal/plan/auto).
CycleMode,
}
/// Esc twice within this window cancels the running turn.
const DOUBLE_ESCAPE_WINDOW: std::time::Duration = std::time::Duration::from_millis(1000);
/// Detects a double Esc press.
#[derive(Default)]
struct DoubleEscape {
last: Option<std::time::Instant>,
}
impl DoubleEscape {
/// Record a press at `now`; true when it completes a double press.
fn press(&mut self, now: std::time::Instant) -> bool {
match self.last.take() {
Some(last) if now.duration_since(last) <= DOUBLE_ESCAPE_WINDOW => true,
_ => {
self.last = Some(now);
false
}
}
}
}
struct Terminal {
want: std::sync::mpsc::Sender<()>,
outstanding: bool,
events: mpsc::UnboundedReceiver<TermInput>,
/// Input read during a turn that is not a queue entry (piped prompts,
/// commands that must wait for the turn to finish).
queued: VecDeque<TermInput>,
/// Messages typed during a turn, waiting to run as later prompts.
/// Editable mid-turn with `/queue` (including removing entries).
messages: queue::MessageQueue,
/// Lines typed during a turn join the message queue (only when stdin is a
/// terminal); piped lines keep the legacy "queued as later prompts" path.
steerable: bool,
/// Where `/settings` saves changes.
config_path: std::path::PathBuf,
/// The line being typed (terminal stdin only).
view: lineedit::SharedView,
/// Recently used models, hoisted in the `/model` type-ahead and saved
/// here on each successful switch.
recents: recents::SharedRecents,
recents_path: std::path::PathBuf,
renderer: std::sync::Arc<ui::Renderer>,
/// Set to make the stdin reader yield the terminal to a foreground picker
/// (a `question`/turn-cap prompt), so the two never race for keystrokes.
suspend: std::sync::Arc<std::sync::atomic::AtomicBool>,
/// Monotonic generation bumped on every `suspend_input()`. It lets a
/// picker's cleanup resume the reader only if no *newer* prompt has since
/// suspended it: a stale auto-away worker must not clear a later prompt's
/// suspension (which would resume stdin under an active dialoguer).
suspend_gen: std::sync::Arc<std::sync::atomic::AtomicU64>,
/// Serialises dialoguer picker workers so at most one ever owns stdin. An
/// auto-away worker that outlived its timeout keeps this held until it
/// exits, so a later question cannot spawn a second stdin reader.
picker_lock: std::sync::Arc<tokio::sync::Mutex<()>>,
}
/// An owned token for one input suspension. Cleanup calls [`InputGate::release`],
/// which resumes the background line reader *only* if this is still the most
/// recent suspension — so a stale worker finishing late cannot clear a newer
/// prompt's gate and resume stdin while that prompt's dialoguer is active.
#[derive(Clone)]
struct InputGate {
flag: std::sync::Arc<std::sync::atomic::AtomicBool>,
generation: std::sync::Arc<std::sync::atomic::AtomicU64>,
token: u64,
}
impl InputGate {
/// Resume the reader iff no later `suspend_input()` has superseded this one.
fn release(&self) {
use std::sync::atomic::Ordering::SeqCst;
if self.generation.load(SeqCst) == self.token {
self.flag.store(false, SeqCst);
}
}
}
impl Terminal {
fn start(
config_path: std::path::PathBuf,
view: lineedit::SharedView,
renderer: std::sync::Arc<ui::Renderer>,
recents: recents::SharedRecents,
recents_path: std::path::PathBuf,
) -> Self {
let (tx, events) = mpsc::unbounded_channel();
let (want, want_rx) = std::sync::mpsc::channel::<()>();
let lines = tx.clone();
let key_mode = io::stdin().is_terminal() && io::stdout().is_terminal();
let reader_view = view.clone();
let suspend = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let reader_suspend = suspend.clone();
std::thread::spawn(move || {
if key_mode {
let mut reader = lineedit::LineReader::with_suspend(reader_suspend);
let send = |key: lineedit::Key| {
let _ = lines.send(match key {
lineedit::Key::Line(line) => TermInput::Line(line),
lineedit::Key::Eof => TermInput::Eof,
lineedit::Key::Interrupt => TermInput::Interrupt,
lineedit::Key::ToggleThinking => TermInput::ToggleThinking,
lineedit::Key::Escape => TermInput::Escape,
lineedit::Key::CycleMode => TermInput::CycleMode,
});
};
for () in want_rx {
let key = reader.read_line(&reader_view, &send);
let eof = matches!(key, lineedit::Key::Eof);
send(key);
if eof {
break;
}
}
return;
}
for () in want_rx {
let mut line = String::new();
match io::stdin().read_line(&mut line) {
Ok(0) | Err(_) => {
let _ = lines.send(TermInput::Eof);
break;
}
Ok(_) => {
if lines.send(TermInput::Line(line)).is_err() {
break;
}
}
}
}
});
tokio::spawn(async move {
while tokio::signal::ctrl_c().await.is_ok() {
if tx.send(TermInput::Interrupt).is_err() {
break;
}
}
});
Self {
want,
outstanding: false,
events,
queued: VecDeque::new(),
messages: queue::MessageQueue::default(),
steerable: io::stdin().is_terminal(),
config_path,
view,
recents,
recents_path,
renderer,
suspend,
suspend_gen: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
picker_lock: std::sync::Arc::new(tokio::sync::Mutex::new(())),
}
}
/// After a successful model switch: hoist the spec in the `/model`
/// type-ahead (persisted), and refresh the config the line editor's
/// argument suggestions read.
fn model_switched(&mut self, agent: &Agent) {
let spec = agent.config().model.clone();
{
let mut recents = self.recents.lock().unwrap();
recents.record(&spec);
recents::save(&self.recents_path, &recents);
}
self.sync_context(agent);
}
/// Refresh the config the line editor's argument suggestions read, after
/// anything that may have changed it (`/settings`, a provider edit).
fn sync_context(&mut self, agent: &Agent) {
let context = self.view.lock().unwrap().context_handle();
context.lock().unwrap().config = agent.config().clone();
}
/// Make the stdin reader yield the terminal so a foreground picker can own
/// it; returns a generation-stamped [`InputGate`] the caller (or an orphaned
/// auto-answer worker) releases once the picker is truly done. The stamp
/// ensures a stale worker cannot resume the reader out from under a newer
/// prompt that has since re-suspended input.
fn suspend_input(&self) -> InputGate {
use std::sync::atomic::Ordering::SeqCst;
let token = self.suspend_gen.fetch_add(1, SeqCst) + 1;
self.suspend.store(true, SeqCst);
InputGate { flag: self.suspend.clone(), generation: self.suspend_gen.clone(), token }
}
/// A clone of the picker serialisation lock (see the field docs).
fn picker_lock(&self) -> std::sync::Arc<tokio::sync::Mutex<()>> {
self.picker_lock.clone()
}
fn request_line(&mut self) {
if !self.outstanding {
self.outstanding = self.want.send(()).is_ok();
}
}
/// The next prompt input at the top-level loop. Deferred terminal input
/// (piped lines, commands typed mid-turn) comes first, then one queued
/// message — and a queued message is handed over only while the agent is
/// not waiting for input: a pending `question`/turn-cap picker owns the
/// terminal, and injecting a queued message there would answer a question
/// the user can see with a prompt they may already have edited or
/// removed. With nothing waiting this reads the terminal.
async fn next(&mut self, agent: &Agent) -> TermInput {
if let Some(input) = self.queued.pop_front() {
return input;
}
if !self.messages.is_empty() && agent.questions().pending().is_none() {
let entry = self.messages.pop().expect("checked non-empty");
self.set_queue_status();
self.renderer.note(&format!("↧ queued #{}: {}", entry.id, entry.text.trim()));
return TermInput::Line(entry.text);
}
self.request_line();
self.recv().await
}
/// Queue a message typed mid-turn; returns its id.
fn queue_message(&mut self, text: &str) -> usize {
let id = self.messages.push(text);
self.set_queue_status();
id
}
/// Apply a `/queue` edit (`remove`/`edit`/`clear`) and say what happened.
fn edit_queue(&mut self, op: &queue::QueueOp) -> String {
let result = queue::apply(&mut self.messages, op);
self.set_queue_status();
result
}
/// Keep the status line's queue count in sync (None hides the segment).
fn set_queue_status(&self) {
let count = self.messages.len();
self.view.lock().unwrap().set_queue_count((count > 0).then_some(count));
}
async fn recv(&mut self) -> TermInput {
let input = self.events.recv().await.unwrap_or(TermInput::Eof);
// These are mid-line events delivered from *inside* an active
// `read_line` (via its `send` callback): the reader keeps running and
// still owns the outstanding read, so they must not clear `outstanding`
// or the next `request_line` would queue a second concurrent reader
// that could race a dialoguer picker for stdin. CycleMode (Shift+Tab)
// is emitted the same way and belongs in this set.
if !matches!(
input,
TermInput::Interrupt | TermInput::ToggleThinking | TermInput::Escape | TermInput::CycleMode
) {
self.outstanding = false;
}
input
}
}
/// Run a turn; typed lines join the message queue (one runs per following
/// turn), `/queue` edits apply immediately, and Ctrl-C or Esc Esc cancels.
async fn run_interactive_turn(agent: &mut Agent, text: &str, terminal: &mut Terminal) -> Result<agent::TurnOutcome> {
let control = agent.control();
let stats = agent.context_stats();
let renderer = terminal.renderer.clone();
if terminal.steerable && ui::verbosity() >= ui::Verbosity::Verbose {
renderer.note("[running: type a message and Enter to queue it (/queue lists, edits, removes), Esc Esc or Ctrl-C to cancel, Ctrl-O to expand thinking]");
}
// Blank line after LLM output (legacy renderer only; the frame renderer
// owns the screen and must not receive stray direct writes).
if !renderer.is_frame() {
println!();
}
renderer.begin_turn();
terminal.view.lock().unwrap().set_mode(lineedit::EditMode::Turn);
let mut escape = DoubleEscape::default();
let outcome = async {
// Grab the broker before the turn future borrows `agent` mutably.
let questions = agent.questions();
let turn = agent.run_turn(None, text);
tokio::pin!(turn);
let mut question_rx = questions.subscribe();
let cap = questions.cap();
let mut cap_rx = cap.subscribe();
loop {
if !terminal.queued.iter().any(|i| matches!(i, TermInput::Eof)) {
terminal.request_line();
}
tokio::select! {
// Poll the turn first: its first poll resets the control, which
// must happen before a buffered cancel or steer is routed to it.
biased;
outcome = &mut turn => break outcome,
// A `question` tool call is waiting for an answer. The handler
// is parked on the blocking pool; answer it here, where we own
// the terminal. While this picker is up the turn future cannot
// resolve, so the message queue is not drained: no queued
// message is injected as a prompt while the agent is asking
// for input (and `Terminal::next` re-checks before draining).
notified = question_rx.changed() => {
if notified.is_err() {
// Broker dropped (agent gone): nothing more to answer.
continue;
}
if let Some(request) = questions.pending() {
// Yield stdin to the picker so the line reader does not
// race it for the answer keystrokes.
let gate = terminal.suspend_input();
// `prompt_question` owns the gate: it clears it while
// holding the picker lock — immediately, or once a
// timed-out auto-answer worker frees stdin — so the
// reader never resumes while a dialoguer is still active.
let answer = prompt_question(request.questions(), &control, &renderer, terminal.picker_lock(), gate).await;
questions.resolve(answer);
// The picker wrote over the owned frame via dialoguer;
// force a full redraw so the next differential render
// isn't computed against stale screen coordinates. If an
// auto-away worker was orphaned it is still parked on
// stdin holding the picker lock, so defer the repaint
// behind that lock: never write the frame while a
// dialoguer still owns the terminal (single-writer).
deferred_frame_resize(&renderer, terminal.picker_lock());
}
}
// The turn hit the cap in normal mode: ask whether to continue.
notified = cap_rx.changed() => {
if notified.is_err() {
continue;
}
let gate = terminal.suspend_input();
let decision = prompt_cap_reached(&renderer, terminal.picker_lock(), gate).await;
cap.decide(decision);
// Repaint after the dialoguer picker clobbered the frame,
// deferred behind the picker lock so it never races an
// orphaned worker still parked on stdin (see above).
deferred_frame_resize(&renderer, terminal.picker_lock());
}
input = terminal.recv() => match input {
TermInput::Interrupt => {
control.cancel();
renderer.urgent_note("[cancelling...]");
}
TermInput::Escape if control.is_cancelled() => {}
TermInput::Escape => {
if escape.press(std::time::Instant::now()) {
control.cancel();
renderer.urgent_note("[cancelling...]");
} else {
renderer.urgent_note("[Esc again to cancel]");
}
}
TermInput::ToggleThinking => {
renderer.toggle_thinking();
}
TermInput::CycleMode => {
// `control` is a shared handle, so this works while the
// turn future holds a `&mut` borrow of the agent.
let mode = control.cycle_mode();
// `Agent::set_mode` (used between turns) also refreshes
// the shared stats; do the equivalent here so the status
// line reflects the new mode immediately, mid-turn.
stats.lock().unwrap().mode = mode;
renderer.event(&agent::AgentEvent::Context);
renderer.note(&format!("[mode: {mode} — {}]", mode.describe()));
}
TermInput::Line(line) if terminal.steerable && !line.trim().is_empty() => {
let text = line.trim();
if let Some(op) = queue_command(text) {
// `/queue` only touches the message queue, so it is
// safe — and most useful — while a turn is running.
match op {
Ok(op) => renderer.note(&terminal.edit_queue(&op)),
Err(usage) => renderer.note(&format!("[{usage}]")),
}
} else if text.starts_with('/') {
renderer.note(&format!("[commands wait for the turn to finish: {text}]"));
terminal.queued.push_back(TermInput::Line(line));
} else {
let id = terminal.queue_message(text);
let n = terminal.messages.len();
renderer.note(&format!("↧ queued #{id} ({n} waiting) — /queue remove {id} to drop"));
}
}
// Blank Enter typed mid-turn is a no-op: dropping it here
// stops it from being deferred into `queued` and replayed
// as an empty line after the turn, which would delay the
// real queued messages/commands behind it.
TermInput::Line(line) if terminal.steerable && line.trim().is_empty() => {}
other => terminal.queued.push_back(other),
},
}
}
}
.await;
renderer.end_turn();
terminal.view.lock().unwrap().set_mode(lineedit::EditMode::Prompt);
let outcome = outcome?;
// A steer typed as the turn finished queues behind what is already
// waiting, unless the turn was cancelled. (The interactive CLI queues
// rather than steers, so this is the ACP path.)
for steer in control.take_pending() {
if outcome.stop_reason == agent::StopReason::Cancelled {
renderer.note(&format!("[steer dropped: {}]", steer.text));
} else {
terminal.queue_message(&steer.text);
}
}
Ok(outcome)
}
/// How long auto mode waits for the user before answering a question itself.
const AUTO_AWAY_SECS: u64 = 15;
/// Render one question and return its answer string, `None` when dismissed.
/// Strip terminal control characters from model-controlled text before it is
/// handed to dialoguer for rendering. `q.question`, option labels and
/// descriptions all originate from the model, so a prompt-injected model could
/// otherwise smuggle ANSI/OSC escape sequences (cursor moves, screen clears,
/// clipboard/title writes) through the interactive picker. Dropping C0/C1
/// control characters — including ESC (0x1B), which begins every such sequence —
/// neutralises them while leaving ordinary printable text intact.
pub(crate) fn sanitize_terminal_text(s: &str) -> String {
s.chars().filter(|c| !c.is_control()).collect()
}
fn ask_one(q: &question::Question) -> Result<Option<String>> {
use dialoguer::{Input, Select};
let mut labels: Vec<String> = q
.options
.iter()
.map(|o| {
let label = sanitize_terminal_text(&o.label);
if o.description.is_empty() {
label
} else {
format!("{} — {}", label, sanitize_terminal_text(&o.description))
}
})
.collect();
let custom_index = if q.custom {
labels.push("Type your own answer".into());
Some(labels.len() - 1)
} else {
None
};
let choice = Select::new()
.with_prompt(sanitize_terminal_text(&q.question))
.items(&labels)
.default(0)
.interact_opt()?;
match choice {
None => Ok(None),
Some(i) if Some(i) == custom_index => {
let text: String = Input::new().with_prompt("Answer").interact_text()?;
Ok(Some(text))
}
Some(i) => Ok(Some(q.options[i].label.clone())),
}
}
/// Force a full frame redraw, but only once the picker lock is free.
///
/// A foreground dialoguer picker owns the terminal while it runs, and an
/// auto-away `prompt_question` can return `Away` while its `spawn_blocking`
/// worker is still parked on stdin holding the picker lock. Repainting the
/// owned frame immediately would write it while that orphaned worker still
/// owns the terminal, violating the single-writer assumption and letting the
/// two tear each other's output. Awaiting the picker lock first defers the
/// repaint until every picker (orphaned or not) has released the terminal; in
/// the common case the lock is already free, so the redraw runs at once.
/// Spawned so the turn loop is never blocked waiting on an away user.
fn deferred_frame_resize(
renderer: &std::sync::Arc<ui::Renderer>,
picker_lock: std::sync::Arc<tokio::sync::Mutex<()>>,
) {
let renderer = renderer.clone();
tokio::spawn(async move {
let _guard = picker_lock.lock_owned().await;
renderer.frame_resize();
});
}
/// Render a pending `question` and return the answer, plus an optional handle
/// to a still-running auto-answer worker. In auto mode the user gets
/// `AUTO_AWAY_SECS` to respond before the question is answered with the away
/// message; a `spawn_blocking` dialoguer worker cannot be aborted, so when the
/// timeout fires the worker is handed back (still parked on stdin) for the
/// caller to await before it resumes the line reader — the two must never read
/// keystrokes at once. Runs on the turn loop, which owns the terminal.
async fn prompt_question(
questions: &[question::Question],
control: &agent::TurnControl,
renderer: &std::sync::Arc<ui::Renderer>,
picker_lock: std::sync::Arc<tokio::sync::Mutex<()>>,
gate: InputGate,
) -> question::QuestionAnswer {
use question::QuestionAnswer;
let auto = control.mode() == mode::AgentMode::Auto;
// The whole prompt run: ask each question, collecting one string each.
// Dismissal (Esc) at any question dismisses the lot.
let ask = |questions: &[question::Question]| -> Result<QuestionAnswer> {
let mut answers = Vec::new();
for q in questions {
match ask_one(q)? {
Some(answer) => answers.push(answer),
None => return Ok(QuestionAnswer::Dismissed),
}
}
Ok(QuestionAnswer::Answers(answers))
};
if !auto {
// Normal mode: the user is present, so it is fine to block until any
// previous (possibly orphaned) worker releases stdin, so only one ever
// reads keystrokes.
let guard = picker_lock.lock_owned().await;
let questions = questions.to_vec();
let asked = tokio::task::spawn_blocking(move || ask(&questions)).await;
// Release the gate while still holding the picker guard, so the input
// reader cannot resume before the next picker (which must take this
// same lock) has re-suspended it. `release()` is generation-aware, so
// it is a no-op if a newer prompt has already re-suspended input.
gate.release();
drop(guard);
return asked.ok().and_then(Result::ok).unwrap_or(QuestionAnswer::Dismissed);
}
// Auto mode: give the user a chance to answer, then answer ourselves.
renderer.note(&format!("[auto: answering for you in {AUTO_AWAY_SECS}s — the user is away]"));
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(AUTO_AWAY_SECS);
// Acquire stdin, but never past the away deadline. If a prior orphaned
// worker still holds the picker lock while the user is away, blocking here
// would keep this question from ever resolving (finding: later questions
// wait on the lock before their own timeout can start). Bound the wait so
// it still answers `Away`, and serialise so only one worker reads stdin.
let guard = match tokio::time::timeout_at(deadline, picker_lock.clone().lock_owned()).await {
Ok(guard) => guard,
Err(_) => {
renderer.note("[auto: no answer — making the best decision]");
// Do not start a competing reader. Keep the caller's input gate
// suspended until the picker frees (the prior orphan exits), then
// release the gate while still holding the guard, so a later picker
// cannot acquire the lock and re-suspend between our lock release
// and the gate release (which would let the reader race stdin).
tokio::spawn(async move {
let guard = picker_lock.lock_owned().await;
gate.release();
drop(guard);
});
return QuestionAnswer::Away;
}
};
let questions = questions.to_vec();
let mut worker = tokio::task::spawn_blocking(move || ask(&questions));
tokio::select! {
joined = &mut worker => {
// Release the gate while still holding the guard, then drop it.
gate.release();
drop(guard);
joined.ok().and_then(Result::ok).unwrap_or(QuestionAnswer::Dismissed)
}
() = tokio::time::sleep_until(deadline) => {
renderer.note("[auto: no answer — making the best decision]");
// Keep the worker alive (it is still blocked on stdin) and hold the
// picker lock until it exits. Release the gate while still holding
// the guard so the reader resumes only once this worker exits, with
// no window for a later picker to slip in between lock and gate
// release.
tokio::spawn(async move {
let _ = worker.await;
gate.release();
drop(guard);
});
QuestionAnswer::Away
}
}
}
/// Ask whether to keep going when the turn cap is reached. Defaults to stop,
/// so an unattended prompt does not run away. Runs on the turn loop.
async fn prompt_cap_reached(
renderer: &std::sync::Arc<ui::Renderer>,
picker_lock: std::sync::Arc<tokio::sync::Mutex<()>>,
gate: InputGate,
) -> question::CapDecision {
use question::CapDecision;
// Serialise with the question picker: this is another dialoguer reader, so
// never run it while an orphaned auto-away worker still holds stdin. Block
// until that worker releases the lock so only one ever reads keystrokes.
let guard = picker_lock.lock_owned().await;
let renderer = renderer.clone();
let decision = tokio::task::spawn_blocking(move || {
let keep_going = dialoguer::Confirm::new()
.with_prompt("Reached the turn cap without a final answer. Keep going?")
.default(false)
.interact_opt()
.ok()
.flatten()
.unwrap_or(false);
if keep_going {
renderer.note("[continuing past the turn cap]");
CapDecision::Continue
} else {
CapDecision::Stop
}
})
.await
.unwrap_or(CapDecision::Stop);
// Release the gate while still holding the guard, so the input reader cannot
// resume before the next picker re-suspends it (generation-aware: a no-op if
// a newer prompt already re-suspended input).
gate.release();
drop(guard);
decision
}
/// The `/queue` edit a line asks for, if it is a `/queue` command. Unlike
/// other commands, `/queue` only touches the message queue, so it runs even
/// while a turn or compaction is in flight — that is when queue edits
/// (removals included) are most useful.
fn queue_command(text: &str) -> Option<std::result::Result<queue::QueueOp, String>> {
let args = match text.strip_prefix("/queue") {
Some(args) if args.is_empty() || args.starts_with(char::is_whitespace) => args,
_ => return None,
};
Some(queue::parse(args))
}
/// Emit a transient diagnostic: through the frame transcript in frame mode (a
/// direct write would corrupt the owned frame), else to stderr as before.
fn diag(renderer: &ui::Renderer, text: &str) {
if renderer.is_frame() {
renderer.note(text);
} else {
eprintln!("{text}");
}
}
/// Run an explicit compaction; Ctrl-C or Esc Esc cancels it.
async fn run_compaction(
agent: &mut Agent,
mode: Option<config::CompactionMode>,
instructions: Option<&str>,
terminal: &mut Terminal,
) -> Result<Option<agent::CompactReport>> {
let control = agent.control();
let stats = agent.context_stats();
let compaction = agent.compact(mode, instructions);
tokio::pin!(compaction);
let mut escape = DoubleEscape::default();
loop {
tokio::select! {
biased;
report = &mut compaction => return report,
input = terminal.recv() => match input {
TermInput::Interrupt => {
control.cancel();
diag(&terminal.renderer, "[cancelling...]");
}
TermInput::Escape if control.is_cancelled() => {}
TermInput::Escape => {
if escape.press(std::time::Instant::now()) {
control.cancel();
diag(&terminal.renderer, "[cancelling...]");
} else {
diag(&terminal.renderer, "[Esc again to cancel]");
}
}
TermInput::ToggleThinking => {
terminal.renderer.toggle_thinking();
}
TermInput::CycleMode => {
let mode = control.cycle_mode();
// Mirror the prompt/turn `CycleMode` path: refresh the
// shared stats and emit a context event so the status line
// reflects the new mode immediately, not just after a later
// refresh while `/compact` is still running.
stats.lock().unwrap().mode = mode;
terminal.renderer.event(&agent::AgentEvent::Context);
diag(&terminal.renderer, &format!("[mode: {mode} — {}]", mode.describe()));
}
other => {
// `/queue` edits apply mid-compaction too; everything else
// waits for the loop to pick it up afterwards.
if let TermInput::Line(line) = &other
&& let Some(op) = queue_command(line.trim())
{
match op {
Ok(op) => {
let result = terminal.edit_queue(&op);
terminal.renderer.note(&result);
}
Err(usage) => terminal.renderer.note(&format!("[{usage}]")),
}
} else {
terminal.queued.push_back(other);
}
}
},
}
}
}
async fn run_command(agent: &mut Agent, cmd: &str, terminal: &mut Terminal) -> Result<bool> {
match cmd {
"/exit" | "/quit" => Ok(false),
"/help" => {
terminal.renderer.print_block(&commands::help_text());
Ok(true)
}
_ if cmd == "/compact" || cmd.starts_with("/compact ") => {
let (mode, focus) = commands::parse_compact_args(&cmd["/compact".len()..]);
let focus = focus.map(str::to_string);
terminal.renderer.print_block("Compacting...");
match run_compaction(agent, mode, focus.as_deref(), terminal).await? {
Some(report) => terminal.renderer.print_block(&format!("Conversation {report}")),
None => terminal.renderer.print_block("Nothing to compact"),
}
Ok(true)
}
"/context" => {
let stats = agent.context_stats().lock().unwrap().clone();
let mut out: Vec<String> = Vec::new();
out.push(format!("Model: {}/{}", stats.provider, stats.model));
out.push(format!(
"Context: {}{} of {} tokens ({:.1}%){}",
if stats.calibrated { "" } else { "~" },
stats.tokens,
stats.window,
stats.percent(),
if stats.calibrated { ", anchored to reported usage" } else { ", estimated" }
));
out.push(format!("Messages: {}", stats.messages));
let system_tokens = crate::context::text_tokens(&agent.system_prompt());
out.push(format!("System prompt: {} tokens", system_tokens));
let files = agent.project_instruction_files();
if files.is_empty() {
out.push("Instructions: none (no AGENTS.md, CLAUDE.md or .github/copilot-instructions.md found)".to_string());
} else {
out.push(format!("Instructions: {}", files.join(", ")));
}
let skills = agent.skills();
if !skills.is_empty() || !skills.warnings.is_empty() {
out.push(format!("Skills: {} (/skills to list them)", skills.skills.len()));
}
if let Some((done, total)) = stats.plan {
out.push(format!("Plan: {done}/{total} done (/plan to show it)"));
}
out.push(format!("Session: {} input, {} output tokens", stats.session_input_tokens, stats.session_output_tokens));
if let Some(aic) = stats.session_aic {
out.push(format!("AI Credits: {aic:.2} used this session"));
}
match stats.auto_compact {
Some(t) => out.push(format!(
"Auto-compact: at {:.0}% (~{} tokens); compacted {} time(s)",
t * 100.0,
(stats.window as f64 * t) as usize,
stats.compactions
)),
None => out.push("Auto-compact: off".to_string()),
}
out.push(format!("Compaction: {} mode (/compact --smart or --standard overrides once)", agent.config().compaction_mode.as_str()));
if stats.history_searches + stats.history_reads > 0 {
out.push(format!("History: {} search(es), {} read(s) this session", stats.history_searches, stats.history_reads));
}
out.push(format!("(context window {})", agent.context_window_with_source().1));
terminal.renderer.print_block(&out.join("\n"));
Ok(true)
}
"/settings" if terminal.outstanding => {
// Typed during a turn: a stdin read is still pending, so an
// interactive editor would race it for keystrokes.
let config = agent.config();
terminal.renderer.print_block(&format!(
"model: {}\ntemperature: {}\nmax_tokens: {}\n(read-only: run /settings again at the prompt to edit)",
config.model, config.temperature, config.max_tokens
));
Ok(true)
}
"/settings" => {
let before = agent.config().model.clone();
settings::run(agent, &terminal.config_path).await?;
// The settings dialog (dialoguer) wrote directly over the owned
// frame; force a full redraw so the frame renderer's next update
// isn't diffed against stale screen coordinates.
terminal.renderer.frame_resize();
// Providers or the model may have changed. A model switched through
// the settings dialog must land in the recents MRU just like one
// switched with `/model`; a provider-only edit just refreshes the
// config the line editor's argument suggestions read.
if agent.config().model != before {
terminal.model_switched(agent);
} else {
terminal.sync_context(agent);
}
Ok(true)
}
"/tools" => {
let mut out = vec!["Available tools:".to_string()];
for def in agent.tool_definitions() {
out.push(format!(" {} - {}", def.name, def.description));
}
terminal.renderer.print_block(&out.join("\n"));
Ok(true)
}
"/skills" => {
let skills = agent.skills();
let mut out: Vec<String> = Vec::new();
if skills.is_empty() {
out.push(format!("No skills found (looked in {}, ai.lock and {}).", agent.config().skills.dirs.join(", "), agent.config().skills.user_dirs.join(", ")));
}
for skill in &skills.skills {
out.push(format!(" {} - {}\n {}", skill.name, skill.description, skill.dir.display()));
}
for warning in &skills.warnings {
out.push(format!("Warning: {warning}"));
}
terminal.renderer.print_block(&out.join("\n"));
Ok(true)
}
"/plan" => {
if agent.plan().is_empty() {
terminal.renderer.print_block("No plan yet. The agent makes one with the plan_add tool.");
} else {
terminal.renderer.print_block(agent.plan().render(true, usize::MAX).trim_end());
}
Ok(true)
}
_ if let Some(op) = queue_command(cmd) => {
// List is the read-only form; the edits (remove/edit/clear) were
// already applied mid-turn when typed then, and apply here at the
// prompt.
match op {
Ok(queue::QueueOp::List) => terminal.renderer.print_block(&queue::describe(&terminal.messages)),
Ok(op) => {
let msg = terminal.edit_queue(&op);
terminal.renderer.print_block(&msg);
}
Err(usage) => terminal.renderer.print_block(&usage),
}
Ok(true)
}
"/model" if terminal.outstanding => {
// Typed during a turn: a stdin read is still pending, so an
// interactive picker would race it for keystrokes.
terminal.renderer.print_block(&format!(
"Model: {} (provider {}, spec {:?})\n(read-only: run /model again at the prompt to switch)",
agent.model_name(), agent.provider_name(), agent.config().model
));
Ok(true)
}
"/model" if !io::stdin().is_terminal() || !io::stderr().is_terminal() => {
// The picker reads keystrokes from stdin and draws on stderr, so it
// needs both to be terminals; piped input/output just gets the
// current model.
terminal.renderer.print_block(&format!("Model: {} (provider {}, spec {:?})", agent.model_name(), agent.provider_name(), agent.config().model));
Ok(true)
}
"/model" => {
let picked = settings::pick_model_interactive(agent).await;
// The picker (dialoguer) wrote directly over the owned frame; force
// a full redraw so the next differential render isn't diffed against
// stale screen coordinates. Do it before propagating any error so
// the frame is repaired on the error path too.
terminal.renderer.frame_resize();
if let Some(spec) = picked? {
agent.set_model(&spec).await?;
terminal.model_switched(agent);
terminal.renderer.print_block(&format!("Model set to {} (provider {})", agent.model_name(), agent.provider_name()));
} else {
terminal.renderer.print_block(&format!("Model unchanged: {} (provider {})", agent.model_name(), agent.provider_name()));
}
Ok(true)
}
_ if cmd.starts_with("/model ") => {
agent.set_model(cmd["/model ".len()..].trim()).await?;
terminal.model_switched(agent);
terminal.renderer.print_block(&format!("Model set to {} (provider {})", agent.model_name(), agent.provider_name()));
Ok(true)
}
"/providers" => {
let (user, default_provider) = agent.config().effective_providers();
let mut out = vec![format!("Providers (default: {default_provider}):")];
for (name, provider) in providers::effective_providers(&user) {
let kind = provider.kind.map(|k| format!("{k:?}").to_lowercase()).unwrap_or_else(|| "?".into());
let key = settings::key_status(&provider);
let url = provider.base_url.unwrap_or_else(|| match provider.kind {
Some(providers::ProviderKind::GithubCopilot) => "(from session token)".into(),
_ => "-".into(),
});
out.push(format!(" {name:<14} {kind:<14} {url:<55} {key}"));
}
terminal.renderer.print_block(&out.join("\n"));
Ok(true)
}
"/session" => {
match (agent.session_id(), agent.session_path()) {
(Some(id), Some(path)) => terminal.renderer.print_block(&format!("Session {id}: {}", path.display())),
_ => terminal.renderer.print_block("Session persistence is disabled"),
}
Ok(true)
}
"/restart" => {
// Start a brand-new session in place: new ID, context reset to
// just the system prompt, empty plan, counters zeroed. The
// previous session log stays on disk and is still resumable.
let id = agent.new_session()?;
terminal.renderer.clear_screen();
if agent.session_path().is_some() {
terminal.renderer.print_block(&format!("Session: {id} (resume with --resume {id})"));
} else {
terminal.renderer.print_block(&format!("Session: {id}"));
}
Ok(true)
}
"/verbosity" => {
let current = ui::verbosity();
let mut out = vec![format!("Verbosity: {current} ({})", current.describe())];
for level in ui::Verbosity::ALL {
out.push(format!(" {:<8} {}", level.to_string(), level.describe()));
}
terminal.renderer.print_block(&out.join("\n"));
Ok(true)
}
"/mode" => {
let current = agent.mode();
let mut out = vec![format!("Mode: {current} ({})", current.describe())];
for mode in mode::AgentMode::ALL {
out.push(format!(" {:<8} {}", mode.to_string(), mode.describe()));
}
out.push("(Shift+Tab cycles; /mode NAME sets it directly)".to_string());
terminal.renderer.print_block(&out.join("\n"));
Ok(true)
}
_ if cmd.starts_with("/mode ") => {
match cmd["/mode ".len()..].parse::<mode::AgentMode>() {
Ok(m) => {
agent.set_mode(m);
terminal.renderer.print_block(&format!("Mode set to {m} ({})", m.describe()));
}
Err(e) => terminal.renderer.print_block(&e),
}
Ok(true)
}
_ if cmd.starts_with("/verbosity ") => {
match cmd["/verbosity ".len()..].parse::<ui::Verbosity>() {
Ok(level) => {
ui::set_verbosity(level);
agent.config_mut().verbosity = level;
terminal.renderer.print_block(&format!("Verbosity set to {level} ({}); /settings saves it", level.describe()));
}
Err(e) => terminal.renderer.print_block(&e),
}
Ok(true)
}
_ => {
let outcome = run_interactive_turn(agent, cmd, terminal).await?;
// In frame mode the turn's response is already rendered from its
// events; re-printing it here would duplicate the answer and
// corrupt the owned frame.
if !terminal.renderer.is_frame() {
if ui::verbosity() == ui::Verbosity::Quiet {
println!("{}", ui::stamp_block(&outcome.response));
} else if outcome.stop_reason == agent::StopReason::Cancelled {
println!("{}", ui::stamp_block(&format!("\x1b[2m{}\x1b[0m", outcome.response)));
} else if outcome.stop_reason == agent::StopReason::MaxTurnRequests {
let last = outcome.response.lines().last().unwrap_or_default();
println!("{}", ui::stamp_block(&format!("\x1b[2m{last}\x1b[0m")));
}
}
Ok(true)
}
}
}
/// Restores the full screen and terminal modes when the interactive loop ends.
struct StatusGuard(Option<std::sync::Arc<status::StatusLine>>);
impl Drop for StatusGuard {
fn drop(&mut self) {
lineedit::restore_terminal();
if let Some(status) = &self.0 {
status.teardown();
}
}
}
struct Args {
acp: bool,
login: Option<String>,
list_models: Option<String>,
model: Option<String>,
resume: Option<String>,
config: Option<std::path::PathBuf>,
verbosity: Option<ui::Verbosity>,
sandbox: Option<sandbox::SandboxMode>,
allow: Vec<String>,
deny: Vec<String>,
}
fn print_version() {
println!("nano-coder {}", env!("CARGO_PKG_VERSION"));
}
fn print_help() {
println!("Usage: nano-coder [--acp] [--model provider/model] [--resume SESSION_ID] [--config PATH]");
println!(" [--verbosity quiet|normal|verbose|debug]");
println!(" [--sandbox off|workspace|read-only] [--allow RULE]... [--deny RULE]...");
println!(" nano-coder --login github-copilot");
println!(" nano-coder --list-models PROVIDER[/model]");
println!(" nano-coder --version");
}
fn parse_args() -> Result<Args> {
// Handle early-exit flags before the value-consuming loop so a preceding
// value-taking option (e.g. `--model --version`) can't swallow them.
for arg in env::args().skip(1) {
match arg.as_str() {
"-V" | "--version" => {
print_version();
std::process::exit(0);
}
"-h" | "--help" => {
print_help();
std::process::exit(0);
}
_ => {}
}
}
let mut args = Args {
acp: false,
login: None,
list_models: None,
model: None,
resume: None,
config: None,
verbosity: None,
sandbox: None,
allow: Vec::new(),
deny: Vec::new(),
};
let mut iter = env::args().skip(1);
while let Some(arg) = iter.next() {
let mut value = |name: &str| iter.next().ok_or_else(|| anyhow::anyhow!("{name} requires a value"));
match arg.as_str() {
"--acp" => args.acp = true,
"--login" => args.login = Some(value("--login")?),
"--list-models" => args.list_models = Some(value("--list-models")?),
"--model" => args.model = Some(value("--model")?),
"--resume" => args.resume = Some(value("--resume")?),
"--config" => args.config = Some(value("--config")?.into()),
"--verbosity" | "-v" => {
args.verbosity = Some(value("--verbosity")?.parse().map_err(|e: String| anyhow::anyhow!(e))?)
}
"--sandbox" => {
args.sandbox = Some(value("--sandbox")?.parse().map_err(|e: String| anyhow::anyhow!(e))?)
}
"--allow" => args.allow.push(value("--allow")?),
"--deny" => args.deny.push(value("--deny")?),
"-V" | "--version" => {
print_version();
std::process::exit(0);
}
"-h" | "--help" => {
print_help();
std::process::exit(0);
}
other => anyhow::bail!("unknown argument {other:?} (see --help)"),
}
}
Ok(args)
}
#[tokio::main]
async fn main() -> Result<()> {
// Detect execution mode from command-line args
let args = parse_args()?;
if let Some(provider) = &args.login {
if provider != "github-copilot" {
anyhow::bail!("--login supports only github-copilot (other providers use API keys)");
}
let path = providers::github_copilot::login().await?;
println!("Saved GitHub Copilot credentials to {}", path.display());
return Ok(());
}
// Load config
let config_mgr = match &args.config {
Some(path) => ConfigManager::from_path(path.clone())?,
None => ConfigManager::new()?,
};
let config_path = config_mgr.config_path().to_path_buf();
let mut config = config_mgr.get().clone();
if let Some(spec) = &args.list_models {
let (user, default_provider) = config.effective_providers();
let client = providers::build_lister(spec, &user, &default_provider)?;
for model in client.list_models().await? {
println!("{}/{model}", client.provider_name());
}
return Ok(());
}
if let Some(model) = args.model.clone().or_else(|| env::var("AGENTIC_HARNESS_MODEL").ok().filter(|m| !m.is_empty())) {
config.model = model;
}
if let Some(level) = args.verbosity {
config.verbosity = level;
}
let env_sandbox = env::var("NANO_CODER_SANDBOX").ok().filter(|m| !m.is_empty());
if let Some(mode) = env_sandbox.map(|m| m.parse::<sandbox::SandboxMode>()).transpose().map_err(|e| anyhow::anyhow!("NANO_CODER_SANDBOX: {e}"))? {
config.sandbox.mode = mode;
}
if let Some(mode) = args.sandbox {
config.sandbox.mode = mode;
}
config.permissions.allow.extend(args.allow.iter().cloned());
config.permissions.deny.extend(args.deny.iter().cloned());
ui::set_verbosity(config.verbosity);
ui::set_timestamps(config.timestamps);
// Create agent with the configured provider
let mut agent = Agent::from_config(config)?;
agent.detect_context_window().await;
// Register tools and hooks
register_builtin_tools(&mut agent);
register_hooks(&mut agent);
if args.acp {
// ACP headless mode; sessions start with session/new or session/load
if let Some(id) = &args.resume {
agent.load_session(id)?;
}
eprintln!("ACP harness ready (provider: {}, model: {})", agent.provider_name(), agent.model_name());
let saw_valid = acp::run_acp(&mut agent).await?;
if !saw_valid {
eprintln!(
"no valid ACP requests received on stdin — is the client speaking ACP (JSON-RPC 2.0, one message per line)?"
);
std::process::exit(2);
}
} else {
// The interactive CLI answers `question` tool calls, but only when a
// real terminal is attached: with piped stdin/stdout the picker cannot
// be driven, so `question` must take the documented headless path
// instead of blocking in dialoguer while `Terminal` also reads stdin.
let interactive = io::stdin().is_terminal() && io::stdout().is_terminal();
agent.questions().set_interactive(interactive);
match &args.resume {
Some(id) => agent.load_session(id)?,
None => {
if agent.config().persist_sessions {
agent.new_session()?;
} else {
agent.apply_project_instructions();
}
}
}
// Interactive CLI mode
println!("nano-coder v{}", env!("CARGO_PKG_VERSION"));
println!("Model: {} (provider: {})", agent.model_name(), agent.provider_name());
if let Some(id) = agent.session_id() {
println!("Session: {id} (resume with --resume {id})");
}
for file in agent.project_instruction_files() {
println!("Instructions: {file}");
}
let skills = agent.skills();
if !skills.is_empty() {
println!("Skills: {}", skills.names().join(", "));
}
for warning in &skills.warnings {
println!("Skills warning: {warning}");
}
println!("Type /help for commands\n");
// Main loop
let status = status::StatusLine::install(agent.context_stats());
let _status_guard = StatusGuard(status.clone());
let previous_hook = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
lineedit::restore_terminal();
previous_hook(info);
}));
let renderer = ui::Renderer::new(status.clone(), agent.config().renderer);
ui::install(renderer.clone());
let sink = renderer.clone();
agent.set_event_sink(Box::new(move |_, event| sink.event(event)));
agent.set_streaming(true);
agent.refresh_stats();
let frame_mode = renderer.is_frame();
let recents_path = recents::default_path();
let recents: recents::SharedRecents = Arc::new(Mutex::new(recents::load(&recents_path)));
let view = {
let context = Arc::new(Mutex::new(lineedit::EditContext { config: agent.config().clone(), recents: recents.clone() }));
lineedit::EditView::shared(status.clone(), context)
};
if frame_mode {
// The app-owned frame renderer draws the editor row itself; route
// every edit through it instead of the inline/scroll-region path.
let renderer = renderer.clone();
view.lock().unwrap().set_edit_hook(Arc::new(move |line: &str, cursor: usize, queued: usize| renderer.set_editor(line, cursor, queued)));
}
if let Ok(mut resized) =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::window_change())
{
let view = view.clone();
let status = status.clone();
let renderer = renderer.clone();
tokio::spawn(async move {
// Debounce a burst of resizes (a window drag) into one render at
// the final size: after a resize, wait for ~40 ms of quiet.
let quiet = std::time::Duration::from_millis(40);
let mut debounce = frame::Debouncer::new(quiet);
while resized.recv().await.is_some() {
if !renderer.is_frame() {
// Legacy: re-anchor immediately, as before.
if let Some(status) = &status {
status.resize();
}
view.lock().unwrap().resize();
continue;
}
debounce.record(std::time::Instant::now());
// Coalesce further resizes arriving within the quiet window.
while debounce.pending() {
tokio::select! {
more = resized.recv() => match more {
Some(()) => debounce.record(std::time::Instant::now()),
None => break,
},
_ = tokio::time::sleep(quiet) => {
if debounce.ready(std::time::Instant::now()) {
debounce.clear();
}
}
}
}
// `frame_resize()` re-renders every row (editor included)
// at the new size in one pass; calling `view.resize()` here
// too would fire the edit hook and emit a second redraw.
renderer.frame_resize();
}
});
}
// A resumed session loads its conversation before the sink is wired,
// so the app-owned frame opens empty. Replay the loaded history now
// (the sink is installed) to reconstruct the transcript into
// `FrameState`; `frame_event` populates user, assistant, tool and plan
// items. Only in frame mode — the legacy renderer would dump the whole
// conversation inline, which it has never done on resume.
if frame_mode && args.resume.is_some() {
agent.replay_history();
}
let mut terminal = Terminal::start(config_path, view, renderer, recents, recents_path);
let mut running = true;
let mut exit_armed = false;
let mut separate = false;
while running {
if let Some(status) = &status
&& !frame_mode
{
status.draw();
}
let prompt = |terminal: &Terminal, separate: bool| {
if terminal.queued.is_empty() && terminal.messages.is_empty() {
let mut view = terminal.view.lock().unwrap();
if frame_mode {
// The frame renderer owns the screen: refresh the editor
// row (and thus the whole frame) instead of writing an
// inline prompt.
let (line, cursor) = view.snapshot();
terminal.renderer.set_editor(&line, cursor, 0);
view.prompt_redrawn();
return;
}
let prompt = view.prompt();
// Serialise the prompt write under the terminal lock so it
// cannot move the cursor mid-way through the SIGWINCH
// anchor's query-to-scroll critical section (status::anchor).
crate::status::with_term_lock(|| {
let mut out = io::stdout().lock();
if separate {
let _ = out.write_all(b"\n");
}
let _ = out.write_all(prompt.as_bytes());
let _ = out.flush();
});
view.prompt_redrawn();
}
};
prompt(&terminal, separate);
separate = false;
let input = loop {
match terminal.next(&agent).await {
TermInput::ToggleThinking => {
if terminal.renderer.toggle_thinking() {
prompt(&terminal, false);
}
}
TermInput::CycleMode => {
let mode = agent.control().cycle_mode();
agent.set_mode(mode);
if frame_mode {
terminal.renderer.note(&format!("Mode: {mode} ({})", mode.describe()));
} else {
println!("\nMode: {mode} ({})", mode.describe());
}
prompt(&terminal, false);
}
other => break other,
}
};
let input = match input {
TermInput::Eof => break,
TermInput::Interrupt if exit_armed => break,
TermInput::Interrupt => {
exit_armed = true;
if frame_mode {
terminal.renderer.note("(Ctrl-C again to exit)");
} else {
println!("\n(Ctrl-C again to exit)");
}
continue;
}
TermInput::ToggleThinking | TermInput::Escape | TermInput::CycleMode => continue,
TermInput::Line(line) => line.trim().to_string(),
};
exit_armed = false;
if input.trim().is_empty() {
continue;
}
if frame_mode && !input.starts_with('/') {
terminal.renderer.frame_user_message(&input);
}
match run_command(&mut agent, &input, &mut terminal).await {
Ok(continue_running) => {
running = continue_running;
}
Err(e) => {
if frame_mode {
terminal.renderer.note(&format!("Error: {:#}", e));
} else {
eprintln!("Error: {:#}", e);
}
}
}
separate = true;
}
println!("\nGoodbye!");
// Repeat the resume instruction on exit so it is still on screen (and
// in scrollback) after a long session has pushed the start-up banner
// away. `/restart` may have swapped the session mid-run, so re-read
// the current ID rather than remembering the start-up one. Gate on
// `session_path()` (not `session_id()`): when persistence is disabled
// `/restart` still assigns a session ID even though nothing is written
// to disk, so printing a `--resume` command there would be unusable.
if agent.session_path().is_some()
&& let Some(id) = agent.session_id()
{
println!("Session: {id} (resume with --resume {id})");
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::{Duration, Instant};
#[test]
fn double_escape_needs_two_presses_within_the_window() {
let mut escape = DoubleEscape::default();
let t = Instant::now();
assert!(!escape.press(t));
assert!(escape.press(t + Duration::from_millis(400)));
// The pair is consumed: the next press starts over.
assert!(!escape.press(t + Duration::from_millis(500)));
// Too slow: the second press re-arms instead of cancelling.
assert!(!escape.press(t + Duration::from_millis(1600)));
assert!(escape.press(t + Duration::from_millis(1700)));
}
#[test]
fn sanitize_terminal_text_strips_control_and_escape_sequences() {
// A prompt-injected ANSI/OSC payload is neutralised, printable text kept.
assert_eq!(sanitize_terminal_text("hi\x1b[2Jthere"), "hi[2Jthere");
assert_eq!(sanitize_terminal_text("a\x07\x00b\tc"), "abc");
assert_eq!(sanitize_terminal_text("plain — label"), "plain — label");
}
}