daat-locus 0.2.0

A long-running local agent runtime with memory, workflows, apps, and sleep-time self-improvement.
use super::*;

pub(super) struct WorkflowMergePlanningInput<'a> {
    pub target_workflow: &'a PrimitiveSpec,
    pub target_reflection: &'a EvaluationArtifactWorkflowReflection,
    pub target_evidence: &'a [PrimitiveRunRecord],
    pub source_workflow: &'a PrimitiveSpec,
    pub source_reflection: &'a EvaluationArtifactWorkflowReflection,
    pub source_evidence: &'a [PrimitiveRunRecord],
}

pub(super) struct WorkflowFrontierReplayInput<'a> {
    pub entry: &'a WorkflowFrontierEntry,
    pub target_workflow: &'a PrimitiveSpec,
    pub target_reflection: Option<&'a EvaluationArtifactWorkflowReflection>,
    pub target_evidence: &'a [PrimitiveRunRecord],
    pub source_workflow: Option<&'a PrimitiveSpec>,
    pub source_reflection: Option<&'a EvaluationArtifactWorkflowReflection>,
    pub source_evidence: &'a [PrimitiveRunRecord],
}

#[async_trait]
pub(super) trait SleepPlannerRuntime: Send + Sync {
    async fn plan_runtime_error_correction(
        &self,
        context: &mut Context,
        runtime_error_cases: &[RuntimeErrorCase],
    ) -> Result<PromptPlanningResult>;

    async fn plan_workflow_improvement(
        &self,
        context: &mut Context,
        workflow: &PrimitiveSpec,
        evidence: &[PrimitiveRunRecord],
    ) -> Result<Option<WorkflowPlanningResult>>;

    async fn plan_workflow_merge(
        &self,
        context: &mut Context,
        input: WorkflowMergePlanningInput<'_>,
    ) -> Result<WorkflowMergePlanningResult>;

    async fn replay_workflow_frontier_entry(
        &self,
        context: &mut Context,
        input: WorkflowFrontierReplayInput<'_>,
    ) -> Result<EvaluationArtifactWorkflowCandidateEvaluation>;
}

pub(super) struct LlmSleepPlannerRuntime;

pub(super) async fn load_sleep_inputs() -> Result<SleepInputs> {
    let runtime_error_cases =
        crate::reasoning::runtime_error::load_runtime_error_case_batch().await?;
    Ok(SleepInputs {
        runtime_error_cases,
    })
}

#[async_trait]
impl SleepPlannerRuntime for LlmSleepPlannerRuntime {
    async fn plan_runtime_error_correction(
        &self,
        context: &mut Context,
        runtime_error_cases: &[RuntimeErrorCase],
    ) -> Result<PromptPlanningResult> {
        if runtime_error_cases.is_empty() {
            return Ok(PromptPlanningResult {
                reflections: Vec::new(),
                candidates: Vec::new(),
                evaluations: Vec::new(),
            });
        }

        let renderer = OpenAIToolRenderer;
        let program = RuntimeErrorCorrectionPlannerProgram;
        let tuning = resolve_program_tuning(context, &program).await;
        let current_additions = context
            .compiled_prompts
            .runtime_system_additions()
            .join("\n");
        let runtime_error_cases_json =
            serde_json::to_string_pretty(runtime_error_cases).into_diagnostic()?;
        let outcome = execute_program_with_ir_report(
            context.judge_llm.as_ref(),
            context,
            &renderer,
            &program,
            program.dataset_ir(current_additions, runtime_error_cases_json),
            &tuning,
            TraceOrigin::Sleep,
        )
        .await?;

        Ok(runtime_error_correction_planning_result_from_output(
            &outcome.output,
            runtime_error_cases,
        ))
    }

    async fn plan_workflow_improvement(
        &self,
        context: &mut Context,
        workflow: &PrimitiveSpec,
        evidence: &[PrimitiveRunRecord],
    ) -> Result<Option<WorkflowPlanningResult>> {
        let renderer = OpenAIToolRenderer;
        let program = WorkflowEvolutionPlannerProgram;
        let tuning = resolve_program_tuning(context, &program).await;
        let workflow_markdown = render_workflow_spec_markdown(workflow);
        let workflow_run_evidence_json = render_workflow_run_evidence_json(evidence)?;
        let outcome = execute_program_with_ir_report(
            context.judge_llm.as_ref(),
            context,
            &renderer,
            &program,
            program.dataset_ir(
                workflow.id.clone(),
                workflow_markdown,
                workflow_run_evidence_json,
            ),
            &tuning,
            TraceOrigin::Sleep,
        )
        .await?;

        Ok(workflow_planning_result_from_output(
            workflow,
            evidence,
            &outcome.output,
        ))
    }

    async fn plan_workflow_merge(
        &self,
        context: &mut Context,
        input: WorkflowMergePlanningInput<'_>,
    ) -> Result<WorkflowMergePlanningResult> {
        let WorkflowMergePlanningInput {
            target_workflow,
            target_reflection,
            target_evidence,
            source_workflow,
            source_reflection,
            source_evidence,
        } = input;

        let renderer = OpenAIToolRenderer;
        let program = WorkflowMergePlannerProgram;
        let tuning = resolve_program_tuning(context, &program).await;
        let outcome = execute_program_with_ir_report(
            context.judge_llm.as_ref(),
            context,
            &renderer,
            &program,
            program.dataset_ir(
                target_workflow.id.clone(),
                render_workflow_spec_markdown(target_workflow),
                serde_json::to_string_pretty(target_reflection).into_diagnostic()?,
                render_workflow_run_evidence_json(target_evidence)?,
                source_workflow.id.clone(),
                render_workflow_spec_markdown(source_workflow),
                serde_json::to_string_pretty(source_reflection).into_diagnostic()?,
                render_workflow_run_evidence_json(source_evidence)?,
            ),
            &tuning,
            TraceOrigin::Sleep,
        )
        .await?;

        Ok(workflow_merge_planning_result_from_output(
            target_workflow,
            source_workflow,
            target_reflection,
            source_reflection,
            target_evidence,
            source_evidence,
            &outcome.output,
        ))
    }

    async fn replay_workflow_frontier_entry(
        &self,
        context: &mut Context,
        input: WorkflowFrontierReplayInput<'_>,
    ) -> Result<EvaluationArtifactWorkflowCandidateEvaluation> {
        let WorkflowFrontierReplayInput {
            entry,
            target_workflow,
            target_reflection,
            target_evidence,
            source_workflow,
            source_reflection,
            source_evidence,
        } = input;

        let renderer = OpenAIToolRenderer;
        let program = WorkflowCandidateRolloutEvaluatorProgram;
        let tuning = resolve_program_tuning(context, &program).await;
        let candidate_json = serde_json::to_string_pretty(entry).into_diagnostic()?;
        let rollout =
            execute_workflow_candidate_rollout(context, entry, target_workflow, source_workflow)
                .await?;
        let target_task_cases = select_workflow_task_cases(
            &target_evidence
                .iter()
                .map(workflow_task_case_from_record)
                .collect::<Vec<_>>(),
            8,
        );
        let source_task_cases = select_workflow_task_cases(
            &source_evidence
                .iter()
                .map(workflow_task_case_from_record)
                .collect::<Vec<_>>(),
            8,
        );
        let case_count = target_task_cases.len().max(source_task_cases.len()).max(1);
        let target_workflow_spec = render_workflow_spec_markdown(&rollout.target_workflow);
        let target_reflection_json =
            serde_json::to_string_pretty(&target_reflection.cloned()).into_diagnostic()?;
        let source_workflow_spec = source_workflow
            .map(render_workflow_spec_markdown)
            .unwrap_or_else(|| "none".to_string());
        let source_reflection_json =
            serde_json::to_string_pretty(&source_reflection.cloned()).into_diagnostic()?;
        let mut outputs = Vec::<WorkflowCandidateRolloutEvaluatorOutput>::new();
        for index in 0..case_count {
            let target_task = target_task_cases
                .get(index)
                .cloned()
                .or_else(|| target_task_cases.last().cloned())
                .unwrap_or_else(blank_workflow_task_case);
            let rolled_out_target_case =
                run_workflow_task_rollout(&rollout.target_workflow, &target_task);
            let source_task = source_task_cases
                .get(index)
                .cloned()
                .or_else(|| source_task_cases.last().cloned())
                .unwrap_or_else(blank_workflow_task_case);
            let rolled_out_source_case =
                source_workflow.map(|workflow| run_workflow_task_rollout(workflow, &source_task));
            let outcome = execute_program_with_ir_report(
                context.judge_llm.as_ref(),
                context,
                &renderer,
                &program,
                program.dataset_ir(
                    entry.candidate_kind.clone(),
                    target_workflow_spec.clone(),
                    format!("{} | {}", rollout.summary, rolled_out_target_case.summary),
                    target_reflection_json.clone(),
                    render_workflow_rollout_case_json(&rolled_out_target_case)?,
                    source_workflow_spec.clone(),
                    source_reflection_json.clone(),
                    if let Some(source_case) = rolled_out_source_case.as_ref() {
                        render_workflow_rollout_case_json(source_case)?
                    } else {
                        "none".to_string()
                    },
                    candidate_json.clone(),
                ),
                &tuning,
                TraceOrigin::Sleep,
            )
            .await?;
            outputs.push(outcome.output);
        }
        Ok(aggregate_workflow_replay_evaluation(
            EvaluationArtifactWorkflowCandidateEvaluation {
                workflow_id: entry.evaluation.workflow_id.clone(),
                candidate_kind: entry.candidate_kind.clone(),
                candidate_title: entry.evaluation.candidate_title.clone(),
                rationale: String::new(),
                score: 0.0,
                accepted: false,
                selected: false,
                source_run_ids: target_evidence
                    .iter()
                    .chain(source_evidence.iter())
                    .map(|record| record.run_id.clone())
                    .collect(),
            },
            &outputs,
        ))
    }
}