systemprompt_runtime/optimization/
iteration.rs1use 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}