1use std::collections::BTreeMap;
28use std::path::PathBuf;
29use std::time::Duration;
30
31use anyhow::{bail, Context, Result};
32use reqwest::StatusCode;
33use serde::{Deserialize, Serialize};
34use tracing::info;
35use workload_spec::WorkloadSpec;
36
37use crate::{ContainerLauncher, ContainerRunSpec};
38
39pub const DEFAULT_SSR_CONTAINER_PORT: u16 = 3000;
43
44#[derive(Debug, Clone, Serialize, Deserialize)]
48pub struct SsrRuntimeSpec {
49 #[serde(default, skip_serializing_if = "Option::is_none")]
54 pub network: Option<String>,
55 #[serde(default, skip_serializing_if = "Option::is_none")]
58 pub network_alias: Option<String>,
59 pub image: String,
61 #[serde(default, skip_serializing_if = "Option::is_none")]
64 pub cmd: Option<Vec<String>>,
65 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
67 pub env: BTreeMap<String, String>,
68 pub host_port: u16,
71 pub container_port: u16,
75 #[serde(default, skip_serializing_if = "Vec::is_empty")]
77 pub volumes: Vec<(PathBuf, PathBuf)>,
78 pub container_name: String,
80 pub container_label: String,
82 #[serde(with = "duration_secs_serde")]
84 pub ready_timeout: Duration,
85 #[serde(default = "default_ready_path")]
89 pub ready_path: String,
90}
91
92fn default_ready_path() -> String {
107 "/readyz".to_string()
108}
109
110mod duration_secs_serde {
111 use serde::{Deserialize, Deserializer, Serialize, Serializer};
112 use std::time::Duration;
113
114 pub fn serialize<S: Serializer>(d: &Duration, s: S) -> Result<S::Ok, S::Error> {
115 d.as_secs().serialize(s)
116 }
117 pub fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result<Duration, D::Error> {
118 let secs = u64::deserialize(d)?;
119 Ok(Duration::from_secs(secs))
120 }
121}
122
123#[derive(Debug, Clone, Serialize, Deserialize)]
126pub struct SsrRuntimeRunning {
127 pub origin_url: String,
130 pub container_name: String,
131}
132
133pub async fn ensure_ssr_runtime_running(
138 runtime: &impl ContainerLauncher,
139 spec: &SsrRuntimeSpec,
140) -> Result<SsrRuntimeRunning> {
141 runtime
142 .ensure_image(&spec.image)
143 .await
144 .with_context(|| format!("pulling {}", spec.image))?;
145
146 let volumes_str: Vec<(PathBuf, String)> = spec
147 .volumes
148 .iter()
149 .map(|(host, target)| (host.clone(), target.display().to_string()))
150 .collect();
151 let network_aliases = match (&spec.network, &spec.network_alias) {
152 (Some(_), Some(alias)) => vec![alias.clone()],
153 _ => vec![],
154 };
155 let run_spec = ContainerRunSpec {
156 name: spec.container_name.clone(),
157 image: spec.image.clone(),
158 label: spec.container_label.clone(),
159 ports: vec![(spec.host_port, spec.container_port)],
160 env: spec.env.clone(),
161 volumes: volumes_str,
162 cmd: spec.cmd.clone().unwrap_or_default(),
163 cap_add: vec![],
164 cgroupns: None,
165 network: spec.network.clone(),
166 network_aliases,
167 extra_hosts: vec![],
168 };
169 runtime
170 .run(&run_spec)
171 .await
172 .with_context(|| format!("starting SSR-runtime container {}", spec.container_name))?;
173
174 let probe_host = crate::pond_probe_host();
177 if !wait_for_port(&probe_host, spec.host_port, spec.ready_timeout).await {
178 let _ = runtime
179 .stop_and_remove(&spec.container_name, Duration::from_secs(2))
180 .await;
181 bail!(
182 "SSR-runtime container did not bind {probe_host}:{} within {:?}",
183 spec.host_port,
184 spec.ready_timeout,
185 );
186 }
187
188 let origin_url = format!("http://127.0.0.1:{}", spec.host_port);
189 let probe_url = format!(
190 "http://{probe_host}:{port}{path}",
191 port = spec.host_port,
192 path = spec.ready_path
193 );
194 if !wait_for_http_ready(&probe_url, spec.ready_timeout).await {
195 let _ = runtime
196 .stop_and_remove(&spec.container_name, Duration::from_secs(2))
197 .await;
198 bail!(
199 "SSR-runtime container at {origin_url} did not pass {probe_url} within {:?}",
200 spec.ready_timeout,
201 );
202 }
203
204 info!(
205 origin_url = %origin_url,
206 container = %spec.container_name,
207 "pond SSR-runtime container ready",
208 );
209
210 Ok(SsrRuntimeRunning {
211 origin_url,
212 container_name: spec.container_name.clone(),
213 })
214}
215
216pub fn lower_workload_spec(
244 ws: &WorkloadSpec,
245 host_port: u16,
246 container_name: String,
247 container_label: String,
248 ready_timeout: Duration,
249) -> Result<SsrRuntimeSpec> {
250 let image_str = compose_image_ref(&ws.image);
251
252 let mut env_map = BTreeMap::new();
253 for var in &ws.env {
254 let (key, value) = lower_env_var(var).with_context(|| {
255 format!("lowering env var for SSR runtime {}", ws.name)
256 })?;
257 env_map.insert(key, value);
258 }
259
260 let volumes = ws
261 .volumes
262 .iter()
263 .filter_map(lower_volume_mount)
264 .collect();
265
266 let container_port = ws
267 .expose
268 .mesh
269 .numbers()
270 .first()
271 .copied()
272 .unwrap_or(DEFAULT_SSR_CONTAINER_PORT);
273
274 Ok(SsrRuntimeSpec {
275 network: None,
276 network_alias: None,
277 image: image_str,
278 cmd: ws.command.clone(),
279 env: env_map,
280 host_port,
281 container_port,
282 volumes,
283 container_name,
284 container_label,
285 ready_timeout,
286 ready_path: ready_path_for(ws),
287 })
288}
289
290fn ready_path_for(ws: &WorkloadSpec) -> String {
292 match ws.healthcheck.as_ref().map(|h| &h.probe) {
293 Some(workload_spec::HealthProbe::HttpGet { path, .. }) => path.clone(),
294 _ => default_ready_path(),
295 }
296}
297
298fn compose_image_ref(image: &workload_spec::ImageRef) -> String {
300 let repo = if image.registry.is_empty() {
301 image.repository.clone()
302 } else {
303 format!("{}/{}", image.registry, image.repository)
304 };
305 if image.digest.is_empty() {
309 format!("{repo}:{}", image.tag)
310 } else {
311 format!("{repo}:{}@{}", image.tag, image.digest)
312 }
313}
314
315fn lower_env_var(var: &workload_spec::EnvVar) -> Result<(String, String)> {
321 match &var.value {
322 workload_spec::EnvValue::Literal { value } => {
323 Ok((var.name.clone(), value.clone()))
324 }
325 workload_spec::EnvValue::FromSecret { .. } => bail!(
326 "env var {:?}: FromSecret not yet supported in pond SSR-runtime lowering",
327 var.name
328 ),
329 workload_spec::EnvValue::FromMesh { .. } => bail!(
330 "env var {:?}: FromMesh not yet supported in pond SSR-runtime lowering",
331 var.name
332 ),
333 }
334}
335
336fn lower_volume_mount(
340 mount: &workload_spec::VolumeMount,
341) -> Option<(PathBuf, PathBuf)> {
342 match &mount.source {
343 workload_spec::VolumeSource::Bind { host_path } => {
344 Some((host_path.clone(), mount.target.clone()))
345 }
346 _ => None,
347 }
348}
349
350async fn wait_for_port(host: &str, port: u16, timeout: Duration) -> bool {
353 let deadline = tokio::time::Instant::now() + timeout;
354 loop {
355 if tokio::net::TcpStream::connect((host, port)).await.is_ok() {
356 return true;
357 }
358 if tokio::time::Instant::now() >= deadline {
359 return false;
360 }
361 tokio::time::sleep(Duration::from_millis(50)).await;
362 }
363}
364
365async fn wait_for_http_ready(url: &str, timeout: Duration) -> bool {
368 let client = reqwest::Client::builder()
369 .timeout(Duration::from_secs(2))
370 .build()
371 .unwrap_or_else(|_| reqwest::Client::new());
372 let deadline = tokio::time::Instant::now() + timeout;
373 loop {
374 if let Ok(resp) = client.get(url).send().await {
375 let s = resp.status();
376 if s != StatusCode::SERVICE_UNAVAILABLE
377 && s != StatusCode::BAD_GATEWAY
378 && s != StatusCode::GATEWAY_TIMEOUT
379 && !s.is_server_error()
380 {
381 return true;
382 }
383 }
384 if tokio::time::Instant::now() >= deadline {
385 return false;
386 }
387 tokio::time::sleep(Duration::from_millis(100)).await;
388 }
389}
390
391#[cfg(test)]
392mod tests {
393 use super::*;
394 use workload_spec::{
395 EnvValue, EnvVar, ExposeSpec, ImageRef, MeshExpose, MeshIdent, ResourceLimits,
396 RestartPolicy, SchemaVersion, StopPolicy, TierTag, VolumeMount, VolumeSource,
397 };
398
399 fn minimal_workload_spec() -> WorkloadSpec {
400 WorkloadSpec {
401 schema_version: SchemaVersion::V1,
402 name: "ssr-runtime".into(),
403 image: ImageRef {
404 registry: "docker.io".into(),
405 repository: "oven/bun".into(),
406 tag: "1".into(),
407 digest: workload_spec::testing::test_digest(),
408 },
409 tier: TierTag("service".into()),
410 tenant: workload_spec::TenantId::singleton(),
411 namespace: workload_spec::NamespaceId::singleton(),
412 replicas: 1,
413 command: Some(vec!["bun".into(), "run".into(), "src/ssr.ts".into()]),
414 entrypoint: None,
415 workdir: None,
416 user: None,
417 env: vec![EnvVar {
418 name: "NODE_ENV".into(),
419 value: EnvValue::Literal {
420 value: "production".into(),
421 },
422 }],
423 secrets: vec![],
424 volumes: vec![VolumeMount {
425 source: VolumeSource::Bind {
426 host_path: PathBuf::from("/host/src"),
427 },
428 target: PathBuf::from("/app/src"),
429 read_only: false,
430 }],
431 resources: ResourceLimits {
432 memory_mb: 256,
433 cpu_millis: 512,
434 ephemeral_storage_mb: 256,
435 },
436 depends_on: vec![],
437 healthcheck: None,
438 restart_policy: RestartPolicy::Always,
439 archetype: None,
440 stop_policy: StopPolicy {
441 signal: 15,
442 grace_period: workload_spec::Millis::from_secs(10),
443 },
444 expose: ExposeSpec {
445 mesh: MeshExpose {
446 identity: MeshIdent("ssr-runtime".into()),
447 ports: MeshExpose::anonymous_ports([3000]),
448 allow_from: vec![],
449 },
450 public: None,
451 operator: None,
452 },
453 labels: Default::default(),
454 annotations: Default::default(),
455 }
456 }
457
458 #[test]
459 fn lower_composes_image_ref_with_tag_and_digest() {
460 let ws = minimal_workload_spec();
461 let spec = lower_workload_spec(
462 &ws,
463 14321,
464 "yah-pond-svc-pond-ssr".into(),
465 "svc:pond:ssr".into(),
466 Duration::from_secs(30),
467 )
468 .unwrap();
469 assert_eq!(
470 spec.image,
471 format!("docker.io/oven/bun:1@{}", workload_spec::testing::TEST_DIGEST)
472 );
473 }
474
475 #[test]
476 fn lower_emits_explicit_digest() {
477 let mut ws = minimal_workload_spec();
478 ws.image.digest = "sha256:abc123".into();
479 let spec = lower_workload_spec(
480 &ws,
481 14321,
482 "yah-pond-svc-pond-ssr".into(),
483 "svc:pond:ssr".into(),
484 Duration::from_secs(30),
485 )
486 .unwrap();
487 assert_eq!(spec.image, "docker.io/oven/bun:1@sha256:abc123");
488 }
489
490 #[test]
491 fn lower_copies_cmd_and_env_literals() {
492 let ws = minimal_workload_spec();
493 let spec = lower_workload_spec(
494 &ws,
495 14321,
496 "yah-pond-svc-pond-ssr".into(),
497 "svc:pond:ssr".into(),
498 Duration::from_secs(30),
499 )
500 .unwrap();
501 assert_eq!(spec.cmd.as_deref(), Some(&["bun".to_string(), "run".to_string(), "src/ssr.ts".to_string()][..]));
502 assert_eq!(spec.env.get("NODE_ENV").map(String::as_str), Some("production"));
503 }
504
505 #[test]
506 fn lower_rejects_from_secret_env() {
507 let mut ws = minimal_workload_spec();
508 ws.env.push(EnvVar {
509 name: "SECRET".into(),
510 value: EnvValue::FromSecret {
511 secret: "stripe-key".into(),
512 key: "value".into(),
513 },
514 });
515 let err = lower_workload_spec(
516 &ws,
517 14321,
518 "n".into(),
519 "l".into(),
520 Duration::from_secs(30),
521 )
522 .unwrap_err();
523 assert!(format!("{err:#}").contains("FromSecret"));
524 }
525
526 #[test]
527 fn lower_uses_expose_mesh_port_for_container_port() {
528 let ws = minimal_workload_spec();
529 let spec = lower_workload_spec(
530 &ws,
531 14321,
532 "n".into(),
533 "l".into(),
534 Duration::from_secs(30),
535 )
536 .unwrap();
537 assert_eq!(spec.container_port, 3000);
538 }
539
540 #[test]
541 fn lower_defaults_container_port_when_no_mesh_port() {
542 let mut ws = minimal_workload_spec();
543 ws.expose.mesh.ports.clear();
544 let spec = lower_workload_spec(
545 &ws,
546 14321,
547 "n".into(),
548 "l".into(),
549 Duration::from_secs(30),
550 )
551 .unwrap();
552 assert_eq!(spec.container_port, DEFAULT_SSR_CONTAINER_PORT);
553 }
554
555 #[test]
556 fn lower_defaults_ready_path_to_readyz() {
557 let ws = minimal_workload_spec();
561 let spec = lower_workload_spec(&ws, 14321, "n".into(), "l".into(), Duration::from_secs(30))
562 .unwrap();
563 assert_eq!(spec.ready_path, "/readyz");
564 }
565
566 #[test]
567 fn lower_takes_ready_path_from_a_declared_http_healthcheck() {
568 let mut ws = minimal_workload_spec();
569 ws.healthcheck = Some(workload_spec::Healthcheck {
570 probe: workload_spec::HealthProbe::HttpGet {
571 path: "/custom/health".into(),
574 port: 9999,
575 expect_status: None,
576 },
577 interval: workload_spec::Millis(1_000),
578 timeout: workload_spec::Millis(1_000),
579 initial_delay: workload_spec::Millis(0),
580 failure_threshold: 3,
581 });
582 let spec = lower_workload_spec(&ws, 14321, "n".into(), "l".into(), Duration::from_secs(30))
583 .unwrap();
584 assert_eq!(spec.ready_path, "/custom/health");
585 assert_eq!(spec.host_port, 14321, "probe port stays the host mapping");
586 }
587
588 #[test]
589 fn lower_falls_back_for_pathless_probe_kinds() {
590 let mut ws = minimal_workload_spec();
591 ws.healthcheck = Some(workload_spec::Healthcheck {
592 probe: workload_spec::HealthProbe::TcpConnect { port: 3000 },
593 interval: workload_spec::Millis(1_000),
594 timeout: workload_spec::Millis(1_000),
595 initial_delay: workload_spec::Millis(0),
596 failure_threshold: 3,
597 });
598 let spec = lower_workload_spec(&ws, 14321, "n".into(), "l".into(), Duration::from_secs(30))
599 .unwrap();
600 assert_eq!(spec.ready_path, "/readyz");
601 }
602
603 #[test]
604 fn lower_host_volume_pairs_through() {
605 let ws = minimal_workload_spec();
606 let spec = lower_workload_spec(
607 &ws,
608 14321,
609 "n".into(),
610 "l".into(),
611 Duration::from_secs(30),
612 )
613 .unwrap();
614 assert_eq!(spec.volumes.len(), 1);
615 assert_eq!(spec.volumes[0].0, PathBuf::from("/host/src"));
616 assert_eq!(spec.volumes[0].1, PathBuf::from("/app/src"));
617 }
618
619 #[test]
620 fn ssr_runtime_spec_serde_roundtrip() {
621 let mut env = BTreeMap::new();
622 env.insert("FOO".into(), "bar".into());
623 let spec = SsrRuntimeSpec {
624 network: None,
625 network_alias: None,
626 image: "oven/bun:1".into(),
627 cmd: Some(vec!["bun".into(), "run".into()]),
628 env,
629 host_port: 14321,
630 container_port: 3000,
631 volumes: vec![(PathBuf::from("/a"), PathBuf::from("/b"))],
632 container_name: "yah-pond-svc-pond-ssr".into(),
633 container_label: "svc:pond:ssr".into(),
634 ready_timeout: Duration::from_secs(30),
635 ready_path: "/".into(),
636 };
637 let s = serde_json::to_string(&spec).unwrap();
638 let round: SsrRuntimeSpec = serde_json::from_str(&s).unwrap();
639 assert_eq!(round.image, spec.image);
640 assert_eq!(round.host_port, spec.host_port);
641 assert_eq!(round.container_port, spec.container_port);
642 assert_eq!(round.ready_path, "/");
643 }
644}