diff --git a/docs/services/ecs.md b/docs/services/ecs.md index ad9d563c462..c5c045e64b6 100644 --- a/docs/services/ecs.md +++ b/docs/services/ecs.md @@ -90,6 +90,11 @@ the agent's reason, `Unable to download firelens s3 config file: unable to downl from bucket : `, instead of leaking a started router. Shared network namespaces (AppConfig agent on `127.0.0.1:2772`) are not implemented. +Container `volumesFrom` entries are also stored and returned. In Docker mode, source containers +are launched before their consumers and their declared volumes are inherited with the requested +read-only or read-write access mode. Startup ordering also respects FireLens router dependencies; +cycles involving both volume inheritance and log routing are rejected before containers start. + ### Tasks | Operation | Description | diff --git a/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerBuilder.java b/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerBuilder.java index facd7da2949..7d1abdb989a 100644 --- a/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerBuilder.java +++ b/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerBuilder.java @@ -9,6 +9,7 @@ import com.github.dockerjava.api.model.Mount; import com.github.dockerjava.api.model.MountType; import com.github.dockerjava.api.model.Volume; +import com.github.dockerjava.api.model.VolumesFrom; import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; @@ -132,6 +133,7 @@ public static class Builder { private String networkMode; private final List mounts = new ArrayList<>(); private final List binds = new ArrayList<>(); + private final List volumesFrom = new ArrayList<>(); private final List extraHosts = new ArrayList<>(); private final Map labels = new HashMap<>(); private LogConfig logConfig; @@ -319,6 +321,14 @@ public Builder withNamedVolume(String volumeName, String containerPath, boolean return this; } + /** + * Inherits every volume declared by another container. + */ + public Builder withVolumesFrom(String sourceContainerId, boolean readOnly) { + volumesFrom.add(new VolumesFrom(sourceContainerId, readOnly ? AccessMode.ro : AccessMode.rw)); + return this; + } + /** * Adds a mount (any type: volume, bind, tmpfs). */ @@ -592,6 +602,7 @@ public ContainerSpec build() { networkMode, List.copyOf(mounts), List.copyOf(binds), + List.copyOf(volumesFrom), List.copyOf(extraHosts), Map.copyOf(labels), logConfig, diff --git a/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManager.java b/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManager.java index 745f932c997..20e5ca7dab5 100644 --- a/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManager.java +++ b/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManager.java @@ -954,6 +954,10 @@ private HostConfig buildHostConfig(ContainerSpec spec) { hostConfig.withBinds(spec.binds().toArray(new Bind[0])); } + if (spec.volumesFrom() != null && !spec.volumesFrom().isEmpty()) { + hostConfig.withVolumesFrom(spec.volumesFrom()); + } + // Docker rejects extra_hosts together with container: network mode. Containers // sharing another container's network namespace already inherit its network path. if (spec.extraHosts() != null && !spec.extraHosts().isEmpty() diff --git a/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerSpec.java b/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerSpec.java index 15b3ba1bcfd..7e59fc28465 100644 --- a/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerSpec.java +++ b/src/main/java/io/github/hectorvent/floci/core/common/docker/ContainerSpec.java @@ -4,6 +4,7 @@ import com.github.dockerjava.api.model.DeviceRequest; import com.github.dockerjava.api.model.LogConfig; import com.github.dockerjava.api.model.Mount; +import com.github.dockerjava.api.model.VolumesFrom; import java.util.List; import java.util.Map; @@ -24,6 +25,7 @@ * @param networkMode Docker network name or mode (null = default bridge) * @param mounts Volume mounts (named volumes, bind mounts, tmpfs) * @param binds Legacy bind mounts (prefer mounts for new code) + * @param volumesFrom Volumes inherited from other containers * @param extraHosts Extra /etc/hosts entries as "hostname:ip" strings * @param labels Container labels merged over the default floci-aws labels * @param logConfig Docker log driver configuration (null = daemon default) @@ -48,6 +50,7 @@ public record ContainerSpec( String networkMode, List mounts, List binds, + List volumesFrom, List extraHosts, Map labels, LogConfig logConfig, @@ -64,11 +67,14 @@ public record ContainerSpec( * All other fields will be null or empty lists. */ public ContainerSpec(String image) { - this(image, null, List.of(), null, null, null, Map.of(), List.of(), List.of(), null, List.of(), List.of(), List.of(), Map.of(), null, false, null, List.of(), null, null, List.of(), List.of()); + this(image, null, List.of(), null, null, null, Map.of(), List.of(), List.of(), null, + List.of(), List.of(), List.of(), List.of(), Map.of(), null, false, null, List.of(), + null, null, List.of(), List.of()); } /** - * Backward-compatible constructor that defaults {@code loopbackPortBindings} to empty. + * Backward-compatible constructor that defaults {@code loopbackPortBindings}, + * {@code volumesFrom}, and {@code deviceRequests} to empty. */ public ContainerSpec( String image, @@ -92,12 +98,15 @@ public ContainerSpec( String user, List groupAdd ) { - this(image, name, env, cmd, entrypoint, memoryBytes, portBindings, List.of(), exposedPorts, networkMode, mounts, binds, extraHosts, labels, logConfig, privileged, cgroupnsMode, dnsServers, workingDir, user, groupAdd, List.of()); + this(image, name, env, cmd, entrypoint, memoryBytes, portBindings, List.of(), exposedPorts, + networkMode, mounts, binds, List.of(), extraHosts, labels, logConfig, privileged, + cgroupnsMode, dnsServers, workingDir, user, groupAdd, List.of()); } /** - * Backward-compatible constructor that defaults {@code deviceRequests} to empty, so a - * caller that predates accelerator support keeps building CPU-only containers. + * Backward-compatible constructor that defaults {@code volumesFrom} and + * {@code deviceRequests} to empty, so callers that predate volume inheritance and + * accelerator support keep their existing behaviour. */ public ContainerSpec( String image, @@ -122,7 +131,42 @@ public ContainerSpec( String user, List groupAdd ) { - this(image, name, env, cmd, entrypoint, memoryBytes, portBindings, loopbackPortBindings, exposedPorts, networkMode, mounts, binds, extraHosts, labels, logConfig, privileged, cgroupnsMode, dnsServers, workingDir, user, groupAdd, List.of()); + this(image, name, env, cmd, entrypoint, memoryBytes, portBindings, loopbackPortBindings, + exposedPorts, networkMode, mounts, binds, List.of(), extraHosts, labels, logConfig, + privileged, cgroupnsMode, dnsServers, workingDir, user, groupAdd, List.of()); + } + + /** + * Backward-compatible constructor that defaults {@code volumesFrom} to empty while + * preserving explicitly requested devices. + */ + public ContainerSpec( + String image, + String name, + List env, + List cmd, + List entrypoint, + Long memoryBytes, + Map portBindings, + List loopbackPortBindings, + List exposedPorts, + String networkMode, + List mounts, + List binds, + List extraHosts, + Map labels, + LogConfig logConfig, + boolean privileged, + String cgroupnsMode, + List dnsServers, + String workingDir, + String user, + List groupAdd, + List deviceRequests + ) { + this(image, name, env, cmd, entrypoint, memoryBytes, portBindings, loopbackPortBindings, + exposedPorts, networkMode, mounts, binds, List.of(), extraHosts, labels, logConfig, + privileged, cgroupnsMode, dnsServers, workingDir, user, groupAdd, deviceRequests); } /** diff --git a/src/main/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandler.java b/src/main/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandler.java index b036f567030..df212e9db35 100644 --- a/src/main/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandler.java +++ b/src/main/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandler.java @@ -36,6 +36,7 @@ import io.github.hectorvent.floci.services.ecs.model.TaskSet; import io.github.hectorvent.floci.services.ecs.model.EfsVolumeConfiguration; import io.github.hectorvent.floci.services.ecs.model.Volume; +import io.github.hectorvent.floci.services.ecs.model.VolumeFrom; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; @@ -1102,6 +1103,17 @@ private ObjectNode containerDefinitionNode(ContainerDefinition def) { n.set("mountPoints", mps); } + if (def.getVolumesFrom() != null && !def.getVolumesFrom().isEmpty()) { + ArrayNode volumesFrom = objectMapper.createArrayNode(); + for (VolumeFrom volumeFrom : def.getVolumesFrom()) { + ObjectNode volumeFromNode = objectMapper.createObjectNode(); + volumeFromNode.put("sourceContainer", volumeFrom.sourceContainer()); + volumeFromNode.put("readOnly", volumeFrom.readOnly()); + volumesFrom.add(volumeFromNode); + } + n.set("volumesFrom", volumesFrom); + } + if (def.getLogConfiguration() != null) { LogConfiguration logConfig = def.getLogConfiguration(); ObjectNode logNode = objectMapper.createObjectNode(); @@ -1496,6 +1508,7 @@ private List parseContainerDefinitions(JsonNode node) { def.setSecrets(parseSecrets(item.path("secrets"))); } def.setMountPoints(parseMountPoints(item.path("mountPoints"))); + def.setVolumesFrom(parseVolumesFrom(item.path("volumesFrom"))); def.setLogConfiguration(parseLogConfiguration(item.path("logConfiguration"))); def.setFirelensConfiguration(parseFirelensConfiguration( item.path("firelensConfiguration"), result.size() + 1)); @@ -1555,6 +1568,19 @@ private List parseSecrets(JsonNode node) { return result; } + private List parseVolumesFrom(JsonNode node) { + List result = new ArrayList<>(); + if (!node.isArray()) { + return result; + } + for (JsonNode item : node) { + result.add(new VolumeFrom( + item.path("sourceContainer").asText(), + item.path("readOnly").asBoolean(false))); + } + return result; + } + private RuntimePlatform parseRuntimePlatform(JsonNode node) { if (node == null || !node.isObject()) { return null; diff --git a/src/main/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManager.java b/src/main/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManager.java index 2a9c0e25c93..cee11dbaa2d 100644 --- a/src/main/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManager.java +++ b/src/main/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManager.java @@ -34,6 +34,7 @@ import io.github.hectorvent.floci.services.ecs.model.Secret; import io.github.hectorvent.floci.services.ecs.model.TaskDefinition; import io.github.hectorvent.floci.services.ecs.model.Volume; +import io.github.hectorvent.floci.services.ecs.model.VolumeFrom; import io.github.hectorvent.floci.services.s3.S3Service; import io.github.hectorvent.floci.services.s3.model.S3Object; import io.github.hectorvent.floci.services.secretsmanager.SecretsManagerService; @@ -150,6 +151,16 @@ public EcsTaskHandle startTask(EcsTask task, TaskDefinition taskDef, Map containerIds = new LinkedHashMap<>(); Map logStreamsByContainerId = new LinkedHashMap<>(); List runtimeContainers = new ArrayList<>(); + ContainerDefinition firelensRouter = findFirelensRouter(taskDef.getContainerDefinitions()); + Map> firelensLogOptions = + awsFirelensLogOptions(taskDef.getContainerDefinitions()); + if (!firelensLogOptions.isEmpty() && firelensRouter == null) { + throw new AwsException("ClientException", + "awsfirelens log driver requires a firelensConfiguration container", 400); + } + + List launchOrder = orderForDependencies( + launchOrder(taskDef.getContainerDefinitions(), firelensRouter), firelensRouter); // Task-level volumes consumed by per-container mountPoints: host volumes map their // name -> absolute host source path; efsVolumeConfiguration volumes map their @@ -173,19 +184,11 @@ public EcsTaskHandle startTask(EcsTask task, TaskDefinition taskDef, Map> envVarsByContainer = new LinkedHashMap<>(); // Resolved before any container is created, so a registry-startup failure can't leak one already started. Map imagesByContainer = new LinkedHashMap<>(); - for (ContainerDefinition def : taskDef.getContainerDefinitions()) { + for (ContainerDefinition def : launchOrder) { envVarsByContainer.put(def, buildEnvVars(def, overridesByName.get(def.getName()), region)); imagesByContainer.put(def, ecrRegistryManager.rewriteImageUri(def.getImage())); } - ContainerDefinition firelensRouter = findFirelensRouter(taskDef.getContainerDefinitions()); - Map> firelensLogOptions = - awsFirelensLogOptions(taskDef.getContainerDefinitions()); - if (!firelensLogOptions.isEmpty() && firelensRouter == null) { - throw new AwsException("ClientException", - "awsfirelens log driver requires a firelensConfiguration container", 400); - } - PreparedNetwork protectedNetwork = prepareNetwork(task, taskDef, region, taskId); String firelensVolumeName = null; @@ -200,8 +203,6 @@ public EcsTaskHandle startTask(EcsTask task, TaskDefinition taskDef, firelensSocketAddress = unixSocketAddress(firelensVolumeName); } - List launchOrder = - launchOrder(taskDef.getContainerDefinitions(), firelensRouter); String fluentHost = null; String networkModeName = taskDef.getNetworkMode() != null ? taskDef.getNetworkMode().name() @@ -323,6 +324,17 @@ public EcsTaskHandle startTask(EcsTask task, TaskDefinition taskDef, } } + if (def.getVolumesFrom() != null) { + for (VolumeFrom volumeFrom : def.getVolumesFrom()) { + String sourceContainerId = containerIds.get(volumeFrom.sourceContainer()); + if (sourceContainerId == null) { + throw new IllegalStateException("ECS volumesFrom source container " + + volumeFrom.sourceContainer() + " has not started"); + } + specBuilder.withVolumesFrom(sourceContainerId, volumeFrom.readOnly()); + } + } + ContainerSpec spec = specBuilder.build(); String dockerId; @@ -379,7 +391,18 @@ public EcsTaskHandle startTask(EcsTask task, TaskDefinition taskDef, throw e; } - task.setContainers(runtimeContainers); + Map runtimeContainersByName = new LinkedHashMap<>(); + for (Container container : runtimeContainers) { + runtimeContainersByName.put(container.getName(), container); + } + List containersInDefinitionOrder = new ArrayList<>(); + for (ContainerDefinition definition : taskDef.getContainerDefinitions()) { + Container container = runtimeContainersByName.get(definition.getName()); + if (container != null) { + containersInDefinitionOrder.add(container); + } + } + task.setContainers(containersInDefinitionOrder); task.setLastStatus(TaskStatus.RUNNING.name()); task.setDesiredStatus(TaskStatus.RUNNING.name()); task.setStartedAt(Instant.now()); @@ -389,6 +412,53 @@ public EcsTaskHandle startTask(EcsTask task, TaskDefinition taskDef, protectedNetwork == null ? null : protectedNetwork.eni().getNetworkInterfaceId(), region); } + private List orderForDependencies(List definitions, + ContainerDefinition firelensRouter) { + Map definitionsByName = new LinkedHashMap<>(); + for (ContainerDefinition definition : definitions) { + definitionsByName.put(definition.getName(), definition); + } + + List ordered = new ArrayList<>(); + Set visiting = new HashSet<>(); + Set visited = new HashSet<>(); + for (ContainerDefinition definition : definitions) { + addAfterDependencies(definition, definitionsByName, firelensRouter, visiting, visited, ordered); + } + return ordered; + } + + private void addAfterDependencies(ContainerDefinition definition, + Map definitionsByName, + ContainerDefinition firelensRouter, + Set visiting, + Set visited, + List ordered) { + if (visited.contains(definition)) { + return; + } + if (!visiting.add(definition)) { + throw new IllegalArgumentException("ECS container dependencies contain a cycle at container " + + definition.getName()); + } + if (isAwsFirelens(definition)) { + addAfterDependencies(firelensRouter, definitionsByName, firelensRouter, visiting, visited, ordered); + } + if (definition.getVolumesFrom() != null) { + for (VolumeFrom volumeFrom : definition.getVolumesFrom()) { + ContainerDefinition source = definitionsByName.get(volumeFrom.sourceContainer()); + if (source == null) { + throw new IllegalArgumentException("ECS volumesFrom references unknown source container " + + volumeFrom.sourceContainer()); + } + addAfterDependencies(source, definitionsByName, firelensRouter, visiting, visited, ordered); + } + } + visiting.remove(definition); + visited.add(definition); + ordered.add(definition); + } + private PreparedNetwork prepareNetwork(EcsTask task, TaskDefinition definition, String region, String taskId) { if (definition.getNetworkMode() != NetworkMode.awsvpc || firewallManager == null || !firewallManager.enabled()) { diff --git a/src/main/java/io/github/hectorvent/floci/services/ecs/model/ContainerDefinition.java b/src/main/java/io/github/hectorvent/floci/services/ecs/model/ContainerDefinition.java index d5c7a66525d..3b30f13cb50 100644 --- a/src/main/java/io/github/hectorvent/floci/services/ecs/model/ContainerDefinition.java +++ b/src/main/java/io/github/hectorvent/floci/services/ecs/model/ContainerDefinition.java @@ -19,6 +19,7 @@ public class ContainerDefinition { private List command; private List entryPoint; private List mountPoints; + private List volumesFrom; private LogConfiguration logConfiguration; private FirelensConfiguration firelensConfiguration; private HealthCheck healthCheck; @@ -59,6 +60,9 @@ public class ContainerDefinition { public List getMountPoints() { return mountPoints; } public void setMountPoints(List mountPoints) { this.mountPoints = mountPoints; } + public List getVolumesFrom() { return volumesFrom; } + public void setVolumesFrom(List volumesFrom) { this.volumesFrom = volumesFrom; } + public LogConfiguration getLogConfiguration() { return logConfiguration; } public void setLogConfiguration(LogConfiguration logConfiguration) { this.logConfiguration = logConfiguration; } diff --git a/src/main/java/io/github/hectorvent/floci/services/ecs/model/VolumeFrom.java b/src/main/java/io/github/hectorvent/floci/services/ecs/model/VolumeFrom.java new file mode 100644 index 00000000000..7f8f7004588 --- /dev/null +++ b/src/main/java/io/github/hectorvent/floci/services/ecs/model/VolumeFrom.java @@ -0,0 +1,11 @@ +package io.github.hectorvent.floci.services.ecs.model; + +import io.quarkus.runtime.annotations.RegisterForReflection; + +/** + * An ECS container volume inheritance reference: + * {@code {"sourceContainer": ..., "readOnly": ...}}. + */ +@RegisterForReflection +public record VolumeFrom(String sourceContainer, boolean readOnly) { +} diff --git a/src/test/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManagerVolumesFromTest.java b/src/test/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManagerVolumesFromTest.java new file mode 100644 index 00000000000..a00aedffe94 --- /dev/null +++ b/src/test/java/io/github/hectorvent/floci/core/common/docker/ContainerLifecycleManagerVolumesFromTest.java @@ -0,0 +1,86 @@ +package io.github.hectorvent.floci.core.common.docker; + +import com.github.dockerjava.api.DockerClient; +import com.github.dockerjava.api.command.CreateContainerCmd; +import com.github.dockerjava.api.command.CreateContainerResponse; +import com.github.dockerjava.api.model.AccessMode; +import com.github.dockerjava.api.model.HostConfig; +import com.github.dockerjava.api.model.VolumesFrom; +import io.github.hectorvent.floci.config.EmulatorConfig; +import io.github.hectorvent.floci.core.common.dns.EmbeddedDnsServer; +import io.github.hectorvent.floci.services.lambda.launcher.ImageCacheService; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.RETURNS_SELF; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class ContainerLifecycleManagerVolumesFromTest { + + @Test + void createPassesInheritedVolumesToDockerHostConfig() { + EmulatorConfig config = mock(EmulatorConfig.class, RETURNS_DEEP_STUBS); + when(config.docker().resourceNamespace()).thenReturn(Optional.empty()); + ContainerBuilder builder = new ContainerBuilder( + config, mock(DockerHostResolver.class), mock(EmbeddedDnsServer.class)); + ContainerSpec spec = builder.newContainer("app:latest") + .withVolumesFrom("readonly-source-id", true) + .withVolumesFrom("readwrite-source-id", false) + .build(); + + HostConfig hostConfig = createHostConfig(config, spec); + VolumesFrom[] inherited = hostConfig.getVolumesFrom(); + assertEquals("readonly-source-id", inherited[0].getContainer()); + assertEquals(AccessMode.ro, inherited[0].getAccessMode()); + assertEquals("readwrite-source-id", inherited[1].getContainer()); + assertEquals(AccessMode.rw, inherited[1].getAccessMode()); + } + + @Test + void createPreservesInheritedVolumesWhenSharingRouterNetworkAndOmitsExtraHosts() { + EmulatorConfig config = mock(EmulatorConfig.class, RETURNS_DEEP_STUBS); + when(config.docker().resourceNamespace()).thenReturn(Optional.empty()); + ContainerBuilder builder = new ContainerBuilder( + config, mock(DockerHostResolver.class), mock(EmbeddedDnsServer.class)); + ContainerSpec spec = builder.newContainer("app:latest") + .withVolumesFrom("source-id", true) + .withNetworkMode("container:router-id") + .withExtraHost("host.docker.internal", "host-gateway") + .build(); + + HostConfig hostConfig = createHostConfig(config, spec); + + assertEquals("container:router-id", hostConfig.getNetworkMode()); + assertEquals(1, hostConfig.getVolumesFrom().length); + assertEquals("source-id", hostConfig.getVolumesFrom()[0].getContainer()); + assertEquals(AccessMode.ro, hostConfig.getVolumesFrom()[0].getAccessMode()); + assertTrue(hostConfig.getExtraHosts() == null || hostConfig.getExtraHosts().length == 0); + } + + private static HostConfig createHostConfig(EmulatorConfig config, ContainerSpec spec) { + DockerClient dockerClient = mock(DockerClient.class); + CreateContainerCmd createCmd = mock(CreateContainerCmd.class, RETURNS_SELF); + when(dockerClient.createContainerCmd("app:latest")).thenReturn(createCmd); + CreateContainerResponse response = mock(CreateContainerResponse.class); + when(response.getId()).thenReturn("app-id"); + when(createCmd.exec()).thenReturn(response); + ImageCacheService imageCacheService = mock(ImageCacheService.class); + when(imageCacheService.ensureImageExists(any())).thenAnswer(invocation -> invocation.getArgument(0)); + ContainerLifecycleManager manager = new ContainerLifecycleManager( + dockerClient, imageCacheService, mock(ContainerDetector.class), mock(PortAllocator.class), config); + + manager.create(spec); + + ArgumentCaptor hostConfig = ArgumentCaptor.forClass(HostConfig.class); + verify(createCmd).withHostConfig(hostConfig.capture()); + return hostConfig.getValue(); + } +} diff --git a/src/test/java/io/github/hectorvent/floci/services/ecs/EcsIntegrationTest.java b/src/test/java/io/github/hectorvent/floci/services/ecs/EcsIntegrationTest.java index 5b0504a6f7e..eab2f56b0d0 100644 --- a/src/test/java/io/github/hectorvent/floci/services/ecs/EcsIntegrationTest.java +++ b/src/test/java/io/github/hectorvent/floci/services/ecs/EcsIntegrationTest.java @@ -639,6 +639,48 @@ void describeTaskDefinitionWithoutRuntimePlatformOrLogConfigurationOmitsBothKeys .body("taskDefinition.containerDefinitions[0]", not(hasKey("logConfiguration"))); } + @Test + @Order(27) + void registerTaskDefinitionRoundTripsVolumesFrom() { + String volumesFromTaskDefArn = ecs("RegisterTaskDefinition") + .body(""" + { + "family": "task-with-volumes-from", + "containerDefinitions": [ + {"name": "source", "image": "sidecar:latest"}, + { + "name": "app", + "image": "app:latest", + "volumesFrom": [ + {"sourceContainer": "source", "readOnly": true} + ] + } + ] + } + """) + .when() + .post("/") + .then() + .statusCode(200) + .body("taskDefinition.containerDefinitions[1].volumesFrom", hasSize(1)) + .body("taskDefinition.containerDefinitions[1].volumesFrom[0].sourceContainer", equalTo("source")) + .body("taskDefinition.containerDefinitions[1].volumesFrom[0].readOnly", equalTo(true)) + .extract() + .path("taskDefinition.taskDefinitionArn"); + + ecs("DescribeTaskDefinition") + .body(""" + {"taskDefinition": "%s"} + """.formatted(volumesFromTaskDefArn)) + .when() + .post("/") + .then() + .statusCode(200) + .body("taskDefinition.containerDefinitions[1].volumesFrom", hasSize(1)) + .body("taskDefinition.containerDefinitions[1].volumesFrom[0].sourceContainer", equalTo("source")) + .body("taskDefinition.containerDefinitions[1].volumesFrom[0].readOnly", equalTo(true)); + } + // ── Services ────────────────────────────────────────────────────────────── @Test diff --git a/src/test/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandlerVolumesFromTest.java b/src/test/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandlerVolumesFromTest.java new file mode 100644 index 00000000000..56272ecd634 --- /dev/null +++ b/src/test/java/io/github/hectorvent/floci/services/ecs/EcsJsonHandlerVolumesFromTest.java @@ -0,0 +1,90 @@ +package io.github.hectorvent.floci.services.ecs; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import io.github.hectorvent.floci.config.EmulatorConfig; +import io.github.hectorvent.floci.services.ecs.container.HostVolumePolicy; +import io.github.hectorvent.floci.services.ecs.model.TaskDefinition; +import jakarta.ws.rs.core.Response; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class EcsJsonHandlerVolumesFromTest { + + private ObjectMapper objectMapper; + private EcsJsonHandler handler; + + @BeforeEach + void setUp() { + objectMapper = new ObjectMapper(); + EcsService service = mock(EcsService.class); + when(service.registerTaskDefinition(anyString(), any(), any(), any(), any(), any(), any(), any(), any(), anyString())) + .thenAnswer(invocation -> { + TaskDefinition taskDefinition = new TaskDefinition(); + taskDefinition.setFamily(invocation.getArgument(0)); + taskDefinition.setRevision(1); + taskDefinition.setStatus("ACTIVE"); + taskDefinition.setContainerDefinitions(invocation.getArgument(1, List.class)); + return taskDefinition; + }); + EmulatorConfig config = mock(EmulatorConfig.class, RETURNS_DEEP_STUBS); + handler = new EcsJsonHandler(service, objectMapper, new HostVolumePolicy(config)); + } + + @Test + void registerTaskDefinitionRoundTripsVolumesFrom() throws Exception { + JsonNode request = objectMapper.readTree(""" + { + "family": "shared-volume-family", + "containerDefinitions": [ + {"name": "source", "image": "sidecar:latest"}, + { + "name": "app", + "image": "app:latest", + "volumesFrom": [ + {"sourceContainer": "source", "readOnly": true}, + {"sourceContainer": "source"} + ] + } + ] + } + """); + + Response response = handler.handle("RegisterTaskDefinition", request, "us-east-1"); + JsonNode app = objectMapper.valueToTree(response.getEntity()) + .path("taskDefinition").path("containerDefinitions").get(1); + + assertEquals("source", app.path("volumesFrom").get(0).path("sourceContainer").asText()); + assertTrue(app.path("volumesFrom").get(0).path("readOnly").asBoolean()); + assertFalse(app.path("volumesFrom").get(1).path("readOnly").asBoolean()); + } + + @Test + void containerDefinitionOmitsVolumesFromWhenNotProvided() throws Exception { + JsonNode request = objectMapper.readTree(""" + { + "family": "plain-family", + "containerDefinitions": [ + {"name": "app", "image": "app:latest"} + ] + } + """); + + Response response = handler.handle("RegisterTaskDefinition", request, "us-east-1"); + JsonNode app = objectMapper.valueToTree(response.getEntity()) + .path("taskDefinition").path("containerDefinitions").get(0); + + assertTrue(app.path("volumesFrom").isMissingNode()); + } +} diff --git a/src/test/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManagerVolumesFromDockerIntegrationTest.java b/src/test/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManagerVolumesFromDockerIntegrationTest.java new file mode 100644 index 00000000000..772f4df71b7 --- /dev/null +++ b/src/test/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManagerVolumesFromDockerIntegrationTest.java @@ -0,0 +1,169 @@ +package io.github.hectorvent.floci.services.ecs.container; + +import com.github.dockerjava.api.DockerClient; +import com.github.dockerjava.api.async.ResultCallback; +import com.github.dockerjava.api.exception.NotFoundException; +import com.github.dockerjava.api.model.Frame; +import com.github.dockerjava.core.command.BuildImageResultCallback; +import com.github.dockerjava.core.command.WaitContainerResultCallback; +import io.github.hectorvent.floci.services.ecs.model.Container; +import io.github.hectorvent.floci.services.ecs.model.ContainerDefinition; +import io.github.hectorvent.floci.services.ecs.model.EcsTask; +import io.github.hectorvent.floci.services.ecs.model.TaskDefinition; +import io.github.hectorvent.floci.services.ecs.model.VolumeFrom; +import io.quarkus.test.junit.QuarkusTest; +import jakarta.inject.Inject; +import org.apache.commons.compress.archivers.tar.TarArchiveEntry; +import org.apache.commons.compress.archivers.tar.TarArchiveOutputStream; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Assumptions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.nio.charset.StandardCharsets; +import java.util.List; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Docker-backed regression coverage for ECS {@code volumesFrom}. The source image declares a + * volume containing the consumer's entrypoint, matching the failure reported in #3806: without + * Docker volume inheritance the app cannot even start because {@code /shared/wrapper} is absent. + */ +@QuarkusTest +class EcsContainerManagerVolumesFromDockerIntegrationTest { + + private static final String BUSYBOX_IMAGE = "public.ecr.aws/docker/library/busybox:latest"; + + @Inject + EcsContainerManager containerManager; + + @Inject + DockerClient dockerClient; + + private String sourceImage; + + @BeforeEach + void requireDocker() { + Assumptions.assumeTrue(isDockerAvailable(), + "Docker daemon must be available for ECS volumesFrom integration tests"); + } + + @AfterEach + void removeSourceImage() { + if (sourceImage != null) { + try { + dockerClient.removeImageCmd(sourceImage).withForce(true).exec(); + } catch (NotFoundException ignored) { + // A failed build may not have produced an image to remove. + } + } + } + + @Test + void appEntrypointCanComeFromSourceContainerVolume() throws Exception { + String suffix = UUID.randomUUID().toString().substring(0, 8); + String sourceImageTag = "floci-ecs-volumes-from-test:" + suffix; + buildSourceImage(sourceImageTag); + sourceImage = sourceImageTag; + + ContainerDefinition app = definition("app", BUSYBOX_IMAGE); + app.setEntryPoint(List.of("/shared/wrapper")); + app.setCommand(List.of("sh", "-c", "test -x /shared/wrapper && echo volumes-from-ok")); + app.setVolumesFrom(List.of(new VolumeFrom("source", true))); + + ContainerDefinition source = definition("source", sourceImage); + source.setCommand(List.of("true")); + + TaskDefinition taskDefinition = new TaskDefinition(); + taskDefinition.setFamily("volumes-from-docker-" + suffix); + // Put the consumer first to prove the manager derives launch order from volumesFrom, + // while the ECS response still preserves task-definition order. + taskDefinition.setContainerDefinitions(List.of(app, source)); + + EcsTask task = new EcsTask(); + task.setTaskArn("arn:aws:ecs:us-east-1:000000000000:task/volumes-from/" + suffix); + + EcsTaskHandle handle = containerManager.startTask(task, taskDefinition, List.of(), "us-east-1"); + try { + String appId = handle.getContainerIds().get("app"); + Integer status = dockerClient.waitContainerCmd(appId) + .exec(new WaitContainerResultCallback()) + .awaitStatusCode(60, TimeUnit.SECONDS); + String output = logs(appId); + + assertEquals(0, status, output); + assertTrue(output.contains("volumes-from-ok"), output); + assertEquals(List.of("app", "source"), + task.getContainers().stream().map(Container::getName).toList()); + } finally { + containerManager.stopTask(handle); + } + } + + private void buildSourceImage(String tag) throws Exception { + String dockerfile = """ + FROM public.ecr.aws/docker/library/busybox:latest + RUN mkdir -p /shared && printf '#!/bin/sh\\nexec "$@"\\n' > /shared/wrapper \ + && chmod +x /shared/wrapper + VOLUME ["/shared"] + """; + + byte[] dockerfileBytes = dockerfile.getBytes(StandardCharsets.UTF_8); + ByteArrayOutputStream context = new ByteArrayOutputStream(); + try (TarArchiveOutputStream tar = new TarArchiveOutputStream(context)) { + TarArchiveEntry entry = new TarArchiveEntry("Dockerfile"); + entry.setSize(dockerfileBytes.length); + tar.putArchiveEntry(entry); + tar.write(dockerfileBytes); + tar.closeArchiveEntry(); + } + + try (BuildImageResultCallback callback = new BuildImageResultCallback()) { + String imageId = dockerClient.buildImageCmd(new ByteArrayInputStream(context.toByteArray())) + .withTags(Set.of(tag)) + .withRemove(true) + .exec(callback) + .awaitImageId(180, TimeUnit.SECONDS); + assertFalse(imageId == null || imageId.isBlank(), "Docker must return the built image ID"); + } + } + + private String logs(String containerId) throws InterruptedException { + StringBuilder output = new StringBuilder(); + dockerClient.logContainerCmd(containerId) + .withStdOut(true) + .withStdErr(true) + .exec(new ResultCallback.Adapter() { + @Override + public void onNext(Frame frame) { + output.append(new String(frame.getPayload(), StandardCharsets.UTF_8)); + } + }) + .awaitCompletion(10, TimeUnit.SECONDS); + return output.toString(); + } + + private boolean isDockerAvailable() { + try { + dockerClient.pingCmd().exec(); + return true; + } catch (Exception e) { + return false; + } + } + + private static ContainerDefinition definition(String name, String image) { + ContainerDefinition definition = new ContainerDefinition(); + definition.setName(name); + definition.setImage(image); + return definition; + } +} diff --git a/src/test/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManagerVolumesFromTest.java b/src/test/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManagerVolumesFromTest.java new file mode 100644 index 00000000000..0f840c6dbc0 --- /dev/null +++ b/src/test/java/io/github/hectorvent/floci/services/ecs/container/EcsContainerManagerVolumesFromTest.java @@ -0,0 +1,280 @@ +package io.github.hectorvent.floci.services.ecs.container; + +import com.github.dockerjava.api.DockerClient; +import com.github.dockerjava.api.command.CopyArchiveToContainerCmd; +import com.github.dockerjava.api.command.InspectVolumeCmd; +import com.github.dockerjava.api.command.InspectVolumeResponse; +import com.github.dockerjava.api.model.LogConfig; +import io.github.hectorvent.floci.config.EmulatorConfig; +import io.github.hectorvent.floci.core.common.AwsException; +import io.github.hectorvent.floci.core.common.RegionResolver; +import io.github.hectorvent.floci.core.common.docker.ContainerBuilder; +import io.github.hectorvent.floci.core.common.docker.ContainerDetector; +import io.github.hectorvent.floci.core.common.docker.ContainerLifecycleManager; +import io.github.hectorvent.floci.core.common.docker.ContainerLifecycleManager.ContainerInfo; +import io.github.hectorvent.floci.core.common.docker.ContainerLogStreamer; +import io.github.hectorvent.floci.core.common.docker.ContainerSpec; +import io.github.hectorvent.floci.core.common.docker.LaunchedContainerAwsEnv; +import io.github.hectorvent.floci.services.ecr.registry.EcrRegistryManager; +import io.github.hectorvent.floci.services.ecs.model.Container; +import io.github.hectorvent.floci.services.ecs.model.ContainerDefinition; +import io.github.hectorvent.floci.services.ecs.model.EcsTask; +import io.github.hectorvent.floci.services.ecs.model.FirelensConfiguration; +import io.github.hectorvent.floci.services.ecs.model.LogConfiguration; +import io.github.hectorvent.floci.services.ecs.model.TaskDefinition; +import io.github.hectorvent.floci.services.ecs.model.VolumeFrom; +import io.github.hectorvent.floci.services.s3.S3Service; +import io.github.hectorvent.floci.services.secretsmanager.SecretsManagerService; +import io.github.hectorvent.floci.services.ssm.SsmService; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.InOrder; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.RETURNS_SELF; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class EcsContainerManagerVolumesFromTest { + + private ContainerBuilder containerBuilder; + private ContainerBuilder.Builder sourceBuilder; + private ContainerBuilder.Builder appBuilder; + private ContainerLifecycleManager lifecycleManager; + private ContainerLogStreamer logStreamer; + private S3Service s3Service; + private EcsContainerManager manager; + + @BeforeEach + void setUp() { + containerBuilder = mock(ContainerBuilder.class); + sourceBuilder = mock(ContainerBuilder.Builder.class, RETURNS_SELF); + appBuilder = mock(ContainerBuilder.Builder.class, RETURNS_SELF); + when(containerBuilder.newContainer("sidecar:latest")).thenReturn(sourceBuilder); + when(containerBuilder.newContainer("app:latest")).thenReturn(appBuilder); + when(sourceBuilder.build()).thenReturn(mock(ContainerSpec.class)); + when(appBuilder.build()).thenReturn(mock(ContainerSpec.class)); + + lifecycleManager = mock(ContainerLifecycleManager.class); + when(lifecycleManager.createAndStart(any())) + .thenReturn(new ContainerInfo("source-id", Map.of())) + .thenReturn(new ContainerInfo("app-id", Map.of())); + + logStreamer = mock(ContainerLogStreamer.class); + ContainerDetector containerDetector = mock(ContainerDetector.class); + EmulatorConfig config = mock(EmulatorConfig.class, RETURNS_DEEP_STUBS); + RegionResolver regionResolver = mock(RegionResolver.class); + LaunchedContainerAwsEnv awsEnv = mock(LaunchedContainerAwsEnv.class); + when(awsEnv.sdkBaselineEnv(any(), any())).thenReturn(List.of()); + when(awsEnv.flociEndpoint()).thenReturn("http://host.docker.internal:4566"); + EcrRegistryManager ecrRegistryManager = mock(EcrRegistryManager.class); + when(ecrRegistryManager.rewriteImageUri(anyString())).thenAnswer(invocation -> invocation.getArgument(0)); + s3Service = mock(S3Service.class); + + manager = new EcsContainerManager(containerBuilder, lifecycleManager, logStreamer, + containerDetector, config, regionResolver, awsEnv, mock(SsmService.class), + mock(SecretsManagerService.class), s3Service, ecrRegistryManager, mock(HostVolumePolicy.class)); + } + + @Test + void volumesFromStartsTheSourceFirstAndUsesItsDockerId() { + ContainerDefinition app = definition("app", "app:latest"); + app.setVolumesFrom(List.of(new VolumeFrom("source", true))); + ContainerDefinition source = definition("source", "sidecar:latest"); + + EcsTask ecsTask = task(); + manager.startTask(ecsTask, taskDefinition(List.of(app, source)), List.of(), "us-east-1"); + + InOrder order = inOrder(containerBuilder); + order.verify(containerBuilder).newContainer("sidecar:latest"); + order.verify(containerBuilder).newContainer("app:latest"); + verify(appBuilder).withVolumesFrom("source-id", true); + assertEquals( + List.of("app", "source"), + ecsTask.getContainers().stream().map(Container::getName).toList()); + } + + @Test + void volumesFromPreservesReadWriteMode() { + ContainerDefinition app = definition("app", "app:latest"); + app.setVolumesFrom(List.of(new VolumeFrom("source", false))); + ContainerDefinition source = definition("source", "sidecar:latest"); + + manager.startTask(task(), taskDefinition(List.of(source, app)), List.of(), "us-east-1"); + + verify(appBuilder).withVolumesFrom("source-id", false); + } + + @Test + void firelensRouterAndApplicationStartAfterTheirVolumeSource() { + ContainerBuilder.Builder routerBuilder = mock(ContainerBuilder.Builder.class, RETURNS_SELF); + when(containerBuilder.newContainer("router:latest")).thenReturn(routerBuilder); + when(routerBuilder.build()).thenReturn(mock(ContainerSpec.class)); + when(lifecycleManager.create(any())).thenReturn("router-id"); + when(lifecycleManager.startCreated(eq("router-id"), any())) + .thenReturn(new ContainerInfo("router-id", Map.of())); + DockerClient dockerClient = mock(DockerClient.class, RETURNS_DEEP_STUBS); + when(lifecycleManager.getDockerClient()).thenReturn(dockerClient); + InspectVolumeCmd inspectVolumeCmd = mock(InspectVolumeCmd.class); + InspectVolumeResponse volume = mock(InspectVolumeResponse.class); + when(dockerClient.inspectVolumeCmd("floci-ecs-firelens-volumesfrom1")).thenReturn(inspectVolumeCmd); + when(inspectVolumeCmd.exec()).thenReturn(volume); + when(volume.getMountpoint()).thenReturn("/var/lib/docker/volumes/floci-ecs-firelens-volumesfrom1/_data"); + CopyArchiveToContainerCmd copyCmd = mock(CopyArchiveToContainerCmd.class, RETURNS_SELF); + when(dockerClient.copyArchiveToContainerCmd("router-id")).thenReturn(copyCmd); + + ContainerDefinition app = definition("app", "app:latest"); + app.setLogConfiguration(new LogConfiguration("awsfirelens", Map.of("Name", "stdout"), null)); + app.setVolumesFrom(List.of(new VolumeFrom("source", true))); + ContainerDefinition router = definition("router", "router:latest"); + router.setFirelensConfiguration(new FirelensConfiguration("fluentbit", Map.of())); + router.setVolumesFrom(List.of(new VolumeFrom("source", false))); + ContainerDefinition source = definition("source", "sidecar:latest"); + + EcsTask ecsTask = task(); + EcsTaskHandle handle = manager.startTask( + ecsTask, taskDefinition(List.of(app, router, source)), List.of(), "us-east-1"); + + InOrder order = inOrder(containerBuilder); + order.verify(containerBuilder).newContainer("sidecar:latest"); + order.verify(containerBuilder).newContainer("router:latest"); + order.verify(containerBuilder).newContainer("app:latest"); + InOrder lifecycleOrder = inOrder(lifecycleManager, copyCmd); + lifecycleOrder.verify(lifecycleManager).createAndStart(any(ContainerSpec.class)); + lifecycleOrder.verify(lifecycleManager).create(any(ContainerSpec.class)); + lifecycleOrder.verify(copyCmd).exec(); + lifecycleOrder.verify(lifecycleManager).startCreated(eq("router-id"), any(ContainerSpec.class)); + lifecycleOrder.verify(lifecycleManager).createAndStart(any(ContainerSpec.class)); + verify(routerBuilder).withVolumesFrom("source-id", false); + verify(appBuilder).withVolumesFrom("source-id", true); + verify(lifecycleManager).ensureVolume("floci-ecs-firelens-volumesfrom1"); + verify(routerBuilder).withNamedVolume("floci-ecs-firelens-volumesfrom1", "/var/run"); + verify(copyCmd).withRemotePath("/fluent-bit/etc"); + verify(routerBuilder, never()).withLoopbackPortBinding(anyInt(), anyInt()); + verify(appBuilder, never()).withNetworkMode(anyString()); + + ArgumentCaptor logConfig = ArgumentCaptor.forClass(LogConfig.class); + verify(appBuilder).withLogConfig(logConfig.capture()); + assertEquals(LogConfig.LoggingType.FLUENTD, logConfig.getValue().getType()); + assertEquals("unix:///var/lib/docker/volumes/floci-ecs-firelens-volumesfrom1/_data/fluent.sock", + logConfig.getValue().getConfig().get("fluentd-address")); + assertEquals("app-firelens-volumesfrom1", logConfig.getValue().getConfig().get("tag")); + verify(appBuilder, never()).withLogRotation(); + verify(sourceBuilder).withLogRotation(); + verify(routerBuilder).withLogRotation(); + verify(logStreamer, never()).attach(eq("app-id"), any(), any(), any(), any()); + verify(logStreamer).attach(eq("source-id"), any(), any(), eq("us-east-1"), any()); + verify(logStreamer).attach(eq("router-id"), any(), any(), eq("us-east-1"), any()); + assertEquals("floci-ecs-firelens-volumesfrom1", handle.getFirelensVolumeName()); + assertEquals(List.of("source", "router", "app"), handle.getContainerIds().keySet().stream().toList()); + assertEquals(List.of("app", "router", "source"), + ecsTask.getContainers().stream().map(Container::getName).toList()); + } + + @Test + void firelensRouterDependingOnItsLoggingApplicationFailsBeforeCreatingContainers() { + ContainerDefinition app = definition("app", "app:latest"); + app.setLogConfiguration(new LogConfiguration("awsfirelens", Map.of(), null)); + ContainerDefinition router = definition("router", "router:latest"); + router.setFirelensConfiguration(new FirelensConfiguration("fluentbit", Map.of())); + router.setVolumesFrom(List.of(new VolumeFrom("app", false))); + ContainerDefinition source = definition("source", "sidecar:latest"); + + assertThrows(IllegalArgumentException.class, + () -> manager.startTask(task(), taskDefinition(List.of(source, router, app)), + List.of(), "us-east-1")); + + verify(containerBuilder, never()).newContainer(anyString()); + verify(lifecycleManager, never()).create(any()); + verify(lifecycleManager, never()).createAndStart(any()); + verify(lifecycleManager, never()).ensureVolume(anyString()); + } + + @Test + void missingFirelensS3ConfigFailsBeforeCreatingItsVolumeSourceOrSocketVolume() { + ContainerDefinition app = definition("app", "app:latest"); + app.setLogConfiguration(new LogConfiguration("awsfirelens", Map.of("Name", "stdout"), null)); + app.setVolumesFrom(List.of(new VolumeFrom("source", true))); + ContainerDefinition router = definition("router", "router:latest"); + router.setFirelensConfiguration(new FirelensConfiguration("fluentbit", Map.of( + "config-file-type", "s3", + "config-file-value", "arn:aws:s3:::firelens-configs/missing.conf"))); + router.setVolumesFrom(List.of(new VolumeFrom("source", false))); + ContainerDefinition source = definition("source", "sidecar:latest"); + when(s3Service.getObject("firelens-configs", "missing.conf")) + .thenThrow(new AwsException("NoSuchKey", "The specified key does not exist.", 404)); + + AwsException failure = assertThrows(AwsException.class, + () -> manager.startTask(task(), taskDefinition(List.of(app, router, source)), + List.of(), "us-east-1")); + + assertEquals("ResourceInitializationError", failure.getErrorCode()); + assertEquals("Unable to download firelens s3 config file: unable to download s3 config " + + "missing.conf from bucket firelens-configs: The specified key does not exist.", + failure.getMessage()); + verify(s3Service).getObject("firelens-configs", "missing.conf"); + verify(containerBuilder, never()).newContainer(anyString()); + verify(lifecycleManager, never()).create(any()); + verify(lifecycleManager, never()).createAndStart(any()); + verify(lifecycleManager, never()).ensureVolume(anyString()); + } + + @Test + void unknownVolumesFromSourceFailsBeforeCreatingContainers() { + ContainerDefinition app = definition("app", "app:latest"); + app.setVolumesFrom(List.of(new VolumeFrom("missing", false))); + + assertThrows(IllegalArgumentException.class, + () -> manager.startTask(task(), taskDefinition(List.of(app)), List.of(), "us-east-1")); + + verify(lifecycleManager, never()).createAndStart(any()); + } + + @Test + void cyclicVolumesFromFailsBeforeCreatingContainers() { + ContainerDefinition app = definition("app", "app:latest"); + app.setVolumesFrom(List.of(new VolumeFrom("source", false))); + ContainerDefinition source = definition("source", "sidecar:latest"); + source.setVolumesFrom(List.of(new VolumeFrom("app", false))); + + assertThrows(IllegalArgumentException.class, + () -> manager.startTask(task(), taskDefinition(List.of(app, source)), List.of(), "us-east-1")); + + verify(lifecycleManager, never()).createAndStart(any()); + } + + private static ContainerDefinition definition(String name, String image) { + ContainerDefinition definition = new ContainerDefinition(); + definition.setName(name); + definition.setImage(image); + return definition; + } + + private static TaskDefinition taskDefinition(List definitions) { + TaskDefinition taskDefinition = new TaskDefinition(); + taskDefinition.setFamily("volumes-from-family"); + taskDefinition.setRevision(1); + taskDefinition.setContainerDefinitions(definitions); + return taskDefinition; + } + + private static EcsTask task() { + EcsTask task = new EcsTask(); + task.setTaskArn("arn:aws:ecs:us-east-1:000000000000:task/test-cluster/volumesfrom1"); + task.setClusterArn("arn:aws:ecs:us-east-1:000000000000:cluster/test-cluster"); + return task; + } +}