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