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