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