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