1use std::path::Path;
13
14use anyhow::{Context, Result, bail, ensure};
15
16use super::worker_binary::{bridge_launch, container_upload_ownership_args, stage_profile};
17use super::{Controller, execute_checked};
18use crate::targets::{self, CommandExecutor, CommandSpec, ProcessExecutor};
19use mj_core::worker_launch::{
20 REVIEWER_DIR, ReviewMcpDelivery, ReviewMcpServer, ReviewerLaunchConfig,
21 reviewer_staging_profile_home,
22};
23
24pub fn reviewer_stager() -> mj_client::session::ReviewerStager {
28 mj_client::session::ReviewerStager::new(ControllerReviewerStager)
29}
30
31struct ControllerReviewerStager;
32
33impl mj_client::session::ReviewerStagerBackend for ControllerReviewerStager {
34 fn stage(
35 &self,
36 _config: mj_core::config::Config,
37 session: mj_core::state::SessionRecord,
38 profile_id: String,
39 generation: u64,
40 cancelled: std::sync::Arc<std::sync::atomic::AtomicBool>,
41 ) -> Result<ReviewerLaunchConfig> {
42 let session_id = session.id.clone();
43 let controller = Controller::load()?;
44 let executor = crate::targets::CancellableProcessExecutor::new(cancelled)
45 .with_deadline(std::time::Duration::from_secs(90));
46 controller.stage_reviewer_profile_controlled(
47 &session_id,
48 &profile_id,
49 generation,
50 &[],
51 &executor,
52 )
53 }
54}
55
56impl Controller {
57 pub fn stage_reviewer_profile(
63 &self,
64 session_id: &str,
65 profile_id: &str,
66 generation: u64,
67 ) -> Result<ReviewerLaunchConfig> {
68 self.stage_reviewer_profile_controlled(
69 session_id,
70 profile_id,
71 generation,
72 &[],
73 &ProcessExecutor,
74 )
75 }
76
77 pub fn stage_reviewer_profile_with_mcp(
85 &self,
86 session_id: &str,
87 profile_id: &str,
88 generation: u64,
89 mcp_servers: &[ReviewMcpServer],
90 dispatch_tool: bool,
91 ) -> Result<ReviewerLaunchConfig> {
92 let mut servers = mcp_servers.to_vec();
93 if dispatch_tool {
94 let (_, worker_root) = self.worker_placement(session_id)?;
95 servers.push(review_dispatch_server(&worker_root));
96 }
97 self.stage_reviewer_profile_controlled(
98 session_id,
99 profile_id,
100 generation,
101 &servers,
102 &ProcessExecutor,
103 )
104 }
105
106 pub fn stage_reviewer_profile_controlled(
107 &self,
108 session_id: &str,
109 profile_id: &str,
110 generation: u64,
111 mcp_servers: &[ReviewMcpServer],
112 executor: &impl CommandExecutor,
113 ) -> Result<ReviewerLaunchConfig> {
114 let profile = self
115 .config
116 .profiles
117 .get(profile_id)
118 .with_context(|| format!("unknown profile {profile_id:?}"))?;
119 ensure!(
120 profile.enabled,
121 "reviewer profile {profile_id:?} is disabled"
122 );
123 if !profile.kind.supports_injected_mcp() {
124 bail!(
125 "Muse Code cannot be a reviewer: muse-acp does not accept the required MCP tools. Select another reviewer profile."
126 );
127 }
128 let session = self
129 .state
130 .sessions
131 .get(session_id)
132 .with_context(|| format!("unknown session {session_id}"))?;
133 let target = session.target_runtime_settings(&self.config)?;
134 let execution_policy = profile
135 .kind
136 .effective_execution_policy(target.execution_policy);
137 let (backend, worker_root) = self.worker_placement(session_id)?;
138
139 let staging = tempfile::tempdir().context("create reviewer staging directory")?;
140 let local = staging.path().join("profile");
141 stage_profile(profile, &local).with_context(|| format!("stage profile {profile_id:?}"))?;
142 if !mcp_servers.is_empty()
146 && ReviewMcpDelivery::for_harness(profile.kind) == ReviewMcpDelivery::HarnessProfile
147 {
148 configure_staged_review_mcp(profile.kind, &local, mcp_servers)
149 .with_context(|| format!("configure reviewer MCP servers for {profile_id:?}"))?;
150 }
151 upload_reviewer_profile(executor, &backend, &worker_root, generation, &local)?;
152
153 let (bridge_command, bridge_args) = bridge_launch(profile.kind, execution_policy);
154 let mut environment = profile.environment.clone();
155 environment.remove(profile.home_env());
159 let excluded_environment = profile.exclude_harness_environment(&mut environment);
162 Ok(ReviewerLaunchConfig {
163 profile_id: profile_id.to_owned(),
164 harness: profile.kind,
165 bridge_command: bridge_command.into(),
166 bridge_args,
167 environment,
168 excluded_environment,
169 execution_policy,
170 model: None,
171 effort: None,
172 fast_mode: None,
173 generation,
174 mcp_servers: match ReviewMcpDelivery::for_harness(profile.kind) {
178 ReviewMcpDelivery::Acp => mcp_servers.to_vec(),
179 ReviewMcpDelivery::HarnessProfile => Vec::new(),
180 },
181 })
182 }
183}
184
185fn configure_staged_review_mcp(
193 harness: mj_core::config::HarnessKind,
194 profile_stage: &Path,
195 servers: &[ReviewMcpServer],
196) -> Result<()> {
197 let Some(file) = harness.mcp_config_file() else {
198 bail!("{harness:?} does not read MCP servers from its profile");
199 };
200 let kimi = harness == mj_core::config::HarnessKind::Kimi;
201 let path = profile_stage.join(file);
202 let mut document = match std::fs::read(&path) {
203 Ok(body) => serde_json::from_slice::<serde_json::Value>(&body)
204 .with_context(|| format!("parse staged reviewer configuration {}", path.display()))?,
205 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
206 serde_json::Value::Object(serde_json::Map::new())
207 }
208 Err(error) => {
209 return Err(error)
210 .with_context(|| format!("read staged reviewer configuration {}", path.display()));
211 }
212 };
213 let root = document.as_object_mut().with_context(|| {
214 format!(
215 "staged reviewer configuration {} must contain a JSON object",
216 path.display()
217 )
218 })?;
219 let configured = root
220 .entry("mcpServers")
221 .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()))
222 .as_object_mut()
223 .with_context(|| {
224 format!(
225 "mcpServers in staged reviewer configuration {} must be a JSON object",
226 path.display()
227 )
228 })?;
229 for server in servers {
230 let mut entry = serde_json::json!({
231 "command": server.command,
232 "args": server.args,
233 });
234 if kimi {
235 let object = entry
236 .as_object_mut()
237 .expect("the server entry is a JSON object");
238 object.insert("transport".into(), "stdio".into());
239 object.insert("runtime_id".into(), "local".into());
240 } else {
241 let object = entry
242 .as_object_mut()
243 .expect("the server entry is a JSON object");
244 object.insert("type".into(), "stdio".into());
245 }
246 configured.insert(server.name.clone(), entry);
247 }
248 let mut body = serde_json::to_vec_pretty(&document)?;
249 body.push(b'\n');
250 mj_core::config::atomic_write(&path, &body)
251 .with_context(|| format!("write staged reviewer configuration {}", path.display()))
252}
253
254fn review_dispatch_server(worker_root: &str) -> ReviewMcpServer {
258 let socket = format!(
259 "{worker_root}/{}/{}",
260 REVIEWER_DIR,
261 mj_core::review::mcp::REVIEW_DISPATCH_SOCKET
262 );
263 ReviewMcpServer {
264 name: mj_core::review::mcp::REVIEW_MCP_SERVER_NAME.to_owned(),
265 command: Path::new(worker_root).join("hel"),
266 args: vec![
267 "worker".to_owned(),
268 "review-mcp".to_owned(),
269 "--socket".to_owned(),
270 socket,
271 ],
272 }
273}
274
275fn reviewer_profile_home(worker_root: &str, generation: u64) -> String {
281 reviewer_staging_profile_home(Path::new(worker_root), generation)
282 .to_string_lossy()
283 .into_owned()
284}
285
286fn upload_reviewer_profile(
292 executor: &impl CommandExecutor,
293 locator: &targets::TargetLocator,
294 worker_root: &str,
295 generation: u64,
296 local: &Path,
297) -> Result<()> {
298 let home = reviewer_profile_home(worker_root, generation);
299 match locator {
300 targets::TargetLocator::LocalBare { .. } => {
301 for command in [
302 CommandSpec::new("rm", ["-rf", "--", &home])
303 .purpose("clear the local reviewer profile"),
304 CommandSpec::new("mkdir", ["-p", &home])
305 .purpose("create the local reviewer profile directory"),
306 CommandSpec::new(
307 "cp",
308 [
309 "-R".to_owned(),
310 format!("{}/.", local.display()),
311 home.clone(),
312 ],
313 )
314 .purpose("install the local reviewer profile"),
315 CommandSpec::new("chmod", ["-R", "go-rwx", &home])
316 .purpose("restrict local reviewer profile permissions"),
317 ] {
318 execute_checked(executor, command)?;
319 }
320 }
321 targets::TargetLocator::LocalPodman { container_id, .. }
322 | targets::TargetLocator::LocalDocker { container_id, .. }
323 | targets::TargetLocator::AppleContainer { container_id, .. } => {
324 let engine = match locator {
325 targets::TargetLocator::LocalPodman { .. } => "podman",
326 targets::TargetLocator::LocalDocker { .. } => "docker",
327 targets::TargetLocator::AppleContainer { .. } => "container",
328 _ => unreachable!("matched local container target"),
329 };
330 for arguments in [
331 vec![
332 "exec".to_owned(),
333 container_id.clone(),
334 "rm".to_owned(),
335 "-rf".to_owned(),
336 "--".to_owned(),
337 home.clone(),
338 ],
339 vec![
340 "exec".to_owned(),
341 container_id.clone(),
342 "mkdir".to_owned(),
343 "-p".to_owned(),
344 home.clone(),
345 ],
346 vec![
347 "cp".to_owned(),
348 format!("{}/.", local.display()),
349 format!("{container_id}:{home}"),
350 ],
351 container_upload_ownership_args(container_id, worker_root, &[&home]),
352 vec![
353 "exec".to_owned(),
354 container_id.clone(),
355 "chmod".to_owned(),
356 "-R".to_owned(),
357 "go-rwx".to_owned(),
358 home.clone(),
359 ],
360 ] {
361 execute_checked(
362 executor,
363 CommandSpec::new(engine, arguments).purpose("stage the reviewer profile"),
364 )?;
365 }
366 }
367 targets::TargetLocator::AwsEc2 { ssh, .. }
368 | targets::TargetLocator::SshBare { ssh, .. } => {
369 let incoming = format!("{home}.incoming");
370 execute_checked(
371 executor,
372 crate::targets::ssh_command(ssh, ["mkdir", "-p", worker_root])
373 .purpose("create the reviewer directory"),
374 )?;
375 execute_checked(
376 executor,
377 crate::targets::ssh_command(ssh, ["rm", "-rf", "--", &incoming, &home])
378 .purpose("clear the reviewer profile"),
379 )?;
380 execute_checked(
381 executor,
382 crate::targets::scp_upload(ssh, local, &incoming, true)
383 .purpose("upload the reviewer profile"),
384 )?;
385 execute_checked(
386 executor,
387 crate::targets::ssh_command(ssh, ["mv", &incoming, &home])
388 .purpose("install the reviewer profile"),
389 )?;
390 execute_checked(
391 executor,
392 crate::targets::ssh_command(ssh, ["chmod", "-R", "go-rwx", &home])
393 .purpose("restrict reviewer profile permissions"),
394 )?;
395 }
396 targets::TargetLocator::SshPodman {
397 ssh, container_id, ..
398 }
399 | targets::TargetLocator::SshDocker {
400 ssh, container_id, ..
401 } => {
402 let engine = match locator {
403 targets::TargetLocator::SshPodman { .. } => "podman",
404 targets::TargetLocator::SshDocker { .. } => "docker",
405 _ => unreachable!("matched remote container target"),
406 };
407 let upload = format!("{worker_root}/.reviewer-upload-{generation}");
408 execute_checked(
409 executor,
410 crate::targets::ssh_command(ssh, ["rm", "-rf", "--", &upload])
411 .purpose("clear remote reviewer staging"),
412 )?;
413 execute_checked(
414 executor,
415 crate::targets::scp_upload(ssh, local, &upload, true)
416 .purpose("upload the remote reviewer profile"),
417 )?;
418 for arguments in [
419 vec![
420 engine.to_owned(),
421 "exec".to_owned(),
422 container_id.clone(),
423 "rm".to_owned(),
424 "-rf".to_owned(),
425 "--".to_owned(),
426 home.clone(),
427 ],
428 vec![
429 engine.to_owned(),
430 "exec".to_owned(),
431 container_id.clone(),
432 "mkdir".to_owned(),
433 "-p".to_owned(),
434 home.clone(),
435 ],
436 vec![
437 engine.to_owned(),
438 "cp".to_owned(),
439 format!("{upload}/."),
440 format!("{container_id}:{home}"),
441 ],
442 std::iter::once(engine.to_owned())
443 .chain(container_upload_ownership_args(
444 container_id,
445 worker_root,
446 &[&home],
447 ))
448 .collect(),
449 vec![
450 engine.to_owned(),
451 "exec".to_owned(),
452 container_id.clone(),
453 "chmod".to_owned(),
454 "-R".to_owned(),
455 "go-rwx".to_owned(),
456 home.clone(),
457 ],
458 ] {
459 execute_checked(
460 executor,
461 crate::targets::ssh_command(ssh, arguments)
462 .purpose("stage the remote reviewer profile"),
463 )?;
464 }
465 execute_checked(
466 executor,
467 crate::targets::ssh_command(ssh, ["rm", "-rf", "--", &upload])
468 .purpose("remove remote reviewer staging"),
469 )?;
470 }
471 }
472 if home.trim().is_empty() {
473 bail!("the reviewer profile home resolved to an empty path");
474 }
475 Ok(())
476}
477
478#[cfg(test)]
479mod tests {
480 use std::cell::RefCell;
481 use std::collections::BTreeMap;
482
483 use super::*;
484 use crate::controller::test_support::checkpoint_test_session;
485 use mj_core::config::{Config, HarnessKind, HarnessProfile, TargetTemplate};
486 use mj_core::state::{SessionState, State};
487
488 use crate::targets::CommandOutput;
489
490 struct RecordingExecutor {
491 commands: RefCell<Vec<CommandSpec>>,
492 }
493
494 impl RecordingExecutor {
495 fn new() -> Self {
496 Self {
497 commands: RefCell::new(Vec::new()),
498 }
499 }
500
501 fn script(&self) -> Vec<String> {
503 self.commands
504 .borrow()
505 .iter()
506 .map(|command| format!("{} {}", command.program, command.args.join(" ")))
507 .collect()
508 }
509 }
510
511 impl CommandExecutor for RecordingExecutor {
512 fn execute(&self, command: &CommandSpec) -> Result<CommandOutput> {
513 self.commands.borrow_mut().push(command.clone());
514 Ok(CommandOutput {
515 status: 0,
516 stdout: Vec::new(),
517 stderr: Vec::new(),
518 })
519 }
520 }
521
522 const SESSION_ID: &str = "0123456789abcdef0123456789abcdef";
525
526 fn fixture(directory: &Path, locator: mj_core::state::TargetLocator) -> (Controller, String) {
527 let session_id = SESSION_ID;
528 let mut session = checkpoint_test_session(session_id);
529 session.target_template_id = "local".into();
530 session.state = SessionState::Running;
531 let template = match &locator {
532 mj_core::state::TargetLocator::LocalPodman { .. } => {
533 serde_json::from_str(r#"{"kind":"local-podman","image":"test"}"#).unwrap()
534 }
535 mj_core::state::TargetLocator::LocalDocker { .. } => {
536 serde_json::from_str(r#"{"kind":"local-docker","image":"test"}"#).unwrap()
537 }
538 _ => TargetTemplate::LocalBare,
539 };
540 session.target = Some(locator);
541 let mut config = Config::default();
542 config.targets.insert("local".into(), template);
543 for (id, kind) in [
544 ("codex", HarnessKind::Codex),
545 ("claude", HarnessKind::Claude),
546 ] {
547 let home = directory.join(id);
548 std::fs::create_dir_all(&home).unwrap();
549 config.profiles.insert(
550 id.to_owned(),
551 HarnessProfile {
552 enabled: true,
553 kind,
554 home,
555 environment: BTreeMap::from([("EXTRA".into(), "1".into())]),
556 context_window_bytes: None,
557 guardian_review_model: None,
558 },
559 );
560 }
561 (
562 Controller {
563 config,
564 state: State {
565 sessions: BTreeMap::from([(session_id.into(), session)]),
566 ..State::default()
567 },
568 },
569 session_id.to_owned(),
570 )
571 }
572
573 #[test]
574 fn staging_copies_the_chosen_profile_into_the_worker_root() {
575 let directory = tempfile::tempdir().unwrap();
576 let worker_root = directory.path().join(SESSION_ID);
577 std::fs::create_dir_all(directory.path().join("claude")).unwrap();
578 std::fs::write(directory.path().join("claude/CLAUDE.md"), b"reviewer").unwrap();
580 let (controller, session_id) = fixture(
581 directory.path(),
582 mj_core::state::TargetLocator::LocalBare {
583 worker_root: worker_root.clone(),
584 },
585 );
586 let executor = RecordingExecutor::new();
587
588 let config = controller
589 .stage_reviewer_profile_controlled(&session_id, "claude", 0, &[], &executor)
590 .unwrap();
591
592 assert_eq!(config.profile_id, "claude");
593 assert_eq!(config.harness, HarnessKind::Claude);
594 assert_eq!(config.generation, 0);
595 assert_eq!(config.model, None);
596 assert_eq!(config.effort, None);
597 assert!(
599 !config
600 .environment
601 .contains_key(HarnessKind::Claude.home_env())
602 );
603 assert_eq!(
604 config.environment.get("EXTRA").map(String::as_str),
605 Some("1")
606 );
607
608 let home = format!("{}/reviewer/profile", worker_root.display());
609 let script = executor.script();
610 let cleared = script
611 .iter()
612 .position(|line| line.starts_with("rm ") && line.contains(&home))
613 .expect("the previous reviewer profile is cleared");
614 let copied = script
615 .iter()
616 .position(|line| line.starts_with("cp ") && line.ends_with(&home))
617 .expect("the staged profile is installed");
618 assert!(
619 cleared < copied,
620 "a stale profile must go before the new one lands: {script:?}"
621 );
622 assert!(
623 script
624 .iter()
625 .any(|line| line.contains("go-rwx") && line.contains(&home)),
626 "the reviewer profile must not be world readable: {script:?}"
627 );
628 }
629
630 #[test]
631 fn local_container_targets_stage_the_reviewer_through_their_engine() {
632 let directory = tempfile::tempdir().unwrap();
633 let container_id = crate::targets::resource_name(SESSION_ID).unwrap();
634 for (locator, engine) in [
635 (
636 mj_core::state::TargetLocator::LocalPodman {
637 borrowed_from: None,
638 container_id: container_id.clone(),
639 workspace_storage: Default::default(),
640 },
641 "podman",
642 ),
643 (
644 mj_core::state::TargetLocator::LocalDocker {
645 borrowed_from: None,
646 container_id: container_id.clone(),
647 },
648 "docker",
649 ),
650 ] {
651 let (controller, session_id) = fixture(directory.path(), locator);
652 let executor = RecordingExecutor::new();
653
654 controller
655 .stage_reviewer_profile_controlled(&session_id, "codex", 3, &[], &executor)
656 .unwrap();
657
658 let script = executor.script();
659 assert!(
660 script
661 .iter()
662 .all(|line| line.starts_with(&format!("{engine} "))),
663 "a container target is reached only through its engine: {script:?}"
664 );
665 let home = script
666 .iter()
667 .find_map(|line| {
668 line.split(' ')
669 .find(|word| word.contains("/reviewer/profile"))
670 })
671 .expect("the reviewer profile is placed")
672 .to_owned();
673 assert!(
674 home.contains(&format!("/{session_id}")),
675 "the reviewer lives under this session's worker root: {home}"
676 );
677 assert!(
679 !script.iter().any(|line| {
680 line.contains("run") || line.contains("git") || line.contains("create")
681 }),
682 "staging a reviewer provisions nothing: {script:?}"
683 );
684 }
685 }
686
687 #[test]
688 fn an_unknown_profile_is_refused_before_anything_is_copied() {
689 let directory = tempfile::tempdir().unwrap();
690 let (controller, session_id) = fixture(
691 directory.path(),
692 mj_core::state::TargetLocator::LocalBare {
693 worker_root: directory.path().join(SESSION_ID),
694 },
695 );
696 let executor = RecordingExecutor::new();
697
698 let error = controller
699 .stage_reviewer_profile_controlled(&session_id, "missing", 0, &[], &executor)
700 .unwrap_err();
701
702 assert!(format!("{error:#}").contains("unknown profile"));
703 assert!(executor.commands.borrow().is_empty());
704 }
705
706 #[test]
707 fn a_new_generation_travels_to_the_worker_so_it_starts_a_fresh_reviewer() {
708 let directory = tempfile::tempdir().unwrap();
709 let (controller, session_id) = fixture(
710 directory.path(),
711 mj_core::state::TargetLocator::LocalBare {
712 worker_root: directory.path().join(SESSION_ID),
713 },
714 );
715 let executor = RecordingExecutor::new();
716
717 let first = controller
718 .stage_reviewer_profile_controlled(&session_id, "codex", 0, &[], &executor)
719 .unwrap();
720 let second = controller
721 .stage_reviewer_profile_controlled(&session_id, "codex", 1, &[], &executor)
722 .unwrap();
723
724 assert!(first.reusable_for(&first));
725 assert!(
726 !first.reusable_for(&second),
727 "a new generation must not reload the old conversation"
728 );
729 }
730}