Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions MEMORY.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,3 +41,6 @@ One install answers at several addresses at once (the loopback port, the tunnel

## Every feature lands on three fronts: levers, defaults, protection
Before building anything, answer all three, and say so. **Levers**: what a person can reach and from where (a node setting, a CLI flag, a graph control, a ctx call, a `metadata.json` key, a word in the language). Name the lever and where it lives; a knob nobody can turn is not a lever, a knob nobody needs is clutter. **Defaults**: what everybody gets without asking, which is the part nobody should have to know exists. Two questions: is it tedious and useful to make somebody fill it in every time, and is there a value right for almost everybody? Yes to both, set it and keep the lever. Never invent a default to paper over a setting that honestly depends on the case: if leaving it unset is legitimate, it stays explicit and optional and the VALIDATION is what makes the bad shape impossible. A hidden default that is right half the time is worse than a refusal naming what is missing. **Protection**: as long as there is a way to build it properly, the bad shape is refused and the message points at the right one, at the earliest place that can see it (compiler, parser, a node's metadata rules, and a runtime floor where nothing earlier can). An obvious footgun (an infinite loop, a run that waits forever on somebody gone) must be impossible to configure, not warned about. The full wording: `docs/src/thinking/how-we-decide.md`.

## Every slow command runs in the background, with a realistic cap
Anything that can take more than a few seconds runs with `run_in_background`: builds, tests, installs, and cloud work too (`gcloud ... create`, a Cloud SQL instance, `terraform apply`, project or machine creation). Never in the foreground, never "just this once"; an unknown duration counts as slow. Every wait on it is capped at how long that thing normally takes (a cargo check about a minute, an e2e test at most 5, a Cloud SQL create about 10), never a catch-all like 25 minutes: at the cap, look at what it is doing (its output, its operation's state) and tell the [user]. A foreground `gcloud sql instances create` once sat silent for many minutes until the [user] backgrounded it by hand.
18 changes: 13 additions & 5 deletions catalog/postgres/database/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,9 +249,15 @@ impl Node for PostgresDatabaseNode {

async fn run(&self, ctx: ExecutionContext) -> WeftResult<()> {
let database: String = ctx.inputs.get("database")?;
let sql = ctx.endpoint("sql").await?;
// Three asks that need nothing from each other, so they go out
// together.
let (sql, credential, published) =
tokio::try_join!(ctx.endpoint("sql"), ctx.endpoint("credential"), ctx.published_access())?;
let (host, port) = sql.host_and_port()?;
let credential = ctx.endpoint("credential").await?;
let published = match published {
Some(access) => Some(ctx.open(&access).await?),
None => None,
};

// The password comes from the connection this node published
// last time when there is one, and from the database itself
Expand All @@ -263,9 +269,9 @@ impl Node for PostgresDatabaseNode {
// Asking the database whether it still knows a password also
// retires it, so a run that gets a yes here has already done
// the retiring this run owes.
let (password, retired) = match ctx.published_access().await? {
Some(published) => {
let held = ctx.open(&published).await?.value("password")?.to_string();
let (password, retired) = match &published {
Some(opened) => {
let held = opened.value("password")?.to_string();
if confirm_stored(&credential, &held).await? {
(held, true)
} else {
Expand All @@ -285,6 +291,8 @@ impl Node for PostgresDatabaseNode {
values.insert("password".to_string(), password.clone());
// Inside the project's own network; the database serves no TLS.
values.insert("sslmode".to_string(), "disable".to_string());
// Published every run, values unchanged or not: publishing also
// brings the connection's recipe and label up to this node's.
let access = ctx.publish_access(values).await?;

// Retire the password now that a connection holds it, unless
Expand Down
22 changes: 21 additions & 1 deletion crates/weft-broker-client/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,15 @@ where
}
}

/// Whether a failed broker call never reached the broker: the connection
/// itself was refused or could not be made, so the request was never
/// sent. Only such a write is safe to send again; a write that reached it
/// and failed some other way may have landed, and sending it twice would
/// apply it twice.
fn never_sent(e: &anyhow::Error) -> bool {
e.chain().any(|cause| cause.downcast_ref::<reqwest::Error>().is_some_and(reqwest::Error::is_connect))
}

/// The first wait before asking again after the broker could not answer
/// a read; each further failure doubles it, up to
/// [`READ_RETRY_LONGEST`].
Expand Down Expand Up @@ -275,6 +284,15 @@ impl BrokerJournalClient {

#[async_trait]
impl JournalClient for BrokerJournalClient {
/// The broker takes a body of at most [`JOURNAL_RECORD_BODY_LIMIT`].
fn parts<'e>(&self, events: &'e [ExecEvent]) -> Result<Vec<&'e [ExecEvent]>> {
record_chunks(events)
}

fn never_reached(&self, error: &anyhow::Error) -> bool {
never_sent(error)
}

async fn record_event(
&self,
event: &ExecEvent,
Expand Down Expand Up @@ -1229,7 +1247,9 @@ mod read_retry_tests {
}

/// `events` cut, in order, into runs whose request bodies stay under
/// [`JOURNAL_RECORD_BODY_LIMIT`]. An event too big on its own goes alone,
/// [`JOURNAL_RECORD_BODY_LIMIT`]: each one request of
/// [`BrokerJournalClient::record_events`], and the parts a writer that
/// sends a write again sends one at a time (`JournalClient::parts`). An event too big on its own goes alone,
/// and the broker refuses it (`413`): the engine refuses an output that
/// big at the node first, so reaching that refusal is a broken contract.
fn record_chunks(events: &[ExecEvent]) -> Result<Vec<&[ExecEvent]>> {
Expand Down
25 changes: 14 additions & 11 deletions crates/weft-broker/src/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -646,16 +646,19 @@ pub async fn task_claim_one(
require_worker(&caller)?;
require_replica_matches(&caller, &req.replica)?;
let filter = req.filter;
if let ClaimFilter::ExecutionId { project_id, .. } = &filter {
scope::require_project_owned_by(&state.scope_cache, &state.pool, &caller, *project_id).await?;
} else {
let ClaimFilter::ExecutionId { project_id, execution_id } = &filter else {
return Err((StatusCode::FORBIDDEN, "a worker claims only the execution it was called for".into()));
}
let task = state
.tasks
.claim_one(&req.replica, filter, held(req.wait_ms))
.await
.map_err(internal)?;
};
scope::require_project_owned_by(&state.scope_cache, &state.pool, &caller, *project_id).await?;
// The worker's next asks (its run's journal first) are scoped by the
// execution, so its scope is read while the claim runs rather than
// after it, on that first ask.
let execution_id = execution_id.clone();
let (task, ()) = tokio::join!(
state.tasks.claim_one(&req.replica, filter, held(req.wait_ms)),
scope::warm_execution_id_scope(&state.scope_cache, &state.pool, &execution_id),
);
let task = task.map_err(internal)?;
// Latest-claim-wins execution ownership is bound IN the claim's own
// transaction by the `task_claim_binds_execution_id_owner` DB trigger:
// claiming an execution-bearing task atomically stamps
Expand Down Expand Up @@ -805,7 +808,7 @@ async fn declared_infra_under(
project: uuid::Uuid,
digest: String,
) -> Result<Option<Arc<weft_core::project::DeclaredInfra>>, (StatusCode, String)> {
if let Some(declared) = state.declared_infra.get(project, &digest) {
if let Some(declared) = state.declared_infra.get(&(project, digest.clone())) {
return Ok(Some(declared));
}
// Read with its own digest, so what is kept is filed under the
Expand All @@ -820,7 +823,7 @@ async fn declared_infra_under(
let definition: weft_core::project::ProjectDefinition = serde_json::from_str(&project_json)
.map_err(|e| internal(anyhow::anyhow!("project {project}: definition: {e}")))?;
let declared = Arc::new(weft_core::project::DeclaredInfra::of(&definition));
state.declared_infra.put(project, digest, declared.clone());
state.declared_infra.put((project, digest), declared.clone());
Ok(Some(declared))
}

Expand Down
10 changes: 10 additions & 0 deletions crates/weft-broker/src/scope.rs
Original file line number Diff line number Diff line change
Expand Up @@ -237,6 +237,16 @@ async fn lookup_project_tenant(
Ok(tenant)
}

/// Read `execution_id`'s scope into the cache ahead of the asks that
/// need it. Only a head start: an execution that cannot be read here is
/// read again, and refused with the reason, by the first ask that needs
/// it, so nothing is lost by not answering here.
pub async fn warm_execution_id_scope(cache: &ScopeCache, pool: &PgPool, execution_id: &str) {
if let Err((_, why)) = lookup_execution_id_scope(cache, pool, execution_id).await {
tracing::debug!(target: "weft_broker::scope", execution_id, why, "could not read an execution's scope ahead of its asks");
}
}

async fn lookup_execution_id_scope(
cache: &ScopeCache,
pool: &PgPool,
Expand Down
2 changes: 1 addition & 1 deletion crates/weft-broker/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ pub struct BrokerState {
/// project and a digest of its definition: a definition never changes
/// under its digest, so an entry is never stale, and a program asking
/// for an endpoint on every run reads its definition once.
pub declared_infra: weft_core::content_cache::ContentCache<weft_core::project::DeclaredInfra>,
pub declared_infra: weft_core::content_cache::ContentCache<(uuid::Uuid, String), weft_core::project::DeclaredInfra>,
/// Where runtime-file bytes live: the install's bucket.
pub object_store: Arc<dyn ObjectStore>,
/// The runtime-file plane (`ctx.storage`): PG metadata + bucket bytes,
Expand Down
2 changes: 1 addition & 1 deletion crates/weft-cli/src/commands/ci.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,7 @@ mod tests {
assert!(first.contains("-frontends/my-app-front:"), "{first}");
write_workflow(dir.path(), "my app", Cloud::Gcp).unwrap();

std::fs::write(&path, first.replace("runs-on: ubuntu-latest", "runs-on: self-hosted")).unwrap();
std::fs::write(&path, first.replace("runs-on: ubuntu-24.04", "runs-on: self-hosted")).unwrap();
let e = write_workflow(dir.path(), "my app", Cloud::Gcp).unwrap_err().to_string();
assert!(e.contains("changed since weft wrote it"), "{e}");
assert!(std::fs::read_to_string(&path).unwrap().contains("self-hosted"), "left as it is");
Expand Down
12 changes: 6 additions & 6 deletions crates/weft-cli/src/commands/executions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
use anyhow::Context;
use weft_core::program::ExecutionPage;

use super::{local_time, Ctx};
use super::{utc_time, Ctx};

/// A value put into a query string. A node id is the author's own
/// spelling, so it can hold anything they typed; only the handful of
Expand Down Expand Up @@ -91,7 +91,7 @@ pub async fn list(ctx: Ctx, filter: ListFilter) -> anyhow::Result<()> {
return Ok(());
}
println!(
"{:<36} {:<9} {:<13} {:<19} {:<36} entry_node tags",
"{:<36} {:<9} {:<13} {:<23} {:<36} entry_node tags",
"execution_id", "status", "phase", "started", "project_id"
);
for row in &page.executions {
Expand All @@ -104,8 +104,8 @@ pub async fn list(ctx: Ctx, filter: ListFilter) -> anyhow::Result<()> {
// Which instance the run is in, when it is in one.
let instance = row.instance.as_ref().map(|m| format!(" (instance {m})")).unwrap_or_default();
println!(
"{execution_id:<36} {status:<9} {phase:<13} {:<19} {project:<36} {entry}{tags}{instance}",
local_time(row.started_at)
"{execution_id:<36} {status:<9} {phase:<13} {:<23} {project:<36} {entry}{tags}{instance}",
utc_time(row.started_at)
);
}
// The server clamps the page size, so a big --limit can come back
Expand Down Expand Up @@ -248,8 +248,8 @@ pub fn event_line(row: &serde_json::Value, full: bool) -> String {
// pulse, a corruption the replay found) have none and get a blank
// of the same width, so the columns still line up.
let at = match row.get("at_unix").and_then(|v| v.as_u64()) {
Some(at) => format!("[{:<19}]", local_time(at)),
None => " ".repeat(21),
Some(at) => format!("[{:<23}]", utc_time(at)),
None => " ".repeat(25),
};
let node = row_node(row).unwrap_or("");
let mut line = format!("{at} {kind:<23} {node}");
Expand Down
30 changes: 16 additions & 14 deletions crates/weft-cli/src/commands/frontend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
//! to the terminal.

use anyhow::{Context, Result};
use weft_core::frontend::{AddFrontendRequest, Frontend, FrontendHost, FrontendWithToken, Repository};
use weft_core::frontend::{AddFrontendRequest, AddedFrontend, Frontend, FrontendHost, FrontendWithToken, Repository};

use super::Ctx;

Expand Down Expand Up @@ -71,9 +71,10 @@ pub async fn run(ctx: Ctx, action: FrontendAction) -> Result<()> {
} else {
made.await
};
let made: FrontendWithToken =
let made: AddedFrontend =
serde_json::from_value(made?).context("read the frontend the install made")?;
hand_over(&ctx, &project, &install_url, made)?;
let token = made.token.zip(made.frontend.token_id);
hand_over(&ctx, &project, &install_url, &made.frontend, token)?;
}
FrontendAction::Token { name, done: None } => {
// A hosted frontend's token is its workflow's: a new one made
Expand All @@ -92,7 +93,7 @@ pub async fn run(ctx: Ctx, action: FrontendAction) -> Result<()> {
serde_json::from_value(client.post_json(&format!("{base}/{name}/token"), &serde_json::json!({})).await?)
.context("read the frontend's new token")?;
let id = renewed.token_id;
hand_over(&ctx, &project, &install_url, renewed)?;
hand_over(&ctx, &project, &install_url, &renewed.frontend, Some((renewed.token, id)))?;
if !ctx.json() {
println!(
"its old token keeps working until the new one is in place; then `weft frontend token {name} --done {id}` retires it"
Expand Down Expand Up @@ -138,15 +139,14 @@ pub async fn run(ctx: Ctx, action: FrontendAction) -> Result<()> {
Ok(())
}

/// Put a fresh token where only this person can read it, and say what
/// goes where. A frontend the install hosts gets no file: its service
/// gets a token of its own from the repository's workflow (`weft target
/// export` mints it, the deploy puts it in place and retires every other),
/// so one written here would only sit on disk until that deploy killed it.
fn hand_over(ctx: &Ctx, project: &str, install_url: &str, made: FrontendWithToken) -> Result<()> {
let f = &made.frontend;
/// Put a fresh `token` (with its id) where only this person can read it,
/// and say what goes where. A frontend the install hosts has none to
/// write: its service gets its token from the repository's workflow
/// (`weft target export` mints it, the deploy puts it in place and retires
/// every other).
fn hand_over(ctx: &Ctx, project: &str, install_url: &str, f: &Frontend, token: Option<(String, uuid::Uuid)>) -> Result<()> {
if let (Some(repo), Some(service)) = (&f.repo, &f.service) {
if ctx.json_out(&serde_json::json!({ "frontend": f, "tokenId": made.token_id }))? {
if ctx.json_out(&serde_json::json!({ "frontend": f }))? {
return Ok(());
}
println!("frontend '{}': the install made its service {service}, and {} may deploy to it", f.name, repo.name);
Expand All @@ -162,13 +162,15 @@ fn hand_over(ctx: &Ctx, project: &str, install_url: &str, made: FrontendWithToke
}
// A frontend running elsewhere reaches the install at its public
// address.
let (token, token_id) =
token.with_context(|| format!("the install made frontend '{}' with no token for it to call with", f.name))?;
let env = vec![
("WEFT_TOKEN", made.token.clone()),
("WEFT_TOKEN", token),
("WEFT_DISPATCHER_URL", install_url.to_string()),
("WEFT_PUBLIC_URL", install_url.to_string()),
];
let file = super::target::write_secrets_file(&format!("{project}-frontend-{}", f.name), &env)?;
if ctx.json_out(&serde_json::json!({ "frontend": f, "tokenId": made.token_id, "tokenFile": file }))? {
if ctx.json_out(&serde_json::json!({ "frontend": f, "tokenId": token_id, "tokenFile": file }))? {
return Ok(());
}
println!("frontend '{}': its token is in {} (readable by you only; shown this once)", f.name, file.display());
Expand Down
4 changes: 2 additions & 2 deletions crates/weft-cli/src/commands/logs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
use anyhow::Context;
use weft_core::program::{ExecutionLogs, ExecutionSummary};

use super::{local_time, resolve_project_id, Ctx};
use super::{utc_time, resolve_project_id, Ctx};

/// The nodes a run skipped, with why, spelled the way the program
/// reads them (`keep.db` for the db of an included file) when the cwd
Expand Down Expand Up @@ -117,7 +117,7 @@ pub async fn run(ctx: Ctx, target: Option<String>, limit: Option<u32>) -> anyhow
.inherited_from
.map(|execution_id| format!(" [inherited from {}]", super::versions::short(&execution_id.to_string())))
.unwrap_or_default();
println!("[{}] {level:>5}{node}{inherited} {msg}", local_time(entry.at_unix));
println!("[{}] {level:>5}{node}{inherited} {msg}", utc_time(entry.at_unix));
}
// The dispatcher answers the tail, so a full page means the run
// may have written more than this; a cut log must never read as
Expand Down
28 changes: 13 additions & 15 deletions crates/weft-cli/src/commands/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -524,37 +524,35 @@ pub fn resolve_project(
))
}

/// A journal unix stamp as the local wall-clock time a person reads
/// (`2026-09-02 21:36:47`), the one rendering every listing verb uses
/// so a run's start, its events and its log lines line up by eye.
/// Zero (a row that never carried a stamp) renders as a dash rather
/// than as 1970.
pub fn local_time(unix_secs: u64) -> String {
/// A journal unix stamp as the UTC time a person reads
/// (`2026-09-02 21:36:47 UTC`), the one rendering every listing verb
/// uses so a run's start, its events and its log lines line up by eye,
/// with the install's own logs too, wherever the install runs. Zero (a
/// row that never carried a stamp) renders as a dash rather than as
/// 1970.
pub fn utc_time(unix_secs: u64) -> String {
if unix_secs == 0 {
return "-".to_string();
}
// A stamp too large for the calendar prints as its number rather
// than as a wrapped-around date.
match i64::try_from(unix_secs).ok().and_then(|s| chrono::DateTime::<chrono::Utc>::from_timestamp(s, 0)) {
Some(t) => t.with_timezone(&chrono::Local).format("%Y-%m-%d %H:%M:%S").to_string(),
Some(t) => t.format("%Y-%m-%d %H:%M:%S UTC").to_string(),
None => unix_secs.to_string(),
}
}

#[cfg(test)]
mod tests {
use super::{local_time, resolve_dispatcher_url, Project};
use super::{utc_time, resolve_dispatcher_url, Project};

/// The stamp renders as a date and a time, and the two sentinels
/// (zero, out of range) never render as a bogus date.
#[test]
fn renders_a_readable_local_time() {
let text = local_time(1_756_838_207);
assert_eq!(text.len(), "2026-09-02 21:36:47".len(), "{text}");
assert_eq!(&text[4..5], "-");
assert_eq!(&text[10..11], " ");
assert_eq!(local_time(0), "-");
assert_eq!(local_time(u64::MAX), u64::MAX.to_string());
fn renders_a_readable_utc_time() {
assert_eq!(utc_time(1_756_838_207), "2025-09-02 18:36:47 UTC");
assert_eq!(utc_time(0), "-");
assert_eq!(utc_time(u64::MAX), u64::MAX.to_string());
}

fn project_with_prod() -> (tempfile::TempDir, Project) {
Expand Down
Loading
Loading