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