use std::fmt;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use tokio::process::Command;
use crate::id::WaveId;
use crate::store::{SharedStore, StoreResult};
use crate::task::TaskObservation;
use crate::wave::runtime::WaveRuntime;
use crate::wave::Wave;
pub const POLL_CADENCE: Duration = Duration::from_secs(10);
#[derive(Debug, Clone)]
pub struct RegistryConfig {
pub store: SharedStore,
pub wave: Wave,
}
pub async fn ensure_wave_row(
store: &SharedStore,
main_repo: &Path,
name: &str,
) -> StoreResult<Wave> {
let existing = store.get_wave_by_name(name).await?;
let is_new = existing.is_none();
let wave = existing.unwrap_or_else(|| {
Wave::new(
WaveId::new(),
name.to_string(),
main_repo.display().to_string(),
)
});
store.create_wave(&wave).await?;
if is_new {
tracing::info!(
wave = name,
wave_id = %wave.id,
"wave was not in the registry; created its row"
);
}
Ok(wave)
}
pub(crate) async fn process_alive(pid: u32) -> bool {
Command::new("kill")
.args(["-0", &pid.to_string()])
.output()
.await
.is_ok_and(|output| output.status.success())
}
pub struct StoreObserver {
runtime: Arc<WaveRuntime>,
store: SharedStore,
wave_id: WaveId,
}
impl fmt::Debug for StoreObserver {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("StoreObserver")
.field("wave_id", &self.wave_id)
.finish()
}
}
impl StoreObserver {
pub fn new(runtime: Arc<WaveRuntime>, store: SharedStore, wave_id: WaveId) -> Self {
Self {
runtime,
store,
wave_id,
}
}
pub async fn run(self: Arc<Self>, cadence: Duration) {
loop {
self.poll_once().await;
tokio::time::sleep(cadence).await;
}
}
pub async fn poll_once(&self) {
self.poll_child_observations().await;
}
async fn poll_child_observations(&self) {
let recipient = crate::child_session::ObservationRecipient::Wave {
wave_id: self.wave_id.clone(),
};
let observations = match self.store.pending_observations(&recipient).await {
Ok(observations) => observations,
Err(error) => {
tracing::debug!(%error, "wave observer child outbox read failed");
return;
}
};
for observation in observations {
let should_ack = match (observation.source, observation.payload) {
(
crate::child_session::ChildRef::Task(session_id),
crate::project_session::ChildEventPayload::Task { event },
) => match self.store.get_task_session(&session_id).await {
Ok(Some(session)) => {
self.runtime.deliver_task_observation(TaskObservation {
session_id,
issue_identifier: session.launch.issue.identifier,
event_id: observation.event_id,
event,
});
true
}
Ok(None) => false,
Err(error) => {
tracing::debug!(%error, %session_id, "wave observer Task read failed");
continue;
}
},
(
crate::child_session::ChildRef::Project(session_id),
crate::project_session::ChildEventPayload::Project { event },
) => match self.store.get_project_session(&session_id).await {
Ok(Some(session)) => {
self.runtime.deliver_project_observation(
crate::project_session::ProjectObservation {
session_id,
project: session.launch.project.slug,
event_id: observation.event_id,
event,
},
);
true
}
Ok(None) => false,
Err(error) => {
tracing::debug!(%error, %session_id, "wave observer Project read failed");
continue;
}
},
_ => {
tracing::warn!(
outbox_id = observation.id,
"child observation shape mismatched its source"
);
false
}
};
if should_ack {
let _ = self.store.mark_observation_delivered(observation.id).await;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::{open_store, StorageConfig};
async fn temp_store(tmp: &std::path::Path) -> SharedStore {
Arc::new(
open_store(&StorageConfig::sqlite(tmp.join("loopflow.db")))
.await
.expect("open sqlite store"),
)
}
#[tokio::test]
async fn boot_with_no_wave_row_creates_one_durable_identity() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let repo = tmp.path().join("repo");
let goal_dir = repo.join("wave/ship");
std::fs::create_dir_all(&goal_dir).expect("wave dir");
std::fs::write(
goal_dir.join("GOAL.md"),
"---\ngoal: keep shipping\n---\nShip.\n",
)
.expect("GOAL.md");
let wave = ensure_wave_row(&store, &repo, "ship")
.await
.expect("row created");
let stored = store
.get_wave_by_name("ship")
.await
.expect("lookup")
.expect("row exists");
assert_eq!(stored.id, wave.id);
assert_eq!(stored.repo(), repo.display().to_string());
let again = ensure_wave_row(&store, &repo, "ship")
.await
.expect("idempotent");
assert_eq!(again.id, wave.id, "ensure reuses the existing row");
}
#[tokio::test]
async fn ensure_wave_row_without_goal_md_uses_defaults() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = ensure_wave_row(&store, tmp.path(), "ship")
.await
.expect("row created");
assert_eq!(wave.name(), "ship");
}
}