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