use std::collections::HashMap;
use std::marker::PhantomData;
use std::sync::atomic::{AtomicBool, AtomicU64};
use std::sync::{Arc, Condvar, Mutex};
use anyhow::Result;
use oxdock_fs::{GuardedPath, WorkspaceFs};
use oxdock_parser::{Step, Value};
use oxdock_process::{BackgroundHandle, CommandContext, ProcessManager};
use super::capture::SpillBuffer;
use super::io::{ExecIo, SlidingWindow};
use super::pipe::KeeperGuard;
pub(super) struct ExecState<P: ProcessManager> {
pub(super) fs: Box<dyn WorkspaceFs>,
pub(super) cargo_target_dir: GuardedPath,
pub(super) cwd: GuardedPath,
pub(super) envs: Arc<HashMap<String, String>>,
pub(super) bg_children: Vec<Box<dyn BackgroundHandle>>,
pub(super) scope_stack: Vec<ScopeSnapshot>,
pub(super) io: ExecIo,
pub(super) assert_windows: Arc<Mutex<HashMap<(usize, usize), SlidingWindow>>>,
pub(super) var_scopes: Vec<HashMap<String, Value>>,
#[allow(dead_code)]
pub(super) cancel_token: Arc<AtomicBool>,
#[allow(dead_code)]
pub(super) active_process: Arc<Mutex<Option<Box<dyn BackgroundHandle>>>>,
#[allow(dead_code)]
pub(super) named_tasks: Arc<Mutex<HashMap<u64, Arc<TaskEntry>>>>,
#[allow(dead_code)]
pub(super) next_task_id: Arc<AtomicU64>,
pub(super) inside_async: bool,
pub(super) keeper_expiry: Option<KeeperExpiry>,
pub(super) cancellable: bool,
pub(super) _marker: PhantomData<P>,
}
pub(super) struct ScopeSnapshot {
pub(super) cwd: GuardedPath,
pub(super) root: GuardedPath,
pub(super) envs: Arc<HashMap<String, String>>,
}
pub(super) enum TaskPhase {
Running,
Awaiting,
Cancelled,
Completed,
}
pub(super) struct TaskEntryState {
pub(super) phase: TaskPhase,
pub(super) handle: Option<Box<dyn BackgroundHandle>>,
pub(super) reaped: bool,
pub(super) sink: Option<Arc<SpillBuffer>>,
}
pub(super) struct TaskEntry {
pub(super) state: Mutex<TaskEntryState>,
pub(super) done: Condvar,
}
impl TaskEntry {
pub(super) fn new_with_sink(handle: Box<dyn BackgroundHandle>, sink: Arc<SpillBuffer>) -> Self {
Self {
state: Mutex::new(TaskEntryState {
phase: TaskPhase::Running,
handle: Some(handle),
reaped: false,
sink: Some(sink),
}),
done: Condvar::new(),
}
}
pub(super) fn take_sink(&self) -> Option<Arc<SpillBuffer>> {
self.state
.lock()
.unwrap_or_else(|e| e.into_inner())
.sink
.take()
}
pub(super) fn wait_reaped(&self) {
let mut guard = self.state.lock().unwrap_or_else(|e| e.into_inner());
while !guard.reaped {
guard = self.done.wait(guard).unwrap_or_else(|e| e.into_inner());
}
}
pub(super) fn finish_teardown(&self) {
{
let mut guard = self.state.lock().unwrap_or_else(|e| e.into_inner());
guard.handle = None;
guard.reaped = true;
}
self.done.notify_all();
}
}
pub(super) struct KeeperExpiry {
body_addr: usize,
body_len: usize,
map: HashMap<usize, Vec<KeeperGuard>>,
}
impl KeeperExpiry {
pub(super) fn new(steps: &[Step], map: HashMap<usize, Vec<KeeperGuard>>) -> Self {
Self {
body_addr: steps.as_ptr() as usize,
body_len: steps.len(),
map,
}
}
pub(super) fn matches(&self, steps: &[Step]) -> bool {
self.body_addr == steps.as_ptr() as usize && self.body_len == steps.len()
}
pub(super) fn expire_step(&mut self, steps: &[Step], idx: usize) -> bool {
if self.matches(steps) {
drop(self.map.remove(&idx));
}
self.map.is_empty()
}
}
impl<P: ProcessManager> ExecState<P> {
pub(super) fn command_ctx(&self) -> Result<CommandContext> {
Ok(CommandContext::new(
&self.cwd.clone().into(),
Arc::clone(&self.envs),
&self.cargo_target_dir,
self.fs.root(),
self.fs.build_context(),
))
}
#[allow(dead_code)]
pub(super) fn fork(&self) -> Self {
Self {
fs: self.fs.clone_box(),
cargo_target_dir: self.cargo_target_dir.clone(),
cwd: self.cwd.clone(),
envs: Arc::clone(&self.envs),
bg_children: Vec::new(),
scope_stack: Vec::new(),
io: self.io.clone(),
assert_windows: Arc::clone(&self.assert_windows),
var_scopes: self.var_scopes.clone(),
cancel_token: Arc::new(AtomicBool::new(false)),
active_process: Arc::new(Mutex::new(None)),
named_tasks: Arc::clone(&self.named_tasks),
next_task_id: Arc::clone(&self.next_task_id),
inside_async: true,
keeper_expiry: None,
cancellable: self.cancellable,
_marker: PhantomData,
}
}
pub(super) fn push_var_scope(&mut self) {
self.var_scopes.push(HashMap::new());
}
pub(super) fn pop_var_scope(&mut self) {
self.var_scopes.pop();
}
pub(super) fn push_scope(&mut self) {
self.scope_stack.push(ScopeSnapshot {
cwd: self.cwd.clone(),
root: self.fs.root().clone(),
envs: Arc::clone(&self.envs),
});
self.push_var_scope();
}
pub(super) fn pop_scope(&mut self) -> Result<()> {
let snapshot = self
.scope_stack
.pop()
.ok_or_else(|| anyhow::anyhow!("scope stack underflow during pop"))?;
self.fs.set_root(&snapshot.root);
self.cwd = snapshot.cwd;
self.envs = snapshot.envs;
self.pop_var_scope();
Ok(())
}
pub(super) fn set_var(&mut self, key: String, value: Value) {
if let Some(scope) = self.var_scopes.last_mut() {
scope.insert(key, value);
}
}
pub(super) fn get_var(&self, key: &str) -> Option<Value> {
for scope in self.var_scopes.iter().rev() {
if let Some(value) = scope.get(key) {
return Some(value.clone());
}
}
None
}
pub(super) fn all_vars(&self) -> HashMap<String, Value> {
let mut result = HashMap::new();
for scope in self.var_scopes.iter().rev() {
for (k, v) in scope {
result.entry(k.clone()).or_insert_with(|| v.clone());
}
}
result
}
}