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