diff --git a/crates/fakecloud-aws/src/arn.rs b/crates/fakecloud-aws/src/arn.rs index 21325ec2a..8b59b69ce 100644 --- a/crates/fakecloud-aws/src/arn.rs +++ b/crates/fakecloud-aws/src/arn.rs @@ -127,6 +127,13 @@ pub fn arn_resource<'a>(arn: &'a str, service: &str) -> Option<&'a str> { rest.strip_prefix(service)?.strip_prefix(':') } +/// The account an ARN names (`arn:::::...`), +/// `None` when `arn` is not an ARN or names no account. +pub fn account_of(arn: &str) -> Option<&str> { + let rest = arn.strip_prefix("arn:")?; + rest.split(':').nth(3).filter(|a| !a.is_empty()) +} + /// The partition an ARN names (`arn::...`), `aws` when it names /// none. pub fn partition_of(arn: &str) -> &str { @@ -338,4 +345,18 @@ mod tests { assert_eq!(partition_for("us-isof-south-1"), "aws-iso-f"); assert_eq!(partition_for("eu-isoe-west-1"), "aws-iso-e"); } + + #[test] + fn account_of_reads_the_account_field() { + assert_eq!( + account_of("arn:aws:iam::123456789012:role/r"), + Some("123456789012") + ); + assert_eq!( + account_of("arn:aws-cn:lambda:cn-north-1:000000000000:function:f:1"), + Some("000000000000") + ); + assert_eq!(account_of("arn:aws:s3:::bucket"), None); + assert_eq!(account_of("not-an-arn"), None); + } } diff --git a/crates/fakecloud-cloudformation/src/extras.rs b/crates/fakecloud-cloudformation/src/extras.rs index f79ecab8b..b69519e9e 100644 --- a/crates/fakecloud-cloudformation/src/extras.rs +++ b/crates/fakecloud-cloudformation/src/extras.rs @@ -3086,6 +3086,7 @@ pub(crate) mod tests { appconfig: shared::(), delivery: Arc::new(DeliveryBus::new()), lambda_runtime: None, + iam_mode: Default::default(), rds_runtime: None, ec2_runtime: None, ecs_runtime: None, diff --git a/crates/fakecloud-cloudformation/src/resource_provisioner/cloudformation.rs b/crates/fakecloud-cloudformation/src/resource_provisioner/cloudformation.rs index 5fda0f074..0b8c67e9f 100644 --- a/crates/fakecloud-cloudformation/src/resource_provisioner/cloudformation.rs +++ b/crates/fakecloud-cloudformation/src/resource_provisioner/cloudformation.rs @@ -128,6 +128,7 @@ impl ResourceProvisioner { cloudformation_state: self.cloudformation_state.clone(), delivery: self.delivery.clone(), lambda_runtime: self.lambda_runtime.clone(), + iam_mode: self.iam_mode, rds_runtime: self.rds_runtime.clone(), ec2_runtime: self.ec2_runtime.clone(), ecs_runtime: self.ecs_runtime.clone(), diff --git a/crates/fakecloud-cloudformation/src/resource_provisioner/lambda.rs b/crates/fakecloud-cloudformation/src/resource_provisioner/lambda.rs index 1e669fcda..d0790b20a 100644 --- a/crates/fakecloud-cloudformation/src/resource_provisioner/lambda.rs +++ b/crates/fakecloud-cloudformation/src/resource_provisioner/lambda.rs @@ -43,6 +43,7 @@ impl ResourceProvisioner { .unwrap_or_else(|| self.physical_name(resource)); let cfg = parse_lambda_function_props(props)?; + self.validate_lambda_execution_role(&cfg.role)?; let function_arn = fakecloud_lambda::function_arn(&self.region, &self.account_id, &function_name); @@ -169,6 +170,21 @@ impl ResourceProvisioner { .with("Version", "$LATEST")) } + /// The same execution-role checks `CreateFunction` runs: a role whose + /// trust policy lets Lambda assume it, in the stack's account under IAM + /// enforcement. + fn validate_lambda_execution_role(&self, role_arn: &str) -> Result<(), String> { + let validator = + fakecloud_iam::pass_role::IamRoleTrustValidator::new(self.iam_state.clone()); + fakecloud_lambda::validate_execution_role( + &self.account_id, + role_arn, + Some(&validator), + self.iam_mode, + ) + .map_err(|e| e.to_string()) + } + /// Apply a CFN template-driven update to an existing Lambda function. /// Mirrors `UpdateFunctionConfiguration` + `UpdateFunctionCode`: /// rewrite mutable configuration fields from the new template, re-hash @@ -184,6 +200,7 @@ impl ResourceProvisioner { let props = &resource.properties; let function_name = existing.physical_id.clone(); let cfg = parse_lambda_function_props(props)?; + self.validate_lambda_execution_role(&cfg.role)?; let new_code_zip = if cfg.code_zip.is_some() { cfg.code_zip.clone() diff --git a/crates/fakecloud-cloudformation/src/resource_provisioner/mod.rs b/crates/fakecloud-cloudformation/src/resource_provisioner/mod.rs index 8e4725079..33b47c393 100644 --- a/crates/fakecloud-cloudformation/src/resource_provisioner/mod.rs +++ b/crates/fakecloud-cloudformation/src/resource_provisioner/mod.rs @@ -524,6 +524,12 @@ fn parse_lambda_function_props(props: &serde_json::Value) -> Result>() }) .unwrap_or_default(); + if let Some(message) = + fakecloud_lambda::runtime::environment::reserved_keys_message(&environment) + { + // Same `: ` shape as the execution-role failures. + return Err(format!("InvalidParameterValueException: {message}")); + } // CFN tags ride as `[{Key, Value}, ...]`; flatten to the map shape // the lambda crate stores tags in. @@ -981,6 +987,8 @@ pub struct ResourceProvisioner { /// images (see `CloudFormationDeps::lambda_runtime`). `None` outside a /// configured runtime (e.g. unit tests). pub lambda_runtime: Option>, + /// IAM enforcement mode (see `CloudFormationDeps::iam_mode`). + pub iam_mode: fakecloud_core::auth::IamMode, /// Container runtimes for stateful services whose CFN-provisioned resources /// must be backed by REAL containers. See `CloudFormationDeps`. `None` /// (no Docker/Podman, e.g. CI/unit tests) keeps metadata-only provisioning. @@ -4268,6 +4276,7 @@ mod tests { )), delivery: Arc::new(DeliveryBus::new()), lambda_runtime: None, + iam_mode: Default::default(), rds_runtime: None, ec2_runtime: None, ecs_runtime: None, @@ -8110,7 +8119,14 @@ mod tests { "MyRole", serde_json::json!({ "RoleName": "my-role", - "AssumeRolePolicyDocument": {"Version": "2012-10-17", "Statement": []} + "AssumeRolePolicyDocument": { + "Version": "2012-10-17", + "Statement": [{ + "Effect": "Allow", + "Principal": {"Service": "lambda.amazonaws.com"}, + "Action": "sts:AssumeRole" + }] + } }), ); let role_sr = prov.create_resource(&role).unwrap(); @@ -8131,6 +8147,62 @@ mod tests { assert!(arn.contains(":function:my-fn")); } + #[test] + fn lambda_function_role_runs_the_create_function_checks() { + let prov = make_provisioner(); + let untrusted = make_resource( + "AWS::IAM::Role", + "Untrusted", + serde_json::json!({ + "RoleName": "untrusted", + "AssumeRolePolicyDocument": {"Version": "2012-10-17", "Statement": []} + }), + ); + let untrusted = prov.create_resource(&untrusted).unwrap(); + let func = |role: &str| { + make_resource( + "AWS::Lambda::Function", + "Fn", + serde_json::json!({ + "FunctionName": "checked-fn", + "Runtime": "python3.12", + "Handler": "index.handler", + "Role": role, + "Code": {"ZipFile": "def handler(e,c): return e"} + }), + ) + }; + let err = prov + .create_resource(&func(&untrusted.physical_id)) + .unwrap_err(); + assert!(err.contains("InvalidParameterValueException"), "{err}"); + + // Another account's role: accepted with IAM enforcement off (the + // default), refused like AWS with it on. + prov.create_resource(&func("arn:aws:iam::999999999999:role/elsewhere")) + .expect("default mode accepts another account's role"); + // Reserved environment keys fail with the Lambda error code too. + let mut reserved = func("arn:aws:iam::123456789012:role/r"); + reserved.properties["FunctionName"] = serde_json::json!("reserved-fn"); + reserved.properties["Environment"] = + serde_json::json!({"Variables": {"AWS_REGION": "eu-west-1"}}); + let err = prov.create_resource(&reserved).unwrap_err(); + assert!( + err.starts_with("InvalidParameterValueException: ") + && err.ends_with("Reserved keys used in this request: AWS_REGION"), + "{err}" + ); + let mut enforcing = make_provisioner(); + enforcing.iam_mode = fakecloud_core::auth::IamMode::Strict; + let err = enforcing + .create_resource(&func("arn:aws:iam::999999999999:role/elsewhere")) + .unwrap_err(); + assert!( + err.contains("AccessDeniedException") && err.contains("Cross-account pass role"), + "{err}" + ); + } + #[test] fn getatt_iam_role_arn_returns_role_arn() { let prov = make_provisioner(); diff --git a/crates/fakecloud-cloudformation/src/service.rs b/crates/fakecloud-cloudformation/src/service.rs index 8137a0b18..3f3bb4105 100644 --- a/crates/fakecloud-cloudformation/src/service.rs +++ b/crates/fakecloud-cloudformation/src/service.rs @@ -580,6 +580,9 @@ pub struct CloudFormationDeps { /// runtime is configured — provisioning still works, the first Invoke just /// falls back to a cold pull. pub lambda_runtime: Option>, + /// IAM enforcement mode. `AWS::Lambda::Function` only refuses a + /// cross-account execution role (as AWS does) while it is on. + pub iam_mode: fakecloud_core::auth::IamMode, /// Container runtimes for the stateful services whose CFN-provisioned /// resources must be backed by REAL containers (not phantom metadata). /// When present, the provisioner inserts the resource record synchronously @@ -1229,6 +1232,7 @@ impl CloudFormationService { cloudformation_state: self.state.clone(), delivery: self.deps.delivery.clone(), lambda_runtime: self.deps.lambda_runtime.clone(), + iam_mode: self.deps.iam_mode, rds_runtime: self.deps.rds_runtime.clone(), ec2_runtime: self.deps.ec2_runtime.clone(), ecs_runtime: self.deps.ecs_runtime.clone(), @@ -4422,6 +4426,7 @@ mod tests { appconfig: mas(), delivery: Arc::new(DeliveryBus::new()), lambda_runtime: None, + iam_mode: Default::default(), rds_runtime: None, ec2_runtime: None, ecs_runtime: None, diff --git a/crates/fakecloud-core/src/auth.rs b/crates/fakecloud-core/src/auth.rs index 3c503713f..37798355f 100644 --- a/crates/fakecloud-core/src/auth.rs +++ b/crates/fakecloud-core/src/auth.rs @@ -618,6 +618,40 @@ pub trait RoleTrustValidator: Send + Sync { ) -> Result<(), PassRoleError>; } +/// Temporary credentials for an assumed-role session, as a compute service +/// hands them to the code it runs (Lambda's execution-role environment, for +/// example). +#[derive(Clone, Debug)] +pub struct SessionCredentials { + pub access_key_id: String, + pub secret_access_key: String, + pub session_token: String, + pub expiration: DateTime, + /// Account the session is registered under, so it can be revoked there. + pub account_id: String, +} + +/// Issues assumed-role session credentials on behalf of a compute service, +/// registered so that requests signed with them resolve to +/// `arn::sts:::assumed-role//` (and verify +/// under `--verify-sigv4`). Implemented over IAM state; services that run +/// user code under a role take it as an optional hook so they stay decoupled +/// from the IAM crate. +pub trait SessionCredentialIssuer: Send + Sync { + /// Mint credentials for `role_arn` with the given session name, valid for + /// `duration`. + fn issue( + &self, + role_arn: &str, + session_name: &str, + duration: chrono::Duration, + ) -> SessionCredentials; + + /// Unregister credentials once the code they were issued to has stopped. + /// Idempotent. + fn revoke(&self, credentials: &SessionCredentials); +} + /// Composite [`ResourcePolicyProvider`] that delegates to a list of /// sub-providers in order, returning the first `Some` hit. /// diff --git a/crates/fakecloud-e2e/tests/lambda_aws_env.rs b/crates/fakecloud-e2e/tests/lambda_aws_env.rs new file mode 100644 index 000000000..817765d3b --- /dev/null +++ b/crates/fakecloud-e2e/tests/lambda_aws_env.rs @@ -0,0 +1,384 @@ +//! A Lambda's container gets the environment real Lambda provides, pointed +//! at fakecloud instead of real AWS. +//! +//! Real Lambda injects region, execution-role credentials and function +//! metadata, and function code relies on them. fakecloud injected none, and +//! nothing pointed the SDK at the emulator, so handler code that called AWS +//! silently targeted the internet: CDK's `BucketDeployment` reported success +//! having copied no files. + +mod helpers; + +use aws_sdk_lambda::primitives::Blob; +use aws_sdk_lambda::types::{Environment, FunctionCode, Runtime}; +use helpers::TestServer; + +fn docker_available() -> bool { + std::process::Command::new("docker") + .arg("info") + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status() + .map(|s| s.success()) + .unwrap_or(false) +} + +fn build_python_handler_zip(body: &str) -> Vec { + use std::io::Write; + let mut buf = Vec::new(); + { + let mut zip = zip::ZipWriter::new(std::io::Cursor::new(&mut buf)); + let opts: zip::write::FileOptions<'_, ()> = + zip::write::FileOptions::default().compression_method(zip::CompressionMethod::Stored); + zip.start_file("index.py", opts).unwrap(); + zip.write_all(body.as_bytes()).unwrap(); + zip.finish().unwrap(); + } + buf +} + +/// SDK config signed with the reserved root-bypass credentials, which skip +/// SigV4 verification and IAM enforcement. +async fn root_config(server: &TestServer) -> aws_config::SdkConfig { + aws_config::defaults(aws_config::BehaviorVersion::latest()) + .endpoint_url(server.endpoint()) + .region(aws_config::Region::new("us-east-1")) + .credentials_provider(aws_sdk_lambda::config::Credentials::new( + "test", "test", None, None, "test", + )) + .load() + .await +} + +/// The handler reports the environment it sees and the identity its SDK +/// calls resolve to. Signature verification is on, so the call only +/// succeeds if the injected credentials are registered with fakecloud. +const PROBE_HANDLER: &str = r#"import os +import boto3 + +KEYS = ( + "AWS_ENDPOINT_URL", "AWS_REGION", "AWS_DEFAULT_REGION", + "AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY", "AWS_SESSION_TOKEN", + "AWS_LAMBDA_FUNCTION_NAME", "AWS_LAMBDA_FUNCTION_VERSION", + "AWS_LAMBDA_FUNCTION_MEMORY_SIZE", "AWS_LAMBDA_LOG_GROUP_NAME", + "AWS_EXECUTION_ENV", "APP_SETTING", +) + +def handler(event, context): + out = {k: os.environ.get(k) for k in KEYS} + out["caller_arn"] = boto3.client("sts").get_caller_identity()["Arn"] + out["context_function_name"] = context.function_name + return out +"#; + +#[tokio::test] +async fn lambda_container_receives_the_aws_environment() { + if !docker_available() { + eprintln!("docker required for Lambda execution; skipping"); + return; + } + let server = TestServer::start_with_env(&[("FAKECLOUD_VERIFY_SIGV4", "true")]).await; + // The test itself drives the control plane with the root-bypass + // credentials; only the function's own calls are verified. + let root = root_config(&server).await; + let lambda = aws_sdk_lambda::Client::new(&root); + + lambda + .create_function() + .function_name("env-probe-fn") + .runtime(Runtime::Python312) + .role("arn:aws:iam::123456789012:role/env-probe-role") + .handler("index.handler") + .memory_size(256) + .timeout(30) + .environment( + Environment::builder() + .variables("APP_SETTING", "kept") + .build(), + ) + .code( + FunctionCode::builder() + .zip_file(build_python_handler_zip(PROBE_HANDLER).into()) + .build(), + ) + .send() + .await + .expect("create_function"); + + let invoked = lambda + .invoke() + .function_name("env-probe-fn") + .payload(Blob::new("{}")) + .send() + .await + .expect("invoke"); + let payload = String::from_utf8( + invoked + .payload() + .map(|b| b.as_ref().to_vec()) + .unwrap_or_default(), + ) + .unwrap_or_default(); + assert!( + invoked.function_error().is_none(), + "handler errored: {payload}" + ); + let env: serde_json::Value = serde_json::from_str(&payload).expect("handler returned JSON"); + let get = |k: &str| env[k].as_str().unwrap_or_default().to_string(); + + // Points back at fakecloud on the host, on the port this server bound, + // never `localhost` (inside the container that is the container itself). + let endpoint = get("AWS_ENDPOINT_URL"); + assert!( + endpoint.ends_with(&format!(":{}", server.port())), + "endpoint {endpoint} should target this server's port {}", + server.port() + ); + assert!( + !endpoint.contains("localhost:") && !endpoint.contains("127.0.0.1"), + "endpoint {endpoint} must use the container's host alias" + ); + + // The SDK inside the function signs as the execution role, in a session + // named after the function, as on AWS. Under --verify-sigv4 this only + // works because the credentials are registered, not placeholders. + assert_eq!( + get("caller_arn"), + "arn:aws:sts::123456789012:assumed-role/env-probe-role/env-probe-fn", + "{payload}" + ); + assert!(!get("AWS_SESSION_TOKEN").is_empty(), "{payload}"); + + assert_eq!(get("AWS_REGION"), "us-east-1", "{payload}"); + assert_eq!(get("AWS_DEFAULT_REGION"), "us-east-1", "{payload}"); + assert_eq!(get("AWS_LAMBDA_FUNCTION_NAME"), "env-probe-fn"); + assert_eq!(get("context_function_name"), "env-probe-fn"); + assert_eq!(get("AWS_LAMBDA_FUNCTION_VERSION"), "$LATEST"); + assert_eq!(get("AWS_LAMBDA_FUNCTION_MEMORY_SIZE"), "256"); + assert_eq!(get("AWS_LAMBDA_LOG_GROUP_NAME"), "/aws/lambda/env-probe-fn"); + assert_eq!(get("AWS_EXECUTION_ENV"), "AWS_Lambda_python3.12"); + assert_eq!(get("APP_SETTING"), "kept"); + + // A configuration change reaches the next invocation instead of being + // masked by the warm container. + lambda + .update_function_configuration() + .function_name("env-probe-fn") + .environment( + Environment::builder() + .variables("APP_SETTING", "changed") + .build(), + ) + .send() + .await + .expect("update_function_configuration"); + let invoked = lambda + .invoke() + .function_name("env-probe-fn") + .payload(Blob::new("{}")) + .send() + .await + .expect("invoke after update"); + let payload = String::from_utf8(invoked.payload().unwrap().as_ref().to_vec()).unwrap(); + let env: serde_json::Value = serde_json::from_str(&payload).expect("handler returned JSON"); + assert_eq!(env["APP_SETTING"], "changed", "{payload}"); +} + +/// Lambda reserves the keys it sets itself; a configuration that sets one is +/// rejected on create and on update, as on AWS. +#[tokio::test] +async fn reserved_environment_keys_are_rejected() { + let server = TestServer::start().await; + let lambda = server.lambda_client().await; + let code = || { + FunctionCode::builder() + .zip_file(build_python_handler_zip("def handler(e, c):\n return 1\n").into()) + .build() + }; + + let err = lambda + .create_function() + .function_name("reserved-env-fn") + .runtime(Runtime::Python312) + .role("arn:aws:iam::123456789012:role/r") + .handler("index.handler") + .environment( + Environment::builder() + .variables("AWS_REGION", "eu-west-1") + .variables("AWS_SESSION_TOKEN", "x") + .build(), + ) + .code(code()) + .send() + .await + .expect_err("reserved keys must be rejected"); + let service_err = err.into_service_error(); + assert!( + service_err.is_invalid_parameter_value_exception(), + "{service_err:?}" + ); + let message = service_err.meta().message().unwrap_or_default().to_string(); + assert!( + message.ends_with("Reserved keys used in this request: AWS_REGION, AWS_SESSION_TOKEN"), + "{message}" + ); + + // Not reserved: a function may point its SDK somewhere else. + lambda + .create_function() + .function_name("reserved-env-fn") + .runtime(Runtime::Python312) + .role("arn:aws:iam::123456789012:role/r") + .handler("index.handler") + .environment( + Environment::builder() + .variables("AWS_ENDPOINT_URL", "http://localhost:4566") + .build(), + ) + .code(code()) + .send() + .await + .expect("AWS_ENDPOINT_URL is not reserved"); + + let err = lambda + .update_function_configuration() + .function_name("reserved-env-fn") + .environment( + Environment::builder() + .variables("AWS_ACCESS_KEY_ID", "AKIA") + .build(), + ) + .send() + .await + .expect_err("reserved keys must be rejected on update"); + assert!(err + .into_service_error() + .is_invalid_parameter_value_exception()); + let cfg = lambda + .get_function_configuration() + .function_name("reserved-env-fn") + .send() + .await + .unwrap(); + let vars = cfg.environment().and_then(|e| e.variables()).unwrap(); + assert!( + vars.contains_key("AWS_ENDPOINT_URL") && !vars.contains_key("AWS_ACCESS_KEY_ID"), + "a rejected update must not apply: {vars:?}" + ); +} + +/// Under IAM enforcement `iam:PassRole` is same-account only, and the role +/// must trust Lambda: both CreateFunction and UpdateFunctionConfiguration +/// enforce it. +#[tokio::test] +async fn execution_role_must_be_same_account_and_trust_lambda() { + let server = TestServer::start_with_env(&[("FAKECLOUD_IAM", "strict")]).await; + let root = root_config(&server).await; + let lambda = aws_sdk_lambda::Client::new(&root); + let iam = aws_sdk_iam::Client::new(&root); + let code = || { + FunctionCode::builder() + .zip_file(build_python_handler_zip("def handler(e, c):\n return 1\n").into()) + .build() + }; + + let err = lambda + .create_function() + .function_name("cross-account-fn") + .runtime(Runtime::Python312) + .role("arn:aws:iam::999999999999:role/elsewhere") + .handler("index.handler") + .code(code()) + .send() + .await + .expect_err("a role in another account must be refused"); + let raw = err.raw_response().map(|r| r.status().as_u16()); + let service_err = err.into_service_error(); + assert_eq!(service_err.meta().code(), Some("AccessDeniedException")); + assert_eq!( + service_err.meta().message(), + Some("Cross-account pass role is not allowed.") + ); + assert_eq!(raw, Some(403)); + + lambda + .create_function() + .function_name("role-checked-fn") + .runtime(Runtime::Python312) + .role("arn:aws:iam::123456789012:role/r") + .handler("index.handler") + .code(code()) + .send() + .await + .expect("same-account role"); + + // A role whose trust policy does not name Lambda cannot be swapped in. + let untrusted = iam + .create_role() + .role_name("not-for-lambda") + .assume_role_policy_document( + r#"{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"Service":"ec2.amazonaws.com"},"Action":"sts:AssumeRole"}]}"#, + ) + .send() + .await + .unwrap(); + let err = lambda + .update_function_configuration() + .function_name("role-checked-fn") + .role(untrusted.role().unwrap().arn()) + .send() + .await + .expect_err("an untrusted role must be refused on update"); + assert!(err + .into_service_error() + .is_invalid_parameter_value_exception()); + + let err = lambda + .update_function_configuration() + .function_name("role-checked-fn") + .role("arn:aws:iam::999999999999:role/elsewhere") + .send() + .await + .expect_err("a cross-account role must be refused on update"); + assert_eq!( + err.into_service_error().meta().code(), + Some("AccessDeniedException") + ); + let cfg = lambda + .get_function_configuration() + .function_name("role-checked-fn") + .send() + .await + .unwrap(); + assert_eq!(cfg.role(), Some("arn:aws:iam::123456789012:role/r")); +} + +/// With IAM enforcement off (the default) a role ARN naming another account, +/// as templates written for other emulators do (`000000000000`), is still +/// accepted on create and update; the trust check applies as before. +#[tokio::test] +async fn default_mode_accepts_another_accounts_role() { + let server = TestServer::start().await; + let lambda = server.lambda_client().await; + lambda + .create_function() + .function_name("foreign-role-fn") + .runtime(Runtime::Python312) + .role("arn:aws:iam::000000000000:role/lambda-role") + .handler("index.handler") + .code( + FunctionCode::builder() + .zip_file(build_python_handler_zip("def handler(e, c):\n return 1\n").into()) + .build(), + ) + .send() + .await + .expect("default mode accepts a role from another account"); + lambda + .update_function_configuration() + .function_name("foreign-role-fn") + .role("arn:aws:iam::111111111111:role/other") + .send() + .await + .expect("and on update"); +} diff --git a/crates/fakecloud-iam/src/sts_service/container_creds.rs b/crates/fakecloud-iam/src/sts_service/container_creds.rs index 0c4840d23..17b925125 100644 --- a/crates/fakecloud-iam/src/sts_service/container_creds.rs +++ b/crates/fakecloud-iam/src/sts_service/container_creds.rs @@ -20,8 +20,10 @@ //! unauthenticated endpoint to roughly one live temp credential per role. use std::collections::HashMap; +use std::sync::Arc; use chrono::{DateTime, Duration, Utc}; +use fakecloud_core::auth::{SessionCredentialIssuer, SessionCredentials}; use parking_lot::Mutex; use super::{ @@ -160,6 +162,25 @@ pub fn mint_container_credentials( default_account_id: &str, role_arn: &str, duration: Duration, +) -> ContainerCredentials { + mint_session_credentials( + iam, + default_account_id, + role_arn, + CONTAINER_CREDENTIALS_SESSION_NAME, + duration, + ) +} + +/// [`mint_container_credentials`] with a caller-chosen role session name, for +/// surfaces whose real counterpart names the session after the workload (a +/// Lambda execution role's session is named after the function). +pub fn mint_session_credentials( + iam: &SharedIamState, + default_account_id: &str, + role_arn: &str, + session_name: &str, + duration: Duration, ) -> ContainerCredentials { let creds = StsCredentials::generate(); let issued_at = Utc::now(); @@ -168,7 +189,6 @@ pub fn mint_container_credentials( let account_id = extract_account_from_arn(role_arn).unwrap_or_else(|| default_account_id.to_string()); let role_name = assumed_role_name(role_arn); - let session_name = CONTAINER_CREDENTIALS_SESSION_NAME; let assumed_role_arn = format_assumed_role_arn(partition, &account_id, role_name, session_name); let user_id = format!("{}:{}", deterministic_role_id(role_arn), session_name); @@ -226,10 +246,72 @@ fn credential_registered(iam: &SharedIamState, creds: &ContainerCredentials) -> /// Remove a minted container credential from IAM state (both maps). Idempotent /// -- a no-op if the key is already gone (e.g. after a reset). fn evict_container_credentials(iam: &SharedIamState, creds: &ContainerCredentials) { - let mut accounts = iam.write(); - let state = accounts.get_or_create(&creds.account_id); - state.credential_identities.remove(&creds.access_key_id); - state.sts_temp_credentials.remove(&creds.access_key_id); + unregister_credentials(iam, &creds.account_id, &creds.access_key_id); +} + +fn unregister_credentials(iam: &SharedIamState, account_id: &str, access_key_id: &str) { + // `get_mut`, not `get_or_create`: after a reset the account may be gone, + // and revoking must not bring it back. + if let Some(state) = iam.write().get_mut(account_id) { + state.credential_identities.remove(access_key_id); + state.sts_temp_credentials.remove(access_key_id); + } +} + +impl From for SessionCredentials { + fn from(c: ContainerCredentials) -> Self { + Self { + access_key_id: c.access_key_id, + secret_access_key: c.secret_access_key, + session_token: c.session_token, + expiration: c.expiration, + account_id: c.account_id, + } + } +} + +/// [`SessionCredentialIssuer`] over IAM state: every set it issues is minted +/// and registered exactly like [`mint_session_credentials`], so the code it +/// is handed to signs requests as the assumed role, and is unregistered again +/// on [`SessionCredentialIssuer::revoke`]. +pub struct IamSessionCredentialIssuer { + iam: SharedIamState, + default_account_id: String, +} + +impl IamSessionCredentialIssuer { + /// `default_account_id` registers sessions for role ARNs that carry no + /// account of their own. + pub fn shared( + iam: SharedIamState, + default_account_id: impl Into, + ) -> Arc { + Arc::new(Self { + iam, + default_account_id: default_account_id.into(), + }) + } +} + +impl SessionCredentialIssuer for IamSessionCredentialIssuer { + fn issue(&self, role_arn: &str, session_name: &str, duration: Duration) -> SessionCredentials { + mint_session_credentials( + &self.iam, + &self.default_account_id, + role_arn, + session_name, + duration, + ) + .into() + } + + fn revoke(&self, credentials: &SessionCredentials) { + unregister_credentials( + &self.iam, + &credentials.account_id, + &credentials.access_key_id, + ); + } } /// The live + recently-superseded credentials for one role. @@ -549,4 +631,35 @@ mod tests { let iso = minted.expiration_iso8601(); assert!(iso.ends_with('Z') && iso.len() == 20, "{iso}"); } + + #[test] + fn issued_session_is_named_after_the_caller_and_revocable() { + let iam = shared(); + let issuer = IamSessionCredentialIssuer::shared(iam.clone(), "123456789012"); + let creds = issuer.issue( + "arn:aws:iam::999999999999:role/service-role/exec", + "my-function", + Duration::hours(1), + ); + assert_eq!(creds.account_id, "999999999999"); + + let resolver = IamCredentialResolver::new(iam.clone()); + let resolved = resolver + .resolve(&creds.access_key_id) + .expect("issued key resolves"); + assert_eq!( + resolved.principal.arn, + "arn:aws:sts::999999999999:assumed-role/exec/my-function" + ); + assert_eq!( + resolved.session_token.as_deref(), + Some(creds.session_token.as_str()) + ); + + issuer.revoke(&creds); + assert!(resolver.resolve(&creds.access_key_id).is_none()); + assert_eq!(temp_cred_count(&iam, "999999999999"), 0); + // Revoking twice is a no-op. + issuer.revoke(&creds); + } } diff --git a/crates/fakecloud-k8s/src/client.rs b/crates/fakecloud-k8s/src/client.rs index 6c996322f..0ede2385a 100644 --- a/crates/fakecloud-k8s/src/client.rs +++ b/crates/fakecloud-k8s/src/client.rs @@ -9,10 +9,10 @@ use std::time::{Duration, Instant}; -use k8s_openapi::api::core::v1::Pod; +use k8s_openapi::api::core::v1::{Pod, Secret}; use k8s_openapi::api::networking::v1::NetworkPolicy; use k8s_openapi::apimachinery::pkg::apis::meta::v1::Status; -use kube::api::{Api, AttachParams, DeleteParams, ListParams, PostParams}; +use kube::api::{Api, AttachParams, DeleteParams, ListParams, Patch, PatchParams, PostParams}; use kube::Client; use tokio::io::AsyncReadExt; @@ -379,6 +379,72 @@ impl K8sClient { reaped } + /// Namespaced Secret API handle. + pub fn secrets(&self) -> Api { + Api::namespaced(self.client.clone(), &self.namespace) + } + + /// Create `secret`, replacing a same-named one left behind by a previous + /// process (delete-then-create, like [`create_pod`](Self::create_pod)). + pub async fn create_secret(&self, secret: &Secret) -> Result<(), K8sError> { + let name = secret + .metadata + .name + .clone() + .ok_or_else(|| K8sError::Other("secret has no metadata.name".into()))?; + let api = self.secrets(); + let _ = api.delete(&name, &DeleteParams::default()).await; + api.create(&PostParams::default(), secret) + .await + .map(|_| ()) + .map_err(K8sError::Kube) + } + + /// Make the Pod `pod` the owner of Secret `secret`, so Kubernetes garbage + /// collection deletes the Secret along with the Pod however the Pod goes + /// away (teardown, a reaper, a manual `kubectl delete`). Best-effort: + /// explicit deletion on teardown still applies if this fails. + pub async fn adopt_secret(&self, secret: &str, pod: &str) { + let uid = match self.pods().get(pod).await { + Ok(p) => p.metadata.uid, + Err(e) => { + tracing::warn!(secret, pod, error = %e, "k8s adopt secret: get pod failed"); + return; + } + }; + let Some(uid) = uid else { + return; + }; + let patch = serde_json::json!({ + "metadata": { + "ownerReferences": [{ + "apiVersion": "v1", + "kind": "Pod", + "name": pod, + "uid": uid, + }] + } + }); + if let Err(e) = self + .secrets() + .patch(secret, &PatchParams::default(), &Patch::Merge(&patch)) + .await + { + tracing::warn!(secret, pod, error = %e, "k8s adopt secret: patch failed"); + } + } + + /// Delete a Secret by name. Idempotent; a `404` is success and other + /// errors are logged, since teardown is best-effort. + pub async fn delete_secret(&self, name: &str) { + if let Err(e) = self.secrets().delete(name, &DeleteParams::default()).await { + if matches!(&e, kube::Error::Api(a) if a.code == 404) { + return; + } + tracing::warn!(secret = %name, namespace = %self.namespace, error = %e, "k8s delete secret failed"); + } + } + /// Namespaced NetworkPolicy API handle. pub fn network_policies(&self) -> Api { Api::namespaced(self.client.clone(), &self.namespace) diff --git a/crates/fakecloud-lambda/src/extras/mod.rs b/crates/fakecloud-lambda/src/extras/mod.rs index be2187388..3f45d8990 100644 --- a/crates/fakecloud-lambda/src/extras/mod.rs +++ b/crates/fakecloud-lambda/src/extras/mod.rs @@ -373,6 +373,23 @@ impl LambdaService { Some(size) => Some(crate::service::validate_ephemeral_storage(size)?), None => None, }; + let environment: Option> = + body["Environment"]["Variables"].as_object().map(|env| { + env.iter() + .filter_map(|(k, v)| v.as_str().map(|s| (k.clone(), s.to_string()))) + .collect() + }); + if let Some(env) = &environment { + crate::service::validate_environment(env)?; + } + if let Some(role) = body["Role"].as_str() { + crate::service::validate_execution_role( + &req.account_id, + role, + self.role_trust_validator.as_deref(), + self.iam_mode, + )?; + } let mut accounts = self.state.write(); // Pre-resolve layer attachments before re-borrowing accounts mutably // for the function. Layer ARNs may live in sibling accounts. @@ -406,11 +423,8 @@ impl LambdaService { if let Some(rt) = body["Runtime"].as_str() { func.runtime = rt.to_string(); } - if let Some(env) = body["Environment"]["Variables"].as_object() { - func.environment = env - .iter() - .filter_map(|(k, v)| v.as_str().map(|s| (k.clone(), s.to_string()))) - .collect(); + if let Some(env) = environment { + func.environment = env; } if let Some(mode) = body["TracingConfig"]["Mode"].as_str() { func.tracing_mode = Some(mode.to_string()); diff --git a/crates/fakecloud-lambda/src/lib.rs b/crates/fakecloud-lambda/src/lib.rs index 02c71beb1..1ff829874 100644 --- a/crates/fakecloud-lambda/src/lib.rs +++ b/crates/fakecloud-lambda/src/lib.rs @@ -7,7 +7,7 @@ pub(crate) mod service; pub(crate) mod state; pub(crate) mod workflows; -pub use service::LambdaService; +pub use service::{validate_execution_role, LambdaService}; pub use state::{ function_arn, layer_arn, qualified_function_arn, AttachedLayer, EventSourceMapping, FunctionAlias, FunctionUrlConfig, LambdaFunction, LambdaInvocation, LambdaSnapshot, diff --git a/crates/fakecloud-lambda/src/runtime/backend.rs b/crates/fakecloud-lambda/src/runtime/backend.rs index 7cacfce29..4d4973b13 100644 --- a/crates/fakecloud-lambda/src/runtime/backend.rs +++ b/crates/fakecloud-lambda/src/runtime/backend.rs @@ -9,6 +9,7 @@ //! fakecloud process. use async_trait::async_trait; +use fakecloud_core::auth::SessionCredentials; use crate::state::LambdaFunction; @@ -84,14 +85,24 @@ pub trait LambdaBackend: Send + Sync + 'static { /// package functions). `layers` are the attached layer ZIPs in /// attach order. `deploy_id` is the facade-computed fingerprint /// used to label resources so reaper logic can correlate. + /// `credentials` are the execution role's session credentials for this + /// instance, exported into its environment (see + /// [`super::environment::function_environment`]). async fn launch( &self, func: &LambdaFunction, code_zip: Option<&[u8]>, layers: &[Vec], deploy_id: &str, + credentials: Option<&SessionCredentials>, ) -> Result; + /// Whether a function's tags shape the instances this backend launches + /// (k8s scheduling overrides), so a tag change needs a fresh instance. + fn launch_uses_tags(&self) -> bool { + false + } + /// Tear down one instance. Must be idempotent — the facade may call /// this against an already-gone instance during cleanup races. async fn terminate(&self, handle: &BackendHandle); diff --git a/crates/fakecloud-lambda/src/runtime/docker.rs b/crates/fakecloud-lambda/src/runtime/docker.rs index a5d33906a..1f5022a00 100644 --- a/crates/fakecloud-lambda/src/runtime/docker.rs +++ b/crates/fakecloud-lambda/src/runtime/docker.rs @@ -12,8 +12,9 @@ use base64::Engine; use tempfile::TempDir; use super::backend::{BackendHandle, LambdaBackend, RuntimeError, WarmInstance}; -use super::env_rewrite::rewrite_localhost_envs; +use super::environment::{function_environment, CREDENTIAL_ENV_KEYS}; use crate::state::LambdaFunction; +use fakecloud_core::auth::SessionCredentials; /// Docker/Podman-based Lambda execution backend. pub struct DockerBackend { @@ -96,6 +97,25 @@ impl DockerBackend { } } + /// Export the function's execution environment into the container. The + /// container reaches fakecloud through the host alias, since `localhost` + /// inside it is the container itself. + fn apply_function_env( + &self, + cmd: &mut tokio::process::Command, + func: &LambdaFunction, + credentials: Option<&SessionCredentials>, + ) { + let endpoint_url = format!("http://{}:{}", self.host_alias, self.server_port); + let (args, child_env) = docker_env_args(function_environment( + func, + &endpoint_url, + &self.host_alias, + credentials, + )); + cmd.args(args).envs(child_env); + } + fn docker_config_path(&self) -> Option { self.docker_config.as_ref().map(|d| d.path().to_path_buf()) } @@ -109,6 +129,7 @@ impl DockerBackend { &self, func: &LambdaFunction, layers: &[Vec], + credentials: Option<&SessionCredentials>, ) -> Result { let image = func.image_uri.as_deref().ok_or_else(|| { RuntimeError::ContainerStartFailed("PackageType=Image function has no ImageUri".into()) @@ -159,11 +180,7 @@ impl DockerBackend { .arg(format!("fakecloud-instance={}", self.instance_id)); self.apply_host_alias(&mut cmd); - for (key, value) in rewrite_localhost_envs(&func.environment, &self.host_alias) { - cmd.arg("-e").arg(format!("{key}={value}")); - } - cmd.arg("-e") - .arg(format!("AWS_LAMBDA_FUNCTION_TIMEOUT={}", func.timeout)); + self.apply_function_env(&mut cmd, func, credentials); let tmpfs_arg = ephemeral_storage_tmpfs_arg(func.ephemeral_storage_size); cmd.arg("--tmpfs").arg(tmpfs_arg); @@ -221,6 +238,7 @@ impl DockerBackend { func: &LambdaFunction, zip_bytes: &[u8], layers: &[Vec], + credentials: Option<&SessionCredentials>, ) -> Result { let image = runtime_to_image(&func.runtime) .ok_or_else(|| RuntimeError::UnsupportedRuntime(func.runtime.clone()))?; @@ -246,12 +264,7 @@ impl DockerBackend { .arg(format!("fakecloud-instance={}", self.instance_id)); self.apply_host_alias(&mut cmd); - for (key, value) in rewrite_localhost_envs(&func.environment, &self.host_alias) { - cmd.arg("-e").arg(format!("{key}={value}")); - } - - cmd.arg("-e") - .arg(format!("AWS_LAMBDA_FUNCTION_TIMEOUT={}", func.timeout)); + self.apply_function_env(&mut cmd, func, credentials); let tmpfs_arg = ephemeral_storage_tmpfs_arg(func.ephemeral_storage_size); cmd.arg("--tmpfs").arg(tmpfs_arg); @@ -444,13 +457,15 @@ impl LambdaBackend for DockerBackend { code_zip: Option<&[u8]>, layers: &[Vec], _deploy_id: &str, + credentials: Option<&SessionCredentials>, ) -> Result { if func.package_type == "Image" { - self.start_image_container(func, layers).await + self.start_image_container(func, layers, credentials).await } else { let bytes = code_zip.ok_or_else(|| RuntimeError::NoCodeZip(func.function_name.clone()))?; - self.start_zip_container(func, bytes, layers).await + self.start_zip_container(func, bytes, layers, credentials) + .await } } @@ -508,6 +523,25 @@ impl LambdaBackend for DockerBackend { } } +/// `-e` arguments for a container environment, plus the variables to set on +/// the container CLI's own process. Credentials go by name only (`-e KEY`), +/// which makes the CLI copy the value from its environment, so secrets never +/// appear in argv (`ps`, audit logs, `docker events`). +fn docker_env_args(env: Vec<(String, String)>) -> (Vec, Vec<(String, String)>) { + let mut args = Vec::with_capacity(env.len() * 2); + let mut child_env = Vec::new(); + for (key, value) in env { + args.push("-e".to_string()); + if CREDENTIAL_ENV_KEYS.contains(&key.as_str()) { + args.push(key.clone()); + child_env.push((key, value)); + } else { + args.push(format!("{key}={value}")); + } + } + (args, child_env) +} + /// Map AWS runtime identifier to a Docker image tag. pub fn runtime_to_image(runtime: &str) -> Option { let (base, tag) = match runtime { @@ -744,4 +778,39 @@ mod tests { assert_eq!(ephemeral_storage_tmpfs_arg(Some(0)), "/tmp:size=64m,exec"); assert_eq!(ephemeral_storage_tmpfs_arg(Some(32)), "/tmp:size=64m,exec"); } + + #[test] + fn credentials_stay_off_the_command_line() { + let env = vec![ + ("AWS_REGION".to_string(), "us-east-1".to_string()), + ("AWS_ACCESS_KEY_ID".to_string(), "FSIAKEY".to_string()), + ("AWS_SECRET_ACCESS_KEY".to_string(), "s3cr3t".to_string()), + ("AWS_SESSION_TOKEN".to_string(), "tok".to_string()), + ]; + let (args, child_env) = docker_env_args(env); + assert_eq!( + args, + vec![ + "-e", + "AWS_REGION=us-east-1", + "-e", + "AWS_ACCESS_KEY_ID", + "-e", + "AWS_SECRET_ACCESS_KEY", + "-e", + "AWS_SESSION_TOKEN" + ] + ); + assert!(args + .iter() + .all(|a| !a.contains("s3cr3t") && !a.contains("tok") && !a.contains("FSIAKEY"))); + assert_eq!( + child_env, + vec![ + ("AWS_ACCESS_KEY_ID".to_string(), "FSIAKEY".to_string()), + ("AWS_SECRET_ACCESS_KEY".to_string(), "s3cr3t".to_string()), + ("AWS_SESSION_TOKEN".to_string(), "tok".to_string()), + ] + ); + } } diff --git a/crates/fakecloud-lambda/src/runtime/environment.rs b/crates/fakecloud-lambda/src/runtime/environment.rs new file mode 100644 index 000000000..362467848 --- /dev/null +++ b/crates/fakecloud-lambda/src/runtime/environment.rs @@ -0,0 +1,365 @@ +//! The environment a function's runtime process starts with. +//! +//! Real Lambda starts every execution environment with a set of reserved +//! variables (region, execution-role credentials, function metadata) that +//! function code and the AWS SDKs rely on: an SDK client constructed with no +//! region or credentials fails before it sends anything. Those keys cannot be +//! set by the function itself (`CreateFunction` rejects them, see +//! [`reserved_keys_in`]), so they always carry the platform's values. +//! +//! On top of that, fakecloud adds `AWS_ENDPOINT_URL` pointing back at the +//! server, the one deliberate deviation from AWS: on real Lambda the SDK's +//! default endpoints are correct, here the handler would otherwise call real +//! AWS. It is not reserved, so a function that sets its own keeps it. + +use std::collections::BTreeMap; + +use chrono::{DateTime, Utc}; +use fakecloud_core::auth::SessionCredentials; + +use super::env_rewrite::rewrite_localhost_envs; +use crate::state::LambdaFunction; + +/// Environment variable keys Lambda reserves for itself. A function +/// configuration that sets any of them is rejected with +/// `InvalidParameterValueException`. +pub const RESERVED_ENV_KEYS: &[&str] = &[ + "_HANDLER", + "_X_AMZN_TRACE_ID", + "AWS_DEFAULT_REGION", + "AWS_REGION", + "AWS_EXECUTION_ENV", + "AWS_LAMBDA_FUNCTION_NAME", + "AWS_LAMBDA_FUNCTION_MEMORY_SIZE", + "AWS_LAMBDA_FUNCTION_VERSION", + "AWS_LAMBDA_INITIALIZATION_TYPE", + "AWS_LAMBDA_LOG_GROUP_NAME", + "AWS_LAMBDA_LOG_STREAM_NAME", + "AWS_ACCESS_KEY", + "AWS_ACCESS_KEY_ID", + "AWS_SECRET_ACCESS_KEY", + "AWS_SESSION_TOKEN", + "AWS_LAMBDA_RUNTIME_API", + "LAMBDA_TASK_ROOT", + "LAMBDA_RUNTIME_DIR", +]; + +/// The reserved keys present in a function's environment, in key order. +pub fn reserved_keys_in(environment: &BTreeMap) -> Vec<&str> { + environment + .keys() + .map(String::as_str) + .filter(|k| RESERVED_ENV_KEYS.contains(k)) + .collect() +} + +/// The `InvalidParameterValueException` message Lambda returns for a +/// configuration that sets reserved keys, or `None` when it sets none. +pub fn reserved_keys_message(environment: &BTreeMap) -> Option { + let reserved = reserved_keys_in(environment); + if reserved.is_empty() { + return None; + } + Some(format!( + "Lambda was unable to configure your environment variables because the environment \ + variables you have provided contains reserved keys that are currently not supported \ + for modification. Reserved keys used in this request: {}", + reserved.join(", ") + )) +} + +/// Region from a function ARN (`arn::lambda:::function:`). +/// Function ARNs are always minted with a region; `us-east-1` only covers a +/// malformed value. +pub fn region_from_function_arn(arn: &str) -> &str { + arn.split(':') + .nth(3) + .filter(|r| !r.is_empty()) + .unwrap_or("us-east-1") +} + +/// Log group the function writes to: `LoggingConfig.LogGroup` when set, +/// otherwise Lambda's default `/aws/lambda/`. +fn log_group_name(func: &LambdaFunction) -> String { + func.logging_config + .as_ref() + .and_then(|c| c["LogGroup"].as_str()) + .filter(|g| !g.is_empty()) + .map(str::to_string) + .unwrap_or_else(|| format!("/aws/lambda/{}", func.function_name)) +} + +/// Lambda's per-execution-environment log stream name: +/// `YYYY/MM/DD/[]<32 hex>`. +fn log_stream_name(version: &str, started_at: DateTime, instance_id: &str) -> String { + format!("{}/[{version}]{instance_id}", started_at.format("%Y/%m/%d")) +} + +/// `AWS_EXECUTION_ENV` for a managed runtime (`AWS_Lambda_python3.12`). Custom +/// runtimes (`provided*`) and container images set none, as on AWS. +fn execution_env(func: &LambdaFunction) -> Option { + if func.package_type == "Image" + || func.runtime.is_empty() + || func.runtime.starts_with("provided") + { + return None; + } + Some(format!("AWS_Lambda_{}", func.runtime)) +} + +/// The environment keys that carry the execution role's credentials. Backends +/// keep their values off command lines and out of Pod specs. +pub const CREDENTIAL_ENV_KEYS: [&str; 4] = [ + "AWS_ACCESS_KEY_ID", + "AWS_ACCESS_KEY", + "AWS_SECRET_ACCESS_KEY", + "AWS_SESSION_TOKEN", +]; + +/// The credential variables for `creds`, keyed by [`CREDENTIAL_ENV_KEYS`]. +pub fn credential_envs(creds: &SessionCredentials) -> [(&'static str, String); 4] { + let [akid, legacy_akid, secret, token] = CREDENTIAL_ENV_KEYS; + [ + (akid, creds.access_key_id.clone()), + (legacy_akid, creds.access_key_id.clone()), + (secret, creds.secret_access_key.clone()), + (token, creds.session_token.clone()), + ] +} + +/// Assemble the full environment for one execution environment of `func`. +/// +/// Returns each key exactly once, resolved by precedence: fakecloud's +/// endpoint override, then the function's own variables (with +/// `localhost`/`127.0.0.1` URLs rewritten to `rewrite_host`, since inside the +/// container those name the container itself), then the reserved variables, +/// which always win. `endpoint_url` is how this backend's containers reach the +/// fakecloud server. `credentials` are the execution role's session +/// credentials; `None` leaves them unset. +pub fn function_environment( + func: &LambdaFunction, + endpoint_url: &str, + rewrite_host: &str, + credentials: Option<&SessionCredentials>, +) -> Vec<(String, String)> { + let started_at = Utc::now(); + let instance_id = uuid::Uuid::new_v4().simple().to_string(); + let region = region_from_function_arn(&func.function_arn); + + let mut env: BTreeMap = BTreeMap::new(); + env.insert("AWS_ENDPOINT_URL".into(), endpoint_url.to_string()); + for (key, value) in rewrite_localhost_envs(&func.environment, rewrite_host) { + if !RESERVED_ENV_KEYS.contains(&key.as_str()) { + env.insert(key, value); + } + } + + let mut reserved: Vec<(&str, String)> = vec![ + ("AWS_REGION", region.to_string()), + ("AWS_DEFAULT_REGION", region.to_string()), + ("AWS_LAMBDA_FUNCTION_NAME", func.function_name.clone()), + ("AWS_LAMBDA_FUNCTION_VERSION", func.version.clone()), + ( + "AWS_LAMBDA_FUNCTION_MEMORY_SIZE", + func.memory_size.to_string(), + ), + ("AWS_LAMBDA_LOG_GROUP_NAME", log_group_name(func)), + ( + "AWS_LAMBDA_LOG_STREAM_NAME", + log_stream_name(&func.version, started_at, &instance_id), + ), + ("AWS_LAMBDA_INITIALIZATION_TYPE", "on-demand".to_string()), + // Not an AWS variable, and not reserved (a function may set it and + // CreateFunction accepts that). The Runtime Interface Emulator + // enforces the invocation timeout from it, so the platform's value + // (the configured Timeout) is the one the container gets. + ("AWS_LAMBDA_FUNCTION_TIMEOUT", func.timeout.to_string()), + ]; + if let Some(exec_env) = execution_env(func) { + reserved.push(("AWS_EXECUTION_ENV", exec_env)); + } + if let Some(creds) = credentials { + reserved.extend(credential_envs(creds)); + } + for (key, value) in reserved { + env.insert(key.to_string(), value); + } + env.into_iter().collect() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn func() -> LambdaFunction { + serde_json::from_value(serde_json::json!({ + "function_name": "probe", + "function_arn": "arn:aws:lambda:eu-west-2:123456789012:function:probe", + "runtime": "python3.12", + "role": "arn:aws:iam::123456789012:role/r", + "handler": "index.handler", + "description": "", + "timeout": 7, + "memory_size": 256, + "code_sha256": "sha", + "code_size": 1, + "version": "$LATEST", + "last_modified": "2020-01-01T00:00:00Z", + "tags": {}, + "environment": {}, + "architectures": ["x86_64"], + "package_type": "Zip", + "code_zip": null, + "policy": null + })) + .expect("build test LambdaFunction") + } + + fn creds() -> SessionCredentials { + SessionCredentials { + access_key_id: "FSIAEXAMPLE".into(), + secret_access_key: "secret".into(), + session_token: "token".into(), + expiration: Utc::now(), + account_id: "123456789012".into(), + } + } + + fn get<'a>(env: &'a [(String, String)], key: &str) -> Option<&'a str> { + env.iter().find(|(k, _)| k == key).map(|(_, v)| v.as_str()) + } + + #[test] + fn carries_the_reserved_lambda_environment() { + let env = function_environment( + &func(), + "http://host.docker.internal:4566", + "host.docker.internal", + Some(&creds()), + ); + assert_eq!( + get(&env, "AWS_ENDPOINT_URL"), + Some("http://host.docker.internal:4566") + ); + assert_eq!(get(&env, "AWS_REGION"), Some("eu-west-2")); + assert_eq!(get(&env, "AWS_DEFAULT_REGION"), Some("eu-west-2")); + assert_eq!(get(&env, "AWS_ACCESS_KEY_ID"), Some("FSIAEXAMPLE")); + assert_eq!(get(&env, "AWS_ACCESS_KEY"), Some("FSIAEXAMPLE")); + assert_eq!(get(&env, "AWS_SECRET_ACCESS_KEY"), Some("secret")); + assert_eq!(get(&env, "AWS_SESSION_TOKEN"), Some("token")); + assert_eq!(get(&env, "AWS_LAMBDA_FUNCTION_NAME"), Some("probe")); + assert_eq!(get(&env, "AWS_LAMBDA_FUNCTION_VERSION"), Some("$LATEST")); + assert_eq!(get(&env, "AWS_LAMBDA_FUNCTION_MEMORY_SIZE"), Some("256")); + assert_eq!(get(&env, "AWS_LAMBDA_FUNCTION_TIMEOUT"), Some("7")); + assert_eq!( + get(&env, "AWS_LAMBDA_LOG_GROUP_NAME"), + Some("/aws/lambda/probe") + ); + assert_eq!( + get(&env, "AWS_LAMBDA_INITIALIZATION_TYPE"), + Some("on-demand") + ); + assert_eq!( + get(&env, "AWS_EXECUTION_ENV"), + Some("AWS_Lambda_python3.12") + ); + let stream = get(&env, "AWS_LAMBDA_LOG_STREAM_NAME").unwrap(); + let (date, rest) = stream.split_at(10); + assert_eq!(date.matches('/').count(), 2, "{stream}"); + let hex = rest.strip_prefix("/[$LATEST]").expect(stream); + assert!( + hex.len() == 32 && hex.chars().all(|c| c.is_ascii_hexdigit()), + "{stream}" + ); + } + + #[test] + fn every_key_appears_once() { + let mut f = func(); + f.environment + .insert("AWS_ENDPOINT_URL".into(), "http://localhost:9999".into()); + let env = function_environment(&f, "http://h:4566", "h", Some(&creds())); + let mut keys: Vec<&str> = env.iter().map(|(k, _)| k.as_str()).collect(); + let total = keys.len(); + keys.dedup(); + assert_eq!(keys.len(), total); + } + + #[test] + fn function_can_override_the_endpoint_but_not_reserved_keys() { + let mut f = func(); + f.environment + .insert("AWS_ENDPOINT_URL".into(), "http://localhost:9999".into()); + f.environment + .insert("AWS_REGION".into(), "ap-south-1".into()); + f.environment + .insert("AWS_LAMBDA_FUNCTION_TIMEOUT".into(), "900".into()); + f.environment.insert("APP_MODE".into(), "test".into()); + let env = function_environment(&f, "http://h:4566", "h", Some(&creds())); + // Its own endpoint wins, rewritten off `localhost` like any URL. + assert_eq!(get(&env, "AWS_ENDPOINT_URL"), Some("http://h:9999")); + assert_eq!(get(&env, "AWS_REGION"), Some("eu-west-2")); + assert_eq!(get(&env, "AWS_LAMBDA_FUNCTION_TIMEOUT"), Some("7")); + assert_eq!(get(&env, "APP_MODE"), Some("test")); + } + + #[test] + fn no_credentials_without_an_issuer() { + let env = function_environment(&func(), "http://h:4566", "h", None); + for key in [ + "AWS_ACCESS_KEY_ID", + "AWS_ACCESS_KEY", + "AWS_SECRET_ACCESS_KEY", + "AWS_SESSION_TOKEN", + ] { + assert_eq!(get(&env, key), None, "{key}"); + } + } + + #[test] + fn version_log_group_and_execution_env_follow_the_configuration() { + let mut f = func(); + f.version = "3".into(); + f.logging_config = Some(serde_json::json!({"LogGroup": "/custom/group"})); + let env = function_environment(&f, "http://h:4566", "h", None); + assert_eq!(get(&env, "AWS_LAMBDA_FUNCTION_VERSION"), Some("3")); + assert_eq!( + get(&env, "AWS_LAMBDA_LOG_GROUP_NAME"), + Some("/custom/group") + ); + assert!(get(&env, "AWS_LAMBDA_LOG_STREAM_NAME") + .unwrap() + .contains("/[3]")); + + f.runtime = "provided.al2023".into(); + let env = function_environment(&f, "http://h:4566", "h", None); + assert_eq!(get(&env, "AWS_EXECUTION_ENV"), None); + + f.runtime = String::new(); + f.package_type = "Image".into(); + let env = function_environment(&f, "http://h:4566", "h", None); + assert_eq!(get(&env, "AWS_EXECUTION_ENV"), None); + } + + #[test] + fn reserved_keys_are_detected() { + let mut env = BTreeMap::new(); + env.insert("APP".to_string(), "x".to_string()); + env.insert("AWS_REGION".to_string(), "x".to_string()); + env.insert("AWS_SESSION_TOKEN".to_string(), "x".to_string()); + env.insert("AWS_ENDPOINT_URL".to_string(), "x".to_string()); + assert_eq!( + reserved_keys_in(&env), + vec!["AWS_REGION", "AWS_SESSION_TOKEN"] + ); + } + + #[test] + fn region_is_read_from_the_function_arn() { + assert_eq!( + region_from_function_arn("arn:aws:lambda:eu-west-2:123456789012:function:f"), + "eu-west-2" + ); + assert_eq!(region_from_function_arn("not-an-arn"), "us-east-1"); + } +} diff --git a/crates/fakecloud-lambda/src/runtime/facade.rs b/crates/fakecloud-lambda/src/runtime/facade.rs index e1a2e0e33..35e13d584 100644 --- a/crates/fakecloud-lambda/src/runtime/facade.rs +++ b/crates/fakecloud-lambda/src/runtime/facade.rs @@ -5,10 +5,12 @@ //! whatever [`LambdaBackend`] it was constructed with. use std::collections::HashMap; -use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, OnceLock}; use std::time::{Duration, Instant}; use base64::Engine; +use fakecloud_core::auth::{SessionCredentialIssuer, SessionCredentials}; use parking_lot::RwLock; use sha2::{Digest, Sha256}; @@ -22,11 +24,20 @@ use crate::state::LambdaFunction; pub(crate) struct WarmEntry { instance: WarmInstance, last_used: RwLock, - /// Combined fingerprint of the function's code SHA-256 plus the - /// SHA-256 of every attached layer's ZIP bytes, joined in attach + /// Combined fingerprint of the function's code SHA-256, the launch + /// configuration baked into the instance (see [`launch_config`]), and + /// the SHA-256 of every attached layer's ZIP bytes, joined in attach /// order. Layers mutate `/opt`, so a layer change invalidates the /// warm instance even when the function code is unchanged. deploy_id: String, + /// The execution role's session credentials exported into this + /// instance's environment; revoked when the instance is torn down. + credentials: Option, + /// Set when the instance must not take another invocation (its + /// credentials were dropped by an IAM reset, or its version deleted). A + /// busy instance finishes its current invocation and is retired once + /// free. + retiring: AtomicBool, /// Held for the duration of a single invocation against this /// instance. The AWS Runtime Interface Emulator (and real Lambda) /// handles exactly one event per execution environment at a time; @@ -56,6 +67,94 @@ const MAX_INVOKE_ATTEMPTS: u32 = 5; /// so the state-machine retry can succeed within its window. const REACHABILITY_PROBE_TIMEOUT: Duration = Duration::from_millis(1500); +/// Lifetime of the execution-role session minted for each instance: the +/// longest an IAM role session can last. +const EXECUTION_CREDENTIALS_LIFETIME: Duration = Duration::from_secs(12 * 60 * 60); + +/// Instances stop taking new invocations this long before their session +/// expires: Lambda's maximum function timeout plus a minute, so an +/// invocation never runs past its credentials. +const INVOCATION_HEADROOM: Duration = Duration::from_secs(16 * 60); + +/// Key of a function's warm pool: its qualified ARN (`:`). +/// Accounts, regions and versions never share or evict each other's instances. +fn pool_key(func: &LambdaFunction) -> String { + format!("{}:{}", func.function_arn, func.version) +} + +/// Function name inside a pool key +/// (`arn:

:lambda:::function::`). +fn pool_key_function_name(key: &str) -> &str { + key.split(':').nth(6).unwrap_or(key) +} + +/// Account inside a pool key. +fn pool_key_account(key: &str) -> &str { + fakecloud_aws::arn::account_of(key).unwrap_or_default() +} + +/// Whether `key` is a pool of the function whose unqualified ARN is +/// `function_arn` (any version). +fn pool_key_belongs_to(key: &str, function_arn: &str) -> bool { + key.strip_prefix(function_arn) + .is_some_and(|rest| rest.starts_with(':')) +} + +/// The role the execution session is minted for: the function's role, in the +/// function's account. The two only differ when IAM is not strict and a +/// template names another account's role (commonly another emulator's +/// default `000000000000`); minting there would send the function's SDK calls +/// to an empty account instead of the one holding its resources. +fn session_role_arn(role_arn: &str, function_arn: &str) -> String { + let (Some(function_account), Some(role_account)) = ( + fakecloud_aws::arn::account_of(function_arn), + fakecloud_aws::arn::account_of(role_arn), + ) else { + return role_arn.to_string(); + }; + let mut parts: Vec<&str> = role_arn.split(':').collect(); + if role_account != function_account { + parts[4] = function_account; + } + parts.join(":") +} + +/// Whether `entry` can take a new invocation of `deploy_id`: same code + +/// launch configuration, not being retired, and credentials that outlive a +/// full invocation (judged from their own expiration, not the launch time). +/// Cheap: no locks beyond the entry's own fields. +fn is_current(entry: &WarmEntry, deploy_id: &str) -> bool { + let headroom = chrono::Duration::from_std(INVOCATION_HEADROOM).expect("headroom fits"); + entry.deploy_id == deploy_id + && !entry.retiring.load(Ordering::Acquire) + && entry + .credentials + .as_ref() + .is_none_or(|c| c.expiration - chrono::Utc::now() > headroom) +} + +/// Execution-role credentials minted for an instance that is still being +/// launched. Revoked on drop (a failed launch, or the launching future being +/// cancelled) unless [`disarm`](Self::disarm)ed once the instance is pooled. +struct PendingCredentials<'a> { + issuer: Option<&'a Arc>, + credentials: Option, +} + +impl PendingCredentials<'_> { + fn disarm(mut self) { + self.credentials = None; + } +} + +impl Drop for PendingCredentials<'_> { + fn drop(&mut self) { + if let (Some(issuer), Some(creds)) = (self.issuer, &self.credentials) { + issuer.revoke(creds); + } + } +} + /// A reserved invocation slot: a warm instance plus the held busy guard /// that grants exclusive use of it until the guard drops. struct Slot { @@ -73,15 +172,45 @@ struct Slot { /// `fakecloud-deploy-id` Pod label; standard base64's `/` would grow an /// extra URL path segment, break the axum route match, and wedge the Pod /// in a cold-start loop for ~49% of deploys (issue #1643). -fn deploy_id_for(func: &LambdaFunction, layers: &[Vec]) -> String { - deploy_id_from(&func.code_sha256, layers) +/// +/// `with_tags` folds the function's tags in, for backends whose instances +/// are shaped by them (k8s scheduling); elsewhere TagResource must not +/// cold-start the function. +fn deploy_id_for(func: &LambdaFunction, layers: &[Vec], with_tags: bool) -> String { + deploy_id_from(&func.code_sha256, &launch_config(func, with_tags), layers) +} + +/// The function configuration an instance is started with: everything that +/// ends up in its environment, command, sandbox limits, or (with `with_tags`) +/// scheduling. A change to any of it must start a fresh instance rather than +/// reuse one configured for something else. Identity (account, version) is +/// the pool key, not part of this. +fn launch_config(func: &LambdaFunction, with_tags: bool) -> String { + serde_json::json!([ + func.role, + func.runtime, + func.handler, + func.timeout, + func.memory_size, + func.environment, + func.package_type, + func.image_uri, + func.image_config, + func.ephemeral_storage_size, + func.logging_config, + func.architectures, + with_tags.then_some(&func.tags), + ]) + .to_string() } /// Pure core of [`deploy_id_for`], split out so the URL-path-safety /// invariant can be tested without constructing a full `LambdaFunction`. -fn deploy_id_from(code_sha256: &str, layers: &[Vec]) -> String { +fn deploy_id_from(code_sha256: &str, launch_config: &str, layers: &[Vec]) -> String { let mut hasher = Sha256::new(); hasher.update(code_sha256.as_bytes()); + hasher.update(b"\0"); + hasher.update(launch_config.as_bytes()); for bytes in layers { let mut layer_hasher = Sha256::new(); layer_hasher.update(bytes); @@ -115,6 +244,9 @@ pub struct LambdaRuntime { starting: RwLock>>>, /// Cap on warm instances per function. max_concurrency: usize, + /// Mints the execution-role credentials each instance is started with. + /// Unset (no credentials exported) until the server wires IAM in. + credential_issuer: OnceLock>, } impl LambdaRuntime { @@ -132,9 +264,100 @@ impl LambdaRuntime { instances: RwLock::new(HashMap::new()), starting: RwLock::new(HashMap::new()), max_concurrency, + credential_issuer: OnceLock::new(), } } + /// Give instances execution-role credentials: each launch mints a + /// session for the function's role (named after the function, as on + /// AWS) and exports it into the instance environment. Set once at + /// server startup; later calls are ignored. + pub fn set_credential_issuer(&self, issuer: Arc) { + let _ = self.credential_issuer.set(issuer); + } + + fn issue_credentials(&self, func: &LambdaFunction) -> Option { + let issuer = self.credential_issuer.get()?; + if func.role.is_empty() { + return None; + } + let lifetime = chrono::Duration::from_std(EXECUTION_CREDENTIALS_LIFETIME) + .expect("session lifetime fits chrono::Duration"); + let role = session_role_arn(&func.role, &func.function_arn); + Some(issuer.issue(&role, &func.function_name, lifetime)) + } + + /// Tear down one instance and revoke the credentials it was given. + async fn retire(&self, entry: &WarmEntry) { + self.backend.terminate(&entry.instance.handle).await; + self.revoke(entry.credentials.as_ref()); + } + + fn revoke(&self, credentials: Option<&SessionCredentials>) { + if let (Some(issuer), Some(creds)) = (self.credential_issuer.get(), credentials) { + issuer.revoke(creds); + } + } + + /// Mark every pooled instance matching `pred` as retiring so it takes no + /// further invocation. With `detach_free`, the free ones are also removed + /// from the pool and returned for teardown; busy ones always stay until + /// released (then the idle sweep or the next launch retires them). + fn mark_retiring( + &self, + pred: impl Fn(&str, &WarmEntry) -> bool, + detach_free: bool, + ) -> Vec> { + let mut map = self.instances.write(); + let mut detached = Vec::new(); + for (key, pool) in map.iter_mut() { + pool.retain(|e| { + if !pred(key, e) { + return true; + } + e.retiring.store(true, Ordering::Release); + if detach_free && e.busy.try_lock().is_ok() { + detached.push(e.clone()); + false + } else { + true + } + }); + } + map.retain(|_, pool| !pool.is_empty()); + detached + } + + /// Synchronously stop handing invocations to every instance holding + /// credentials registered in `account_id` (all accounts when `None`). + /// Called in step with an IAM reset that drops those credentials, so no + /// invocation can start with them afterwards; tear the instances down + /// with [`Self::retire_released`]. + pub fn mark_credentials_revoked(&self, account_id: Option<&str>) { + self.mark_retiring( + |_, e| { + e.credentials + .as_ref() + .is_some_and(|c| account_id.is_none_or(|a| c.account_id == a)) + }, + false, + ); + } + + /// Retire every free instance already marked retiring. + pub async fn retire_released(&self) { + let free = self.mark_retiring(|_, e| e.retiring.load(Ordering::Acquire), true); + self.terminate_instances(free).await; + } + + /// DeleteFunction with a Qualifier: stop the version's instances without + /// cutting off an in-flight invocation. Free instances are detached and + /// returned for teardown; busy ones are marked retiring and go once free. + pub(crate) fn retire_version(&self, function_arn: &str, version: &str) -> Vec> { + let key = format!("{function_arn}:{version}"); + self.mark_retiring(|k, _| k == key, true) + } + /// Auto-detect Docker or Podman. Returns `None` if neither is available. /// Override with `FAKECLOUD_CONTAINER_CLI` env var. pub fn auto_detect_docker(server_port: u16) -> Option { @@ -257,7 +480,7 @@ impl LambdaRuntime { { let entry = slot.entry.clone(); drop(slot); - self.evict_entry(&func.function_name, &entry).await; + self.evict_entry(&pool_key(func), &entry).await; if attempt < MAX_INVOKE_ATTEMPTS { tracing::warn!( function = %func.function_name, @@ -305,7 +528,7 @@ impl LambdaRuntime { // function already ran and may have side effects. let entry = slot.entry.clone(); drop(slot); - self.evict_entry(&func.function_name, &entry).await; + self.evict_entry(&pool_key(func), &entry).await; Err(RuntimeError::InvocationFailed(e.to_string())) } }; @@ -321,7 +544,7 @@ impl LambdaRuntime { // duplicate invoke. let entry = slot.entry.clone(); drop(slot); - self.evict_entry(&func.function_name, &entry).await; + self.evict_entry(&pool_key(func), &entry).await; if attempt < MAX_INVOKE_ATTEMPTS && e.is_connect() { tracing::warn!( function = %func.function_name, @@ -367,7 +590,7 @@ impl LambdaRuntime { { let entry = slot.entry.clone(); drop(slot); - self.evict_entry(&func.function_name, &entry).await; + self.evict_entry(&pool_key(func), &entry).await; if attempt < MAX_INVOKE_ATTEMPTS { continue; } @@ -405,7 +628,7 @@ impl LambdaRuntime { // half-run handler isn't invoked twice. let entry = slot.entry.clone(); drop(slot); - self.evict_entry(&func.function_name, &entry).await; + self.evict_entry(&pool_key(func), &entry).await; if attempt < MAX_INVOKE_ATTEMPTS && e.is_connect() { continue; } @@ -434,106 +657,144 @@ impl LambdaRuntime { return Err(RuntimeError::NoCodeZip(func.function_name.clone())); } - let deploy_id = deploy_id_for(func, layers); + let deploy_id = deploy_id_for(func, layers, self.backend.launch_uses_tags()); + let key = pool_key(func); - // (1) Fast path: a free instance already running the right deploy. - if let Some(slot) = self.try_take_free(&func.function_name, &deploy_id) { - return Ok(slot); - } + loop { + // (1) Fast path: a free instance already running the right deploy. + if let Some(slot) = self.try_take_free(&key, &deploy_id) { + return Ok(slot); + } - // Serialize launch decisions per function so a burst of cold - // invokes doesn't each push the pool past the cap. - let startup_lock = { - let mut starting = self.starting.write(); - starting - .entry(func.function_name.clone()) - .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(()))) - .clone() - }; - let startup_guard = startup_lock.lock().await; + // Serialize launch decisions per pool so a burst of cold + // invokes doesn't each push the pool past the cap. + let startup_lock = { + let mut starting = self.starting.write(); + starting + .entry(key.clone()) + .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(()))) + .clone() + }; + let startup_guard = startup_lock.lock().await; - // Re-check under the startup lock — another task may have freed or - // launched an instance while we waited. - if let Some(slot) = self.try_take_free(&func.function_name, &deploy_id) { - return Ok(slot); - } + // Re-check under the startup lock: another task may have freed or + // launched an instance while we waited. + if let Some(slot) = self.try_take_free(&key, &deploy_id) { + return Ok(slot); + } - // Tear down any instances left over from a previous deploy. - self.evict_stale_deploy(&func.function_name, &deploy_id) - .await; + // Tear down free instances left over from a previous deploy. + self.evict_stale_deploy(&key, &deploy_id).await; - let pool_len = self - .instances - .read() - .get(&func.function_name) - .map_or(0, |v| v.len()); - - // (2) Room to grow: launch a fresh instance and reserve it. - if pool_len < self.max_concurrency { - let instance = self - .backend - .launch(func, func.code_zip.as_deref(), layers, &deploy_id) - .await?; - let entry = Arc::new(WarmEntry { - instance, - last_used: RwLock::new(Instant::now()), - deploy_id, - busy: Arc::new(tokio::sync::Mutex::new(())), + // Size against current instances only: stale ones still busy + // finishing an invocation are on their way out. + let current_len = self.instances.read().get(&key).map_or(0, |pool| { + pool.iter().filter(|e| is_current(e, &deploy_id)).count() }); - let guard = entry - .busy - .clone() - .try_lock_owned() - .expect("freshly created busy lock is uncontended"); - self.instances - .write() - .entry(func.function_name.clone()) - .or_default() - .push(entry.clone()); - return Ok(Slot { entry, guard }); - } - // (3) At capacity: release the startup lock and wait for whichever - // current instance frees up *first*. Racing every instance's lock - // (rather than blocking on a fixed one) avoids convoying every - // queued caller onto pool[0] while a different instance goes idle. - drop(startup_guard); - let candidates: Vec> = { - let map = self.instances.read(); - map.get(&func.function_name) - .map(|pool| { - pool.iter() - .filter(|e| e.deploy_id == deploy_id) - .cloned() - .collect() + // (2) Room to grow: launch a fresh instance and reserve it. + if current_len < self.max_concurrency { + return self.launch_slot(func, layers, &key, deploy_id).await; + } + + // (3) At capacity: release the startup lock and wait for whichever + // current instance frees up *first*. Racing every instance's lock + // (rather than blocking on a fixed one) avoids convoying every + // queued caller onto pool[0] while a different instance goes idle. + drop(startup_guard); + let candidates: Vec> = { + let map = self.instances.read(); + map.get(&key) + .map(|pool| { + pool.iter() + .filter(|e| is_current(e, &deploy_id)) + .cloned() + .collect() + }) + .unwrap_or_default() + }; + if candidates.is_empty() { + // Every instance went away while we sized the pool; start over. + continue; + } + let waiters = candidates.into_iter().map(|entry| { + Box::pin(async move { + let guard = entry.busy.clone().lock_owned().await; + Slot { entry, guard } }) - .unwrap_or_default() - }; - if candidates.is_empty() { - return Err(RuntimeError::InvocationFailed(format!( - "no warm instance available for {}", - func.function_name - ))); + }); + let (slot, _idx, _rest) = futures_util::future::select_all(waiters).await; + // The instance may have been retired or evicted while we waited. + if !is_current(&slot.entry, &deploy_id) || !self.in_pool(&key, &slot.entry) { + continue; + } + *slot.entry.last_used.write() = Instant::now(); + return Ok(slot); } - let waiters = candidates.into_iter().map(|entry| { - Box::pin(async move { - let guard = entry.busy.clone().lock_owned().await; - Slot { entry, guard } - }) + } + + /// Launch a fresh instance for `func`, with freshly minted execution-role + /// credentials, and reserve it for the caller. + async fn launch_slot( + &self, + func: &LambdaFunction, + layers: &[Vec], + key: &str, + deploy_id: String, + ) -> Result { + let pending = PendingCredentials { + issuer: self.credential_issuer.get(), + credentials: self.issue_credentials(func), + }; + let instance = self + .backend + .launch( + func, + func.code_zip.as_deref(), + layers, + &deploy_id, + pending.credentials.as_ref(), + ) + .await?; + let entry = Arc::new(WarmEntry { + instance, + last_used: RwLock::new(Instant::now()), + deploy_id, + credentials: pending.credentials.clone(), + retiring: AtomicBool::new(false), + busy: Arc::new(tokio::sync::Mutex::new(())), }); - let (slot, _idx, _rest) = futures_util::future::select_all(waiters).await; - *slot.entry.last_used.write() = Instant::now(); - Ok(slot) + let guard = entry + .busy + .clone() + .try_lock_owned() + .expect("freshly created busy lock is uncontended"); + self.instances + .write() + .entry(key.to_string()) + .or_default() + .push(entry.clone()); + // Pooled: the entry owns the credentials now and revokes them on + // retirement. + pending.disarm(); + Ok(Slot { entry, guard }) + } + + fn in_pool(&self, key: &str, entry: &Arc) -> bool { + self.instances + .read() + .get(key) + .is_some_and(|pool| pool.iter().any(|e| Arc::ptr_eq(e, entry))) } /// Try to reserve a free, current-deploy instance without launching. /// Returns `None` if every matching instance is busy (or there are /// none). - fn try_take_free(&self, function_name: &str, deploy_id: &str) -> Option { + fn try_take_free(&self, key: &str, deploy_id: &str) -> Option { let map = self.instances.read(); - let pool = map.get(function_name)?; + let pool = map.get(key)?; for entry in pool { - if entry.deploy_id != deploy_id { + if !is_current(entry, deploy_id) { continue; } if let Ok(guard) = entry.busy.clone().try_lock_owned() { @@ -549,17 +810,17 @@ impl LambdaRuntime { /// Remove one specific instance from a function's pool and terminate /// it. Used when an invocation finds the instance unreachable. - async fn evict_entry(&self, function_name: &str, target: &Arc) { + async fn evict_entry(&self, key: &str, target: &Arc) { let removed = { let mut map = self.instances.write(); - match map.get_mut(function_name) { + match map.get_mut(key) { Some(pool) => { let removed = pool .iter() .position(|e| Arc::ptr_eq(e, target)) .map(|pos| pool.remove(pos)); if pool.is_empty() { - map.remove(function_name); + map.remove(key); } removed } @@ -568,24 +829,26 @@ impl LambdaRuntime { }; if let Some(entry) = removed { tracing::info!( - function = %function_name, + pool = %key, handle = ?entry.instance.handle, "evicting unreachable Lambda runtime instance" ); - self.backend.terminate(&entry.instance.handle).await; + self.retire(&entry).await; } } - /// Tear down every instance in a function's pool whose `deploy_id` - /// no longer matches the current code+layers fingerprint. - async fn evict_stale_deploy(&self, function_name: &str, deploy_id: &str) { + /// Tear down every *free* instance in a pool that can no longer take an + /// invocation of the current deploy (see [`is_current`]). A busy one is + /// mid-invocation and is never cut off: it is left in place, takes no + /// new work, and goes on a later sweep once free. + async fn evict_stale_deploy(&self, key: &str, deploy_id: &str) { let stale: Vec> = { let mut map = self.instances.write(); - match map.get_mut(function_name) { + match map.get_mut(key) { Some(pool) => { let mut stale = Vec::new(); pool.retain(|e| { - if e.deploy_id == deploy_id { + if is_current(e, deploy_id) || e.busy.try_lock().is_err() { true } else { stale.push(e.clone()); @@ -593,7 +856,7 @@ impl LambdaRuntime { } }); if pool.is_empty() { - map.remove(function_name); + map.remove(key); } stale } @@ -602,11 +865,11 @@ impl LambdaRuntime { }; for entry in stale { tracing::info!( - function = %function_name, + pool = %key, handle = ?entry.instance.handle, "stopping stale-deploy Lambda runtime instance" ); - self.backend.terminate(&entry.instance.handle).await; + self.retire(&entry).await; } } @@ -617,11 +880,27 @@ impl LambdaRuntime { /// identically) is not reaped by the deferred stop. Synchronous so the /// caller can take the snapshot while still ordered before any recreate /// (bug-hunt 2026-06-13, finding 4.2). - pub(crate) fn take_warm_instances(&self, function_name: &str) -> Vec> { - self.instances - .write() - .remove(function_name) - .unwrap_or_default() + /// + /// `function_arn` is the unqualified function ARN; `version` limits the + /// snapshot to one version's pool (DeleteFunction with a Qualifier), + /// `None` takes every version's. + pub(crate) fn take_warm_instances( + &self, + function_arn: &str, + version: Option<&str>, + ) -> Vec> { + let mut map = self.instances.write(); + let keys: Vec = map + .keys() + .filter(|k| match version { + Some(v) => k.as_str() == format!("{function_arn}:{v}"), + None => pool_key_belongs_to(k, function_arn), + }) + .cloned() + .collect(); + keys.into_iter() + .flat_map(|k| map.remove(&k).unwrap_or_default()) + .collect() } /// Terminate a previously-snapshotted set of warm instances. Pairs with @@ -632,13 +911,14 @@ impl LambdaRuntime { handle = ?entry.instance.handle, "stopping Lambda runtime instance" ); - self.backend.terminate(&entry.instance.handle).await; + self.retire(&entry).await; } } - /// Stop and remove every warm instance for a specific function. - pub async fn stop_container(&self, function_name: &str) { - let pool = self.take_warm_instances(function_name); + /// Stop and remove every warm instance (all versions) of the function + /// with unqualified ARN `function_arn`. + pub async fn stop_container(&self, function_arn: &str) { + let pool = self.take_warm_instances(function_arn, None); self.terminate_instances(pool).await; } @@ -653,7 +933,7 @@ impl LambdaRuntime { handle = ?entry.instance.handle, "stopping Lambda runtime instance (cleanup)" ); - self.backend.terminate(&entry.instance.handle).await; + self.retire(&entry).await; } } } @@ -668,10 +948,12 @@ impl LambdaRuntime { let entries = self.instances.read(); let accounts = lambda_state.read(); let mut rows = Vec::new(); - for (name, pool) in entries.iter() { + for (key, pool) in entries.iter() { + let name = pool_key_function_name(key); let runtime = accounts - .iter() - .find_map(|(_, state)| state.functions.get(name).map(|f| f.runtime.clone())) + .get(pool_key_account(key)) + .and_then(|state| state.functions.get(name)) + .map(|f| f.runtime.clone()) .unwrap_or_default(); for entry in pool { let idle_secs = entry.last_used.read().elapsed().as_secs(); @@ -700,14 +982,21 @@ impl LambdaRuntime { rows } - /// Evict (stop and remove) every warm instance for a specific - /// function. Returns true if at least one instance was evicted. + /// Evict (stop and remove) every warm instance of every function named + /// `function_name` (any account, any version). Returns true if at least + /// one instance was evicted. pub async fn evict_container(&self, function_name: &str) -> bool { - let pool = self - .instances - .write() - .remove(function_name) - .unwrap_or_default(); + let pool: Vec> = { + let mut map = self.instances.write(); + let keys: Vec = map + .keys() + .filter(|k| pool_key_function_name(k) == function_name) + .cloned() + .collect(); + keys.into_iter() + .flat_map(|k| map.remove(&k).unwrap_or_default()) + .collect() + }; let found = !pool.is_empty(); for entry in pool { tracing::info!( @@ -715,7 +1004,7 @@ impl LambdaRuntime { handle = ?entry.instance.handle, "evicting Lambda runtime instance via simulation API" ); - self.backend.terminate(&entry.instance.handle).await; + self.retire(&entry).await; } found } @@ -740,7 +1029,8 @@ impl LambdaRuntime { for (name, pool) in map.iter_mut() { let mut i = 0; while i < pool.len() { - let idle = pool[i].last_used.read().elapsed() > ttl; + let idle = pool[i].last_used.read().elapsed() > ttl + || pool[i].retiring.load(Ordering::Acquire); let free = pool[i].busy.try_lock().is_ok(); if idle && free { out.push((name.clone(), pool.remove(i))); @@ -754,7 +1044,7 @@ impl LambdaRuntime { }; for (name, entry) in expired { tracing::info!(function = %name, "stopping idle Lambda runtime instance"); - self.backend.terminate(&entry.instance.handle).await; + self.retire(&entry).await; } } } @@ -779,7 +1069,7 @@ mod tests { } else { vec![] }; - let id = deploy_id_from(&code_sha256, &layers); + let id = deploy_id_from(&code_sha256, "{}", &layers); assert!( !id.contains('/') && !id.contains('+') && !id.contains('='), "deploy id {id:?} (seed {i}) is not URL-path-safe" @@ -792,11 +1082,12 @@ mod tests { #[test] fn deploy_id_is_stable() { let layers = vec![b"layer-a".to_vec(), b"layer-b".to_vec()]; - let a = deploy_id_from("abc123", &layers); - let b = deploy_id_from("abc123", &layers); + let a = deploy_id_from("abc123", "cfg", &layers); + let b = deploy_id_from("abc123", "cfg", &layers); assert_eq!(a, b); - assert_ne!(a, deploy_id_from("abc124", &layers)); - assert_ne!(a, deploy_id_from("abc123", &[])); + assert_ne!(a, deploy_id_from("abc124", "cfg", &layers)); + assert_ne!(a, deploy_id_from("abc123", "cfg", &[])); + assert_ne!(a, deploy_id_from("abc123", "cfg2", &layers)); } // ---- warm-pool concurrency + eviction (issue #1644) ---- @@ -804,6 +1095,7 @@ mod tests { use super::LambdaRuntime; use crate::runtime::backend::{BackendHandle, LambdaBackend, RuntimeError, WarmInstance}; use crate::state::LambdaFunction; + use fakecloud_core::auth::{SessionCredentialIssuer, SessionCredentials}; use parking_lot::RwLock; use std::collections::{HashMap, VecDeque}; use std::sync::atomic::{AtomicUsize, Ordering::SeqCst}; @@ -819,6 +1111,8 @@ mod tests { default_endpoint: String, launches: AtomicUsize, terminates: AtomicUsize, + /// Access key id of the credentials each launch was handed. + launched_with: StdMutex>>, } impl CountingBackend { @@ -828,6 +1122,7 @@ mod tests { default_endpoint: default_endpoint.into(), launches: AtomicUsize::new(0), terminates: AtomicUsize::new(0), + launched_with: StdMutex::new(Vec::new()), }) } fn with_queue(default_endpoint: impl Into, queue: Vec) -> Arc { @@ -836,6 +1131,7 @@ mod tests { default_endpoint: default_endpoint.into(), launches: AtomicUsize::new(0), terminates: AtomicUsize::new(0), + launched_with: StdMutex::new(Vec::new()), }) } } @@ -851,7 +1147,12 @@ mod tests { _code_zip: Option<&[u8]>, _layers: &[Vec], _deploy_id: &str, + credentials: Option<&SessionCredentials>, ) -> Result { + self.launched_with + .lock() + .unwrap() + .push(credentials.map(|c| c.access_key_id.clone())); let n = self.launches.fetch_add(1, SeqCst); let endpoint = self .endpoints @@ -921,9 +1222,18 @@ mod tests { instances: RwLock::new(HashMap::new()), starting: RwLock::new(HashMap::new()), max_concurrency, + credential_issuer: std::sync::OnceLock::new(), }) } + fn arn(name: &str) -> String { + format!("arn:aws:lambda:us-east-1:123456789012:function:{name}") + } + + fn key(name: &str) -> String { + format!("{}:$LATEST", arn(name)) + } + fn test_func(name: &str, sha: &str) -> LambdaFunction { serde_json::from_value(serde_json::json!({ "function_name": name, @@ -1114,7 +1424,7 @@ mod tests { "the stale-deploy instance should have been torn down" ); // Exactly one current instance remains in the pool. - let pool_len = rt.instances.read().get("upd").map_or(0, |v| v.len()); + let pool_len = rt.instances.read().get(&key("upd")).map_or(0, |v| v.len()); assert_eq!(pool_len, 1); } @@ -1135,25 +1445,23 @@ mod tests { }, last_used: RwLock::new(std::time::Instant::now()), deploy_id: "d".to_string(), + credentials: None, + retiring: std::sync::atomic::AtomicBool::new(false), busy: Arc::new(tokio::sync::Mutex::new(())), }) }; // A function "f" with one warm instance. - rt.instances - .write() - .insert("f".to_string(), vec![mk("old")]); + rt.instances.write().insert(key("f"), vec![mk("old")]); // Delete snapshots the pool synchronously and removes it from the map. - let snapshot = rt.take_warm_instances("f"); + let snapshot = rt.take_warm_instances(&arn("f"), None); assert_eq!(snapshot.len(), 1); - assert!(rt.instances.read().get("f").is_none()); + assert!(rt.instances.read().get(&key("f")).is_none()); // A recreate + warm-up of the same name wins the race ahead of the // deferred terminate. - rt.instances - .write() - .insert("f".to_string(), vec![mk("new")]); + rt.instances.write().insert(key("f"), vec![mk("new")]); // Terminating the snapshot must touch only the old instance. rt.terminate_instances(snapshot).await; @@ -1165,7 +1473,381 @@ mod tests { // The recreated function keeps its fresh warm instance. let pool = rt.instances.read(); - let f = pool.get("f").expect("recreated function pool must survive"); + let f = pool + .get(&key("f")) + .expect("recreated function pool must survive"); assert_eq!(f.len(), 1); } + + /// Issuer double: mints sequential keys with a configurable lifetime + /// and records what was issued and revoked. + struct RecordingIssuer { + lifetime: chrono::Duration, + issued: StdMutex>, + revoked: StdMutex>, + } + + impl RecordingIssuer { + fn new(lifetime: chrono::Duration) -> Arc { + Arc::new(Self { + lifetime, + issued: StdMutex::new(Vec::new()), + revoked: StdMutex::new(Vec::new()), + }) + } + } + + impl SessionCredentialIssuer for RecordingIssuer { + fn issue( + &self, + role_arn: &str, + session_name: &str, + _duration: chrono::Duration, + ) -> SessionCredentials { + let mut issued = self.issued.lock().unwrap(); + let key = format!("KEY{}", issued.len()); + issued.push((key.clone(), role_arn.to_string(), session_name.to_string())); + SessionCredentials { + access_key_id: key, + secret_access_key: "s".into(), + session_token: "t".into(), + expiration: chrono::Utc::now() + self.lifetime, + account_id: "123456789012".into(), + } + } + fn revoke(&self, credentials: &SessionCredentials) { + self.revoked + .lock() + .unwrap() + .push(credentials.access_key_id.clone()); + } + } + + /// Each instance is launched with a session for the function's execution + /// role, named after the function, and the session is revoked when the + /// instance goes away. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn instances_get_execution_role_credentials_revoked_on_teardown() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(5), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + let issuer = RecordingIssuer::new(chrono::Duration::hours(1)); + rt.set_credential_issuer(issuer.clone()); + + let func = test_func("creds", "sha-A"); + rt.invoke(&func, b"{}", &[]).await.unwrap(); + rt.invoke(&func, b"{}", &[]).await.unwrap(); + + assert_eq!( + *issuer.issued.lock().unwrap(), + vec![( + "KEY0".to_string(), + "arn:aws:iam::123456789012:role/r".to_string(), + "creds".to_string() + )], + "a warm instance keeps its credentials across invocations" + ); + assert_eq!( + *backend.launched_with.lock().unwrap(), + vec![Some("KEY0".to_string())] + ); + + rt.stop_container(&arn("creds")).await; + assert_eq!(*issuer.revoked.lock().unwrap(), vec!["KEY0".to_string()]); + } + + /// An IAM reset drops the credentials warm instances hold: free ones are + /// retired (and replaced on the next invoke), busy ones finish their + /// invocation untouched and take no new work. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn iam_reset_retires_free_instances_and_spares_busy_ones() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(300), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + let issuer = RecordingIssuer::new(chrono::Duration::hours(12)); + rt.set_credential_issuer(issuer.clone()); + + // Two instances: one left free, one kept busy across the reset. + let func = test_func("reset", "sha-A"); + let (a, b) = tokio::join!(rt.invoke(&func, b"{}", &[]), rt.invoke(&func, b"{}", &[])); + a.unwrap(); + b.unwrap(); + assert_eq!(backend.launches.load(SeqCst), 2); + let busy = { + let rt = rt.clone(); + let func = func.clone(); + tokio::spawn(async move { rt.invoke(&func, b"{}", &[]).await }) + }; + tokio::time::sleep(Duration::from_millis(100)).await; + + // Marking is synchronous: from here on no invocation can land on + // either instance, before any teardown has run. + rt.mark_credentials_revoked(Some("123456789012")); + assert!(rt + .instances + .read() + .get(&key("reset")) + .unwrap() + .iter() + .all(|e| e.retiring.load(std::sync::atomic::Ordering::Acquire))); + assert_eq!(backend.terminates.load(SeqCst), 0); + // Another account's reset marks nothing further. + rt.mark_credentials_revoked(Some("999999999999")); + rt.retire_released().await; + assert_eq!( + backend.terminates.load(SeqCst), + 1, + "only the free instance is retired" + ); + busy.await + .unwrap() + .expect("the in-flight invocation completes"); + + // The next invoke never lands on the retiring instance: it gets a + // fresh one with new credentials, and the now-free survivor goes. + rt.invoke(&func, b"{}", &[]).await.unwrap(); + assert_eq!(backend.launches.load(SeqCst), 3); + assert_eq!( + backend + .launched_with + .lock() + .unwrap() + .last() + .cloned() + .flatten(), + Some("KEY2".to_string()) + ); + assert_eq!(backend.terminates.load(SeqCst), 2); + assert_eq!(issuer.revoked.lock().unwrap().len(), 2); + } + + /// A deploy change never terminates an instance mid-invocation: the busy + /// stale instance finishes and is only reaped once free. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn deploy_change_spares_busy_instance() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(300), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + + let old = test_func("busy", "sha-A"); + let in_flight = { + let rt = rt.clone(); + tokio::spawn(async move { rt.invoke(&old, b"{}", &[]).await }) + }; + tokio::time::sleep(Duration::from_millis(100)).await; + rt.invoke(&test_func("busy", "sha-B"), b"{}", &[]) + .await + .unwrap(); + assert_eq!( + backend.terminates.load(SeqCst), + 0, + "the busy stale instance must not be terminated" + ); + in_flight + .await + .unwrap() + .expect("in-flight invocation completes"); + + // Next launch-path sweep reaps the now-free stale instance. + rt.evict_stale_deploy( + &key("busy"), + &super::deploy_id_for(&test_func("busy", "sha-B"), &[], false), + ) + .await; + assert_eq!(backend.terminates.load(SeqCst), 1); + } + + /// Versions (and accounts) have their own pools and never evict each + /// other, even with identical code. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn versions_do_not_evict_each_other() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(5), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + + let latest = test_func("ver", "sha-A"); + let mut v1 = latest.clone(); + v1.version = "1".into(); + for _ in 0..2 { + rt.invoke(&latest, b"{}", &[]).await.unwrap(); + rt.invoke(&v1, b"{}", &[]).await.unwrap(); + } + assert_eq!(backend.launches.load(SeqCst), 2); + assert_eq!(backend.terminates.load(SeqCst), 0); + + // Deleting one version's pool leaves the other's. + let taken = rt.take_warm_instances(&arn("ver"), Some("1")); + assert_eq!(taken.len(), 1); + assert!(rt.instances.read().contains_key(&key("ver"))); + assert_eq!(rt.take_warm_instances(&arn("ver"), None).len(), 1); + } + + /// A configuration change that reaches the instance environment (here + /// the function's variables) starts a fresh instance even though the + /// code is unchanged. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn configuration_change_replaces_the_warm_instance() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(5), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + + let mut func = test_func("cfg", "sha-A"); + rt.invoke(&func, b"{}", &[]).await.unwrap(); + func.environment.insert("MODE".into(), "b".into()); + rt.invoke(&func, b"{}", &[]).await.unwrap(); + + assert_eq!(backend.launches.load(SeqCst), 2); + assert_eq!(backend.terminates.load(SeqCst), 1); + } + + #[test] + fn session_is_minted_in_the_function_account() { + let f = "arn:aws:lambda:us-east-1:123456789012:function:f"; + assert_eq!( + super::session_role_arn("arn:aws:iam::123456789012:role/r", f), + "arn:aws:iam::123456789012:role/r" + ); + assert_eq!( + super::session_role_arn("arn:aws:iam::000000000000:role/path/r", f), + "arn:aws:iam::123456789012:role/path/r" + ); + assert_eq!(super::session_role_arn("not-an-arn", f), "not-an-arn"); + } + + /// Backend double whose launches fail, or never finish. + struct BrokenBackend { + hang: bool, + } + + #[async_trait::async_trait] + impl LambdaBackend for BrokenBackend { + fn name(&self) -> &str { + "broken" + } + async fn launch( + &self, + _func: &LambdaFunction, + _code_zip: Option<&[u8]>, + _layers: &[Vec], + _deploy_id: &str, + _credentials: Option<&SessionCredentials>, + ) -> Result { + if self.hang { + std::future::pending::<()>().await; + } + Err(RuntimeError::ContainerStartFailed("boom".into())) + } + async fn terminate(&self, _handle: &BackendHandle) {} + } + + fn broken_runtime(hang: bool) -> Arc { + Arc::new(LambdaRuntime { + backend: Arc::new(BrokenBackend { hang }), + instances: RwLock::new(HashMap::new()), + starting: RwLock::new(HashMap::new()), + max_concurrency: 2, + credential_issuer: std::sync::OnceLock::new(), + }) + } + + /// Credentials minted for a launch that fails, or whose launching future + /// is dropped mid-launch, are revoked rather than left registered. + #[tokio::test] + async fn credentials_of_an_unfinished_launch_are_revoked() { + let func = test_func("broken", "sha-A"); + + let rt = broken_runtime(false); + let issuer = RecordingIssuer::new(chrono::Duration::hours(12)); + rt.set_credential_issuer(issuer.clone()); + assert!(rt.invoke(&func, b"{}", &[]).await.is_err()); + assert_eq!(*issuer.revoked.lock().unwrap(), vec!["KEY0".to_string()]); + + let rt = broken_runtime(true); + let issuer = RecordingIssuer::new(chrono::Duration::hours(12)); + rt.set_credential_issuer(issuer.clone()); + let cancelled = + tokio::time::timeout(Duration::from_millis(100), rt.invoke(&func, b"{}", &[])).await; + assert!(cancelled.is_err(), "the launch never finishes"); + assert_eq!(*issuer.revoked.lock().unwrap(), vec!["KEY0".to_string()]); + } + + /// An instance stops taking invocations once its credentials could lapse + /// during one, judged from the credentials' own expiration. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn instance_whose_credentials_near_expiry_is_replaced() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(5), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + // Minted with less than the invocation headroom left. + let issuer = RecordingIssuer::new(chrono::Duration::minutes(10)); + rt.set_credential_issuer(issuer.clone()); + + let func = test_func("aging", "sha-A"); + rt.invoke(&func, b"{}", &[]).await.unwrap(); + rt.invoke(&func, b"{}", &[]).await.unwrap(); + assert_eq!(backend.launches.load(SeqCst), 2); + // The first (free) one was retired on the second launch. + assert_eq!(*issuer.revoked.lock().unwrap(), vec!["KEY0".to_string()]); + } + + /// Deleting a version stops its free instances now and leaves a busy one + /// to finish, retiring it once released; other versions are untouched. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn deleting_a_version_spares_its_busy_instance() { + let peak = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_rie(Duration::from_millis(300), peak).await; + let backend = CountingBackend::new(endpoint); + let rt = runtime_with(backend.clone(), 2); + + let latest = test_func("delv", "sha-A"); + let mut v1 = latest.clone(); + v1.version = "1".into(); + rt.invoke(&latest, b"{}", &[]).await.unwrap(); + let in_flight = { + let rt = rt.clone(); + let v1 = v1.clone(); + tokio::spawn(async move { rt.invoke(&v1, b"{}", &[]).await }) + }; + tokio::time::sleep(Duration::from_millis(100)).await; + + let detached = rt.retire_version(&arn("delv"), "1"); + assert!(detached.is_empty(), "the busy instance is not detached"); + rt.terminate_instances(detached).await; + assert_eq!(backend.terminates.load(SeqCst), 0); + in_flight + .await + .unwrap() + .expect("in-flight invocation completes"); + + rt.retire_released().await; + assert_eq!(backend.terminates.load(SeqCst), 1); + assert!(rt + .instances + .read() + .get(&format!("{}:1", arn("delv"))) + .is_none()); + assert!(rt.instances.read().contains_key(&key("delv"))); + } + + /// Tags only change the deploy fingerprint for backends that read them. + #[test] + fn tags_change_the_deploy_only_when_the_backend_uses_them() { + let plain = test_func("tagged", "sha-A"); + let mut tagged = plain.clone(); + tagged.tags.insert("team".into(), "a".into()); + assert_eq!( + super::deploy_id_for(&plain, &[], false), + super::deploy_id_for(&tagged, &[], false) + ); + assert_ne!( + super::deploy_id_for(&plain, &[], true), + super::deploy_id_for(&tagged, &[], true) + ); + } } diff --git a/crates/fakecloud-lambda/src/runtime/k8s/mod.rs b/crates/fakecloud-lambda/src/runtime/k8s/mod.rs index a32fd43cf..a3a6a762f 100644 --- a/crates/fakecloud-lambda/src/runtime/k8s/mod.rs +++ b/crates/fakecloud-lambda/src/runtime/k8s/mod.rs @@ -18,11 +18,12 @@ pub mod spec; use std::time::Duration; use async_trait::async_trait; +use fakecloud_core::auth::SessionCredentials; use fakecloud_k8s::{K8sClient, K8sEnv, K8sEnvError, K8sPodConfig, K8sPodConfigError}; use super::backend::{BackendHandle, LambdaBackend, RuntimeError, WarmInstance}; use crate::state::LambdaFunction; -use spec::{build_pod_spec, unique_pod_name, PodSpecContext}; +use spec::{build_credentials_secret, build_pod_spec, unique_pod_name, PodSpecContext}; /// Which `fakecloud-service` label Lambda Pods carry, so reaping only /// touches Lambda Pods. @@ -109,7 +110,7 @@ impl K8sBackend { /// Extract the account ID from a function ARN /// (`arn:aws:lambda:::function:[:]`). fn account_id_from_arn(arn: &str) -> &str { - arn.split(':').nth(4).unwrap_or("000000000000") + fakecloud_aws::arn::account_of(arn).unwrap_or("000000000000") } #[async_trait] @@ -118,14 +119,25 @@ impl LambdaBackend for K8sBackend { "kubernetes" } + /// Per-function tags carry scheduling overrides (`K8sPodConfig::from_tags`). + fn launch_uses_tags(&self) -> bool { + true + } + async fn launch( &self, func: &LambdaFunction, _code_zip: Option<&[u8]>, _layers: &[Vec], deploy_id: &str, + credentials: Option<&SessionCredentials>, ) -> Result { let account_id = account_id_from_arn(&func.function_arn); + // A per-launch unique name (instead of the deterministic + // function+deploy one) so concurrent instances of the same function + // don't collide and a terminating Pod never blocks its replacement + // (see `unique_pod_name`). The credentials Secret shares it. + let pod_name = unique_pod_name(&func.function_name, deploy_id); let ctx = PodSpecContext { instance_id: self.client.instance_id(), namespace: self.client.namespace(), @@ -136,14 +148,10 @@ impl LambdaBackend for K8sBackend { internal_token: &self.internal_token, account_id, pull_secret: self.pull_secret.as_deref(), + credentials_secret: credentials.map(|_| pod_name.as_str()), }; let mut pod = build_pod_spec(func, deploy_id, &ctx).map_err(RuntimeError::ContainerStartFailed)?; - // Override the deterministic function+deploy name with a per-launch - // unique one so concurrent instances of the same function don't collide - // and a terminating Pod never blocks its replacement (see - // `unique_pod_name`). - let pod_name = unique_pod_name(&func.function_name, deploy_id); pod.metadata.name = Some(pod_name.clone()); // Apply operator-configured scheduling/metadata: global + @@ -154,10 +162,26 @@ impl LambdaBackend for K8sBackend { .merge(K8sPodConfig::from_tags(&func.tags)) .apply(&mut pod); - self.client - .create_pod(&pod) - .await - .map_err(|e| RuntimeError::ContainerStartFailed(format!("k8s create pod: {e}")))?; + // The Secret goes first so the Pod's secretKeyRefs resolve at start. + if let Some(creds) = credentials { + let secret = build_credentials_secret(&pod_name, creds, self.client.instance_id()); + self.client.create_secret(&secret).await.map_err(|e| { + RuntimeError::ContainerStartFailed(format!( + "k8s create credentials secret (the ServiceAccount needs \ + create/delete/patch on secrets): {e}" + )) + })?; + } + if let Err(e) = self.client.create_pod(&pod).await { + self.client.delete_secret(&pod_name).await; + return Err(RuntimeError::ContainerStartFailed(format!( + "k8s create pod: {e}" + ))); + } + if credentials.is_some() { + // Owned by the Pod: garbage-collected with it however it goes. + self.client.adopt_secret(&pod_name, &pod_name).await; + } // Tear the Pod down again if it never becomes ready, so a failed // launch doesn't leak a Pod. @@ -169,6 +193,7 @@ impl LambdaBackend for K8sBackend { Ok(ip) => ip, Err(e) => { self.client.delete_pod(&pod_name).await; + self.client.delete_secret(&pod_name).await; return Err(RuntimeError::ContainerStartFailed(e.to_string())); } }; @@ -176,6 +201,7 @@ impl LambdaBackend for K8sBackend { // is listening yet — TCP-handshake the invoke port like Docker. if let Err(e) = K8sClient::wait_for_tcp(&pod_ip, 8080, Duration::from_secs(10)).await { self.client.delete_pod(&pod_name).await; + self.client.delete_secret(&pod_name).await; return Err(RuntimeError::ContainerStartFailed(format!( "RIE on {pod_ip}:8080 not ready: {e}" ))); @@ -200,7 +226,12 @@ impl LambdaBackend for K8sBackend { async fn terminate(&self, handle: &BackendHandle) { match handle { - BackendHandle::Pod { name, .. } => self.client.delete_pod(name).await, + BackendHandle::Pod { name, .. } => { + self.client.delete_pod(name).await; + // Also garbage-collected with the Pod; deleting explicitly + // drops the credentials without waiting on the collector. + self.client.delete_secret(name).await; + } // Docker handles aren't ours to manage — defensive no-op. BackendHandle::Container { .. } => {} } diff --git a/crates/fakecloud-lambda/src/runtime/k8s/spec.rs b/crates/fakecloud-lambda/src/runtime/k8s/spec.rs index 00902825d..fae44da48 100644 --- a/crates/fakecloud-lambda/src/runtime/k8s/spec.rs +++ b/crates/fakecloud-lambda/src/runtime/k8s/spec.rs @@ -6,8 +6,8 @@ use std::collections::BTreeMap; use k8s_openapi::api::core::v1::{ - Container, EmptyDirVolumeSource, EnvVar, LocalObjectReference, Pod, PodSpec, - ResourceRequirements, Volume, VolumeMount, + Container, EmptyDirVolumeSource, EnvVar, EnvVarSource, LocalObjectReference, Pod, PodSpec, + ResourceRequirements, Secret, SecretKeySelector, Volume, VolumeMount, }; use k8s_openapi::apimachinery::pkg::api::resource::Quantity; use k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta; @@ -15,8 +15,9 @@ use k8s_openapi::apimachinery::pkg::apis::meta::v1::ObjectMeta; use fakecloud_k8s::names::label_safe; use super::super::docker::runtime_to_image; -use super::super::env_rewrite::rewrite_localhost_envs; +use super::super::environment::{credential_envs, function_environment, CREDENTIAL_ENV_KEYS}; use crate::state::LambdaFunction; +use fakecloud_core::auth::SessionCredentials; /// Inputs that don't come from the function itself — instance identity, /// the in-cluster fakecloud URL, the bearer token the init container @@ -45,6 +46,9 @@ pub struct PodSpecContext<'a> { /// `kubernetes.io/dockerconfigjson` used as `imagePullSecrets` for /// container-image functions. pub pull_secret: Option<&'a str>, + /// Name of the Secret holding this instance's execution-role + /// credentials (see [`build_credentials_secret`]); `None` exports none. + pub credentials_secret: Option<&'a str>, } /// Resource cap for the `/tmp` `emptyDir` (`medium: Memory`) — sized @@ -96,9 +100,14 @@ pub fn build_pod_spec( label_safe(&deploy_id[..deploy_id.len().min(40)]), ); - // Env vars — user-supplied (with localhost rewritten to in-cluster - // fakecloud host) + the AWS_LAMBDA_FUNCTION_TIMEOUT the RIE honors. - let mut env: Vec = rewrite_localhost_envs(&func.environment, ctx.self_host) + // The function's execution environment: the reserved Lambda variables + // plus the user's own (localhost rewritten to the in-cluster fakecloud + // host), with SDK calls pointed at the same in-cluster fakecloud URL the + // init container fetches code from. + // Credentials never sit in the Pod spec: they are read from the + // per-Pod Secret (see [`build_credentials_secret`]). + let endpoint_url = ctx.self_url.trim_end_matches('/'); + let mut env: Vec = function_environment(func, endpoint_url, ctx.self_host, None) .into_iter() .map(|(k, v)| EnvVar { name: k, @@ -106,11 +115,20 @@ pub fn build_pod_spec( value_from: None, }) .collect(); - env.push(EnvVar { - name: "AWS_LAMBDA_FUNCTION_TIMEOUT".into(), - value: Some(func.timeout.to_string()), - value_from: None, - }); + if let Some(secret) = ctx.credentials_secret { + env.extend(CREDENTIAL_ENV_KEYS.iter().map(|key| EnvVar { + name: (*key).to_string(), + value: None, + value_from: Some(EnvVarSource { + secret_key_ref: Some(SecretKeySelector { + name: secret.to_string(), + key: (*key).to_string(), + optional: Some(false), + }), + ..EnvVarSource::default() + }), + })); + } // Shared volumes — init container writes /var/task, /opt, layers; // main container reads them. `/tmp` is sized from ephemeral_storage. @@ -259,6 +277,39 @@ pub fn build_pod_spec( }) } +/// The Secret carrying one instance's execution-role credentials, named +/// `name` (the instance's Pod name) and labelled like the Pod so the same +/// ownership rules apply. The Pod references it through `secretKeyRef`, so +/// the credentials never appear in the Pod spec. +pub fn build_credentials_secret( + name: &str, + credentials: &SessionCredentials, + instance_id: &str, +) -> Secret { + let mut labels = BTreeMap::new(); + labels.insert( + fakecloud_k8s::labels::MANAGED_BY.into(), + fakecloud_k8s::labels::MANAGED_BY_VALUE.into(), + ); + labels.insert(fakecloud_k8s::labels::INSTANCE.into(), instance_id.into()); + labels.insert(fakecloud_k8s::labels::SERVICE.into(), super::SERVICE.into()); + Secret { + metadata: ObjectMeta { + name: Some(name.to_string()), + labels: Some(labels), + ..ObjectMeta::default() + }, + type_: Some("Opaque".into()), + string_data: Some( + credential_envs(credentials) + .into_iter() + .map(|(k, v)| (k.to_string(), v)) + .collect(), + ), + ..Secret::default() + } +} + /// Build a deterministic, DNS-1123-safe Pod name for the given /// function + deploy_id. Truncated/lowercased so it fits the 63-char /// label limit. The suffix is a stable hash of `deploy_id` so a new @@ -322,6 +373,7 @@ mod tests { internal_token: "secret-token-xyz", account_id: "000000000000", pull_secret: None, + credentials_secret: None, } } @@ -471,6 +523,81 @@ mod tests { assert_eq!(t, "30"); } + #[test] + fn main_container_gets_the_aws_environment_pointed_at_the_cluster() { + let f = zip_function("my-fn"); + let mut c = ctx(); + c.credentials_secret = Some("my-fn-pod"); + let pod = build_pod_spec(&f, "d", &c).unwrap(); + let env = pod.spec.unwrap().containers[0].env.clone().unwrap(); + let get = |k: &str| { + env.iter() + .find(|e| e.name == k) + .and_then(|e| e.value.clone()) + }; + assert_eq!( + get("AWS_ENDPOINT_URL").as_deref(), + Some("http://fakecloud.fakecloud.svc.cluster.local:4566") + ); + assert_eq!(get("AWS_REGION").as_deref(), Some("us-east-1")); + assert_eq!(get("AWS_LAMBDA_FUNCTION_NAME").as_deref(), Some("my-fn")); + + // Credentials come from the per-Pod Secret, never inline values. + for key in CREDENTIAL_ENV_KEYS { + let var = env.iter().find(|e| e.name == key).expect(key); + assert_eq!(var.value, None, "{key} must not be inline"); + let r = var + .value_from + .as_ref() + .and_then(|v| v.secret_key_ref.as_ref()) + .expect("secretKeyRef"); + assert_eq!(r.name, "my-fn-pod"); + assert_eq!(r.key, key); + } + // Kubernetes keeps duplicate names ambiguous; each is set once. + let mut names: Vec<&str> = env.iter().map(|e| e.name.as_str()).collect(); + let total = names.len(); + names.sort_unstable(); + names.dedup(); + assert_eq!(names.len(), total); + } + + #[test] + fn no_credential_vars_without_a_secret() { + let pod = build_pod_spec(&zip_function("my-fn"), "d", &ctx()).unwrap(); + let env = pod.spec.unwrap().containers[0].env.clone().unwrap(); + assert!(env + .iter() + .all(|e| !CREDENTIAL_ENV_KEYS.contains(&e.name.as_str()))); + } + + #[test] + fn credentials_secret_carries_every_referenced_key() { + let creds = SessionCredentials { + access_key_id: "FSIAEXAMPLE".into(), + secret_access_key: "secret".into(), + session_token: "token".into(), + expiration: Utc::now(), + account_id: "000000000000".into(), + }; + let secret = build_credentials_secret("my-fn-pod", &creds, "fakecloud-1234"); + assert_eq!(secret.metadata.name.as_deref(), Some("my-fn-pod")); + let labels = secret.metadata.labels.unwrap(); + assert_eq!( + labels + .get(fakecloud_k8s::labels::INSTANCE) + .map(String::as_str), + Some("fakecloud-1234") + ); + let data = secret.string_data.unwrap(); + for key in CREDENTIAL_ENV_KEYS { + assert!(data.contains_key(key), "{key}"); + } + assert_eq!(data["AWS_ACCESS_KEY_ID"], "FSIAEXAMPLE"); + assert_eq!(data["AWS_SECRET_ACCESS_KEY"], "secret"); + assert_eq!(data["AWS_SESSION_TOKEN"], "token"); + } + #[test] fn ephemeral_storage_maps_to_emptydir_memory_with_size_limit() { let f = zip_function("my-fn"); diff --git a/crates/fakecloud-lambda/src/runtime/mod.rs b/crates/fakecloud-lambda/src/runtime/mod.rs index b604a58cc..9295b3953 100644 --- a/crates/fakecloud-lambda/src/runtime/mod.rs +++ b/crates/fakecloud-lambda/src/runtime/mod.rs @@ -9,6 +9,7 @@ pub(crate) mod backend; pub(crate) mod docker; pub(crate) mod env_rewrite; +pub mod environment; pub(crate) mod facade; pub mod k8s; diff --git a/crates/fakecloud-lambda/src/service/functions.rs b/crates/fakecloud-lambda/src/service/functions.rs index 5e6e44f0c..334b5e03b 100644 --- a/crates/fakecloud-lambda/src/service/functions.rs +++ b/crates/fakecloud-lambda/src/service/functions.rs @@ -26,21 +26,14 @@ impl LambdaService { )); } - // PassRole trust-policy check: the supplied execution role must - // have a trust policy that allows lambda.amazonaws.com to call - // sts:AssumeRole. Real AWS rejects with InvalidParameterValueException - // when the trust policy doesn't include the service principal. - if let Some(ref validator) = self.role_trust_validator { - if let Err(err) = - validator.validate(&req.account_id, &input.role, "lambda.amazonaws.com") - { - return Err(AwsServiceError::aws_error( - StatusCode::BAD_REQUEST, - "InvalidParameterValueException", - err.to_string(), - )); - } - } + // PassRole: a role whose trust policy lets Lambda assume it, and (under + // IAM enforcement) in the caller's account. + super::validate_execution_role( + &req.account_id, + &input.role, + self.role_trust_validator.as_deref(), + self.iam_mode, + )?; let mut accounts = self.state.write(); // Pre-resolve layer attachments before re-borrowing accounts mutably. @@ -335,10 +328,22 @@ impl LambdaService { if let Some(list) = state.function_versions.get_mut(function_name) { list.retain(|v| v != q); } + let live_arn = state + .functions + .get(function_name) + .map(|f| f.function_arn.clone()); + drop(accounts); + // Stop the deleted version's warm instances: free ones now, busy + // ones once their in-flight invocation completes. + if let (Some(runtime), Some(live_arn)) = (&self.runtime, live_arn) { + let rt = runtime.clone(); + let pool = rt.retire_version(&live_arn, q); + tokio::spawn(async move { rt.terminate_instances(pool).await }); + } return Ok(AwsResponse::json(StatusCode::NO_CONTENT, "")); } - if state.functions.remove(function_name).is_none() { + let Some(removed) = state.functions.remove(function_name) else { return Err(AwsServiceError::aws_error( StatusCode::NOT_FOUND, "ResourceNotFoundException", @@ -347,7 +352,7 @@ impl LambdaService { function_arn(region, &account_id_owned, function_name) ), )); - } + }; // Drop all numbered versions + their snapshots so the function // is gone end-to-end (AWS deletes everything when no Qualifier // is supplied). @@ -374,7 +379,7 @@ impl LambdaService { // function keeps its new container. if let Some(ref runtime) = self.runtime { let rt = runtime.clone(); - let pool = rt.take_warm_instances(function_name); + let pool = rt.take_warm_instances(&removed.function_arn, None); tokio::spawn(async move { rt.terminate_instances(pool).await }); } diff --git a/crates/fakecloud-lambda/src/service/mod.rs b/crates/fakecloud-lambda/src/service/mod.rs index a36265b38..57b9ab50c 100644 --- a/crates/fakecloud-lambda/src/service/mod.rs +++ b/crates/fakecloud-lambda/src/service/mod.rs @@ -562,6 +562,75 @@ pub(crate) fn validate_ephemeral_storage(size: i64) -> Result, +) -> Result<(), AwsServiceError> { + match crate::runtime::environment::reserved_keys_message(environment) { + None => Ok(()), + Some(message) => Err(AwsServiceError::aws_error( + StatusCode::BAD_REQUEST, + "InvalidParameterValueException", + message, + )), + } +} + +/// The execution-role checks `CreateFunction` and `UpdateFunctionConfiguration` +/// run before accepting a role, shared with CloudFormation's +/// `AWS::Lambda::Function`: +/// +/// - `iam:PassRole` is same-account only on AWS. Under `--iam strict` a role +/// owned by another account is refused with `AccessDeniedException`, +/// whatever its trust policy says; `--iam soft` logs the would-be denial to +/// the IAM audit target and allows it; with IAM off (the default) it is +/// accepted, so templates carrying another emulator's default account +/// (`000000000000`) keep working. The execution session is then minted in +/// the function's account. +/// - The role's trust policy must let `lambda.amazonaws.com` assume it, +/// looked up in the caller's (function's) account, where the session is +/// minted. Always applied, as it always was on `CreateFunction`. +pub fn validate_execution_role( + caller_account: &str, + role_arn: &str, + validator: Option<&dyn fakecloud_core::auth::RoleTrustValidator>, + iam_mode: fakecloud_core::auth::IamMode, +) -> Result<(), AwsServiceError> { + let cross_account = + fakecloud_aws::arn::account_of(role_arn).is_some_and(|a| a != caller_account); + if cross_account && iam_mode.is_enabled() { + tracing::warn!( + target: "fakecloud::iam::audit", + action = "iam:PassRole", + resource = %role_arn, + account = %caller_account, + mode = %iam_mode, + "cross-account pass role denied" + ); + if iam_mode.is_strict() { + return Err(AwsServiceError::aws_error( + StatusCode::FORBIDDEN, + "AccessDeniedException", + "Cross-account pass role is not allowed.", + )); + } + } + if let Some(validator) = validator { + validator + .validate(caller_account, role_arn, "lambda.amazonaws.com") + .map_err(|err| { + AwsServiceError::aws_error( + StatusCode::BAD_REQUEST, + "InvalidParameterValueException", + err.to_string(), + ) + })?; + } + Ok(()) +} + /// All fields of a `CreateFunction` request, already parsed and /// defaulted. The code zip (if any) is eagerly base64-decoded so the /// caller can hash it without doing the decode again. @@ -623,6 +692,7 @@ impl CreateFunctionInput { .collect() }) .unwrap_or_default(); + validate_environment(&environment)?; let architectures = body["Architectures"] .as_array() @@ -1006,6 +1076,9 @@ pub struct LambdaService { snapshot_lock: Arc>, pub(crate) delivery_bus: Option>, pub(crate) role_trust_validator: Option>, + /// IAM enforcement mode; a cross-account execution role is refused only + /// while it is on (see [`validate_execution_role`]). + pub(crate) iam_mode: fakecloud_core::auth::IamMode, pub(crate) s3_delivery: Option>, /// Per-account-per-function in-flight invocation count, used to /// gate `Invoke` against `PutFunctionConcurrency`'s @@ -1030,6 +1103,7 @@ impl LambdaService { snapshot_lock: Arc::new(AsyncMutex::new(())), delivery_bus: None, role_trust_validator: None, + iam_mode: fakecloud_core::auth::IamMode::Off, s3_delivery: None, inflight_invocations: Arc::new(parking_lot::RwLock::new(BTreeMap::new())), } @@ -1063,6 +1137,12 @@ impl LambdaService { self } + /// Apply the server's IAM enforcement mode to the execution-role checks. + pub fn with_iam_mode(mut self, mode: fakecloud_core::auth::IamMode) -> Self { + self.iam_mode = mode; + self + } + async fn save_snapshot(&self) { save_lambda_snapshot( &self.state, diff --git a/crates/fakecloud-lambda/src/service_tests.rs b/crates/fakecloud-lambda/src/service_tests.rs index 4cb91be87..07b9f0cde 100644 --- a/crates/fakecloud-lambda/src/service_tests.rs +++ b/crates/fakecloud-lambda/src/service_tests.rs @@ -3487,6 +3487,7 @@ impl crate::runtime::LambdaBackend for PrepullSpy { _code_zip: Option<&[u8]>, _layers: &[Vec], _deploy_id: &str, + _credentials: Option<&fakecloud_core::auth::SessionCredentials>, ) -> Result { unreachable!("prepull regression test never invokes launch") } @@ -4340,3 +4341,60 @@ async fn china_region_arns_use_the_aws_cn_partition() { resp.message() ); } + +/// Trust validator double that records the account each lookup ran in. +#[derive(Default)] +struct RecordingTrustValidator { + accounts: parking_lot::Mutex>, +} + +impl fakecloud_core::auth::RoleTrustValidator for RecordingTrustValidator { + fn validate( + &self, + account_id: &str, + _role_arn: &str, + _service_principal: &str, + ) -> Result<(), fakecloud_core::auth::PassRoleError> { + self.accounts.lock().push(account_id.to_string()); + Ok(()) + } +} + +#[test] +fn cross_account_execution_role_is_refused_only_under_strict_iam() { + use fakecloud_core::auth::IamMode; + let foreign = "arn:aws:iam::000000000000:role/r"; + for mode in [IamMode::Off, IamMode::Soft] { + validate_execution_role("123456789012", foreign, None, mode) + .unwrap_or_else(|e| panic!("{mode} must allow (soft only audits): {e}")); + } + let err = validate_execution_role("123456789012", foreign, None, IamMode::Strict) + .expect_err("strict refuses a cross-account role"); + assert_eq!( + err.to_string(), + "AccessDeniedException: Cross-account pass role is not allowed." + ); + // A same-account role is fine in every mode. + validate_execution_role( + "123456789012", + "arn:aws:iam::123456789012:role/r", + None, + IamMode::Strict, + ) + .unwrap(); +} + +#[test] +fn trust_check_runs_in_the_function_account() { + let validator = RecordingTrustValidator::default(); + validate_execution_role( + "123456789012", + "arn:aws:iam::000000000000:role/r", + Some(&validator), + fakecloud_core::auth::IamMode::Off, + ) + .unwrap(); + // The session is minted in the function's account, so that is where + // the role's trust policy is looked up. + assert_eq!(*validator.accounts.lock(), vec!["123456789012".to_string()]); +} diff --git a/crates/fakecloud-server/src/main.rs b/crates/fakecloud-server/src/main.rs index 619450057..80c83d232 100644 --- a/crates/fakecloud-server/src/main.rs +++ b/crates/fakecloud-server/src/main.rs @@ -376,6 +376,15 @@ async fn main() { // runtime is resolved) instead of scattering five separate log lines. let mut degraded_runtimes: Vec<&str> = Vec::new(); if let Some(ref rt) = container_runtime { + // Function code runs with its execution role's credentials, minted + // and registered like an AssumeRole session so its SDK calls resolve + // to the role (and verify under --verify-sigv4). + rt.set_credential_issuer( + fakecloud_iam::sts_service::container_creds::IamSessionCredentialIssuer::shared( + iam_state.clone(), + cli.account_id.clone(), + ), + ); tracing::info!(backend = rt.cli_name(), "Lambda execution enabled"); } else { degraded_runtimes.push("Lambda (Invoke returns errors for functions with code)"); @@ -1567,6 +1576,7 @@ async fn main() { appconfig: appconfig_state.clone(), delivery: delivery_for_cf, lambda_runtime: container_runtime.clone(), + iam_mode: cli.iam_mode(), rds_runtime: rds_runtime.clone(), ec2_runtime: ec2_runtime.clone(), ecs_runtime: ecs_runtime.clone(), @@ -2091,6 +2101,7 @@ async fn main() { lambda_service = lambda_service.with_role_trust_validator( fakecloud_iam::pass_role::IamRoleTrustValidator::shared(iam_state.clone()), ); + lambda_service = lambda_service.with_iam_mode(cli.iam_mode()); lambda_service = lambda_service.with_s3_delivery(s3_delivery_for_logs.clone()); if let Some(ref rt) = container_runtime { lambda_service = lambda_service.with_runtime(rt.clone()); diff --git a/crates/fakecloud-server/src/reset.rs b/crates/fakecloud-server/src/reset.rs index 13dfe15bf..a0d416d64 100644 --- a/crates/fakecloud-server/src/reset.rs +++ b/crates/fakecloud-server/src/reset.rs @@ -58,7 +58,17 @@ impl ResetState { pub(crate) fn reset_service(&self, service: &str) -> Result<(), String> { match service { "iam" | "sts" => { + // The reset drops the execution-role sessions warm Lambda + // instances hold: stop handing them invocations in the same + // step, then tear the free ones down in the background. + if let Some(ref rt) = self.container_runtime { + rt.mark_credentials_revoked(None); + } self.iam.write().reset(); + if let Some(ref rt) = self.container_runtime { + let rt = rt.clone(); + tokio::spawn(async move { rt.retire_released().await }); + } } "sqs" => { self.sqs.write().reset(); @@ -235,10 +245,19 @@ impl ResetState { ) -> Result<(), String> { match service { "iam" | "sts" => { - let mut mas = self.iam.write(); - let region = mas.region().to_string(); - if let Some(state) = mas.get_mut(account_id) { - state.reset(®ion); + if let Some(ref rt) = self.container_runtime { + rt.mark_credentials_revoked(Some(account_id)); + } + { + let mut mas = self.iam.write(); + let region = mas.region().to_string(); + if let Some(state) = mas.get_mut(account_id) { + state.reset(®ion); + } + } + if let Some(ref rt) = self.container_runtime { + let rt = rt.clone(); + tokio::spawn(async move { rt.retire_released().await }); } } "sqs" => { diff --git a/website/content/docs/guides/kubernetes-backend.md b/website/content/docs/guides/kubernetes-backend.md index 0d0c83673..91311315d 100644 --- a/website/content/docs/guides/kubernetes-backend.md +++ b/website/content/docs/guides/kubernetes-backend.md @@ -48,7 +48,7 @@ Optional env vars: ## RBAC -The fakecloud Pod's ServiceAccount needs permission to create / list / watch / delete Pods in the configured namespace. The ElastiCache backend additionally execs into cache Pods (`redis-cli` for CONFIG/ACL and snapshot SAVE), so it needs `pods/exec` too. +The fakecloud Pod's ServiceAccount needs permission to create / list / watch / delete Pods in the configured namespace. The ElastiCache backend additionally execs into cache Pods (`redis-cli` for CONFIG/ACL and snapshot SAVE), so it needs `pods/exec` too. The Lambda backend hands each function Pod its execution-role credentials through a per-Pod Secret (owned by the Pod, deleted with it), so it needs `create` / `patch` / `delete` on Secrets. ```yaml apiVersion: v1 @@ -66,6 +66,10 @@ rules: - apiGroups: [""] resources: ["pods"] verbs: ["create", "get", "list", "watch", "delete"] + # Required for the Lambda backend (per-Pod execution-role credentials). + - apiGroups: [""] + resources: ["secrets"] + verbs: ["create", "patch", "delete"] # Required only for the ElastiCache backend (exec into cache Pods). - apiGroups: [""] resources: ["pods/exec"] diff --git a/website/content/docs/services/lambda.md b/website/content/docs/services/lambda.md index 24aa0277f..805eb698f 100644 --- a/website/content/docs/services/lambda.md +++ b/website/content/docs/services/lambda.md @@ -13,14 +13,14 @@ fakecloud implements **73 of 73** Lambda operations at 100% Smithy conformance. - **23 runtimes** — Node.js (16/18/20/22/24), Python (3.8/3.9/3.10/3.11/3.12/3.13/3.14), Java (11/17/21/25), Go (1.x), Ruby (3.3/3.4), .NET (8/10), `provided.al2`, `provided.al2023` - **Event source mappings** — SQS, Kinesis, DynamoDB Streams polling loops with **`FilterCriteria`** (EventBridge-style JSON pattern, exists/prefix/suffix/equals-ignore-case/anything-but/numeric operators, SQS body decode), **`StartingPosition`** (`TRIM_HORIZON` / `LATEST` / `AT_TIMESTAMP` for Kinesis, `TRIM_HORIZON` / `LATEST` for DDB Streams), **`MaximumBatchingWindowInSeconds`** (SQS), and **`FunctionResponseTypes=[ReportBatchItemFailures]`** for SQS partial-batch failure semantics - **Layers** — create, publish, attach to functions; layer ZIP content is extracted into `/opt` of the runtime container at invoke time, so Python `import`, Node `require`, and `LD_LIBRARY_PATH` lookups resolve against attached layers exactly as on real AWS -- **Environment variables** — passed to the container +- **Environment variables**: the function's own variables plus the environment real Lambda provides: `AWS_REGION`/`AWS_DEFAULT_REGION` (the function's region), execution-role credentials (`AWS_ACCESS_KEY_ID`/`AWS_SECRET_ACCESS_KEY`/`AWS_SESSION_TOKEN`, a registered session for the function's role named after the function, so SDK calls sign as `assumed-role//` and verify under `--verify-sigv4`), and `AWS_LAMBDA_FUNCTION_NAME`/`_VERSION`/`_MEMORY_SIZE`, `AWS_LAMBDA_LOG_GROUP_NAME`/`_LOG_STREAM_NAME`, `AWS_EXECUTION_ENV`. `AWS_ENDPOINT_URL` points the SDK back at fakecloud (overridable by setting it on the function). Reserved keys are rejected on `CreateFunction`/`UpdateFunctionConfiguration` with `InvalidParameterValueException`, as on AWS. The execution role must trust `lambda.amazonaws.com`; under `--iam strict` it must also be in the caller's account (`AccessDeniedException: Cross-account pass role is not allowed.` otherwise; `--iam soft` logs the would-be denial), while the default mode accepts another account's role (e.g. `000000000000` templates) and mints the session in the function's account. Applied on create, update and CloudFormation alike - **Aliases and versions** — publish, point aliases at versions; alias-based weighted routing (`RoutingConfig`) is enforced at invoke time so traffic splits between versions exactly as on AWS - **Concurrency controls** — reserved concurrency enforced at invocation time: per-function reservation caps in-flight invocations and excess requests are rejected with `TooManyRequestsException` (HTTP 429) and `Reason=ReservedFunctionConcurrentInvocationLimitExceeded` - **`UpdateFunctionCode` from S3** — `S3Bucket`/`S3Key`/`S3ObjectVersion` fetches the ZIP from the fakecloud S3 implementation; the stored `CodeSha256` is the real SHA-256 of the fetched bytes - **CloudWatch metrics** — every invoke publishes `Invocations`, `Errors`, `Duration`, `Throttles`, and `ConcurrentExecutions` to the `AWS/Lambda` namespace, queryable via `GetMetricStatistics` / `GetMetricData` - **Resource-based policies** — the statement API (`AddPermission` / `RemovePermission` / `GetPolicy`) and the document API (`PutResourcePolicy` / `GetResourcePolicy` / `DeleteResourcePolicy`) address one policy per function or qualifier: a `RevisionId` acts as an optimistic-concurrency precondition (`PreconditionFailedException` on mismatch) and an `Allow` open to every principal with no `Condition` is rejected with `PublicPolicyException` - **`GetAccountSettings`** — returns real `AccountUsage` counters (`FunctionCount`, `TotalCodeSize`) and `AccountLimit` so SDKs that pre-flight account quotas see live values -- **Warm container reuse** — subsequent invocations of the same function reuse the container +- **Warm container reuse**: subsequent invocations of the same function reuse the container; each version has its own warm pool, and a configuration change (environment, memory, timeout, role, handler, tags) starts a fresh instance once the old one finishes any in-flight invocation - **Async invoke destinations** — `OnSuccess` / `OnFailure` routes the invocation result to SQS, SNS, EventBridge, or another Lambda by ARN scheme; record matches the AWS destinations schema (`requestContext`, `requestPayload`, `responseContext`, `responsePayload`) - **`InvocationType` honored** — `Event` returns 202 and runs in the background, `RequestResponse` blocks for the result, `DryRun` validates without executing