1use 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#[derive(serde::Deserialize, serde::Serialize)]
31struct ApiResponse<T> {
32 data: T,
33}
34
35#[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 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 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 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 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 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 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 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}