loopflow 0.12.32

Run steps and flows with coding agents
Documentation
//! Durable Wave identity, authored context, and placement.

pub mod config;
pub mod context;
pub mod metrics;
pub mod relocate;

use serde::Serialize;
use time::OffsetDateTime;

use crate::id::WaveId;
use crate::repository::{CanonicalRepo, CanonicalRepoError};

/// A Wave's mutable readable address inside one canonical repository.
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct WaveLocator {
    repo: CanonicalRepo,
    slug: String,
}

#[derive(Debug, thiserror::Error)]
pub enum WaveLocatorError {
    #[error(transparent)]
    Repository(#[from] CanonicalRepoError),
    #[error("invalid Wave slug {0:?}")]
    InvalidSlug(String),
}

impl WaveLocator {
    pub fn discover(repo: &std::path::Path, slug: &str) -> Result<Self, WaveLocatorError> {
        Self::new(CanonicalRepo::discover(repo)?, slug)
    }

    pub fn new(repo: CanonicalRepo, slug: &str) -> Result<Self, WaveLocatorError> {
        let slug = crate::ops::util::normalize_wave_name(slug)
            .ok_or_else(|| WaveLocatorError::InvalidSlug(slug.to_string()))?;
        if slug.contains('\\')
            || slug
                .split('/')
                .any(|component| component.is_empty() || matches!(component, "." | ".."))
        {
            return Err(WaveLocatorError::InvalidSlug(slug));
        }
        Ok(Self { repo, slug })
    }

    pub fn repo(&self) -> &CanonicalRepo {
        &self.repo
    }

    pub fn slug(&self) -> &str {
        &self.slug
    }
}

#[derive(Debug, Clone, Serialize)]
pub struct Wave {
    id: WaveId,
    name: String,
    /// Derived on read from this Wave and its ancestors; never persisted.
    #[serde(skip)]
    slug: String,
    /// Current canonical repository half of the mutable locator.
    repo: String,
    #[serde(with = "time::serde::rfc3339::option")]
    created_at: Option<OffsetDateTime>,
    /// Directory parent; absent for a top-level Wave.
    parent_wave_id: Option<WaveId>,
    /// The first completed promotion occurrence. Ancestry alone leaves this
    /// absent so an older parent link cannot manufacture a new wake.
    #[serde(with = "time::serde::rfc3339::option")]
    promoted_at: Option<OffsetDateTime>,
    #[serde(with = "time::serde::rfc3339::option")]
    retired_at: Option<OffsetDateTime>,
    superseded_by_wave_id: Option<WaveId>,
    retirement_reason: Option<String>,
}

impl Wave {
    pub fn new(id: WaveId, name: String, repo: String) -> Self {
        Self {
            id,
            slug: name.clone(),
            name,
            repo,
            created_at: Some(OffsetDateTime::now_utc()),
            parent_wave_id: None,
            promoted_at: None,
            retired_at: None,
            superseded_by_wave_id: None,
            retirement_reason: None,
        }
    }

    /// Establish an initial directory parent without recording a promotion occurrence.
    /// An existing parent always wins.
    pub fn with_parent(mut self, parent: WaveId) -> Self {
        self.parent_wave_id.get_or_insert(parent);
        self
    }

    #[allow(clippy::too_many_arguments)] // Exact Wave row shape; named accessors expose the domain API.
    pub(crate) fn from_stored_parts(
        id: WaveId,
        name: String,
        repo: String,
        created_at: OffsetDateTime,
        parent_wave_id: Option<WaveId>,
        promoted_at: Option<OffsetDateTime>,
        retired_at: Option<OffsetDateTime>,
        superseded_by_wave_id: Option<WaveId>,
        retirement_reason: Option<String>,
    ) -> Self {
        Self {
            id,
            slug: name.clone(),
            name,
            repo,
            created_at: Some(created_at),
            parent_wave_id,
            promoted_at,
            retired_at,
            superseded_by_wave_id,
            retirement_reason,
        }
    }

    pub fn id(&self) -> &WaveId {
        &self.id
    }

    /// Directory parent, `None` for a root Wave.
    pub fn parent_wave_id(&self) -> Option<&WaveId> {
        self.parent_wave_id.as_ref()
    }

    /// The first completed promotion occurrence, distinct from ancestry.
    pub fn promoted_at(&self) -> Option<OffsetDateTime> {
        self.promoted_at
    }

    pub fn retired_at(&self) -> Option<OffsetDateTime> {
        self.retired_at
    }

    pub fn superseded_by_wave_id(&self) -> Option<&WaveId> {
        self.superseded_by_wave_id.as_ref()
    }

    pub fn retirement_reason(&self) -> Option<&str> {
        self.retirement_reason.as_deref()
    }

    pub fn is_retired(&self) -> bool {
        self.retired_at.is_some()
    }

    pub fn slug(&self) -> &str {
        &self.slug
    }

    pub(crate) fn with_slug(mut self, slug: String) -> Self {
        self.slug = slug;
        self
    }

    pub fn name(&self) -> &str {
        &self.name
    }

    /// The repo this wave targets.
    pub fn repo(&self) -> &str {
        &self.repo
    }

    pub fn created_at(&self) -> Option<OffsetDateTime> {
        self.created_at
    }
}

pub async fn ensure_wave_row(
    store: &crate::store::Store,
    main_repo: &std::path::Path,
    name: &str,
) -> crate::store::StoreResult<Wave> {
    use crate::store::StoreError;
    let locator = WaveLocator::discover(main_repo, name)
        .map_err(|error| StoreError::InvalidData(error.to_string()))?;
    let repo = locator.repo().to_string();
    let mut parent: Option<WaveId> = None;
    let mut prefix = String::new();
    for segment in locator.slug().split('/') {
        if !prefix.is_empty() {
            prefix.push('/');
        }
        prefix.push_str(segment);
        let address = WaveLocator::new(locator.repo().clone(), &prefix)
            .map_err(|error| StoreError::InvalidData(error.to_string()))?;
        let config = config::try_read_wave_config(main_repo, &prefix)
            .map_err(|error| StoreError::InvalidData(error.to_string()))?;
        let at_address = store.get_wave_at(&address).await?;
        let id = config
            .as_ref()
            .and_then(|config| config.id.clone())
            .or_else(|| at_address.as_ref().map(|wave| wave.id().clone()))
            .unwrap_or_default();
        if at_address.as_ref().is_some_and(|wave| wave.id() != &id) {
            return Err(StoreError::InvalidData(format!(
                "Wave {prefix} is already registered with another id"
            )));
        }
        if let Some(existing) = store.get_wave(&id).await? {
            if existing.is_retired() || existing.repo() != repo {
                return Err(StoreError::InvalidData(format!(
                    "Wave {id} belongs to another repository or is retired"
                )));
            }
            if existing.slug() != prefix {
                let previous = config::try_read_wave_config(main_repo, existing.slug())
                    .map_err(|error| StoreError::InvalidData(error.to_string()))?;
                if previous.and_then(|config| config.id).as_ref() == Some(&id) {
                    return Err(StoreError::InvalidData(format!(
                        "Wave {id} appears at both {} and {prefix}",
                        existing.slug()
                    )));
                }
            }
            store
                .reconcile_wave_directory(&id, segment, parent.as_ref())
                .await?;
        } else {
            let mut wave = Wave::new(id.clone(), segment.to_string(), repo.clone());
            if let Some(parent) = &parent {
                wave = wave.with_parent(parent.clone());
            }
            store.create_wave(&wave).await?;
        }
        if config.as_ref().is_none_or(|config| config.id.is_none()) {
            config::update_wave_goal_config(main_repo, &prefix, |map| {
                map.insert(
                    serde_yaml_ng::Value::String("id".into()),
                    serde_yaml_ng::Value::String(id.to_string()),
                );
                Ok(())
            })
            .map_err(StoreError::InvalidData)?;
        }
        parent = Some(id);
    }
    store.get_wave_at(&locator).await?.ok_or_else(|| {
        StoreError::InvalidData(format!(
            "Wave {name} disappeared during directory discovery"
        ))
    })
}

#[cfg(test)]
mod tests {
    use super::{Wave, WaveLocator};
    use crate::id::WaveId;
    use time::OffsetDateTime;

    #[tokio::test]
    async fn directory_discovery_rejects_duplicate_ids_and_preserves_sibling_names() {
        let repo = loopflow_test_support::TestRepo::new();
        let home = tempfile::tempdir().unwrap();
        let store = crate::store::open_ephemeral_store(&crate::store::StorageConfig::sqlite(
            home.path().join("registry.db"),
        ))
        .await
        .unwrap();
        let first = super::ensure_wave_row(&store, repo.path(), "infra/release")
            .await
            .unwrap();
        let other = super::ensure_wave_row(&store, repo.path(), "product/release")
            .await
            .unwrap();
        assert_ne!(first.id(), other.id());
        assert_eq!(first.name(), other.name());
        std::fs::copy(
            repo.path().join("wave/infra/release/GOAL.md"),
            repo.path().join("wave/product/release/GOAL.md"),
        )
        .unwrap();
        assert!(
            super::ensure_wave_row(&store, repo.path(), "product/release")
                .await
                .is_err()
        );
        assert_eq!(
            store.get_wave(other.id()).await.unwrap().unwrap().slug(),
            "product/release"
        );
        assert!(store
            .reconcile_wave_directory(first.parent_wave_id().unwrap(), "infra", Some(first.id()))
            .await
            .is_err());
        assert_eq!(
            store.get_wave(first.id()).await.unwrap().unwrap().slug(),
            "infra/release"
        );
    }

    #[test]
    fn reconstructing_identity_cannot_copy_or_reparent_a_promotion() {
        let parent = WaveId::new();
        let promoted = Wave::from_stored_parts(
            WaveId::new(),
            "ship".to_string(),
            "/repo".to_string(),
            OffsetDateTime::now_utc(),
            Some(parent.clone()),
            Some(OffsetDateTime::now_utc()),
            None,
            None,
            None,
        );

        let rebuilt = Wave::new(
            promoted.id().clone(),
            promoted.name().to_string(),
            promoted.repo().to_string(),
        );
        assert_eq!(rebuilt.parent_wave_id(), None);
        assert_eq!(rebuilt.promoted_at(), None);

        let other_parent = WaveId::new();
        let unchanged = promoted.with_parent(other_parent);
        assert_eq!(unchanged.parent_wave_id(), Some(&parent));
        assert!(unchanged.promoted_at().is_some());
    }

    #[test]
    fn locator_slugs_allow_safe_nesting_but_never_path_traversal() {
        let repo = crate::repository::CanonicalRepo::discover(
            &std::env::current_dir().expect("current directory"),
        )
        .unwrap();
        assert_eq!(
            WaveLocator::new(repo.clone(), "goals/release")
                .unwrap()
                .slug(),
            "goals/release"
        );
        for unsafe_slug in [
            "../escape",
            "goals/../escape",
            "goals//release",
            "goals\\release",
        ] {
            assert!(WaveLocator::new(repo.clone(), unsafe_slug).is_err());
        }
    }
}