1use std::io::{self, BufRead, IsTerminal, Write};
2use std::path::{Path, PathBuf};
3use std::sync::mpsc;
4
5use serde::Deserialize;
6use serde_json::{Map, Value};
7
8use crate::cancellation::CancellationToken;
9use crate::config::{AuthProvider, Config, LlmSettings};
10use crate::context::{resolve_boot_context_with_api_key_env, InstructionSource, SkillEntry};
11use crate::model::{estimate_context_tokens, estimate_message_tokens, ChatMessage, ChatToolCall};
12use crate::protocol::{EventSink, ProtocolEvent, ProtocolWriter};
13use crate::provider::{Provider, ProviderStreamEvent, ProviderTurn};
14use crate::redaction::{
15 conflicts_with_protected_literal, conflicts_with_tui_literal, is_structural_key, redact_secret,
16 redaction_marker,
17};
18use crate::session::Session;
19
20#[derive(Debug)]
21struct CliOptions {
22 session: Option<String>,
23 list_sessions: bool,
24 jsonl: bool,
25 tui: bool,
26 version: bool,
27 command: Option<CliCommand>,
28}
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq)]
31enum CliCommand {
32 CodexLogin,
33 CodexLogout,
34}
35
36#[derive(Debug, Deserialize)]
37struct InputRecord {
38 #[serde(rename = "type")]
39 record_type: String,
40 text: Option<String>,
41}
42
43const USER_CANCEL_REASON: &str = "user_cancelled";
44const PROVIDER_PHASE: &str = "provider_stream";
45const COMMAND_PHASE: &str = "cmd";
46const AUTO_COMPACTION_THRESHOLD_PERCENT: usize = 95;
47const COMPACTION_KEEP_RECENT_TOKENS: usize = 20_000;
48const COMPACTION_SYSTEM_PROMPT: &str = "You are compacting a coding-agent conversation. Produce a concise, factual continuation summary. Preserve the user's goals, explicit decisions, constraints, files and code changes, commands and results, current implementation state, unresolved work, and exact identifiers that future turns need. Do not invent facts. Return only the summary text; do not call tools.";
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub enum FrontendMode {
52 Jsonl,
53 Tui,
54}
55
56pub fn run_cli<R, W, E>(args: &[String], input: R, output: W, diagnostics: E) -> i32
57where
58 R: BufRead + Send + 'static,
59 W: Write,
60 E: Write,
61{
62 let options = match parse_args(args) {
63 Ok(options) => options,
64 Err(error) => {
65 let mut diagnostics = diagnostics;
66 write_diagnostic(&mut diagnostics, &error);
67 return 2;
68 }
69 };
70 if options.version {
71 if let Err(error) = write_version(output) {
72 let mut diagnostics = diagnostics;
73 write_diagnostic(
74 &mut diagnostics,
75 &format!("unable to write version: {error}"),
76 );
77 return 1;
78 }
79 return 0;
80 }
81
82 let home = match home_directory() {
83 Ok(home) => home,
84 Err(error) => {
85 let mut diagnostics = diagnostics;
86 write_diagnostic(&mut diagnostics, &error);
87 return 1;
88 }
89 };
90 let cwd = match std::env::current_dir() {
91 Ok(cwd) => cwd,
92 Err(_error) => {
93 let mut diagnostics = diagnostics;
94 write_diagnostic(&mut diagnostics, "unable to resolve cwd");
95 return 1;
96 }
97 };
98 run_cli_at_home_with_terminals(
99 args,
100 input,
101 output,
102 diagnostics,
103 &home,
104 &cwd,
105 io::stdin().is_terminal(),
106 io::stdout().is_terminal(),
107 )
108}
109
110pub fn run_cli_at_home<R, W, E>(
111 args: &[String],
112 input: R,
113 output: W,
114 diagnostics: E,
115 home: &Path,
116 cwd: &Path,
117) -> i32
118where
119 R: BufRead + Send + 'static,
120 W: Write,
121 E: Write,
122{
123 run_cli_at_home_with_terminals(args, input, output, diagnostics, home, cwd, false, false)
126}
127
128#[allow(clippy::too_many_arguments)]
129fn run_cli_at_home_with_terminals<R, W, E>(
130 args: &[String],
131 input: R,
132 output: W,
133 mut diagnostics: E,
134 home: &Path,
135 cwd: &Path,
136 stdin_is_tty: bool,
137 stdout_is_tty: bool,
138) -> i32
139where
140 R: BufRead + Send + 'static,
141 W: Write,
142 E: Write,
143{
144 let options = match parse_args(args) {
145 Ok(options) => options,
146 Err(error) => {
147 let mut diagnostics = diagnostics;
148 write_diagnostic(&mut diagnostics, &error);
149 return 2;
150 }
151 };
152 if options.version {
153 if let Err(error) = write_version(output) {
154 write_diagnostic(
155 &mut diagnostics,
156 &format!("unable to write version: {error}"),
157 );
158 return 1;
159 }
160 return 0;
161 }
162 if let Some(command) = options.command {
163 return run_codex_command(command, home, output, &mut diagnostics);
164 }
165 let mode = match resolve_mode(args, stdin_is_tty, stdout_is_tty) {
166 Ok(mode) => mode,
167 Err(error) => {
168 write_diagnostic(&mut diagnostics, &error);
169 return 2;
170 }
171 };
172
173 if options.list_sessions {
174 let mut protocol = ProtocolWriter::new(output);
175 if let Err(error) = Config::ensure_exists(home) {
176 write_diagnostic(&mut diagnostics, &error.to_string());
177 return 1;
178 }
179 let codex_secret = Config::load_or_create(home)
180 .ok()
181 .and_then(|config| config.resolved_auth().ok())
182 .and_then(|auth| configured_codex_secret(home, auth.provider));
183 return match Session::list_with_secret(home, codex_secret.as_deref()) {
184 Ok(sessions) => {
185 for session in sessions {
186 if let Err(error) = protocol.emit_serializable(&session) {
187 write_diagnostic(
188 &mut diagnostics,
189 &format!("unable to write session metadata: {error}"),
190 );
191 return 1;
192 }
193 }
194 0
195 }
196 Err(error) => {
197 write_diagnostic(&mut diagnostics, &error.to_string());
198 1
199 }
200 };
201 }
202
203 let (session, provider, resumed, attached_agents) = if let Some(id) = options.session.as_deref()
204 {
205 let mut session = match Session::resume(home, id) {
206 Ok(session) => session,
207 Err(error) => {
208 write_diagnostic(&mut diagnostics, &error.to_string());
209 return 1;
210 }
211 };
212 let config = match Config::load_or_create(home) {
213 Ok(config) => config,
214 Err(error) => {
215 write_diagnostic(&mut diagnostics, &error.to_string());
216 return 1;
217 }
218 };
219 let auth = match config.resolved_auth() {
220 Ok(auth) => auth,
221 Err(error) => {
222 write_diagnostic(&mut diagnostics, &error.to_string());
223 return 1;
224 }
225 };
226 if let Some(secret) = configured_codex_secret(home, auth.provider) {
227 session = match Session::resume_with_secret(home, id, Some(&secret)) {
228 Ok(session) => session,
229 Err(error) => {
230 write_diagnostic_safe(&mut diagnostics, &error.to_string(), Some(&secret));
231 return 1;
232 }
233 };
234 }
235 let mut selected = match config.resolved_llm() {
236 Ok(settings) => settings,
237 Err(error) => {
238 write_diagnostic_safe(
239 &mut diagnostics,
240 &error.to_string(),
241 configured_api_key(&config).as_deref(),
242 );
243 return 1;
244 }
245 };
246 apply_auth_to_settings(&mut selected, auth.provider);
247 session.llm.model = selected.model;
248 session.llm.effort = selected.effort;
249 session.llm.api_key_env = selected.api_key_env;
250 let provider = match provider_for_settings(home, &session.llm) {
251 Ok(provider) => provider,
252 Err(error) => {
253 write_diagnostic(&mut diagnostics, &error.to_string());
254 return 1;
255 }
256 };
257 if let Err(error) =
258 session.append_provider_settings(session.llm.model.clone(), session.llm.effort.clone())
259 {
260 write_diagnostic_safe(
261 &mut diagnostics,
262 &error.to_string(),
263 Some(&provider.api_key()),
264 );
265 return 1;
266 }
267 if mode == FrontendMode::Tui && conflicts_with_tui_literal(&provider.api_key()) {
268 write_diagnostic_safe(
269 &mut diagnostics,
270 "API key conflicts with terminal UI literals",
271 Some(&provider.api_key()),
272 );
273 return 1;
274 }
275 (session, provider, true, Vec::new())
276 } else {
277 let config = match Config::load_or_create(home) {
278 Ok(config) => config,
279 Err(error) => {
280 write_diagnostic(&mut diagnostics, &error.to_string());
281 return 1;
282 }
283 };
284 let auth = match config.resolved_auth() {
285 Ok(auth) => auth,
286 Err(error) => {
287 write_diagnostic(&mut diagnostics, &error.to_string());
288 return 1;
289 }
290 };
291 let configured_secret = configured_api_key(&config);
292 let api_key_env = auth.api_key_env.clone();
293 let mut llm = match config.resolved_llm() {
294 Ok(llm) => llm,
295 Err(error) => {
296 write_diagnostic_safe(
297 &mut diagnostics,
298 &error.to_string(),
299 configured_secret.as_deref(),
300 );
301 return 1;
302 }
303 };
304 apply_auth_to_settings(&mut llm, auth.provider);
305 let provider = match provider_for_settings(home, &llm) {
306 Ok(provider) => provider,
307 Err(error) => {
308 write_diagnostic_safe(
309 &mut diagnostics,
310 &error.to_string(),
311 configured_secret.as_deref(),
312 );
313 return 1;
314 }
315 };
316 if mode == FrontendMode::Tui && conflicts_with_tui_literal(&provider.api_key()) {
317 write_diagnostic_safe(
318 &mut diagnostics,
319 "API key conflicts with terminal UI literals",
320 Some(&provider.api_key()),
321 );
322 return 1;
323 }
324 let safe_cwd = match std::fs::canonicalize(cwd) {
325 Ok(cwd) if !cwd.display().to_string().contains(&provider.api_key()) => cwd,
326 Ok(_) => {
327 write_diagnostic_safe(
328 &mut diagnostics,
329 "session header rejected",
330 Some(&provider.api_key()),
331 );
332 return 1;
333 }
334 Err(_) => {
335 write_diagnostic_safe(
336 &mut diagnostics,
337 "unable to resolve session cwd",
338 Some(&provider.api_key()),
339 );
340 return 1;
341 }
342 };
343 let context = match resolve_boot_context_with_api_key_env(
344 home,
345 &safe_cwd,
346 &config.system_prompt,
347 api_key_env.as_deref(),
348 ) {
349 Ok(context) => context,
350 Err(error) => {
351 write_diagnostic_safe(
352 &mut diagnostics,
353 &error.to_string(),
354 configured_secret.as_deref(),
355 );
356 return 1;
357 }
358 };
359 let boot_system_prompt = redact_secret(&context.system_prompt, Some(&provider.api_key()));
360 let attached_agents = attached_agents(context.instruction_files, &provider.api_key());
361 let skills = redact_skills(context.skills, &provider.api_key());
362 let session = match Session::create_with_skills_and_secret(
363 home,
364 &safe_cwd,
365 boot_system_prompt,
366 llm,
367 skills,
368 Some(&provider.api_key()),
369 ) {
370 Ok(session) => session,
371 Err(error) => {
372 write_diagnostic_safe(
373 &mut diagnostics,
374 &error.to_string(),
375 Some(&provider.api_key()),
376 );
377 return 1;
378 }
379 };
380 (session, provider, false, attached_agents)
381 };
382
383 let provider = provider.with_session_id(&session.id);
384 let harness = Harness {
385 home: home.to_path_buf(),
386 session,
387 provider,
388 context_window: None,
389 attached_agents,
390 background_commands: crate::command::BackgroundCommands::default(),
391 };
392 if mode == FrontendMode::Tui {
393 return match crate::tui::run(harness, resumed, output) {
394 Ok(()) => 0,
395 Err(error) => {
396 write_diagnostic(&mut diagnostics, &error);
397 1
398 }
399 };
400 }
401
402 let mut protocol = ProtocolWriter::new(output);
403 let mut harness = harness;
404 if let Err(error) = protocol.session(&harness.session.id, resumed) {
405 write_diagnostic_safe(
406 &mut diagnostics,
407 &format!("unable to write session event: {error}"),
408 Some(harness.provider.api_key().as_str()),
409 );
410 return 1;
411 }
412
413 let (input_tx, input_rx) = mpsc::channel();
414 std::thread::spawn(move || {
415 for line in input.lines() {
416 if input_tx.send(line).is_err() {
417 break;
418 }
419 }
420 });
421 let mut input_closed = false;
422 loop {
423 if harness.has_completed_background_commands() {
424 if let Err(error) = harness.handle_background_completions(&mut protocol, None) {
425 let error = redact_secret(&error, Some(harness.provider.api_key().as_str()));
426 if protocol.error(&error).is_err() {
427 return 1;
428 }
429 }
430 continue;
431 }
432 if input_closed {
433 if harness.has_active_background_commands() {
434 std::thread::sleep(std::time::Duration::from_millis(25));
435 continue;
436 }
437 break;
438 }
439 let line = match input_rx.recv_timeout(std::time::Duration::from_millis(25)) {
440 Ok(Ok(line)) => line,
441 Ok(Err(error)) => {
442 write_diagnostic_safe(
443 &mut diagnostics,
444 &format!("unable to read stdin: {error}"),
445 Some(harness.provider.api_key().as_str()),
446 );
447 return 1;
448 }
449 Err(mpsc::RecvTimeoutError::Timeout) => continue,
450 Err(mpsc::RecvTimeoutError::Disconnected) => {
451 input_closed = true;
452 continue;
453 }
454 };
455 if line.trim().is_empty() {
456 continue;
457 }
458 let text = match parse_input_message(&line) {
459 Ok(text) => text,
460 Err(error) => {
461 let error = redact_secret(&error, Some(harness.provider.api_key().as_str()));
462 if let Err(write_error) = protocol.error(&error) {
463 write_diagnostic_safe(
464 &mut diagnostics,
465 &format!("unable to write protocol error: {write_error}"),
466 Some(harness.provider.api_key().as_str()),
467 );
468 return 1;
469 }
470 continue;
471 }
472 };
473 if let Err(error) = harness.handle_message(&text, &mut protocol, None) {
474 let error = redact_secret(&error, Some(harness.provider.api_key().as_str()));
475 if let Err(write_error) = protocol.error(&error) {
476 write_diagnostic_safe(
477 &mut diagnostics,
478 &format!("unable to write protocol error: {write_error}"),
479 Some(harness.provider.api_key().as_str()),
480 );
481 return 1;
482 }
483 }
484 }
485 0
486}
487
488pub fn resolve_mode(
489 args: &[String],
490 stdin_is_tty: bool,
491 stdout_is_tty: bool,
492) -> Result<FrontendMode, String> {
493 let options = parse_args(args)?;
494 if options.list_sessions {
495 if options.tui {
496 return Err("--tui cannot be combined with --list-sessions".to_owned());
497 }
498 return Ok(FrontendMode::Jsonl);
499 }
500 if options.tui && !(stdin_is_tty && stdout_is_tty) {
501 return Err("--tui requires a terminal on stdin and stdout".to_owned());
502 }
503 if options.tui {
504 Ok(FrontendMode::Tui)
505 } else if options.jsonl || !(stdin_is_tty && stdout_is_tty) {
506 Ok(FrontendMode::Jsonl)
507 } else {
508 Ok(FrontendMode::Tui)
509 }
510}
511
512pub(crate) struct Harness {
513 pub(crate) home: PathBuf,
514 pub(crate) session: Session,
515 pub(crate) provider: Provider,
516 pub(crate) context_window: Option<usize>,
520 pub(crate) attached_agents: Vec<String>,
523 background_commands: crate::command::BackgroundCommands,
524}
525
526fn should_compact_context(context_tokens: usize, context_window: usize) -> bool {
527 context_window > 0
528 && context_tokens as u128 * 100
529 >= context_window as u128 * AUTO_COMPACTION_THRESHOLD_PERCENT as u128
530}
531
532fn find_compaction_boundary(
533 messages: &[ChatMessage],
534 previous_boundary: Option<usize>,
535) -> Option<usize> {
536 let user_starts = messages
537 .iter()
538 .enumerate()
539 .filter_map(|(index, message)| (message.role == "user").then_some(index))
540 .collect::<Vec<_>>();
541 let mut start = *user_starts.last()?;
542 let end = messages.len();
543 let mut kept_tokens = messages[start..end]
544 .iter()
545 .map(estimate_message_tokens)
546 .sum::<usize>();
547
548 while kept_tokens < COMPACTION_KEEP_RECENT_TOKENS {
549 let Some(previous_start) = user_starts
550 .iter()
551 .copied()
552 .rev()
553 .find(|candidate| *candidate < start)
554 else {
555 break;
556 };
557 start = previous_start;
558 kept_tokens = messages[start..end]
559 .iter()
560 .map(estimate_message_tokens)
561 .sum::<usize>();
562 }
563
564 (start > 0 && previous_boundary.is_none_or(|previous| start > previous)).then_some(start)
565}
566
567impl Harness {
568 pub(crate) fn apply_settings(
569 &mut self,
570 home: &Path,
571 model: String,
572 effort: Option<String>,
573 ) -> Result<(), String> {
574 let config = Config::load_or_create(home).map_err(|error| error.to_string())?;
575 let mut settings = config.resolved_llm().map_err(|error| error.to_string())?;
576 settings.model = model.trim().to_owned();
577 settings.effort = effort
578 .map(|value| value.trim().to_owned())
579 .filter(|value| !value.is_empty());
580 settings.base_url = self.session.llm.base_url.clone();
582 settings.api_key_env = self.session.llm.api_key_env.clone();
583 apply_auth_to_settings(&mut settings, auth_provider_for_settings(&self.session.llm));
584 let provider = provider_for_settings(home, &settings)
585 .map_err(|error| error.to_string())?
586 .with_session_id(&self.session.id);
587 Config::save_selection(home, &settings.model, settings.effort.as_deref())
589 .map_err(|error| error.to_string())?;
590 self.session
591 .append_provider_settings(settings.model.clone(), settings.effort.clone())
592 .map_err(|error| error.to_string())?;
593 self.session.llm = settings;
594 self.provider = provider;
595 self.context_window = self.provider.context_window();
596 Ok(())
597 }
598
599 fn should_compact(&self, messages: &[ChatMessage]) -> bool {
600 self.context_window
601 .is_some_and(|window| should_compact_context(estimate_context_tokens(messages), window))
602 }
603
604 fn compaction_boundary(&self) -> Option<usize> {
605 let latest_boundary = self
606 .session
607 .history
608 .iter()
609 .rev()
610 .find_map(|record| match record {
611 crate::session::SessionHistoryRecord::Compaction(compaction) => {
612 Some(compaction.first_kept_message)
613 }
614 _ => None,
615 });
616 find_compaction_boundary(&self.session.messages, latest_boundary)
617 }
618
619 fn compact_context<S: EventSink>(
620 &mut self,
621 sink: &mut S,
622 cancellation: Option<&crate::cancellation::CancellationToken>,
623 tokens_before: usize,
624 ) -> Result<(), String> {
625 let Some(boundary) = self.compaction_boundary() else {
626 return Err("context cannot be compacted without an earlier complete turn".to_owned());
627 };
628 let Some(cancellation) = cancellation else {
629 return Err("context compaction requires a cancellable turn".to_owned());
630 };
631 sink.compaction_started()
632 .map_err(|error| format!("unable to emit compaction state: {error}"))?;
633 let context_messages = self.session.provider_messages();
634 let mut summary_messages = Vec::with_capacity(context_messages.len() + 1);
635 summary_messages.push(ChatMessage::system(self.session.boot_system_prompt.clone()));
636 summary_messages.push(ChatMessage::system(COMPACTION_SYSTEM_PROMPT.to_owned()));
637 summary_messages.extend(context_messages.into_iter().skip(1));
638 let summary = match self.provider.summarize(&summary_messages, cancellation) {
639 Ok(summary) => redact_secret(&summary, Some(self.provider.api_key().as_str())),
640 Err(error) if cancellation.is_cancelled() || error.is_cancelled() => {
641 return self.interrupt(sink, PROVIDER_PHASE, "", &[], Vec::new());
642 }
643 Err(error) => return Err(format!("unable to compact context: {error}")),
644 };
645 self.session
646 .append_compaction(summary, boundary, tokens_before)
647 .map_err(|error| format!("unable to persist context compaction: {error}"))?;
648 let tokens_after = estimate_context_tokens(&self.session.provider_messages());
649 sink.compaction_finished(tokens_before, tokens_after)
650 .map_err(|error| format!("unable to emit compaction state: {error}"))?;
651 Ok(())
652 }
653
654 pub(crate) fn handle_message<S: EventSink>(
655 &mut self,
656 text: &str,
657 sink: &mut S,
658 cancellation: Option<&crate::cancellation::CancellationToken>,
659 ) -> Result<(), String> {
660 if cancellation.is_some_and(CancellationToken::is_cancelled) {
661 return self.interrupt(sink, PROVIDER_PHASE, "", &[], Vec::new());
662 }
663 let secret = self.provider.api_key();
664 let expanded = expand_skill_invocation(text, &self.session.skills)?;
665 let user_message = ChatMessage::user(redact_secret(&expanded.text, Some(&secret)));
666 if let Err(error) = self.session.append_message(user_message) {
667 if cancellation.is_some_and(|token| token.is_cancelled()) {
668 let interruption = self.interrupt(sink, PROVIDER_PHASE, "", &[], Vec::new());
669 return interruption
670 .map_err(|interrupt_error| format!("{error}; {interrupt_error}"));
671 }
672 return Err(error.to_string());
673 }
674 if let Some(name) = expanded.attached_skill.as_deref() {
675 sink.skill_instruction_attached(name)
676 .map_err(|error| format!("unable to emit skill attachment state: {error}"))?;
677 }
678
679 self.continue_turn(sink, cancellation)
680 }
681
682 pub(crate) fn has_active_background_commands(&self) -> bool {
683 self.background_commands.has_active()
684 }
685
686 pub(crate) fn has_completed_background_commands(&self) -> bool {
687 self.background_commands.has_completed()
688 }
689
690 pub(crate) fn handle_background_completions<S: EventSink>(
691 &mut self,
692 sink: &mut S,
693 cancellation: Option<&crate::cancellation::CancellationToken>,
694 ) -> Result<bool, String> {
695 if !self.append_background_completions()? {
696 return Ok(false);
697 }
698 self.continue_turn(sink, cancellation)?;
699 Ok(true)
700 }
701
702 fn append_background_completions(&mut self) -> Result<bool, String> {
703 let completions = self.background_commands.take_completions();
704 if completions.is_empty() {
705 return Ok(false);
706 }
707 for completion in completions {
708 let result = serde_json::json!({
709 "background_id": completion.id,
710 "status": "completed",
711 "result": completion.result,
712 });
713 let content = format!(
714 "Lucy background command completed. Treat this as the automatic result for the previously registered background command:
715{}",
716 serde_json::to_string(&result)
717 .map_err(|error| format!("unable to encode background cmd result: {error}"))?
718 );
719 self.session
720 .append_message(ChatMessage::system(content))
721 .map_err(|error| error.to_string())?;
722 }
723 Ok(true)
724 }
725
726 fn continue_turn<S: EventSink>(
727 &mut self,
728 sink: &mut S,
729 cancellation: Option<&crate::cancellation::CancellationToken>,
730 ) -> Result<(), String> {
731 let secret = self.provider.api_key();
732 let mut compacted_for_turn = false;
733 loop {
734 self.append_background_completions()?;
735 if cancellation.is_some_and(CancellationToken::is_cancelled) {
736 return self.interrupt(sink, PROVIDER_PHASE, "", &[], Vec::new());
737 }
738 let mut messages = self.session.provider_messages();
739 let tokens_before = estimate_context_tokens(&messages);
740 if !compacted_for_turn && self.should_compact(&messages) {
741 self.compact_context(sink, cancellation, tokens_before)?;
742 compacted_for_turn = true;
743 messages = self.session.provider_messages();
744 }
745 sink.context_usage(estimate_context_tokens(&messages))
746 .map_err(|error| format!("unable to emit context usage: {error}"))?;
747 let mut raw_content = String::new();
748 let mut redactor = SecretRedactor::new(&secret);
749 let mut reasoning_active = false;
750 let stream_result = {
751 let mut on_event = |event: ProviderStreamEvent| -> io::Result<()> {
752 match event {
753 ProviderStreamEvent::ReasoningStarted => {
754 if !reasoning_active {
755 reasoning_active = true;
756 sink.reasoning_started()?;
757 }
758 Ok(())
759 }
760 ProviderStreamEvent::Text(delta) => {
761 if reasoning_active {
762 reasoning_active = false;
763 sink.reasoning_completed()?;
764 }
765 raw_content.push_str(&delta);
766 redactor.push(&delta, |safe_delta| {
767 sink.emit_event(&ProtocolEvent::AssistantDelta {
768 text: safe_delta.to_owned(),
769 })
770 })
771 }
772 }
773 };
774 match cancellation {
775 Some(token) => self
776 .provider
777 .stream_chat_cancellable_with_options_and_events(
778 &messages,
779 &mut on_event,
780 token,
781 true,
782 ),
783 None => self.provider.stream_chat(&messages, &mut |delta| {
784 raw_content.push_str(delta);
785 redactor.push(delta, |safe_delta| {
786 sink.emit_event(&ProtocolEvent::AssistantDelta {
787 text: safe_delta.to_owned(),
788 })
789 })
790 }),
791 }
792 };
793 redactor
794 .finish(|safe_delta| {
795 sink.emit_event(&ProtocolEvent::AssistantDelta {
796 text: safe_delta.to_owned(),
797 })
798 })
799 .map_err(|error| format!("unable to write assistant delta: {error}"))?;
800 let turn = match stream_result {
801 Ok(turn) => {
802 if reasoning_active {
803 sink.reasoning_completed()
804 .map_err(|error| format!("unable to emit reasoning state: {error}"))?;
805 }
806 turn
807 }
808 Err(error)
809 if cancellation.is_some_and(|token| token.is_cancelled())
810 || error.is_cancelled() =>
811 {
812 if reasoning_active {
813 sink.reasoning_completed()
814 .map_err(|error| format!("unable to emit reasoning state: {error}"))?;
815 }
816 let partial = error.partial_turn().cloned().unwrap_or(ProviderTurn {
817 content: raw_content,
818 tool_calls: Vec::new(),
819 reasoning_details: Vec::new(),
820 });
821 return self.interrupt(
822 sink,
823 PROVIDER_PHASE,
824 &partial.content,
825 &partial.tool_calls,
826 Vec::new(),
827 );
828 }
829 Err(error) => {
830 if reasoning_active {
831 sink.reasoning_completed()
832 .map_err(|error| format!("unable to emit reasoning state: {error}"))?;
833 }
834 return Err(error.to_string());
835 }
836 };
837 let canceled_after_stream = cancellation.is_some_and(|token| token.is_cancelled());
838
839 if turn
840 .tool_calls
841 .iter()
842 .any(|call| !matches!(call.name.as_str(), "cmd"))
843 {
844 if canceled_after_stream {
845 return self.interrupt(sink, PROVIDER_PHASE, &turn.content, &[], Vec::new());
846 }
847 return Err("provider requested an unsupported tool".to_owned());
848 }
849 let safe_tool_calls = turn
850 .tool_calls
851 .iter()
852 .map(|call| safe_tool_call(call, &secret))
853 .collect::<Vec<_>>();
854 let assistant_content = redact_secret(&turn.content, Some(&secret));
855 let safe_reasoning_details = redact_reasoning_details(&turn.reasoning_details, &secret);
856 let mut assistant =
857 ChatMessage::assistant(assistant_content.clone(), safe_tool_calls.clone());
858 assistant.reasoning_details = safe_reasoning_details;
859 if let Err(error) = self.session.append_message(assistant) {
860 if cancellation.is_some_and(|token| token.is_cancelled()) {
861 let interruption = self.interrupt(
862 sink,
863 PROVIDER_PHASE,
864 &assistant_content,
865 &turn.tool_calls,
866 Vec::new(),
867 );
868 return interruption
869 .map_err(|interrupt_error| format!("{error}; {interrupt_error}"));
870 }
871 return Err(error.to_string());
872 }
873
874 if safe_tool_calls.is_empty() {
875 if canceled_after_stream
876 || cancellation.is_some_and(CancellationToken::is_cancelled)
877 {
878 return self.interrupt(sink, PROVIDER_PHASE, "", &[], Vec::new());
879 }
880 if self.append_background_completions()? {
881 continue;
882 }
883 if cancellation.is_some_and(|token| !token.try_complete()) {
884 return self.interrupt(sink, PROVIDER_PHASE, "", &[], Vec::new());
885 }
886 sink.context_usage(estimate_context_tokens(&self.session.provider_messages()))
887 .map_err(|error| format!("unable to emit context usage: {error}"))?;
888 sink.emit_event(&ProtocolEvent::TurnEnd)
889 .map_err(|error| format!("unable to write turn end: {error}"))?;
890 return Ok(());
891 }
892
893 for safe_call in &safe_tool_calls {
894 sink.emit_event(&ProtocolEvent::ToolCall {
895 id: safe_call.id.clone(),
896 name: safe_call.name.clone(),
897 arguments: safe_call.arguments.clone(),
898 })
899 .map_err(|error| format!("unable to write tool call: {error}"))?;
900 }
901 for (index, raw_call) in turn.tool_calls.iter().enumerate() {
902 let safe_call = &safe_tool_calls[index];
903 let result = if cancellation.is_some_and(|token| token.is_cancelled()) {
904 serde_json::to_value(crate::command::canceled_result(
905 &safe_call.arguments,
906 &secret,
907 ))
908 .map_err(|error| format!("unable to encode cmd result: {error}"))?
909 } else {
910 crate::command::execute_managed(
911 &raw_call.arguments,
912 &self.session.cwd,
913 self.provider.api_key_env(),
914 Some(&secret),
915 cancellation,
916 &mut self.background_commands,
917 )
918 };
919 let result = redact_json_value(result, &secret);
920 let tool_content = serde_json::to_string(&result)
921 .map_err(|error| format!("unable to encode tool result: {error}"))?;
922 let tool_message = ChatMessage::tool(
923 safe_call.id.clone(),
924 safe_call.name.clone(),
925 redact_secret(&tool_content, Some(&secret)),
926 );
927 let observation = crate::session::SessionToolResult {
928 id: safe_call.id.clone(),
929 name: safe_call.name.clone(),
930 result: result.clone(),
931 };
932 if let Err(error) = self.session.append_message(tool_message) {
933 if cancellation.is_some_and(|token| token.is_cancelled()) {
934 let interruption =
935 self.interrupt(sink, COMMAND_PHASE, "", &[], vec![observation]);
936 return interruption
937 .map_err(|interrupt_error| format!("{error}; {interrupt_error}"));
938 }
939 return Err(error.to_string());
940 }
941 sink.emit_event(&ProtocolEvent::ToolResult {
942 id: safe_call.id.clone(),
943 name: safe_call.name.clone(),
944 result: result.clone(),
945 })
946 .map_err(|error| format!("unable to write tool result: {error}"))?;
947 if cancellation.is_some_and(|token| token.is_cancelled()) {
948 for pending_call in safe_tool_calls.iter().skip(index + 1) {
949 let pending_result = redact_json_value(
950 serde_json::to_value(crate::command::canceled_result(
951 &pending_call.arguments,
952 &secret,
953 ))
954 .map_err(|error| format!("unable to encode cmd result: {error}"))?,
955 &secret,
956 );
957 let pending_content = serde_json::to_string(&pending_result)
958 .map_err(|error| format!("unable to encode tool result: {error}"))?;
959 let pending_message = ChatMessage::tool(
960 pending_call.id.clone(),
961 pending_call.name.clone(),
962 redact_secret(&pending_content, Some(&secret)),
963 );
964 let pending_observation = crate::session::SessionToolResult {
965 id: pending_call.id.clone(),
966 name: pending_call.name.clone(),
967 result: pending_result.clone(),
968 };
969 if let Err(error) = self.session.append_message(pending_message) {
970 if cancellation.is_some_and(|token| token.is_cancelled()) {
971 let interruption = self.interrupt(
972 sink,
973 COMMAND_PHASE,
974 "",
975 &[],
976 vec![pending_observation],
977 );
978 return interruption.map_err(|interrupt_error| {
979 format!("{error}; {interrupt_error}")
980 });
981 }
982 return Err(error.to_string());
983 }
984 sink.emit_event(&ProtocolEvent::ToolResult {
985 id: pending_call.id.clone(),
986 name: pending_call.name.clone(),
987 result: pending_result.clone(),
988 })
989 .map_err(|error| format!("unable to write tool result: {error}"))?;
990 }
991 return self.interrupt(sink, COMMAND_PHASE, "", &[], Vec::new());
992 }
993 }
994 if cancellation.is_some_and(CancellationToken::is_cancelled) {
995 return self.interrupt(sink, COMMAND_PHASE, "", &[], Vec::new());
996 }
997 }
998 }
999
1000 fn interrupt<S: EventSink>(
1001 &mut self,
1002 sink: &mut S,
1003 phase: &str,
1004 assistant_text: &str,
1005 tool_calls: &[ChatToolCall],
1006 tool_results: Vec<crate::session::SessionToolResult>,
1007 ) -> Result<(), String> {
1008 let secret = self.provider.api_key();
1009 let safe_tool_calls = tool_calls
1010 .iter()
1011 .filter(|call| call.name == "cmd")
1012 .map(|call| safe_partial_tool_call(call, &secret))
1013 .collect::<Vec<_>>();
1014 let safe_tool_results = tool_results.clone();
1015 let interruption = crate::session::InterruptionRecord {
1016 timestamp: 0,
1017 reason: USER_CANCEL_REASON.to_owned(),
1018 phase: phase.to_owned(),
1019 assistant_text: redact_secret(assistant_text, Some(&secret)),
1020 tool_calls: safe_tool_calls.clone(),
1021 tool_results,
1022 };
1023 let persistence_error = self.session.append_interruption(interruption).err();
1024 let mut event_error = None;
1025 for call in &safe_tool_calls {
1026 if let Err(error) = sink.emit_event(&ProtocolEvent::ToolCall {
1027 id: call.id.clone(),
1028 name: call.name.clone(),
1029 arguments: call.arguments.clone(),
1030 }) {
1031 event_error.get_or_insert(error);
1032 }
1033 }
1034 for observation in &safe_tool_results {
1035 if let Err(error) = sink.emit_event(&ProtocolEvent::ToolResult {
1036 id: observation.id.clone(),
1037 name: observation.name.clone(),
1038 result: observation.result.clone(),
1039 }) {
1040 event_error.get_or_insert(error);
1041 }
1042 }
1043 if let Err(error) = sink.emit_event(&ProtocolEvent::TurnInterrupted {
1044 reason: USER_CANCEL_REASON.to_owned(),
1045 phase: phase.to_owned(),
1046 }) {
1047 event_error.get_or_insert(error);
1048 }
1049 match (persistence_error, event_error) {
1050 (None, None) => Ok(()),
1051 (Some(error), None) => Err(format!("unable to persist interruption: {error}")),
1052 (None, Some(error)) => Err(format!("unable to write interruption event: {error}")),
1053 (Some(persistence), Some(event)) => Err(format!(
1054 "unable to persist interruption: {persistence}; unable to write interruption event: {event}"
1055 )),
1056 }
1057 }
1058}
1059
1060struct SecretRedactor {
1061 secret_text: String,
1062 secret: Vec<char>,
1063 marker: String,
1064 pending: String,
1065}
1066
1067impl SecretRedactor {
1068 fn new(secret: &str) -> Self {
1069 Self {
1070 secret_text: secret.to_owned(),
1071 secret: secret.chars().collect(),
1072 marker: redaction_marker(secret).unwrap_or_default(),
1073 pending: String::new(),
1074 }
1075 }
1076
1077 fn push<F>(&mut self, text: &str, mut emit: F) -> io::Result<()>
1078 where
1079 F: FnMut(&str) -> io::Result<()>,
1080 {
1081 if self.secret.is_empty() {
1082 return emit(text);
1083 }
1084
1085 let mut output = String::new();
1086 for character in text.chars() {
1087 self.pending.push(character);
1088 if self.pending.chars().eq(self.secret.iter().copied()) {
1089 self.pending.clear();
1090 output.push_str(&self.marker);
1091 continue;
1092 }
1093 if self.pending_is_secret_prefix() {
1094 continue;
1095 }
1096
1097 let pending = self.pending.chars().collect::<Vec<_>>();
1098 let suffix_len = (1..pending.len())
1099 .rev()
1100 .find(|length| {
1101 pending[pending.len() - length..].iter().copied().eq(self
1102 .secret
1103 .iter()
1104 .copied()
1105 .take(*length))
1106 })
1107 .unwrap_or(0);
1108 let safe_len = pending.len() - suffix_len;
1109 output.extend(pending[..safe_len].iter());
1110 self.pending = pending[safe_len..].iter().collect();
1111 }
1112
1113 if output.is_empty() {
1114 Ok(())
1115 } else {
1116 let safe_output = redact_secret(&output, Some(&self.secret_text));
1117 emit(&safe_output)
1118 }
1119 }
1120
1121 fn finish<F>(&mut self, mut emit: F) -> io::Result<()>
1122 where
1123 F: FnMut(&str) -> io::Result<()>,
1124 {
1125 let pending = std::mem::take(&mut self.pending);
1126 if pending.is_empty() {
1127 return Ok(());
1128 }
1129 let safe_pending = redact_secret(&pending, Some(&self.secret_text));
1130 emit(&safe_pending)
1131 }
1132
1133 fn pending_is_secret_prefix(&self) -> bool {
1134 let length = self.pending.chars().count();
1135 length < self.secret.len()
1136 && self
1137 .pending
1138 .chars()
1139 .zip(self.secret.iter().copied())
1140 .all(|(pending, secret)| pending == secret)
1141 }
1142}
1143
1144fn attached_agents(instruction_files: Vec<InstructionSource>, secret: &str) -> Vec<String> {
1147 instruction_files
1148 .into_iter()
1149 .filter(|source| {
1150 source
1151 .path
1152 .file_name()
1153 .is_some_and(|name| name == "AGENTS.md")
1154 })
1155 .map(|source| redact_secret(&source.path.display().to_string(), Some(secret)))
1156 .collect()
1157}
1158
1159fn escape_xml_attribute(text: &str) -> String {
1162 text.replace('&', "&")
1163 .replace('<', "<")
1164 .replace('>', ">")
1165 .replace('\"', """)
1166 .replace('\'', "'")
1167}
1168
1169fn redact_skills(skills: Vec<SkillEntry>, secret: &str) -> Vec<SkillEntry> {
1170 skills
1171 .into_iter()
1172 .map(|skill| SkillEntry {
1173 name: redact_secret(&skill.name, Some(secret)),
1174 description: redact_secret(&skill.description, Some(secret)),
1175 path: std::path::PathBuf::from(redact_secret(
1176 &skill.path.display().to_string(),
1177 Some(secret),
1178 )),
1179 contents: redact_secret(&skill.contents, Some(secret)),
1180 model_invocable: skill.model_invocable,
1181 })
1182 .collect()
1183}
1184
1185#[derive(Debug)]
1188struct ExpandedSkillInvocation {
1189 text: String,
1190 attached_skill: Option<String>,
1191}
1192
1193fn expand_skill_invocation(
1197 text: &str,
1198 skills: &[SkillEntry],
1199) -> Result<ExpandedSkillInvocation, String> {
1200 let Some(invocation) = text.strip_prefix('/') else {
1201 return Ok(ExpandedSkillInvocation {
1202 text: text.to_owned(),
1203 attached_skill: None,
1204 });
1205 };
1206 let mut pieces = invocation.splitn(2, char::is_whitespace);
1207 let name = pieces.next().unwrap_or_default();
1208 if name.is_empty() {
1209 return Err("skill command requires a skill name: /<name> [args]".to_owned());
1210 }
1211 let Some(skill) = skills.iter().find(|skill| skill.name == name) else {
1212 return Err(format!("unknown skill: {name}"));
1213 };
1214 let arguments = pieces.next().unwrap_or_default().trim();
1215 let mut message = format!(
1216 "<skill name=\"{}\" location=\"{}\">\n{}\n</skill>",
1217 escape_xml_attribute(&skill.name),
1218 escape_xml_attribute(&skill.path.display().to_string()),
1219 skill.contents.trim()
1220 );
1221 if !arguments.is_empty() {
1222 message.push_str("\n\nUser: ");
1223 message.push_str(arguments);
1224 }
1225 Ok(ExpandedSkillInvocation {
1226 text: message,
1227 attached_skill: Some(skill.name.clone()),
1228 })
1229}
1230
1231#[cfg(test)]
1232fn redact_tool_arguments(arguments: &str, secret: &str) -> String {
1233 safe_tool_call(
1234 &ChatToolCall {
1235 id: String::new(),
1236 name: "cmd".to_owned(),
1237 arguments: arguments.to_owned(),
1238 },
1239 secret,
1240 )
1241 .arguments
1242}
1243
1244fn safe_tool_call(call: &ChatToolCall, secret: &str) -> ChatToolCall {
1245 let valid = match call.name.as_str() {
1246 "cmd" => serde_json::from_str::<Value>(&call.arguments)
1247 .ok()
1248 .and_then(|value| value.as_object().cloned())
1249 .is_some_and(|object| {
1250 (object.len() == 1 || object.len() == 2)
1251 && object.get("command").is_some_and(Value::is_string)
1252 && object.get("background").is_none_or(Value::is_boolean)
1253 && object
1254 .keys()
1255 .all(|key| matches!(key.as_str(), "command" | "background"))
1256 }),
1257 _ => false,
1258 };
1259 let arguments = if valid {
1260 serde_json::to_string(&redact_json_value(
1261 serde_json::from_str(&call.arguments).unwrap_or(Value::Null),
1262 secret,
1263 ))
1264 .unwrap_or_else(|_| "{}".to_owned())
1265 } else {
1266 "{}".to_owned()
1267 };
1268 ChatToolCall {
1269 id: redact_secret(&call.id, Some(secret)),
1270 name: redact_secret(&call.name, Some(secret)),
1271 arguments,
1272 }
1273}
1274
1275fn safe_partial_tool_call(call: &ChatToolCall, secret: &str) -> ChatToolCall {
1276 let arguments = if serde_json::from_str::<Value>(&call.arguments)
1277 .ok()
1278 .and_then(|value| value.as_object().cloned())
1279 .is_some_and(|object| {
1280 (object.len() == 1 || object.len() == 2)
1281 && object.contains_key("command")
1282 && object
1283 .keys()
1284 .all(|key| matches!(key.as_str(), "command" | "background"))
1285 }) {
1286 safe_tool_call(call, secret).arguments
1287 } else {
1288 "{}".to_owned()
1292 };
1293 ChatToolCall {
1294 id: redact_secret(&call.id, Some(secret)),
1295 name: redact_secret(&call.name, Some(secret)),
1296 arguments,
1297 }
1298}
1299
1300fn redact_json_value(value: Value, secret: &str) -> Value {
1301 match value {
1302 Value::String(text) => Value::String(redact_secret(&text, Some(secret))),
1303 Value::Array(values) => Value::Array(
1304 values
1305 .into_iter()
1306 .map(|value| redact_json_value(value, secret))
1307 .collect(),
1308 ),
1309 Value::Object(object) => {
1310 let marker = redaction_marker(secret).unwrap_or_default();
1311 let mut redacted = Map::new();
1312 for (key, value) in object {
1313 let mut safe_key = if is_structural_key(&key) {
1314 key
1315 } else {
1316 redact_secret(&key, Some(secret))
1317 };
1318 if redacted.contains_key(&safe_key) {
1319 if marker.is_empty() {
1320 continue;
1321 }
1322 while redacted.contains_key(&safe_key) {
1323 safe_key.push_str(&marker);
1324 }
1325 }
1326 redacted.insert(safe_key, redact_json_value(value, secret));
1327 }
1328 Value::Object(redacted)
1329 }
1330 value => value,
1331 }
1332}
1333
1334fn redact_reasoning_details(details: &[Value], secret: &str) -> Option<Vec<Value>> {
1335 if details.is_empty() {
1336 return None;
1337 }
1338 match redact_json_value(Value::Array(details.to_vec()), secret) {
1339 Value::Array(details) => Some(details),
1340 _ => None,
1341 }
1342}
1343
1344fn write_version<W: Write>(mut output: W) -> io::Result<()> {
1345 writeln!(output, "lucy {}", env!("CARGO_PKG_VERSION"))
1346}
1347
1348fn parse_args(args: &[String]) -> Result<CliOptions, String> {
1349 let mut options = CliOptions {
1350 session: None,
1351 list_sessions: false,
1352 jsonl: false,
1353 tui: false,
1354 version: false,
1355 command: None,
1356 };
1357 if args.len() == 2 && args[0] == "codex" {
1358 options.command = Some(match args[1].as_str() {
1359 "login" => CliCommand::CodexLogin,
1360 "logout" => CliCommand::CodexLogout,
1361 _ => return Err("usage: lucy codex <login|logout>".to_owned()),
1362 });
1363 return Ok(options);
1364 }
1365 if args.first().is_some_and(|arg| arg == "codex") {
1366 return Err("usage: lucy codex <login|logout>".to_owned());
1367 }
1368 let mut index = 0;
1369 while index < args.len() {
1370 match args[index].as_str() {
1371 "--session" => {
1372 if options.list_sessions || options.session.is_some() {
1373 return Err("--session cannot be combined or repeated".to_owned());
1374 }
1375 index += 1;
1376 let Some(id) = args.get(index) else {
1377 return Err("--session requires an id".to_owned());
1378 };
1379 options.session = Some(id.clone());
1380 }
1381 "--list-sessions" => {
1382 if options.session.is_some() || options.list_sessions {
1383 return Err("--list-sessions cannot be combined or repeated".to_owned());
1384 }
1385 options.list_sessions = true;
1386 }
1387 "--jsonl" => {
1388 if options.jsonl || options.tui {
1389 return Err("--jsonl cannot be combined or repeated".to_owned());
1390 }
1391 options.jsonl = true;
1392 }
1393 "--tui" => {
1394 if options.tui || options.jsonl {
1395 return Err("--tui cannot be combined or repeated".to_owned());
1396 }
1397 options.tui = true;
1398 }
1399 "--version" => {
1400 if options.version {
1401 return Err("--version cannot be repeated".to_owned());
1402 }
1403 options.version = true;
1404 }
1405 "--help" | "-h" => {
1406 return Err(
1407 "usage: lucy [--version] [--jsonl|--tui] [--session <id>] [--list-sessions] | lucy codex <login|logout>"
1408 .to_owned(),
1409 );
1410 }
1411 _ => return Err("unknown argument".to_owned()),
1412 }
1413 index += 1;
1414 }
1415 Ok(options)
1416}
1417
1418fn parse_input_message(line: &str) -> Result<String, String> {
1419 let record: InputRecord = serde_json::from_str(line)
1420 .map_err(|_| "input must be a JSONL message record".to_owned())?;
1421 if record.record_type != "message" {
1422 return Err("input record type must be message".to_owned());
1423 }
1424 record
1425 .text
1426 .ok_or_else(|| "message record requires a text string".to_owned())
1427}
1428
1429fn home_directory() -> Result<PathBuf, String> {
1430 std::env::var_os("HOME")
1431 .map(PathBuf::from)
1432 .ok_or_else(|| "HOME is not set; Lucy needs a user home directory".to_owned())
1433}
1434
1435fn configured_api_key_env(config: &Config) -> Option<String> {
1436 config.resolved_auth().ok()?.api_key_env
1437}
1438
1439fn configured_api_key(config: &Config) -> Option<String> {
1440 configured_api_key_env(config)
1441 .and_then(|api_key_env| std::env::var(api_key_env).ok())
1442 .filter(|secret| !secret.is_empty())
1443}
1444
1445fn run_codex_command<W: Write, E: Write>(
1446 command: CliCommand,
1447 home: &Path,
1448 mut output: W,
1449 diagnostics: &mut E,
1450) -> i32 {
1451 match command {
1452 CliCommand::CodexLogin => match crate::auth::login(home) {
1453 Ok(_) => {
1454 let _ = writeln!(output, "Codex login successful");
1455 0
1456 }
1457 Err(error) => {
1458 write_diagnostic(diagnostics, &error.to_string());
1459 1
1460 }
1461 },
1462 CliCommand::CodexLogout => match crate::auth::AuthStore::for_home(home).logout() {
1463 Ok(true) => {
1464 let _ = writeln!(output, "Codex logout successful");
1465 0
1466 }
1467 Ok(false) => {
1468 let _ = writeln!(output, "Codex was not logged in");
1469 0
1470 }
1471 Err(error) => {
1472 write_diagnostic(diagnostics, &error.to_string());
1473 1
1474 }
1475 },
1476 }
1477}
1478
1479fn apply_auth_to_settings(settings: &mut LlmSettings, provider: AuthProvider) {
1480 if provider == AuthProvider::CodexSubscription {
1481 settings.api_key_env = crate::codex_provider::CODEX_ENV_SENTINEL.to_owned();
1482 }
1483}
1484
1485fn auth_provider_for_settings(settings: &LlmSettings) -> AuthProvider {
1486 if settings.api_key_env == crate::codex_provider::CODEX_ENV_SENTINEL {
1487 AuthProvider::CodexSubscription
1488 } else {
1489 AuthProvider::Openrouter
1490 }
1491}
1492
1493fn provider_for_settings(
1494 home: &Path,
1495 settings: &LlmSettings,
1496) -> Result<Provider, crate::provider::ProviderError> {
1497 match auth_provider_for_settings(settings) {
1498 AuthProvider::CodexSubscription => Provider::new_codex(home, settings),
1499 AuthProvider::Openrouter => Provider::new(settings),
1500 }
1501}
1502
1503fn configured_codex_secret(home: &Path, provider: AuthProvider) -> Option<String> {
1504 if provider != AuthProvider::CodexSubscription {
1505 return None;
1506 }
1507 crate::auth::AuthStore::for_home(home)
1508 .load()
1509 .ok()
1510 .flatten()
1511 .map(|credentials| credentials.access)
1512 .filter(|secret| !secret.is_empty())
1513}
1514
1515fn write_diagnostic_safe<W: Write>(diagnostics: &mut W, message: &str, secret: Option<&str>) {
1516 write_diagnostic_safe_with_environment(
1517 diagnostics,
1518 message,
1519 secret,
1520 std::env::vars().map(|(_, value)| value),
1521 );
1522}
1523
1524fn write_diagnostic_safe_with_environment<W, I>(
1525 diagnostics: &mut W,
1526 message: &str,
1527 secret: Option<&str>,
1528 environment_values: I,
1529) where
1530 W: Write,
1531 I: IntoIterator<Item = String>,
1532{
1533 let mut safe_line = format!("!: {message}");
1534 safe_line = redact_secret(&safe_line, secret);
1535 let mut environment_secrets = environment_values
1536 .into_iter()
1537 .filter(|value| !value.is_empty() && !conflicts_with_protected_literal(value))
1538 .collect::<Vec<_>>();
1539 environment_secrets.sort_by_key(|value| std::cmp::Reverse(value.len()));
1540 for environment_secret in environment_secrets {
1541 safe_line = redact_secret(&safe_line, Some(&environment_secret));
1542 }
1543 let _ = writeln!(diagnostics, "{safe_line}");
1544}
1545
1546fn write_diagnostic<W: Write>(diagnostics: &mut W, message: &str) {
1547 write_diagnostic_safe(diagnostics, message, None);
1548}
1549
1550#[cfg(test)]
1551mod tests {
1552 use super::*;
1553 use crate::cancellation::CancellationToken;
1554 use std::io::{Cursor, Read, Write};
1555 use std::net::TcpListener;
1556 use std::thread;
1557
1558 #[test]
1559 fn codex_subcommands_parse_without_entering_a_session() {
1560 assert_eq!(
1561 parse_args(&["codex".to_owned(), "login".to_owned()])
1562 .expect("codex login")
1563 .command,
1564 Some(CliCommand::CodexLogin)
1565 );
1566 assert_eq!(
1567 parse_args(&["codex".to_owned(), "logout".to_owned()])
1568 .expect("codex logout")
1569 .command,
1570 Some(CliCommand::CodexLogout)
1571 );
1572 assert_eq!(
1573 parse_args(&["codex".to_owned(), "status".to_owned()])
1574 .expect_err("unknown codex command"),
1575 "usage: lucy codex <login|logout>"
1576 );
1577 }
1578
1579 #[test]
1580 fn codex_logout_is_idempotent_and_does_not_bootstrap_a_session() {
1581 let home = std::env::temp_dir().join(format!("lucy-codex-logout-{}", std::process::id()));
1582 let _ = std::fs::remove_dir_all(&home);
1583 let cwd = std::env::current_dir().expect("cwd");
1584 let mut output = Vec::new();
1585 let mut diagnostics = Vec::new();
1586 let exit = run_cli_at_home(
1587 &["codex".to_owned(), "logout".to_owned()],
1588 Cursor::new(Vec::<u8>::new()),
1589 &mut output,
1590 &mut diagnostics,
1591 &home,
1592 &cwd,
1593 );
1594 assert_eq!(exit, 0);
1595 assert!(String::from_utf8_lossy(&output).contains("not logged in"));
1596 assert!(diagnostics.is_empty());
1597 assert!(!home.exists());
1598 }
1599
1600 #[test]
1601 fn auto_compaction_triggers_at_or_above_ninety_five_percent_only() {
1602 assert!(!should_compact_context(94, 100));
1603 assert!(should_compact_context(95, 100));
1604 assert!(should_compact_context(96, 100));
1605 assert!(!should_compact_context(100, 0));
1606 }
1607
1608 #[test]
1609 fn compaction_boundary_keeps_complete_recent_turns() {
1610 let messages = [
1611 ChatMessage::user("old request".to_owned()),
1612 ChatMessage::assistant("old answer".to_owned(), Vec::new()),
1613 ChatMessage::user("recent request".to_owned()),
1614 ChatMessage::assistant("recent answer ".repeat(8_000), Vec::new()),
1615 ];
1616
1617 assert_eq!(find_compaction_boundary(&messages, None), Some(2));
1618 assert_eq!(find_compaction_boundary(&messages, Some(2)), None);
1619 }
1620
1621 #[test]
1622 fn mid_turn_compaction_summarizes_without_tools_then_continues_original_request() {
1623 let listener = TcpListener::bind(("127.0.0.1", 0)).expect("compaction listener");
1624 let address = listener.local_addr().expect("compaction address");
1625 let responses = ["summary", "continued"];
1626 let server = thread::spawn(move || {
1627 let mut requests = Vec::new();
1628 for response_text in responses {
1629 let (mut stream, _) = listener.accept().expect("compaction request");
1630 let mut request = String::new();
1631 let mut reader = std::io::BufReader::new(stream.try_clone().expect("clone"));
1632 let mut content_length = 0usize;
1633 loop {
1634 let mut line = String::new();
1635 reader.read_line(&mut line).expect("request header");
1636 if line == "\r\n" {
1637 break;
1638 }
1639 if let Some((name, value)) = line.split_once(':') {
1640 if name.eq_ignore_ascii_case("content-length") {
1641 content_length = value.trim().parse().expect("content length");
1642 }
1643 }
1644 }
1645 let mut body = vec![0u8; content_length];
1646 reader.read_exact(&mut body).expect("request body");
1647 request.push_str(std::str::from_utf8(&body).expect("request JSON"));
1648 requests.push(serde_json::from_str::<Value>(&request).expect("request value"));
1649 let payload = serde_json::json!({
1650 "choices": [{
1651 "delta": {"content": response_text},
1652 "finish_reason": null
1653 }]
1654 });
1655 let body = format!("data: {payload}\n\ndata: [DONE]\n\n");
1656 let header = format!(
1657 "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
1658 body.len()
1659 );
1660 stream
1661 .write_all(header.as_bytes())
1662 .expect("response header");
1663 stream.write_all(body.as_bytes()).expect("response body");
1664 stream.flush().expect("response flush");
1665 }
1666 requests
1667 });
1668
1669 let key_env = format!("LUCY_COMPACTION_APP_KEY_{}", std::process::id());
1670 std::env::set_var(&key_env, "provider-secret");
1671 let settings = crate::config::LlmSettings {
1672 base_url: format!("http://{address}/v1"),
1673 model: "model".to_owned(),
1674 api_key_env: key_env.clone(),
1675 effort: None,
1676 };
1677 let provider = Provider::new(&settings).expect("provider");
1678 let home = std::env::temp_dir().join(format!("lucy-app-compaction-{}", std::process::id()));
1679 let _ = std::fs::remove_dir_all(&home);
1680 std::fs::create_dir(&home).expect("temp home");
1681 let cwd = std::env::current_dir().expect("cwd");
1682 let mut session = Session::create_with_secret(
1683 &home,
1684 &cwd,
1685 "prompt".to_owned(),
1686 settings,
1687 Some("provider-secret"),
1688 )
1689 .expect("session");
1690 session
1691 .append_message(ChatMessage::user("old request".to_owned()))
1692 .expect("old user");
1693 session
1694 .append_message(ChatMessage::assistant("old answer".to_owned(), Vec::new()))
1695 .expect("old answer");
1696 session
1697 .append_message(ChatMessage::user("recent request".to_owned()))
1698 .expect("recent user");
1699 session
1700 .append_message(ChatMessage::assistant(
1701 "recent answer ".repeat(8_000),
1702 Vec::new(),
1703 ))
1704 .expect("recent answer");
1705
1706 struct Sink {
1707 events: Vec<ProtocolEvent>,
1708 compaction_started: bool,
1709 compaction_finished: bool,
1710 }
1711 impl EventSink for Sink {
1712 fn emit_event(&mut self, event: &ProtocolEvent) -> io::Result<()> {
1713 self.events.push(event.clone());
1714 Ok(())
1715 }
1716 fn compaction_started(&mut self) -> io::Result<()> {
1717 self.compaction_started = true;
1718 Ok(())
1719 }
1720 fn compaction_finished(&mut self, _: usize, _: usize) -> io::Result<()> {
1721 self.compaction_finished = true;
1722 Ok(())
1723 }
1724 }
1725
1726 let provider = provider.with_session_id(&session.id);
1727 let mut harness = Harness {
1728 home: std::env::temp_dir(),
1729 session,
1730 provider,
1731 context_window: Some(1),
1732 attached_agents: Vec::new(),
1733 background_commands: crate::command::BackgroundCommands::default(),
1734 };
1735 let cancellation = CancellationToken::new();
1736 let mut sink = Sink {
1737 events: Vec::new(),
1738 compaction_started: false,
1739 compaction_finished: false,
1740 };
1741 harness
1742 .handle_message("continue", &mut sink, Some(&cancellation))
1743 .expect("continued turn");
1744
1745 let requests = server.join().expect("server");
1746 assert_eq!(requests.len(), 2);
1747 assert!(requests[0].get("tools").is_none());
1748 assert!(requests[1].get("tools").is_some());
1749 assert!(requests
1751 .iter()
1752 .all(|request| request.get("session_id").is_none()));
1753 assert!(sink.compaction_started);
1754 assert!(sink.compaction_finished);
1755 assert!(sink.events.iter().any(
1756 |event| matches!(event, ProtocolEvent::AssistantDelta { text } if text == "continued")
1757 ));
1758 assert!(harness
1759 .session
1760 .history
1761 .iter()
1762 .any(|record| matches!(record, crate::session::SessionHistoryRecord::Compaction(_))));
1763 let provider_text = harness
1764 .session
1765 .provider_messages()
1766 .iter()
1767 .filter_map(|message| message.content.as_deref())
1768 .collect::<Vec<_>>()
1769 .join("\n");
1770 assert!(!provider_text.contains("old request"));
1771 assert!(provider_text.contains("continue"));
1772
1773 std::env::remove_var(key_env);
1774 std::fs::remove_dir_all(home).expect("cleanup");
1775 }
1776
1777 #[test]
1778 fn parses_only_message_records() {
1779 assert_eq!(
1780 parse_input_message(r#"{"type":"message","text":"hello"}"#).expect("message"),
1781 "hello"
1782 );
1783 assert!(parse_input_message(r#"{"type":"event","text":"hello"}"#).is_err());
1784 assert_eq!(
1785 parse_input_message(r#"{"type":"message","text":""}"#).expect("empty message"),
1786 ""
1787 );
1788 }
1789
1790 #[test]
1791 fn resolves_terminal_and_forced_modes() {
1792 assert_eq!(
1793 resolve_mode(&[], true, true).expect("default TUI"),
1794 FrontendMode::Tui
1795 );
1796 assert_eq!(
1797 resolve_mode(&[], true, false).expect("automatic JSONL"),
1798 FrontendMode::Jsonl
1799 );
1800 assert_eq!(
1801 resolve_mode(&["--jsonl".to_owned()], true, true).expect("forced JSONL"),
1802 FrontendMode::Jsonl
1803 );
1804 assert!(resolve_mode(&["--tui".to_owned()], true, false).is_err());
1805 }
1806
1807 #[test]
1808 fn redactor_does_not_leak_a_secret_across_deltas() {
1809 let mut redactor = SecretRedactor::new("secret");
1810 let mut output = Vec::new();
1811 redactor
1812 .push("prefix sec", |text| {
1813 output.push(text.to_owned());
1814 Ok(())
1815 })
1816 .expect("push");
1817 redactor
1818 .push("ret suffix", |text| {
1819 output.push(text.to_owned());
1820 Ok(())
1821 })
1822 .expect("push");
1823 redactor
1824 .finish(|text| {
1825 output.push(text.to_owned());
1826 Ok(())
1827 })
1828 .expect("finish");
1829 let output = output.join("");
1830 assert_eq!(
1831 output,
1832 format!("prefix {} suffix", redaction_marker("secret").unwrap())
1833 );
1834 assert!(!output.contains("secret"));
1835 }
1836
1837 #[test]
1838 fn redactor_handles_secrets_introduced_by_protocol_json_escaping() {
1839 let mut redactor = SecretRedactor::new("n0");
1840 let mut output = String::new();
1841 redactor
1842 .push("\n0", |text| {
1843 output.push_str(text);
1844 Ok(())
1845 })
1846 .expect("push");
1847 redactor
1848 .finish(|text| {
1849 output.push_str(text);
1850 Ok(())
1851 })
1852 .expect("finish");
1853 assert!(!output.contains("n0"));
1854 assert_eq!(output, redaction_marker("n0").unwrap());
1855 }
1856
1857 #[test]
1858 fn redactor_does_not_emit_a_secret_when_it_completes_at_a_delta_boundary() {
1859 let mut redactor = SecretRedactor::new("secret");
1860 let mut output = Vec::new();
1861 redactor
1862 .push("xsecre", |text| {
1863 output.push(text.to_owned());
1864 Ok(())
1865 })
1866 .expect("first delta");
1867 redactor
1868 .push("t", |text| {
1869 output.push(text.to_owned());
1870 Ok(())
1871 })
1872 .expect("second delta");
1873 redactor
1874 .finish(|text| {
1875 output.push(text.to_owned());
1876 Ok(())
1877 })
1878 .expect("finish");
1879 let output = output.join("");
1880 assert_eq!(output, format!("x{}", redaction_marker("secret").unwrap()));
1881 assert!(!output.contains("secret"));
1882 }
1883
1884 #[test]
1885 fn streaming_redaction_handles_marker_collision_keys_at_delta_boundaries() {
1886 for secret in ["REDACTED", "[REDACTED]"] {
1887 let mut redactor = SecretRedactor::new(secret);
1888 let split = secret.len() / 2;
1889 let (first, second) = secret.split_at(split);
1890 let mut output = String::new();
1891 redactor
1892 .push(first, |text| {
1893 output.push_str(text);
1894 Ok(())
1895 })
1896 .expect("first delta");
1897 redactor
1898 .push(second, |text| {
1899 output.push_str(text);
1900 Ok(())
1901 })
1902 .expect("second delta");
1903 redactor
1904 .finish(|text| {
1905 output.push_str(text);
1906 Ok(())
1907 })
1908 .expect("finish");
1909 assert!(!output.contains(secret));
1910 assert!(output.len() <= secret.len());
1911 }
1912 }
1913
1914 #[test]
1915 fn malformed_tool_arguments_use_a_safe_copy() {
1916 let secret = "provider-secret";
1917 let escaped = secret
1918 .chars()
1919 .map(|character| format!(r#"\u{:04x}"#, character as u32))
1920 .collect::<String>();
1921 let arguments = format!(r#"{{"command":"{escaped}""#);
1922 let safe = redact_tool_arguments(&arguments, secret);
1923 assert_eq!(safe, "{}");
1924 serde_json::from_str::<Value>(&safe).expect("safe arguments JSON");
1925 assert!(!safe.contains(secret));
1926 assert!(!safe.contains(&escaped));
1927 for invalid in ["[]", "{\"command\":1}", "{\"other\":\"value\"}"] {
1928 assert_eq!(redact_tool_arguments(invalid, secret), "{}");
1929 }
1930 assert_eq!(
1931 redact_tool_arguments(r#"{"command":"printf ordinary","background":true}"#, secret,),
1932 r#"{"background":true,"command":"printf ordinary"}"#
1933 );
1934 }
1935
1936 #[test]
1937 fn structured_redaction_preserves_tool_and_result_schema_keys() {
1938 let secret = "provider-secret";
1939 let value = serde_json::json!({
1940 "command": "printf provider-secret",
1941 "stdout": "provider-secret",
1942 "stderr": "ordinary",
1943 "exit_code": 0,
1944 "timed_out": false,
1945 "stdout_truncated": false,
1946 "stderr_truncated": false,
1947 "unknown-provider-secret": "provider-secret"
1948 });
1949 let redacted = redact_json_value(value, secret);
1950 for key in [
1951 "command",
1952 "stdout",
1953 "stderr",
1954 "exit_code",
1955 "timed_out",
1956 "stdout_truncated",
1957 "stderr_truncated",
1958 ] {
1959 assert!(redacted.get(key).is_some(), "missing schema key: {key}");
1960 }
1961 let encoded = serde_json::to_string(&redacted).expect("redacted JSON");
1962 assert!(!encoded.contains(secret));
1963 assert!(redacted.get("unknown-provider-secret").is_none());
1964 }
1965
1966 #[test]
1967 fn structured_redaction_preserves_typed_values_even_for_a_pathological_key() {
1968 let value = serde_json::json!({
1969 "exit_code": 0,
1970 "timed_out": false,
1971 "stdout_truncated": true,
1972 "error": null,
1973 });
1974 let redacted = redact_json_value(value, "0");
1975 assert!(redacted["exit_code"].is_number());
1976 assert!(redacted["timed_out"].is_boolean());
1977 assert!(redacted["stdout_truncated"].is_boolean());
1978 assert!(redacted["error"].is_null());
1979 }
1980
1981 #[test]
1982 fn reasoning_details_are_recursively_redacted_before_persistence() {
1983 let details = vec![serde_json::json!({
1984 "type": "reasoning.text",
1985 "text": "provider-secret",
1986 "nested": [{"value": "provider-secret"}],
1987 "provider-secret": "provider-secret"
1988 })];
1989 let redacted = redact_reasoning_details(&details, "provider-secret")
1990 .expect("non-empty reasoning details");
1991 let redacted = Value::Array(redacted);
1992 let encoded = serde_json::to_string(&redacted).expect("reasoning details JSON");
1993 assert!(!encoded.contains("provider-secret"));
1994 assert_eq!(redacted[0]["type"], "reasoning.text");
1995 assert_eq!(redacted[0]["text"], "[REDACTED]");
1996 assert_eq!(redacted[0]["nested"][0]["value"], "[REDACTED]");
1997 assert!(redacted[0].get("provider-secret").is_none());
1998 }
1999
2000 #[test]
2001 fn malformed_input_error_does_not_echo_secret_bearing_input() {
2002 let error =
2003 parse_input_message(r#"{"type":"message","text":"provider-secret","unexpected":}"#)
2004 .expect_err("invalid input");
2005 assert!(!error.contains("provider-secret"));
2006 }
2007
2008 #[test]
2009 fn malformed_input_is_an_error_event_and_not_diagnostic_json() {
2010 let mut output = Vec::new();
2011 let error = parse_input_message("not json").expect_err("invalid input");
2012 let mut protocol = ProtocolWriter::new(&mut output);
2013 protocol.error(&error).expect("error event");
2014 assert_eq!(String::from_utf8_lossy(&output).lines().count(), 1);
2015 let _ = Cursor::new("");
2016 }
2017
2018 #[test]
2019 fn early_diagnostic_scrubbing_removes_short_values_from_the_complete_line() {
2020 let secret = "lucy";
2021 let mut diagnostics = Vec::new();
2022 write_diagnostic_safe_with_environment(
2023 &mut diagnostics,
2024 secret,
2025 None,
2026 vec![secret.to_owned()],
2027 );
2028 let diagnostics = String::from_utf8(diagnostics).expect("diagnostic UTF-8");
2029 assert!(!diagnostics.contains(secret));
2030 }
2031 #[test]
2032 fn attached_agents_keeps_only_agents_files_and_redacts_their_paths() {
2033 let sources = vec![
2034 InstructionSource {
2035 path: std::path::PathBuf::from("/project/AGENTS.md"),
2036 contents: "agents".to_owned(),
2037 },
2038 InstructionSource {
2039 path: std::path::PathBuf::from("/project/CLAUDE.md"),
2040 contents: "claude".to_owned(),
2041 },
2042 InstructionSource {
2043 path: std::path::PathBuf::from("/private-secret/AGENTS.md"),
2044 contents: "agents".to_owned(),
2045 },
2046 ];
2047
2048 assert_eq!(
2049 attached_agents(sources, "secret"),
2050 vec!["/project/AGENTS.md", "/private-!/AGENTS.md"]
2051 );
2052 }
2053
2054 #[test]
2055 fn expands_slash_prefixed_skill_names_and_keeps_ordinary_messages() {
2056 let skill = SkillEntry {
2057 name: "release-notes".to_owned(),
2058 description: "Writes release notes".to_owned(),
2059 path: std::path::PathBuf::from("/skills/release-notes/SKILL.md"),
2060 contents: "# Release notes\nUse the template.".to_owned(),
2061 model_invocable: true,
2062 };
2063 let expanded = expand_skill_invocation("/release-notes v1.2", std::slice::from_ref(&skill))
2064 .expect("skill command");
2065 assert!(expanded.text.contains("# Release notes"));
2066 assert!(expanded.text.contains("User: v1.2"));
2067 assert_eq!(expanded.attached_skill.as_deref(), Some("release-notes"));
2068 let ordinary = expand_skill_invocation("ordinary message", &[]).expect("ordinary message");
2069 assert_eq!(ordinary.text, "ordinary message");
2070 assert_eq!(ordinary.attached_skill, None);
2071 assert_eq!(
2072 expand_skill_invocation("/missing", &[]).unwrap_err(),
2073 "unknown skill: missing"
2074 );
2075 assert_eq!(
2076 expand_skill_invocation("/skill:release-notes", &[skill]).unwrap_err(),
2077 "unknown skill: skill:release-notes"
2078 );
2079 }
2080}