Skip to main content

lucy/
app.rs

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    // The generic test/library entry point has no terminal handles. The real
124    // binary uses run_cli, which supplies the actual stdio terminal state.
125    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    /// Model context metadata resolved by the interactive frontend; `None`
516    /// keeps compaction disabled when an OpenAI-compatible provider exposes no
517    /// context-window metadata.
518    pub(crate) context_window: Option<usize>,
519    /// AGENTS.md sources selected for this newly created session's boot context.
520    /// The TUI uses these only while its first-boot welcome is visible.
521    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        // Endpoint and credential remain the session's established provider boundary.
580        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        // Validate the candidate before changing the user-owned source of truth.
585        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
1141/// Return the AGENTS.md files selected for the current new-session boot context.
1142/// Paths are secret-redacted before they can reach the terminal UI.
1143fn 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
1156/// Store a secret-safe skill snapshot with the session. The source is read
1157/// once during secure context discovery; later invocations never follow paths.
1158fn escape_xml_attribute(text: &str) -> String {
1159    text.replace('&', "&amp;")
1160        .replace('<', "&lt;")
1161        .replace('>', "&gt;")
1162        .replace('\"', "&quot;")
1163        .replace('\'', "&apos;")
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/// The message delivered to the provider and the optional name of the saved
1183/// skill snapshot that was attached to it.
1184#[derive(Debug)]
1185struct ExpandedSkillInvocation {
1186    text: String,
1187    attached_skill: Option<String>,
1188}
1189
1190/// Expand slash-prefixed skill names into the user message sent to the
1191/// provider. This deliberately adds no model-facing tool: skills are context,
1192/// not an executable capability of their own.
1193fn 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        // An incomplete argument fragment is an observation only. Do not
1286        // preserve malformed provider JSON: decoding it later could expose a
1287        // credential that was hidden by the outer JSON string.
1288        "{}".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}