loopflow 0.9.11

Run steps and flows with coding agents
Documentation
use crate::lfd::id::LfdId;
use crate::lfd::store::SharedStore;
use crate::lfd::types::{
    AttentionItem, AttentionKind, AttentionStatus, QueueBlock, QueueBlockReason, Wave, WaveRun,
    WaveRunStatus,
};
use serde_json::json;
use time::OffsetDateTime;

/// Stable ID for queue-block attention items: reuse the run_id so upserts converge.
pub fn attention_id_for_queue_block(run_id: &LfdId) -> LfdId {
    run_id.clone()
}

/// Daemon-created algedonic attention: repair chain exhausted.
pub async fn create_step_failure_attention(
    store: &SharedStore,
    wave: &Wave,
    run: &WaveRun,
    step_name: &str,
    error: &str,
) -> Result<AttentionItem, String> {
    let item = AttentionItem {
        id: LfdId::new(),
        wave_id: wave.id().clone(),
        run_id: Some(run.id.clone()),
        kind: AttentionKind::Algedonic,
        status: AttentionStatus::Surfaced,
        title: format!("Step failed: {step_name}"),
        summary: error.to_string(),
        context: json!({
            "step": step_name,
            "error": error,
        }),
        surfaced_at: OffsetDateTime::now_utc(),
        viewed_at: None,
        resolved_at: None,
    };
    store
        .upsert_attention_item(&item)
        .await
        .map_err(|err| format!("upsert step failure attention failed: {err}"))?;
    Ok(item)
}

/// Resolve an attention item by ID.
pub async fn resolve_attention_item(
    store: &SharedStore,
    attention_id: &LfdId,
) -> Result<Option<AttentionItem>, String> {
    let Some(mut item) = store
        .get_attention_item(attention_id)
        .await
        .map_err(|err| format!("get attention item failed: {err}"))?
    else {
        return Ok(None);
    };
    if item.status == AttentionStatus::Resolved {
        return Ok(Some(item));
    }
    item.status = AttentionStatus::Resolved;
    item.resolved_at = Some(OffsetDateTime::now_utc());
    store
        .upsert_attention_item(&item)
        .await
        .map_err(|err| format!("resolve attention item failed: {err}"))?;
    Ok(Some(item))
}

pub async fn mark_attention_viewed(
    store: &SharedStore,
    attention_id: &LfdId,
) -> Result<Option<AttentionItem>, String> {
    let Some(mut item) = store
        .get_attention_item(attention_id)
        .await
        .map_err(|err| format!("get attention item failed: {err}"))?
    else {
        return Ok(None);
    };
    if item.status == AttentionStatus::Resolved {
        return Ok(Some(item));
    }
    if item.status == AttentionStatus::Surfaced {
        item.status = AttentionStatus::Viewed;
        item.viewed_at = Some(OffsetDateTime::now_utc());
        store
            .upsert_attention_item(&item)
            .await
            .map_err(|err| format!("update attention item failed: {err}"))?;
    }
    Ok(Some(item))
}

/// Reconciliation: auto-resolve stale attention items based on context.
///
/// For algedonic items, checks whether the underlying condition has cleared:
/// - Queue blocks with `reason` in context: never auto-resolve (requires manual intervention)
/// - Items with `error` in context: resolve when the run is no longer failed or a newer run exists
/// - All others: resolve when the wave has a newer run
///
/// For interactive items: resolve when the wave has a newer run (the human moved on).
pub async fn reconcile_attention_items(store: &SharedStore) -> Result<Vec<AttentionItem>, String> {
    let mut resolved = Vec::new();
    let items = store
        .list_attention_items(None, None)
        .await
        .map_err(|err| format!("list attention items failed: {err}"))?;

    for mut item in items {
        if item.status == AttentionStatus::Resolved {
            continue;
        }
        let should_resolve = if is_queue_block(&item) {
            // Queue blocks require manual intervention.
            false
        } else if has_step_failure(&item) {
            should_resolve_step_failure(store, &item).await?
        } else {
            // Interactive items and other algedonic items resolve when the wave restarts.
            should_resolve_when_wave_restarted(store, &item).await?
        };
        if should_resolve {
            item.status = AttentionStatus::Resolved;
            item.resolved_at = Some(OffsetDateTime::now_utc());
            store
                .upsert_attention_item(&item)
                .await
                .map_err(|err| format!("resolve attention item failed: {err}"))?;
            resolved.push(item);
        }
    }

    Ok(resolved)
}

/// Queue blocks have a `reason` field in context (from QueueBlockReason).
fn is_queue_block(item: &AttentionItem) -> bool {
    item.context.get("reason").is_some() && item.context.get("conflict_files").is_some()
}

/// Step failures have an `error` field and a `step` field but no `reason` field.
fn has_step_failure(item: &AttentionItem) -> bool {
    item.context.get("error").is_some() && item.context.get("reason").is_none()
}

async fn should_resolve_step_failure(
    store: &SharedStore,
    item: &AttentionItem,
) -> Result<bool, String> {
    let Some(run_id) = item.run_id.as_ref() else {
        return Ok(true);
    };
    let Some(run) = store
        .get_wave_run(run_id)
        .await
        .map_err(|err| format!("get wave run failed: {err}"))?
    else {
        return Ok(true);
    };
    if run.status != WaveRunStatus::Failed {
        return Ok(true);
    }
    let latest = store
        .get_latest_wave_run(&item.wave_id)
        .await
        .map_err(|err| format!("get latest wave run failed: {err}"))?;
    Ok(latest.is_some_and(|latest| latest.id != run.id))
}

async fn should_resolve_when_wave_restarted(
    store: &SharedStore,
    item: &AttentionItem,
) -> Result<bool, String> {
    let latest = store
        .get_latest_wave_run(&item.wave_id)
        .await
        .map_err(|err| format!("get latest wave run failed: {err}"))?;
    Ok(latest.is_some_and(|run| {
        item.run_id
            .as_ref()
            .is_none_or(|prev_id| &run.id != prev_id)
    }))
}

// Queue block <-> attention item conversion helpers.

pub fn queue_block_from_attention(item: &AttentionItem) -> Result<Option<QueueBlock>, String> {
    let Some(run_id) = item.run_id.clone() else {
        return Ok(None);
    };
    let reason = item
        .context
        .get("reason")
        .and_then(serde_json::Value::as_str)
        .unwrap_or(QueueBlockReason::PromotionFailed.as_str())
        .parse()
        .map_err(|err| format!("invalid queue block reason: {err}"))?;
    let conflict_files = item
        .context
        .get("conflict_files")
        .and_then(serde_json::Value::as_array)
        .map(|files| {
            files
                .iter()
                .filter_map(serde_json::Value::as_str)
                .map(ToString::to_string)
                .collect::<Vec<_>>()
        })
        .unwrap_or_default();
    let error = item
        .context
        .get("error")
        .and_then(serde_json::Value::as_str)
        .map(ToString::to_string);
    Ok(Some(QueueBlock {
        wave_id: item.wave_id.clone(),
        run_id,
        reason,
        attempted_at: item.surfaced_at,
        conflict_files,
        error,
    }))
}

pub fn queue_block_attention_item(block: &QueueBlock) -> AttentionItem {
    AttentionItem {
        id: attention_id_for_queue_block(&block.run_id),
        wave_id: block.wave_id.clone(),
        run_id: Some(block.run_id.clone()),
        kind: AttentionKind::Algedonic,
        status: AttentionStatus::Surfaced,
        title: format!("Queue blocked: {}", block.reason.as_str().replace('_', " ")),
        summary: block
            .error
            .clone()
            .unwrap_or_else(|| "Queue requires attention before it can advance.".to_string()),
        context: json!({
            "reason": block.reason.as_str(),
            "conflict_files": block.conflict_files,
            "error": block.error,
        }),
        surfaced_at: block.attempted_at,
        viewed_at: None,
        resolved_at: None,
    }
}

pub fn queue_block_attention_item_from_existing(
    block: &QueueBlock,
    existing: Option<&AttentionItem>,
) -> AttentionItem {
    let mut item = queue_block_attention_item(block);
    let Some(existing) = existing else {
        return item;
    };
    if existing.kind != AttentionKind::Algedonic || existing.status == AttentionStatus::Resolved {
        return item;
    }
    item.status = existing.status;
    item.surfaced_at = existing.surfaced_at;
    item.viewed_at = existing.viewed_at;
    item
}

#[cfg(test)]
mod tests {
    use super::{
        queue_block_attention_item, queue_block_attention_item_from_existing,
        queue_block_from_attention, resolve_attention_item,
    };
    use crate::lfd::id::LfdId;
    use crate::lfd::store::{open_store, SharedStore, StorageConfig};
    use crate::lfd::types::{
        AttentionItem, AttentionKind, AttentionStatus, QueueBlock, QueueBlockReason,
    };
    use std::sync::Arc;
    use time::OffsetDateTime;

    #[test]
    fn queue_block_helpers_round_trip() {
        let block = QueueBlock {
            wave_id: LfdId::new(),
            run_id: LfdId::new(),
            reason: QueueBlockReason::RebaseConflict,
            attempted_at: OffsetDateTime::now_utc(),
            conflict_files: vec!["src/lib.rs".to_string()],
            error: Some("merge failed".to_string()),
        };

        let item = queue_block_attention_item(&block);
        assert_eq!(item.kind, AttentionKind::Algedonic);
        let restored = queue_block_from_attention(&item)
            .expect("queue block context parses")
            .expect("queue block exists");

        assert_eq!(restored.wave_id, block.wave_id);
        assert_eq!(restored.run_id, block.run_id);
        assert_eq!(restored.reason, block.reason);
        assert_eq!(restored.attempted_at, block.attempted_at);
        assert_eq!(restored.conflict_files, block.conflict_files);
        assert_eq!(restored.error, block.error);
    }

    #[test]
    fn queue_block_attention_preserves_open_lifecycle_fields() {
        let block = QueueBlock {
            wave_id: LfdId::new(),
            run_id: LfdId::new(),
            reason: QueueBlockReason::ScratchDirty,
            attempted_at: OffsetDateTime::now_utc(),
            conflict_files: Vec::new(),
            error: None,
        };
        let mut existing = queue_block_attention_item(&block);
        existing.status = AttentionStatus::Viewed;
        existing.viewed_at = Some(existing.surfaced_at + time::Duration::minutes(1));
        existing.surfaced_at -= time::Duration::hours(1);

        let refreshed = queue_block_attention_item_from_existing(&block, Some(&existing));

        assert_eq!(refreshed.status, AttentionStatus::Viewed);
        assert_eq!(refreshed.surfaced_at, existing.surfaced_at);
        assert_eq!(refreshed.viewed_at, existing.viewed_at);
    }

    async fn test_store() -> SharedStore {
        let db_path = std::env::temp_dir().join(format!("lfd-attention-test-{}.db", LfdId::new()));
        let config = StorageConfig::sqlite(db_path);
        Arc::new(open_store(&config).await.expect("sqlite store"))
    }

    #[tokio::test]
    async fn resolve_attention_item_by_id() {
        use crate::lfd::types::Wave;
        let store = test_store().await;
        let wave_id = LfdId::new();
        let wave = Wave::new(wave_id.clone(), "test".to_string(), "/tmp/repo".to_string());
        store.create_wave(&wave).await.unwrap();
        let item = AttentionItem {
            id: LfdId::new(),
            wave_id,
            run_id: None,
            kind: AttentionKind::Interactive,
            status: AttentionStatus::Surfaced,
            title: "test".to_string(),
            summary: "test".to_string(),
            context: serde_json::json!({}),
            surfaced_at: OffsetDateTime::now_utc(),
            viewed_at: None,
            resolved_at: None,
        };
        store.upsert_attention_item(&item).await.unwrap();

        let resolved = resolve_attention_item(&store, &item.id)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(resolved.status, AttentionStatus::Resolved);
        assert!(resolved.resolved_at.is_some());

        // Resolving again returns the already-resolved item
        let again = resolve_attention_item(&store, &item.id)
            .await
            .unwrap()
            .unwrap();
        assert_eq!(again.status, AttentionStatus::Resolved);
    }
}