use std::collections::BTreeMap;
use std::path::PathBuf;
use anyhow::{Context, Result};
use tokio::sync::oneshot;
use tracing::{info, warn};
use local_driver::pond_minio::ensure_bucket_public;
use super::pond::{
ensure_sim_port_free, run_miniflare_supervisor, spawn_miniflare_child, DoorScripts,
PondOptions,
};
use super::{into_running, slot_field_u16, ReconcileCtx, RunningWorkload};
use crate::capability::Capability;
use crate::config::Provider;
use crate::reconciler::s3_driver::{running_endpoint, DEV_ACCESS_KEY, DEV_SECRET_KEY};
const STATIC_SLOT: &str = "static";
pub const DEFAULT_DEV_DOOR_PORT: u16 = 4321;
pub fn binds_dev_store(ctx: &ReconcileCtx<'_>) -> bool {
ctx.mirror
.driver(Capability::S3)
.and_then(|slot| slot.inline_kind())
== Some(Provider::LocalS3Fs)
}
fn door_dir(ctx: &ReconcileCtx<'_>) -> PathBuf {
ctx.workspace_root
.join(".yah/infra/state/dev/door")
.join(format!("{}-{}", ctx.service.name, ctx.component.id))
}
fn door_scripts(ctx: &ReconcileCtx<'_>) -> Result<DoorScripts> {
let dir = door_dir(ctx);
std::fs::create_dir_all(&dir).with_context(|| format!("creating {}", dir.display()))?;
Ok(DoorScripts {
worker_js: dir.join("worker.js"),
miniflare_shim: dir.join("miniflare-sim.mjs"),
})
}
fn resolve_bucket(ctx: &ReconcileCtx<'_>, fields: &BTreeMap<String, toml::Value>) -> String {
fields
.get("bucket")
.and_then(|v| v.as_str())
.unwrap_or(ctx.service.name.as_str())
.to_string()
}
fn store_endpoint(ctx: &ReconcileCtx<'_>) -> Result<String> {
running_endpoint(ctx.workspace_root).with_context(|| {
let coords = super::s3_driver::coords_path(ctx.workspace_root);
match super::s3_driver::recorded_bringup_error(ctx.workspace_root) {
Some(why) => format!(
"the dev-tier s3 driver is not running — {} has no coordinates, so there is \
nowhere to publish {}/{}'s assets. Its last bring-up failed: {why}",
coords.display(),
ctx.service.name,
ctx.component.id,
),
None => format!(
"the dev-tier s3 driver is not running — {} has no coordinates, so there is \
nowhere to publish {}/{}'s assets. The camp brings this driver up: start \
`yah camp` for this workspace (or attach it in the desktop).",
coords.display(),
ctx.service.name,
ctx.component.id,
),
}
})
}
async fn publish(ctx: &ReconcileCtx<'_>, endpoint: &str, bucket: &str) -> Result<usize> {
ensure_bucket_public(endpoint, bucket, DEV_ACCESS_KEY, DEV_SECRET_KEY)
.await
.with_context(|| format!("ensuring dev bucket {bucket} exists and is public-read"))?;
let out_dir = super::mesofact_static::read_workload_out_dir(&ctx.workload_dir())
.unwrap_or_else(|| "dist".to_string());
let dist_dir = ctx.workload_dir().join(&out_dir);
if !dist_dir.exists() {
warn!(
dist = %dist_dir.display(),
"no built dist to publish — serving existing bucket contents",
);
return Ok(0);
}
let report = super::pond_publish::publish_to_pond(
&dist_dir,
endpoint,
bucket,
DEV_ACCESS_KEY,
DEV_SECRET_KEY,
None,
)
.await
.with_context(|| format!("publishing {} to dev bucket {bucket}", dist_dir.display()))?;
Ok(report.uploaded.len())
}
pub async fn up_dev_door(
ctx: &ReconcileCtx<'_>,
options: &PondOptions,
static_fields: &BTreeMap<String, toml::Value>,
worker_script: &str,
) -> Result<RunningWorkload> {
let door_port = slot_field_u16(static_fields, "port").unwrap_or(DEFAULT_DEV_DOOR_PORT);
let endpoint = store_endpoint(ctx)?;
let bucket = resolve_bucket(ctx, static_fields);
let scripts = door_scripts(ctx)?;
let uploaded = publish(ctx, &endpoint, &bucket).await?;
info!(
uploaded,
bucket = %bucket,
endpoint = %endpoint,
"dev publish complete",
);
let worker_mode =
super::mesofact_static::parse_worker_mode(&ctx.component.kind, static_fields);
let asset_origin = format!("{}/{}", endpoint.trim_end_matches('/'), bucket);
let (child, log_buf) = spawn_miniflare_child(
ctx,
options,
door_port,
&scripts,
&asset_origin,
worker_script,
&worker_mode,
)
.await?;
let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
let supervisor = tokio::spawn(async move {
run_miniflare_supervisor(child, shutdown_rx).await;
Ok(())
});
let dev_url = format!("http://127.0.0.1:{door_port}");
info!(
dev_url = %dev_url,
asset_origin = %asset_origin,
"dev door ready (miniflare + yah-s3-fs)",
);
Ok(into_running(
super::mesofact_static::WORKLOAD_KIND,
STATIC_SLOT,
Some(dev_url),
None,
Some(log_buf),
shutdown_tx,
supervisor,
))
}
pub async fn sync_dev_door(
ctx: &ReconcileCtx<'_>,
static_fields: &BTreeMap<String, toml::Value>,
) -> Result<usize> {
let endpoint = store_endpoint(ctx)?;
let bucket = resolve_bucket(ctx, static_fields);
let uploaded = publish(ctx, &endpoint, &bucket).await?;
info!(uploaded, bucket = %bucket, "dev re-publish complete");
Ok(uploaded)
}
pub async fn ensure_door_port_free(port: u16) -> Result<()> {
ensure_sim_port_free(port).await
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::{MirrorConfig, ServiceComponent, ServiceConfig};
use crate::reconciler::ProviderScope;
use std::path::Path;
fn mirror(src: &str) -> MirrorConfig {
toml::from_str(src).expect("parse mirror")
}
fn service() -> ServiceConfig {
toml::from_str("schema_version = 1\nname = \"yah-marketing\"\n[address]\nkind = \"front-door\"\ndomain = \"yah.dev\"\n")
.expect("parse service")
}
fn component() -> ServiceComponent {
toml::from_str(
"id = \"site\"\nkind = \"mesofact-static\"\npath = \"app/yah/web\"\nrole = \"static\"\n",
)
.expect("parse component")
}
fn ctx<'a>(
root: &'a Path,
service: &'a ServiceConfig,
component: &'a ServiceComponent,
mirror: &'a MirrorConfig,
) -> ReconcileCtx<'a> {
ReconcileCtx {
workspace_root: root,
service,
component,
mirror,
env: "dev",
scope: ProviderScope::singleton(),
}
}
#[test]
fn the_dev_arm_is_selected_by_the_s3_driver_binding() {
let svc = service();
let comp = component();
let bound = mirror(
"schema_version = 1\nshape = \"local\"\n\n\
[providers.static]\nkind = \"miniflare-native\"\nport = 4321\n\n\
[drivers.s3]\nkind = \"local-s3-fs\"\n",
);
assert!(binds_dev_store(&ctx(
Path::new("/tmp"),
&svc,
&comp,
&bound
)));
let unbound = mirror(
"schema_version = 1\nshape = \"local\"\n\n\
[providers.static]\nkind = \"miniflare-native\"\nport = 4321\n",
);
assert!(!binds_dev_store(&ctx(
Path::new("/tmp"),
&svc,
&comp,
&unbound
)));
}
#[test]
fn the_bucket_defaults_to_the_service_name_the_camp_created() {
let svc = service();
let comp = component();
let m = mirror("schema_version = 1\nshape = \"local\"\n");
let c = ctx(Path::new("/tmp"), &svc, &comp, &m);
assert_eq!(resolve_bucket(&c, &BTreeMap::new()), "yah-marketing");
let mut fields = BTreeMap::new();
fields.insert("bucket".to_string(), toml::Value::String("other".into()));
assert_eq!(resolve_bucket(&c, &fields), "other");
}
#[tokio::test]
async fn a_missing_driver_names_the_coords_file_and_the_camp() {
let dir = tempfile::tempdir().expect("tempdir");
let svc = service();
let comp = component();
let m = mirror("schema_version = 1\nshape = \"local\"\n");
let err = store_endpoint(&ctx(dir.path(), &svc, &comp, &m))
.expect_err("no coords written")
.to_string();
assert!(err.contains("coords.json"), "{err}");
assert!(err.contains("yah camp"), "{err}");
}
#[tokio::test]
async fn a_recorded_bringup_failure_is_quoted_instead_of_telling_a_running_camp_to_start() {
let dir = tempfile::tempdir().expect("tempdir");
let failure = super::super::s3_driver::bringup_error_path(dir.path());
std::fs::create_dir_all(failure.parent().unwrap()).unwrap();
std::fs::write(&failure, "port 49786 (PORT_S3) is held by pid 80519\n").unwrap();
let svc = service();
let comp = component();
let m = mirror("schema_version = 1\nshape = \"local\"\n");
let err = format!(
"{:#}",
store_endpoint(&ctx(dir.path(), &svc, &comp, &m)).expect_err("no coords written")
);
assert!(err.contains("coords.json"), "{err}");
assert!(err.contains("pid 80519"), "{err}");
assert!(!err.contains("start `yah camp`"), "{err}");
}
}