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