Skip to main content

ironflow_worker/
api_store.rs

1//! HTTP-based [`RunStore`] that talks to the ironflow API internal routes.
2
3use std::future::Future;
4use std::pin::Pin;
5
6use std::time::Duration;
7
8use chrono::{DateTime, Utc};
9use reqwest::{Client, StatusCode};
10use serde_json::{Value, from_str, json};
11use uuid::Uuid;
12
13use ironflow_engine::error::CONCURRENCY_CONFLICT_CODE;
14use ironflow_store::api_key_store::ApiKeyStore;
15use ironflow_store::approval_delegation_store::ApprovalDelegationStore;
16use ironflow_store::artifact_store::ArtifactStore;
17use ironflow_store::audit_log_store::AuditLogStore;
18use ironflow_store::entities::{
19    ApiKey, ApiKeyUpdate, ApprovalDelegation, Artifact, ArtifactLookup, AuditLogEntry,
20    AuditLogFilter, ConcurrencyGroupBacklog, DelegationFilter, KeyVersionStatus, LeaseRequest,
21    LogEntry, LogFilter, NewApiKey, NewApprovalDelegation, NewArtifact, NewAuditLogEntry,
22    NewLogEntries, NewRun, NewSchedule, NewStep, NewStepDependency, NewUser, Page, PurgePolicy,
23    PurgeableRun, ReapedRun, RotationBatch, RotationRequest, Run, RunCreation, RunFilter, RunStats,
24    RunStatus, RunUpdate, Schedule, ScheduleFiring, ScheduleNext, ScheduleUpdate, Secret,
25    SecretMetadata, StatsHistoryBucket, StatsHistoryFilter, Step, StepApproval, StepDependency,
26    StepUpdate, User, WorkerCapabilities,
27};
28use ironflow_store::entities::{
29    NewProviderAccount, NewProviderAccountObservation, ProviderAccount, ProviderAccountCandidate,
30    ProviderAccountUpdate, ProviderAccountUsagePoint, ProviderAccountWindow,
31};
32use ironflow_store::entities::{
33    NewRefreshToken, NewSignal, Signal, SignalFilter, SignalInsert, SignalStepResolution,
34};
35use ironflow_store::error::StoreError;
36use ironflow_store::log_store::LogStore;
37use ironflow_store::provider_account_store::ProviderAccountStore;
38use ironflow_store::schedule_store::ScheduleStore;
39use ironflow_store::secret_store::SecretStore;
40use ironflow_store::signal_store::SignalStore;
41use ironflow_store::store::RunStore;
42use ironflow_store::user_store::UserStore;
43
44type StoreFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, StoreError>> + Send + 'a>>;
45
46/// API response envelope.
47#[derive(serde::Deserialize, serde::Serialize)]
48struct ApiResponse<T> {
49    data: T,
50}
51
52/// RunStore implementation that communicates with the API server via HTTP.
53#[derive(Debug, Clone)]
54pub struct ApiRunStore {
55    client: Client,
56    base_url: String,
57    token: String,
58}
59
60impl ApiRunStore {
61    pub fn new(base_url: &str, token: &str) -> Self {
62        let client = Client::builder()
63            .timeout(Duration::from_secs(30))
64            .connect_timeout(Duration::from_secs(5))
65            .build()
66            .expect("failed to build HTTP client");
67
68        Self {
69            client,
70            base_url: base_url.trim_end_matches('/').to_string(),
71            token: token.to_string(),
72        }
73    }
74
75    fn internal(&self, path: &str) -> String {
76        format!("{}/api/v1/internal{}", self.base_url, path)
77    }
78
79    fn err(e: reqwest::Error) -> StoreError {
80        StoreError::Database(format!("worker HTTP error: {e}"))
81    }
82
83    fn status_err(body: &str) -> StoreError {
84        StoreError::Database(format!("worker API error: {body}"))
85    }
86}
87
88impl RunStore for ApiRunStore {
89    fn create_run(&self, req: NewRun) -> StoreFuture<'_, RunCreation> {
90        Box::pin(async move {
91            let resp = self
92                .client
93                .post(self.internal("/runs"))
94                .bearer_auth(&self.token)
95                .json(&req)
96                .send()
97                .await
98                .map_err(Self::err)?;
99
100            // 409 means another active run holds the concurrency key: the
101            // engine turns it into a recorded conflict, never a failure.
102            if resp.status() == StatusCode::CONFLICT {
103                let body = resp.text().await.unwrap_or_default();
104                return Err(match concurrency_conflict(&body) {
105                    Some(conflict) => conflict,
106                    None => Self::status_err(&body),
107                });
108            }
109            if !resp.status().is_success() {
110                let body = resp.text().await.unwrap_or_default();
111                return Err(Self::status_err(&body));
112            }
113
114            let api_resp: ApiResponse<Run> = resp.json().await.map_err(Self::err)?;
115            // The internal endpoint is not idempotent: a success is always a creation.
116            Ok(RunCreation::Created(api_resp.data))
117        })
118    }
119
120    fn find_run_by_idempotency_key(&self, _key: &str) -> StoreFuture<'_, Option<Run>> {
121        Box::pin(async move {
122            Err(StoreError::Database(
123                "find_run_by_idempotency_key not supported via worker API".to_string(),
124            ))
125        })
126    }
127
128    fn get_run(&self, id: Uuid) -> StoreFuture<'_, Option<Run>> {
129        Box::pin(async move {
130            let resp = self
131                .client
132                .get(self.internal(&format!("/runs/{id}")))
133                .bearer_auth(&self.token)
134                .send()
135                .await
136                .map_err(Self::err)?;
137
138            if resp.status() == reqwest::StatusCode::NOT_FOUND {
139                return Ok(None);
140            }
141            if !resp.status().is_success() {
142                let body = resp.text().await.unwrap_or_default();
143                return Err(Self::status_err(&body));
144            }
145
146            #[derive(serde::Deserialize)]
147            struct RunDetail {
148                run: Run,
149            }
150
151            let api_resp: ApiResponse<RunDetail> = resp.json().await.map_err(Self::err)?;
152            Ok(Some(api_resp.data.run))
153        })
154    }
155
156    fn list_runs(
157        &self,
158        _filter: RunFilter,
159        _page: u32,
160        _per_page: u32,
161    ) -> StoreFuture<'_, Page<Run>> {
162        Box::pin(async move {
163            Err(StoreError::Database(
164                "list_runs not supported via worker API".to_string(),
165            ))
166        })
167    }
168
169    fn update_run_status(&self, id: Uuid, new_status: RunStatus) -> StoreFuture<'_, ()> {
170        Box::pin(async move {
171            let resp = self
172                .client
173                .put(self.internal(&format!("/runs/{id}/status")))
174                .bearer_auth(&self.token)
175                .json(&serde_json::json!({ "status": new_status }))
176                .send()
177                .await
178                .map_err(Self::err)?;
179
180            if !resp.status().is_success() {
181                let body = resp.text().await.unwrap_or_default();
182                return Err(Self::status_err(&body));
183            }
184            Ok(())
185        })
186    }
187
188    fn update_run(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, ()> {
189        Box::pin(async move {
190            let resp = self
191                .client
192                .put(self.internal(&format!("/runs/{id}")))
193                .bearer_auth(&self.token)
194                .json(&update)
195                .send()
196                .await
197                .map_err(Self::err)?;
198
199            if !resp.status().is_success() {
200                let body = resp.text().await.unwrap_or_default();
201                return Err(Self::status_err(&body));
202            }
203            Ok(())
204        })
205    }
206
207    fn list_active_descendants(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Run>> {
208        Box::pin(async move {
209            let resp = self
210                .client
211                .get(self.internal(&format!("/runs/{run_id}/descendants")))
212                .bearer_auth(&self.token)
213                .send()
214                .await
215                .map_err(Self::err)?;
216
217            if !resp.status().is_success() {
218                let body = resp.text().await.unwrap_or_default();
219                return Err(Self::status_err(&body));
220            }
221
222            let api_resp: ApiResponse<Vec<Run>> = resp.json().await.map_err(Self::err)?;
223            Ok(api_resp.data)
224        })
225    }
226
227    fn pick_next_pending_for(
228        &self,
229        lease: Option<LeaseRequest>,
230        capabilities: Option<WorkerCapabilities>,
231    ) -> StoreFuture<'_, Option<Run>> {
232        Box::pin(async move {
233            let mut request = self
234                .client
235                .get(self.internal("/runs/next"))
236                .bearer_auth(&self.token);
237
238            if let Some(lease) = lease {
239                request = request.query(&[
240                    ("worker_id", lease.worker_id),
241                    ("lease_ttl_secs", lease.ttl.as_secs().to_string()),
242                ]);
243            }
244            if let Some(capabilities) = capabilities {
245                request = request.query(&capability_query(&capabilities));
246            }
247
248            let resp = request.send().await.map_err(Self::err)?;
249
250            if !resp.status().is_success() {
251                let body = resp.text().await.unwrap_or_default();
252                return Err(Self::status_err(&body));
253            }
254
255            let api_resp: ApiResponse<Option<Run>> = resp.json().await.map_err(Self::err)?;
256            Ok(api_resp.data)
257        })
258    }
259
260    fn renew_lease(&self, id: Uuid, lease: LeaseRequest) -> StoreFuture<'_, DateTime<Utc>> {
261        Box::pin(async move {
262            #[derive(serde::Serialize)]
263            struct RenewLeaseBody {
264                worker_id: String,
265                lease_ttl_secs: u64,
266            }
267
268            #[derive(serde::Deserialize)]
269            struct RenewLeaseData {
270                lease_expires_at: DateTime<Utc>,
271            }
272
273            let resp = self
274                .client
275                .post(self.internal(&format!("/runs/{id}/lease")))
276                .bearer_auth(&self.token)
277                .json(&RenewLeaseBody {
278                    worker_id: lease.worker_id.clone(),
279                    lease_ttl_secs: lease.ttl.as_secs(),
280                })
281                .send()
282                .await
283                .map_err(Self::err)?;
284
285            // 409 means another worker owns the run now, or the run left the
286            // Running state: the caller must abandon it, not retry.
287            if resp.status() == StatusCode::CONFLICT {
288                return Err(StoreError::LeaseLost {
289                    run_id: id,
290                    held_by: None,
291                });
292            }
293            if resp.status() == StatusCode::NOT_FOUND {
294                return Err(StoreError::RunNotFound(id));
295            }
296            if !resp.status().is_success() {
297                let body = resp.text().await.unwrap_or_default();
298                return Err(Self::status_err(&body));
299            }
300
301            let api_resp: ApiResponse<RenewLeaseData> = resp.json().await.map_err(Self::err)?;
302            Ok(api_resp.data.lease_expires_at)
303        })
304    }
305
306    fn count_blocked_runs_by_group(&self) -> StoreFuture<'_, Vec<ConcurrencyGroupBacklog>> {
307        Box::pin(async move {
308            Err(StoreError::Database(
309                "count_blocked_runs_by_group not supported via worker API".to_string(),
310            ))
311        })
312    }
313
314    fn reap_expired_leases(&self, _limit: u32) -> StoreFuture<'_, Vec<ReapedRun>> {
315        // Recovery is an API-server responsibility: the worker has no route for
316        // it and must never requeue runs it does not own.
317        Box::pin(async move { Ok(Vec::new()) })
318    }
319
320    fn claim_due_approval_deadlines(&self, _limit: u32) -> StoreFuture<'_, Vec<Step>> {
321        // Escalation is an API-server responsibility: the worker has no route for
322        // it and must never resolve a gate it does not own.
323        Box::pin(async move { Ok(Vec::new()) })
324    }
325
326    fn claim_due_sleeping_runs(&self, _limit: u32) -> StoreFuture<'_, Vec<Run>> {
327        // Waking sleeping runs is an API-server responsibility: the worker has
328        // no route for it and picks the requeued runs up like any pending run.
329        Box::pin(async move { Ok(Vec::new()) })
330    }
331
332    fn list_purgeable_runs(
333        &self,
334        _policy: &PurgePolicy,
335        _batch_size: u32,
336    ) -> StoreFuture<'_, Vec<PurgeableRun>> {
337        // Purging is an API-server responsibility.
338        Box::pin(async move { Ok(Vec::new()) })
339    }
340
341    fn delete_run(&self, _id: Uuid) -> StoreFuture<'_, Vec<String>> {
342        // Purging is an API-server responsibility.
343        Box::pin(async move { Ok(Vec::new()) })
344    }
345
346    fn create_step(&self, req: NewStep) -> StoreFuture<'_, Step> {
347        Box::pin(async move {
348            let resp = self
349                .client
350                .post(self.internal("/steps"))
351                .bearer_auth(&self.token)
352                .json(&req)
353                .send()
354                .await
355                .map_err(Self::err)?;
356
357            if !resp.status().is_success() {
358                let body = resp.text().await.unwrap_or_default();
359                return Err(Self::status_err(&body));
360            }
361
362            let api_resp: ApiResponse<Step> = resp.json().await.map_err(Self::err)?;
363            Ok(api_resp.data)
364        })
365    }
366
367    fn update_step(&self, id: Uuid, update: StepUpdate) -> StoreFuture<'_, ()> {
368        Box::pin(async move {
369            let resp = self
370                .client
371                .put(self.internal(&format!("/steps/{id}")))
372                .bearer_auth(&self.token)
373                .json(&update)
374                .send()
375                .await
376                .map_err(Self::err)?;
377
378            if !resp.status().is_success() {
379                let body = resp.text().await.unwrap_or_default();
380                return Err(Self::status_err(&body));
381            }
382            Ok(())
383        })
384    }
385
386    fn get_step(&self, _id: Uuid) -> StoreFuture<'_, Option<Step>> {
387        // The worker never reads a step back through its store — step lookup
388        // lives on the API side. Return None so this trait method stays total
389        // without adding a dedicated HTTP route the worker doesn't use.
390        Box::pin(async move { Ok(None) })
391    }
392
393    fn record_step_approval(
394        &self,
395        _step_id: Uuid,
396        _approval: StepApproval,
397    ) -> StoreFuture<'_, Step> {
398        Box::pin(async move {
399            Err(StoreError::Database(
400                "record_step_approval not available in worker".to_string(),
401            ))
402        })
403    }
404
405    fn list_steps(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Step>> {
406        Box::pin(async move {
407            let resp = self
408                .client
409                .get(self.internal(&format!("/runs/{run_id}")))
410                .bearer_auth(&self.token)
411                .send()
412                .await
413                .map_err(Self::err)?;
414
415            if !resp.status().is_success() {
416                let body = resp.text().await.unwrap_or_default();
417                return Err(Self::status_err(&body));
418            }
419
420            #[derive(serde::Deserialize)]
421            struct RunDetail {
422                steps: Vec<Step>,
423            }
424
425            let api_resp: ApiResponse<RunDetail> = resp.json().await.map_err(Self::err)?;
426            Ok(api_resp.data.steps)
427        })
428    }
429
430    fn get_stats(&self, _filter: RunFilter) -> StoreFuture<'_, RunStats> {
431        Box::pin(async move {
432            Err(StoreError::Database(
433                "get_stats not supported via worker API".to_string(),
434            ))
435        })
436    }
437
438    fn get_stats_history(
439        &self,
440        _filter: StatsHistoryFilter,
441    ) -> StoreFuture<'_, Vec<StatsHistoryBucket>> {
442        Box::pin(async move {
443            Err(StoreError::Database(
444                "get_stats_history not supported via worker API".to_string(),
445            ))
446        })
447    }
448
449    fn create_step_dependencies(&self, deps: Vec<NewStepDependency>) -> StoreFuture<'_, ()> {
450        Box::pin(async move {
451            if deps.is_empty() {
452                return Ok(());
453            }
454
455            let resp = self
456                .client
457                .post(self.internal("/step-dependencies"))
458                .bearer_auth(&self.token)
459                .json(&deps)
460                .send()
461                .await
462                .map_err(Self::err)?;
463
464            if !resp.status().is_success() {
465                let body = resp.text().await.unwrap_or_default();
466                return Err(Self::status_err(&body));
467            }
468
469            Ok(())
470        })
471    }
472
473    fn list_step_dependencies(&self, _run_id: Uuid) -> StoreFuture<'_, Vec<StepDependency>> {
474        Box::pin(async move {
475            Err(StoreError::Database(
476                "list_step_dependencies not supported via worker API".to_string(),
477            ))
478        })
479    }
480}
481
482impl UserStore for ApiRunStore {
483    fn create_user(&self, _req: NewUser) -> StoreFuture<'_, User> {
484        Box::pin(async move {
485            Err(StoreError::Database(
486                "UserStore not available in worker".to_string(),
487            ))
488        })
489    }
490
491    fn find_user_by_email(&self, _email: &str) -> StoreFuture<'_, Option<User>> {
492        Box::pin(async move { Ok(None) })
493    }
494
495    fn find_user_by_username(&self, _username: &str) -> StoreFuture<'_, Option<User>> {
496        Box::pin(async move { Ok(None) })
497    }
498
499    fn find_user_by_id(&self, _id: Uuid) -> StoreFuture<'_, Option<User>> {
500        Box::pin(async move { Ok(None) })
501    }
502
503    fn count_users(&self) -> StoreFuture<'_, u64> {
504        Box::pin(async move {
505            Err(StoreError::Database(
506                "UserStore not available in worker".to_string(),
507            ))
508        })
509    }
510
511    fn list_users(&self, _page: u32, _per_page: u32) -> StoreFuture<'_, Page<User>> {
512        Box::pin(async move {
513            Err(StoreError::Database(
514                "UserStore not available in worker".to_string(),
515            ))
516        })
517    }
518
519    fn delete_user(&self, _id: Uuid) -> StoreFuture<'_, ()> {
520        Box::pin(async move {
521            Err(StoreError::Database(
522                "UserStore not available in worker".to_string(),
523            ))
524        })
525    }
526
527    fn update_user_role(&self, _id: Uuid, _is_admin: bool) -> StoreFuture<'_, User> {
528        Box::pin(async move {
529            Err(StoreError::Database(
530                "UserStore not available in worker".to_string(),
531            ))
532        })
533    }
534
535    fn update_user_password(&self, _id: Uuid, _password_hash: String) -> StoreFuture<'_, ()> {
536        Box::pin(async move {
537            Err(StoreError::Database(
538                "UserStore not available in worker".to_string(),
539            ))
540        })
541    }
542
543    fn list_user_groups(&self, _user_id: Uuid) -> StoreFuture<'_, Vec<String>> {
544        Box::pin(async move {
545            Err(StoreError::Database(
546                "UserStore not available in worker".to_string(),
547            ))
548        })
549    }
550
551    fn set_user_groups(
552        &self,
553        _user_id: Uuid,
554        _groups: Vec<String>,
555    ) -> StoreFuture<'_, Vec<String>> {
556        Box::pin(async move {
557            Err(StoreError::Database(
558                "UserStore not available in worker".to_string(),
559            ))
560        })
561    }
562
563    fn revoke_user_sessions(&self, _id: Uuid) -> StoreFuture<'_, i64> {
564        Box::pin(async move {
565            Err(StoreError::Database(
566                "UserStore not available in worker".to_string(),
567            ))
568        })
569    }
570
571    fn store_refresh_token(&self, _token: NewRefreshToken) -> StoreFuture<'_, ()> {
572        Box::pin(async move {
573            Err(StoreError::Database(
574                "UserStore not available in worker".to_string(),
575            ))
576        })
577    }
578
579    fn consume_refresh_token(&self, _token_hash: &str) -> StoreFuture<'_, Option<Uuid>> {
580        Box::pin(async move {
581            Err(StoreError::Database(
582                "UserStore not available in worker".to_string(),
583            ))
584        })
585    }
586}
587
588impl ApiKeyStore for ApiRunStore {
589    fn create_api_key(&self, _req: NewApiKey) -> StoreFuture<'_, ApiKey> {
590        Box::pin(async move {
591            Err(StoreError::Database(
592                "ApiKeyStore not available in worker".to_string(),
593            ))
594        })
595    }
596
597    fn find_api_key_by_prefix(&self, _prefix: &str) -> StoreFuture<'_, Option<ApiKey>> {
598        Box::pin(async move { Ok(None) })
599    }
600
601    fn find_api_key_by_id(&self, _id: Uuid) -> StoreFuture<'_, Option<ApiKey>> {
602        Box::pin(async move { Ok(None) })
603    }
604
605    fn list_api_keys_by_user(&self, _user_id: Uuid) -> StoreFuture<'_, Vec<ApiKey>> {
606        Box::pin(async move {
607            Err(StoreError::Database(
608                "ApiKeyStore not available in worker".to_string(),
609            ))
610        })
611    }
612
613    fn update_api_key(&self, _id: Uuid, _update: ApiKeyUpdate) -> StoreFuture<'_, ()> {
614        Box::pin(async move {
615            Err(StoreError::Database(
616                "ApiKeyStore not available in worker".to_string(),
617            ))
618        })
619    }
620
621    fn touch_api_key(&self, _id: Uuid) -> StoreFuture<'_, ()> {
622        Box::pin(async move {
623            Err(StoreError::Database(
624                "ApiKeyStore not available in worker".to_string(),
625            ))
626        })
627    }
628
629    fn delete_api_key(&self, _id: Uuid) -> StoreFuture<'_, ()> {
630        Box::pin(async move {
631            Err(StoreError::Database(
632                "ApiKeyStore not available in worker".to_string(),
633            ))
634        })
635    }
636}
637
638impl AuditLogStore for ApiRunStore {
639    fn append_audit_log(&self, _entry: NewAuditLogEntry) -> StoreFuture<'_, AuditLogEntry> {
640        Box::pin(async move {
641            Err(StoreError::Database(
642                "AuditLogStore not available in worker".to_string(),
643            ))
644        })
645    }
646
647    fn list_audit_logs(
648        &self,
649        _filter: AuditLogFilter,
650        _page: u32,
651        _per_page: u32,
652    ) -> StoreFuture<'_, Page<AuditLogEntry>> {
653        Box::pin(async move {
654            Err(StoreError::Database(
655                "AuditLogStore not available in worker".to_string(),
656            ))
657        })
658    }
659}
660
661impl SecretStore for ApiRunStore {
662    fn get_secret(&self, key: &str) -> StoreFuture<'_, Option<Secret>> {
663        let key = key.to_string();
664        Box::pin(async move {
665            let resp = self
666                .client
667                .get(self.internal(&format!("/secrets/{key}")))
668                .bearer_auth(&self.token)
669                .send()
670                .await
671                .map_err(Self::err)?;
672
673            if resp.status() == reqwest::StatusCode::NOT_FOUND {
674                return Ok(None);
675            }
676
677            if !resp.status().is_success() {
678                let body = resp.text().await.unwrap_or_default();
679                return Err(Self::status_err(&body));
680            }
681
682            let api_resp: ApiResponse<Secret> = resp.json().await.map_err(Self::err)?;
683            Ok(Some(api_resp.data))
684        })
685    }
686
687    fn set_secret(&self, _key: &str, _value: &str) -> StoreFuture<'_, Secret> {
688        Box::pin(async move {
689            Err(StoreError::Database(
690                "SecretStore not available in worker".to_string(),
691            ))
692        })
693    }
694
695    fn delete_secret(&self, _key: &str) -> StoreFuture<'_, bool> {
696        Box::pin(async move {
697            Err(StoreError::Database(
698                "SecretStore not available in worker".to_string(),
699            ))
700        })
701    }
702
703    fn list_secret_keys(&self, _prefix: &str) -> StoreFuture<'_, Vec<String>> {
704        Box::pin(async move {
705            Err(StoreError::Database(
706                "SecretStore not available in worker".to_string(),
707            ))
708        })
709    }
710
711    fn list_secrets(
712        &self,
713        _prefix: &str,
714        _page: u32,
715        _per_page: u32,
716    ) -> StoreFuture<'_, Page<SecretMetadata>> {
717        Box::pin(async move {
718            Err(StoreError::Database(
719                "SecretStore not available in worker".to_string(),
720            ))
721        })
722    }
723
724    fn secret_key_status(&self) -> StoreFuture<'_, KeyVersionStatus> {
725        Box::pin(async move {
726            Err(StoreError::Database(
727                "SecretStore not available in worker".to_string(),
728            ))
729        })
730    }
731
732    fn rotate_secrets(&self, _request: RotationRequest) -> StoreFuture<'_, RotationBatch> {
733        Box::pin(async move {
734            Err(StoreError::Database(
735                "SecretStore not available in worker".to_string(),
736            ))
737        })
738    }
739}
740
741impl LogStore for ApiRunStore {
742    fn append_logs(&self, _entries: NewLogEntries) -> StoreFuture<'_, ()> {
743        // Log persistence is handled by the API server via push_logs.
744        Box::pin(async move { Ok(()) })
745    }
746
747    fn get_logs(
748        &self,
749        _run_id: Uuid,
750        _filter: LogFilter,
751        _cursor: Option<Uuid>,
752        _limit: u32,
753    ) -> StoreFuture<'_, Vec<LogEntry>> {
754        Box::pin(async move {
755            Err(StoreError::Database(
756                "LogStore not available in worker".to_string(),
757            ))
758        })
759    }
760}
761
762impl ArtifactStore for ApiRunStore {
763    fn create_artifact(&self, artifact: NewArtifact) -> StoreFuture<'_, Artifact> {
764        Box::pin(async move {
765            // The worker records an artifact by uploading its bytes, which the
766            // API writes and registers in one call. There is no metadata-only
767            // route, so this must never be reached from the worker.
768            let _ = artifact;
769            Err(StoreError::Database(
770                "create_artifact not supported via worker API — upload the bytes instead"
771                    .to_string(),
772            ))
773        })
774    }
775
776    fn get_artifact(&self, _step_id: Uuid, _name: &str) -> StoreFuture<'_, Option<Artifact>> {
777        // The worker resolves artifacts through find_artifact_for_input, which
778        // needs the run to scope the search. Keep this total rather than add a
779        // route the worker never calls.
780        Box::pin(async move { Ok(None) })
781    }
782
783    fn list_artifacts_for_run(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Artifact>> {
784        Box::pin(async move {
785            let resp = self
786                .client
787                .get(self.internal(&format!("/runs/{run_id}/artifacts")))
788                .bearer_auth(&self.token)
789                .send()
790                .await
791                .map_err(Self::err)?;
792
793            if !resp.status().is_success() {
794                let body = resp.text().await.unwrap_or_default();
795                return Err(Self::status_err(&body));
796            }
797
798            let api_resp: ApiResponse<Vec<Artifact>> = resp.json().await.map_err(Self::err)?;
799            Ok(api_resp.data)
800        })
801    }
802
803    fn find_artifact_for_input(&self, lookup: ArtifactLookup) -> StoreFuture<'_, Option<Artifact>> {
804        Box::pin(async move {
805            // Resolved worker-side from the run's steps and artifacts, so the
806            // matching rule stays in one place instead of being duplicated in a
807            // dedicated route.
808            let steps = self.list_steps(lookup.run_id).await?;
809
810            let Some(producer) = steps
811                .iter()
812                .filter(|step| {
813                    step.attempt == lookup.attempt
814                        && step.name == lookup.step_name
815                        && step.position < lookup.before_position
816                })
817                .max_by_key(|step| step.position)
818            else {
819                return Ok(None);
820            };
821
822            let artifacts = self.list_artifacts_for_run(lookup.run_id).await?;
823            Ok(artifacts
824                .into_iter()
825                .find(|artifact| artifact.step_id == producer.id && artifact.name == lookup.name))
826        })
827    }
828
829    fn find_artifact_by_sha256(&self, _sha256: &str) -> StoreFuture<'_, Option<Artifact>> {
830        Box::pin(async move { Ok(None) })
831    }
832
833    fn count_artifacts_by_storage_key(&self, _storage_key: &str) -> StoreFuture<'_, u64> {
834        Box::pin(async move { Ok(0) })
835    }
836
837    fn list_all_storage_keys(&self) -> StoreFuture<'_, Vec<String>> {
838        Box::pin(async move { Ok(Vec::new()) })
839    }
840}
841
842impl ScheduleStore for ApiRunStore {
843    fn create_schedule(&self, _req: NewSchedule) -> StoreFuture<'_, Schedule> {
844        Box::pin(async move {
845            Err(StoreError::Database(
846                "ScheduleStore not available in worker".to_string(),
847            ))
848        })
849    }
850
851    fn find_schedule_by_id(&self, _id: Uuid) -> StoreFuture<'_, Option<Schedule>> {
852        Box::pin(async move { Ok(None) })
853    }
854
855    fn list_schedules(&self, _page: u32, _per_page: u32) -> StoreFuture<'_, Page<Schedule>> {
856        Box::pin(async move {
857            Err(StoreError::Database(
858                "ScheduleStore not available in worker".to_string(),
859            ))
860        })
861    }
862
863    fn update_schedule(&self, _id: Uuid, _update: ScheduleUpdate) -> StoreFuture<'_, Schedule> {
864        Box::pin(async move {
865            Err(StoreError::Database(
866                "ScheduleStore not available in worker".to_string(),
867            ))
868        })
869    }
870
871    fn delete_schedule(&self, _id: Uuid) -> StoreFuture<'_, ()> {
872        Box::pin(async move {
873            Err(StoreError::Database(
874                "ScheduleStore not available in worker".to_string(),
875            ))
876        })
877    }
878
879    fn list_due_schedules(&self) -> StoreFuture<'_, Vec<Schedule>> {
880        Box::pin(async move {
881            Err(StoreError::Database(
882                "ScheduleStore not available in worker".to_string(),
883            ))
884        })
885    }
886
887    fn fire_due_schedule(
888        &self,
889        _id: Uuid,
890        _occurrence: DateTime<Utc>,
891        _next: ScheduleNext,
892    ) -> StoreFuture<'_, Option<ScheduleFiring>> {
893        Box::pin(async move {
894            Err(StoreError::Database(
895                "ScheduleStore not available in worker".to_string(),
896            ))
897        })
898    }
899}
900
901impl ApprovalDelegationStore for ApiRunStore {
902    fn create_delegation(
903        &self,
904        _req: NewApprovalDelegation,
905    ) -> StoreFuture<'_, ApprovalDelegation> {
906        Box::pin(async move {
907            Err(StoreError::Database(
908                "ApprovalDelegationStore not available in worker".to_string(),
909            ))
910        })
911    }
912
913    fn find_delegation_by_id(&self, _id: Uuid) -> StoreFuture<'_, Option<ApprovalDelegation>> {
914        Box::pin(async move { Ok(None) })
915    }
916
917    fn list_active_delegations(
918        &self,
919        _filter: DelegationFilter,
920        _page: u32,
921        _per_page: u32,
922    ) -> StoreFuture<'_, Page<ApprovalDelegation>> {
923        Box::pin(async move {
924            Err(StoreError::Database(
925                "ApprovalDelegationStore not available in worker".to_string(),
926            ))
927        })
928    }
929
930    fn find_active_delegation(
931        &self,
932        _from_user_id: Uuid,
933        _to_user_id: Uuid,
934        _workflow_name: &str,
935    ) -> StoreFuture<'_, Option<ApprovalDelegation>> {
936        Box::pin(async move { Ok(None) })
937    }
938
939    fn delete_delegation(&self, _id: Uuid) -> StoreFuture<'_, ()> {
940        Box::pin(async move {
941            Err(StoreError::Database(
942                "ApprovalDelegationStore not available in worker".to_string(),
943            ))
944        })
945    }
946}
947
948/// Error for the account administration methods the worker never needs.
949fn account_method_unavailable(method: &str) -> StoreError {
950    StoreError::Database(format!(
951        "ProviderAccountStore::{method} not available in worker"
952    ))
953}
954
955impl ProviderAccountStore for ApiRunStore {
956    fn create_provider_account(
957        &self,
958        _req: NewProviderAccount,
959    ) -> StoreFuture<'_, ProviderAccount> {
960        Box::pin(async { Err(account_method_unavailable("create_provider_account")) })
961    }
962
963    fn get_provider_account(&self, _id: Uuid) -> StoreFuture<'_, Option<ProviderAccount>> {
964        Box::pin(async { Err(account_method_unavailable("get_provider_account")) })
965    }
966
967    fn list_provider_accounts_by_ids(
968        &self,
969        _ids: Vec<Uuid>,
970    ) -> StoreFuture<'_, Vec<ProviderAccount>> {
971        Box::pin(async { Err(account_method_unavailable("list_provider_accounts_by_ids")) })
972    }
973
974    fn find_provider_account_by_name(
975        &self,
976        _name: &str,
977    ) -> StoreFuture<'_, Option<ProviderAccount>> {
978        Box::pin(async { Err(account_method_unavailable("find_provider_account_by_name")) })
979    }
980
981    fn list_provider_accounts(
982        &self,
983        _kind: Option<String>,
984        _page: u32,
985        _per_page: u32,
986    ) -> StoreFuture<'_, Page<ProviderAccount>> {
987        Box::pin(async { Err(account_method_unavailable("list_provider_accounts")) })
988    }
989
990    fn update_provider_account(
991        &self,
992        _id: Uuid,
993        _update: ProviderAccountUpdate,
994    ) -> StoreFuture<'_, ProviderAccount> {
995        Box::pin(async { Err(account_method_unavailable("update_provider_account")) })
996    }
997
998    fn delete_provider_account(&self, _id: Uuid) -> StoreFuture<'_, bool> {
999        Box::pin(async { Err(account_method_unavailable("delete_provider_account")) })
1000    }
1001
1002    fn list_provider_account_windows(
1003        &self,
1004        _ids: Vec<Uuid>,
1005    ) -> StoreFuture<'_, Vec<ProviderAccountWindow>> {
1006        Box::pin(async { Err(account_method_unavailable("list_provider_account_windows")) })
1007    }
1008
1009    fn list_provider_account_usage(
1010        &self,
1011        _id: Uuid,
1012        _since: DateTime<Utc>,
1013    ) -> StoreFuture<'_, Vec<ProviderAccountUsagePoint>> {
1014        Box::pin(async { Err(account_method_unavailable("list_provider_account_usage")) })
1015    }
1016
1017    fn record_provider_account_observation(
1018        &self,
1019        id: Uuid,
1020        observation: NewProviderAccountObservation,
1021    ) -> StoreFuture<'_, Vec<ProviderAccountWindow>> {
1022        Box::pin(async move {
1023            let resp = self
1024                .client
1025                .post(self.internal(&format!("/provider-accounts/{id}/observations")))
1026                .bearer_auth(&self.token)
1027                .json(&observation)
1028                .send()
1029                .await
1030                .map_err(Self::err)?;
1031
1032            if resp.status() == StatusCode::NOT_FOUND {
1033                return Err(StoreError::ProviderAccountNotFound(id));
1034            }
1035            if !resp.status().is_success() {
1036                let body = resp.text().await.unwrap_or_default();
1037                return Err(Self::status_err(&body));
1038            }
1039
1040            let api_resp: ApiResponse<Vec<ProviderAccountWindow>> =
1041                resp.json().await.map_err(Self::err)?;
1042            Ok(api_resp.data)
1043        })
1044    }
1045
1046    fn purge_provider_account_usage(&self, _before: DateTime<Utc>) -> StoreFuture<'_, u64> {
1047        Box::pin(async { Err(account_method_unavailable("purge_provider_account_usage")) })
1048    }
1049
1050    fn list_provider_account_candidates(
1051        &self,
1052        kind: String,
1053    ) -> StoreFuture<'_, Vec<ProviderAccountCandidate>> {
1054        Box::pin(async move {
1055            let resp = self
1056                .client
1057                .get(self.internal("/provider-accounts/candidates"))
1058                .bearer_auth(&self.token)
1059                .query(&[("kind", kind)])
1060                .send()
1061                .await
1062                .map_err(Self::err)?;
1063
1064            if !resp.status().is_success() {
1065                let body = resp.text().await.unwrap_or_default();
1066                return Err(Self::status_err(&body));
1067            }
1068
1069            let api_resp: ApiResponse<Vec<ProviderAccountCandidate>> =
1070                resp.json().await.map_err(Self::err)?;
1071            Ok(api_resp.data)
1072        })
1073    }
1074}
1075
1076/// The [`StoreError::ConcurrencyConflict`] carried by a 409 error body, if it
1077/// is one.
1078///
1079/// Any other 409 body yields `None`, so the caller reports it as is.
1080/// Query parameters carrying a worker's capabilities, as comma-separated
1081/// lists. `workflows` is only sent when the worker restricts them; `tags` is
1082/// always sent, even empty, so the API knows the worker routes by tags.
1083pub(crate) fn capability_query(capabilities: &WorkerCapabilities) -> Vec<(&'static str, String)> {
1084    let mut query = Vec::with_capacity(2);
1085    if let Some(ref workflows) = capabilities.workflows {
1086        query.push(("workflows", workflows.join(",")));
1087    }
1088    query.push(("tags", capabilities.tags.join(",")));
1089    query
1090}
1091
1092fn concurrency_conflict(body: &str) -> Option<StoreError> {
1093    #[derive(serde::Deserialize)]
1094    struct Details {
1095        key: String,
1096        run_id: Uuid,
1097    }
1098    #[derive(serde::Deserialize)]
1099    struct Error {
1100        code: String,
1101        details: Details,
1102    }
1103    #[derive(serde::Deserialize)]
1104    struct Envelope {
1105        error: Error,
1106    }
1107
1108    let envelope: Envelope = from_str(body).ok()?;
1109    if envelope.error.code != CONCURRENCY_CONFLICT_CODE {
1110        return None;
1111    }
1112    let Details { key, run_id } = envelope.error.details;
1113    Some(StoreError::ConcurrencyConflict { key, run_id })
1114}
1115
1116/// Error for the signal methods only the API server runs.
1117fn signal_method_unavailable(method: &str) -> StoreError {
1118    StoreError::Database(format!("SignalStore::{method} not available in worker"))
1119}
1120
1121impl SignalStore for ApiRunStore {
1122    fn insert_signal(&self, _signal: NewSignal) -> StoreFuture<'_, SignalInsert> {
1123        Box::pin(async { Err(signal_method_unavailable("insert_signal")) })
1124    }
1125
1126    fn list_signals(
1127        &self,
1128        _filter: SignalFilter,
1129        _page: u32,
1130        _per_page: u32,
1131    ) -> StoreFuture<'_, Page<Signal>> {
1132        Box::pin(async { Err(signal_method_unavailable("list_signals")) })
1133    }
1134
1135    fn list_signals_for_key(
1136        &self,
1137        name: &str,
1138        key: &str,
1139        since: DateTime<Utc>,
1140    ) -> StoreFuture<'_, Vec<Signal>> {
1141        let query = [
1142            ("name", name.to_string()),
1143            ("key", key.to_string()),
1144            ("since", since.to_rfc3339()),
1145        ];
1146        Box::pin(async move {
1147            let resp = self
1148                .client
1149                .get(self.internal("/signals"))
1150                .bearer_auth(&self.token)
1151                .query(&query)
1152                .send()
1153                .await
1154                .map_err(Self::err)?;
1155
1156            if !resp.status().is_success() {
1157                let body = resp.text().await.unwrap_or_default();
1158                return Err(Self::status_err(&body));
1159            }
1160
1161            let api_resp: ApiResponse<Vec<Signal>> = resp.json().await.map_err(Self::err)?;
1162            Ok(api_resp.data)
1163        })
1164    }
1165
1166    fn list_signal_waiters(&self, _name: &str, _key: &str) -> StoreFuture<'_, Vec<Step>> {
1167        Box::pin(async { Err(signal_method_unavailable("list_signal_waiters")) })
1168    }
1169
1170    fn resolve_signal_step(
1171        &self,
1172        step_id: Uuid,
1173        output: Value,
1174    ) -> StoreFuture<'_, SignalStepResolution> {
1175        Box::pin(async move {
1176            let resp = self
1177                .client
1178                .post(self.internal(&format!("/steps/{step_id}/signal-resolution")))
1179                .bearer_auth(&self.token)
1180                .json(&json!({ "output": output }))
1181                .send()
1182                .await
1183                .map_err(Self::err)?;
1184
1185            if resp.status() == StatusCode::NOT_FOUND {
1186                return Err(StoreError::StepNotFound(step_id));
1187            }
1188            if !resp.status().is_success() {
1189                let body = resp.text().await.unwrap_or_default();
1190                return Err(Self::status_err(&body));
1191            }
1192
1193            let api_resp: ApiResponse<SignalStepResolution> =
1194                resp.json().await.map_err(Self::err)?;
1195            Ok(api_resp.data)
1196        })
1197    }
1198
1199    fn suspend_run_on_signal(
1200        &self,
1201        run_id: Uuid,
1202        step_id: Uuid,
1203        deadline_at: DateTime<Utc>,
1204    ) -> StoreFuture<'_, bool> {
1205        Box::pin(async move {
1206            let resp = self
1207                .client
1208                .post(self.internal(&format!("/runs/{run_id}/signal-suspension")))
1209                .bearer_auth(&self.token)
1210                .json(&json!({ "step_id": step_id, "deadline_at": deadline_at }))
1211                .send()
1212                .await
1213                .map_err(Self::err)?;
1214
1215            if !resp.status().is_success() {
1216                let body = resp.text().await.unwrap_or_default();
1217                return Err(Self::status_err(&body));
1218            }
1219
1220            let api_resp: ApiResponse<bool> = resp.json().await.map_err(Self::err)?;
1221            Ok(api_resp.data)
1222        })
1223    }
1224
1225    fn purge_signals(&self, _before: DateTime<Utc>) -> StoreFuture<'_, u64> {
1226        Box::pin(async { Err(signal_method_unavailable("purge_signals")) })
1227    }
1228}
1229
1230#[cfg(test)]
1231mod tests {
1232    use std::collections::HashMap;
1233    use std::sync::Arc;
1234
1235    use super::*;
1236
1237    use axum::serve;
1238    use ironflow_api::routes::{RouterConfig, create_router};
1239    use ironflow_api::state::AppState;
1240    use ironflow_auth::jwt::JwtConfig;
1241    use ironflow_core::providers::claude::ClaudeCodeProvider;
1242    use ironflow_engine::engine::Engine;
1243    use ironflow_engine::notify::Event;
1244    use ironflow_store::entities::{PARENT_RUN_ID_LABEL, TriggerKind};
1245    use ironflow_store::memory::InMemoryStore;
1246    use tokio::net::TcpListener;
1247    use tokio::spawn;
1248    use tokio::sync::broadcast;
1249
1250    use serde_json::json;
1251
1252    #[tokio::test]
1253    async fn create_run_returns_error_on_unreachable_server() {
1254        let store = ApiRunStore::new("http://127.0.0.1:1", "token");
1255        let req = NewRun {
1256            created_by: None,
1257            workflow_name: "test".to_string(),
1258            trigger: TriggerKind::Manual,
1259            payload: json!({}),
1260            max_retries: 0,
1261            handler_version: None,
1262            labels: HashMap::new(),
1263            scheduled_at: None,
1264            idempotency_key: None,
1265            concurrency_key: None,
1266            priority: 0,
1267            concurrency_limits: Vec::new(),
1268            max_cost_usd: None,
1269            worker_tags: Vec::new(),
1270        };
1271        let result = store.create_run(req).await;
1272        assert!(result.is_err());
1273    }
1274
1275    #[tokio::test]
1276    async fn insert_signal_not_available_in_worker() {
1277        let store = ApiRunStore::new("http://localhost:3000", "token");
1278        let result = store
1279            .insert_signal(NewSignal {
1280                name: "demo.done".to_string(),
1281                key: "k1".to_string(),
1282                payload: json!({}),
1283                idempotency_id: None,
1284            })
1285            .await;
1286        match result {
1287            Err(StoreError::Database(msg)) => assert!(msg.contains("not available in worker")),
1288            other => panic!("expected an unavailable-method error, got {other:?}"),
1289        }
1290    }
1291
1292    #[tokio::test]
1293    async fn find_run_by_idempotency_key_not_supported() {
1294        let store = ApiRunStore::new("http://localhost:3000", "token");
1295        let result = store.find_run_by_idempotency_key("github:abc").await;
1296        match result {
1297            Err(StoreError::Database(msg)) => assert!(msg.contains("not supported")),
1298            other => panic!("expected an unsupported-operation error, got {other:?}"),
1299        }
1300    }
1301
1302    #[tokio::test]
1303    async fn list_runs_not_supported() {
1304        let store = ApiRunStore::new("http://localhost:3000", "token");
1305        let result = store
1306            .list_runs(ironflow_store::entities::RunFilter::default(), 0, 10)
1307            .await;
1308        assert!(result.is_err());
1309        match result {
1310            Err(StoreError::Database(msg)) => {
1311                assert!(msg.contains("not supported"));
1312            }
1313            _ => panic!("expected Database error"),
1314        }
1315    }
1316
1317    #[tokio::test]
1318    async fn get_stats_not_supported() {
1319        let store = ApiRunStore::new("http://localhost:3000", "token");
1320        let result = store.get_stats(RunFilter::default()).await;
1321        assert!(result.is_err());
1322        match result {
1323            Err(StoreError::Database(msg)) => {
1324                assert!(msg.contains("not supported"));
1325            }
1326            _ => panic!("expected Database error"),
1327        }
1328    }
1329
1330    #[test]
1331    fn api_run_store_clone() {
1332        let store = ApiRunStore::new("http://localhost:3000", "token");
1333        let store2 = store.clone();
1334        assert_eq!(store.base_url, store2.base_url);
1335        assert_eq!(store.token, store2.token);
1336    }
1337
1338    #[test]
1339    fn api_run_store_with_trailing_slash() {
1340        let store = ApiRunStore::new("http://localhost:3000/", "token");
1341        assert_eq!(store.base_url, "http://localhost:3000");
1342    }
1343
1344    #[test]
1345    fn api_run_store_without_trailing_slash() {
1346        let store = ApiRunStore::new("http://localhost:3000", "token");
1347        assert_eq!(store.base_url, "http://localhost:3000");
1348    }
1349
1350    #[test]
1351    fn api_run_store_builds_internal_url() {
1352        let store = ApiRunStore::new("http://localhost:3000", "token");
1353        let url = store.internal("/runs/123");
1354        assert_eq!(url, "http://localhost:3000/api/v1/internal/runs/123");
1355    }
1356    /// Serve the real API router over TCP and return its base URL.
1357    async fn spawn_api() -> String {
1358        let store = Arc::new(InMemoryStore::new());
1359        let engine = Engine::new(store.clone(), Arc::new(ClaudeCodeProvider::new()));
1360        let jwt_config = Arc::new(JwtConfig {
1361            secret: "test-secret".to_string(),
1362            access_token_ttl_secs: 900,
1363            refresh_token_ttl_secs: 604800,
1364            cookie_domain: None,
1365            cookie_secure: false,
1366        });
1367        let (event_sender, _) = broadcast::channel::<Event>(1);
1368        let state = AppState::new(
1369            store,
1370            Arc::new(engine),
1371            jwt_config,
1372            "test-worker-token".to_string(),
1373            event_sender,
1374        );
1375        let config = RouterConfig {
1376            rate_limit_auth: None,
1377            rate_limit_general: None,
1378            ..RouterConfig::default()
1379        };
1380        let router = create_router(state, config);
1381        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1382        let addr = listener.local_addr().unwrap();
1383        spawn(async move {
1384            serve(listener, router).await.unwrap();
1385        });
1386        format!("http://{addr}")
1387    }
1388
1389    fn keyed_run(key: &str) -> NewRun {
1390        NewRun {
1391            created_by: None,
1392            workflow_name: "child".to_string(),
1393            trigger: TriggerKind::Workflow,
1394            payload: json!({}),
1395            max_retries: 0,
1396            handler_version: None,
1397            labels: HashMap::new(),
1398            scheduled_at: None,
1399            idempotency_key: None,
1400            concurrency_key: Some(key.to_string()),
1401            priority: 0,
1402            concurrency_limits: Vec::new(),
1403            max_cost_usd: None,
1404            worker_tags: Vec::new(),
1405        }
1406    }
1407
1408    #[tokio::test]
1409    async fn list_active_descendants_goes_through_the_internal_route() {
1410        let store = ApiRunStore::new(&spawn_api().await, "test-worker-token");
1411        let mut root = keyed_run("root");
1412        (root.trigger, root.concurrency_key) = (TriggerKind::Manual, None);
1413        let root = store.create_run(root).await.unwrap().into_run();
1414        let mut child = keyed_run("child");
1415        child.labels = HashMap::from([(PARENT_RUN_ID_LABEL.to_string(), root.id.to_string())]);
1416        let child = store.create_run(child).await.unwrap().into_run();
1417
1418        let found = store.list_active_descendants(root.id).await.unwrap();
1419        assert_eq!(found.iter().map(|r| r.id).collect::<Vec<_>>(), [child.id]);
1420        assert_eq!(found[0].status.state, RunStatus::Pending);
1421        assert!(
1422            store
1423                .list_active_descendants(Uuid::now_v7())
1424                .await
1425                .unwrap()
1426                .is_empty()
1427        );
1428    }
1429
1430    #[tokio::test]
1431    async fn create_run_maps_409_concurrency_conflict() {
1432        let store = ApiRunStore::new(&spawn_api().await, "test-worker-token");
1433
1434        let holder = store
1435            .create_run(keyed_run("issue:12"))
1436            .await
1437            .expect("the first run takes the key")
1438            .into_run();
1439        assert_eq!(holder.concurrency_key.as_deref(), Some("issue:12"));
1440
1441        match store.create_run(keyed_run("issue:12")).await {
1442            Err(StoreError::ConcurrencyConflict { key, run_id }) => {
1443                assert_eq!(key, "issue:12");
1444                assert_eq!(run_id, holder.id);
1445            }
1446            other => panic!("expected a concurrency conflict, got {other:?}"),
1447        }
1448    }
1449
1450    #[test]
1451    fn concurrency_conflict_reads_the_api_error_envelope() {
1452        let run_id = Uuid::now_v7();
1453        let body = json!({
1454            "error": {
1455                "code": "CONCURRENCY_CONFLICT",
1456                "message": "concurrency key \"issue:12\" is held by active run",
1457                "details": { "key": "issue:12", "run_id": run_id }
1458            }
1459        })
1460        .to_string();
1461
1462        match concurrency_conflict(&body) {
1463            Some(StoreError::ConcurrencyConflict { key, run_id: held }) => {
1464                assert_eq!(key, "issue:12");
1465                assert_eq!(held, run_id);
1466            }
1467            other => panic!("expected a concurrency conflict, got {other:?}"),
1468        }
1469    }
1470
1471    #[test]
1472    fn concurrency_conflict_ignores_other_409_bodies() {
1473        let other_code = json!({
1474            "error": {
1475                "code": "CONFLICT",
1476                "message": "something else",
1477                "details": { "key": "issue:12", "run_id": Uuid::now_v7() }
1478            }
1479        })
1480        .to_string();
1481        assert!(concurrency_conflict(&other_code).is_none());
1482        assert!(concurrency_conflict("not json").is_none());
1483        assert!(concurrency_conflict("").is_none());
1484    }
1485
1486    #[test]
1487    fn capability_query_sends_workflows_and_tags() {
1488        let caps = WorkerCapabilities::new(
1489            Some(vec!["build".to_string(), "deploy".to_string()]),
1490            vec!["gpu".to_string(), "region:eu".to_string()],
1491        );
1492        assert_eq!(
1493            capability_query(&caps),
1494            vec![
1495                ("workflows", "build,deploy".to_string()),
1496                ("tags", "gpu,region:eu".to_string()),
1497            ]
1498        );
1499    }
1500
1501    #[test]
1502    fn capability_query_omits_workflows_when_unrestricted() {
1503        let caps = WorkerCapabilities::new(None, vec!["gpu".to_string()]);
1504        assert_eq!(capability_query(&caps), vec![("tags", "gpu".to_string())]);
1505    }
1506
1507    #[test]
1508    fn capability_query_sends_empty_tags() {
1509        let caps = WorkerCapabilities::new(Some(Vec::new()), Vec::new());
1510        assert_eq!(
1511            capability_query(&caps),
1512            vec![("workflows", String::new()), ("tags", String::new())]
1513        );
1514    }
1515}