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