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