use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread;
use ohno::AppError;
use crate::constants::CONNECT_TIMEOUT;
use crate::outbox::Outbox;
use crate::pal::error::PalError;
use crate::pal::ids::{ConnId, JobId, ListenerId};
use crate::pal::processes::Processes;
use crate::pal::pseudoconsole::Pseudoconsole;
use crate::pal::session_store::SessionStore;
use crate::pal::transport::Transport;
use crate::protocol::Message;
use crate::supervisor::record_writer::RecordWriter;
use crate::supervisor::shared::{Client, FirstAttach, Shared, preamble_messages};
use crate::supervisor::startup::Initialized;
use crate::{PalFailedError, StoreError};
#[cfg_attr(test, mutants::skip)]
pub(super) 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 attached_generation = Arc::new(AtomicU64::default());
let record_writer = RecordWriter::start(store, session_id, Arc::clone(&attached_generation));
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 = record_writer.set_attached();
thread::spawn({
let shared = Arc::clone(&shared);
let transport = transport.clone();
move || accept_loop(&shared, &transport, listener, store_flag)
});
let output_failed = AtomicBool::new(false);
let job_closed = AtomicBool::new(false);
thread::scope(|scope| {
let pty_pump = scope.spawn({
let shared = Arc::clone(&shared);
let output_failed = &output_failed;
let job_closed = &job_closed;
move || {
let result = pty_output_loop(&shared);
if result.is_err() {
output_failed.store(true, Ordering::SeqCst);
close_job_once(processes, job, job_closed);
}
result
}
});
let waited = processes.wait_app(app);
if waited.is_ok() && !output_failed.load(Ordering::SeqCst) {
shared.await_first_attach();
}
transport.close_listener(listener);
close_job_once(processes, job, &job_closed);
pty_host.finish(pty);
let pumped = pty_pump
.join()
.expect("the output pump contains no panic-capable callbacks");
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 pumped.is_ok()
&& let Ok(status) = &waited
{
client.outbox.send(Message::AppExited { status: *status });
}
client.outbox.finish();
}
record_writer.finish();
let deleted = store.delete_owned_by(session_id, &identity);
if let Some(client) = client {
client.outbox.wait_for_writer();
}
let status = waited.map_err(PalFailedError::caused_by)?;
pumped.map_err(PalFailedError::caused_by)?;
deleted.map_err(StoreError::caused_by)?;
Ok(status)
})
}
#[cfg_attr(test, mutants::skip)]
fn close_job_once<P: Processes>(processes: &P, job: JobId, closed: &AtomicBool) {
if closed
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.is_ok()
{
processes.close_job(job);
}
}
#[cfg_attr(test, mutants::skip)]
pub(super) 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)]
pub(super) 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_timeout(conn, CONNECT_TIMEOUT) {
Ok(Message::Attach { size }) => {
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();
outbox.send(Message::Attached {
session_id: shared.session_id,
});
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_claimed();
if let Some(old) = previous {
old.outbox.send(Message::Displaced);
old.outbox.finish();
}
_ = shared.pty_host.resize(shared.pty, size);
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 { size } => {
_ = shared.pty_host.resize(shared.pty, size);
}
_ => 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)]
pub(super) fn pty_output_loop<T, C>(shared: &Shared<T, C>) -> Result<(), PalError>
where
T: Transport + Clone,
C: Pseudoconsole,
{
while !shared.stopping.load(Ordering::SeqCst) {
let Some(bytes) = shared.pty_host.read_output(shared.pty)? else {
return Ok(());
};
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));
}
}
Ok(())
}