use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use tokio::io::AsyncReadExt;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use tracing::{debug, warn};
use crate::client::completion::Verdict;
use crate::orchestration::execution_facts::{
EpisodeObserver, EpisodeTerminal, EpisodeTransition, EpisodeTransitionKind, ExecutionFactsStore,
};
use crate::web::remote_control_api::dto::{
ChangeExecutionState, ExecutionEventFile, ExecutionEventType, ExecutionSinkCapability,
ExecutionSinkSpec, ProposalSubscriptionCapability, EXECUTION_EVENT_SCHEMA_VERSION,
};
use crate::web::remote_control_api::ExecutionContractHandle;
pub const MAX_COMMAND_ARGS: usize = 16;
pub const MAX_COMMAND_ARG_LEN: usize = 4096;
pub const CALLBACK_TIMEOUT: Duration = Duration::from_secs(20);
pub const MAX_CALLBACK_OUTPUT_BYTES: usize = 8 * 1024;
const SHUTDOWN_DEADLINE: Duration = Duration::from_secs(40);
const DRAIN_GRACE: Duration = Duration::from_secs(2);
const MAX_EVIDENCE_BYTES: usize = 512;
const VERIFY_ATTEMPTS: usize = 5;
const VERIFY_RETRY_INTERVAL: Duration = Duration::from_millis(200);
const VERIFY_ROUND_BUDGET: Duration = Duration::from_secs(20);
pub const MAX_PROPOSAL_TARGETS: usize = 64;
pub fn capability() -> ExecutionSinkCapability {
ExecutionSinkCapability {
available: true,
max_command_args: MAX_COMMAND_ARGS,
max_command_arg_len: MAX_COMMAND_ARG_LEN,
callback_timeout_ms: CALLBACK_TIMEOUT.as_millis() as u64,
max_callback_output_bytes: MAX_CALLBACK_OUTPUT_BYTES,
}
}
pub fn proposal_capability() -> ProposalSubscriptionCapability {
ProposalSubscriptionCapability {
available: true,
max_targets: MAX_PROPOSAL_TARGETS,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SinkRefusal {
UnknownExecution,
BindingMismatch {
actual_change_id: String,
},
InstanceMismatch,
InvalidCommand(String),
}
#[derive(Debug, Clone, Default)]
struct ProposalEntry {
sink: Option<ExecutionSinkSpec>,
latest: Option<String>,
}
impl ProposalEntry {
fn is_empty(&self) -> bool {
self.sink.is_none() && self.latest.is_none()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProposalSubscriptionView {
pub change_id: String,
pub sink: Option<ExecutionSinkSpec>,
pub execution_id: Option<String>,
pub terminal_dispatched: bool,
pub delivered_events: Vec<ExecutionEventType>,
}
#[derive(Debug, Clone)]
struct Entry {
change_id: String,
sink: Option<ExecutionSinkSpec>,
proposal_bound: bool,
terminal: Option<EpisodeTerminal>,
terminal_attempted: bool,
terminal_dispatched: bool,
blocked_active: bool,
delivered: Vec<ExecutionEventType>,
stopping_attempted: bool,
}
impl Entry {
fn new(change_id: String) -> Self {
Self {
change_id,
sink: None,
proposal_bound: false,
terminal: None,
terminal_attempted: false,
terminal_dispatched: false,
blocked_active: false,
delivered: Vec::new(),
stopping_attempted: false,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SinkView {
pub change_id: String,
pub sink: Option<ExecutionSinkSpec>,
pub terminal_dispatched: bool,
pub delivered_events: Vec<ExecutionEventType>,
}
#[derive(Debug)]
enum Task {
Episode(EpisodeTransition),
Registered { execution_id: String },
Stopping(tokio::sync::oneshot::Sender<()>),
}
#[derive(Debug, Clone, Copy)]
struct Limits {
callback_timeout: Duration,
shutdown_deadline: Duration,
}
impl Default for Limits {
fn default() -> Self {
Self {
callback_timeout: CALLBACK_TIMEOUT,
shutdown_deadline: SHUTDOWN_DEADLINE,
}
}
}
pub struct CompletionSinkRegistry {
instance_id: String,
entries: Mutex<HashMap<String, Entry>>,
proposals: Mutex<HashMap<String, ProposalEntry>>,
facts: Arc<ExecutionFactsStore>,
contract: Arc<ExecutionContractHandle>,
repo_root: Mutex<Option<PathBuf>>,
event_dir: Mutex<Option<PathBuf>>,
tasks: mpsc::UnboundedSender<Task>,
stopping: AtomicBool,
cancel: CancellationToken,
limits: Mutex<Limits>,
reap_gate: Mutex<Option<Arc<ReapGate>>>,
}
#[derive(Debug)]
pub struct ReapGate {
reached: tokio::sync::Semaphore,
released: tokio::sync::Notify,
}
impl Default for ReapGate {
fn default() -> Self {
Self {
reached: tokio::sync::Semaphore::new(0),
released: tokio::sync::Notify::new(),
}
}
}
impl ReapGate {
#[cfg_attr(not(test), allow(dead_code))]
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
}
#[cfg_attr(not(test), allow(dead_code))]
pub async fn reached(&self) {
self.reached
.acquire()
.await
.expect("the reap gate is never closed")
.forget();
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn release(&self) {
self.released.notify_one();
}
async fn hold(&self) {
self.reached.add_permits(1);
self.released.notified().await;
}
}
impl std::fmt::Debug for CompletionSinkRegistry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("CompletionSinkRegistry")
.field("instance_id", &self.instance_id)
.finish_non_exhaustive()
}
}
struct Wiring {
registry: Arc<CompletionSinkRegistry>,
tasks: mpsc::UnboundedReceiver<Task>,
}
impl CompletionSinkRegistry {
fn build(
instance_id: String,
facts: Arc<ExecutionFactsStore>,
contract: Arc<ExecutionContractHandle>,
) -> Wiring {
let (tx, rx) = mpsc::unbounded_channel();
Wiring {
registry: Arc::new(Self {
instance_id,
entries: Mutex::new(HashMap::new()),
proposals: Mutex::new(HashMap::new()),
facts,
contract,
repo_root: Mutex::new(None),
event_dir: Mutex::new(None),
tasks: tx,
stopping: AtomicBool::new(false),
cancel: CancellationToken::new(),
limits: Mutex::new(Limits::default()),
reap_gate: Mutex::new(None),
}),
tasks: rx,
}
}
pub fn start(
instance_id: String,
facts: Arc<ExecutionFactsStore>,
contract: Arc<ExecutionContractHandle>,
) -> Arc<Self> {
let Wiring {
registry,
mut tasks,
} = Self::build(instance_id, facts.clone(), contract);
facts.bind_episode_observer(registry.clone());
let dispatcher = registry.clone();
tokio::spawn(async move {
while let Some(task) = tasks.recv().await {
dispatcher.handle(task).await;
}
});
registry
}
pub fn bind_repo_root(&self, repo_root: PathBuf) {
*self.lock_repo_root() = Some(repo_root);
}
fn lock_repo_root(&self) -> MutexGuard<'_, Option<PathBuf>> {
self.repo_root
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn lock(&self) -> MutexGuard<'_, HashMap<String, Entry>> {
self.entries
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn lock_proposals(&self) -> MutexGuard<'_, HashMap<String, ProposalEntry>> {
self.proposals
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn instance_id(&self) -> &str {
&self.instance_id
}
pub fn execution_state(&self, change_id: &str) -> ChangeExecutionState {
ChangeExecutionState::from_shared(self.facts.change(change_id).execution_state)
}
pub fn view(
&self,
execution_id: &str,
instance_id: &str,
change_id: &str,
) -> Result<SinkView, SinkRefusal> {
let entries = self.lock();
let entry = self.resolve(&entries, execution_id, instance_id, change_id)?;
Ok(SinkView {
change_id: entry.change_id.clone(),
sink: entry.sink.clone(),
terminal_dispatched: entry.terminal_dispatched,
delivered_events: entry.delivered.clone(),
})
}
fn resolve<'a>(
&self,
entries: &'a HashMap<String, Entry>,
execution_id: &str,
instance_id: &str,
change_id: &str,
) -> Result<&'a Entry, SinkRefusal> {
let entry = entries
.get(execution_id)
.ok_or(SinkRefusal::UnknownExecution)?;
if instance_id != self.instance_id || change_id != entry.change_id {
return Err(SinkRefusal::BindingMismatch {
actual_change_id: entry.change_id.clone(),
});
}
Ok(entry)
}
pub fn set_sink(
&self,
execution_id: &str,
instance_id: &str,
change_id: &str,
spec: ExecutionSinkSpec,
) -> Result<SinkView, SinkRefusal> {
validate_command(&spec.command)?;
let view = {
let mut entries = self.lock();
self.resolve(&entries, execution_id, instance_id, change_id)?;
let entry = entries
.get_mut(execution_id)
.ok_or(SinkRefusal::UnknownExecution)?;
entry.sink = Some(spec);
entry.proposal_bound = false;
SinkView {
change_id: entry.change_id.clone(),
sink: entry.sink.clone(),
terminal_dispatched: entry.terminal_dispatched,
delivered_events: entry.delivered.clone(),
}
};
if !self.stopping.load(Ordering::SeqCst) {
let _ = self.tasks.send(Task::Registered {
execution_id: execution_id.to_string(),
});
}
Ok(view)
}
pub fn clear_sink(
&self,
execution_id: &str,
instance_id: &str,
change_id: &str,
) -> Result<SinkView, SinkRefusal> {
let mut entries = self.lock();
self.resolve(&entries, execution_id, instance_id, change_id)?;
let entry = entries
.get_mut(execution_id)
.ok_or(SinkRefusal::UnknownExecution)?;
entry.sink = None;
entry.proposal_bound = false;
Ok(SinkView {
change_id: entry.change_id.clone(),
sink: None,
terminal_dispatched: entry.terminal_dispatched,
delivered_events: entry.delivered.clone(),
})
}
pub fn set_proposal_subscription(
&self,
change_id: &str,
instance_id: &str,
spec: ExecutionSinkSpec,
) -> Result<ProposalSubscriptionView, SinkRefusal> {
validate_command(&spec.command)?;
if instance_id != self.instance_id {
return Err(SinkRefusal::InstanceMismatch);
}
let (view, bound) = {
let mut proposals = self.lock_proposals();
let proposal = proposals.entry(change_id.to_string()).or_default();
proposal.sink = Some(spec.clone());
let latest = proposal.latest.clone();
let mut entries = self.lock();
let bound = latest.as_ref().filter(|execution_id| {
match entries.get_mut(execution_id.as_str()) {
Some(entry) if entry.terminal_dispatched => false,
Some(entry) if entry.sink.is_some() && !entry.proposal_bound => false,
Some(entry) => {
entry.sink = Some(spec.clone());
entry.proposal_bound = true;
true
}
None => false,
}
});
let bound = bound.cloned();
(
Self::project_proposal(change_id, proposals.get(change_id), &entries),
bound,
)
};
if let Some(execution_id) = bound {
if !self.stopping.load(Ordering::SeqCst) {
let _ = self.tasks.send(Task::Registered { execution_id });
}
}
Ok(view)
}
pub fn view_proposal_subscription(
&self,
change_id: &str,
instance_id: &str,
) -> Result<ProposalSubscriptionView, SinkRefusal> {
if instance_id != self.instance_id {
return Err(SinkRefusal::InstanceMismatch);
}
let proposals = self.lock_proposals();
let entries = self.lock();
Ok(Self::project_proposal(
change_id,
proposals.get(change_id),
&entries,
))
}
pub fn clear_proposal_subscription(
&self,
change_id: &str,
instance_id: &str,
) -> Result<ProposalSubscriptionView, SinkRefusal> {
if instance_id != self.instance_id {
return Err(SinkRefusal::InstanceMismatch);
}
let mut proposals = self.lock_proposals();
let mut entries = self.lock();
if let Some(proposal) = proposals.get_mut(change_id) {
proposal.sink = None;
if let Some(entry) = proposal
.latest
.as_ref()
.and_then(|execution_id| entries.get_mut(execution_id.as_str()))
{
if entry.proposal_bound {
entry.sink = None;
entry.proposal_bound = false;
}
}
}
let view = Self::project_proposal(change_id, proposals.get(change_id), &entries);
if proposals
.get(change_id)
.is_some_and(ProposalEntry::is_empty)
{
proposals.remove(change_id);
}
Ok(view)
}
fn project_proposal(
change_id: &str,
proposal: Option<&ProposalEntry>,
entries: &HashMap<String, Entry>,
) -> ProposalSubscriptionView {
let sink = proposal.and_then(|proposal| proposal.sink.clone());
let execution_id = proposal.and_then(|proposal| proposal.latest.clone());
let episode = execution_id
.as_ref()
.and_then(|execution_id| entries.get(execution_id.as_str()));
ProposalSubscriptionView {
change_id: change_id.to_string(),
sink,
execution_id,
terminal_dispatched: episode.is_some_and(|entry| entry.terminal_dispatched),
delivered_events: episode
.map(|entry| entry.delivered.clone())
.unwrap_or_default(),
}
}
fn bind_episode(&self, change_id: &str, execution_id: &str) {
let mut proposals = self.lock_proposals();
let mut entries = self.lock();
let entry = entries
.entry(execution_id.to_string())
.or_insert_with(|| Entry::new(change_id.to_string()));
let proposal = proposals.entry(change_id.to_string()).or_default();
if entry.sink.is_none() || entry.proposal_bound {
if let Some(spec) = proposal.sink.clone() {
entry.sink = Some(spec);
entry.proposal_bound = true;
}
}
let superseded = proposal.latest.replace(execution_id.to_string());
if let Some(superseded) = superseded {
if superseded != execution_id {
let disposable = entries
.get(&superseded)
.is_some_and(|entry| entry.proposal_bound || entry.sink.is_none());
if disposable {
entries.remove(&superseded);
}
}
}
}
pub async fn owner_stopping(&self) {
self.stopping.store(true, Ordering::SeqCst);
let (done, mut wait) = tokio::sync::oneshot::channel();
if self.tasks.send(Task::Stopping(done)).is_err() {
self.retain_events(
"the dispatcher was gone before it could be asked to acknowledge a reap",
);
return;
}
let deadline = self.limits().shutdown_deadline;
let acknowledged = match tokio::time::timeout(deadline, &mut wait).await {
Ok(Ok(())) => true,
Ok(Err(_)) => {
self.retain_events(
"the dispatcher dropped its acknowledgement before the shutdown deadline",
);
false
}
Err(_) => {
self.cancel.cancel();
match wait.await {
Ok(()) => true,
Err(_) => {
self.retain_events(
"the dispatcher dropped its acknowledgement after shutdown \
cancellation",
);
false
}
}
}
};
if acknowledged {
self.remove_events();
}
}
fn remove_events(&self) {
let Some(dir) = self.lock_event_dir().take() else {
return;
};
if let Err(error) = std::fs::remove_dir_all(&dir) {
if error.kind() != std::io::ErrorKind::NotFound {
debug!(
event_dir = %dir.display(),
error = %error,
"the acknowledged owner-private event directory could not be removed"
);
}
}
}
fn retain_events(&self, reason: &'static str) {
let Some(dir) = self.lock_event_dir().clone() else {
return;
};
warn!(
event_dir = %dir.display(),
reason,
"callback reap was not acknowledged, so the owner-private event directory is retained"
);
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn set_reap_gate(&self, gate: Arc<ReapGate>) {
*self.lock_reap_gate() = Some(gate);
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn set_callback_timeout(&self, timeout: Duration) {
self.lock_limits().callback_timeout = timeout;
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn set_shutdown_deadline(&self, deadline: Duration) {
self.lock_limits().shutdown_deadline = deadline;
}
fn limits(&self) -> Limits {
*self.lock_limits()
}
fn lock_limits(&self) -> MutexGuard<'_, Limits> {
self.limits
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn lock_event_dir(&self) -> MutexGuard<'_, Option<PathBuf>> {
self.event_dir
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn lock_reap_gate(&self) -> MutexGuard<'_, Option<Arc<ReapGate>>> {
self.reap_gate
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn reap_gate(&self) -> Option<Arc<ReapGate>> {
self.lock_reap_gate().clone()
}
async fn handle(&self, task: Task) {
match task {
Task::Episode(transition) => self.handle_episode(transition).await,
Task::Registered { execution_id } => self.handle_terminal(&execution_id).await,
Task::Stopping(done) => {
self.handle_stopping().await;
let _ = done.send(());
}
}
}
async fn handle_episode(&self, transition: EpisodeTransition) {
match transition.kind {
EpisodeTransitionKind::Started => {
self.bind_episode(&transition.change_id, &transition.execution_id);
}
EpisodeTransitionKind::BlockedEntered => {
let should_deliver = {
let mut entries = self.lock();
match entries.get_mut(&transition.execution_id) {
Some(entry) => {
entry.blocked_active = true;
entry.sink.as_ref().is_some_and(|sink| sink.notify_blocked)
&& !entry.terminal_dispatched
}
None => false,
}
};
if should_deliver {
self.deliver(&transition.execution_id, ExecutionEventType::Blocked, None)
.await;
}
}
EpisodeTransitionKind::BlockedLeft => {
if let Some(entry) = self.lock().get_mut(&transition.execution_id) {
entry.blocked_active = false;
}
}
EpisodeTransitionKind::Terminal(terminal) => {
if let Some(entry) = self.lock().get_mut(&transition.execution_id) {
entry.terminal = Some(terminal);
entry.blocked_active = false;
}
self.handle_terminal(&transition.execution_id).await;
}
}
}
fn claim_terminal(&self, execution_id: &str) -> Option<(EpisodeTerminal, String)> {
let mut entries = self.lock();
let entry = entries.get_mut(execution_id)?;
let terminal = entry.terminal?;
if entry.terminal_attempted || entry.sink.is_none() {
return None;
}
entry.terminal_attempted = true;
Some((terminal, entry.change_id.clone()))
}
async fn handle_terminal(&self, execution_id: &str) {
let Some((terminal, change_id)) = self.claim_terminal(execution_id) else {
return;
};
let (event_type, evidence) = match terminal {
EpisodeTerminal::Failed => (ExecutionEventType::Failed, None),
EpisodeTerminal::Stopped => (ExecutionEventType::Stopped, None),
EpisodeTerminal::Completed => match self.certify(&change_id).await {
Some(evidence) => (ExecutionEventType::Completed, Some(evidence)),
None => {
warn!(
change_id = %change_id,
execution_id = %execution_id,
"execution reported terminal success but repository evidence did not \
prove the owner's terminal mode; no completion event was dispatched"
);
return;
}
},
};
self.deliver(execution_id, event_type, evidence).await;
}
async fn handle_stopping(&self) {
let live: Vec<String> = {
let mut entries = self.lock();
entries
.iter_mut()
.filter(|(_, entry)| {
entry.sink.is_some() && !entry.terminal_dispatched && !entry.stopping_attempted
})
.map(|(id, entry)| {
entry.stopping_attempted = true;
id.clone()
})
.collect()
};
for execution_id in live {
if self.cancel.is_cancelled() {
debug!(
remaining = %execution_id,
"the shutdown deadline passed, so no further owner_stopping delivery starts"
);
break;
}
self.deliver(&execution_id, ExecutionEventType::OwnerStopping, None)
.await;
}
}
async fn certify(&self, change_id: &str) -> Option<String> {
let repo_root = self.lock_repo_root().clone();
let Some(repo_root) = repo_root else {
debug!(
change_id = %change_id,
"no repository root is bound, so a claimed completion cannot be certified"
);
return None;
};
let Some(contract) = self.contract.resolve(Some(change_id)) else {
debug!(
change_id = %change_id,
"this owner published no execution contract, so nothing would prove completion"
);
return None;
};
for attempt in 0..VERIFY_ATTEMPTS {
if self.cancel.is_cancelled() {
debug!(
change_id = %change_id,
"the shutdown deadline passed, so completion verification stops"
);
return None;
}
let deadline = crate::bounded_git::GitDeadline::Operation(
tokio::time::Instant::now() + VERIFY_ROUND_BUDGET,
);
match crate::client::completion::certify(change_id, &repo_root, &contract, deadline)
.await
{
Verdict::Completed { evidence } => return Some(evidence),
Verdict::Broken { detail } => {
debug!(change_id = %change_id, detail = %detail, "completion evidence is unusable");
return None;
}
Verdict::Unsupported { detail } => {
debug!(change_id = %change_id, detail = %detail, "terminal mode has no repository proof");
return None;
}
Verdict::NotCompleted { detail } => {
debug!(
change_id = %change_id,
attempt = attempt + 1,
detail = %detail,
"completion evidence is not yet present"
);
}
Verdict::DeadlineExpired { stage } => {
debug!(change_id = %change_id, stage = ?stage, "repository verification exceeded its budget");
}
}
if attempt + 1 < VERIFY_ATTEMPTS {
tokio::select! {
_ = tokio::time::sleep(VERIFY_RETRY_INTERVAL) => {}
_ = self.cancel.cancelled() => {}
}
}
}
None
}
async fn deliver(
&self,
execution_id: &str,
event_type: ExecutionEventType,
evidence: Option<String>,
) {
if self.cancel.is_cancelled() {
return;
}
let Some((change_id, sink)) = ({
let entries = self.lock();
entries.get(execution_id).and_then(|entry| {
entry
.sink
.clone()
.map(|sink| (entry.change_id.clone(), sink))
})
}) else {
return;
};
let payload = ExecutionEventFile {
schema_version: EXECUTION_EVENT_SCHEMA_VERSION,
event_type,
instance_id: self.instance_id.clone(),
execution_id: execution_id.to_string(),
change_id: change_id.clone(),
emitted_at: chrono::Utc::now().to_rfc3339(),
terminal: event_type.is_terminal(),
terminal_mode: self
.contract
.resolve(Some(&change_id))
.map(|contract| contract.terminal_mode),
evidence: evidence.map(|evidence| truncate(&evidence, MAX_EVIDENCE_BYTES)),
};
let path = match self.write_event(execution_id, event_type, &payload) {
Ok(path) => path,
Err(error) => {
warn!(
change_id = %change_id,
execution_id = %execution_id,
error = %error,
"the completion event file could not be written; no callback was started"
);
return;
}
};
{
let mut entries = self.lock();
if let Some(entry) = entries.get_mut(execution_id) {
entry.delivered.push(event_type);
if event_type.is_terminal() {
entry.terminal_dispatched = true;
}
}
}
let report = run_callback(
&sink.command,
&path,
&payload,
self.limits().callback_timeout,
&self.cancel,
self.reap_gate(),
)
.await;
let _ = std::fs::remove_file(&path);
let truncated = report.truncated();
match &report.outcome {
Ok(()) => debug!(
change_id = %change_id,
execution_id = %execution_id,
event = event_type.as_str(),
stdout_bytes = report.stdout.total,
stderr_bytes = report.stderr.total,
output_truncated = truncated,
"completion callback finished"
),
Err(detail) => warn!(
change_id = %change_id,
execution_id = %execution_id,
event = event_type.as_str(),
detail = %detail,
stdout_bytes = report.stdout.total,
stderr_bytes = report.stderr.total,
output_truncated = truncated,
"completion callback failed"
),
}
}
fn write_event(
&self,
execution_id: &str,
event_type: ExecutionEventType,
payload: &ExecutionEventFile,
) -> std::io::Result<PathBuf> {
if self.cancel.is_cancelled() {
return Err(std::io::Error::other(
"the shutdown deadline passed, so no event directory or artifact is created",
));
}
let dir = {
let mut slot = self.lock_event_dir();
match slot.as_ref() {
Some(dir) => dir.clone(),
None => {
let created = tempfile::Builder::new().prefix("cflx-events-").tempdir()?;
restrict(created.path(), 0o700)?;
let path = created.keep();
*slot = Some(path.clone());
path
}
}
};
let path = dir.join(format!("{execution_id}-{}.json", event_type.as_str()));
let body = serde_json::to_vec_pretty(payload)
.map_err(|error| std::io::Error::other(error.to_string()))?;
write_owner_only(&path, &body)?;
Ok(path)
}
}
impl EpisodeObserver for CompletionSinkRegistry {
fn observe_episode(&self, transition: &EpisodeTransition) {
if self.stopping.load(Ordering::SeqCst) {
return;
}
let _ = self.tasks.send(Task::Episode(transition.clone()));
}
}
#[cfg(unix)]
fn restrict(path: &Path, mode: u32) -> std::io::Result<()> {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(mode))
}
#[cfg(not(unix))]
fn restrict(_path: &Path, _mode: u32) -> std::io::Result<()> {
Ok(())
}
#[cfg(unix)]
fn write_owner_only(path: &Path, body: &[u8]) -> std::io::Result<()> {
use std::io::Write;
use std::os::unix::fs::OpenOptionsExt;
match std::fs::remove_file(path) {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(error),
}
let mut file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o400)
.open(path)?;
file.write_all(body)
}
#[cfg(not(unix))]
fn write_owner_only(path: &Path, body: &[u8]) -> std::io::Result<()> {
std::fs::write(path, body)
}
pub fn validate_command(command: &[String]) -> Result<(), SinkRefusal> {
if command.is_empty() {
return Err(SinkRefusal::InvalidCommand(
"a completion sink needs at least a program to run".to_string(),
));
}
if command.len() > MAX_COMMAND_ARGS {
return Err(SinkRefusal::InvalidCommand(format!(
"a completion sink accepts at most {MAX_COMMAND_ARGS} argv elements"
)));
}
for argument in command {
if argument.len() > MAX_COMMAND_ARG_LEN {
return Err(SinkRefusal::InvalidCommand(format!(
"one argv element exceeds {MAX_COMMAND_ARG_LEN} bytes"
)));
}
if argument.bytes().any(|byte| byte < 0x20 && byte != b'\t') {
return Err(SinkRefusal::InvalidCommand(
"argv elements may not contain control characters".to_string(),
));
}
}
if command[0].trim().is_empty() {
return Err(SinkRefusal::InvalidCommand(
"the program name is empty".to_string(),
));
}
Ok(())
}
#[derive(Debug, Default)]
struct Drained {
retained: Vec<u8>,
total: usize,
}
impl Drained {
fn truncated(&self) -> bool {
self.total > self.retained.len()
}
fn text(&self) -> String {
let text = String::from_utf8_lossy(&self.retained);
match self.truncated() {
true => format!("{text}… ({} bytes total)", self.total),
false => text.to_string(),
}
}
}
#[derive(Debug)]
struct CallbackReport {
outcome: Result<(), String>,
stdout: Drained,
stderr: Drained,
}
impl CallbackReport {
fn unstarted(detail: String) -> Self {
Self {
outcome: Err(detail),
stdout: Drained::default(),
stderr: Drained::default(),
}
}
fn truncated(&self) -> bool {
self.stdout.truncated() || self.stderr.truncated()
}
}
async fn drain<R>(mut stream: R, limit: usize) -> Drained
where
R: tokio::io::AsyncRead + Unpin,
{
let mut drained = Drained {
retained: Vec::new(),
total: 0,
};
let mut chunk = [0u8; 8 * 1024];
loop {
match stream.read(&mut chunk).await {
Ok(0) | Err(_) => break,
Ok(read) => {
drained.total = drained.total.saturating_add(read);
if drained.retained.len() < limit {
let room = limit - drained.retained.len();
drained.retained.extend_from_slice(&chunk[..read.min(room)]);
}
}
}
}
drained
}
async fn collect(handle: tokio::task::JoinHandle<Drained>) -> Drained {
let mut handle = handle;
match tokio::time::timeout(DRAIN_GRACE, &mut handle).await {
Ok(Ok(drained)) => drained,
Ok(Err(_)) => Drained::default(),
Err(_) => {
handle.abort();
Drained::default()
}
}
}
async fn run_callback(
command: &[String],
event_path: &Path,
payload: &ExecutionEventFile,
timeout: Duration,
cancel: &CancellationToken,
reap_gate: Option<Arc<ReapGate>>,
) -> CallbackReport {
enum Ended {
Exited(std::process::ExitStatus),
Broken(String),
TimedOut,
Cancelled,
}
let mut child = tokio::process::Command::new(&command[0]);
child
.args(&command[1..])
.env_clear()
.env("CFLX_EVENT_PATH", event_path)
.env("CFLX_EVENT_TYPE", payload.event_type.as_str())
.env("CFLX_EXECUTION_ID", &payload.execution_id)
.env("CFLX_CHANGE_ID", &payload.change_id)
.env("CFLX_INSTANCE_ID", &payload.instance_id)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
let mut spawned = match child.spawn() {
Ok(spawned) => spawned,
Err(error) => return CallbackReport::unstarted(error.to_string()),
};
let stdout = spawned.stdout.take();
let stderr = spawned.stderr.take();
let stdout = tokio::spawn(async move {
match stdout {
Some(stream) => drain(stream, MAX_CALLBACK_OUTPUT_BYTES).await,
None => Drained::default(),
}
});
let stderr = tokio::spawn(async move {
match stderr {
Some(stream) => drain(stream, MAX_CALLBACK_OUTPUT_BYTES).await,
None => Drained::default(),
}
});
let mut ended = tokio::select! {
result = spawned.wait() => match result {
Ok(status) => Ended::Exited(status),
Err(error) => Ended::Broken(error.to_string()),
},
_ = tokio::time::sleep(timeout) => Ended::TimedOut,
_ = cancel.cancelled() => Ended::Cancelled,
};
if matches!(ended, Ended::TimedOut | Ended::Cancelled) {
if matches!(ended, Ended::Cancelled) {
if let Some(gate) = &reap_gate {
gate.hold().await;
}
}
if let Err(error) = spawned.kill().await {
ended = Ended::Broken(format!("the callback could not be terminated: {error}"));
}
}
let stdout = collect(stdout).await;
let stderr = collect(stderr).await;
let outcome = match ended {
Ended::Exited(status) if status.success() => Ok(()),
Ended::Exited(status) => Err(format!(
"exit {:?}: {}",
status.code(),
truncate(&stderr.text(), MAX_CALLBACK_OUTPUT_BYTES)
)),
Ended::Broken(error) => Err(error),
Ended::TimedOut => Err(format!(
"the callback did not finish within {}ms and was terminated",
timeout.as_millis()
)),
Ended::Cancelled => Err(
"the owner reached its shutdown deadline, so the callback was terminated".to_string(),
),
};
CallbackReport {
outcome,
stdout,
stderr,
}
}
fn truncate(value: &str, limit: usize) -> String {
if value.len() <= limit {
return value.to_string();
}
let mut end = limit;
while end > 0 && !value.is_char_boundary(end) {
end -= 1;
}
format!("{}…", &value[..end])
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn an_argv_must_be_bounded_and_free_of_control_bytes() {
assert!(validate_command(&["/bin/true".to_string()]).is_ok());
assert!(matches!(
validate_command(&[]),
Err(SinkRefusal::InvalidCommand(_))
));
assert!(matches!(
validate_command(&["".to_string()]),
Err(SinkRefusal::InvalidCommand(_))
));
assert!(matches!(
validate_command(&["/bin/echo".to_string(), "a\nb".to_string()]),
Err(SinkRefusal::InvalidCommand(_))
));
let too_many: Vec<String> = (0..MAX_COMMAND_ARGS + 1).map(|i| i.to_string()).collect();
assert!(matches!(
validate_command(&too_many),
Err(SinkRefusal::InvalidCommand(_))
));
assert!(matches!(
validate_command(&["x".repeat(MAX_COMMAND_ARG_LEN + 1)]),
Err(SinkRefusal::InvalidCommand(_))
));
}
#[test]
fn every_terminal_event_type_is_terminal_and_the_others_are_not() {
assert!(ExecutionEventType::Completed.is_terminal());
assert!(ExecutionEventType::Failed.is_terminal());
assert!(ExecutionEventType::Stopped.is_terminal());
assert!(!ExecutionEventType::Blocked.is_terminal());
assert!(!ExecutionEventType::OwnerStopping.is_terminal());
}
#[test]
fn truncation_keeps_a_character_boundary() {
assert_eq!(truncate("abc", 8), "abc");
assert_eq!(truncate("abcdef", 3), "abc…");
assert_eq!(truncate("aé", 2), "a…");
}
#[tokio::test]
async fn a_drain_retains_at_most_the_limit_and_still_consumes_the_rest() {
let produced = vec![b'a'; MAX_CALLBACK_OUTPUT_BYTES * 4 + 7];
let drained = drain(produced.as_slice(), MAX_CALLBACK_OUTPUT_BYTES).await;
assert_eq!(
drained.retained.len(),
MAX_CALLBACK_OUTPUT_BYTES,
"owner memory stays inside the configured bound"
);
assert_eq!(
drained.total,
produced.len(),
"and every byte was still read, so the writer is never blocked"
);
assert!(drained.truncated());
let text = drained.text();
assert!(
text.ends_with(&format!("({} bytes total)", produced.len())),
"{text}"
);
}
#[tokio::test]
async fn a_drain_under_the_limit_retains_everything_and_reports_no_truncation() {
let produced = b"one bounded line\n".to_vec();
let drained = drain(produced.as_slice(), MAX_CALLBACK_OUTPUT_BYTES).await;
assert_eq!(drained.retained, produced);
assert_eq!(drained.total, produced.len());
assert!(!drained.truncated());
assert_eq!(drained.text(), "one bounded line\n");
}
fn proposal_registry() -> (Arc<CompletionSinkRegistry>, mpsc::UnboundedReceiver<Task>) {
let Wiring { registry, tasks } = CompletionSinkRegistry::build(
"i-1".to_string(),
Arc::new(ExecutionFactsStore::new()),
Arc::new(crate::web::remote_control_api::ExecutionContractHandle::default()),
);
(registry, tasks)
}
fn spec(program: &str) -> ExecutionSinkSpec {
ExecutionSinkSpec {
command: vec![program.to_string()],
notify_blocked: false,
}
}
fn start_episode(registry: &CompletionSinkRegistry, change_id: &str, execution_id: &str) {
registry.bind_episode(change_id, execution_id);
}
fn settle(registry: &CompletionSinkRegistry, execution_id: &str, terminal: EpisodeTerminal) {
let mut entries = registry.lock();
let entry = entries
.get_mut(execution_id)
.expect("the episode must exist before it settles");
entry.terminal = Some(terminal);
}
fn queued_registrations(
tasks: &mut mpsc::UnboundedReceiver<Task>,
execution_id: &str,
) -> usize {
let mut count = 0;
while let Ok(task) = tasks.try_recv() {
if matches!(&task, Task::Registered { execution_id: queued } if queued == execution_id)
{
count += 1;
}
}
count
}
#[test]
fn a_subscription_binds_the_owner_incarnation_it_names() {
let (registry, _tasks) = proposal_registry();
assert!(matches!(
registry.set_proposal_subscription("alpha", "i-2", spec("/bin/true")),
Err(SinkRefusal::InstanceMismatch)
));
assert!(matches!(
registry.view_proposal_subscription("alpha", "i-2"),
Err(SinkRefusal::InstanceMismatch)
));
assert!(matches!(
registry.clear_proposal_subscription("alpha", "i-2"),
Err(SinkRefusal::InstanceMismatch)
));
let view = registry
.view_proposal_subscription("alpha", "i-1")
.expect("this incarnation");
assert!(view.sink.is_none());
}
#[test]
fn an_unacceptable_argv_is_refused_before_anything_is_recorded() {
let (registry, _tasks) = proposal_registry();
let refusal = registry.set_proposal_subscription(
"alpha",
"i-1",
ExecutionSinkSpec {
command: Vec::new(),
notify_blocked: false,
},
);
assert!(matches!(refusal, Err(SinkRefusal::InvalidCommand(_))));
assert!(registry
.view_proposal_subscription("alpha", "i-1")
.unwrap()
.sink
.is_none());
}
#[test]
fn a_subscription_registered_before_admission_binds_the_next_episode() {
let (registry, mut tasks) = proposal_registry();
let view = registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/true"))
.expect("a pre-admission subscription is legal");
assert!(view.sink.is_some());
assert_eq!(view.execution_id, None, "no episode is synthesized");
assert_eq!(queued_registrations(&mut tasks, "e-1"), 0);
start_episode(®istry, "alpha", "e-1");
let view = registry.view_proposal_subscription("alpha", "i-1").unwrap();
assert_eq!(view.execution_id.as_deref(), Some("e-1"));
settle(®istry, "e-1", EpisodeTerminal::Completed);
assert!(
registry.claim_terminal("e-1").is_some(),
"the bound episode owes exactly one delivery attempt"
);
}
#[test]
fn a_late_subscription_claims_the_retained_terminal_episode_once() {
let (registry, mut tasks) = proposal_registry();
start_episode(®istry, "alpha", "e-1");
settle(®istry, "e-1", EpisodeTerminal::Failed);
assert!(registry.claim_terminal("e-1").is_none());
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/true"))
.expect("a late subscription");
assert_eq!(
queued_registrations(&mut tasks, "e-1"),
1,
"the late registration asks the dispatcher to look at the settled episode"
);
assert!(registry.claim_terminal("e-1").is_some());
assert!(
registry.claim_terminal("e-1").is_none(),
"and only once, however many times it is asked"
);
}
#[test]
fn replacing_or_clearing_a_subscription_never_replays_a_delivered_terminal() {
for replay in ["replace", "clear then set"] {
let (registry, mut tasks) = proposal_registry();
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/true"))
.unwrap();
start_episode(®istry, "alpha", "e-1");
settle(®istry, "e-1", EpisodeTerminal::Stopped);
assert!(registry.claim_terminal("e-1").is_some(), "{replay}");
if replay == "clear then set" {
registry
.clear_proposal_subscription("alpha", "i-1")
.unwrap();
}
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/false"))
.unwrap();
assert_eq!(queued_registrations(&mut tasks, "e-1"), 1, "{replay}");
assert!(
registry.claim_terminal("e-1").is_none(),
"{replay} must not replay e-1"
);
start_episode(®istry, "alpha", "e-2");
settle(®istry, "e-2", EpisodeTerminal::Stopped);
assert!(registry.claim_terminal("e-2").is_some(), "{replay}");
}
}
#[test]
fn re_admission_creates_a_distinct_episode_under_the_same_subscription() {
let (registry, _tasks) = proposal_registry();
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/true"))
.unwrap();
start_episode(®istry, "alpha", "e-1");
settle(®istry, "e-1", EpisodeTerminal::Failed);
assert!(registry.claim_terminal("e-1").is_some());
start_episode(®istry, "alpha", "e-2");
let view = registry.view_proposal_subscription("alpha", "i-1").unwrap();
assert_eq!(view.execution_id.as_deref(), Some("e-2"));
assert!(view.sink.is_some(), "the subscription outlives the episode");
assert!(
!registry.lock().contains_key("e-1"),
"a superseded proposal-bound episode is discarded"
);
settle(®istry, "e-2", EpisodeTerminal::Failed);
assert!(
registry.claim_terminal("e-2").is_some(),
"e-1's dedupe must not suppress e-2"
);
}
#[test]
fn clearing_removes_only_the_named_proposal() {
let (registry, _tasks) = proposal_registry();
for change_id in ["alpha", "beta", "gamma"] {
registry
.set_proposal_subscription(change_id, "i-1", spec("/bin/true"))
.unwrap();
}
registry
.clear_proposal_subscription("alpha", "i-1")
.unwrap();
registry
.clear_proposal_subscription("gamma", "i-1")
.unwrap();
for cleared in ["alpha", "gamma"] {
let view = registry.view_proposal_subscription(cleared, "i-1").unwrap();
assert!(view.sink.is_none(), "{cleared}");
}
assert!(registry
.view_proposal_subscription("beta", "i-1")
.unwrap()
.sink
.is_some());
}
#[test]
fn clearing_cancels_an_unstarted_delivery_for_the_live_episode() {
let (registry, _tasks) = proposal_registry();
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/true"))
.unwrap();
start_episode(®istry, "alpha", "e-1");
settle(®istry, "e-1", EpisodeTerminal::Completed);
registry
.clear_proposal_subscription("alpha", "i-1")
.unwrap();
assert!(
registry.claim_terminal("e-1").is_none(),
"a cleared subscription has nobody to deliver to"
);
let view = registry.view_proposal_subscription("alpha", "i-1").unwrap();
assert_eq!(view.execution_id.as_deref(), Some("e-1"));
}
#[test]
fn a_proposal_subscription_never_disturbs_an_execution_scoped_registration() {
let (registry, _tasks) = proposal_registry();
start_episode(®istry, "alpha", "e-1");
registry
.set_sink("e-1", "i-1", "alpha", spec("/bin/direct"))
.expect("a direct registration");
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/proposal"))
.unwrap();
let view = registry.view("e-1", "i-1", "alpha").expect("still there");
assert_eq!(
view.sink.expect("the direct registration stands").command,
vec!["/bin/direct".to_string()],
"a standing proposal rule must not replace a registration that named this episode"
);
registry
.clear_proposal_subscription("alpha", "i-1")
.unwrap();
let view = registry.view("e-1", "i-1", "alpha").expect("still there");
assert_eq!(
view.sink.expect("the direct registration survives").command,
vec!["/bin/direct".to_string()],
"clearing a proposal must not detach an execution-scoped sink"
);
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/proposal"))
.unwrap();
start_episode(®istry, "alpha", "e-2");
assert!(
registry.lock().contains_key("e-1"),
"an episode carrying a direct registration is not pruned by a retry"
);
assert_eq!(
registry
.lock()
.get("e-2")
.unwrap()
.sink
.as_ref()
.unwrap()
.command,
vec!["/bin/proposal".to_string()]
);
}
#[test]
fn a_subscription_replaced_after_delivery_applies_to_the_next_episode_only() {
let (registry, _tasks) = proposal_registry();
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/first"))
.unwrap();
start_episode(®istry, "alpha", "e-1");
settle(®istry, "e-1", EpisodeTerminal::Completed);
assert!(registry.claim_terminal("e-1").is_some());
registry.lock().get_mut("e-1").unwrap().terminal_dispatched = true;
registry
.set_proposal_subscription("alpha", "i-1", spec("/bin/second"))
.unwrap();
assert_eq!(
registry
.lock()
.get("e-1")
.unwrap()
.sink
.as_ref()
.unwrap()
.command,
vec!["/bin/first".to_string()],
"a closed episode keeps the argv it was delivered with"
);
start_episode(®istry, "alpha", "e-2");
assert_eq!(
registry
.lock()
.get("e-2")
.unwrap()
.sink
.as_ref()
.unwrap()
.command,
vec!["/bin/second".to_string()],
"the replacement applies to the next episode"
);
}
#[test]
fn the_published_proposal_capability_matches_the_enforced_bound() {
let capability = proposal_capability();
assert!(capability.available);
assert_eq!(capability.max_targets, MAX_PROPOSAL_TARGETS);
}
#[test]
fn the_published_capability_matches_the_enforced_limits() {
let capability = capability();
assert!(capability.available);
assert_eq!(capability.max_command_args, MAX_COMMAND_ARGS);
assert_eq!(capability.max_command_arg_len, MAX_COMMAND_ARG_LEN);
assert_eq!(
capability.callback_timeout_ms,
CALLBACK_TIMEOUT.as_millis() as u64
);
assert_eq!(
capability.max_callback_output_bytes,
MAX_CALLBACK_OUTPUT_BYTES
);
}
}