Skip to main content

mj_controller/
session_manager.rs

1//! Multiplexed controller-side ownership of durable ACP relay sessions.
2
3use std::collections::{BTreeMap, VecDeque};
4use std::path::PathBuf;
5use std::sync::{Arc, Mutex};
6use std::time::{Duration, Instant};
7
8use anyhow::{Context, Result, bail, ensure};
9use tokio::sync::{mpsc, oneshot, watch};
10
11use crate::database::{
12    ProjectionApplyOutcome, ProjectionIntegrityError, apply_projection_page,
13    save_materialized_session,
14};
15use crate::worker_client::{RelayClient, RelayEventPage, RelayRejected, RelayTransportDead};
16use mj_checkpoint::archive::verify_archive_streaming;
17use mj_core::credentials::{CredentialSyncSignal, relay_event_credential_sync_reason};
18use mj_core::elicitation::ElicitationResponse;
19use mj_core::state::{ManagedSessionSnapshot, MaterializedSession};
20use mj_transcript::projection::{
21    ProjectionIndex, apply_committed_projection_event_indexed, materialized_session_from_canonical,
22    project_relay_event_indexed,
23};
24
25use crate::targets::{
26    CancellableProcessExecutor, CommandExecutor, CommandPlan, CommandSpec, TargetLocator,
27    TargetRecoveryOutcome, TargetRecoveryPlan, ensure_recovery_target_running,
28};
29use mj_core::relay::{RelayCommand, RelayCursor, RelayOperationalState};
30
31pub use mj_client::session::{
32    ManagedSessionView, ReviewerAction, ReviewerOutcome, ViewError, new_command_id,
33};
34#[cfg(test)]
35use mj_core::worker_launch::ReviewerLaunchConfig;
36
37/// Sync cadence while the worker owns work, so streamed events show promptly.
38const SESSION_SYNC_INTERVAL: Duration = Duration::from_millis(150);
39/// Sync cadence while the worker is quiet. See `actor::sync_delay`.
40const QUIET_SESSION_SYNC_INTERVAL: Duration = Duration::from_secs(2);
41/// Release SQLite's single writer between bounded pieces of a large relay
42/// catch-up. One transport page can contain thousands of terminal events and
43/// must not prevent every other session actor from publishing its view.
44const PROJECTION_TRANSACTION_EVENT_BUDGET: usize = 128;
45const RECONNECT_INTERVAL: Duration = Duration::from_secs(1);
46/// Ceiling for reconnect backoff. A worker that exited stays gone until the
47/// user acts, so retrying it every second only burns process spawns.
48const RECONNECT_BACKOFF_CEILING: Duration = Duration::from_secs(30);
49const UNREACHABLE_FAILURE_THRESHOLD: u32 = 2;
50const WORKER_RESTART_TIMEOUT: Duration = Duration::from_secs(30);
51const WORKER_RESTART_COOLDOWN: Duration = Duration::from_secs(60);
52const SESSION_MANAGER_SHUTDOWN_GRACE: Duration = Duration::from_millis(750);
53
54mod types;
55pub use types::*;
56mod recovery;
57pub(crate) use recovery::*;
58mod delegation;
59pub(crate) use delegation::*;
60mod channels;
61pub use channels::*;
62mod handle;
63pub use handle::*;
64mod client_backend;
65use client_backend::*;
66mod actor_types;
67use actor_types::*;
68pub use actor_types::{RelayConnectionJob, RelayJobDeferred};
69mod remote;
70pub use remote::*;
71mod spawn;
72pub use spawn::*;
73mod actor;
74use actor::*;
75mod standalone;
76pub use standalone::*;
77
78#[cfg(test)]
79mod tests;