use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use crate::child::{AbandonIntent, ChildRef, ObservationRecipient};
pub use crate::durable::ProjectId;
use crate::durable::WorkStatus;
use crate::id::WaveId;
use crate::planning::ProjectPlan;
use crate::task::{TaskEventKind, TaskId};
pub(crate) mod answer;
pub mod runner;
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ProjectDataError {
#[error("invalid Project id: {0}")]
InvalidId(String),
#[error("invalid Project: {0}")]
InvalidInvariant(String),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Project {
pub id: ProjectId,
pub plan: ProjectPlan,
pub wave_id: WaveId,
pub iteration: u32,
pub observation_cursor: i64,
pub last_state_fingerprint: Option<String>,
pub agent: String,
pub provider: String,
pub provider_session_id: Option<String>,
pub abandon_intent: Option<AbandonIntent>,
pub created_at: OffsetDateTime,
pub updated_at: OffsetDateTime,
}
impl Project {
pub fn supervisor_restart_bar(&self) -> Option<String> {
self.abandon_intent.as_ref().map(|intent| {
format!(
"Project {} is being abandoned: {}",
self.plan.slug, intent.reason
)
})
}
pub fn validate(&self) -> Result<(), ProjectDataError> {
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ProjectEventKind {
Started,
BodyHandedOff {
handoff: crate::child::ChildBodyHandoff,
},
TaskObserved {
task_id: TaskId,
task_event_id: i64,
event: Box<TaskEventKind>,
},
IterationCompleted {
iteration: u32,
summary: String,
},
Completed {
summary: String,
},
Failed {
error: String,
resumable: bool,
},
}
impl ProjectEventKind {
pub fn is_wave_observable(&self) -> bool {
!matches!(self, Self::Started | Self::TaskObserved { .. })
}
fn resumable_failure_reason(&self) -> Option<&str> {
match self {
Self::Failed {
error,
resumable: true,
} => Some(error),
_ => None,
}
}
}
pub(crate) fn status_reason(status: &WorkStatus, latest: Option<&ProjectEventKind>) -> String {
if matches!(status, WorkStatus::Ready | WorkStatus::Waiting { .. }) {
if let Some(reason) = latest.and_then(ProjectEventKind::resumable_failure_reason) {
return reason.to_string();
}
}
match status {
WorkStatus::Running { run_id } => format!("Run {run_id} is active"),
WorkStatus::Waiting { .. } => "waiting for input or an event".to_string(),
WorkStatus::Ready => "ready".to_string(),
WorkStatus::Done => "done".to_string(),
WorkStatus::Abandoned => "abandoned".to_string(),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ProjectEvent {
pub id: i64,
pub project_id: ProjectId,
pub kind: ProjectEventKind,
pub created_at: OffsetDateTime,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ProjectObservation {
pub project_id: ProjectId,
pub project: String,
pub event_id: i64,
pub event: ProjectEventKind,
}
impl ProjectObservation {
pub fn inbox_id(&self) -> String {
format!("project-{}-{}", self.project_id, self.event_id)
}
pub fn prompt(&self) -> String {
let payload = serde_json::to_string(&self.event)
.expect("Project observation always serializes to structured JSON");
format!(
"<project_observation project_id=\"{}\" project=\"{}\" event_id=\"{}\">\n{}\n</project_observation>",
self.project_id, self.project, self.event_id, payload
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ChildEventPayload {
Project { event: ProjectEventKind },
Task { event: TaskEventKind },
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ObservationOutboxRow {
pub id: i64,
pub recipient: ObservationRecipient,
pub source: ChildRef,
pub event_id: i64,
pub payload: ChildEventPayload,
pub delivered_at: Option<OffsetDateTime>,
}