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