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 crate::{ContainerRunSpec, LocalRuntime};
pub const WORKER_SCRIPT_CONTAINER_PATH: &str = "/work/worker.js";
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MiniflareSpec {
pub image: String,
pub container_name: String,
pub container_label: String,
#[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 port: u16,
pub worker_script: String,
pub state_dir: PathBuf,
pub asset_origin: String,
#[serde(default = "default_worker_mode")]
pub worker_mode: String,
#[serde(default)]
pub ssr_origin: String,
#[serde(default)]
pub ssr_prefixes: Vec<String>,
#[serde(with = "duration_secs_serde")]
pub ready_timeout: Duration,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub extra_env: BTreeMap<String, String>,
}
fn default_worker_mode() -> String {
"static".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 MiniflareRunning {
pub dev_url: String,
pub container_name: String,
}
pub async fn ensure_miniflare_running(
runtime: &LocalRuntime,
spec: &MiniflareSpec,
) -> Result<MiniflareRunning> {
std::fs::create_dir_all(&spec.state_dir)
.with_context(|| format!("creating miniflare state dir {}", spec.state_dir.display()))?;
let worker_path = spec.state_dir.join("worker.js");
std::fs::write(&worker_path, spec.worker_script.as_bytes())
.with_context(|| format!("writing {}", worker_path.display()))?;
runtime
.ensure_image(&spec.image)
.await
.with_context(|| format!("pulling {}", spec.image))?;
let ssr_prefixes_json =
serde_json::to_string(&spec.ssr_prefixes).unwrap_or_else(|_| "[]".to_string());
let mut env = BTreeMap::new();
env.insert("MF_PORT".into(), spec.port.to_string());
env.insert("MF_HOST".into(), "0.0.0.0".into());
env.insert("MF_SCRIPT".into(), WORKER_SCRIPT_CONTAINER_PATH.into());
env.insert("ASSET_ORIGIN".into(), spec.asset_origin.clone());
env.insert("WORKER_MODE".into(), spec.worker_mode.clone());
env.insert("SSR_ORIGIN".into(), spec.ssr_origin.clone());
env.insert("SSR_PREFIXES".into(), ssr_prefixes_json);
for (k, v) in &spec.extra_env {
env.insert(k.clone(), v.clone());
}
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.port, spec.port)],
env,
volumes: vec![(worker_path.clone(), WORKER_SCRIPT_CONTAINER_PATH.into())],
cmd: vec![],
cap_add: vec![],
cgroupns: None,
network: spec.network.clone(),
network_aliases,
extra_hosts: vec![],
};
runtime
.run(&run_spec)
.await
.with_context(|| format!("starting miniflare container {}", spec.container_name))?;
let probe_host = crate::pond_probe_host();
if !wait_for_port(&probe_host, spec.port, spec.ready_timeout).await {
let _ = runtime
.stop_and_remove(&spec.container_name, Duration::from_secs(2))
.await;
bail!(
"miniflare container did not bind {probe_host}:{} within {:?}",
spec.port,
spec.ready_timeout,
);
}
let probe_url = format!("http://{probe_host}:{}", spec.port);
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!(
"miniflare container at {probe_url} did not respond within {:?}",
spec.ready_timeout,
);
}
let dev_url = format!("http://127.0.0.1:{}", spec.port);
info!(
dev_url = %dev_url,
container = %spec.container_name,
"pond miniflare container ready",
);
Ok(MiniflareRunning {
dev_url,
container_name: spec.container_name.clone(),
})
}
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::*;
fn fixture() -> MiniflareSpec {
MiniflareSpec {
image: "ghcr.io/yah-ai/yah-miniflare:latest".into(),
container_name: "yah-pond-svc-pond-static".into(),
container_label: "svc:pond:static".into(),
network: Some("yah-pond-svc-pond".into()),
network_alias: Some("miniflare".into()),
port: 4322,
worker_script: "export default { fetch() { return new Response('ok'); } }".into(),
state_dir: PathBuf::from("/tmp/yah-pond/svc-pond/miniflare"),
asset_origin: "http://minio:9000/yah-dev".into(),
worker_mode: "static".into(),
ssr_origin: String::new(),
ssr_prefixes: vec![],
ready_timeout: Duration::from_secs(30),
extra_env: BTreeMap::new(),
}
}
#[test]
fn spec_round_trips_through_serde() {
let spec = fixture();
let s = serde_json::to_string(&spec).unwrap();
let round: MiniflareSpec = serde_json::from_str(&s).unwrap();
assert_eq!(round.image, spec.image);
assert_eq!(round.network.as_deref(), Some("yah-pond-svc-pond"));
assert_eq!(round.network_alias.as_deref(), Some("miniflare"));
assert_eq!(round.port, 4322);
assert_eq!(round.asset_origin, "http://minio:9000/yah-dev");
assert_eq!(round.ready_timeout, Duration::from_secs(30));
}
#[test]
fn spec_round_trips_without_bridge_fields_for_legacy_peers() {
let legacy = r#"{
"image": "ghcr.io/yah-ai/yah-miniflare:latest",
"container_name": "yah-pond-svc-pond-static",
"container_label": "svc:pond:static",
"port": 4322,
"worker_script": "// noop",
"state_dir": "/tmp/yah-pond/svc-pond/miniflare",
"asset_origin": "http://127.0.0.1:9000/yah-dev",
"ready_timeout": 30
}"#;
let spec: MiniflareSpec = serde_json::from_str(legacy).unwrap();
assert!(spec.network.is_none());
assert!(spec.network_alias.is_none());
assert_eq!(spec.worker_mode, "static");
assert!(spec.ssr_prefixes.is_empty());
}
}