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