use std::time::Duration;
use crate::child::ChildRef;
use crate::durable::WorkRef;
use crate::store::SharedStore;
use crate::work::task::Task;
use super::{OpsError, OpsResult};
pub(crate) const CHILD_STARTUP_GRACE: Duration = Duration::from_secs(10);
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum WorkControlReceipt {
Interrupt { work: WorkRef },
Resume { work: WorkRef },
}
impl WorkControlReceipt {
pub fn label(&self) -> String {
match self {
Self::Interrupt { work } => work.id().to_string(),
Self::Resume { work } => work.id().to_string(),
}
}
pub fn action(&self) -> &'static str {
match self {
Self::Interrupt { .. } => "interrupted",
Self::Resume { .. } => "resumed",
}
}
}
pub(crate) async fn resume_task(store: &SharedStore, mut task: Task) -> OpsResult<WorkRef> {
let label = format!("Task {}", task.plan.identifier);
if let Some(intent) = &task.abandon_intent {
return Err(child_error(format!(
"{label} is being abandoned: {}",
intent.reason
)));
}
let work = store
.work_for_child(&ChildRef::Task(task.id.clone()))
.await
.map_err(child_error)?;
super::task::resume_inactive_process(store, &mut task).await?;
Ok(work)
}
pub(crate) async fn inject_live_steers(
store: &SharedStore,
task_id: &crate::work::task::TaskId,
harness: &mut dyn crate::harness::Harness,
cursor: &mut i64,
) -> Vec<crate::durable::Steer> {
let mut delivered = Vec::new();
let steers = match store.task_steers(task_id).await {
Ok(steers) => steers,
Err(error) => {
tracing::warn!(%error, "failed to read steers for live injection");
return delivered;
}
};
for steer in &steers {
if steer.id <= *cursor {
continue;
}
match harness.send_current(&steer.text).await {
crate::harness::SendCurrentOutcome::Sent { .. } => {
*cursor = steer.id;
delivered.push(steer.clone());
}
crate::harness::SendCurrentOutcome::NotSteerable => break,
crate::harness::SendCurrentOutcome::Failed { error }
| crate::harness::SendCurrentOutcome::Unknown { error, .. } => {
tracing::warn!(
%error,
steer = steer.id,
"live steer injection failed; deferring to the next boundary"
);
break;
}
}
}
delivered
}
pub(crate) async fn observe_interrupt(
store: &SharedStore,
work: &WorkRef,
harness: &mut dyn crate::harness::Harness,
cursor: &mut i64,
) {
match store.latest_interrupt_id(work).await {
Ok(id) if id > *cursor => {
*cursor = id;
if let Err(error) = harness.interrupt().await {
tracing::warn!(%error, "interrupt request failed to reach the provider");
}
}
Ok(_) => {}
Err(error) => tracing::warn!(%error, "failed to read interrupt requests"),
}
}
fn child_error(error: impl std::fmt::Display) -> OpsError {
OpsError::Message(error.to_string())
}