loopflow 0.12.31

Run steps and flows with coding agents
Documentation
use time::OffsetDateTime;

use crate::durable::{
    FlowSession, TaskFlowBlocker, TaskId, TaskWorkerClaim, TaskWorkerClaimOutcome, TaskWorkerOwner,
};
use crate::engine::ExecutionCursor;

use super::{run_sqlite, Store, StoreResult};

impl Store {
    pub async fn flow_inventory(
        &self,
        filter: &crate::durable::FlowFilter,
        after: Option<&str>,
        limit: std::num::NonZeroU32,
    ) -> StoreResult<crate::durable::FlowPage> {
        let filter = filter.clone();
        let after = after.map(str::to_owned);
        run_sqlite(&self.sqlite, move |store| {
            store.flow_inventory(&filter, after.as_deref(), limit)
        })
        .await
    }

    pub async fn flow_detail(
        &self,
        selector: &str,
    ) -> StoreResult<Option<crate::durable::FlowDetail>> {
        let selector = selector.to_string();
        run_sqlite(&self.sqlite, move |store| store.flow_detail(&selector)).await
    }

    pub async fn create_flow(&self, flow: FlowSession) -> StoreResult<FlowSession> {
        run_sqlite(&self.sqlite, move |store| store.create_flow(&flow)).await
    }

    pub async fn start_task_flow(
        &self,
        task_id: &TaskId,
        flow: FlowSession,
    ) -> StoreResult<FlowSession> {
        let task_id = task_id.clone();
        run_sqlite(&self.sqlite, move |store| {
            store.start_task_flow(&task_id, &flow)
        })
        .await
    }

    pub async fn flow(&self, id: &str) -> StoreResult<Option<FlowSession>> {
        let id = id.to_string();
        run_sqlite(&self.sqlite, move |store| store.flow(&id)).await
    }

    pub async fn task_flow(&self, task_id: &TaskId) -> StoreResult<Option<FlowSession>> {
        let task_id = task_id.clone();
        run_sqlite(&self.sqlite, move |store| store.task_flow(&task_id)).await
    }

    pub async fn recover_flow(
        &self,
        id: &str,
        claim: Option<&TaskWorkerClaim>,
    ) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let claim = claim.cloned();
        run_sqlite(&self.sqlite, move |store| {
            store.recover_flow(&id, claim.as_ref())
        })
        .await
    }

    pub async fn reserve_attempt(
        &self,
        id: &str,
        version: u64,
        claim: Option<&TaskWorkerClaim>,
        exec: Option<&crate::id::ExecId>,
    ) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let claim = claim.cloned();
        let exec = exec.cloned();
        run_sqlite(&self.sqlite, move |store| {
            store.reserve_attempt(&id, version, claim.as_ref(), exec.as_ref())
        })
        .await
    }

    pub async fn publish_attempt(
        &self,
        id: &str,
        version: u64,
        captured: i64,
        claim: Option<&TaskWorkerClaim>,
        provider: &str,
        model: Option<&str>,
    ) -> StoreResult<()> {
        let id = id.to_string();
        let claim = claim.cloned();
        let provider = provider.to_string();
        let model = model.map(str::to_owned);
        run_sqlite(&self.sqlite, move |store| {
            store.publish_attempt(
                &id,
                version,
                captured,
                claim.as_ref(),
                &provider,
                model.as_deref(),
            )
        })
        .await
    }

    pub async fn retry_flow(&self, id: &str, direction: Option<&str>) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let direction = direction.map(str::to_owned);
        run_sqlite(&self.sqlite, move |store| {
            store.retry_flow(&id, direction.as_deref())
        })
        .await
    }

    pub async fn release_flow(
        &self,
        id: &str,
        version: u64,
        claim: Option<&TaskWorkerClaim>,
    ) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let claim = claim.cloned();
        run_sqlite(&self.sqlite, move |store| {
            store.release_flow(&id, version, claim.as_ref())
        })
        .await
    }

    pub(crate) async fn reset_flow_input(
        &self,
        id: &str,
        version: u64,
        claim: Option<&TaskWorkerClaim>,
    ) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let claim = claim.cloned();
        run_sqlite(&self.sqlite, move |store| {
            store.reset_flow_input(&id, version, claim.as_ref())
        })
        .await
    }

    pub async fn checkpoint_flow(
        &self,
        id: &str,
        version: u64,
        cursor: &ExecutionCursor,
        claim: Option<&TaskWorkerClaim>,
        progress: Option<&str>,
    ) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let cursor = cursor.clone();
        let claim = claim.cloned();
        let progress = progress.map(str::to_owned);
        run_sqlite(&self.sqlite, move |store| {
            store.checkpoint_flow(&id, version, &cursor, claim.as_ref(), progress.as_deref())
        })
        .await
    }

    pub async fn fail_flow(
        &self,
        id: &str,
        version: u64,
        claim: Option<&TaskWorkerClaim>,
        failure: &TaskFlowBlocker,
    ) -> StoreResult<FlowSession> {
        let id = id.to_string();
        let claim = claim.cloned();
        let failure = failure.clone();
        run_sqlite(&self.sqlite, move |store| {
            store.fail_flow(&id, version, claim.as_ref(), &failure)
        })
        .await
    }

    pub async fn end_flow(
        &self,
        id: &str,
        version: u64,
        claim: Option<&TaskWorkerClaim>,
        summary: &str,
    ) -> StoreResult<()> {
        let id = id.to_string();
        let claim = claim.cloned();
        let summary = summary.to_string();
        run_sqlite(&self.sqlite, move |store| {
            store.end_flow(&id, version, claim.as_ref(), &summary)
        })
        .await
    }

    pub async fn reserve_task_review(&self, id: &str, version: u64) -> StoreResult<FlowSession> {
        let id = id.to_string();
        run_sqlite(&self.sqlite, move |store| {
            store.reserve_task_review(&id, version)
        })
        .await
    }

    pub async fn complete_task_review(
        &self,
        task_id: &TaskId,
        expected: &FlowSession,
        summary: &str,
    ) -> StoreResult<()> {
        let task_id = task_id.clone();
        let expected = expected.clone();
        let summary = summary.to_string();
        run_sqlite(&self.sqlite, move |store| {
            store.complete_task_review(&task_id, &expected, &summary)
        })
        .await
    }

    pub async fn claim_task_worker(
        &self,
        task_id: &TaskId,
        expected_invocation: &str,
        expected_version: u64,
        owner: &TaskWorkerOwner,
        claimed_at: OffsetDateTime,
    ) -> StoreResult<TaskWorkerClaimOutcome> {
        let task_id = task_id.clone();
        let expected_invocation = expected_invocation.to_string();
        let owner = owner.clone();
        run_sqlite(&self.sqlite, move |store| {
            store.claim_task_worker(
                &task_id,
                &expected_invocation,
                expected_version,
                &owner,
                claimed_at,
            )
        })
        .await
    }

    pub async fn reclaim_task_worker(
        &self,
        task_id: &TaskId,
        expected: &TaskWorkerClaim,
        owner: &TaskWorkerOwner,
        claimed_at: OffsetDateTime,
    ) -> StoreResult<TaskWorkerClaim> {
        let task_id = task_id.clone();
        let expected = expected.clone();
        let owner = owner.clone();
        run_sqlite(&self.sqlite, move |store| {
            store.reclaim_task_worker(&task_id, &expected, &owner, claimed_at)
        })
        .await
    }

    pub async fn waiting_review(&self, session_id: &str) -> StoreResult<FlowSession> {
        let session_id = session_id.to_string();
        run_sqlite(&self.sqlite, move |store| store.waiting_review(&session_id)).await
    }
}