magi-code 0.77.1

Repository-aware CLI coding agent for terminal work
Documentation
use super::*;

fn queued_messages(coordinator: &Coordinator, connection: &str) -> Vec<Value> {
    coordinator.connections[connection]
        .queue
        .iter()
        .map(|record| serde_json::from_str(record).unwrap())
        .collect()
}

#[test]
fn claim_snapshot_and_terminal_delivery_cover_both_sequence_boundaries() {
    for terminal_before_claim in [false, true] {
        let temp = TempDir::new().unwrap();
        std::fs::write(temp.path().join("evidence.txt"), "synthetic evidence").unwrap();
        let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
        let (release, wait) = bounded(1);
        let (cleanup, cleanup_wait) = bounded(1);
        install_tool_worker(
            &mut coordinator,
            wait,
            Arc::new(AtomicUsize::new(0)),
            cleanup_wait,
        );
        let first = connect(&mut coordinator);
        let session = create(&mut coordinator, &first);
        let grant = claim(&mut coordinator, &first, &session);
        let start = request(
            &coordinator,
            &first,
            "turn.start",
            Some(&session),
            Some(grant),
            Some("boundary-turn"),
            json!({"prompt":"Read the synthetic fixture"}),
        );
        let accepted = invoke(&mut coordinator, &first, start);
        assert_eq!(accepted["payload"]["status"], "accepted");
        coordinator.disconnect(&first);
        release.send(()).unwrap();
        run_until(&mut coordinator, |state| {
            state.actors[&session].phase == "cleaning"
        });
        if terminal_before_claim {
            cleanup.send(()).unwrap();
            run_until(&mut coordinator, |state| {
                !state.actors.contains_key(&session)
            });
        }
        let second = connect(&mut coordinator);
        let claim_request = request(
            &coordinator,
            &second,
            "session.claim",
            Some(&session),
            None,
            Some("boundary-claim"),
            json!({}),
        );
        // Keep the response queued so its order relative to terminal delivery is tested.
        coordinator
            .submit(
                &second,
                &serde_json::to_vec(&claim_request).unwrap(),
                Instant::now(),
            )
            .unwrap();
        if !terminal_before_claim {
            cleanup.send(()).unwrap();
            run_until(&mut coordinator, |state| {
                state.actors[&session].phase == "idle"
            });
        }
        coordinator.tick(Instant::now());
        let messages = queued_messages(&coordinator, &second);
        let reply = &messages[0];
        assert_eq!(reply["request_id"], claim_request.request_id);
        assert!(reply["error"].is_null(), "{reply}");
        let snapshot = &reply["payload"]["snapshot"];
        let boundary = snapshot["live_sequence"].as_u64().unwrap();
        let terminal = &coordinator.terminals[&session];
        if terminal_before_claim {
            assert_eq!(
                messages.len(),
                1,
                "settled terminal must not be redelivered"
            );
            assert_eq!(snapshot["phase"], "idle");
            assert!(snapshot["turn"].is_null());
            assert_eq!(snapshot["terminal"], terminal.payload);
            assert_eq!(boundary, terminal.sequence);
        } else {
            assert_eq!(
                messages.len(),
                2,
                "claim must deliver exactly one later terminal"
            );
            assert_eq!(snapshot["phase"], "cleaning");
            assert!(snapshot["terminal"].is_null());
            let event = &messages[1];
            assert_eq!(event["event"], "turn.terminal");
            assert_eq!(event["live_sequence"], boundary + 1);
            assert_eq!(event["live_sequence"], terminal.sequence);
            assert_eq!(event["grant_generation"], reply["payload"]["generation"]);
            assert_eq!(event["payload"]["turn_id"], snapshot["turn"]["turn_id"]);
            assert_eq!(
                event["payload"]["assistant_text"],
                "before tool\n\nafter tool"
            );
            assert_eq!(event["payload"]["status"], "completed");
        }
        coordinator.disconnect(&second);
    }
}

#[test]
fn committed_login_survives_initiating_disconnect_and_lookup_reports_cleanup() {
    let temp = TempDir::new().unwrap();
    let runtime = runtime(&temp);
    let paths = runtime.config.paths.clone();
    let worker_paths = paths.clone();
    let mut coordinator = Coordinator::new(runtime).unwrap();
    let (committed, commit_observed) = bounded(1);
    let (release, wait) = bounded(1);
    coordinator.execution.set_login_worker(Arc::new(move |job| {
        let generation = crate::config::read_auth_store(&worker_paths)
            .unwrap()
            .provider_generation("openai-codex");
        crate::login::complete_service_codex_login(
            &worker_paths,
            generation,
            &job.cancel,
            || {
                Ok(crate::config::NormalizedToken {
                    access: "synthetic-access".into(),
                    refresh: Some("synthetic-refresh".into()),
                    expires: Some(i64::MAX),
                    account_id: "synthetic-account".into(),
                })
            },
            || {
                committed.send(Arc::clone(&job.cancel)).unwrap();
                wait.recv_timeout(Duration::from_secs(5)).unwrap();
            },
        )
    }));
    let owner = connect(&mut coordinator);
    let start = request(
        &coordinator,
        &owner,
        "auth.login.start",
        None,
        None,
        Some("committed-login"),
        json!({"provider_id":"openai-codex"}),
    );
    assert!(invoke(&mut coordinator, &owner, start)["error"].is_null());
    let cancel = commit_observed
        .recv_timeout(Duration::from_secs(2))
        .unwrap();
    let expected_credentials = crate::config::AuthProviderRecord::OAuth {
        access: "synthetic-access".into(),
        refresh: Some("synthetic-refresh".into()),
        expires: Some(i64::MAX),
        account_id: Some("synthetic-account".into()),
    };
    assert_eq!(
        crate::config::read_auth(&paths)
            .unwrap()
            .providers
            .get("openai-codex"),
        Some(&expected_credentials),
    );
    coordinator.disconnect(&owner);
    run_until(&mut coordinator, |_| cancel.load(Ordering::Relaxed));
    assert_eq!(
        coordinator.operations.lookup(
            &coordinator.instance,
            &coordinator.instance,
            "committed-login"
        )["state"],
        "accepted",
        "commit alone does not establish cleanup"
    );
    release.send(()).unwrap();
    run_until(&mut coordinator, |state| state.login.is_none());
    let observer = connect(&mut coordinator);
    let lookup = request(
        &coordinator,
        &observer,
        "operation.lookup",
        None,
        None,
        None,
        json!({"target_instance_id":coordinator.instance,"operation_id":"committed-login"}),
    );
    let reply = invoke(&mut coordinator, &observer, lookup);
    assert!(reply["error"].is_null(), "{reply}");
    assert_eq!(reply["payload"]["state"], "terminal");
    assert_eq!(reply["payload"]["result"]["status"], "completed");
    assert_eq!(reply["payload"]["result"]["cleanup_complete"], true);
    assert!(reply["payload"]["error"].is_null());
    assert_eq!(
        crate::config::read_auth(&paths)
            .unwrap()
            .providers
            .get("openai-codex"),
        Some(&expected_credentials),
        "disconnect after commit must not remove the stored login",
    );
}

#[test]
fn standalone_writer_waits_for_daemon_cleanup_and_control_release() {
    let temp = TempDir::new().unwrap();
    std::fs::write(temp.path().join("evidence.txt"), "synthetic evidence").unwrap();
    let runtime = runtime(&temp);
    let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
    let (release, wait) = bounded(1);
    let (cleanup, cleanup_wait) = bounded(1);
    install_tool_worker(
        &mut coordinator,
        wait,
        Arc::new(AtomicUsize::new(0)),
        cleanup_wait,
    );
    let owner = connect(&mut coordinator);
    let session = create(&mut coordinator, &owner);
    let separate = create(&mut coordinator, &owner);
    let disk = runtime.session_manager.open_existing(&session).unwrap();
    let separate_disk = runtime.session_manager.open_existing(&separate).unwrap();
    let grant = claim(&mut coordinator, &owner, &session);
    assert!(
        disk.try_frontend_writer().unwrap().is_none(),
        "idle claim owns writer"
    );
    let separate_writer = separate_disk.try_frontend_writer().unwrap().unwrap();
    let start = request(
        &coordinator,
        &owner,
        "turn.start",
        Some(&session),
        Some(grant),
        Some("writer-handoff"),
        json!({"prompt":"Read the synthetic fixture"}),
    );
    assert_eq!(
        invoke(&mut coordinator, &owner, start)["payload"]["status"],
        "accepted"
    );
    assert_eq!(coordinator.actors[&session].phase, "running");
    assert!(disk.try_frontend_writer().unwrap().is_none());
    release.send(()).unwrap();
    run_until(&mut coordinator, |state| {
        state.actors[&session].phase == "cleaning"
    });
    assert!(disk.try_frontend_writer().unwrap().is_none());
    cleanup.send(()).unwrap();
    run_until(&mut coordinator, |state| {
        state.actors[&session].phase == "idle"
    });
    assert!(
        disk.try_frontend_writer().unwrap().is_none(),
        "cleanup does not revoke control"
    );
    coordinator.disconnect(&owner);
    assert!(!coordinator.actors.contains_key(&session));
    let standalone_writer = disk.try_frontend_writer().unwrap().unwrap();
    assert!(separate_disk.try_frontend_writer().unwrap().is_none());
    drop(standalone_writer);
    drop(separate_writer);
}