loopflow 0.12.29

Run steps and flows with coding agents
Documentation
//! Task context and live input for the ordinary skill command.
use crate::child::ChildRef;
use crate::durable::{Steer, WorkRef};
use crate::harness::Harness;
use crate::planning::ProjectPlan;
use crate::store::SharedStore;
use crate::work::task::Task;
use anyhow::{anyhow, Result};
use std::io::BufRead;
use std::sync::Arc;
use std::time::{Duration, Instant};

/// Message and control cursors captured together before the skill starts.
pub(crate) struct TaskSeed {
    pub task: Task,
    pub message: String,
    pub steer: i64,
    pub interrupt: i64,
}

pub(crate) async fn prepare(store: &SharedStore, task: &Task, wave: &str) -> Result<TaskSeed> {
    let consumed = match crate::ops::flow_run::token()? {
        Some(token) => {
            let flow = store
                .flow(&token.invocation)
                .await?
                .ok_or_else(|| anyhow!("Flow {} is missing", token.invocation))?;
            anyhow::ensure!(
                flow.version == token.version && flow.task_id.as_ref() == Some(&task.id),
                "Task input belongs to a stale or different Flow"
            );
            store.sqlite.completed_step_steer_id(&flow)?
        }
        None => 0,
    };
    let steers: Vec<_> = store
        .task_steers(&task.id)
        .await?
        .into_iter()
        .filter(|steer| steer.id > consumed)
        .collect();
    let interrupt = store
        .latest_interrupt_id(&WorkRef::Task(task.id.clone()))
        .await?;
    let pr = store
        .task_prs(&task.id)
        .await?
        .pop()
        .ok_or_else(|| anyhow!("Task {} has no recorded PR", task.id))?;
    let project = store
        .get_project(&task.project_id)
        .await?
        .ok_or_else(|| anyhow!("Task Project is missing"))?;
    let message = format!(
        "{}\n\n{}",
        task_seed(task, &project.plan, &pr, wave, &steers),
        crate::ops::task::task_workspace_context(task, &pr)?
    );
    Ok(TaskSeed {
        task: task.clone(),
        message,
        steer: steers.last().map_or(consumed, |steer| steer.id),
        interrupt,
    })
}

/// A command's live Task input survives its provider's transient retries.
#[derive(Clone)]
pub struct TaskInput(Arc<tokio::sync::Mutex<Controls>>);

impl std::fmt::Debug for TaskInput {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.write_str("TaskInput")
    }
}

struct Controls {
    store: SharedStore,
    task: Task,
    steer: i64,
    interrupt: i64,
    receiver: Option<tokio::sync::mpsc::UnboundedReceiver<String>>,
    next_steer: Instant,
}

impl TaskInput {
    pub(crate) fn new(store: SharedStore, seed: TaskSeed) -> Self {
        Self(Arc::new(tokio::sync::Mutex::new(Controls {
            store,
            task: seed.task,
            steer: seed.steer,
            interrupt: seed.interrupt,
            receiver: None,
            next_steer: Instant::now(),
        })))
    }

    pub(crate) async fn record_seed(&self, capture: &crate::session_record::CaptureHandle) {
        capture.record_input("steer_seed_through", &self.0.lock().await.steer.to_string());
    }

    pub(crate) fn refresh(&self) -> CommentRefresh {
        let input = self.clone();
        CommentRefresh(tokio::spawn(async move {
            let (store, task) = {
                let controls = input.0.lock().await;
                (controls.store.clone(), controls.task.clone())
            };
            loop {
                tokio::time::sleep(Duration::from_secs(15)).await;
                if let Err(error) =
                    crate::ops::linear_observe::refresh_task_comments(&store, &task).await
                {
                    tracing::warn!(%error, "Linear comment refresh failed; retaining confirmed Task direction");
                }
            }
        }))
    }

    pub(crate) async fn poll(
        &self,
        harness: &mut dyn Harness,
        capture: Option<&crate::session_record::CaptureHandle>,
    ) -> Result<()> {
        let mut controls = self.0.lock().await;
        let Controls {
            store,
            task,
            steer,
            interrupt,
            receiver,
            next_steer,
        } = &mut *controls;
        let receiver = receiver.get_or_insert_with(|| {
            let (sender, receiver) = tokio::sync::mpsc::unbounded_channel();
            std::thread::spawn(move || {
                for line in std::io::stdin().lock().lines() {
                    let Ok(line) = line else { break };
                    if sender.send(line).is_err() {
                        break;
                    }
                }
            });
            println!(
                "task {}> attached; /status, /interrupt, /detach, or type a message/instruction",
                task.plan.identifier
            );
            receiver
        });
        while let Ok(line) = receiver.try_recv() {
            handle_attachment(store, task, harness, line).await?;
        }
        if Instant::now() >= *next_steer {
            for delivered in
                crate::ops::child::inject_live_steers(store, &task.id, harness, steer).await
            {
                if let Some(capture) = capture {
                    capture.record_input(
                        &format!("steer_transport_accepted:{}", delivered.id),
                        &delivered.text,
                    );
                }
            }
            crate::ops::child::observe_interrupt(
                store,
                &WorkRef::Task(task.id.clone()),
                harness,
                interrupt,
            )
            .await;
            *next_steer = Instant::now() + Duration::from_secs(5);
        }
        Ok(())
    }
}

pub(crate) struct CommentRefresh(tokio::task::JoinHandle<()>);
impl Drop for CommentRefresh {
    fn drop(&mut self) {
        self.0.abort();
    }
}

async fn handle_attachment(
    store: &SharedStore,
    task: &Task,
    harness: &mut dyn Harness,
    line: String,
) -> Result<()> {
    let line = line.trim();
    if line.is_empty() {
        return Ok(());
    }
    if line == "/status" {
        let work = store
            .work_for_child(&ChildRef::Task(task.id.clone()))
            .await?;
        println!(
            "{}  {:?}",
            task.plan.identifier,
            store.work_status(&work).await?
        );
        return Ok(());
    }
    if line == "/detach" {
        let _ = std::process::Command::new("tmux")
            .args(["detach-client"])
            .status();
        return Ok(());
    }
    if line == "/interrupt" {
        harness.interrupt().await?;
        println!("interrupted active provider turn");
    } else {
        let comment_id = crate::ops::linear_observe::publish_task_steer(store, task, line).await?;
        println!("posted to Linear {comment_id}");
    }
    Ok(())
}

pub(crate) fn task_seed(
    task: &Task,
    project: &ProjectPlan,
    pr: &crate::work::task::TaskPr,
    wave_name: &str,
    steers: &[Steer],
) -> String {
    let context = crate::ops::render_task_context(task, project, pr, wave_name, steers);
    format!(
        "{context}\n\nTyped PR and Task operations own publication, landing, rotation, and completion. `lf pr abandon` discards only this PR. If this PR already merged out of band and follow-up work remains, `lf pr next [slug]` rotates to the next serial PR, carrying committed and uncommitted follow-up forward."
    )
}