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