Skip to main content

alien_core/resources/
worker.rs

1use crate::error::{ErrorData, Result};
2use crate::resource::{ResourceDefinition, ResourceOutputsDefinition, ResourceRef, ResourceType};
3use crate::{PublicEndpointOutput, WorkerPublicEndpoint, APEX_HOST_LABEL};
4use alien_error::AlienError;
5use bon::Builder;
6use serde::{Deserialize, Serialize};
7use std::any::Any;
8use std::collections::HashMap;
9use std::fmt::Debug;
10
11/// Specifies the source of the worker's executable code.
12/// This can be a pre-built container image or source code that the system will build.
13#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
14#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
15#[serde(rename_all = "camelCase", tag = "type")]
16pub enum WorkerCode {
17    /// Container image.
18    #[serde(rename_all = "camelCase")]
19    Image {
20        /// Container image (e.g., `ghcr.io/myorg/myimage:latest`).
21        image: String,
22    },
23    /// Source code to be built.
24    #[serde(rename_all = "camelCase")]
25    Source {
26        /// The source directory to build from
27        src: String,
28        /// Toolchain configuration with type-safe options
29        toolchain: ToolchainConfig,
30    },
31}
32
33/// Configuration for different programming language toolchains.
34/// Each toolchain provides type-safe build configuration and auto-detection capabilities.
35#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
36#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
37#[serde(rename_all = "lowercase", tag = "type")]
38pub enum ToolchainConfig {
39    /// Rust with Cargo build system
40    #[serde(rename_all = "camelCase")]
41    Rust {
42        /// Name of the binary to build and run
43        binary_name: String,
44    },
45    /// TypeScript/JavaScript compiled to single executable with Bun
46    #[serde(rename_all = "camelCase")]
47    TypeScript {
48        /// Name of the compiled binary (defaults to package.json name if not specified)
49        #[serde(default, skip_serializing_if = "Option::is_none")]
50        binary_name: Option<String>,
51    },
52    /// Docker build from Dockerfile
53    #[serde(rename_all = "camelCase")]
54    Docker {
55        /// Dockerfile path relative to src (default: "Dockerfile")
56        #[serde(skip_serializing_if = "Option::is_none")]
57        dockerfile: Option<String>,
58        /// Build arguments for docker build
59        #[serde(skip_serializing_if = "Option::is_none")]
60        build_args: Option<HashMap<String, String>>,
61        /// Multi-stage build target
62        #[serde(skip_serializing_if = "Option::is_none")]
63        target: Option<String>,
64    },
65}
66
67/// Defines what triggers a worker execution.
68#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
69#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
70#[serde(tag = "type", rename_all = "camelCase")]
71pub enum WorkerTrigger {
72    /// Worker triggered by queue messages (always 1 message per invocation)
73    Queue {
74        /// Reference to the queue resource
75        queue: ResourceRef,
76    },
77    /// Worker triggered by storage events (object created, deleted, etc.)
78    Storage {
79        /// Reference to the storage resource
80        storage: ResourceRef,
81        /// Events to trigger on (e.g., ["created", "deleted"])
82        events: Vec<String>,
83    },
84    /// Worker triggered on a schedule (cron expression)
85    Schedule {
86        /// Cron expression for scheduling (standard 5-field unix cron)
87        cron: String,
88    },
89}
90
91impl WorkerTrigger {
92    /// Creates a queue trigger for the specified queue resource.
93    /// The worker will be automatically invoked when messages arrive in the queue.
94    /// Each message is processed individually (batch size of 1).
95    pub fn queue<R: ?Sized>(queue: &R) -> Self
96    where
97        for<'a> &'a R: Into<ResourceRef>,
98    {
99        let queue_ref: ResourceRef = queue.into();
100        WorkerTrigger::Queue { queue: queue_ref }
101    }
102
103    /// Creates a storage trigger for the specified storage resource.
104    /// The worker will be invoked when matching events occur on the storage resource.
105    pub fn storage<R: ?Sized>(storage: &R, events: Vec<String>) -> Self
106    where
107        for<'a> &'a R: Into<ResourceRef>,
108    {
109        let storage_ref: ResourceRef = storage.into();
110        WorkerTrigger::Storage {
111            storage: storage_ref,
112            events,
113        }
114    }
115
116    /// Creates a schedule trigger with the specified cron expression.
117    /// Uses standard 5-field unix cron format (minute hour day-of-month month day-of-week).
118    pub fn schedule<S: Into<String>>(cron: S) -> Self {
119        WorkerTrigger::Schedule { cron: cron.into() }
120    }
121}
122
123/// Represents a serverless worker that executes code in response to triggers or direct invocations.
124/// Workers are the primary compute resource in serverless applications, designed to be stateless and ephemeral.
125#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Builder)]
126#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
127#[serde(rename_all = "camelCase", deny_unknown_fields)]
128#[builder(start_fn = new)]
129pub struct Worker {
130    /// Identifier for the worker. Must contain only alphanumeric characters, hyphens, and underscores ([A-Za-z0-9-_]).
131    /// Maximum 64 characters.
132    #[builder(start_fn)]
133    pub id: String,
134
135    /// List of resource references this worker depends on.
136    // TODO: We need to verify that the same link isn't added multiple times.
137    #[builder(field)]
138    pub links: Vec<ResourceRef>,
139
140    /// List of triggers that define what events automatically invoke this worker.
141    /// If empty, the worker is only invokable directly via HTTP calls or platform-specific invocation APIs.
142    /// When configured, the worker will be automatically invoked when any of the specified trigger conditions are met.
143    #[builder(field)]
144    pub triggers: Vec<WorkerTrigger>,
145
146    /// Public endpoints exposed by this worker.
147    #[builder(field)]
148    #[serde(default, skip_serializing_if = "Vec::is_empty")]
149    pub public_endpoints: Vec<WorkerPublicEndpoint>,
150
151    /// Permission profile name that defines the permissions granted to this worker.
152    /// This references a profile defined in the stack's permission definitions.
153    pub permissions: String,
154
155    /// Code for the worker, either a pre-built image or source code to be built.
156    pub code: WorkerCode,
157
158    /// Memory allocated to the worker in megabytes (MB).
159    /// Default: 512
160    ///
161    /// Platform-specific constraints:
162    /// - **AWS Lambda**: 128–10240 MB in 1 MB increments
163    /// - **GCP Cloud Run**: 128–32768 MB
164    /// - **Azure Container Apps**: fixed CPU/memory pairs — 512, 1024, 1536, 2048, 2560,
165    ///   3072, 3584, 4096 MB. Values below 512 are automatically rounded up at deploy time.
166    #[builder(default = default_memory_mb())]
167    #[serde(default = "default_memory_mb")]
168    #[cfg_attr(feature = "openapi", schema(default = default_memory_mb))]
169    pub memory_mb: u32,
170
171    /// Maximum execution time for the worker in seconds.
172    /// Constraints: 1‑3600 seconds (platform-specific limits may apply)
173    /// Default: 180
174    #[builder(
175        default = default_timeout_seconds(),
176        with = |timeout_seconds: u32| -> crate::Result<_> {
177            validate_timeout_seconds(timeout_seconds)
178        }
179    )]
180    #[serde(
181        default = "default_timeout_seconds",
182        deserialize_with = "deserialize_timeout_seconds"
183    )]
184    #[cfg_attr(
185        feature = "openapi",
186        schema(default = default_timeout_seconds, minimum = 1, maximum = 3600)
187    )]
188    pub timeout_seconds: u32,
189
190    /// Key-value pairs to set as environment variables for the worker.
191    #[builder(default)]
192    #[serde(default)]
193    pub environment: HashMap<String, String>,
194
195    /// Whether the worker can receive remote commands via the Commands protocol.
196    /// When enabled, the platform pushes commands into the Worker runtime,
197    /// which executes registered handlers.
198    #[builder(default = default_commands_enabled())]
199    #[serde(default = "default_commands_enabled")]
200    #[cfg_attr(feature = "openapi", schema(default = default_commands_enabled))]
201    pub commands_enabled: bool,
202
203    /// Maximum number of concurrent executions allowed for the worker.
204    /// None means platform default applies.
205    pub concurrency_limit: Option<u32>,
206
207    /// Whether this worker hosts sandbox sessions.
208    ///
209    /// Set by preflight, not by an application: on GCP a sandbox is a subprocess of the Cloud Run
210    /// instance running the app, and the instance can only launch one if its container declares
211    /// it. Declaring it by hand would be a permission the workload does not need.
212    #[builder(default)]
213    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
214    pub sandbox_launcher: bool,
215
216    /// Optional readiness probe configuration.
217    /// Only applicable for workers with Public ingress.
218    /// When configured, the probe will be executed after provisioning/update to verify the worker is ready.
219    pub readiness_probe: Option<ReadinessProbe>,
220}
221
222impl Worker {
223    /// The resource type identifier for Workers
224    pub const RESOURCE_TYPE: ResourceType = ResourceType::from_static("worker");
225
226    /// Returns the permission profile name for this worker.
227    pub fn get_permissions(&self) -> &str {
228        &self.permissions
229    }
230
231    fn validate_public_endpoints(&self) -> Result<()> {
232        let mut endpoint_names = std::collections::HashSet::new();
233        let mut apex_endpoint_name: Option<&str> = None;
234        for endpoint in &self.public_endpoints {
235            endpoint.validate_for_resource(&self.id)?;
236            if !endpoint_names.insert(endpoint.name.as_str()) {
237                return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
238                    resource_id: self.id.clone(),
239                    reason: format!("duplicate public endpoint name '{}'", endpoint.name),
240                }));
241            }
242            if endpoint.host_label.as_deref() == Some(APEX_HOST_LABEL) {
243                if let Some(existing_name) = apex_endpoint_name {
244                    return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
245                        resource_id: self.id.clone(),
246                        reason: format!(
247                            "only one apex public endpoint is allowed per resource; '{}' already uses hostLabel '@'",
248                            existing_name
249                        ),
250                    }));
251                }
252                apex_endpoint_name = Some(endpoint.name.as_str());
253            }
254        }
255
256        Ok(())
257    }
258}
259
260fn default_memory_mb() -> u32 {
261    512
262}
263
264fn default_timeout_seconds() -> u32 {
265    180
266}
267
268/// Longest Worker execution supported by every Commands delivery path.
269pub const MAX_WORKER_TIMEOUT_SECONDS: u32 = 3600;
270
271fn deserialize_timeout_seconds<'de, D>(deserializer: D) -> std::result::Result<u32, D::Error>
272where
273    D: serde::Deserializer<'de>,
274{
275    let value = u32::deserialize(deserializer)?;
276    validate_timeout_seconds(value).map_err(serde::de::Error::custom)
277}
278
279fn validate_timeout_seconds(timeout_seconds: u32) -> Result<u32> {
280    if (1..=MAX_WORKER_TIMEOUT_SECONDS).contains(&timeout_seconds) {
281        return Ok(timeout_seconds);
282    }
283
284    Err(AlienError::new(ErrorData::WorkerTimeoutInvalid {
285        timeout_seconds,
286        max_timeout_seconds: MAX_WORKER_TIMEOUT_SECONDS,
287    }))
288}
289
290fn default_commands_enabled() -> bool {
291    false
292}
293
294/// HTTP method for readiness probe requests.
295#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
296#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
297#[serde(rename_all = "UPPERCASE")]
298#[derive(Default)]
299pub enum HttpMethod {
300    #[default]
301    Get,
302    Post,
303    Put,
304    Delete,
305    Head,
306    Options,
307    Patch,
308}
309
310/// Configuration for HTTP-based readiness probe.
311/// This probe is executed after worker provisioning/update to verify the worker is ready to serve traffic.
312/// Only works with workers that have Public ingress.
313#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
314#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
315#[serde(rename_all = "camelCase")]
316pub struct ReadinessProbe {
317    /// HTTP method to use for the probe request.
318    /// Default: GET
319    #[serde(default)]
320    pub method: HttpMethod,
321
322    /// Path to request for the probe (e.g., "/health", "/ready").
323    /// Default: "/"
324    #[serde(default = "default_probe_path")]
325    pub path: String,
326}
327
328fn default_probe_path() -> String {
329    "/".to_string()
330}
331
332impl Default for ReadinessProbe {
333    fn default() -> Self {
334        Self {
335            method: HttpMethod::default(),
336            path: default_probe_path(),
337        }
338    }
339}
340
341use crate::resources::worker::worker_builder::State;
342
343impl<S: State> WorkerBuilder<S> {
344    /// Links the worker to another resource with specified permissions.
345    /// Accepts a reference to any type `R` where `&R` can be converted into `ResourceRef`.
346    pub fn link<R: ?Sized>(mut self, resource: &R) -> Self
347    where
348        for<'a> &'a R: Into<ResourceRef>, // Use Higher-Rank Trait Bound (HRTB)
349    {
350        // Perform the conversion from &R to ResourceRef using .into()
351        let resource_ref: ResourceRef = resource.into();
352        self.links.push(resource_ref);
353        self
354    }
355
356    /// Adds a trigger to the worker. Workers can have multiple triggers.
357    /// Each trigger will independently invoke the worker when its conditions are met.
358    ///
359    /// # Examples
360    /// ```rust
361    /// # use alien_core::{Worker, WorkerTrigger, WorkerCode, Queue};
362    /// # let queue1 = Queue::new("queue-1".to_string()).build();
363    /// # let queue2 = Queue::new("queue-2".to_string()).build();
364    /// let worker = Worker::new("my-worker".to_string())
365    ///     .code(WorkerCode::Image { image: "test:latest".to_string() })
366    ///     .permissions("execution".to_string())
367    ///     .trigger(WorkerTrigger::queue(&queue1))
368    ///     .trigger(WorkerTrigger::queue(&queue2))
369    ///     .build();
370    /// ```
371    pub fn trigger(mut self, trigger: WorkerTrigger) -> Self {
372        self.triggers.push(trigger);
373        self
374    }
375
376    /// Exposes a named public endpoint for the worker.
377    pub fn public_endpoint(mut self, endpoint: WorkerPublicEndpoint) -> Self {
378        self.public_endpoints.push(endpoint);
379        self
380    }
381}
382
383// Implementation of ResourceDefinition trait for Worker
384impl ResourceDefinition for Worker {
385    fn get_resource_type(&self) -> ResourceType {
386        Self::RESOURCE_TYPE
387    }
388
389    fn id(&self) -> &str {
390        &self.id
391    }
392
393    fn get_dependencies(&self) -> Vec<ResourceRef> {
394        let mut dependencies = self.links.clone();
395
396        // Add trigger dependencies
397        for trigger in &self.triggers {
398            match trigger {
399                WorkerTrigger::Queue { queue } => {
400                    dependencies.push(queue.clone());
401                }
402                WorkerTrigger::Storage { storage, .. } => {
403                    dependencies.push(storage.clone());
404                }
405                WorkerTrigger::Schedule { .. } => {
406                    // Schedule triggers don't depend on other resources
407                }
408            }
409        }
410
411        dependencies
412    }
413
414    fn get_permissions(&self) -> Option<&str> {
415        Some(&self.permissions)
416    }
417
418    fn validate_update(&self, new_config: &dyn ResourceDefinition) -> Result<()> {
419        // Downcast to Worker type to use the existing validate_update method
420        let new_worker = new_config
421            .as_any()
422            .downcast_ref::<Worker>()
423            .ok_or_else(|| {
424                AlienError::new(ErrorData::UnexpectedResourceType {
425                    resource_id: self.id.clone(),
426                    expected: Self::RESOURCE_TYPE,
427                    actual: new_config.get_resource_type(),
428                })
429            })?;
430
431        if self.id != new_worker.id {
432            return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
433                resource_id: self.id.clone(),
434                reason: "the 'id' field is immutable".to_string(),
435            }));
436        }
437        self.validate_public_endpoints()?;
438        new_worker.validate_public_endpoints()?;
439        if self.public_endpoints != new_worker.public_endpoints {
440            return Err(AlienError::new(ErrorData::InvalidResourceUpdate {
441                resource_id: self.id.clone(),
442                reason: "the 'publicEndpoints' field is immutable".to_string(),
443            }));
444        }
445        Ok(())
446    }
447
448    fn as_any(&self) -> &dyn Any {
449        self
450    }
451
452    fn as_any_mut(&mut self) -> &mut dyn Any {
453        self
454    }
455
456    fn box_clone(&self) -> Box<dyn ResourceDefinition> {
457        Box::new(self.clone())
458    }
459
460    fn resource_eq(&self, other: &dyn ResourceDefinition) -> bool {
461        other.as_any().downcast_ref::<Worker>() == Some(self)
462    }
463
464    fn to_json_value(&self) -> serde_json::Result<serde_json::Value> {
465        serde_json::to_value(self)
466    }
467}
468
469/// Outputs generated by a successfully provisioned Worker.
470#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
471#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
472#[serde(rename_all = "camelCase")]
473pub struct WorkerOutputs {
474    /// The platform-specific worker name.
475    pub worker_name: String,
476    /// Public endpoints resolved for this worker.
477    #[serde(default, skip_serializing_if = "HashMap::is_empty")]
478    pub public_endpoints: HashMap<String, PublicEndpointOutput>,
479    /// The ARN or platform-specific identifier.
480    #[serde(skip_serializing_if = "Option::is_none")]
481    pub identifier: Option<String>,
482    /// Push target for commands delivery. Platform-specific:
483    /// - AWS: Lambda function name or ARN
484    /// - GCP: Full Pub/Sub topic path (projects/{project}/topics/{topic})
485    /// - Azure: Service Bus "{namespace}/{queue}"
486    #[serde(default, skip_serializing_if = "Option::is_none")]
487    pub commands_push_target: Option<String>,
488}
489
490impl ResourceOutputsDefinition for WorkerOutputs {
491    fn get_resource_type(&self) -> ResourceType {
492        Worker::RESOURCE_TYPE.clone()
493    }
494
495    fn as_any(&self) -> &dyn Any {
496        self
497    }
498
499    fn box_clone(&self) -> Box<dyn ResourceOutputsDefinition> {
500        Box::new(self.clone())
501    }
502
503    fn outputs_eq(&self, other: &dyn ResourceOutputsDefinition) -> bool {
504        other.as_any().downcast_ref::<WorkerOutputs>() == Some(self)
505    }
506
507    fn to_json_value(&self) -> serde_json::Result<serde_json::Value> {
508        serde_json::to_value(self)
509    }
510}
511
512#[cfg(test)]
513mod tests {
514    use super::*;
515    use crate::Storage;
516
517    #[test]
518    fn test_worker_builder_direct_refs() {
519        let dummy_storage = Storage::new("test-storage".to_string()).build();
520        let dummy_storage_2 = Storage::new("test-storage-2".to_string()).build();
521
522        let worker = Worker::new("my-worker".to_string())
523            .code(WorkerCode::Image {
524                image: "test-image".to_string(),
525            })
526            .permissions("execution".to_string())
527            .link(&dummy_storage) // Pass reference directly
528            .link(&dummy_storage_2) // Add a second link
529            .build();
530
531        assert_eq!(worker.id, "my-worker");
532        assert_eq!(
533            worker.code,
534            WorkerCode::Image {
535                image: "test-image".to_string()
536            }
537        );
538
539        // Verify permissions was set correctly
540        assert_eq!(worker.permissions, "execution");
541
542        // Verify links were added correctly
543        assert!(worker
544            .links
545            .contains(&ResourceRef::new(Storage::RESOURCE_TYPE, "test-storage")));
546        assert!(worker
547            .links
548            .contains(&ResourceRef::new(Storage::RESOURCE_TYPE, "test-storage-2")));
549        assert_eq!(worker.links.len(), 2); // Expect 2 links now
550    }
551
552    #[test]
553    fn test_worker_with_readiness_probe() {
554        let probe = ReadinessProbe {
555            method: HttpMethod::Post,
556            path: "/health".to_string(),
557        };
558
559        let worker = Worker::new("my-worker".to_string())
560            .code(WorkerCode::Image {
561                image: "test-image".to_string(),
562            })
563            .permissions("execution".to_string())
564            .public_endpoint(WorkerPublicEndpoint {
565                name: "api".to_string(),
566                host_label: None,
567                wildcard_subdomains: false,
568            })
569            .readiness_probe(probe.clone())
570            .build();
571
572        assert_eq!(worker.id, "my-worker");
573        assert_eq!(worker.public_endpoints[0].name, "api");
574        assert_eq!(worker.readiness_probe, Some(probe));
575    }
576
577    #[test]
578    fn test_readiness_probe_defaults() {
579        let probe = ReadinessProbe::default();
580        assert_eq!(probe.method, HttpMethod::Get);
581        assert_eq!(probe.path, "/");
582    }
583
584    #[test]
585    fn test_worker_with_rust_toolchain() {
586        let worker = Worker::new("my-rust-worker".to_string())
587            .code(WorkerCode::Source {
588                src: "./".to_string(),
589                toolchain: ToolchainConfig::Rust {
590                    binary_name: "my-app".to_string(),
591                },
592            })
593            .permissions("execution".to_string())
594            .build();
595
596        assert_eq!(worker.id, "my-rust-worker");
597
598        match &worker.code {
599            WorkerCode::Source { src, toolchain } => {
600                assert_eq!(src, "./");
601                assert_eq!(
602                    toolchain,
603                    &ToolchainConfig::Rust {
604                        binary_name: "my-app".to_string(),
605                    }
606                );
607            }
608            _ => panic!("Expected Source code"),
609        }
610    }
611
612    #[test]
613    fn test_worker_with_typescript_toolchain() {
614        let worker = Worker::new("my-ts-worker".to_string())
615            .code(WorkerCode::Source {
616                src: "./".to_string(),
617                toolchain: ToolchainConfig::TypeScript {
618                    binary_name: Some("my-ts-worker".to_string()),
619                },
620            })
621            .permissions("execution".to_string())
622            .build();
623
624        assert_eq!(worker.id, "my-ts-worker");
625
626        match &worker.code {
627            WorkerCode::Source { src, toolchain } => {
628                assert_eq!(src, "./");
629                assert_eq!(
630                    toolchain,
631                    &ToolchainConfig::TypeScript {
632                        binary_name: Some("my-ts-worker".to_string())
633                    }
634                );
635            }
636            _ => panic!("Expected Source code"),
637        }
638    }
639
640    #[test]
641    fn test_worker_with_queue_trigger() {
642        use crate::Queue;
643
644        let queue = Queue::new("test-queue".to_string()).build();
645
646        let worker = Worker::new("triggered-worker".to_string())
647            .code(WorkerCode::Image {
648                image: "test-image".to_string(),
649            })
650            .permissions("execution".to_string())
651            .trigger(WorkerTrigger::queue(&queue))
652            .build();
653
654        assert_eq!(worker.triggers.len(), 1);
655        if let WorkerTrigger::Queue { queue: queue_ref } = &worker.triggers[0] {
656            assert_eq!(queue_ref.resource_type, Queue::RESOURCE_TYPE);
657            assert_eq!(queue_ref.id, "test-queue");
658        } else {
659            panic!("Expected queue trigger");
660        }
661    }
662
663    #[test]
664    fn test_worker_trigger_dependencies() {
665        use crate::Queue;
666
667        let queue = Queue::new("test-queue".to_string()).build();
668        let storage = Storage::new("test-storage".to_string()).build();
669
670        let worker = Worker::new("triggered-worker".to_string())
671            .code(WorkerCode::Image {
672                image: "test-image".to_string(),
673            })
674            .permissions("execution".to_string())
675            .link(&storage) // regular link dependency
676            .trigger(WorkerTrigger::queue(&queue)) // trigger dependency
677            .build();
678
679        let dependencies = worker.get_dependencies();
680
681        // Should have both link and trigger dependencies
682        assert_eq!(dependencies.len(), 2);
683        assert!(dependencies.contains(&ResourceRef::new(Storage::RESOURCE_TYPE, "test-storage")));
684        assert!(dependencies.contains(&ResourceRef::new(Queue::RESOURCE_TYPE, "test-queue")));
685    }
686
687    #[test]
688    fn test_worker_trigger_helper_methods() {
689        use crate::Queue;
690
691        let queue = Queue::new("my-queue".to_string()).build();
692
693        // Test the helper method
694        let trigger = WorkerTrigger::queue(&queue);
695
696        if let WorkerTrigger::Queue { queue: queue_ref } = trigger {
697            assert_eq!(queue_ref.resource_type, Queue::RESOURCE_TYPE);
698            assert_eq!(queue_ref.id, "my-queue");
699        } else {
700            panic!("Expected queue trigger");
701        }
702    }
703
704    #[test]
705    fn test_worker_with_multiple_triggers() {
706        use crate::Queue;
707
708        let queue1 = Queue::new("queue-1".to_string()).build();
709        let queue2 = Queue::new("queue-2".to_string()).build();
710
711        let worker = Worker::new("multi-triggered-worker".to_string())
712            .code(WorkerCode::Image {
713                image: "test-image".to_string(),
714            })
715            .permissions("execution".to_string())
716            .trigger(WorkerTrigger::queue(&queue1))
717            .trigger(WorkerTrigger::queue(&queue2))
718            .trigger(WorkerTrigger::schedule("0 * * * *".to_string()))
719            .build();
720
721        assert_eq!(worker.triggers.len(), 3);
722
723        // Check first queue trigger
724        if let WorkerTrigger::Queue { queue: queue_ref } = &worker.triggers[0] {
725            assert_eq!(queue_ref.id, "queue-1");
726        } else {
727            panic!("Expected first trigger to be queue-1");
728        }
729
730        // Check second queue trigger
731        if let WorkerTrigger::Queue { queue: queue_ref } = &worker.triggers[1] {
732            assert_eq!(queue_ref.id, "queue-2");
733        } else {
734            panic!("Expected second trigger to be queue-2");
735        }
736
737        // Check schedule trigger
738        if let WorkerTrigger::Schedule { cron } = &worker.triggers[2] {
739            assert_eq!(cron, "0 * * * *");
740        } else {
741            panic!("Expected third trigger to be schedule");
742        }
743
744        // Check dependencies include both queues
745        let dependencies = worker.get_dependencies();
746        assert_eq!(dependencies.len(), 2); // Only queues, schedule has no dependency
747        assert!(dependencies.contains(&ResourceRef::new(Queue::RESOURCE_TYPE, "queue-1")));
748        assert!(dependencies.contains(&ResourceRef::new(Queue::RESOURCE_TYPE, "queue-2")));
749    }
750
751    #[test]
752    fn test_worker_with_commands_enabled() {
753        let worker = Worker::new("cmd-worker".to_string())
754            .code(WorkerCode::Image {
755                image: "test-image".to_string(),
756            })
757            .permissions("execution".to_string())
758            .commands_enabled(true)
759            .build();
760
761        assert_eq!(worker.id, "cmd-worker");
762        assert!(worker.public_endpoints.is_empty());
763        assert_eq!(worker.commands_enabled, true);
764    }
765
766    #[test]
767    fn test_worker_defaults() {
768        let worker = Worker::new("default-worker".to_string())
769            .code(WorkerCode::Image {
770                image: "test-image".to_string(),
771            })
772            .permissions("execution".to_string())
773            .build();
774
775        // Test that defaults are applied correctly
776        assert!(worker.public_endpoints.is_empty());
777        assert_eq!(worker.commands_enabled, false);
778        assert_eq!(worker.memory_mb, 512);
779        assert_eq!(worker.timeout_seconds, 180);
780    }
781
782    #[test]
783    fn worker_deserialization_rejects_timeout_outside_supported_range() {
784        let worker = Worker::new("timeout-worker".to_string())
785            .code(WorkerCode::Image {
786                image: "test-image".to_string(),
787            })
788            .permissions("execution".to_string())
789            .build();
790        let mut value = serde_json::to_value(worker).expect("serialize worker");
791
792        value["timeoutSeconds"] = serde_json::json!(0);
793        assert!(serde_json::from_value::<Worker>(value.clone()).is_err());
794
795        value["timeoutSeconds"] = serde_json::json!(MAX_WORKER_TIMEOUT_SECONDS + 1);
796        assert!(serde_json::from_value::<Worker>(value).is_err());
797    }
798
799    #[test]
800    fn worker_builder_rejects_zero_timeout() {
801        let Err(error) = Worker::new("timeout-worker".to_string()).timeout_seconds(0) else {
802            panic!("zero timeout must be rejected");
803        };
804
805        assert_eq!(error.code, "WORKER_TIMEOUT_INVALID");
806        assert_eq!(error.http_status_code, Some(400));
807    }
808
809    #[test]
810    fn worker_builder_rejects_timeout_above_maximum() {
811        let Err(error) = Worker::new("timeout-worker".to_string())
812            .timeout_seconds(MAX_WORKER_TIMEOUT_SECONDS + 1)
813        else {
814            panic!("timeout above maximum must be rejected");
815        };
816
817        assert_eq!(error.code, "WORKER_TIMEOUT_INVALID");
818    }
819
820    #[test]
821    fn worker_builder_accepts_minimum_timeout() {
822        let worker = Worker::new("timeout-worker".to_string())
823            .timeout_seconds(1)
824            .expect("minimum timeout is valid")
825            .code(WorkerCode::Image {
826                image: "test-image".to_string(),
827            })
828            .permissions("execution".to_string())
829            .build();
830
831        assert_eq!(worker.timeout_seconds, 1);
832    }
833
834    #[test]
835    fn worker_builder_accepts_maximum_timeout() {
836        let worker = Worker::new("timeout-worker".to_string())
837            .timeout_seconds(MAX_WORKER_TIMEOUT_SECONDS)
838            .expect("maximum timeout is valid")
839            .code(WorkerCode::Image {
840                image: "test-image".to_string(),
841            })
842            .permissions("execution".to_string())
843            .build();
844
845        assert_eq!(worker.timeout_seconds, MAX_WORKER_TIMEOUT_SECONDS);
846    }
847
848    #[test]
849    fn test_worker_public_ingress_with_commands() {
850        let worker = Worker::new("public-cmd-worker".to_string())
851            .code(WorkerCode::Image {
852                image: "test-image".to_string(),
853            })
854            .permissions("execution".to_string())
855            .public_endpoint(WorkerPublicEndpoint {
856                name: "api".to_string(),
857                host_label: None,
858                wildcard_subdomains: false,
859            })
860            .commands_enabled(true)
861            .build();
862
863        assert_eq!(worker.public_endpoints[0].name, "api");
864        assert_eq!(worker.commands_enabled, true);
865    }
866
867    #[test]
868    fn worker_rejects_multiple_apex_public_endpoints() {
869        let worker = Worker::new("apex-worker".to_string())
870            .code(WorkerCode::Image {
871                image: "test-image".to_string(),
872            })
873            .permissions("execution".to_string())
874            .public_endpoint(WorkerPublicEndpoint {
875                name: "api".to_string(),
876                host_label: Some(APEX_HOST_LABEL.to_string()),
877                wildcard_subdomains: false,
878            })
879            .public_endpoint(WorkerPublicEndpoint {
880                name: "admin".to_string(),
881                host_label: Some(APEX_HOST_LABEL.to_string()),
882                wildcard_subdomains: false,
883            })
884            .build();
885
886        assert!(worker.validate_public_endpoints().is_err());
887    }
888}