Skip to main content

systemprompt_runtime/optimization/
iteration.rs

1//! Bounded development iterations reuse frozen execution settings and the
2//! campaign budget. Holdout cases never enter automatic improvement runs.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use systemprompt_evaluation::campaigns::diagnostics::{DiagnosticCode, DiagnosticStage};
8use systemprompt_evaluation::campaigns::suggestions::RetainedSuggestion;
9use systemprompt_evaluation::campaigns::{CampaignPolicy, OptimizationObjective};
10use systemprompt_evaluation::experiments::records::{ExperimentDetail, ExperimentStatus};
11use systemprompt_evaluation::experiments::resources::{Partition, ResourceContent};
12use systemprompt_evaluation::experiments::{ExperimentSpec, Objective};
13use systemprompt_evaluation::models::{CampaignStatus, SuggestionStatus};
14use systemprompt_evaluation::repository::experiments::{
15    CampaignExperiment, ManagedWorkspaceRegistration,
16};
17use systemprompt_identifiers::{EvalCampaignId, EvalExperimentId, ResourceRevisionId, UserId};
18
19use super::diagnostics::DiagnosticContext;
20use super::{OptimizationError, SkillOptimizationOrchestrator};
21
22impl SkillOptimizationOrchestrator {
23    pub(super) async fn launch_inner(
24        &self,
25        owner: &UserId,
26        actor: &UserId,
27        input: &CampaignExperiment,
28    ) -> Result<EvalExperimentId, OptimizationError> {
29        self.validate_campaign_spec(owner, &input.campaign_id, &input.spec)
30            .await?;
31        Ok(self
32            .evaluations
33            .experiments
34            .create_for_campaign(owner, actor, input)
35            .await?)
36    }
37
38
39    pub async fn attach(
40        &self,
41        owner: &UserId,
42        actor: &UserId,
43        campaign_id: &EvalCampaignId,
44        experiment_id: &EvalExperimentId,
45    ) -> Result<(), OptimizationError> {
46        let detail = self
47            .evaluations
48            .experiments
49            .get(owner, experiment_id)
50            .await?;
51        self.validate_campaign_spec(owner, campaign_id, &detail.experiment.spec)
52            .await?;
53        Ok(self
54            .evaluations
55            .campaigns
56            .attach_experiment(owner, actor, campaign_id, experiment_id)
57            .await?)
58    }
59
60    async fn validate_campaign_spec(
61        &self,
62        owner: &UserId,
63        campaign_id: &EvalCampaignId,
64        spec: &ExperimentSpec,
65    ) -> Result<(), OptimizationError> {
66        let campaign = self.evaluations.campaigns.get(owner, campaign_id).await?;
67        let baseline = self
68            .managed
69            .get_revision_bundle(owner, &campaign.policy.baseline_revision_id)
70            .await?;
71        if spec.variants.len() != 2
72            || spec.variants[0].skill_bundle_digest != baseline.digest()?.as_str()
73        {
74            return Err(OptimizationError::Source(
75                "Campaign runs require the retained baseline and one candidate".to_owned(),
76            ));
77        }
78        let workspace = self
79            .evaluations
80            .evidence
81            .get_managed_workspace(owner, &spec.variants[1].skill_bundle_digest)
82            .await?;
83        if self
84            .managed
85            .revision_resource(owner, &workspace.managed_revision_id)
86            .await?
87            != campaign.policy.resource_id
88        {
89            return Err(OptimizationError::Source(
90                "Candidate must belong to the campaign resource".to_owned(),
91            ));
92        }
93        Ok(())
94    }
95
96    pub(super) async fn advance_inner(
97        &self,
98        owner: &UserId,
99        actor: &UserId,
100        id: &EvalCampaignId,
101    ) -> Result<Option<EvalExperimentId>, OptimizationError> {
102        let campaign = self.evaluations.campaigns.get(owner, id).await?;
103        let experiments = self
104            .evaluations
105            .campaigns
106            .list_experiments(owner, id)
107            .await?;
108        if campaign.status != CampaignStatus::Active || !campaign.policy.automatic {
109            return Ok(None);
110        }
111        let ctx = DiagnosticContext {
112            owner,
113            actor,
114            campaign: id,
115            key: "automatic",
116            stage: DiagnosticStage::AutomaticFollowup,
117        };
118        if experiments.len() >= campaign.policy.maximum_iterations as usize {
119            self.blocked(&ctx, DiagnosticCode::IterationLimit).await?;
120            return Ok(None);
121        }
122        let Some(previous) = experiments.last() else {
123            self.blocked(&ctx, DiagnosticCode::MissingTemplate).await?;
124            return Ok(None);
125        };
126        let Some((detail, suggestion)) = self.select_followup(&ctx, previous).await? else {
127            return Ok(None);
128        };
129        let revision = self
130            .apply_suggestion(
131                owner,
132                &campaign.policy,
133                &detail.experiment.spec.variants[1],
134                &suggestion,
135            )
136            .await?;
137        let digest = self.register_workspace(owner, &revision).await?;
138        let spec = self
139            .development_followup_spec(owner, detail.experiment.spec, digest, &campaign.policy)
140            .await?;
141        let input = CampaignExperiment {
142            campaign_id: id.clone(),
143            idempotency_key: format!("suggestion:{}", suggestion.id),
144            spec,
145        };
146        Ok(Some(self.launch(owner, actor, &input).await?))
147    }
148
149    async fn select_followup(
150        &self,
151        ctx: &DiagnosticContext<'_>,
152        previous: &EvalExperimentId,
153    ) -> Result<Option<(ExperimentDetail, RetainedSuggestion)>, OptimizationError> {
154        let detail = self
155            .evaluations
156            .experiments
157            .get(ctx.owner, previous)
158            .await?;
159        if detail.experiment.status != ExperimentStatus::Completed
160            || detail.experiment.spec.variants.len() != 2
161            || detail.experiment.spec.claim_independent_improvement
162        {
163            return Ok(None);
164        }
165        let suggestions = self
166            .evaluations
167            .lifecycle
168            .list_suggestions(ctx.owner, previous)
169            .await?;
170        let Some(suggestion) = suggestions
171            .into_iter()
172            .find(|suggestion| suggestion.status == SuggestionStatus::Draft)
173        else {
174            self.blocked(ctx, DiagnosticCode::MissingSuggestion).await?;
175            return Ok(None);
176        };
177        if !self
178            .evaluations
179            .experiments
180            .execution_availability(&detail.experiment.spec)
181            .admitted
182        {
183            self.blocked(ctx, DiagnosticCode::UnsupportedCapability)
184                .await?;
185            return Ok(None);
186        }
187        Ok(Some((detail, suggestion)))
188    }
189
190    async fn development_followup_spec(
191        &self,
192        owner: &UserId,
193        mut spec: ExperimentSpec,
194        candidate_digest: String,
195        policy: &CampaignPolicy,
196    ) -> Result<ExperimentSpec, OptimizationError> {
197        let mut development = Vec::new();
198        for case in &spec.cases {
199            if let ResourceContent::Case(content) =
200                self.evaluations.revisions.get(owner, case).await?
201                && content.partition == Partition::Development
202            {
203                development.push(case.clone());
204            }
205        }
206        if development.is_empty() {
207            return Err(OptimizationError::Source(
208                "No development cases available for automatic iteration".to_owned(),
209            ));
210        }
211        spec.cases = development;
212        spec.claim_independent_improvement = false;
213        spec.variants[1].skill_bundle_digest = candidate_digest;
214        spec.objective = match policy.objective {
215            OptimizationObjective::Quality => Objective::Quality,
216            OptimizationObjective::Tokens => Objective::Tokens,
217            OptimizationObjective::Cost => Objective::Cost,
218            OptimizationObjective::Latency => Objective::Latency,
219        };
220        Ok(spec)
221    }
222
223
224    pub async fn register_workspace(
225        &self,
226        owner: &UserId,
227        revision: &ResourceRevisionId,
228    ) -> Result<String, OptimizationError> {
229        let bundle = self.managed.get_revision_bundle(owner, revision).await?;
230        let digest = bundle.digest()?;
231        match self
232            .evaluations
233            .evidence
234            .get_managed_workspace(owner, digest.as_str())
235            .await
236        {
237            Ok(_) => return Ok(digest.as_str().to_owned()),
238            Err(systemprompt_evaluation::EvaluationError::ResourceNotFound(_)) => {},
239            Err(error) => return Err(error.into()),
240        }
241        let files: Vec<_> = bundle
242            .revisions
243            .values()
244            .flat_map(|manifest| manifest.files.values())
245            .collect();
246        let bytes = files
247            .iter()
248            .try_fold(0usize, |sum, file| {
249                usize::try_from(file.bytes)
250                    .ok()
251                    .and_then(|bytes| sum.checked_add(bytes))
252            })
253            .ok_or_else(|| OptimizationError::Source("Bundle size overflow".to_owned()))?;
254        self.evaluations
255            .evidence
256            .register_managed_workspace(
257                owner,
258                &ManagedWorkspaceRegistration {
259                    managed_revision_id: revision,
260                    publication_generation: None,
261                    manifest: &bundle,
262                    expected_digest: digest.as_str(),
263                    file_count: files.len(),
264                    byte_count: bytes,
265                },
266            )
267            .await?;
268        Ok(digest.as_str().to_owned())
269    }
270}