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