use std::collections::BTreeMap;
use std::path::PathBuf;
use std::time::Duration;
use anyhow::{bail, Context, Result};
use reqwest::StatusCode;
use serde::{Deserialize, Serialize};
use tracing::info;
use workload_spec::WorkloadSpec;
use crate::{ContainerLauncher, ContainerRunSpec};
pub const DEFAULT_SSR_CONTAINER_PORT: u16 = 3000;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SsrRuntimeSpec {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub network: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub network_alias: Option<String>,
pub image: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cmd: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub env: BTreeMap<String, String>,
pub host_port: u16,
pub container_port: u16,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub volumes: Vec<(PathBuf, PathBuf)>,
pub container_name: String,
pub container_label: String,
#[serde(with = "duration_secs_serde")]
pub ready_timeout: Duration,
#[serde(default = "default_ready_path")]
pub ready_path: String,
}
fn default_ready_path() -> String {
"/readyz".to_string()
}
mod duration_secs_serde {
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::time::Duration;
pub fn serialize<S: Serializer>(d: &Duration, s: S) -> Result<S::Ok, S::Error> {
d.as_secs().serialize(s)
}
pub fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result<Duration, D::Error> {
let secs = u64::deserialize(d)?;
Ok(Duration::from_secs(secs))
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SsrRuntimeRunning {
pub origin_url: String,
pub container_name: String,
}
pub async fn ensure_ssr_runtime_running(
runtime: &impl ContainerLauncher,
spec: &SsrRuntimeSpec,
) -> Result<SsrRuntimeRunning> {
runtime
.ensure_image(&spec.image)
.await
.with_context(|| format!("pulling {}", spec.image))?;
let volumes_str: Vec<(PathBuf, String)> = spec
.volumes
.iter()
.map(|(host, target)| (host.clone(), target.display().to_string()))
.collect();
let network_aliases = match (&spec.network, &spec.network_alias) {
(Some(_), Some(alias)) => vec![alias.clone()],
_ => vec![],
};
let run_spec = ContainerRunSpec {
name: spec.container_name.clone(),
image: spec.image.clone(),
label: spec.container_label.clone(),
ports: vec![(spec.host_port, spec.container_port)],
env: spec.env.clone(),
volumes: volumes_str,
cmd: spec.cmd.clone().unwrap_or_default(),
cap_add: vec![],
cgroupns: None,
network: spec.network.clone(),
network_aliases,
extra_hosts: vec![],
};
runtime
.run(&run_spec)
.await
.with_context(|| format!("starting SSR-runtime container {}", spec.container_name))?;
let probe_host = crate::pond_probe_host();
if !wait_for_port(&probe_host, spec.host_port, spec.ready_timeout).await {
let _ = runtime
.stop_and_remove(&spec.container_name, Duration::from_secs(2))
.await;
bail!(
"SSR-runtime container did not bind {probe_host}:{} within {:?}",
spec.host_port,
spec.ready_timeout,
);
}
let origin_url = format!("http://127.0.0.1:{}", spec.host_port);
let probe_url = format!(
"http://{probe_host}:{port}{path}",
port = spec.host_port,
path = spec.ready_path
);
if !wait_for_http_ready(&probe_url, spec.ready_timeout).await {
let _ = runtime
.stop_and_remove(&spec.container_name, Duration::from_secs(2))
.await;
bail!(
"SSR-runtime container at {origin_url} did not pass {probe_url} within {:?}",
spec.ready_timeout,
);
}
info!(
origin_url = %origin_url,
container = %spec.container_name,
"pond SSR-runtime container ready",
);
Ok(SsrRuntimeRunning {
origin_url,
container_name: spec.container_name.clone(),
})
}
pub fn lower_workload_spec(
ws: &WorkloadSpec,
host_port: u16,
container_name: String,
container_label: String,
ready_timeout: Duration,
) -> Result<SsrRuntimeSpec> {
let image_str = compose_image_ref(&ws.image);
let mut env_map = BTreeMap::new();
for var in &ws.env {
let (key, value) = lower_env_var(var).with_context(|| {
format!("lowering env var for SSR runtime {}", ws.name)
})?;
env_map.insert(key, value);
}
let volumes = ws
.volumes
.iter()
.filter_map(lower_volume_mount)
.collect();
let container_port = ws
.expose
.mesh
.ports
.first()
.copied()
.unwrap_or(DEFAULT_SSR_CONTAINER_PORT);
Ok(SsrRuntimeSpec {
network: None,
network_alias: None,
image: image_str,
cmd: ws.command.clone(),
env: env_map,
host_port,
container_port,
volumes,
container_name,
container_label,
ready_timeout,
ready_path: ready_path_for(ws),
})
}
fn ready_path_for(ws: &WorkloadSpec) -> String {
match ws.healthcheck.as_ref().map(|h| &h.probe) {
Some(workload_spec::HealthProbe::HttpGet { path, .. }) => path.clone(),
_ => default_ready_path(),
}
}
fn compose_image_ref(image: &workload_spec::ImageRef) -> String {
let repo = if image.registry.is_empty() {
image.repository.clone()
} else {
format!("{}/{}", image.registry, image.repository)
};
if image.digest.is_empty() {
format!("{repo}:{}", image.tag)
} else {
format!("{repo}:{}@{}", image.tag, image.digest)
}
}
fn lower_env_var(var: &workload_spec::EnvVar) -> Result<(String, String)> {
match &var.value {
workload_spec::EnvValue::Literal { value } => {
Ok((var.name.clone(), value.clone()))
}
workload_spec::EnvValue::FromSecret { .. } => bail!(
"env var {:?}: FromSecret not yet supported in pond SSR-runtime lowering",
var.name
),
workload_spec::EnvValue::FromMesh { .. } => bail!(
"env var {:?}: FromMesh not yet supported in pond SSR-runtime lowering",
var.name
),
}
}
fn lower_volume_mount(
mount: &workload_spec::VolumeMount,
) -> Option<(PathBuf, PathBuf)> {
match &mount.source {
workload_spec::VolumeSource::Bind { host_path } => {
Some((host_path.clone(), mount.target.clone()))
}
_ => None,
}
}
async fn wait_for_port(host: &str, port: u16, timeout: Duration) -> bool {
let deadline = tokio::time::Instant::now() + timeout;
loop {
if tokio::net::TcpStream::connect((host, port)).await.is_ok() {
return true;
}
if tokio::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
async fn wait_for_http_ready(url: &str, timeout: Duration) -> bool {
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(2))
.build()
.unwrap_or_else(|_| reqwest::Client::new());
let deadline = tokio::time::Instant::now() + timeout;
loop {
if let Ok(resp) = client.get(url).send().await {
let s = resp.status();
if s != StatusCode::SERVICE_UNAVAILABLE
&& s != StatusCode::BAD_GATEWAY
&& s != StatusCode::GATEWAY_TIMEOUT
&& !s.is_server_error()
{
return true;
}
}
if tokio::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use workload_spec::{
EnvValue, EnvVar, ExposeSpec, ImageRef, MeshExpose, MeshIdent, ResourceLimits,
RestartPolicy, SchemaVersion, StopPolicy, TierTag, VolumeMount, VolumeSource,
};
fn minimal_workload_spec() -> WorkloadSpec {
WorkloadSpec {
schema_version: SchemaVersion::V1,
name: "ssr-runtime".into(),
image: ImageRef {
registry: "docker.io".into(),
repository: "oven/bun".into(),
tag: "1".into(),
digest: workload_spec::testing::test_digest(),
},
tier: TierTag("service".into()),
tenant: workload_spec::TenantId::singleton(),
namespace: workload_spec::NamespaceId::singleton(),
replicas: 1,
command: Some(vec!["bun".into(), "run".into(), "src/ssr.ts".into()]),
entrypoint: None,
workdir: None,
user: None,
env: vec![EnvVar {
name: "NODE_ENV".into(),
value: EnvValue::Literal {
value: "production".into(),
},
}],
secrets: vec![],
volumes: vec![VolumeMount {
source: VolumeSource::Bind {
host_path: PathBuf::from("/host/src"),
},
target: PathBuf::from("/app/src"),
read_only: false,
}],
resources: ResourceLimits {
memory_mb: 256,
cpu_millis: 512,
ephemeral_storage_mb: 256,
},
depends_on: vec![],
healthcheck: None,
restart_policy: RestartPolicy::Always,
archetype: None,
stop_policy: StopPolicy {
signal: 15,
grace_period: workload_spec::Millis::from_secs(10),
},
expose: ExposeSpec {
mesh: MeshExpose {
identity: MeshIdent("ssr-runtime".into()),
ports: vec![3000],
allow_from: vec![],
},
public: None,
operator: None,
},
labels: Default::default(),
annotations: Default::default(),
}
}
#[test]
fn lower_composes_image_ref_with_tag_and_digest() {
let ws = minimal_workload_spec();
let spec = lower_workload_spec(
&ws,
14321,
"yah-pond-svc-pond-ssr".into(),
"svc:pond:ssr".into(),
Duration::from_secs(30),
)
.unwrap();
assert_eq!(
spec.image,
format!("docker.io/oven/bun:1@{}", workload_spec::testing::TEST_DIGEST)
);
}
#[test]
fn lower_emits_explicit_digest() {
let mut ws = minimal_workload_spec();
ws.image.digest = "sha256:abc123".into();
let spec = lower_workload_spec(
&ws,
14321,
"yah-pond-svc-pond-ssr".into(),
"svc:pond:ssr".into(),
Duration::from_secs(30),
)
.unwrap();
assert_eq!(spec.image, "docker.io/oven/bun:1@sha256:abc123");
}
#[test]
fn lower_copies_cmd_and_env_literals() {
let ws = minimal_workload_spec();
let spec = lower_workload_spec(
&ws,
14321,
"yah-pond-svc-pond-ssr".into(),
"svc:pond:ssr".into(),
Duration::from_secs(30),
)
.unwrap();
assert_eq!(spec.cmd.as_deref(), Some(&["bun".to_string(), "run".to_string(), "src/ssr.ts".to_string()][..]));
assert_eq!(spec.env.get("NODE_ENV").map(String::as_str), Some("production"));
}
#[test]
fn lower_rejects_from_secret_env() {
let mut ws = minimal_workload_spec();
ws.env.push(EnvVar {
name: "SECRET".into(),
value: EnvValue::FromSecret {
secret: "stripe-key".into(),
key: "value".into(),
},
});
let err = lower_workload_spec(
&ws,
14321,
"n".into(),
"l".into(),
Duration::from_secs(30),
)
.unwrap_err();
assert!(format!("{err:#}").contains("FromSecret"));
}
#[test]
fn lower_uses_expose_mesh_port_for_container_port() {
let ws = minimal_workload_spec();
let spec = lower_workload_spec(
&ws,
14321,
"n".into(),
"l".into(),
Duration::from_secs(30),
)
.unwrap();
assert_eq!(spec.container_port, 3000);
}
#[test]
fn lower_defaults_container_port_when_no_mesh_port() {
let mut ws = minimal_workload_spec();
ws.expose.mesh.ports.clear();
let spec = lower_workload_spec(
&ws,
14321,
"n".into(),
"l".into(),
Duration::from_secs(30),
)
.unwrap();
assert_eq!(spec.container_port, DEFAULT_SSR_CONTAINER_PORT);
}
#[test]
fn lower_defaults_ready_path_to_readyz() {
let ws = minimal_workload_spec();
let spec = lower_workload_spec(&ws, 14321, "n".into(), "l".into(), Duration::from_secs(30))
.unwrap();
assert_eq!(spec.ready_path, "/readyz");
}
#[test]
fn lower_takes_ready_path_from_a_declared_http_healthcheck() {
let mut ws = minimal_workload_spec();
ws.healthcheck = Some(workload_spec::Healthcheck {
probe: workload_spec::HealthProbe::HttpGet {
path: "/custom/health".into(),
port: 9999,
expect_status: None,
},
interval: workload_spec::Millis(1_000),
timeout: workload_spec::Millis(1_000),
initial_delay: workload_spec::Millis(0),
failure_threshold: 3,
});
let spec = lower_workload_spec(&ws, 14321, "n".into(), "l".into(), Duration::from_secs(30))
.unwrap();
assert_eq!(spec.ready_path, "/custom/health");
assert_eq!(spec.host_port, 14321, "probe port stays the host mapping");
}
#[test]
fn lower_falls_back_for_pathless_probe_kinds() {
let mut ws = minimal_workload_spec();
ws.healthcheck = Some(workload_spec::Healthcheck {
probe: workload_spec::HealthProbe::TcpConnect { port: 3000 },
interval: workload_spec::Millis(1_000),
timeout: workload_spec::Millis(1_000),
initial_delay: workload_spec::Millis(0),
failure_threshold: 3,
});
let spec = lower_workload_spec(&ws, 14321, "n".into(), "l".into(), Duration::from_secs(30))
.unwrap();
assert_eq!(spec.ready_path, "/readyz");
}
#[test]
fn lower_host_volume_pairs_through() {
let ws = minimal_workload_spec();
let spec = lower_workload_spec(
&ws,
14321,
"n".into(),
"l".into(),
Duration::from_secs(30),
)
.unwrap();
assert_eq!(spec.volumes.len(), 1);
assert_eq!(spec.volumes[0].0, PathBuf::from("/host/src"));
assert_eq!(spec.volumes[0].1, PathBuf::from("/app/src"));
}
#[test]
fn ssr_runtime_spec_serde_roundtrip() {
let mut env = BTreeMap::new();
env.insert("FOO".into(), "bar".into());
let spec = SsrRuntimeSpec {
network: None,
network_alias: None,
image: "oven/bun:1".into(),
cmd: Some(vec!["bun".into(), "run".into()]),
env,
host_port: 14321,
container_port: 3000,
volumes: vec![(PathBuf::from("/a"), PathBuf::from("/b"))],
container_name: "yah-pond-svc-pond-ssr".into(),
container_label: "svc:pond:ssr".into(),
ready_timeout: Duration::from_secs(30),
ready_path: "/".into(),
};
let s = serde_json::to_string(&spec).unwrap();
let round: SsrRuntimeSpec = serde_json::from_str(&s).unwrap();
assert_eq!(round.image, spec.image);
assert_eq!(round.host_port, spec.host_port);
assert_eq!(round.container_port, spec.container_port);
assert_eq!(round.ready_path, "/");
}
}