mj_controller/
session_manager.rs1use 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
37const SESSION_SYNC_INTERVAL: Duration = Duration::from_millis(150);
39const QUIET_SESSION_SYNC_INTERVAL: Duration = Duration::from_secs(2);
41const PROJECTION_TRANSACTION_EVENT_BUDGET: usize = 128;
45const RECONNECT_INTERVAL: Duration = Duration::from_secs(1);
46const 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;