Skip to main content

lucy/
app.rs

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