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