use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex, MutexGuard};
use std::thread;
use ohno::AppError;
use crate::constants::{
CONNECT_TIMEOUT, DEFAULT_PTY_COLS, DEFAULT_PTY_ROWS, MAX_CLIENT_BACKLOG_BYTES,
MAX_OUTPUT_CHUNK_BYTES,
};
use crate::outbox::Outbox;
use crate::pal::error::{PalError, PalErrorKind};
use crate::pal::ids::{AppId, ConnId, JobId, ListenerId, PtyId};
use crate::pal::processes::{AppSpawn, Processes};
use crate::pal::pseudoconsole::{Pseudoconsole, WindowSize};
use crate::pal::session_store::SessionStore;
use crate::pal::transport::Transport;
use crate::protocol::Message;
use crate::session_id::SessionId;
use crate::session_record::{ProcessIdentity, SessionRecord};
use crate::wall_clock::unix_now_ms;
use crate::{BreakawayDeniedError, PalFailedError, StartupFailedError, StoreError};
const DEFAULT_PTY_SIZE: WindowSize = WindowSize {
cols: DEFAULT_PTY_COLS,
rows: DEFAULT_PTY_ROWS,
};
struct InitGuard<'a, P: Processes, S: SessionStore, T: Transport, C: Pseudoconsole> {
processes: &'a P,
store: &'a S,
transport: &'a T,
pty_host: &'a C,
job: Option<JobId>,
pty: Option<PtyId>,
listener: Option<ListenerId>,
session: Option<(SessionId, ProcessIdentity)>,
committed: bool,
}
impl<P: Processes, S: SessionStore, T: Transport, C: Pseudoconsole> Drop
for InitGuard<'_, P, S, T, C>
{
fn drop(&mut self) {
if self.committed {
return;
}
if let Some(listener) = self.listener {
self.transport.close_listener(listener);
}
if let Some((id, owner)) = self.session {
_ = self.store.delete_owned_by(id, &owner);
}
if let Some(job) = self.job {
self.processes.close_job(job);
}
if let Some(pty) = self.pty {
self.pty_host.close(pty);
}
}
}
#[cfg_attr(test, mutants::skip)]
pub(crate) fn run_supervisor<P, S, T, C>(
processes: &P,
store: &S,
transport: &T,
pty_host: &C,
startup_pipe: &str,
launch_directory: PathBuf,
command: Vec<String>,
) -> Result<i32, AppError>
where
P: Processes,
S: SessionStore + Clone,
T: Transport + Clone + Send + Sync + 'static,
C: Pseudoconsole + Clone + Send + Sync + 'static,
{
let startup = transport
.connect(startup_pipe, CONNECT_TIMEOUT)
.map_err(|_error| StartupFailedError::new())?;
let mut guard = InitGuard {
processes,
store,
transport,
pty_host,
job: None,
pty: None,
listener: None,
session: None,
committed: false,
};
let result = initialize(
&mut guard,
processes,
store,
transport,
pty_host,
launch_directory,
command,
);
let initialized = match result {
Ok(initialized) => initialized,
Err(error) => {
_ = transport.send(startup, &Message::StartupErr);
transport.disconnect(startup);
return Err(error);
}
};
let startup_ok = Message::StartupOk {
session_id: initialized.session_id,
durability: processes.durability(),
};
if transport.send(startup, &startup_ok).is_err() {
transport.disconnect(startup);
return Err(StartupFailedError::new().into());
}
let committed = transport.recv_timeout(startup, CONNECT_TIMEOUT);
if !matches!(committed, Ok(Message::StartupCommit)) {
transport.disconnect(startup);
return Err(StartupFailedError::new().into());
}
guard.committed = true;
let status = serve(processes, store, transport, pty_host, &initialized, startup)?;
Ok(status)
}
#[derive(Clone, Copy)]
struct Initialized {
session_id: SessionId,
identity: ProcessIdentity,
listener: ListenerId,
pty: PtyId,
job: JobId,
app: AppId,
}
fn initialize<P, S, T, C>(
guard: &mut InitGuard<'_, P, S, T, C>,
processes: &P,
store: &S,
transport: &T,
pty_host: &C,
launch_directory: PathBuf,
command: Vec<String>,
) -> Result<Initialized, AppError>
where
P: Processes,
S: SessionStore,
T: Transport,
C: Pseudoconsole,
{
let job = processes
.create_lifetime_job()
.map_err(|error| map_startup(&error))?;
guard.job = Some(job);
let pty = pty_host
.create(DEFAULT_PTY_SIZE)
.map_err(|error| map_startup(&error))?;
guard.pty = Some(pty);
let app = processes
.spawn_app(&AppSpawn {
command: command.clone(),
launch_directory: launch_directory.clone(),
pty,
job,
})
.map_err(|error| map_startup(&error))?;
let nonce = processes.random_nonce();
let pipe_name = transport.pipe_name(&nonce);
let listener = transport
.listen(&pipe_name)
.map_err(|error| map_startup(&error))?;
guard.listener = Some(listener);
let identity = processes
.current_identity()
.map_err(|error| map_startup(&error))?;
let session_id = store
.allocate_id(&identity)
.map_err(|_error| StoreError::new())?;
guard.session = Some((session_id, identity));
let record = SessionRecord {
id: session_id.get(),
supervisor_pid: identity.pid,
supervisor_creation_time: identity.creation_time,
pipe_name,
launch_directory,
command,
started_at_unix_ms: unix_now_ms(),
attached: false,
};
store.publish(&record).map_err(|_error| StoreError::new())?;
Ok(Initialized {
session_id,
identity,
listener,
pty,
job,
app,
})
}
fn map_startup(error: &PalError) -> AppError {
match error.kind() {
PalErrorKind::BreakawayDenied => BreakawayDeniedError::new().into(),
_ => StartupFailedError::new().into(),
}
}
struct Shared<T: Transport, C> {
transport: T,
pty_host: C,
pty: PtyId,
session_id: SessionId,
client: Mutex<Option<Client<T>>>,
attach: Mutex<()>,
attached_generation: Arc<AtomicU64>,
preamble: Mutex<Option<Vec<u8>>>,
first_attach: Mutex<FirstAttach>,
first_attach_changed: Condvar,
stopping: AtomicBool,
}
struct Client<T: Transport> {
conn: ConnId,
outbox: Arc<Outbox<T>>,
}
impl<T: Transport> Clone for Client<T> {
fn clone(&self) -> Self {
Self {
conn: self.conn,
outbox: Arc::clone(&self.outbox),
}
}
}
#[derive(Debug, Default)]
struct FirstAttach {
attached: bool,
initiator_gone: bool,
}
impl<T: Transport, C> Shared<T, C> {
fn first_attach(&self) -> MutexGuard<'_, FirstAttach> {
self.first_attach
.lock()
.expect("first-attach flags are only set, never held across a panic")
}
#[cfg_attr(test, mutants::skip)]
fn note_attached(&self) {
self.first_attach().attached = true;
self.first_attach_changed.notify_all();
}
#[cfg_attr(test, mutants::skip)]
fn note_initiator_gone(&self) {
self.first_attach().initiator_gone = true;
self.first_attach_changed.notify_all();
}
#[cfg_attr(test, mutants::skip)]
fn await_first_attach(&self) {
let mut state = self.first_attach();
while !state.attached && !state.initiator_gone {
state = self
.first_attach_changed
.wait(state)
.expect("first-attach flags are only set, never held across a panic");
}
}
fn client(&self) -> MutexGuard<'_, Option<Client<T>>> {
self.client
.lock()
.expect("client slot is only copied or replaced, never held across a panic")
}
#[cfg_attr(test, mutants::skip)]
fn next_attached_generation(&self) -> u64 {
let previous = self
.attached_generation
.try_update(Ordering::SeqCst, Ordering::SeqCst, |generation| {
generation.checked_add(1)
})
.expect("the process cannot perform enough ownership changes to exhaust u64");
previous
.checked_add(1)
.expect("try_update only succeeds when the next generation exists")
}
fn hold_for_first_client(&self, bytes: &[u8]) {
let mut preamble = self
.preamble
.lock()
.expect("the preamble is only appended to or taken, never held across a panic");
let Some(held) = preamble.as_mut() else {
return;
};
let free = MAX_CLIENT_BACKLOG_BYTES.saturating_sub(held.len());
held.extend(bytes.iter().take(free));
}
fn take_preamble(&self) -> Option<Vec<u8>> {
self.preamble
.lock()
.expect("the preamble is only appended to or taken, never held across a panic")
.take()
.filter(|held| !held.is_empty())
}
}
fn preamble_messages(held: &[u8]) -> impl Iterator<Item = Message> + use<'_> {
held.chunks(MAX_OUTPUT_CHUNK_BYTES.get())
.map(|chunk| Message::Output(chunk.to_vec()))
}
#[cfg_attr(test, mutants::skip)]
fn serve<P, S, T, C>(
processes: &P,
store: &S,
transport: &T,
pty_host: &C,
initialized: &Initialized,
startup: ConnId,
) -> Result<i32, AppError>
where
P: Processes,
S: SessionStore + Clone,
T: Transport + Clone + Send + Sync + 'static,
C: Pseudoconsole + Clone + Send + Sync + 'static,
{
let Initialized {
session_id,
identity,
listener,
pty,
job,
app,
} = *initialized;
let record_live = Arc::new(Mutex::new(true));
let attached_generation = Arc::new(AtomicU64::default());
let shared = Arc::new(Shared {
transport: transport.clone(),
pty_host: pty_host.clone(),
pty,
session_id,
client: Mutex::new(None),
attach: Mutex::new(()),
attached_generation: Arc::clone(&attached_generation),
preamble: Mutex::new(Some(Vec::new())),
first_attach: Mutex::new(FirstAttach::default()),
first_attach_changed: Condvar::new(),
stopping: AtomicBool::new(false),
});
thread::spawn({
let shared = Arc::clone(&shared);
let transport = transport.clone();
move || {
_ = transport.recv(startup);
transport.disconnect(startup);
shared.note_initiator_gone();
}
});
let store_flag = store_attached_flag(
store,
session_id,
Arc::clone(&record_live),
attached_generation,
);
thread::spawn({
let shared = Arc::clone(&shared);
let transport = transport.clone();
move || accept_loop(&shared, &transport, listener, store_flag)
});
let pty_pump = thread::spawn({
let shared = Arc::clone(&shared);
move || pty_output_loop(&shared)
});
let waited = processes.wait_app(app);
if waited.is_ok() {
shared.await_first_attach();
}
transport.close_listener(listener);
processes.close_job(job);
pty_host.finish(pty);
_ = pty_pump.join();
pty_host.close(pty);
let client = {
let _attach = shared
.attach
.lock()
.expect("the attach lock guards no data, so it is never poisoned by its guard");
shared.stopping.store(true, Ordering::SeqCst);
shared.client().take()
};
if let Some(client) = &client {
if let Ok(status) = &waited {
client.outbox.send(Message::AppExited { status: *status });
}
client.outbox.finish();
}
let deleted = {
let mut live = record_live
.lock()
.expect("record_live is only set false here, never held across a panic");
*live = false;
store.delete_owned_by(session_id, &identity)
};
let status = waited.map_err(|_error| PalFailedError::new())?;
deleted.map_err(|_error| StoreError::new())?;
if let Some(client) = client {
client.outbox.flush();
}
Ok(status)
}
#[cfg_attr(test, mutants::skip)]
fn store_attached_flag<S: SessionStore + Clone>(
store: &S,
id: SessionId,
record_live: Arc<Mutex<bool>>,
current_generation: Arc<AtomicU64>,
) -> impl Fn(u64, bool) + Clone + Send + 'static {
let store = store.clone();
move |generation: u64, attached: bool| {
let live = record_live
.lock()
.expect("record_live is only set false at delete, never held across a panic");
if !*live || current_generation.load(Ordering::SeqCst) != generation {
return;
}
if let Ok(Some(mut record)) = store.read(id) {
record.attached = attached;
_ = store.publish(&record);
}
}
}
#[cfg_attr(test, mutants::skip)]
fn accept_loop<T, C>(
shared: &Arc<Shared<T, C>>,
transport: &T,
listener: ListenerId,
set_attached: impl Fn(u64, bool) + Clone + Send + 'static,
) where
T: Transport + Clone,
C: Pseudoconsole,
{
while !shared.stopping.load(Ordering::SeqCst) {
let Ok(conn) = transport.accept(listener) else {
break;
};
thread::spawn({
let shared = Arc::clone(shared);
let set_attached = set_attached.clone();
move || client_loop(&shared, conn, &set_attached)
});
}
}
#[cfg_attr(test, mutants::skip)]
fn client_loop<T, C>(shared: &Shared<T, C>, conn: ConnId, set_attached: &impl Fn(u64, bool))
where
T: Transport + Clone,
C: Pseudoconsole,
{
match shared.transport.recv(conn) {
Ok(Message::Attach { cols, rows }) => {
let _attach = shared
.attach
.lock()
.expect("the attach lock guards no data, so it is never poisoned by its guard");
if shared.stopping.load(Ordering::SeqCst) {
shared.transport.disconnect(conn);
return;
}
let outbox = Outbox::start(shared.transport.clone(), conn);
let (previous, generation) = {
let mut slot = shared.client();
if shared
.transport
.send(
conn,
&Message::Attached {
session_id: shared.session_id,
},
)
.is_err()
{
drop(slot);
outbox.abandon();
return;
}
if let Some(held) = shared.take_preamble() {
for message in preamble_messages(&held) {
outbox.send(message);
}
}
let previous = slot.replace(Client {
conn,
outbox: Arc::clone(&outbox),
});
let generation = shared.next_attached_generation();
(previous, generation)
};
shared.note_attached();
if let Some(old) = previous {
old.outbox.send(Message::Displaced);
old.outbox.finish();
}
_ = shared
.pty_host
.resize(shared.pty, WindowSize { cols, rows });
drop(_attach);
set_attached(generation, true);
}
_ => {
shared.transport.disconnect(conn);
return;
}
}
while let Ok(message) = shared.transport.recv(conn) {
let slot = shared.client();
if slot.as_ref().map(|client| client.conn) != Some(conn) {
break;
}
match message {
Message::Input(data) => {
_ = shared.pty_host.write_input(shared.pty, &data);
}
Message::Resize { cols, rows } => {
_ = shared
.pty_host
.resize(shared.pty, WindowSize { cols, rows });
}
_ => break,
}
drop(slot);
}
let (departing, generation) = {
let mut slot = shared.client();
if slot.as_ref().map(|client| client.conn) == Some(conn) {
let departing = slot.take();
let generation = shared.next_attached_generation();
(departing, Some(generation))
} else {
(None, None)
}
};
if let Some(generation) = generation {
set_attached(generation, false);
}
if let Some(departing) = departing {
departing.outbox.abandon();
}
}
#[cfg_attr(test, mutants::skip)]
fn pty_output_loop<T, C>(shared: &Shared<T, C>)
where
T: Transport + Clone,
C: Pseudoconsole,
{
while !shared.stopping.load(Ordering::SeqCst) {
let Ok(bytes) = shared.pty_host.read_output(shared.pty) else {
break;
};
let client = {
let slot = shared.client();
match slot.as_ref() {
Some(client) => Some(client.clone()),
None => {
shared.hold_for_first_client(&bytes);
None
}
}
};
if let Some(client) = client {
client.outbox.send(Message::Output(bytes));
}
}
}
#[cfg(test)]
#[cfg_attr(coverage_nightly, coverage(off))]
mod tests {
use std::path::PathBuf;
use std::sync::{Arc, Condvar, Mutex, mpsc};
use std::thread;
use memory_session_store::MemorySessionStore;
use testing::{WatchdogPhaseReporter, with_watchdog_phases};
use super::*;
use crate::durability::Durability;
use crate::pal::ids::{AppId, ConnId, JobId, ListenerId};
use crate::pal::processes::MockProcesses;
use crate::pal::pseudoconsole::MemoryPseudoconsole;
use crate::pal::transport::MemoryTransport;
use crate::protocol::{Message, encode, payload_len_ok};
use crate::session_record::ProcessIdentity;
const SAMPLE_APP_EXIT: i32 = 7;
const ORDINARY_ATTACH: Message = Message::Attach { cols: 80, rows: 24 };
fn commit_startup(
transport: &MemoryTransport,
listener: ListenerId,
phase_reporter: &WatchdogPhaseReporter,
) -> (ConnId, SessionId, Durability) {
phase_reporter.report("waiting for the supervisor startup connection");
let conn = transport.accept(listener).unwrap();
phase_reporter.report("waiting for the supervisor startup response");
let Message::StartupOk {
session_id,
durability,
} = transport.recv(conn).unwrap()
else {
panic!("expected startup ok");
};
transport.send(conn, &Message::StartupCommit).unwrap();
(conn, session_id, durability)
}
fn mock_processes(exit: Arc<(Mutex<bool>, Condvar)>) -> MockProcesses {
mock_processes_with(exit, Durability::Durable, AppWait::Reports)
}
#[derive(Clone, Copy)]
enum AppWait {
Reports,
Fails,
}
fn mock_processes_with(
exit: Arc<(Mutex<bool>, Condvar)>,
durability: Durability,
wait: AppWait,
) -> MockProcesses {
mock_processes_with_close_job(exit, durability, wait, |_| {})
}
fn mock_processes_with_close_job(
exit: Arc<(Mutex<bool>, Condvar)>,
durability: Durability,
wait: AppWait,
close_job: impl Fn(JobId) + Send + Sync + 'static,
) -> MockProcesses {
let mut processes = MockProcesses::new();
processes.expect_durability().returning(move || durability);
processes
.expect_create_lifetime_job()
.returning(|| Ok(JobId(1)));
processes.expect_spawn_app().returning(|_| Ok(AppId(1)));
processes.expect_current_identity().returning(|| {
Ok(ProcessIdentity {
pid: 10,
creation_time: 100,
})
});
processes
.expect_random_nonce()
.returning(|| "nonce".to_string());
processes.expect_close_job().returning(close_job);
processes.expect_wait_app().returning(move |_| {
let (lock, cvar) = &*exit;
let mut done = lock.lock().expect("exit lock");
while !*done {
done = cvar.wait(done).expect("exit wait");
}
match wait {
AppWait::Reports => Ok(SAMPLE_APP_EXIT),
AppWait::Fails => Err(PalError::new(PalErrorKind::Other)),
}
});
processes
}
#[test]
#[cfg_attr(miri, ignore)]
fn steal_displaces_first_client() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(false), Condvar::new()));
let processes = mock_processes(Arc::clone(&exit));
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let store = store.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
transport.disconnect(startup_conn);
let pipe = transport.pipe_name("nonce");
let first = transport.connect(&pipe, CONNECT_TIMEOUT).unwrap();
transport.send(first, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the first attach acknowledgement");
assert!(matches!(
transport.recv(first).unwrap(),
Message::Attached { session_id: id } if id == session_id
));
let second = transport.connect(&pipe, CONNECT_TIMEOUT).unwrap();
transport.send(second, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the second attach acknowledgement");
assert!(matches!(
transport.recv(second).unwrap(),
Message::Attached { .. }
));
phase_reporter.report("waiting for the displacement notice");
assert!(matches!(transport.recv(first).unwrap(), Message::Displaced));
{
let (lock, cvar) = &*exit;
*lock.lock().expect("exit lock") = true;
cvar.notify_all();
}
phase_reporter.report("waiting for supervisor shutdown");
assert_eq!(supervisor.join().unwrap().unwrap(), SAMPLE_APP_EXIT);
assert!(store.list().unwrap().is_empty());
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn final_output_arrives_before_the_exit_status() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(false), Condvar::new()));
let processes = mock_processes(Arc::clone(&exit));
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let pty = pty.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
transport.disconnect(startup_conn);
let client = transport
.connect(&transport.pipe_name("nonce"), CONNECT_TIMEOUT)
.unwrap();
transport.send(client, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the attach acknowledgement");
assert!(matches!(
transport.recv(client).unwrap(),
Message::Attached { .. }
));
pty.withhold_output(PtyId(1));
pty.push_output(PtyId(1), b"bye");
{
let (lock, cvar) = &*exit;
*lock.lock().expect("exit lock") = true;
cvar.notify_all();
}
phase_reporter.report("waiting for final app output");
assert_eq!(
transport.recv(client).unwrap(),
Message::Output(b"bye".to_vec())
);
phase_reporter.report("waiting for the app exit status");
assert_eq!(
transport.recv(client).unwrap(),
Message::AppExited {
status: SAMPLE_APP_EXIT
}
);
phase_reporter.report("waiting for supervisor shutdown");
assert_eq!(supervisor.join().unwrap().unwrap(), SAMPLE_APP_EXIT);
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn an_app_that_exits_before_anyone_attaches_still_reports_its_status() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(true), Condvar::new()));
let processes = mock_processes(Arc::clone(&exit));
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
let client = transport
.connect(&transport.pipe_name("nonce"), CONNECT_TIMEOUT)
.unwrap();
transport.send(client, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the attach acknowledgement");
assert!(matches!(
transport.recv(client).unwrap(),
Message::Attached { .. }
));
phase_reporter.report("waiting for the app exit status");
assert_eq!(
transport.recv(client).unwrap(),
Message::AppExited {
status: SAMPLE_APP_EXIT
}
);
phase_reporter.report("waiting for supervisor shutdown");
assert_eq!(supervisor.join().unwrap().unwrap(), SAMPLE_APP_EXIT);
transport.disconnect(startup_conn);
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_stalled_attached_flag_does_not_delay_the_exit_status() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(true), Condvar::new()));
let processes =
mock_processes_with_close_job(exit, Durability::Durable, AppWait::Reports, {
let store = store.clone();
move |_| store.wait_for_stalled_publish()
});
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let store = store.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
store.stall_publishes();
let client = transport
.connect(&transport.pipe_name("nonce"), CONNECT_TIMEOUT)
.unwrap();
transport.send(client, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the attach acknowledgement");
assert!(matches!(
transport.recv(client).unwrap(),
Message::Attached { .. }
));
phase_reporter.report("waiting for the attached-flag publication to stall");
store.wait_for_stalled_publish();
phase_reporter.report("waiting for the app exit status");
assert_eq!(
transport.recv(client).unwrap(),
Message::AppExited {
status: SAMPLE_APP_EXIT
}
);
store.resume_publishes();
transport.disconnect(client);
phase_reporter.report("waiting for supervisor shutdown");
assert_eq!(supervisor.join().unwrap().unwrap(), SAMPLE_APP_EXIT);
transport.disconnect(startup_conn);
assert!(store.list().unwrap().is_empty());
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_session_nobody_comes_for_ends_when_its_initiator_gives_up() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(true), Condvar::new()));
let processes = mock_processes(Arc::clone(&exit));
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let store = store.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
transport.disconnect(startup_conn);
phase_reporter.report("waiting for supervisor shutdown");
assert_eq!(supervisor.join().unwrap().unwrap(), SAMPLE_APP_EXIT);
assert!(store.list().unwrap().is_empty());
});
}
enum RejectedStartup {
Disconnect,
Message(Message),
Timeout,
}
fn assert_rejected_startup_rolls_back(rejection: RejectedStartup) {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(true), Condvar::new()));
let processes = mock_processes(exit);
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let store = store.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
phase_reporter.report("waiting for the supervisor startup connection");
let startup_conn = transport.accept(startup).unwrap();
phase_reporter.report("waiting for the supervisor startup response");
assert!(matches!(
transport.recv(startup_conn).unwrap(),
Message::StartupOk { .. }
));
let records = store.list().unwrap();
let [record] = records.as_slice() else {
panic!("expected exactly one record");
};
let session_pipe = record.pipe_name.clone();
match rejection {
RejectedStartup::Disconnect => transport.disconnect(startup_conn),
RejectedStartup::Message(acknowledgement) => {
transport.send(startup_conn, &acknowledgement).unwrap();
}
RejectedStartup::Timeout => transport.timeout_next_recv(),
}
phase_reporter.report("waiting for startup rollback");
let error = supervisor.join().unwrap().unwrap_err();
assert!(error.find_source::<StartupFailedError>().is_some());
assert!(store.list().unwrap().is_empty());
transport
.connect(&session_pipe, CONNECT_TIMEOUT)
.unwrap_err();
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_session_without_startup_commit_is_rolled_back() {
assert_rejected_startup_rolls_back(RejectedStartup::Disconnect);
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_session_with_an_invalid_startup_commit_is_rolled_back() {
assert_rejected_startup_rolls_back(RejectedStartup::Message(Message::StartupErr));
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_startup_commit_timeout_is_rolled_back() {
assert_rejected_startup_rolls_back(RejectedStartup::Timeout);
}
#[test]
#[cfg_attr(miri, ignore)]
fn failure_to_send_startup_ok_rolls_back() {
with_watchdog_phases("running the rejected startup", |_phase_reporter| {
let transport = MemoryTransport::new();
transport.fail_next_send();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(true), Condvar::new()));
let processes = mock_processes(exit);
let _startup = transport.listen("startup").unwrap();
let error = run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
.unwrap_err();
assert!(error.find_source::<StartupFailedError>().is_some());
assert!(store.list().unwrap().is_empty());
});
}
#[test]
fn init_failure_sends_startup_err_and_closes_job() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let mut processes = MockProcesses::new();
processes
.expect_durability()
.returning(|| Durability::Durable);
processes
.expect_create_lifetime_job()
.returning(|| Ok(JobId(1)));
processes
.expect_spawn_app()
.returning(|_| Err(PalError::new(PalErrorKind::Other)));
processes.expect_close_job().times(1).returning(|_| ());
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let store = store.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
phase_reporter.report("waiting for the supervisor startup connection");
let startup_conn = transport.accept(startup).unwrap();
phase_reporter.report("waiting for the startup error");
assert_eq!(transport.recv(startup_conn).unwrap(), Message::StartupErr);
phase_reporter.report("waiting for startup rollback");
supervisor.join().unwrap().unwrap_err();
assert!(store.list().unwrap().is_empty());
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_supervisor_that_cannot_outlive_its_launcher_says_so_on_the_startup_pipe() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let exit = Arc::new((Mutex::new(false), Condvar::new()));
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let processes = mock_processes_with(
Arc::clone(&exit),
Durability::TiedToLauncher,
AppWait::Reports,
);
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, durability) =
commit_startup(&transport, startup, &phase_reporter);
assert_eq!(durability, Durability::TiedToLauncher);
transport.disconnect(startup_conn);
{
let (lock, cvar) = &*exit;
*lock.lock().expect("exit lock") = true;
cvar.notify_all();
}
phase_reporter.report("waiting for supervisor shutdown");
supervisor.join().unwrap().unwrap();
});
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_wait_that_fails_still_takes_the_session_off_the_host() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let exit = Arc::new((Mutex::new(false), Condvar::new()));
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let processes =
mock_processes_with(Arc::clone(&exit), Durability::Durable, AppWait::Fails);
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let store = store.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
assert_eq!(store.list().unwrap().len(), 1, "the session was published");
transport.disconnect(startup_conn);
{
let (lock, cvar) = &*exit;
*lock.lock().expect("exit lock") = true;
cvar.notify_all();
}
phase_reporter.report("waiting for failed-wait cleanup");
supervisor.join().unwrap().unwrap_err();
assert!(
store.list().unwrap().is_empty(),
"the record outlived the session"
);
});
}
#[test]
fn breakaway_denial_keeps_its_identity() {
let denied = map_startup(&PalError::new(PalErrorKind::BreakawayDenied));
assert!(denied.find_source::<BreakawayDeniedError>().is_some());
let other = map_startup(&PalError::new(PalErrorKind::Other));
assert!(other.find_source::<StartupFailedError>().is_some());
}
#[test]
#[cfg_attr(miri, ignore)]
fn output_produced_before_the_first_attach_reaches_that_client() {
with_watchdog_phases("setting up the supervisor", |phase_reporter| {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let store = MemorySessionStore::new();
let exit = Arc::new((Mutex::new(false), Condvar::new()));
let processes = mock_processes(Arc::clone(&exit));
let startup = transport.listen("startup").unwrap();
let supervisor = thread::spawn({
let transport = transport.clone();
let pty = pty.clone();
move || {
run_supervisor(
&processes,
&store,
&transport,
&pty,
"startup",
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
}
});
let (startup_conn, _session_id, _durability) =
commit_startup(&transport, startup, &phase_reporter);
pty.push_output(PtyId(1), b"hello");
transport.disconnect(startup_conn);
let client = transport
.connect(&transport.pipe_name("nonce"), CONNECT_TIMEOUT)
.unwrap();
transport.send(client, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the attach acknowledgement");
assert!(matches!(
transport.recv(client).unwrap(),
Message::Attached { .. }
));
{
let (lock, cvar) = &*exit;
*lock.lock().expect("exit lock") = true;
cvar.notify_all();
}
let mut received = Vec::new();
loop {
phase_reporter.report("waiting for held output or the app exit status");
match transport.recv(client).unwrap() {
Message::Output(bytes) => received.extend(bytes),
Message::AppExited { status } => {
assert_eq!(status, SAMPLE_APP_EXIT);
break;
}
other => panic!("unexpected message {other:?}"),
}
}
assert_eq!(received, b"hello");
phase_reporter.report("waiting for supervisor shutdown");
assert_eq!(supervisor.join().unwrap().unwrap(), SAMPLE_APP_EXIT);
});
}
#[test]
fn only_the_first_client_receives_what_was_held_for_it() {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty);
shared.hold_for_first_client(b"hello");
assert_eq!(shared.take_preamble(), Some(b"hello".to_vec()));
shared.hold_for_first_client(b"unheard");
assert_eq!(shared.take_preamble(), None);
}
#[test]
fn nothing_held_is_nothing_to_deliver() {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty);
assert_eq!(shared.take_preamble(), None);
}
#[test]
#[cfg_attr(miri, ignore)]
fn what_is_held_for_the_first_client_is_capped() {
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty);
let chunk = vec![b'x'; 64 * 1024];
let rounds = MAX_CLIENT_BACKLOG_BYTES.div_euclid(chunk.len()) + 2;
for _ in 0..rounds {
shared.hold_for_first_client(&chunk);
}
assert_eq!(
shared.take_preamble().map(|held| held.len()),
Some(MAX_CLIENT_BACKLOG_BYTES)
);
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_preamble_too_large_for_one_frame_is_split_into_frames_a_receiver_accepts() {
let held = vec![b'x'; MAX_OUTPUT_CHUNK_BYTES.get().saturating_add(1)];
const {
assert!(
MAX_CLIENT_BACKLOG_BYTES > MAX_OUTPUT_CHUNK_BYTES.get(),
"a hold that cannot outgrow one frame would make this test vacuous"
);
}
let messages: Vec<Message> = preamble_messages(&held).collect();
assert!(messages.len() > 1, "the hold was not split");
let mut rejoined = Vec::new();
for message in &messages {
let frame = encode(message);
let prefix: [u8; 4] = frame
.get(..4)
.expect("a frame carries a length prefix")
.try_into()
.unwrap();
assert!(
payload_len_ok(u32::from_le_bytes(prefix)),
"a receiver would reject this frame"
);
match message {
Message::Output(bytes) => rejoined.extend_from_slice(bytes),
other => panic!("the hold must be relayed as output, got {other:?}"),
}
}
assert_eq!(rejoined, held);
}
fn shared_session(
transport: &MemoryTransport,
pty_host: &MemoryPseudoconsole,
) -> Shared<MemoryTransport, MemoryPseudoconsole> {
Shared {
transport: transport.clone(),
pty_host: pty_host.clone(),
pty: pty_host.create(DEFAULT_PTY_SIZE).unwrap(),
session_id: SessionId::from_u32(1).unwrap(),
client: Mutex::new(None),
attach: Mutex::new(()),
attached_generation: Arc::new(AtomicU64::default()),
preamble: Mutex::new(Some(Vec::new())),
first_attach: Mutex::new(FirstAttach::default()),
first_attach_changed: Condvar::new(),
stopping: AtomicBool::new(false),
}
}
fn client_conn(shared: &Shared<MemoryTransport, MemoryPseudoconsole>) -> Option<ConnId> {
shared.client().as_ref().map(|client| client.conn)
}
fn connected_pair(transport: &MemoryTransport, name: &str) -> (ConnId, ConnId) {
let listener = transport.listen(name).unwrap();
let client = transport.connect(name, CONNECT_TIMEOUT).unwrap();
let supervisor = transport.accept(listener).unwrap();
(supervisor, client)
}
fn attach_recorder() -> (Arc<Mutex<Vec<bool>>>, impl Fn(u64, bool)) {
let flags = Arc::new(Mutex::new(Vec::new()));
let recorder = {
let flags = Arc::clone(&flags);
move |_generation: u64, attached: bool| {
flags.lock().expect("flag lock").push(attached);
}
};
(flags, recorder)
}
#[test]
fn acknowledgement_failure_leaves_the_slot_empty() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty_host);
let (supervisor, client) = connected_pair(&transport, "pipe");
transport.send(client, &ORDINARY_ATTACH).unwrap();
transport.disconnect(client);
let (flags, recorder) = attach_recorder();
client_loop(&shared, supervisor, &recorder);
assert!(client_conn(&shared).is_none());
assert!(flags.lock().unwrap().is_empty());
}
#[test]
fn a_stalled_attach_acknowledgement_does_not_signal_attachment() {
with_watchdog_phases("setting up the client relay", |phase_reporter| {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = Arc::new(shared_session(&transport, &pty_host));
let (supervisor, client) = connected_pair(&transport, "pipe");
let (attached_tx, attached_rx) = mpsc::channel();
transport.send(client, &ORDINARY_ATTACH).unwrap();
transport.stall(supervisor);
let relay = thread::spawn({
let shared = Arc::clone(&shared);
move || {
client_loop(&shared, supervisor, &|_generation, attached| {
attached_tx.send(attached).unwrap();
});
}
});
phase_reporter.report("waiting for the attach acknowledgement to stall");
transport.wait_for_stalled_send(supervisor);
assert!(!shared.first_attach().attached);
transport.resume(supervisor);
phase_reporter.report("waiting for the attach acknowledgement");
assert!(matches!(
transport.recv(client).unwrap(),
Message::Attached { .. }
));
phase_reporter.report("waiting for the attached-flag update");
assert!(attached_rx.recv().unwrap());
assert!(shared.first_attach().attached);
transport.send(client, &Message::StartupErr).unwrap();
phase_reporter.report("waiting for the client relay to stop");
relay.join().unwrap();
assert!(!attached_rx.recv().unwrap());
});
}
#[test]
fn a_stalled_detach_update_does_not_block_or_overwrite_a_steal() {
with_watchdog_phases("setting up the client relays", |phase_reporter| {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = Arc::new(shared_session(&transport, &pty_host));
let store = MemorySessionStore::new();
let owner = ProcessIdentity::for_test(1);
let id = store.allocate_id(&owner).unwrap();
assert_eq!(id, shared.session_id);
store
.publish(&SessionRecord {
id: id.get(),
supervisor_pid: owner.pid,
supervisor_creation_time: owner.creation_time,
pipe_name: "pipe".to_string(),
launch_directory: PathBuf::from("/work"),
command: vec!["app.exe".to_string()],
started_at_unix_ms: 1,
attached: false,
})
.unwrap();
let record_live = Arc::new(Mutex::new(true));
let set_attached = store_attached_flag(
&store,
id,
record_live,
Arc::clone(&shared.attached_generation),
);
let (updated_tx, updated_rx) = mpsc::channel();
let observe_update = move |generation, attached| {
set_attached(generation, attached);
updated_tx.send(attached).unwrap();
};
let (first_supervisor, first_client) = connected_pair(&transport, "first");
let first_relay = thread::spawn({
let shared = Arc::clone(&shared);
let observe_update = observe_update.clone();
move || client_loop(&shared, first_supervisor, &observe_update)
});
transport.send(first_client, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the first attach acknowledgement");
assert!(matches!(
transport.recv(first_client).unwrap(),
Message::Attached { .. }
));
phase_reporter.report("waiting for the first attached-flag update");
assert!(updated_rx.recv().unwrap());
store.stall_publishes();
transport.disconnect(first_client);
phase_reporter.report("waiting for the detached-flag update to stall");
store.wait_for_stalled_publish();
let (second_supervisor, second_client) = connected_pair(&transport, "second");
let second_relay = thread::spawn({
let shared = Arc::clone(&shared);
move || client_loop(&shared, second_supervisor, &observe_update)
});
transport.send(second_client, &ORDINARY_ATTACH).unwrap();
phase_reporter.report("waiting for the stealing attach acknowledgement");
assert!(matches!(
transport.recv(second_client).unwrap(),
Message::Attached { .. }
));
store.resume_publishes();
phase_reporter.report("waiting for the stalled detach and newer attach updates");
let completed = [updated_rx.recv().unwrap(), updated_rx.recv().unwrap()];
assert!(completed.contains(&false));
assert!(completed.contains(&true));
first_relay.join().unwrap();
assert!(store.read(id).unwrap().unwrap().attached);
transport.disconnect(second_client);
phase_reporter.report("waiting for the final detach update");
assert!(!updated_rx.recv().unwrap());
second_relay.join().unwrap();
assert!(!store.read(id).unwrap().unwrap().attached);
});
}
#[test]
fn an_attach_that_arrives_after_teardown_claims_the_slot_is_refused() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty_host);
let (supervisor, client) = connected_pair(&transport, "pipe");
shared.stopping.store(true, Ordering::SeqCst);
transport.send(client, &ORDINARY_ATTACH).unwrap();
let (flags, recorder) = attach_recorder();
client_loop(&shared, supervisor, &recorder);
assert!(client_conn(&shared).is_none());
assert!(flags.lock().unwrap().is_empty());
transport.recv(client).unwrap_err();
}
#[test]
fn a_client_that_does_not_attach_first_is_dropped() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty_host);
let (supervisor, client) = connected_pair(&transport, "pipe");
transport
.send(client, &Message::Input(b"x".to_vec()))
.unwrap();
let (flags, recorder) = attach_recorder();
client_loop(&shared, supervisor, &recorder);
assert!(client_conn(&shared).is_none());
assert!(flags.lock().unwrap().is_empty());
assert!(pty_host.take_input(shared.pty).is_empty());
transport.recv(client).unwrap_err();
}
#[test]
fn the_relay_forwards_input_and_resize_until_the_client_stops() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty_host);
let (supervisor, client) = connected_pair(&transport, "pipe");
transport.send(client, &ORDINARY_ATTACH).unwrap();
transport
.send(
client,
&Message::Resize {
cols: 120,
rows: 40,
},
)
.unwrap();
transport
.send(client, &Message::Input(b"hi".to_vec()))
.unwrap();
transport.send(client, &Message::StartupErr).unwrap();
let (flags, recorder) = attach_recorder();
client_loop(&shared, supervisor, &recorder);
assert_eq!(
pty_host.size(shared.pty),
Some(WindowSize {
cols: 120,
rows: 40
})
);
assert_eq!(pty_host.take_input(shared.pty), b"hi");
assert!(client_conn(&shared).is_none());
assert_eq!(*flags.lock().unwrap(), vec![true, false]);
}
#[test]
fn a_displaced_relay_leaves_the_new_client_installed() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = Arc::new(shared_session(&transport, &pty_host));
let (supervisor, client) = connected_pair(&transport, "pipe");
let (successor, _successor_client) = connected_pair(&transport, "successor");
let successor_outbox = Outbox::start(transport.clone(), successor);
let displaced = Arc::new(Mutex::new(None));
let steal = {
let shared = Arc::clone(&shared);
let successor_outbox = Arc::clone(&successor_outbox);
let displaced = Arc::clone(&displaced);
move |_generation: u64, attached: bool| {
if attached {
let previous = shared.client().replace(Client {
conn: successor,
outbox: Arc::clone(&successor_outbox),
});
if let Some(previous) = previous {
previous.outbox.finish();
*displaced.lock().unwrap() = Some(previous.outbox);
}
}
}
};
transport.send(client, &ORDINARY_ATTACH).unwrap();
transport
.send(client, &Message::Input(b"x".to_vec()))
.unwrap();
client_loop(&shared, supervisor, &steal);
assert!(pty_host.take_input(shared.pty).is_empty());
assert_eq!(client_conn(&shared), Some(successor));
displaced
.lock()
.unwrap()
.take()
.expect("the steal displaced the first client")
.flush();
successor_outbox.finish();
successor_outbox.flush();
}
#[test]
fn output_is_relayed_to_the_installed_client() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty_host);
let (supervisor, client) = connected_pair(&transport, "pipe");
let outbox = Outbox::start(transport.clone(), supervisor);
*shared.client() = Some(Client {
conn: supervisor,
outbox: Arc::clone(&outbox),
});
pty_host.push_output(shared.pty, b"out");
pty_host.finish(shared.pty);
pty_output_loop(&shared);
outbox.finish();
outbox.flush();
assert!(matches!(
transport.recv(client).unwrap(),
Message::Output(bytes) if bytes == b"out",
));
}
#[test]
fn output_for_a_client_that_is_gone_is_discarded() {
let transport = MemoryTransport::new();
let pty_host = MemoryPseudoconsole::new();
let shared = shared_session(&transport, &pty_host);
let (supervisor, client) = connected_pair(&transport, "pipe");
let outbox = Outbox::start(transport.clone(), supervisor);
*shared.client() = Some(Client {
conn: supervisor,
outbox: Arc::clone(&outbox),
});
transport.disconnect(client);
pty_host.push_output(shared.pty, b"out");
pty_host.finish(shared.pty);
pty_output_loop(&shared);
outbox.flush();
}
#[test]
#[cfg_attr(miri, ignore)]
fn a_failure_after_id_allocation_releases_the_id() {
let store = MemorySessionStore::new();
store.fail_next_publish();
let transport = MemoryTransport::new();
let pty = MemoryPseudoconsole::new();
let exit = Arc::new((Mutex::new(true), Condvar::new()));
let processes = mock_processes(exit);
{
let mut guard = InitGuard {
processes: &processes,
store: &store,
transport: &transport,
pty_host: &pty,
job: None,
pty: None,
listener: None,
session: None,
committed: false,
};
let error = initialize(
&mut guard,
&processes,
&store,
&transport,
&pty,
PathBuf::from("/work"),
vec!["app.exe".to_string()],
)
.err()
.expect("record publication fails");
assert!(error.find_source::<StoreError>().is_some());
}
assert!(store.list_reservations().unwrap().is_empty());
}
#[test]
fn attached_flag_publishes_only_the_current_generation_while_the_record_lives() {
let store = MemorySessionStore::new();
let id = store.allocate_id(&ProcessIdentity::for_test(1)).unwrap();
store
.publish(&SessionRecord {
id: id.get(),
supervisor_pid: 10,
supervisor_creation_time: 100,
pipe_name: "pipe".to_string(),
launch_directory: PathBuf::from("/work"),
command: vec!["app.exe".to_string()],
started_at_unix_ms: 1,
attached: false,
})
.unwrap();
let record_live = Arc::new(Mutex::new(true));
let current_generation = Arc::new(AtomicU64::new(2));
let set_attached = store_attached_flag(
&store,
id,
Arc::clone(&record_live),
Arc::clone(¤t_generation),
);
set_attached(1, true);
assert!(!store.read(id).unwrap().unwrap().attached);
set_attached(2, true);
assert!(store.read(id).unwrap().unwrap().attached);
*record_live.lock().expect("record_live lock") = false;
current_generation.store(3, Ordering::SeqCst);
set_attached(3, false);
assert!(store.read(id).unwrap().unwrap().attached);
}
mod memory_session_store;
}