use std::path::PathBuf;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use crate::id::WaveId;
use crate::project::ProjectId;
use crate::task::TaskId;
macro_rules! prefixed_uuid_id {
($name:ident, $prefix:literal, $error:ty, $invalid:path) => {
#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
#[serde(transparent)]
pub struct $name(String);
impl $name {
pub fn new() -> Self {
Self(format!("{}{}", $prefix, uuid::Uuid::new_v4().simple()))
}
pub fn parse(value: &str) -> Result<Self, $error> {
let suffix = value
.strip_prefix($prefix)
.ok_or_else(|| $invalid(format!("expected {} id", $prefix)))?;
uuid::Uuid::parse_str(suffix).map_err(|error| $invalid(error.to_string()))?;
Ok(Self::from_raw(value.to_string()))
}
pub fn as_str(&self) -> &str {
&self.0
}
pub(crate) fn from_raw(value: impl Into<String>) -> Self {
Self(value.into())
}
}
impl Default for $name {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Display for $name {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(&self.0)
}
}
impl std::str::FromStr for $name {
type Err = $error;
fn from_str(value: &str) -> Result<Self, Self::Err> {
Self::parse(value)
}
}
};
}
pub(crate) use prefixed_uuid_id;
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ChildDataError {
#[error("invalid child-session id: {0}")]
InvalidId(String),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ObservationRecipient {
Wave { wave_id: WaveId },
Project { project_id: ProjectId },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ChildBodyHandoffRequest {
pub agent: String,
pub provider: String,
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ChildBodyHandoff {
pub from_agent: String,
pub to_agent: String,
pub from_provider: String,
pub to_provider: String,
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ChildExecutionContext {
pub lf_bin: PathBuf,
pub db_path: PathBuf,
pub lf_home: PathBuf,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AbandonIntent {
pub requested_at: OffsetDateTime,
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", content = "id", rename_all = "snake_case")]
pub enum ChildRef {
Project(ProjectId),
Task(TaskId),
}
impl ChildRef {
pub fn target_kind(&self) -> &'static str {
match self {
Self::Project(_) => "project",
Self::Task(_) => "task",
}
}
pub fn target_id(&self) -> &str {
match self {
Self::Project(id) => id.as_str(),
Self::Task(id) => id.as_str(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum BodyCategory {
Working,
Stalled,
Recovering,
NeedsInput,
Stopped,
Terminal,
Unobservable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum BodyOwner {
Work,
Loopflow,
User,
Nobody,
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum BodyControl {
Attach,
Steer,
Interrupt,
Stop,
Extend,
Resume,
Decide,
Abandon,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BodyObservation {
pub category: BodyCategory,
pub reason: String,
pub owner: BodyOwner,
pub controls: Vec<BodyControl>,
pub progress_age_secs: Option<u64>,
pub deadline_in_secs: Option<i64>,
pub step: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BodyIntent {
Active,
Waiting,
Terminal,
}
#[derive(Debug, Clone)]
pub struct BodyEvidence {
pub intent: BodyIntent,
pub observable: bool,
pub process_alive: bool,
pub progress_age: Duration,
pub step: Option<String>,
pub reason: String,
}
pub const DEFAULT_STALL_AFTER: Duration = Duration::from_secs(30 * 60);
pub(crate) fn body_progress_age(
latest_event_at: Option<OffsetDateTime>,
status_at: OffsetDateTime,
now: OffsetDateTime,
) -> Duration {
let progress_at = latest_event_at.map_or(status_at, |event_at| event_at.max(status_at));
let seconds = (now - progress_at).whole_seconds().max(0);
Duration::from_secs(seconds as u64)
}
pub fn observe(evidence: &BodyEvidence, stall_after: Duration) -> BodyObservation {
let make = |category, reason: &str, owner, controls, progress, deadline| BodyObservation {
category,
reason: reason.to_string(),
owner,
controls,
progress_age_secs: progress,
deadline_in_secs: deadline,
step: evidence.step.clone(),
};
match evidence.intent {
BodyIntent::Terminal => make(
BodyCategory::Terminal,
&evidence.reason,
BodyOwner::Nobody,
vec![],
None,
None,
),
BodyIntent::Waiting => make(
BodyCategory::NeedsInput,
&evidence.reason,
BodyOwner::User,
vec![
BodyControl::Decide,
BodyControl::Resume,
BodyControl::Abandon,
],
None,
None,
),
BodyIntent::Active => {
if !evidence.observable {
return make(
BodyCategory::Unobservable,
"this machine cannot observe the body",
BodyOwner::Unknown,
vec![],
None,
None,
);
}
if !evidence.process_alive {
return make(
BodyCategory::Stopped,
"no live body for active intent; a wake will adopt or start one",
BodyOwner::Loopflow,
vec![BodyControl::Resume, BodyControl::Stop],
None,
None,
);
}
let progress = Some(evidence.progress_age.as_secs());
let remaining = stall_after.as_secs() as i64 - evidence.progress_age.as_secs() as i64;
if evidence.progress_age > stall_after {
make(
BodyCategory::Stalled,
"alive but no meaningful progress past the deadline",
BodyOwner::Loopflow,
vec![
BodyControl::Attach,
BodyControl::Extend,
BodyControl::Interrupt,
BodyControl::Stop,
],
progress,
Some(remaining),
)
} else {
make(
BodyCategory::Working,
&evidence.reason,
BodyOwner::Work,
vec![
BodyControl::Attach,
BodyControl::Steer,
BodyControl::Interrupt,
BodyControl::Stop,
],
progress,
Some(remaining),
)
}
}
}
}
#[cfg(test)]
mod tests {
use super::{
body_progress_age, observe, BodyCategory, BodyControl, BodyEvidence, BodyIntent, BodyOwner,
Duration, DEFAULT_STALL_AFTER,
};
fn evidence(intent: BodyIntent, alive: bool, progress: Duration) -> BodyEvidence {
BodyEvidence {
intent,
observable: true,
process_alive: alive,
progress_age: progress,
step: Some("task_pursue".to_string()),
reason: "running".to_string(),
}
}
#[test]
fn a_live_body_that_just_progressed_is_working() {
let obs = observe(
&evidence(BodyIntent::Active, true, Duration::from_secs(60)),
DEFAULT_STALL_AFTER,
);
assert_eq!(obs.category, BodyCategory::Working);
assert_eq!(obs.owner, BodyOwner::Work);
assert_eq!(obs.progress_age_secs, Some(60));
assert!(obs.controls.contains(&BodyControl::Steer));
assert!(obs.deadline_in_secs.unwrap() > 0);
}
#[test]
fn a_live_body_past_its_deadline_is_stalled() {
let stalled = observe(
&evidence(BodyIntent::Active, true, Duration::from_secs(31 * 60)),
DEFAULT_STALL_AFTER,
);
assert_eq!(stalled.category, BodyCategory::Stalled);
assert_eq!(stalled.owner, BodyOwner::Loopflow);
assert_eq!(stalled.progress_age_secs, Some(31 * 60));
assert!(stalled.deadline_in_secs.unwrap() < 0);
assert!(stalled.controls.contains(&BodyControl::Extend));
}
#[test]
fn the_stall_boundary_is_the_threshold_exactly() {
let clock = Duration::from_secs(10);
assert_eq!(
observe(&evidence(BodyIntent::Active, true, clock), clock).category,
BodyCategory::Working,
);
assert_eq!(
observe(
&evidence(BodyIntent::Active, true, clock + Duration::from_secs(1)),
clock,
)
.category,
BodyCategory::Stalled,
);
}
#[test]
fn progress_age_uses_the_freshest_durable_evidence() {
let status_at = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(10);
let event_at = status_at + time::Duration::seconds(5);
let now = status_at + time::Duration::seconds(20);
assert_eq!(
body_progress_age(Some(event_at), status_at, now),
Duration::from_secs(15)
);
assert_eq!(
body_progress_age(Some(now + time::Duration::seconds(5)), status_at, now),
Duration::ZERO,
);
}
#[test]
fn active_intent_with_no_live_body_is_stopped_not_gone() {
let obs = observe(
&evidence(BodyIntent::Active, false, Duration::from_secs(5)),
DEFAULT_STALL_AFTER,
);
assert_eq!(obs.category, BodyCategory::Stopped);
assert_eq!(obs.owner, BodyOwner::Loopflow);
assert_eq!(obs.progress_age_secs, None);
}
#[test]
fn an_unobservable_body_is_never_asserted_gone() {
let mut ev = evidence(BodyIntent::Active, false, Duration::from_secs(5));
ev.observable = false;
let obs = observe(&ev, DEFAULT_STALL_AFTER);
assert_eq!(obs.category, BodyCategory::Unobservable);
assert_eq!(obs.owner, BodyOwner::Unknown);
assert!(obs.controls.is_empty());
}
#[test]
fn waiting_intent_needs_user_input() {
let obs = observe(
&evidence(BodyIntent::Waiting, false, Duration::from_secs(5)),
DEFAULT_STALL_AFTER,
);
assert_eq!(obs.category, BodyCategory::NeedsInput);
assert_eq!(obs.owner, BodyOwner::User);
assert!(obs.controls.contains(&BodyControl::Decide));
}
#[test]
fn terminal_intent_owns_nobody_and_offers_no_controls() {
let obs = observe(
&evidence(BodyIntent::Terminal, false, Duration::from_secs(0)),
DEFAULT_STALL_AFTER,
);
assert_eq!(obs.category, BodyCategory::Terminal);
assert_eq!(obs.owner, BodyOwner::Nobody);
assert!(obs.controls.is_empty());
}
}