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