1#![allow(clippy::missing_errors_doc)]
2mod args;
5mod client_state;
6mod clipboard;
7mod config;
8mod display_preview;
9mod environment;
10mod error;
11pub mod launcher;
12mod local_store;
13mod oauth;
14mod profiles;
15mod prompt_input;
16mod rpc;
17mod runner;
18pub(crate) mod runtime_coordinator;
19pub(crate) mod service;
20mod slash_commands;
21mod tui;
22mod update_check;
23
24use std::env;
25
26pub use args::{Cli, CliCommand, OutputMode, RpcCommand, RpcTransport, SessionCommand};
27pub use config::{CliConfig, ConfigResolver};
28pub use error::{CliError, CliResult};
29pub use local_store::{
30 DisplayReplayWindow, LocalSessionStore, LocalStore, LocalStreamArchive, TrimReport,
31};
32pub use service::CliService;
33pub use slash_commands::SlashCommandDefinition;
34
35pub fn run_from_env() -> CliResult<()> {
37 run(env::args())
38}
39
40pub fn run(args: impl IntoIterator<Item = String>) -> CliResult<()> {
42 let cli = match args::parse(args) {
43 Ok(cli) => cli,
44 Err(CliError::Display(output)) => {
45 print!("{output}");
46 return Ok(());
47 }
48 Err(error) => return Err(error),
49 };
50 let config = ConfigResolver::default().resolve(&cli)?;
51 if let Some(CliCommand::Rpc(command)) = &cli.command {
52 return rpc::run(&config, command);
53 }
54 let output = command_output_from_parts(cli, config)?;
55 print!("{output}");
56 Ok(())
57}
58
59pub fn run_rpc_server(command: &RpcCommand, store: Option<String>) -> CliResult<()> {
61 let cli = rpc_cli(store, command.clone());
62 let config = ConfigResolver::default().resolve(&cli)?;
63 rpc::run(&config, command)
64}
65
66pub fn command_output(args: impl IntoIterator<Item = String>) -> CliResult<String> {
68 let cli = match args::parse(args) {
69 Ok(cli) => cli,
70 Err(CliError::Display(output)) => return Ok(output),
71 Err(error) => return Err(error),
72 };
73 let config = ConfigResolver::default().resolve(&cli)?;
74 command_output_from_parts(cli, config)
75}
76
77const fn rpc_cli(store: Option<String>, command: RpcCommand) -> Cli {
78 Cli {
79 prompt: None,
80 session: None,
81 continue_session: false,
82 new_session: false,
83 run: None,
84 branch_from: None,
85 profile: None,
86 worker: None,
87 worker_label: None,
88 worktree: None,
89 worktree_name: None,
90 branch: None,
91 output: None,
92 hitl: None,
93 store,
94 command: Some(CliCommand::Rpc(command)),
95 }
96}
97
98fn command_output_from_parts(cli: Cli, config: CliConfig) -> CliResult<String> {
99 let show_update_hint = should_show_update_hint(&cli, &config);
100 if show_update_hint {
101 update_check::spawn_update_check_if_due(&config);
102 }
103 let hint = show_update_hint.then(|| update_check::update_hint(&config));
104 let service = CliService::open(config)?;
105 let mut output = service.execute(cli)?;
106 if let Some(Some(hint)) = hint {
107 output.push_str(&hint);
108 }
109 Ok(output)
110}
111
112const fn should_show_update_hint(cli: &Cli, config: &CliConfig) -> bool {
113 matches!(config.default_output, OutputMode::Text | OutputMode::Silent)
114 && matches!(
115 &cli.command,
116 None | Some(
117 CliCommand::Version
118 | CliCommand::Diagnostics
119 | CliCommand::ReplayCheck
120 | CliCommand::Run(_),
121 )
122 )
123}
124
125#[cfg(test)]
126mod tests {
127 #![allow(clippy::unwrap_used)]
128
129 use std::{ffi::OsString, io, path::Path};
130
131 use super::*;
132
133 fn output(root: &Path, raw_args: &[&str]) -> CliResult<String> {
134 let mut command_args = vec!["starweaver-cli".to_string()];
135 command_args.extend(raw_args.iter().map(|arg| (*arg).to_string()));
136 let cli = args::parse(command_args)?;
137 let config = ConfigResolver::for_tests(root).resolve(&cli)?;
138 CliService::open(config)?.execute(cli)
139 }
140
141 fn write_trim_test_config(root: &Path, auto_after_run: bool) {
142 let global = root.join("global");
143 std::fs::create_dir_all(&global).unwrap();
144 std::fs::write(
145 global.join("config.toml"),
146 format!(
147 r#"
148[general]
149model = "local_echo"
150
151[trim]
152auto_after_run = {auto_after_run}
153current_session_keep_recent_runs = 1
154all_sessions_keep_recent_runs = 1
155all_sessions_keep_days = 1
156all_sessions_interval_hours = 1
157"#
158 ),
159 )
160 .unwrap();
161 }
162
163 fn first_session_from_list_json(output: &str) -> serde_json::Value {
164 serde_json::from_str::<serde_json::Value>(output).unwrap()["sessions"][0].clone()
165 }
166
167 fn session_run_count(root: &Path, session_id: &str) -> usize {
168 let show = output(root, &["session", "show", session_id, "--output", "json"]).unwrap();
169 serde_json::from_str::<serde_json::Value>(&show).unwrap()["runs"]
170 .as_array()
171 .unwrap()
172 .len()
173 }
174
175 fn stored_run_dir_count(config: &CliConfig, session_id: &str) -> usize {
176 config
177 .file_store_path
178 .join("sessions")
179 .join(session_id)
180 .join("runs")
181 .read_dir()
182 .unwrap()
183 .filter(|entry| entry.as_ref().unwrap().file_type().unwrap().is_dir())
184 .count()
185 }
186
187 fn create_current_retention_session(root: &Path) -> String {
188 write_trim_test_config(root, true);
189 output(root, &["-p", "one", "--output", "silent"]).unwrap();
190 output(root, &["-p", "two", "--continue", "--output", "silent"]).unwrap();
191 output(root, &["-p", "three", "--continue", "--output", "silent"]).unwrap();
192 let session = first_session_from_list_json(
193 &output(root, &["session", "list", "--output", "json"]).unwrap(),
194 )["session_id"]
195 .as_str()
196 .unwrap()
197 .to_string();
198 assert_eq!(session_run_count(root, &session), 1);
199 session
200 }
201
202 fn create_archive_retention_session(root: &Path) -> String {
203 write_trim_test_config(root, false);
204 output(
205 root,
206 &["-p", "archive-one", "--new-session", "--output", "silent"],
207 )
208 .unwrap();
209 let session = first_session_from_list_json(
210 &output(root, &["session", "list", "--output", "json"]).unwrap(),
211 )["session_id"]
212 .as_str()
213 .unwrap()
214 .to_string();
215 output(
216 root,
217 &[
218 "-p",
219 "archive-two",
220 "--session",
221 &session,
222 "--output",
223 "silent",
224 ],
225 )
226 .unwrap();
227 output(
228 root,
229 &[
230 "-p",
231 "archive-three",
232 "--session",
233 &session,
234 "--output",
235 "silent",
236 ],
237 )
238 .unwrap();
239 session
240 }
241
242 fn mark_session_as_retention_eligible(config: &CliConfig, session_id: &str) {
243 let old = chrono::DateTime::parse_from_rfc3339("2000-01-01T00:00:00+00:00")
244 .unwrap()
245 .to_rfc3339();
246 let conn = rusqlite::Connection::open(&config.database_path).unwrap();
247 conn.execute(
248 "UPDATE runs SET updated_at = ?1 WHERE session_id = ?2",
249 rusqlite::params![old, session_id],
250 )
251 .unwrap();
252 conn.execute(
253 "UPDATE sessions SET updated_at = ?1 WHERE session_id = ?2",
254 rusqlite::params![old, session_id],
255 )
256 .unwrap();
257 drop(conn);
258
259 let state_path = config.project_dir.join("state.json");
260 let mut state = serde_json::from_str::<serde_json::Value>(
261 &std::fs::read_to_string(&state_path).unwrap(),
262 )
263 .unwrap();
264 state["last_retention_maintenance_at"] = serde_json::json!(old);
265 std::fs::write(&state_path, serde_json::to_vec_pretty(&state).unwrap()).unwrap();
266 }
267
268 #[test]
269 #[allow(clippy::too_many_lines)]
270 fn args_and_error_helpers_cover_edge_branches() {
271 let run = args::RunCommand {
272 prompt: Some(" explicit ".to_string()),
273 prompt_parts: vec!["ignored".to_string()],
274 continue_session: false,
275 session: None,
276 new_session: false,
277 run: None,
278 branch_from: None,
279 profile: None,
280 output: None,
281 hitl: None,
282 goal: None,
283 worker: None,
284 worker_label: None,
285 worktree: None,
286 worktree_name: None,
287 branch: None,
288 session_affinity_id: None,
289 environment_attachments: Vec::new(),
290 };
291 assert_eq!(run.prompt_text().unwrap(), " explicit ");
292
293 let joined = args::RunCommand {
294 prompt: None,
295 prompt_parts: vec!["hello".to_string(), "world".to_string()],
296 continue_session: false,
297 session: None,
298 new_session: false,
299 run: None,
300 branch_from: None,
301 profile: None,
302 output: None,
303 hitl: None,
304 goal: None,
305 worker: None,
306 worker_label: None,
307 worktree: None,
308 worktree_name: None,
309 branch: None,
310 session_affinity_id: None,
311 environment_attachments: Vec::new(),
312 };
313 assert_eq!(joined.prompt_text().unwrap(), "hello world");
314
315 let empty = args::RunCommand {
316 prompt: Some(" ".to_string()),
317 prompt_parts: Vec::new(),
318 continue_session: false,
319 session: None,
320 new_session: false,
321 run: None,
322 branch_from: None,
323 profile: None,
324 output: None,
325 hitl: None,
326 goal: None,
327 worker: None,
328 worker_label: None,
329 worktree: None,
330 worktree_name: None,
331 branch: None,
332 session_affinity_id: None,
333 environment_attachments: Vec::new(),
334 };
335 assert!(
336 matches!(empty.prompt_text(), Err(CliError::Usage(message)) if message.contains("run -p"))
337 );
338
339 let parsed = args::parse_os([
340 OsString::from("starweaver-cli"),
341 OsString::from("run"),
342 OsString::from("hello"),
343 ])
344 .unwrap();
345 assert!(matches!(parsed.command, Some(args::CliCommand::Run(_))));
346
347 let parsed = args::parse_os([
348 OsString::from("starweaver-cli"),
349 OsString::from("-p"),
350 OsString::from("hello"),
351 OsString::from("-s"),
352 OsString::from("session_test"),
353 OsString::from("--profile"),
354 OsString::from("coding"),
355 OsString::from("--worker"),
356 OsString::from("off"),
357 OsString::from("--worktree"),
358 OsString::from("feature"),
359 OsString::from("--branch"),
360 OsString::from("feature/work"),
361 ])
362 .unwrap();
363 assert_eq!(parsed.session.as_deref(), Some("session_test"));
364 assert_eq!(parsed.profile.as_deref(), Some("coding"));
365 assert_eq!(parsed.worker.as_deref(), Some("off"));
366 assert_eq!(parsed.worktree.as_deref(), Some("feature"));
367 assert_eq!(parsed.branch.as_deref(), Some("feature/work"));
368
369 let parsed = args::parse_os([
370 OsString::from("starweaver-cli"),
371 OsString::from("-p"),
372 OsString::from("hello"),
373 OsString::from("--worker"),
374 OsString::from("-w"),
375 OsString::from("--worker-label"),
376 OsString::from("executor"),
377 OsString::from("--worktree-name"),
378 OsString::from("feature"),
379 ])
380 .unwrap();
381 assert_eq!(parsed.worker.as_deref(), Some("true"));
382 assert_eq!(parsed.worker_label.as_deref(), Some("executor"));
383 assert_eq!(parsed.worktree.as_deref(), Some("true"));
384 assert_eq!(parsed.worktree_name.as_deref(), Some("feature"));
385
386 let parse_error =
387 args::parse_os([OsString::from("starweaver-cli"), OsString::from("--bad")]);
388 assert!(
389 matches!(parse_error, Err(CliError::Usage(message)) if message.contains("unexpected argument"))
390 );
391
392 assert!(
393 format!(
394 "{}",
395 CliError::from(serde_json::from_str::<serde_json::Value>("{").unwrap_err())
396 )
397 .contains("serialization error")
398 );
399 assert!(
400 format!(
401 "{}",
402 CliError::from(toml::from_str::<toml::Value>("=").unwrap_err())
403 )
404 .contains("configuration error")
405 );
406 assert!(
407 format!(
408 "{}",
409 CliError::from(toml::to_string(&f64::NAN).unwrap_err())
410 )
411 .contains("configuration error")
412 );
413 let io_error = error::io_error(
414 "/tmp/missing",
415 io::Error::new(io::ErrorKind::NotFound, "gone"),
416 );
417 assert!(format!("{io_error}").contains("filesystem error at /tmp/missing"));
418 }
419
420 #[test]
421 fn version_and_diagnostics_work() {
422 let temp = tempfile::tempdir().unwrap();
423 assert_eq!(
424 output(temp.path(), &["version"]).unwrap(),
425 "starweaver-agent-sdk\n"
426 );
427 let diagnostics = output(temp.path(), &["diagnostics"]).unwrap();
428 assert!(diagnostics.contains("sdk=starweaver-agent-sdk"));
429 assert!(diagnostics.contains("database_path="));
430 assert!(diagnostics.contains("model_profiles="));
431 assert!(diagnostics.contains("wal=true"));
432 }
433
434 #[test]
435 fn config_model_profiles_work() {
436 let temp = tempfile::tempdir().unwrap();
437 let global = temp.path().join("global");
438 std::fs::create_dir_all(&global).unwrap();
439 std::fs::write(
440 global.join("config.toml"),
441 r#"
442[general]
443model = "homelab@openai-responses:gpt-5.5"
444model_settings = "openai_responses_high"
445model_cfg = "gpt5_270k"
446
447[model_profiles.codex-subs]
448label = "Codex Subs"
449model = "oauth@codex:gpt-5.5"
450model_settings = "openai_responses_high"
451model_cfg = "gpt5_270k"
452
453[providers.homelab]
454base_url = "https://gateway.example/v1"
455max_tokens_parameter = "omit"
456
457[oauth_refresh]
458enabled = true
459interval_seconds = 42
460failure_retry_seconds = 7
461refresh_on_startup = false
462
463[env]
464HOMELAB_API_KEY = "test-key"
465"#,
466 )
467 .unwrap();
468 let diagnostics = output(temp.path(), &["diagnostics"]).unwrap();
469 assert!(diagnostics.contains("profile=default_model"));
470 assert!(diagnostics.contains("model_profiles=1"));
471 assert_eq!(
472 output(
473 temp.path(),
474 &["config", "get", "providers.homelab.max_tokens_parameter"]
475 )
476 .unwrap(),
477 "omit\n"
478 );
479 assert_eq!(
480 output(
481 temp.path(),
482 &["config", "get", "oauth_refresh.interval_seconds"]
483 )
484 .unwrap(),
485 "42\n"
486 );
487 assert_eq!(
488 output(
489 temp.path(),
490 &["config", "get", "oauth_refresh.refresh_on_startup"]
491 )
492 .unwrap(),
493 "false\n"
494 );
495 let profiles = output(temp.path(), &["profile", "list"]).unwrap();
496 assert!(profiles.contains("default_model"));
497 assert!(profiles.contains("codex-subs"));
498 let default_profile = output(temp.path(), &["profile", "show", "default_model"]).unwrap();
499 assert!(default_profile.contains("model_id: homelab@openai-responses:gpt-5.5"));
500 assert!(default_profile.contains("settings_preset: openai_responses_high"));
501 assert!(default_profile.contains("config_preset: gpt5_270k"));
502 assert!(default_profile.contains("# source: config"));
503 }
504
505 #[test]
506 fn configured_slash_commands_layer_aliases_and_redact_unmapped_metadata() {
507 let temp = tempfile::tempdir().unwrap();
508 let global = temp.path().join("global");
509 let project = temp.path().join("project/.starweaver");
510 std::fs::create_dir_all(&global).unwrap();
511 std::fs::create_dir_all(&project).unwrap();
512 std::fs::write(
513 global.join("config.toml"),
514 r#"
515[commands.review]
516description = "Global review"
517aliases = ["rv", "bad alias", "model"]
518prompt = "global secret prompt"
519
520[commands.other]
521aliases = ["review"]
522prompt = "Other command"
523"#,
524 )
525 .unwrap();
526 std::fs::write(
527 project.join("config.toml"),
528 r#"
529[commands.review]
530description = "Project review"
531aliases = ["pr"]
532prompt = "Project review prompt"
533
534[commands.bad_name]
535prompt = "ignored because underscore is valid"
536
537[commands."bad name"]
538prompt = "ignored invalid name"
539"#,
540 )
541 .unwrap();
542
543 let cli = args::parse(["starweaver-cli".to_string(), "diagnostics".to_string()]).unwrap();
544 let config = ConfigResolver::for_tests(temp.path())
545 .resolve(&cli)
546 .unwrap();
547 let review = config.slash_commands.get("review").unwrap();
548 assert_eq!(review.prompt, "Project review prompt");
549 assert_eq!(review.aliases, vec!["pr".to_string()]);
550 assert!(config.slash_commands.contains_key("pr"));
551 assert!(!config.slash_commands.contains_key("rv"));
552 assert!(!config.slash_commands.contains_key("bad alias"));
553 assert!(!config.slash_commands.contains_key("model"));
554 assert!(config.slash_commands.contains_key("bad_name"));
555 assert!(!config.slash_commands.contains_key("bad name"));
556 let unmapped = output(temp.path(), &["config", "get", "metadata.unmapped"]).unwrap();
557 assert!(!unmapped.contains("global secret prompt"));
558 assert!(!unmapped.contains("Project review prompt"));
559 assert!(!unmapped.contains("commands"));
560 }
561
562 #[test]
563 fn configured_subagent_inherits_profile_model() {
564 let temp = tempfile::tempdir().unwrap();
565 let global = temp.path().join("global");
566 let project = temp.path().join("project/.starweaver");
567 std::fs::create_dir_all(global.join("subagents")).unwrap();
568 std::fs::write(
569 global.join("config.toml"),
570 r#"
571[general]
572model = "local_echo"
573
574[subagents]
575dirs = ["subagents"]
576"#,
577 )
578 .unwrap();
579 std::fs::write(
580 global.join("subagents/helper.md"),
581 r"---
582name: helper
583description: Helper subagent
584model: inherit
585---
586You are a helper.
587",
588 )
589 .unwrap();
590
591 let cli = args::parse([
592 "starweaver-cli".to_string(),
593 "-p".to_string(),
594 "hello".to_string(),
595 "--profile".to_string(),
596 "default_model".to_string(),
597 ])
598 .unwrap();
599 let config = ConfigResolver::for_tests(temp.path())
600 .resolve(&cli)
601 .unwrap();
602 assert_eq!(config.project_dir, project);
603 let profile = crate::profiles::resolve_profile(&config, Some("default_model")).unwrap();
604 let agent = profile.build_agent().unwrap();
605 let tools = agent.tools().names();
606 assert!(tools.contains(&"delegate".to_string()));
607 assert!(tools.contains(&"subagent_info".to_string()));
608
609 let run = output(
610 temp.path(),
611 &[
612 "-p",
613 "hello",
614 "--profile",
615 "default_model",
616 "--output",
617 "silent",
618 ],
619 )
620 .unwrap();
621 assert!(run.contains("status=completed"));
622 }
623
624 #[test]
625 fn headless_run_expands_configured_slash_commands() {
626 let temp = tempfile::tempdir().unwrap();
627 let global = temp.path().join("global");
628 std::fs::create_dir_all(&global).unwrap();
629 std::fs::write(
630 global.join("config.toml"),
631 r#"
632[general]
633model = "local_echo"
634
635[commands.review]
636description = "Review the current changes"
637aliases = ["rv"]
638prompt = "Review carefully."
639"#,
640 )
641 .unwrap();
642
643 let run = output(
644 temp.path(),
645 &[
646 "-p",
647 "/rv staged diff",
648 "--profile",
649 "default_model",
650 "--output",
651 "text",
652 ],
653 )
654 .unwrap();
655 assert!(run.contains("local echo: Review carefully."));
656 assert!(run.contains("User instruction: staged diff"));
657
658 let sessions = output(temp.path(), &["session", "list"]).unwrap();
659 let session: serde_json::Value =
660 serde_json::from_str(sessions.lines().next().unwrap()).unwrap();
661 let session_id = session["session_id"].as_str().unwrap();
662 let cli = args::parse([
663 "starweaver-cli".to_string(),
664 "session".to_string(),
665 "list".to_string(),
666 ])
667 .unwrap();
668 let config = ConfigResolver::for_tests(temp.path())
669 .resolve(&cli)
670 .unwrap();
671 let store = LocalStore::open(&config).unwrap();
672 let run_id = session["head_run_id"].as_str().unwrap();
673 let run_record = store.load_run(session_id, run_id).unwrap();
674 let run_value = serde_json::to_value(&run_record).unwrap();
675 assert_eq!(
676 run_value["input"][0]["text"],
677 "Review carefully.\n\nUser instruction: staged diff"
678 );
679 assert_eq!(run_value["metadata"]["cli.slash_command.name"], "review");
680 assert_eq!(run_value["metadata"]["cli.slash_command.invoked"], "rv");
681 }
682
683 #[test]
684 fn headless_run_creates_session_and_run() {
685 let temp = tempfile::tempdir().unwrap();
686 let first = output(temp.path(), &["-p", "hello", "--output", "display-jsonl"]).unwrap();
687 let first_message: serde_json::Value =
688 serde_json::from_str(first.lines().next().unwrap()).unwrap();
689 assert_eq!(first_message["schema"], "starweaver.display.v1");
690 assert_eq!(first_message["type"], "RUN_QUEUED");
691 let agui_temp = tempfile::tempdir().unwrap();
692 let agui = output(agui_temp.path(), &["-p", "hello", "--output", "agui-jsonl"]).unwrap();
693 let agui_events = agui
694 .lines()
695 .map(|line| serde_json::from_str::<serde_json::Value>(line).unwrap())
696 .collect::<Vec<_>>();
697 assert!(
698 agui_events
699 .iter()
700 .any(|event| event["type"] == "RUN_STARTED")
701 );
702 assert!(
703 agui_events
704 .iter()
705 .any(|event| event["type"] == "TEXT_MESSAGE_CHUNK")
706 );
707 assert!(
708 agui_events
709 .iter()
710 .any(|event| event["type"] == "RUN_FINISHED")
711 );
712 let sessions = output(temp.path(), &["session", "list"]).unwrap();
713 assert!(sessions.contains("session_"));
714 let value: serde_json::Value =
715 serde_json::from_str(sessions.lines().next().unwrap()).unwrap();
716 assert_eq!(value["run_count"], 1);
717 assert_eq!(value["head_success_run_id"], value["head_run_id"]);
718 }
719
720 #[test]
721 fn continue_appends_run_under_existing_session() {
722 let temp = tempfile::tempdir().unwrap();
723 output(temp.path(), &["-p", "one"]).unwrap();
724 output(temp.path(), &["-p", "two", "--continue"]).unwrap();
725 let sessions = output(temp.path(), &["session", "list"]).unwrap();
726 let value: serde_json::Value =
727 serde_json::from_str(sessions.lines().next().unwrap()).unwrap();
728 assert_eq!(value["run_count"], 2);
729 let session_id = value["session_id"].as_str().unwrap();
730 let show = output(temp.path(), &["session", "show", session_id]).unwrap();
731 assert_eq!(show.lines().count(), 3);
732 }
733
734 #[test]
735 fn display_replay_window_uses_scoped_cursors() {
736 let temp = tempfile::tempdir().unwrap();
737 output(temp.path(), &["-p", "one"]).unwrap();
738 output(temp.path(), &["-p", "two", "--continue"]).unwrap();
739 let sessions = output(temp.path(), &["session", "list"]).unwrap();
740 let value: serde_json::Value =
741 serde_json::from_str(sessions.lines().next().unwrap()).unwrap();
742 let session_id = value["session_id"].as_str().unwrap();
743 let cli = args::parse([
744 "starweaver-cli".to_string(),
745 "session".to_string(),
746 "list".to_string(),
747 ])
748 .unwrap();
749 let config = ConfigResolver::for_tests(temp.path())
750 .resolve(&cli)
751 .unwrap();
752 let store = LocalStore::open(&config).unwrap();
753
754 let session_window = store.replay_display_window(session_id, None, None).unwrap();
755 assert_eq!(
756 session_window.scope,
757 starweaver_stream::ReplayScope::session(session_id)
758 );
759 assert_eq!(session_window.next_sequence, session_window.events.len());
760 for (sequence, event) in session_window.events.iter().enumerate() {
761 assert_eq!(
762 event.scope,
763 starweaver_stream::ReplayScope::session(session_id)
764 );
765 assert_eq!(event.sequence, sequence);
766 }
767
768 let session_cursor = starweaver_stream::ReplayCursor::new(session_window.scope.clone(), 0);
769 let session_tail = store
770 .replay_display_window(session_id, None, Some(&session_cursor))
771 .unwrap();
772 assert!(session_tail.events.iter().all(|event| event.sequence > 0));
773 assert_eq!(session_tail.next_sequence, session_window.next_sequence);
774
775 let first_run = store.list_runs(session_id, 10).unwrap().remove(0);
776 let run_window = store
777 .replay_display_window(session_id, Some(&first_run.run_id), None)
778 .unwrap();
779 assert_eq!(
780 run_window.scope,
781 starweaver_stream::ReplayScope::run(&first_run.run_id)
782 );
783 assert!(run_window.next_sequence > 0);
784 assert!(
785 run_window
786 .events
787 .iter()
788 .all(|event| event.scope == starweaver_stream::ReplayScope::run(&first_run.run_id))
789 );
790 }
791
792 #[test]
793 fn local_stream_archive_implements_shared_stream_contract() {
794 use starweaver_stream::StreamArchive as _;
795
796 let temp = tempfile::tempdir().unwrap();
797 output(temp.path(), &["-p", "archive"]).unwrap();
798 let sessions = output(temp.path(), &["session", "list"]).unwrap();
799 let value: serde_json::Value =
800 serde_json::from_str(sessions.lines().next().unwrap()).unwrap();
801 let session_id = value["session_id"].as_str().unwrap();
802 let cli = args::parse([
803 "starweaver-cli".to_string(),
804 "session".to_string(),
805 "list".to_string(),
806 ])
807 .unwrap();
808 let config = ConfigResolver::for_tests(temp.path())
809 .resolve(&cli)
810 .unwrap();
811 let store = LocalStore::open(&config).unwrap();
812 let run = store.list_runs(session_id, 10).unwrap().remove(0);
813 let archive = LocalStreamArchive::new(config);
814 let runtime = tokio::runtime::Runtime::new().unwrap();
815 let run_scope = starweaver_stream::ReplayScope::run(&run.run_id);
816 let session_scope = starweaver_stream::ReplayScope::session(session_id);
817
818 let run_messages = runtime
819 .block_on(archive.replay_display_after(&run_scope, None))
820 .unwrap();
821 assert!(!run_messages.is_empty());
822 let run_range = runtime.block_on(archive.cursor_range(&run_scope)).unwrap();
823 assert!(run_range.is_some());
824
825 let session_messages = runtime
826 .block_on(archive.replay_display_after(
827 &session_scope,
828 Some(starweaver_stream::ReplayCursor::new(
829 session_scope.clone(),
830 0,
831 )),
832 ))
833 .unwrap();
834 assert!(session_messages.len() < run_messages.len());
835 let session_range = runtime
836 .block_on(archive.cursor_range(&session_scope))
837 .unwrap()
838 .unwrap();
839 assert_eq!(session_range.0.sequence, 0);
840
841 let extra = starweaver_stream::DisplayMessage::new(
842 999,
843 starweaver_core::SessionId::from_string(session_id),
844 starweaver_core::RunId::from_string(&run.run_id),
845 starweaver_stream::DisplayMessageKind::HostEvent,
846 )
847 .with_preview("extra archive message");
848 runtime
849 .block_on(archive.append_display_messages(run_scope.clone(), vec![extra]))
850 .unwrap();
851 let appended = runtime
852 .block_on(archive.replay_display_after(
853 &run_scope,
854 Some(starweaver_stream::ReplayCursor::new(run_scope.clone(), 998)),
855 ))
856 .unwrap();
857 assert_eq!(appended.len(), 1);
858 assert_eq!(
859 appended[0].preview.as_deref(),
860 Some("extra archive message")
861 );
862
863 let snapshot = starweaver_stream::ReplaySnapshot {
864 scope: Some(run_scope.clone()),
865 revision: 7,
866 cursor: Some(starweaver_stream::ReplayCursor::new(run_scope.clone(), 999)),
867 display_messages: appended,
868 metadata: serde_json::Map::default(),
869 };
870 runtime
871 .block_on(archive.append_snapshot(run_scope.clone(), snapshot.clone()))
872 .unwrap();
873 let latest = runtime
874 .block_on(archive.latest_snapshot(&run_scope))
875 .unwrap()
876 .unwrap();
877 assert_eq!(latest.revision, snapshot.revision);
878
879 let _raw_records = runtime
880 .block_on(archive.replay_raw_after(
881 &starweaver_core::SessionId::from_string(session_id),
882 &starweaver_core::RunId::from_string(&run.run_id),
883 None,
884 ))
885 .unwrap();
886 }
887
888 #[test]
889 fn local_session_and_stream_adapters_back_agent_runtime_builder() {
890 use std::sync::Arc;
891
892 use starweaver_session::SessionStore as _;
893 use starweaver_stream::StreamArchive as _;
894
895 let temp = tempfile::tempdir().unwrap();
896 let cli = args::parse([
897 "starweaver-cli".to_string(),
898 "session".to_string(),
899 "list".to_string(),
900 ])
901 .unwrap();
902 let config = ConfigResolver::for_tests(temp.path())
903 .resolve(&cli)
904 .unwrap();
905 let session_id = starweaver_core::SessionId::from_string("session_local_runtime");
906 let session_store = Arc::new(LocalSessionStore::new(config.clone()));
907 let stream_archive = Arc::new(LocalStreamArchive::new(config));
908 let runtime = tokio::runtime::Runtime::new().unwrap();
909 let mut agent_runtime = starweaver_agent::AgentRuntimeBuilder::new(Arc::new(
910 starweaver_agent::TestModel::with_text("ok"),
911 ))
912 .durable_session_id(session_id.clone())
913 .session_store(session_store.clone())
914 .stream_archive(stream_archive.clone())
915 .build();
916
917 let result = runtime.block_on(agent_runtime.run_stream("hello")).unwrap();
918 assert_eq!(result.result.output, "ok");
919
920 let runs = runtime
921 .block_on(session_store.list_runs(&session_id))
922 .unwrap();
923 assert_eq!(runs.len(), 1);
924 assert_eq!(runs[0].status, starweaver_session::RunStatus::Completed);
925 let run_scope = starweaver_stream::ReplayScope::run(runs[0].run_id.as_str());
926 let display_messages = runtime
927 .block_on(stream_archive.replay_display_after(&run_scope, None))
928 .unwrap();
929 assert!(!display_messages.is_empty());
930 let trace = runtime
931 .block_on(session_store.compact_session_trace(&session_id))
932 .unwrap();
933 assert_eq!(trace.runs, 1);
934 }
935
936 #[test]
937 fn automatic_retention_prunes_old_runs_without_deleting_sessions() {
938 let temp = tempfile::tempdir().unwrap();
939 let current_session = create_current_retention_session(temp.path());
940 let archive_session = create_archive_retention_session(temp.path());
941 let config = ConfigResolver::for_tests(temp.path())
942 .resolve(&args::parse(["starweaver-cli".to_string()]).unwrap())
943 .unwrap();
944 mark_session_as_retention_eligible(&config, &archive_session);
945 assert_eq!(session_run_count(temp.path(), &archive_session), 3);
946 assert_eq!(stored_run_dir_count(&config, &archive_session), 3);
947
948 write_trim_test_config(temp.path(), true);
949 output(
950 temp.path(),
951 &[
952 "-p",
953 "trigger",
954 "--session",
955 ¤t_session,
956 "--output",
957 "silent",
958 ],
959 )
960 .unwrap();
961
962 assert_eq!(session_run_count(temp.path(), ¤t_session), 1);
963 assert_eq!(session_run_count(temp.path(), &archive_session), 1);
964 assert_eq!(stored_run_dir_count(&config, &archive_session), 1);
965 let sessions = output(temp.path(), &["session", "list", "--output", "json"]).unwrap();
966 let sessions = serde_json::from_str::<serde_json::Value>(&sessions).unwrap()["sessions"]
967 .as_array()
968 .unwrap()
969 .len();
970 assert_eq!(sessions, 2);
971 assert!(
972 std::fs::read_to_string(config.project_dir.join("state.json"))
973 .unwrap()
974 .contains("last_retention_maintenance_at")
975 );
976 }
977
978 #[test]
979 fn replay_and_trim_work() {
980 let temp = tempfile::tempdir().unwrap();
981 output(temp.path(), &["-p", "one"]).unwrap();
982 output(temp.path(), &["-p", "two", "--continue"]).unwrap();
983 output(temp.path(), &["-p", "three", "--continue"]).unwrap();
984 let sessions = output(temp.path(), &["session", "list"]).unwrap();
985 let session: serde_json::Value =
986 serde_json::from_str(sessions.lines().next().unwrap()).unwrap();
987 let session_id = session["session_id"].as_str().unwrap();
988 let replay = output(temp.path(), &["session", "replay", session_id]).unwrap();
989 assert!(replay.contains("RUN_FINISHED"));
990 let dry = output(
991 temp.path(),
992 &[
993 "session",
994 "trim",
995 "--session",
996 session_id,
997 "--keep-runs",
998 "1",
999 "--dry-run",
1000 ],
1001 )
1002 .unwrap();
1003 let report: serde_json::Value = serde_json::from_str(dry.trim()).unwrap();
1004 assert_eq!(report["runs_to_trim"], 2);
1005 output(
1006 temp.path(),
1007 &[
1008 "session",
1009 "trim",
1010 "--session",
1011 session_id,
1012 "--keep-runs",
1013 "1",
1014 ],
1015 )
1016 .unwrap();
1017 let show = output(temp.path(), &["session", "show", session_id]).unwrap();
1018 assert_eq!(show.lines().count(), 2);
1019 }
1020}