Skip to main content

fakecloud_ecs/
service.rs

1use std::sync::Arc;
2
3use async_trait::async_trait;
4use chrono::Utc;
5use http::StatusCode;
6use serde_json::{json, Map, Value};
7use tokio::sync::Mutex as AsyncMutex;
8
9use fakecloud_aws::arn::Arn;
10use fakecloud_core::service::{AwsRequest, AwsResponse, AwsService, AwsServiceError};
11use fakecloud_persistence::SnapshotStore;
12
13use crate::state::{
14    Attribute, AttributeRef, AwsLogsConfig, CapacityProvider, CircuitBreakerConfig, Cluster,
15    Container, ContainerInstance, Deployment, EcsSnapshot, EcsState, Service, SharedEcsState,
16    TagEntry, Task, TaskDefinition, TaskSet, ECS_SNAPSHOT_SCHEMA_VERSION,
17};
18
19const SUPPORTED_ACTIONS: &[&str] = &[
20    "CreateCluster",
21    "DescribeClusters",
22    "DeleteCluster",
23    "ListClusters",
24    "UpdateCluster",
25    "UpdateClusterSettings",
26    "PutClusterCapacityProviders",
27    "RegisterTaskDefinition",
28    "DescribeTaskDefinition",
29    "DeregisterTaskDefinition",
30    "DeleteTaskDefinitions",
31    "ListTaskDefinitions",
32    "ListTaskDefinitionFamilies",
33    "TagResource",
34    "UntagResource",
35    "ListTagsForResource",
36    "PutAccountSetting",
37    "PutAccountSettingDefault",
38    "DeleteAccountSetting",
39    "ListAccountSettings",
40    "RunTask",
41    "StartTask",
42    "StopTask",
43    "DescribeTasks",
44    "ListTasks",
45    "CreateService",
46    "UpdateService",
47    "DeleteService",
48    "DescribeServices",
49    "ListServices",
50    "ListServicesByNamespace",
51    "RegisterContainerInstance",
52    "DeregisterContainerInstance",
53    "DescribeContainerInstances",
54    "ListContainerInstances",
55    "UpdateContainerAgent",
56    "UpdateContainerInstancesState",
57    "PutAttributes",
58    "DeleteAttributes",
59    "ListAttributes",
60    "CreateCapacityProvider",
61    "DeleteCapacityProvider",
62    "DescribeCapacityProviders",
63    "UpdateCapacityProvider",
64    "GetTaskProtection",
65    "UpdateTaskProtection",
66    "CreateTaskSet",
67    "UpdateTaskSet",
68    "DeleteTaskSet",
69    "DescribeTaskSets",
70    "UpdateServicePrimaryTaskSet",
71    "ExecuteCommand",
72    "SubmitContainerStateChange",
73    "SubmitTaskStateChange",
74    "SubmitAttachmentStateChanges",
75    "DiscoverPollEndpoint",
76    "StopServiceDeployment",
77    "ContinueServiceDeployment",
78    "ListServiceDeployments",
79    "DescribeServiceDeployments",
80    "DescribeServiceRevisions",
81    "RegisterDaemonTaskDefinition",
82    "DescribeDaemonTaskDefinition",
83    "DeleteDaemonTaskDefinition",
84    "ListDaemonTaskDefinitions",
85    "CreateDaemon",
86    "DescribeDaemon",
87    "UpdateDaemon",
88    "DeleteDaemon",
89    "ListDaemons",
90    "DescribeDaemonDeployments",
91    "ListDaemonDeployments",
92    "DescribeDaemonRevisions",
93    "CreateExpressGatewayService",
94    "DescribeExpressGatewayService",
95    "UpdateExpressGatewayService",
96    "DeleteExpressGatewayService",
97];
98
99pub struct EcsService {
100    state: SharedEcsState,
101    snapshot_store: Option<Arc<dyn SnapshotStore>>,
102    snapshot_lock: Arc<AsyncMutex<()>>,
103    runtime: Option<Arc<crate::runtime::EcsRuntime>>,
104    role_trust_validator: Option<Arc<dyn fakecloud_core::auth::RoleTrustValidator>>,
105}
106
107impl EcsService {
108    pub fn new(state: SharedEcsState) -> Self {
109        Self {
110            state,
111            snapshot_store: None,
112            snapshot_lock: Arc::new(AsyncMutex::new(())),
113            runtime: None,
114            role_trust_validator: None,
115        }
116    }
117
118    pub fn with_snapshot_store(mut self, store: Arc<dyn SnapshotStore>) -> Self {
119        self.snapshot_store = Some(store);
120        self
121    }
122
123    pub fn with_runtime(mut self, runtime: Arc<crate::runtime::EcsRuntime>) -> Self {
124        self.runtime = Some(runtime);
125        self
126    }
127
128    pub fn with_role_trust_validator(
129        mut self,
130        validator: Arc<dyn fakecloud_core::auth::RoleTrustValidator>,
131    ) -> Self {
132        self.role_trust_validator = Some(validator);
133        self
134    }
135
136    fn check_pass_role(&self, account_id: &str, role_arn: &str) -> Result<(), AwsServiceError> {
137        let Some(ref validator) = self.role_trust_validator else {
138            return Ok(());
139        };
140        if let Err(err) = validator.validate(account_id, role_arn, "ecs-tasks.amazonaws.com") {
141            return Err(AwsServiceError::aws_error(
142                StatusCode::BAD_REQUEST,
143                "InvalidParameterException",
144                err.to_string(),
145            ));
146        }
147        Ok(())
148    }
149
150    pub fn state_handle(&self) -> &SharedEcsState {
151        &self.state
152    }
153
154    async fn save_snapshot(&self) {
155        save_ecs_snapshot(
156            &self.state,
157            self.snapshot_store.clone(),
158            &self.snapshot_lock,
159        )
160        .await;
161    }
162
163    /// Build a hook that persists the current state when invoked, or `None` in
164    /// memory mode. The CloudFormation provisioner mutates `state` directly and
165    /// uses this to write a CFN-provisioned resource through to disk.
166    pub fn snapshot_hook(&self) -> Option<fakecloud_persistence::SnapshotHook> {
167        let store = self.snapshot_store.clone()?;
168        let state = self.state.clone();
169        let lock = self.snapshot_lock.clone();
170        Some(Arc::new(move || {
171            let state = state.clone();
172            let store = store.clone();
173            let lock = lock.clone();
174            Box::pin(async move {
175                save_ecs_snapshot(&state, Some(store), &lock).await;
176            })
177        }))
178    }
179
180    /// Reconcile persisted task state with reality after a fakecloud
181    /// restart. Same bug class as RDS #1338: tasks were persisted with
182    /// `lastStatus = RUNNING` but their docker containers are gone, so
183    /// DescribeTasks reported phantom-running tasks and `docker exec`
184    /// (for ECS Exec) would fail.
185    ///
186    /// Mark every non-STOPPED task as STOPPED with a `stoppedReason`
187    /// that explains the restart, and reset service `runningCount` /
188    /// `pendingCount` to zero. The service scheduler ticker then
189    /// reconciles desiredCount and launches fresh tasks. Standalone
190    /// `RunTask` tasks aren't auto-respawned because we don't have
191    /// enough context to replay them safely; callers re-invoke.
192    pub async fn reconcile_persisted_tasks(&self) {
193        let touched = {
194            let mut accounts = self.state.write();
195            let mut touched_tasks = 0usize;
196            let mut touched_services = 0usize;
197            let now = chrono::Utc::now();
198            for (_, state) in accounts.iter_mut() {
199                for task in state.tasks.values_mut() {
200                    if task.last_status != "STOPPED" {
201                        task.last_status = "STOPPED".to_string();
202                        task.desired_status = "STOPPED".to_string();
203                        task.stop_code = Some("EssentialContainerExited".to_string());
204                        task.stopped_reason =
205                            Some("fakecloud restart - backing container lost".to_string());
206                        if task.stopping_at.is_none() {
207                            task.stopping_at = Some(now);
208                        }
209                        if task.stopped_at.is_none() {
210                            task.stopped_at = Some(now);
211                        }
212                        for container in task.containers.iter_mut() {
213                            container.last_status = "STOPPED".to_string();
214                        }
215                        touched_tasks += 1;
216                    }
217                }
218                for service in state.services.values_mut() {
219                    if service.running_count != 0 || service.pending_count != 0 {
220                        service.running_count = 0;
221                        service.pending_count = 0;
222                        for deployment in service.deployments.iter_mut() {
223                            deployment.running_count = 0;
224                            deployment.pending_count = 0;
225                        }
226                        touched_services += 1;
227                    }
228                }
229            }
230            (touched_tasks, touched_services)
231        };
232        if touched.0 + touched.1 > 0 {
233            tracing::info!(
234                tasks = touched.0,
235                services = touched.1,
236                "reconciled persisted ecs tasks / service counts after restart",
237            );
238            self.save_snapshot().await;
239        }
240    }
241
242    /// Converge every ACTIVE service toward its `desiredCount`, launching
243    /// replacement tasks for any shortfall. Real ECS maintains `desiredCount`
244    /// autonomously; `reconcile_persisted_tasks` STOPs all tasks and zeroes
245    /// counts on restart "trusting a scheduler ticker," but no such ticker
246    /// existed, so services with `desiredCount=N` stayed at `runningCount=0`
247    /// forever. This is that ticker. It also re-launches tasks lost to crashed
248    /// containers during normal operation. bug-audit 2026-06-15, 4.7.
249    ///
250    /// One pass: for each service, count its non-STOPPED tasks (matched by the
251    /// `ecs-svc/<name>` `started_by` tag, same as `recompute_service_counts`);
252    /// if that count is below `desiredCount`, spawn the difference via the
253    /// shared `spawn_service_tasks` path and hand the new task IDs to the
254    /// runtime. Persist only when something was launched.
255    pub async fn reconcile_service_desired_counts(&self) {
256        // Phase 1 (write lock): find shortfalls, spawn PENDING task rows,
257        // and refresh each service's running/pending counts so DescribeServices
258        // is accurate immediately. Collect the new task IDs to launch after the
259        // lock is released.
260        let mut to_launch: Vec<(String, String)> = Vec::new(); // (account_id, task_id)
261        {
262            let mut accounts = self.state.write();
263            for (account_id, state) in accounts.iter_mut() {
264                // Snapshot the services we may need to scale so we don't hold
265                // an aliasing borrow of `state` while mutating it.
266                let scalable: Vec<Service> = state
267                    .services
268                    .values()
269                    .filter(|s| {
270                        s.status == "ACTIVE"
271                            && s.scheduling_strategy != "DAEMON"
272                            && s.desired_count > 0
273                    })
274                    .cloned()
275                    .collect();
276                for service in scalable {
277                    let service_tag = format!("ecs-svc/{}", service.service_name);
278                    let mut active = 0i32;
279                    for t in state.tasks.values() {
280                        if t.started_by.as_deref() == Some(service_tag.as_str())
281                            && t.cluster_name == service.cluster_name
282                            && matches!(
283                                t.last_status.as_str(),
284                                "RUNNING" | "PENDING" | "PROVISIONING"
285                            )
286                        {
287                            active += 1;
288                        }
289                    }
290                    let shortfall = service.desired_count - active;
291                    if shortfall <= 0 {
292                        continue;
293                    }
294                    let launch_type = if service.launch_type.is_empty() {
295                        "FARGATE"
296                    } else {
297                        &service.launch_type
298                    };
299                    // Background re-spawn has no request region; keep the new
300                    // tasks in the region the service's stored ARN already
301                    // carries (stamped from req.region at CreateService time).
302                    let svc_region = service
303                        .cluster_arn
304                        .split(':')
305                        .nth(3)
306                        .map(str::to_string)
307                        .unwrap_or_else(|| state.region.clone());
308                    let ids = spawn_service_tasks(
309                        state,
310                        &svc_region,
311                        &service,
312                        shortfall,
313                        service.role_arn.as_deref().unwrap_or(""),
314                        launch_type,
315                        None,
316                    );
317                    if ids.is_empty() {
318                        continue;
319                    }
320                    tracing::info!(
321                        service = %service.service_name,
322                        cluster = %service.cluster_name,
323                        launched = ids.len(),
324                        desired = service.desired_count,
325                        "ecs scheduler: launching tasks to converge to desiredCount",
326                    );
327                    // Reflect the new PENDING tasks in the stored counts.
328                    // Services are keyed by "cluster/name", not the bare name —
329                    // looking up by service_name returned None, so these counts
330                    // were never updated (stale snapshot until DescribeServices
331                    // recomputed them).
332                    let svc_key = crate::state::EcsState::service_key(
333                        &service.cluster_name,
334                        &service.service_name,
335                    );
336                    if let Some(svc) = state.services.get_mut(&svc_key) {
337                        svc.pending_count = svc.pending_count.saturating_add(ids.len() as i32);
338                        for d in svc.deployments.iter_mut() {
339                            if d.status == "PRIMARY" {
340                                d.pending_count = d.pending_count.saturating_add(ids.len() as i32);
341                            }
342                        }
343                    }
344                    for id in ids {
345                        to_launch.push((account_id.to_string(), id));
346                    }
347                }
348            }
349        }
350
351        if to_launch.is_empty() {
352            return;
353        }
354
355        // Phase 2 (no lock): hand new tasks to the runtime, which advances them
356        // PENDING -> RUNNING in the background. Without a runtime (docker/podman
357        // absent) mark them STOPPED so they don't linger PENDING forever.
358        if let Some(rt) = &self.runtime {
359            for (account_id, task_id) in &to_launch {
360                rt.clone()
361                    .run_task(self.state.clone(), task_id.clone(), account_id.clone());
362            }
363        } else {
364            let mut accounts = self.state.write();
365            for (account_id, task_id) in &to_launch {
366                if let Some(state) = accounts.get_mut(account_id) {
367                    if let Some(t) = state.tasks.get_mut(task_id) {
368                        t.last_status = "STOPPED".into();
369                        t.desired_status = "STOPPED".into();
370                        t.stop_code = Some("TaskFailedToStart".into());
371                        t.stopped_reason = Some(format!(
372                            "No container runtime available (docker/podman not installed). {}",
373                            fakecloud_core::container_net::CONTAINER_RUNTIME_HINT
374                        ));
375                        t.stopped_at = Some(chrono::Utc::now());
376                        for c in t.containers.iter_mut() {
377                            c.last_status = "STOPPED".into();
378                        }
379                    }
380                }
381            }
382        }
383
384        self.save_snapshot().await;
385    }
386}
387
388pub async fn save_ecs_snapshot(
389    state: &SharedEcsState,
390    store: Option<Arc<dyn SnapshotStore>>,
391    lock: &AsyncMutex<()>,
392) {
393    let Some(store) = store else {
394        return;
395    };
396    let _guard = lock.lock().await;
397    let snapshot = EcsSnapshot {
398        schema_version: ECS_SNAPSHOT_SCHEMA_VERSION,
399        accounts: Some(state.read().clone()),
400    };
401    let join = tokio::task::spawn_blocking(move || -> std::io::Result<()> {
402        let bytes = serde_json::to_vec(&snapshot)
403            .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e.to_string()))?;
404        store.save(&bytes)
405    })
406    .await;
407    match join {
408        Ok(Ok(())) => {}
409        Ok(Err(err)) => tracing::error!(%err, "failed to write ecs snapshot"),
410        Err(err) => tracing::error!(%err, "ecs snapshot task panicked"),
411    }
412}
413
414/// Background scheduler ticker: periodically converges ECS services toward
415/// `desiredCount`. Wired in `fakecloud-server` like the other tickers. Real
416/// ECS does this continuously; we tick every few seconds. bug-audit 4.7.
417pub async fn run_scheduler_ticker(service: Arc<EcsService>, interval: std::time::Duration) {
418    let mut ticker = tokio::time::interval(interval);
419    ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
420    // Skip the immediate first tick; restart reconciliation runs separately.
421    ticker.tick().await;
422    loop {
423        ticker.tick().await;
424        service.reconcile_service_desired_counts().await;
425    }
426}
427
428#[async_trait]
429impl AwsService for EcsService {
430    fn service_name(&self) -> &str {
431        "ecs"
432    }
433
434    async fn handle(&self, request: AwsRequest) -> Result<AwsResponse, AwsServiceError> {
435        let mutates = is_mutating(request.action.as_str());
436        let result = match request.action.as_str() {
437            "CreateCluster" => self.create_cluster(&request),
438            "DescribeClusters" => self.describe_clusters(&request),
439            "DeleteCluster" => self.delete_cluster(&request),
440            "ListClusters" => self.list_clusters(&request),
441            "UpdateCluster" => self.update_cluster(&request),
442            "UpdateClusterSettings" => self.update_cluster_settings(&request),
443            "PutClusterCapacityProviders" => self.put_cluster_capacity_providers(&request),
444            "RegisterTaskDefinition" => self.register_task_definition(&request),
445            "DescribeTaskDefinition" => self.describe_task_definition(&request),
446            "DeregisterTaskDefinition" => self.deregister_task_definition(&request),
447            "DeleteTaskDefinitions" => self.delete_task_definitions(&request),
448            "ListTaskDefinitions" => self.list_task_definitions(&request),
449            "ListTaskDefinitionFamilies" => self.list_task_definition_families(&request),
450            "TagResource" => self.tag_resource(&request),
451            "UntagResource" => self.untag_resource(&request),
452            "ListTagsForResource" => self.list_tags_for_resource(&request),
453            "PutAccountSetting" => self.put_account_setting(&request),
454            "PutAccountSettingDefault" => self.put_account_setting_default(&request),
455            "DeleteAccountSetting" => self.delete_account_setting(&request),
456            "ListAccountSettings" => self.list_account_settings(&request),
457            "RunTask" => self.run_task(&request),
458            "StartTask" => self.start_task(&request),
459            "StopTask" => self.stop_task(&request).await,
460            "DescribeTasks" => self.describe_tasks(&request),
461            "ListTasks" => self.list_tasks(&request),
462            "CreateService" => self.create_service(&request),
463            "UpdateService" => self.update_service(&request),
464            "DeleteService" => self.delete_service(&request).await,
465            "DescribeServices" => self.describe_services(&request),
466            "ListServices" => self.list_services(&request),
467            "ListServicesByNamespace" => self.list_services_by_namespace(&request),
468            "RegisterContainerInstance" => self.register_container_instance(&request),
469            "DeregisterContainerInstance" => self.deregister_container_instance(&request),
470            "DescribeContainerInstances" => self.describe_container_instances(&request),
471            "ListContainerInstances" => self.list_container_instances(&request),
472            "UpdateContainerAgent" => self.update_container_agent(&request),
473            "UpdateContainerInstancesState" => self.update_container_instances_state(&request),
474            "PutAttributes" => self.put_attributes(&request),
475            "DeleteAttributes" => self.delete_attributes(&request),
476            "ListAttributes" => self.list_attributes(&request),
477            "CreateCapacityProvider" => self.create_capacity_provider(&request),
478            "DeleteCapacityProvider" => self.delete_capacity_provider(&request),
479            "DescribeCapacityProviders" => self.describe_capacity_providers(&request),
480            "UpdateCapacityProvider" => self.update_capacity_provider(&request),
481            "GetTaskProtection" => self.get_task_protection(&request),
482            "UpdateTaskProtection" => self.update_task_protection(&request),
483            "CreateTaskSet" => self.create_task_set(&request),
484            "UpdateTaskSet" => self.update_task_set(&request),
485            "DeleteTaskSet" => self.delete_task_set(&request),
486            "DescribeTaskSets" => self.describe_task_sets(&request),
487            "UpdateServicePrimaryTaskSet" => self.update_service_primary_task_set(&request),
488            "ExecuteCommand" => self.execute_command(&request).await,
489            "SubmitContainerStateChange" => self.submit_container_state_change(&request),
490            "SubmitTaskStateChange" => self.submit_task_state_change(&request),
491            "SubmitAttachmentStateChanges" => self.submit_attachment_state_changes(&request),
492            "DiscoverPollEndpoint" => self.discover_poll_endpoint(&request),
493            "StopServiceDeployment" => self.stop_service_deployment(&request),
494            "ContinueServiceDeployment" => self.continue_service_deployment(&request),
495            "ListServiceDeployments" => self.list_service_deployments(&request),
496            "DescribeServiceDeployments" => self.describe_service_deployments(&request),
497            "DescribeServiceRevisions" => self.describe_service_revisions(&request),
498            "RegisterDaemonTaskDefinition" => self.register_daemon_task_definition(&request),
499            "DescribeDaemonTaskDefinition" => self.describe_daemon_task_definition(&request),
500            "DeleteDaemonTaskDefinition" => self.delete_daemon_task_definition(&request),
501            "ListDaemonTaskDefinitions" => self.list_daemon_task_definitions(&request),
502            "CreateDaemon" => self.create_daemon(&request),
503            "DescribeDaemon" => self.describe_daemon(&request),
504            "UpdateDaemon" => self.update_daemon(&request),
505            "DeleteDaemon" => self.delete_daemon(&request),
506            "ListDaemons" => self.list_daemons(&request),
507            "DescribeDaemonDeployments" => self.describe_daemon_deployments(&request),
508            "ListDaemonDeployments" => self.list_daemon_deployments(&request),
509            "DescribeDaemonRevisions" => self.describe_daemon_revisions(&request),
510            "CreateExpressGatewayService" => self.create_express_gateway_service(&request),
511            "DescribeExpressGatewayService" => self.describe_express_gateway_service(&request),
512            "UpdateExpressGatewayService" => self.update_express_gateway_service(&request),
513            "DeleteExpressGatewayService" => self.delete_express_gateway_service(&request),
514            _ => Err(AwsServiceError::action_not_implemented(
515                "ecs",
516                &request.action,
517            )),
518        };
519        if mutates && matches!(result.as_ref(), Ok(resp) if resp.status.is_success()) {
520            self.save_snapshot().await;
521        }
522        result
523    }
524
525    fn supported_actions(&self) -> &[&str] {
526        SUPPORTED_ACTIONS
527    }
528}
529
530// -------- helpers --------
531
532// -------- operations: clusters --------
533
534// -------- operations: task definitions --------
535
536// -------- operations: tagging --------
537
538// -------- operations: account settings --------
539
540// -------- operations: tasks --------
541
542// -------- operations: services --------
543
544// -------- operations: container instances, attributes, capacity providers, task sets, task protection, ExecuteCommand, agent surface --------
545
546#[path = "service_clusters.rs"]
547mod service_clusters;
548
549#[path = "service_task_definitions.rs"]
550mod service_task_definitions;
551
552#[path = "service_tagging.rs"]
553mod service_tagging;
554
555#[path = "service_account_settings.rs"]
556mod service_account_settings;
557
558#[path = "service_tasks.rs"]
559mod service_tasks;
560
561#[path = "service_services_resource.rs"]
562mod service_services_resource;
563
564#[path = "service_container_instances_etc.rs"]
565mod service_container_instances_etc;
566
567#[path = "service_daemons.rs"]
568mod service_daemons;
569
570#[path = "service_express_gateway.rs"]
571mod service_express_gateway;
572
573#[path = "helpers.rs"]
574mod helpers;
575pub(crate) use helpers::*;
576
577#[cfg(test)]
578#[path = "service_tests.rs"]
579mod tests;