durable-actors 0.7.15

Standalone regional durable-actors control plane, host, and durability runtime
use axum::{Router, extract::WebSocketUpgrade, response::Response, routing::get};
use futures_util::{SinkExt, StreamExt};
use tokio::sync::oneshot;
use tokio_tungstenite::{connect_async, tungstenite::Message};

use super::*;
use std::collections::HashMap;

#[test]
fn configures_gateway_connection_limit() -> Result<()> {
    let mut values = process_environment();
    let parse = |values: &HashMap<&str, &str>| {
        ControlPlaneProcessConfig::from_lookup(|name| values.get(name).map(|value| (*value).into()))
    };
    assert_eq!(parse(&values)?.max_socket_connections, 32768);
    values.insert("DURABLE_ACTORS_SOCKET_MAX_CONNECTIONS", "4096");
    let config = parse(&values)?;
    assert_eq!(config.max_socket_connections, 4096);
    values.insert("DURABLE_ACTORS_SOCKET_MAX_CONNECTIONS", "0");
    assert!(parse(&values).is_err());
    Ok(())
}

#[tokio::test]
async fn server_carries_websocket_upgrades() -> Result<()> {
    let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
    let address = listener.local_addr()?;
    let routes = tonic::service::Routes::from(Router::new().route("/socket", get(echo_websocket)));
    let (shutdown_tx, shutdown_rx) = oneshot::channel();
    let server = tokio::spawn(serve_routes(listener, routes, async {
        let _ = shutdown_rx.await;
    }));

    let (mut socket, _) = connect_async(format!("ws://{address}/socket")).await?;
    socket.send(Message::Text("hello".into())).await?;
    assert_eq!(
        socket.next().await.transpose()?,
        Some(Message::Text("hello".into()))
    );
    socket.close(None).await?;
    let _ = shutdown_tx.send(());
    server.await??;
    Ok(())
}

#[test]
fn parses_the_minimal_storage_configuration() -> Result<()> {
    let values = process_environment();
    let config = ControlPlaneProcessConfig::from_lookup(|name| {
        values.get(name).map(|value| (*value).into())
    })?;
    assert_eq!(config.storage.bucket, "actor-state-test");
    let resources = &config.sandbox_provider.pool.resources;
    assert_eq!((resources.cpu_millis, resources.memory_mib), (1000, 256));
    assert_eq!(resources, &crate::sandbox::ResourceLimits::default());
    assert_eq!(config.sandbox_provider.runtime.host_idle_timeout_ms, 10_000);
    assert_eq!(config.jwt_max_lifetime, Duration::from_secs(86_400));
    assert_eq!(config.api_key.as_deref(), Some("api-key"));
    Ok(())
}

#[test]
fn pool_capacity_is_configurable_and_validated() -> Result<()> {
    let mut values = process_environment();
    let parse = |values: &HashMap<&str, &str>| {
        ControlPlaneProcessConfig::from_lookup(|name| values.get(name).map(|v| (*v).into()))
    };
    let defaults = parse(&values)?.sandbox_provider.pool;
    assert_eq!(
        (defaults.idle, defaults.fleet_maximum, defaults.max_starting),
        (64, 256, 32)
    );
    values.extend([
        ("DURABLE_ACTORS_SPARE_IDLE", "128"),
        ("DURABLE_ACTORS_SPARE_FLEET_MAX", "512"),
        ("DURABLE_ACTORS_SPARE_MAX_STARTING", "4"),
        ("DURABLE_ACTORS_HOST_CPU_MILLIS", "500"),
    ]);
    let configured = parse(&values)?.sandbox_provider.pool;
    assert_eq!(configured.resources.cpu_millis, 500);
    assert_eq!(
        (
            configured.idle,
            configured.fleet_maximum,
            configured.max_starting
        ),
        (128, 512, 4)
    );
    values.insert("DURABLE_ACTORS_SPARE_FLEET_MAX", "1");
    assert!(parse(&values).is_err());
    values.insert("DURABLE_ACTORS_SPARE_FLEET_MAX", "512");
    values.insert("DURABLE_ACTORS_SPARE_MAX_STARTING", "0");
    assert!(parse(&values).is_err());
    Ok(())
}

#[test]
fn host_idle_timeout_is_configurable_and_bounded() -> Result<()> {
    for value in ["1", "120000", "86400000"] {
        let mut values = process_environment();
        values.insert("DURABLE_ACTORS_HOST_IDLE_TIMEOUT_MS", value);
        let config = ControlPlaneProcessConfig::from_lookup(|name| {
            values.get(name).map(|value| (*value).into())
        })?;
        assert_eq!(
            config.sandbox_provider.runtime.host_idle_timeout_ms,
            value.parse::<u64>()?
        );
    }
    for value in ["0", "-1", "1.5", "86400001", "not-a-number"] {
        let mut values = process_environment();
        values.insert("DURABLE_ACTORS_HOST_IDLE_TIMEOUT_MS", value);
        assert!(
            ControlPlaneProcessConfig::from_lookup(|name| {
                values.get(name).map(|value| (*value).into())
            })
            .is_err()
        );
    }
    Ok(())
}

#[test]
fn server_configuration_allows_an_unset_secret_on_any_address() -> Result<()> {
    for bind in [
        "127.0.0.1:7100",
        "[::1]:7100",
        "0.0.0.0:7100",
        "[::]:7100",
        "192.168.1.1:7100",
    ] {
        let mut values = process_environment();
        values.remove("DURABLE_ACTORS_SECRET");
        values.insert("DURABLE_ACTORS_CONTROL_PLANE_BIND", bind);
        let config = ControlPlaneProcessConfig::from_lookup(|name| {
            values.get(name).map(|value| (*value).into())
        })?;
        assert!(config.api_key.is_none());
        assert_eq!(config.bind, bind.parse::<SocketAddr>()?);
    }
    Ok(())
}

#[test]
fn server_rejects_an_empty_or_untrimmed_shared_secret() {
    for secret in ["", " ", " key", "key "] {
        let mut values = process_environment();
        values.insert("DURABLE_ACTORS_SECRET", secret);
        let result = ControlPlaneProcessConfig::from_lookup(|name| {
            values.get(name).map(|value| (*value).into())
        });
        assert!(result.is_err(), "server accepted an invalid secret");
    }
}

#[test]
fn authentication_warning_depends_on_the_listening_address_and_secret() -> Result<()> {
    for (bind, exposed) in [
        ("127.0.0.1:7100", false),
        ("127.0.0.2:7100", false),
        ("[::1]:7100", false),
        ("0.0.0.0:7100", true),
        ("[::]:7100", true),
        ("192.168.1.1:7100", true),
    ] {
        for secret in [None, Some("configured-secret")] {
            let output = tempfile::NamedTempFile::new()?;
            let subscriber = tracing_subscriber::fmt()
                .without_time()
                .with_ansi(false)
                .with_max_level(tracing::Level::WARN)
                .with_writer(output.reopen()?)
                .finish();
            tracing::subscriber::with_default(subscriber, || {
                warn_if_authentication_disabled(bind.parse().unwrap(), secret);
            });
            let logs = std::fs::read_to_string(output.path())?;
            if exposed && secret.is_none() {
                assert!(
                    logs.contains("Authentication is disabled"),
                    "{bind}: {logs}"
                );
                assert!(logs.contains(bind));
                assert!(logs.contains("DURABLE_ACTORS_SECRET"));
            } else {
                assert!(logs.is_empty(), "unexpected warning at {bind}: {logs}");
            }
        }
    }
    Ok(())
}

#[test]
fn production_defaults_to_two_rapid_zones_and_gke() -> Result<()> {
    let values = process_environment();
    let config = ControlPlaneProcessConfig::from_lookup(|name| {
        values.get(name).map(|value| (*value).into())
    })?;
    assert!(matches!(
        config.storage.persistence,
        crate::bucket::PersistenceConfig::Rapid { .. }
    ));
    assert_eq!(
        config.sandbox_provider.gke.zones["north-america-west"],
        vec!["us-west4-a"]
    );
    Ok(())
}

#[test]
fn compute_region_accepts_multiple_zones_and_rejects_empty_or_mismatched_sets() -> Result<()> {
    let mut values = process_environment();
    values.insert(
        "DURABLE_ACTORS_GKE_ZONES",
        r#"{"north-america-west":["us-west4-a","us-west4-b","us-west4-c"]}"#,
    );
    let parse = |values: &HashMap<&str, &str>| {
        ControlPlaneProcessConfig::from_lookup(|name| values.get(name).map(|value| (*value).into()))
    };
    parse(&values)?;
    for zones in [
        r#"{"north-america-west":[]}"#,
        r#"{"north-america-west":["us-west4-a","us-east4-b"]}"#,
        r#"{"north-america-west":["us-west4-a","us-west4-a"]}"#,
    ] {
        values.insert("DURABLE_ACTORS_GKE_ZONES", zones);
        assert!(parse(&values).is_err(), "{zones}");
    }
    Ok(())
}

#[test]
fn configures_socket_events_without_a_separate_key() -> Result<()> {
    let mut complete = HashMap::from([(
        "DURABLE_ACTORS_SOCKET_EVENT_URL",
        "https://api.example.com/events",
    )]);
    let sink =
        socket_event_sink_config(&mut |name| complete.get(name).map(|value| (*value).into()))?
            .context("socket event sink was not configured")?;
    assert_eq!(sink.url, "https://api.example.com/events");
    complete.remove("DURABLE_ACTORS_SOCKET_EVENT_URL");
    assert!(
        socket_event_sink_config(&mut |name| complete.get(name).map(|value| (*value).into()))?
            .is_none()
    );
    Ok(())
}

async fn echo_websocket(upgrade: WebSocketUpgrade) -> Response {
    upgrade.on_upgrade(async |mut socket| {
        if let Some(Ok(message)) = socket.recv().await {
            let _ = socket.send(message).await;
        }
    })
}

fn process_environment() -> HashMap<&'static str, &'static str> {
    HashMap::from([
        ("DURABLE_ACTORS_GATEWAY_ROUTE", "http://10.0.0.1:7100"),
        (
            "DURABLE_ACTORS_GOOGLE_SERVICE_ACCOUNT",
            "test@project.iam.gserviceaccount.com",
        ),
        ("DURABLE_ACTORS_JWT_SIGNING_KEY", "c2lnbmluZw=="),
        ("DURABLE_ACTORS_SECRET", "api-key"),
        ("DURABLE_ACTORS_BUCKET", "actor-state-test"),
        (
            "DURABLE_ACTORS_RUNTIME_IMAGE",
            "registry.example/runtime@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
        ),
        ("DURABLE_ACTORS_ARTIFACT_BUCKET", "customer-code"),
        ("DURABLE_ACTORS_ARCHIVE_BUCKET", "actor-archive"),
        (
            "DURABLE_ACTORS_RAPID_BUCKETS",
            r#"[{"bucket":"rapid-test-a","zone":"us-west4-a"},{"bucket":"rapid-test-b","zone":"us-west4-b"}]"#,
        ),
        (
            "DURABLE_ACTORS_GKE_ZONES",
            r#"{"north-america-west":"us-west4-a"}"#,
        ),
        ("DURABLE_ACTORS_PUBLIC_URL", "https://actors.example.com"),
        (
            "DURABLE_ACTORS_CONTROL_PLANE_URL",
            "https://objects.example.com",
        ),
        (
            "DURABLE_ACTORS_POSTGRES_URL",
            "postgresql://localhost/actors",
        ),
    ])
}

#[test]
fn analytics_retention_is_configurable_and_bounded() -> Result<()> {
    let mut values = process_environment();
    let parse = |values: &HashMap<&str, &str>| {
        ControlPlaneProcessConfig::from_lookup(|name| values.get(name).map(|value| (*value).into()))
    };
    assert_eq!(
        parse(&values)?.storage.trace_retention,
        Duration::from_secs(30 * 86400)
    );
    values.insert("DURABLE_ACTORS_ANALYTICS_RETENTION_DAYS", "7");
    assert_eq!(
        parse(&values)?.storage.trace_retention,
        Duration::from_secs(7 * 86400)
    );
    for invalid in ["0", "-1", "1.5", "3651", ""] {
        values.insert("DURABLE_ACTORS_ANALYTICS_RETENTION_DAYS", invalid);
        assert!(parse(&values).is_err());
    }
    Ok(())
}

#[test]
fn archive_batch_triggers_are_configurable_and_positive() -> Result<()> {
    let mut values = process_environment();
    let parse = |values: &HashMap<&str, &str>| {
        ControlPlaneProcessConfig::from_lookup(|name| values.get(name).map(|value| (*value).into()))
    };
    let defaults = parse(&values)?.storage.persistence;
    let crate::bucket::PersistenceConfig::Rapid { archive_batch, .. } = defaults.clone() else {
        panic!()
    };
    assert_eq!(archive_batch.bytes, 16 * 1024 * 1024);
    assert_eq!(archive_batch.interval_ms, 10_000);
    values.insert("DURABLE_ACTORS_ARCHIVE_BATCH_BYTES", "2097152");
    values.insert("DURABLE_ACTORS_ARCHIVE_BATCH_INTERVAL_MS", "250");
    let configured = parse(&values)?.storage.persistence;
    let crate::bucket::PersistenceConfig::Rapid { archive_batch, .. } = configured.clone() else {
        panic!()
    };
    assert_eq!(archive_batch.bytes, 2 * 1024 * 1024);
    assert_eq!(archive_batch.interval_ms, 250);
    assert!(configured.same_backend(&defaults));
    assert_eq!(
        serde_json::from_value::<crate::bucket::PersistenceConfig>(serde_json::to_value(
            &configured
        )?)?,
        configured
    );
    for name in [
        "DURABLE_ACTORS_ARCHIVE_BATCH_BYTES",
        "DURABLE_ACTORS_ARCHIVE_BATCH_INTERVAL_MS",
    ] {
        for invalid in ["0", "-1", "1.5", ""] {
            let mut bad = values.clone();
            bad.insert(name, invalid);
            assert!(parse(&bad).is_err(), "accepted {name}={invalid}");
        }
    }
    Ok(())
}

#[test]
fn default_region_requires_a_configured_compute_zone() -> Result<()> {
    let mut values = process_environment();
    let parse = |values: &HashMap<&str, &str>| {
        ControlPlaneProcessConfig::from_lookup(|name| values.get(name).map(|value| (*value).into()))
    };
    assert_eq!(
        parse(&values)?.region.as_deref(),
        Some("north-america-west")
    );
    values.insert("DURABLE_ACTORS_REGION", "north-america-east");
    assert!(parse(&values).is_err());
    values.remove("DURABLE_ACTORS_REGION");
    values.insert(
        "DURABLE_ACTORS_GKE_ZONES",
        r#"{"north-america-west":"us-west4-b"}"#,
    );
    assert!(parse(&values).is_ok());
    Ok(())
}