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