Compare commits
17
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bb274d08c6 | ||
|
|
7e07c389c6 | ||
|
|
389b41f8e6 | ||
|
|
ac6bf72943 | ||
|
|
0d9498ec6e | ||
|
|
15e7608e4a | ||
|
|
5d98fcf44a | ||
|
|
4ff4e6f7ee | ||
|
|
deb60be98d | ||
|
|
758b2dbd96 | ||
|
|
37fac288d2 | ||
|
|
5232175c88 | ||
|
|
ac47dcbe94 | ||
|
|
ad89ef94cd | ||
|
|
9de2cf34e4 | ||
|
|
3124fd3c8f | ||
|
|
cf076bd8ea |
Generated
+12
@@ -970,6 +970,7 @@ dependencies = [
|
||||
"serde_yaml",
|
||||
"sha2",
|
||||
"sqlx",
|
||||
"tar",
|
||||
"tempfile",
|
||||
"thiserror 2.0.18",
|
||||
"time",
|
||||
@@ -5029,6 +5030,17 @@ dependencies = [
|
||||
"windows",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "tar"
|
||||
version = "0.4.46"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3f6221d9a6003c78398e3b239969f352578258df48c8eb051caadae0015bc840"
|
||||
dependencies = [
|
||||
"filetime",
|
||||
"libc",
|
||||
"xattr",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "tempfile"
|
||||
version = "3.27.0"
|
||||
|
||||
@@ -36,6 +36,9 @@ publish = false
|
||||
# Shared dependency versions; crates opt in via { workspace = true }.
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
# Streaming tar for mission copy-in/copy-out (no compression: the payload is
|
||||
# a git checkout on a local socket, so CPU spent zipping buys nothing).
|
||||
tar = "0.4"
|
||||
thiserror = "2"
|
||||
uuid = { version = "1", features = ["v7", "serde"] }
|
||||
proptest = "1"
|
||||
|
||||
@@ -33,6 +33,7 @@ cm-config = { path = "../cm-config" }
|
||||
cm-db = { path = "../cm-db" }
|
||||
cm-domain = { path = "../cm-domain" }
|
||||
cm-files = { path = "../cm-files" }
|
||||
tar = { workspace = true }
|
||||
cm-llm = { path = "../cm-llm" }
|
||||
cm-orchestrator = { path = "../cm-orchestrator", features = ["provider"] }
|
||||
cm-runtime = { path = "../cm-runtime" }
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
//! Merging a delivered branch into the base, when that is provably safe.
|
||||
//!
|
||||
//! Every mission type delivers to a branch and never to `main`. For most that
|
||||
//! is where it should stop — a human reads the code and merges. But some
|
||||
//! missions only ever *add* files in a folder they own: a paper catalogue, a
|
||||
//! benchmark record. Those branches carry no judgement call, and leaving them
|
||||
//! to pile up unmerged means the work is done but not actually in the vault.
|
||||
//!
|
||||
//! # Additive-only is a property, not a preference
|
||||
//!
|
||||
//! The gate is not "is this mission type trusted". It is measured from the
|
||||
//! diff: if the branch modifies or deletes anything that already existed, it
|
||||
//! does not qualify, whatever its template says. A research harvest that
|
||||
//! somehow rewrote a hand-written note would be refused by the same check
|
||||
//! that lets its new notes through.
|
||||
//!
|
||||
//! Three conditions, all required:
|
||||
//!
|
||||
//! 1. the mission type declares [`MergePolicy::AdditiveOnly`]
|
||||
//! 2. verification passed — a run that did not prove its work does not merge
|
||||
//! 3. the diff against the base contains only additions
|
||||
//!
|
||||
//! Anything else lands as a branch for a human, which is the existing
|
||||
//! behaviour and the safe default.
|
||||
|
||||
use std::path::Path;
|
||||
|
||||
/// What a mission type is allowed to do with its own branch.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum MergePolicy {
|
||||
/// Always leave the branch for a human. Correct for anything that touches
|
||||
/// code: `refactor`, `research_and_code`, security patches.
|
||||
Never,
|
||||
/// Merge automatically when the diff is provably additive and the run
|
||||
/// verified. Correct for catalogues and recorded measurements.
|
||||
AdditiveOnly,
|
||||
}
|
||||
|
||||
impl MergePolicy {
|
||||
/// Parse a template's `merge_policy`. Unknown values fall back to `Never`
|
||||
/// and say so: a typo must not silently grant auto-merge.
|
||||
pub fn parse(raw: Option<&str>) -> MergePolicy {
|
||||
match raw.map(str::trim) {
|
||||
Some("additive_only") => MergePolicy::AdditiveOnly,
|
||||
Some("never") | None => MergePolicy::Never,
|
||||
Some(other) => {
|
||||
eprintln!(
|
||||
"auto_merge: unknown merge_policy {other:?} — refusing to auto-merge"
|
||||
);
|
||||
MergePolicy::Never
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Why a branch was or was not merged. The reason is always recorded: a
|
||||
/// branch that silently did not merge is indistinguishable from one that was
|
||||
/// never delivered.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct MergeOutcome {
|
||||
pub merged: bool,
|
||||
pub reason: String,
|
||||
}
|
||||
|
||||
impl MergeOutcome {
|
||||
fn refused(reason: impl Into<String>) -> MergeOutcome {
|
||||
MergeOutcome {
|
||||
merged: false,
|
||||
reason: reason.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Classify a `git diff --name-status` body.
|
||||
///
|
||||
/// Returns the offending entries, empty when every change is an addition.
|
||||
/// Split out so the rule is testable without a repository.
|
||||
pub fn non_additive_changes(name_status: &str) -> Vec<String> {
|
||||
name_status
|
||||
.lines()
|
||||
.filter(|l| !l.trim().is_empty())
|
||||
.filter(|l| {
|
||||
// Status is the first field: A/M/D/R###/C###.
|
||||
!matches!(l.chars().next(), Some('A'))
|
||||
})
|
||||
.map(|l| l.trim().to_string())
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn git(repo: &Path, args: &[&str]) -> Result<String, String> {
|
||||
let out = tokio::process::Command::new("git")
|
||||
.arg("-C")
|
||||
.arg(repo)
|
||||
.args(["-c", &format!("safe.directory={}", repo.display())])
|
||||
.args(args)
|
||||
.env("GIT_AUTHOR_NAME", crate::mission_delivery::commit_identity().0)
|
||||
.env("GIT_AUTHOR_EMAIL", crate::mission_delivery::commit_identity().1)
|
||||
.env(
|
||||
"GIT_COMMITTER_NAME",
|
||||
crate::mission_delivery::commit_identity().0,
|
||||
)
|
||||
.env(
|
||||
"GIT_COMMITTER_EMAIL",
|
||||
crate::mission_delivery::commit_identity().1,
|
||||
)
|
||||
.output()
|
||||
.await
|
||||
.map_err(|e| format!("spawn git: {e}"))?;
|
||||
if !out.status.success() {
|
||||
return Err(format!(
|
||||
"git {} → {}: {}",
|
||||
args.first().copied().unwrap_or("?"),
|
||||
out.status,
|
||||
crate::mission_workspace::redact_token(&String::from_utf8_lossy(&out.stderr))
|
||||
.chars()
|
||||
.take(300)
|
||||
.collect::<String>()
|
||||
));
|
||||
}
|
||||
Ok(String::from_utf8_lossy(&out.stdout).into_owned())
|
||||
}
|
||||
|
||||
/// Merge `branch` into `base` and push, if all three conditions hold.
|
||||
///
|
||||
/// Never returns `Err` for a refusal — a refusal is a normal outcome with a
|
||||
/// reason. `Err` is reserved for the merge itself going wrong after we decided
|
||||
/// to attempt it.
|
||||
pub async fn try_merge(
|
||||
repo: &Path,
|
||||
push_url: &str,
|
||||
branch: &str,
|
||||
base: &str,
|
||||
policy: MergePolicy,
|
||||
verified: bool,
|
||||
) -> Result<MergeOutcome, String> {
|
||||
if policy != MergePolicy::AdditiveOnly {
|
||||
return Ok(MergeOutcome::refused(
|
||||
"merge_policy is not additive_only; left for a human",
|
||||
));
|
||||
}
|
||||
if !verified {
|
||||
return Ok(MergeOutcome::refused(
|
||||
"run did not verify; refusing to merge unproven work",
|
||||
));
|
||||
}
|
||||
|
||||
// Compare against the base as the REMOTE has it, not a local ref that may
|
||||
// be stale. `...` gives changes on the branch since it diverged, so an
|
||||
// unrelated commit landing on main meanwhile is not misread as ours.
|
||||
git(repo, &["fetch", push_url, base]).await?;
|
||||
let diff = git(
|
||||
repo,
|
||||
&["diff", "--name-status", &format!("FETCH_HEAD...{branch}")],
|
||||
)
|
||||
.await?;
|
||||
|
||||
let offending = non_additive_changes(&diff);
|
||||
if !offending.is_empty() {
|
||||
return Ok(MergeOutcome::refused(format!(
|
||||
"diff is not additive ({} non-add change(s), first: {}); left for a human",
|
||||
offending.len(),
|
||||
offending.first().map(String::as_str).unwrap_or("?")
|
||||
)));
|
||||
}
|
||||
if diff.trim().is_empty() {
|
||||
return Ok(MergeOutcome::refused("branch adds nothing"));
|
||||
}
|
||||
|
||||
// Merge onto the freshly fetched base rather than a local branch.
|
||||
git(repo, &["checkout", "-B", base, "FETCH_HEAD"]).await?;
|
||||
if let Err(e) = git(
|
||||
repo,
|
||||
&["merge", "--no-ff", "-m", &format!("auto-merge {branch}"), branch],
|
||||
)
|
||||
.await
|
||||
{
|
||||
// Leave the repo clean so the next run is not fighting a wedged merge.
|
||||
let _ = git(repo, &["merge", "--abort"]).await;
|
||||
return Ok(MergeOutcome::refused(format!(
|
||||
"merge conflicted ({e}); left for a human"
|
||||
)));
|
||||
}
|
||||
|
||||
git(repo, &["push", push_url, &format!("HEAD:refs/heads/{base}")]).await?;
|
||||
Ok(MergeOutcome {
|
||||
merged: true,
|
||||
reason: format!("additive-only and verified; merged into {base}"),
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn only_pure_additions_qualify() {
|
||||
assert!(non_additive_changes("A\t60 Papers/a.md\nA\t60 Papers/b.md\n").is_empty());
|
||||
|
||||
// A modification disqualifies the whole branch.
|
||||
let m = non_additive_changes("A\t60 Papers/a.md\nM\tREADME.md\n");
|
||||
assert_eq!(m.len(), 1);
|
||||
assert!(m[0].contains("README.md"));
|
||||
|
||||
// So do deletes and renames — a rename is a delete plus an add, and
|
||||
// the delete half can destroy hand-written work.
|
||||
assert_eq!(non_additive_changes("D\tnotes/old.md\n").len(), 1);
|
||||
assert_eq!(non_additive_changes("R100\ta.md\tb.md\n").len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_unknown_policy_never_grants_auto_merge() {
|
||||
assert_eq!(MergePolicy::parse(None), MergePolicy::Never);
|
||||
assert_eq!(MergePolicy::parse(Some("never")), MergePolicy::Never);
|
||||
assert_eq!(
|
||||
MergePolicy::parse(Some("additive_only")),
|
||||
MergePolicy::AdditiveOnly
|
||||
);
|
||||
// A typo must fail closed, not open.
|
||||
assert_eq!(MergePolicy::parse(Some("aditive_only")), MergePolicy::Never);
|
||||
assert_eq!(MergePolicy::parse(Some("always")), MergePolicy::Never);
|
||||
}
|
||||
}
|
||||
@@ -291,6 +291,35 @@ pub async fn unseen(
|
||||
.collect())
|
||||
}
|
||||
|
||||
/// How many NEW sources a mission contributed.
|
||||
///
|
||||
/// The verification predicate for a continuous research mission. `record`
|
||||
/// never reassigns `mission_id` on conflict, so the first mission to find a
|
||||
/// source keeps the credit and a rerun cannot inflate its own count by
|
||||
/// re-recording what an earlier run already had.
|
||||
///
|
||||
/// A mission whose answer is zero produced nothing, whatever its transcript
|
||||
/// says — which is the check the 0030-0044 generation of this feature lacked.
|
||||
pub async fn contributed(
|
||||
pool: &sqlx::PgPool,
|
||||
workspace_id: Uuid,
|
||||
corpus_id: &str,
|
||||
mission_id: Uuid,
|
||||
) -> Result<i64, String> {
|
||||
let row: (i64,) = sqlx::query_as(
|
||||
"SELECT count(*) FROM corpus_items
|
||||
WHERE workspace_id = $1 AND corpus_id = $2 AND mission_id = $3
|
||||
AND kind = 'source'",
|
||||
)
|
||||
.bind(workspace_id)
|
||||
.bind(corpus_id)
|
||||
.bind(mission_id)
|
||||
.fetch_one(pool)
|
||||
.await
|
||||
.map_err(|e| format!("contributed({mission_id}): {e}"))?;
|
||||
Ok(row.0)
|
||||
}
|
||||
|
||||
/// Index every note in a checkout. Idempotent by construction.
|
||||
pub async fn index_vault(
|
||||
pool: &sqlx::PgPool,
|
||||
|
||||
@@ -16,12 +16,15 @@ mod mcp_door;
|
||||
mod mcp_skills;
|
||||
pub mod mission_orchestrator;
|
||||
pub mod mission_refiner;
|
||||
pub mod auto_merge;
|
||||
pub mod corpus;
|
||||
pub mod harvest;
|
||||
pub mod library;
|
||||
pub mod mission_delivery;
|
||||
pub mod mission_fs;
|
||||
pub mod papers;
|
||||
pub mod phase_config;
|
||||
pub mod session_executor;
|
||||
pub mod runtime_preflight;
|
||||
pub mod mission_runtime;
|
||||
pub mod mission_workspace;
|
||||
|
||||
@@ -35,6 +35,11 @@ pub struct LibraryRun {
|
||||
/// papers but could not push still has the PDFs and the checkmarks; the
|
||||
/// notes are simply not on the forge yet.
|
||||
pub pushed: bool,
|
||||
/// Whether the branch was auto-merged into `main`.
|
||||
pub merged: bool,
|
||||
/// Always populated — a branch that quietly did not merge is
|
||||
/// indistinguishable from one that was never delivered.
|
||||
pub merge_reason: String,
|
||||
pub error: Option<String>,
|
||||
}
|
||||
|
||||
@@ -160,6 +165,8 @@ pub async fn run_to_vault(
|
||||
harvest: total,
|
||||
branch,
|
||||
pushed: false,
|
||||
merged: false,
|
||||
merge_reason: "nothing new to push".into(),
|
||||
error: None,
|
||||
});
|
||||
}
|
||||
@@ -181,16 +188,41 @@ pub async fn run_to_vault(
|
||||
let auth = mission_workspace::with_ambient_auth(clone_url);
|
||||
let refspec = format!("HEAD:refs/heads/{branch}");
|
||||
match git(&vault, &["push", &auth, &refspec]).await {
|
||||
Ok(_) => Ok(LibraryRun {
|
||||
harvest: total,
|
||||
branch,
|
||||
pushed: true,
|
||||
error: None,
|
||||
}),
|
||||
Ok(_) => {
|
||||
// A catalogue branch only ever adds notes under `60 Papers/`, so
|
||||
// it qualifies for auto-merge — but the check is measured from the
|
||||
// diff, not assumed from the mission type. Verified here means the
|
||||
// run shelved something and errored on nothing.
|
||||
let verified = total.healthy() && !total.shelved.is_empty();
|
||||
let merge = crate::auto_merge::try_merge(
|
||||
&vault,
|
||||
&auth,
|
||||
&branch,
|
||||
"main",
|
||||
crate::auto_merge::MergePolicy::AdditiveOnly,
|
||||
verified,
|
||||
)
|
||||
.await
|
||||
.unwrap_or_else(|e| crate::auto_merge::MergeOutcome {
|
||||
merged: false,
|
||||
reason: format!("merge attempt failed: {e}"),
|
||||
});
|
||||
eprintln!("library: branch {branch} — {}", merge.reason);
|
||||
Ok(LibraryRun {
|
||||
harvest: total,
|
||||
branch,
|
||||
pushed: true,
|
||||
merged: merge.merged,
|
||||
merge_reason: merge.reason,
|
||||
error: None,
|
||||
})
|
||||
}
|
||||
Err(e) => Ok(LibraryRun {
|
||||
harvest: total,
|
||||
branch,
|
||||
pushed: false,
|
||||
merged: false,
|
||||
merge_reason: "not pushed, so not merged".into(),
|
||||
error: Some(e),
|
||||
}),
|
||||
}
|
||||
|
||||
@@ -583,6 +583,7 @@ pub async fn commit_phase_work(
|
||||
String::new()
|
||||
}
|
||||
);
|
||||
clear_stale_commit_editmsg(repo);
|
||||
git(repo, &["commit", "--no-verify", "-m", &message]).await?;
|
||||
}
|
||||
|
||||
@@ -878,6 +879,34 @@ impl TestOutcome {
|
||||
}
|
||||
}
|
||||
|
||||
/// Remove a `COMMIT_EDITMSG` the agent left behind as root.
|
||||
///
|
||||
/// The checkout is shared between the server (uid 65532) and the agent
|
||||
/// container (root). `core.sharedRepository` makes git create *objects and
|
||||
/// refs* group-writable — `.git/index` lands as 0666, which is why commits
|
||||
/// work at all — but it does not cover `COMMIT_EDITMSG`, which git writes
|
||||
/// with the default umask. An agent that runs `git commit` itself leaves that
|
||||
/// file owned by root at 0644, and the server's next commit dies with:
|
||||
///
|
||||
/// ```text
|
||||
/// git commit → exit 128: could not open '.git/COMMIT_EDITMSG': Permission denied
|
||||
/// ```
|
||||
///
|
||||
/// Observed on mission `019fcd0c`, which produced correct work — a reviewed,
|
||||
/// tested function plus a REVIEW.md quoting a real `cargo test` summary — and
|
||||
/// then delivered none of it.
|
||||
///
|
||||
/// Unlinking works where overwriting does not: removing a file requires write
|
||||
/// permission on the *directory*, and `.git/` is owned by the server. Silent
|
||||
/// on failure by design — if the file is absent or cannot be removed, the
|
||||
/// commit below reports the real error rather than this speculative cleanup.
|
||||
fn clear_stale_commit_editmsg(repo: &Path) {
|
||||
let msg = repo.join(".git/COMMIT_EDITMSG");
|
||||
if msg.exists() {
|
||||
let _ = std::fs::remove_file(&msg);
|
||||
}
|
||||
}
|
||||
|
||||
/// Mark a phase as impossible to capture, so it stops being selected.
|
||||
///
|
||||
/// A phase whose checkout has already been reaped can never be captured. It
|
||||
|
||||
@@ -0,0 +1,273 @@
|
||||
//! Move a mission's checkout in and out of its container, instead of sharing it.
|
||||
//!
|
||||
//! Today the checkout lives on the host and is bind-mounted into the mission
|
||||
//! container. That single directory is written by **two users** — cm-api as
|
||||
//! uid 65532 and the agent as root — and every bug that pattern can produce,
|
||||
//! it has produced:
|
||||
//!
|
||||
//! | Symptom | Fix that was needed |
|
||||
//! |---|---|
|
||||
//! | `.git/objects` permission denied | `core.sharedRepository=0777` |
|
||||
//! | capture base overwritten each phase | advance the base after commit |
|
||||
//! | `.git/COMMIT_EDITMSG` root-owned | unlink before commit |
|
||||
//! | `reset --hard` deleting a prior phase | `.git/clawmates-in-use` marker |
|
||||
//!
|
||||
//! Four fixes, one cause. `core.sharedRepository` was never a general
|
||||
//! solution — it covers objects and refs, and every *other* file git touches
|
||||
//! is a fresh opportunity.
|
||||
//!
|
||||
//! Copy-in/copy-out removes the cause: the agent owns its filesystem
|
||||
//! completely, as root, with no other writer. Nothing on the host is shared,
|
||||
//! so nothing on the host can collide.
|
||||
//!
|
||||
//! # Cost
|
||||
//!
|
||||
//! Measured on gw-04 against a real 65 MB checkout of this repository:
|
||||
//! **0.23s in, 0.18s out**. That was the one open risk in the plan — a
|
||||
//! monorepo copied per phase — and it is not a risk at this size. Measure
|
||||
//! again before assuming it holds for a repository an order of magnitude
|
||||
//! larger.
|
||||
//!
|
||||
//! No compression: the payload crosses a local Docker socket, so gzip would
|
||||
//! spend CPU to save nothing.
|
||||
|
||||
use std::path::Path;
|
||||
|
||||
use bollard::Docker;
|
||||
|
||||
/// Where a mission's checkout lives inside its container.
|
||||
pub const CONTAINER_MISSION_DIR: &str = "/mission";
|
||||
|
||||
/// Pack a host directory into an uncompressed tar.
|
||||
///
|
||||
/// `name_in_archive` is the top-level entry, so unpacking at
|
||||
/// [`CONTAINER_MISSION_DIR`] yields `/mission/<name>`. Kept separate from the
|
||||
/// upload so the packing is testable without Docker.
|
||||
pub fn pack_dir(root: &Path, name_in_archive: &str) -> Result<Vec<u8>, String> {
|
||||
let mut builder = tar::Builder::new(Vec::new());
|
||||
// Follow no symlinks: a checkout can contain a link pointing outside the
|
||||
// tree, and dereferencing it would pull host files into the container.
|
||||
builder.follow_symlinks(false);
|
||||
builder
|
||||
.append_dir_all(name_in_archive, root)
|
||||
.map_err(|e| format!("pack {}: {e}", root.display()))?;
|
||||
builder
|
||||
.into_inner()
|
||||
.map_err(|e| format!("finish archive for {}: {e}", root.display()))
|
||||
}
|
||||
|
||||
/// Unpack a tar into a host directory.
|
||||
///
|
||||
/// `tar` refuses entries whose paths escape the destination, which is the
|
||||
/// property that matters here: the archive comes back from a container the
|
||||
/// agent controls as root, so it is untrusted input. A `../../etc` entry must
|
||||
/// not be able to write outside the collection directory.
|
||||
pub fn unpack_into(archive: &[u8], dest: &Path) -> Result<(), String> {
|
||||
std::fs::create_dir_all(dest).map_err(|e| format!("mkdir {}: {e}", dest.display()))?;
|
||||
let mut ar = tar::Archive::new(archive);
|
||||
ar.set_overwrite(true);
|
||||
// Ownership in the archive is the container's root; re-applying it on the
|
||||
// host would recreate the very uid split this module exists to remove.
|
||||
ar.set_preserve_permissions(false);
|
||||
ar.unpack(dest)
|
||||
.map_err(|e| format!("unpack into {}: {e}", dest.display()))
|
||||
}
|
||||
|
||||
/// Copy a host directory into a running container at [`CONTAINER_MISSION_DIR`].
|
||||
pub async fn copy_in(
|
||||
docker: &Docker,
|
||||
container: &str,
|
||||
host_dir: &Path,
|
||||
name_in_archive: &str,
|
||||
) -> Result<(), String> {
|
||||
let archive = pack_dir(host_dir, name_in_archive)?;
|
||||
let opts = bollard::query_parameters::UploadToContainerOptionsBuilder::default()
|
||||
.path(CONTAINER_MISSION_DIR)
|
||||
.build();
|
||||
docker
|
||||
.upload_to_container(container, Some(opts), bollard::body_full(archive.into()))
|
||||
.await
|
||||
.map_err(|e| format!("copy into {container}:{CONTAINER_MISSION_DIR}: {e}"))
|
||||
}
|
||||
|
||||
/// Copy a directory back out of a container onto the host.
|
||||
pub async fn copy_out(
|
||||
docker: &Docker,
|
||||
container: &str,
|
||||
container_path: &str,
|
||||
dest: &Path,
|
||||
) -> Result<(), String> {
|
||||
use futures::StreamExt;
|
||||
|
||||
let opts = bollard::query_parameters::DownloadFromContainerOptionsBuilder::default()
|
||||
.path(container_path)
|
||||
.build();
|
||||
let mut stream = docker.download_from_container(container, Some(opts));
|
||||
let mut archive = Vec::new();
|
||||
while let Some(chunk) = stream.next().await {
|
||||
let bytes = chunk.map_err(|e| format!("copy out of {container}:{container_path}: {e}"))?;
|
||||
archive.extend_from_slice(&bytes);
|
||||
}
|
||||
unpack_into(&archive, dest)
|
||||
}
|
||||
|
||||
/// Is the copy-in/copy-out filesystem model enabled?
|
||||
///
|
||||
/// Opt-in. The bind-mount path is what production has run since the beginning,
|
||||
/// and silently changing how every mission receives its code is exactly the
|
||||
/// class of change that should require someone to have typed it.
|
||||
pub fn copy_mode() -> bool {
|
||||
matches!(
|
||||
std::env::var("CLAWMATES_MISSION_FS").as_deref(),
|
||||
Ok("copy")
|
||||
)
|
||||
}
|
||||
|
||||
/// Host directory holding a mission's checkout.
|
||||
fn host_repo(mission_id: uuid::Uuid) -> std::path::PathBuf {
|
||||
crate::mission_workspace::checkout_path(mission_id)
|
||||
}
|
||||
|
||||
/// Push the host checkout into the container before a phase runs.
|
||||
///
|
||||
/// No-op when the mission has no repo — research-only missions have no
|
||||
/// checkout, and that must not fail a phase launch.
|
||||
pub async fn sync_in(container: &str, mission_id: uuid::Uuid) -> Result<(), String> {
|
||||
let repo = host_repo(mission_id);
|
||||
if !repo.is_dir() {
|
||||
return Ok(());
|
||||
}
|
||||
let docker = crate::container_exec::connect()?;
|
||||
copy_in(&docker, container, &repo, "repo").await
|
||||
}
|
||||
|
||||
/// Pull the agent's work back onto the host after a phase.
|
||||
///
|
||||
/// Unpacks over the SAME host path the checkout came from, so the host
|
||||
/// directory stays a server-owned staging area with exactly one writer — and
|
||||
/// `mission_delivery::capture_phase_diff_at` needs no change at all, because
|
||||
/// it still finds a normal checkout exactly where it always has.
|
||||
pub async fn sync_out(container: &str, mission_id: uuid::Uuid) -> Result<(), String> {
|
||||
let repo = host_repo(mission_id);
|
||||
if !repo.is_dir() {
|
||||
return Ok(());
|
||||
}
|
||||
let parent = repo
|
||||
.parent()
|
||||
.ok_or_else(|| format!("{} has no parent", repo.display()))?;
|
||||
let docker = crate::container_exec::connect()?;
|
||||
copy_out(&docker, container, "/mission/repo", parent).await
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn seed(root: &Path) {
|
||||
std::fs::create_dir_all(root.join("src")).unwrap();
|
||||
std::fs::create_dir_all(root.join(".git")).unwrap();
|
||||
std::fs::write(root.join("src/lib.rs"), "pub fn x() {}\n").unwrap();
|
||||
std::fs::write(root.join(".git/HEAD"), "ref: refs/heads/main\n").unwrap();
|
||||
}
|
||||
|
||||
/// A checkout must survive the round trip intact — including `.git`,
|
||||
/// without which the whole delivery path (diff, commit, push) is dead.
|
||||
#[test]
|
||||
fn a_checkout_round_trips_with_its_git_dir() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let src = tmp.path().join("repo");
|
||||
seed(&src);
|
||||
|
||||
let archive = pack_dir(&src, "repo").unwrap();
|
||||
let dest = tmp.path().join("out");
|
||||
unpack_into(&archive, &dest).unwrap();
|
||||
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(dest.join("repo/src/lib.rs")).unwrap(),
|
||||
"pub fn x() {}\n"
|
||||
);
|
||||
assert!(
|
||||
dest.join("repo/.git/HEAD").exists(),
|
||||
"the .git dir must survive or delivery has nothing to diff"
|
||||
);
|
||||
}
|
||||
|
||||
/// The archive comes back from a container the agent controls as root, so
|
||||
/// it is untrusted. An entry that climbs out of the destination must not
|
||||
/// be able to write to the host.
|
||||
#[test]
|
||||
fn an_archive_cannot_escape_the_destination() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let dest = tmp.path().join("dest");
|
||||
let canary = tmp.path().join("ESCAPED");
|
||||
|
||||
// The path has to be written into the header bytes directly: the tar
|
||||
// crate refuses to BUILD an entry containing `..`, which is itself
|
||||
// reassuring but means a hostile archive cannot be produced through
|
||||
// the safe API. A real attacker writes the bytes, so the test does.
|
||||
let body = b"pwned\n";
|
||||
let mut header = tar::Header::new_gnu();
|
||||
header.set_size(body.len() as u64);
|
||||
header.set_mode(0o644);
|
||||
header.set_entry_type(tar::EntryType::Regular);
|
||||
{
|
||||
let gnu = header.as_gnu_mut().expect("gnu header");
|
||||
let evil = b"../ESCAPED";
|
||||
gnu.name[..evil.len()].copy_from_slice(evil);
|
||||
}
|
||||
header.set_cksum();
|
||||
|
||||
let mut archive = Vec::new();
|
||||
archive.extend_from_slice(header.as_bytes());
|
||||
let mut block = [0u8; 512];
|
||||
block[..body.len()].copy_from_slice(body);
|
||||
archive.extend_from_slice(&block);
|
||||
archive.extend_from_slice(&[0u8; 1024]); // end-of-archive marker
|
||||
|
||||
let _ = unpack_into(&archive, &dest);
|
||||
assert!(
|
||||
!canary.exists(),
|
||||
"a ../ entry wrote outside the destination"
|
||||
);
|
||||
}
|
||||
|
||||
/// A symlink pointing at the host filesystem must be packed as a link,
|
||||
/// not followed and inlined — otherwise copy-in would smuggle host files
|
||||
/// into the container.
|
||||
#[test]
|
||||
fn symlinks_are_not_dereferenced_into_the_archive() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let src = tmp.path().join("repo");
|
||||
seed(&src);
|
||||
let secret = tmp.path().join("host-secret");
|
||||
std::fs::write(&secret, "TOP SECRET\n").unwrap();
|
||||
std::os::unix::fs::symlink(&secret, src.join("link")).unwrap();
|
||||
|
||||
let archive = pack_dir(&src, "repo").unwrap();
|
||||
let haystack = String::from_utf8_lossy(&archive);
|
||||
assert!(
|
||||
!haystack.contains("TOP SECRET"),
|
||||
"symlink target contents were inlined into the archive"
|
||||
);
|
||||
}
|
||||
|
||||
/// The switch must be explicit — a near-miss value leaves production on
|
||||
/// the proven bind-mount path rather than silently changing it.
|
||||
#[test]
|
||||
fn copy_mode_requires_the_exact_word() {
|
||||
for wrong in ["Copy", "copies", "bind", "1", "true", ""] {
|
||||
assert_ne!(wrong, "copy", "{wrong:?} must not enable copy mode");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_empty_directory_packs_without_error() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let src = tmp.path().join("empty");
|
||||
std::fs::create_dir_all(&src).unwrap();
|
||||
let archive = pack_dir(&src, "repo").unwrap();
|
||||
let dest = tmp.path().join("out");
|
||||
unpack_into(&archive, &dest).unwrap();
|
||||
assert!(dest.join("repo").is_dir());
|
||||
}
|
||||
}
|
||||
@@ -83,10 +83,28 @@ impl RuntimeAuth {
|
||||
///
|
||||
/// The other three are unrelated providers with no subscription equivalent, so
|
||||
/// they forward in both modes.
|
||||
///
|
||||
/// In subscription mode `CLAUDE_CODE_OAUTH_TOKEN` forwards instead. The
|
||||
/// original design assumed a persisted `claude /login` under a bind-mounted
|
||||
/// `$HOME`, but a *mission* container gets its own data dir and therefore no
|
||||
/// login — so the token has to travel. Missing it is not a loud failure:
|
||||
/// `claude -p` simply hangs with no credential, which is what a phase stuck
|
||||
/// at `running` for ten minutes looked like when this was first switched on.
|
||||
pub fn forwarded_provider_keys(auth: RuntimeAuth) -> Vec<&'static str> {
|
||||
let mut keys = vec!["GEMINI_API_KEY", "GROQ_API_KEY", "OPENAI_API_KEY"];
|
||||
if auth == RuntimeAuth::ApiKey {
|
||||
keys.push("ANTHROPIC_API_KEY");
|
||||
// ZAI/KIMI reach their backends through the SAME `claude` binary via
|
||||
// ANTHROPIC_BASE_URL, so a mission that selects one needs its key present
|
||||
// in the container. They are unrelated to the Anthropic credential and
|
||||
// forward in both auth modes.
|
||||
let mut keys = vec![
|
||||
"GEMINI_API_KEY",
|
||||
"GROQ_API_KEY",
|
||||
"OPENAI_API_KEY",
|
||||
"ZAI_API_KEY",
|
||||
"KIMI_API_KEY",
|
||||
];
|
||||
match auth {
|
||||
RuntimeAuth::ApiKey => keys.push("ANTHROPIC_API_KEY"),
|
||||
RuntimeAuth::Subscription => keys.push("CLAUDE_CODE_OAUTH_TOKEN"),
|
||||
}
|
||||
keys
|
||||
}
|
||||
@@ -142,6 +160,146 @@ const MISSIONS_HOST_ROOT: &str = "/var/lib/clawmates-missions";
|
||||
/// mission so this rarely bites. Long-term: copy-on-write per mission.
|
||||
const DEFAULT_SEED_DIR: &str = "/root/clawmates-runtime/data";
|
||||
|
||||
|
||||
|
||||
/// What a mission gets its own copy of.
|
||||
///
|
||||
/// An allow-list, not the whole directory. The seed dir is **1.7 GB** on
|
||||
/// gw-04 and 1.5 GB of that is a vestigial `.rustup` — a Rust toolchain that
|
||||
/// installed itself into the data dir back when `HOME=/zeroclaw-data` and the
|
||||
/// image had no toolchain. The image now ships Rust at `/usr/local/cargo`,
|
||||
/// which is what the container's PATH actually resolves (verified live), so
|
||||
/// that copy is dead weight. Copying it per mission would cost tens of
|
||||
/// seconds and ~17 GB across ten concurrent missions.
|
||||
///
|
||||
/// So: copy what carries per-mission identity or secrets, and leave the
|
||||
/// caches and toolchains behind.
|
||||
const SEEDED_PATHS: &[&str] = &[
|
||||
// The whole point: config.toml carries the §15 door bearer token, and
|
||||
// data/ holds sessions.db + devices.db. ~26 MB.
|
||||
".zeroclaw",
|
||||
// Door MCP config — also a bearer token.
|
||||
"clawmates-mcp.json",
|
||||
// Claude Code's own state and credentials (~16 MB). Per-mission so a
|
||||
// token refresh or project state in one mission cannot leak into another.
|
||||
".claude",
|
||||
".claude.json",
|
||||
// Per-CLI state for the alternate backends; small.
|
||||
".kimi-code",
|
||||
"glm-home",
|
||||
// The seeded agent library.
|
||||
"agents",
|
||||
];
|
||||
|
||||
/// Deliberately NOT copied — caches and toolchains, no secrets, expensive:
|
||||
/// `.rustup` (1.5 GB, vestigial), `.npm` (85 MB), `.cargo`, `.cache`,
|
||||
/// `.local`. A mission that needs them reads the image's copies.
|
||||
fn copy_script() -> String {
|
||||
let mut out = String::from("set -e\n");
|
||||
for p in SEEDED_PATHS {
|
||||
// Missing entries are normal — a fresh deployment has no .kimi-code
|
||||
// until Kimi is first used — so absence must not fail the copy.
|
||||
out.push_str(&format!(
|
||||
"if [ -e '/seed/{p}' ]; then cp -a '/seed/{p}' /dst/; fi\n"
|
||||
));
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
/// Give this mission its own copy of the runtime seed data.
|
||||
///
|
||||
/// Every per-mission container used to bind-mount the SAME host seed dir as
|
||||
/// `/zeroclaw-data` — shared with each other and with the singleton runtime.
|
||||
/// That directory holds `config.toml`, which carries the §15 door bearer
|
||||
/// token, plus `sessions.db` and `devices.db`. So a mission could read another
|
||||
/// mission's credential, and anything it wrote there was inherited by every
|
||||
/// later mission. Teardown never cleaned it, because teardown only removes
|
||||
/// `/var/lib/clawmates-missions/{id}`.
|
||||
///
|
||||
/// The code already knew: the comment on `DEFAULT_SEED_DIR` names the sqlite
|
||||
/// race and calls copy-on-write per mission the long-term fix. This is that.
|
||||
///
|
||||
/// The copy runs in a throwaway container because cm-api cannot see the seed
|
||||
/// dir — it hands that host path to Docker but never mounts it itself. The
|
||||
/// runtime image is reused so nothing extra is pulled.
|
||||
///
|
||||
/// Failure is fatal to container creation on purpose. Falling back to the
|
||||
/// shared mount would silently restore the credential-sharing this removes,
|
||||
/// and a silent fallback to a weaker posture is the failure mode this
|
||||
/// codebase keeps paying for.
|
||||
async fn seed_runtime_data(
|
||||
docker: &Docker,
|
||||
image: &str,
|
||||
seed_dir: &str,
|
||||
dest_dir: &str,
|
||||
) -> Result<(), String> {
|
||||
let name = format!("cm-seed-{}", Uuid::now_v7().simple());
|
||||
let config = ContainerCreateBody {
|
||||
image: Some(image.to_string()),
|
||||
entrypoint: Some(vec!["/bin/sh".to_string()]),
|
||||
cmd: Some(vec!["-c".to_string(), copy_script()]),
|
||||
host_config: Some(HostConfig {
|
||||
mounts: Some(vec![
|
||||
Mount {
|
||||
target: Some("/seed".to_string()),
|
||||
source: Some(seed_dir.to_string()),
|
||||
typ: Some(MountTypeEnum::BIND),
|
||||
read_only: Some(true),
|
||||
..Default::default()
|
||||
},
|
||||
Mount {
|
||||
target: Some("/dst".to_string()),
|
||||
source: Some(dest_dir.to_string()),
|
||||
typ: Some(MountTypeEnum::BIND),
|
||||
read_only: Some(false),
|
||||
..Default::default()
|
||||
},
|
||||
]),
|
||||
auto_remove: Some(true),
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
docker
|
||||
.create_container(
|
||||
Some(CreateContainerOptions {
|
||||
name: Some(name.clone()),
|
||||
..Default::default()
|
||||
}),
|
||||
config,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| format!("create seed copier: {e}"))?;
|
||||
docker
|
||||
.start_container(&name, None::<StartContainerOptions>)
|
||||
.await
|
||||
.map_err(|e| format!("start seed copier: {e}"))?;
|
||||
|
||||
// `auto_remove` means the container disappears the moment it exits, so
|
||||
// poll for absence rather than waiting on it.
|
||||
for _ in 0..120 {
|
||||
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
|
||||
match docker
|
||||
.inspect_container(&name, None::<InspectContainerOptions>)
|
||||
.await
|
||||
{
|
||||
Err(_) => return Ok(()),
|
||||
Ok(info) => {
|
||||
let running = info
|
||||
.state
|
||||
.as_ref()
|
||||
.and_then(|st| st.running)
|
||||
.unwrap_or(false);
|
||||
if !running {
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(format!("seed copy into {dest_dir} did not finish in 30s"))
|
||||
}
|
||||
|
||||
/// Deterministic docker container name for a mission's runtime.
|
||||
/// Uses the full UUID hex — UUIDv7 encodes time in the leading bytes,
|
||||
/// so a short prefix isn't guaranteed unique across missions minted
|
||||
@@ -239,29 +397,40 @@ impl MissionRuntimeProvisioner {
|
||||
let _ = tokio::fs::create_dir_all(&mission_dir).await;
|
||||
let seed_dir = std::env::var("CLAWMATES_RUNTIME_SEED_DIR")
|
||||
.unwrap_or_else(|_| DEFAULT_SEED_DIR.to_string());
|
||||
let mounts = vec![
|
||||
// Mount just this mission's directory. Agents can navigate
|
||||
// its `/repo` subdir but never see other missions'.
|
||||
Mount {
|
||||
// Per-mission copy of the seed data. See `seed_runtime_data`: sharing
|
||||
// one directory meant sharing the door token and letting any mission
|
||||
// poison every later one.
|
||||
let runtime_data_dir = format!("{mission_dir}/runtime-data");
|
||||
let _ = tokio::fs::create_dir_all(&runtime_data_dir).await;
|
||||
seed_runtime_data(&self.docker, &self.image, &seed_dir, &runtime_data_dir).await?;
|
||||
let mut mounts = Vec::new();
|
||||
// In copy mode the checkout is pushed in and pulled back out, so the
|
||||
// container gets its OWN filesystem and the host directory has exactly
|
||||
// one writer (the server). Binding it here would put two uids back on
|
||||
// one directory — the cause of four separate work-loss bugs.
|
||||
if !crate::mission_fs::copy_mode() {
|
||||
mounts.push(Mount {
|
||||
target: Some("/mission".to_string()),
|
||||
source: Some(mission_dir.clone()),
|
||||
typ: Some(MountTypeEnum::BIND),
|
||||
read_only: Some(false),
|
||||
..Default::default()
|
||||
},
|
||||
// Share the shared-runtime data dir so this gateway inherits
|
||||
// the seeded agent library (`claw_*` templates). We then
|
||||
// mint a per-mission pairing code via /admin/paircode/new
|
||||
// below — the mint writes into the shared devices.db but
|
||||
// the resulting token is unique to this mission.
|
||||
});
|
||||
}
|
||||
mounts.extend([
|
||||
// This mission's OWN copy of the seeded agent library
|
||||
// (`claw_*` templates). Copied rather than shared, so its
|
||||
// config.toml — which carries the door bearer token — and its
|
||||
// sqlite files belong to this mission alone and are removed with
|
||||
// it by `teardown_container`.
|
||||
Mount {
|
||||
target: Some("/zeroclaw-data".to_string()),
|
||||
source: Some(seed_dir),
|
||||
source: Some(runtime_data_dir),
|
||||
typ: Some(MountTypeEnum::BIND),
|
||||
read_only: Some(false),
|
||||
..Default::default()
|
||||
},
|
||||
];
|
||||
]);
|
||||
|
||||
let host_config = HostConfig {
|
||||
mounts: Some(mounts),
|
||||
@@ -802,9 +971,38 @@ mod tests {
|
||||
it silently bills the API. Forwarded: {keys:?}"
|
||||
);
|
||||
// Unrelated providers have no subscription equivalent and must survive.
|
||||
for k in ["GEMINI_API_KEY", "GROQ_API_KEY", "OPENAI_API_KEY"] {
|
||||
for k in [
|
||||
"GEMINI_API_KEY",
|
||||
"GROQ_API_KEY",
|
||||
"OPENAI_API_KEY",
|
||||
"ZAI_API_KEY",
|
||||
"KIMI_API_KEY",
|
||||
] {
|
||||
assert!(keys.contains(&k), "{k} should still be forwarded");
|
||||
}
|
||||
// And the subscription credential MUST travel. A mission container
|
||||
// has its own data dir, so unlike the shared runtime it has no
|
||||
// persisted `claude /login` to fall back on. Without this the CLI
|
||||
// has no credential and simply hangs — a phase stuck at `running`
|
||||
// with nothing in the logs, which is exactly how this was found.
|
||||
assert!(
|
||||
keys.contains(&"CLAUDE_CODE_OAUTH_TOKEN"),
|
||||
"subscription mode must forward the token; without it `claude -p` \
|
||||
hangs with no credential. Forwarded: {keys:?}"
|
||||
);
|
||||
}
|
||||
|
||||
/// The two credentials must never travel together: Claude Code would pick
|
||||
/// the API key and bill it while the deployment believes it is on the
|
||||
/// subscription.
|
||||
#[test]
|
||||
fn the_two_anthropic_credentials_are_mutually_exclusive() {
|
||||
for mode in [RuntimeAuth::ApiKey, RuntimeAuth::Subscription] {
|
||||
let keys = forwarded_provider_keys(mode);
|
||||
let both = keys.contains(&"ANTHROPIC_API_KEY")
|
||||
&& keys.contains(&"CLAUDE_CODE_OAUTH_TOKEN");
|
||||
assert!(!both, "{mode:?} forwards both credentials: {keys:?}");
|
||||
}
|
||||
}
|
||||
|
||||
/// Default behaviour is unchanged, so a deployment that never opts in keeps
|
||||
@@ -820,6 +1018,10 @@ mod tests {
|
||||
] {
|
||||
assert!(keys.contains(&k), "{k} should be forwarded in api_key mode");
|
||||
}
|
||||
assert!(
|
||||
!keys.contains(&"CLAUDE_CODE_OAUTH_TOKEN"),
|
||||
"api_key mode must not also ship the subscription token"
|
||||
);
|
||||
}
|
||||
|
||||
/// An unset or misspelled value must fall back to the *existing* behaviour.
|
||||
@@ -939,4 +1141,75 @@ allowed_tools = ["file_read", "file_edit"]
|
||||
let url = endpoint_url("cm-runtime-mission-abc");
|
||||
assert_eq!(url, "http://cm-runtime-mission-abc:42617");
|
||||
}
|
||||
|
||||
/// The seed data path must be per-mission, not the shared seed dir.
|
||||
///
|
||||
/// Every mission container used to bind the SAME host directory as
|
||||
/// `/zeroclaw-data`. It holds `config.toml`, which carries the §15 door
|
||||
/// bearer token, plus `sessions.db`/`devices.db`. Sharing it meant one
|
||||
/// mission could read another's credential, and anything written there
|
||||
/// was inherited by every later mission — `teardown_container` only
|
||||
/// removes `/var/lib/clawmates-missions/{id}`, so the shared dir was
|
||||
/// never cleaned.
|
||||
///
|
||||
/// Asserting the path shape is what keeps this from silently regressing:
|
||||
/// a future edit that points the mount back at the seed dir restores the
|
||||
/// credential sharing with no other visible symptom.
|
||||
#[test]
|
||||
fn runtime_data_is_scoped_to_one_mission() {
|
||||
let a = Uuid::now_v7();
|
||||
let b = Uuid::now_v7();
|
||||
let path = |id: Uuid| format!("{MISSIONS_HOST_ROOT}/{id}/runtime-data");
|
||||
|
||||
assert_ne!(path(a), path(b), "two missions must not share runtime data");
|
||||
assert!(
|
||||
path(a).starts_with(&format!("{MISSIONS_HOST_ROOT}/{a}")),
|
||||
"runtime data must live under the mission dir so teardown removes it"
|
||||
);
|
||||
assert_ne!(
|
||||
path(a),
|
||||
DEFAULT_SEED_DIR,
|
||||
"the mount must never be the shared seed dir itself"
|
||||
);
|
||||
assert!(
|
||||
!path(a).starts_with(DEFAULT_SEED_DIR),
|
||||
"runtime data must not live inside the shared seed dir either"
|
||||
);
|
||||
}
|
||||
|
||||
/// The copy must be an allow-list, and must include the secret-bearing
|
||||
/// paths while excluding the expensive ones.
|
||||
///
|
||||
/// Measured on gw-04: the seed dir is 1.7 GB, of which 1.5 GB is a
|
||||
/// vestigial `.rustup` that is not even on the container's PATH (the
|
||||
/// image ships Rust at /usr/local/cargo). Copying everything per mission
|
||||
/// would cost tens of seconds and ~17 GB across ten concurrent missions —
|
||||
/// which is what the first version of this did.
|
||||
#[test]
|
||||
fn the_seed_copy_takes_secrets_and_skips_caches() {
|
||||
// The two paths that carry the door bearer token MUST be copied, or
|
||||
// this whole change accomplishes nothing.
|
||||
assert!(SEEDED_PATHS.contains(&".zeroclaw"));
|
||||
assert!(SEEDED_PATHS.contains(&"clawmates-mcp.json"));
|
||||
// Claude Code's credentials and state.
|
||||
assert!(SEEDED_PATHS.contains(&".claude"));
|
||||
|
||||
// The expensive, secret-free ones must NOT be.
|
||||
for cache in [".rustup", ".npm", ".cargo", ".cache"] {
|
||||
assert!(
|
||||
!SEEDED_PATHS.contains(&cache),
|
||||
"{cache} is a cache and must not be copied per mission"
|
||||
);
|
||||
}
|
||||
|
||||
let script = copy_script();
|
||||
// A missing entry is normal on a fresh deployment (no .kimi-code
|
||||
// until Kimi is first used) and must not fail the copy.
|
||||
assert!(script.contains("if [ -e "), "absent paths must be tolerated");
|
||||
assert!(script.contains("/seed/.zeroclaw"));
|
||||
assert!(!script.contains("/seed/.rustup"));
|
||||
for p in SEEDED_PATHS {
|
||||
assert!(script.contains(p), "{p} missing from the copy script");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -112,6 +112,23 @@ async fn capture_finished_coding_phases(pool: &PgPool) -> Result<(), String> {
|
||||
use sqlx::Row;
|
||||
let phase_id: Uuid = row.get("id");
|
||||
let mission_id: Uuid = row.get("mission_id");
|
||||
// Pull the agent's work back onto the host before capturing it.
|
||||
// Unpacks over the same checkout path, so capture below is unchanged.
|
||||
if crate::mission_fs::copy_mode() {
|
||||
let container = crate::mission_runtime::container_name(mission_id);
|
||||
if let Err(e) = crate::mission_fs::sync_out(&container, mission_id).await {
|
||||
// Loud, and skip capture: capturing now would diff a stale
|
||||
// host tree and record "no changes" for work that exists —
|
||||
// reporting success for nothing, which is the failure this
|
||||
// codebase keeps paying for.
|
||||
eprintln!(
|
||||
"phase_runner: could NOT collect work from {container} for phase \
|
||||
{phase_id} ({e}) — skipping capture so a stale tree is not \
|
||||
recorded as an empty diff"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
match crate::mission_delivery::capture_phase_diff(pool, mission_id, phase_id).await {
|
||||
Ok(Some(_)) => {}
|
||||
Ok(None) => {
|
||||
@@ -299,6 +316,15 @@ async fn launch_phase(pool: &PgPool, p: PhaseLaunch<'_>) -> Result<(), String> {
|
||||
match prov.ensure_container(mission_id).await {
|
||||
Ok(ec) => {
|
||||
let name = crate::mission_runtime::container_name(mission_id);
|
||||
// Push the checkout into the container. A no-op in bind mode;
|
||||
// in copy mode it is how the agent gets the code at all, so a
|
||||
// failure must fail the launch rather than silently starting a
|
||||
// phase against an empty directory.
|
||||
if crate::mission_fs::copy_mode() {
|
||||
if let Err(e) = crate::mission_fs::sync_in(&name, mission_id).await {
|
||||
return Err(format!("copy checkout into {name}: {e}"));
|
||||
}
|
||||
}
|
||||
if let Err(e) = cm_db::repo::missions::set_runtime_binding(
|
||||
pool,
|
||||
mission_id,
|
||||
@@ -346,6 +372,26 @@ async fn launch_phase(pool: &PgPool, p: PhaseLaunch<'_>) -> Result<(), String> {
|
||||
_ => task,
|
||||
};
|
||||
|
||||
// Direct-session executor: run the whole phase as ONE `claude -p` session
|
||||
// against the mission checkout, instead of driving turns through ZeroClaw.
|
||||
//
|
||||
// Measured on the same task against a real checkout: 7s direct versus
|
||||
// minutes per turn through the adapter, and the adapter needed three
|
||||
// rounds of config before it worked at all — a hang, a timeout, and a
|
||||
// mission that COMPLETED having written nothing. With claude_cli the
|
||||
// adapter is a WebSocket-to-subprocess shim whose own controls (risk
|
||||
// profiles, tool gating, memory) never reach the subprocess, so it adds
|
||||
// failure modes without adding governance.
|
||||
//
|
||||
// It still creates one `topology_runs` row. That is deliberate: the whole
|
||||
// downstream lifecycle — close_finished_phases, evaluation, capture,
|
||||
// delivery — keys off those rows, and inventing a second completion path
|
||||
// would mean two ways for a phase to finish and one of them untested.
|
||||
if crate::session_executor::direct_mode() {
|
||||
return launch_direct_session(pool, mission_id, phase_id, workspace_id, iteration, &task)
|
||||
.await;
|
||||
}
|
||||
|
||||
// Purge prior failed / cancelled runs for this phase so the card
|
||||
// starts fresh on re-attempts. Completed runs are kept for
|
||||
// auditability (a mission that succeeded once and got re-run
|
||||
@@ -408,6 +454,94 @@ async fn launch_phase(pool: &PgPool, p: PhaseLaunch<'_>) -> Result<(), String> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Launch a phase as a single headless session.
|
||||
///
|
||||
/// Returns as soon as the session is spawned: `launch_phase` runs inside the
|
||||
/// sweep loop, and blocking it for the length of a coding session would stall
|
||||
/// every other mission.
|
||||
async fn launch_direct_session(
|
||||
pool: &PgPool,
|
||||
mission_id: Uuid,
|
||||
phase_id: Uuid,
|
||||
workspace_id: Uuid,
|
||||
iteration: i32,
|
||||
task: &str,
|
||||
) -> Result<(), String> {
|
||||
sqlx::query(
|
||||
"DELETE FROM topology_runs
|
||||
WHERE mission_phase_id = $1 AND status IN ('failed', 'cancelled')",
|
||||
)
|
||||
.bind(phase_id)
|
||||
.execute(pool)
|
||||
.await
|
||||
.map_err(|e| format!("purge prior runs for phase {phase_id}: {e}"))?;
|
||||
|
||||
let run_id = Uuid::now_v7();
|
||||
sqlx::query(
|
||||
"INSERT INTO topology_runs
|
||||
(id, workspace_id, task, kind, status, graph, tier,
|
||||
mission_id, mission_phase_id, iteration)
|
||||
VALUES ($1, $2, $3, 'run', 'running', $4, 'session', $5, $6, $7)",
|
||||
)
|
||||
.bind(run_id)
|
||||
.bind(workspace_id)
|
||||
.bind(task)
|
||||
.bind(serde_json::json!({ "nodes": [], "edges": [], "executor": "session" }))
|
||||
.bind(mission_id)
|
||||
.bind(phase_id)
|
||||
.bind(iteration)
|
||||
.execute(pool)
|
||||
.await
|
||||
.map_err(|e| format!("enqueue session run for phase {phase_id}: {e}"))?;
|
||||
|
||||
sqlx::query(
|
||||
"UPDATE mission_phases
|
||||
SET status = 'running', started_at = now()
|
||||
WHERE id = $1 AND status = 'pending'",
|
||||
)
|
||||
.bind(phase_id)
|
||||
.execute(pool)
|
||||
.await
|
||||
.map_err(|e| format!("mark phase {phase_id} running: {e}"))?;
|
||||
|
||||
let container = crate::mission_runtime::container_name(mission_id);
|
||||
let task = task.to_string();
|
||||
let pool = pool.clone();
|
||||
tokio::spawn(async move {
|
||||
let repo = "/mission/repo";
|
||||
let branch = crate::session_executor::session_branch(mission_id);
|
||||
let (summary, exit) =
|
||||
match crate::session_executor::run_session(&container, repo, &task, &branch).await {
|
||||
Ok(v) => v,
|
||||
Err(e) => (format!("session failed to start: {e}"), None),
|
||||
};
|
||||
// The agent's own account is diagnostic only. Whether the phase
|
||||
// succeeded is decided downstream by capture + delivery against the
|
||||
// repository, never by this text.
|
||||
let ok = exit == Some(0);
|
||||
eprintln!(
|
||||
"phase_runner: session for mission {mission_id} phase {phase_id} exited {exit:?} — {}",
|
||||
summary.chars().take(200).collect::<String>()
|
||||
);
|
||||
let status = if ok { "completed" } else { "failed" };
|
||||
if let Err(e) = sqlx::query(
|
||||
"UPDATE topology_runs SET status = $2, updated_at = now() WHERE id = $1",
|
||||
)
|
||||
.bind(run_id)
|
||||
.bind(status)
|
||||
.execute(&pool)
|
||||
.await
|
||||
{
|
||||
eprintln!("phase_runner: could not close session run {run_id}: {e}");
|
||||
}
|
||||
});
|
||||
|
||||
eprintln!(
|
||||
"phase_runner: mission {mission_id} phase {phase_id} launched as a DIRECT SESSION"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn phase_task_text(
|
||||
kind: &str,
|
||||
title: &str,
|
||||
|
||||
@@ -26,6 +26,11 @@ pub struct RunRequest {
|
||||
/// library can otherwise pull hundreds of PDFs in one go.
|
||||
#[serde(default)]
|
||||
pub per_topic: Option<usize>,
|
||||
/// Attribute this run to a mission, so the mission can later be asked
|
||||
/// what it contributed. `corpus_items.mission_id` has existed since the
|
||||
/// table landed; without this field nothing could ever populate it.
|
||||
#[serde(default, rename = "missionId")]
|
||||
pub mission_id: Option<uuid::Uuid>,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
@@ -37,6 +42,8 @@ pub struct RunResponse {
|
||||
pub notes: Vec<String>,
|
||||
pub branch: String,
|
||||
pub pushed: bool,
|
||||
pub merged: bool,
|
||||
pub merge_reason: String,
|
||||
pub error: Option<String>,
|
||||
/// A run that errored on nothing. Reported explicitly so a caller does not
|
||||
/// have to infer health from an empty `shelved` list — a quiet week and a
|
||||
@@ -78,7 +85,7 @@ pub async fn run(
|
||||
&work_root,
|
||||
&topics,
|
||||
per_topic,
|
||||
None,
|
||||
req.mission_id,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
@@ -102,6 +109,8 @@ pub async fn run(
|
||||
healthy: out.harvest.healthy(),
|
||||
branch: out.branch,
|
||||
pushed: out.pushed,
|
||||
merged: out.merged,
|
||||
merge_reason: out.merge_reason,
|
||||
error: out.error,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -20,12 +20,20 @@ pub fn claw_alias(claw_id: Uuid) -> String {
|
||||
|
||||
/// Map a claw's chosen model to a configured provider alias.
|
||||
///
|
||||
/// v0.8.3 fold: `claude_cli.*` and `kimi_cli.*` families were deleted
|
||||
/// upstream; every alias now lives under a real provider family
|
||||
/// (`anthropic`, `groq`, `gemini`, ...). Our compose currently
|
||||
/// configures `anthropic.default`, `anthropic.door`, `groq.default`,
|
||||
/// and `gemini.default`, so unknown models resolve to
|
||||
/// `anthropic.default` — the workspace's high-quality baseline.
|
||||
/// Claude models resolve to `claude_cli.default`, which spawns the real
|
||||
/// `claude` binary against the Max subscription rather than posting to the
|
||||
/// raw API with Claude Code identity headers. The API-key path still exists
|
||||
/// and the judge uses it deliberately (see below), but agent work — which is
|
||||
/// ~99% of the tokens — belongs on the subscription and on the supported
|
||||
/// client.
|
||||
///
|
||||
/// The judge stays on `anthropic.judge`/API key on purpose: if the
|
||||
/// subscription throttles, missions degrade but verification keeps working.
|
||||
/// Putting both on one credential would mean a single limit blinds the
|
||||
/// verifier at exactly the moment there is most to verify.
|
||||
///
|
||||
/// Non-Claude families are unchanged: `groq.default`, `gemini.default`, and
|
||||
/// the GLM/Kimi substitution below.
|
||||
pub fn provider_alias_for(model: &str) -> &'static str {
|
||||
let m = model.trim().to_ascii_lowercase();
|
||||
// Prefix families first (covers claude-sonnet-5, claude-opus-4-8,
|
||||
@@ -33,7 +41,7 @@ pub fn provider_alias_for(model: &str) -> &'static str {
|
||||
// decides what "its own family" means, so the two can't drift apart.
|
||||
if is_exact_provider_match(&m) {
|
||||
if m.starts_with("claude") {
|
||||
return "anthropic.default";
|
||||
return "claude_cli.default";
|
||||
}
|
||||
if m.starts_with("gemini") {
|
||||
return "gemini.default";
|
||||
@@ -54,18 +62,18 @@ pub fn provider_alias_for(model: &str) -> &'static str {
|
||||
| "kimi" | "kimi-k2" | "kimi-for-coding" => {
|
||||
eprintln!(
|
||||
"runtime_provision: model {m:?} has no provider family configured — \
|
||||
substituting anthropic.default, which spends ANTHROPIC_API_KEY"
|
||||
substituting claude_cli.default, which spends the Claude subscription"
|
||||
);
|
||||
"anthropic.default"
|
||||
"claude_cli.default"
|
||||
}
|
||||
_ => {
|
||||
if !m.is_empty() {
|
||||
eprintln!(
|
||||
"runtime_provision: unrecognised model {m:?} — defaulting to \
|
||||
anthropic.default"
|
||||
claude_cli.default"
|
||||
);
|
||||
}
|
||||
"anthropic.default"
|
||||
"claude_cli.default"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -342,14 +350,15 @@ mod tests {
|
||||
|
||||
/// The GLM/Kimi substitution is intentional but must be reported as a
|
||||
/// substitution, because its consequence is that a user who picked a
|
||||
/// non-Anthropic model is spending the Anthropic key.
|
||||
/// non-Anthropic model is spending someone else's budget — now the
|
||||
/// Claude subscription rather than the Anthropic API key.
|
||||
#[test]
|
||||
fn substituted_families_are_not_reported_as_exact_matches() {
|
||||
for m in ["kimi", "glm-4.7", "glm5", "kimi-k2", "something-unknown"] {
|
||||
assert_eq!(super::provider_alias_for(m), "anthropic.default");
|
||||
assert_eq!(super::provider_alias_for(m), "claude_cli.default");
|
||||
assert!(
|
||||
!super::is_exact_provider_match(m),
|
||||
"{m} resolves to anthropic.default by substitution, not by family"
|
||||
"{m} resolves to claude_cli.default by substitution, not by family"
|
||||
);
|
||||
}
|
||||
for m in [
|
||||
@@ -395,19 +404,20 @@ mod tests {
|
||||
fn provider_alias_mapping() {
|
||||
assert_eq!(provider_alias_for("gemini"), "gemini.default");
|
||||
assert_eq!(provider_alias_for("gemini-2.0-flash"), "gemini.default");
|
||||
// v0.8.3: glm/kimi families fall back to anthropic until their
|
||||
// own provider tables are configured in the runtime template.
|
||||
assert_eq!(provider_alias_for("GLM-4.7"), "anthropic.default");
|
||||
assert_eq!(provider_alias_for("kimi"), "anthropic.default");
|
||||
// glm/kimi families fall back to Claude until their own provider
|
||||
// tables are configured in the runtime template.
|
||||
assert_eq!(provider_alias_for("GLM-4.7"), "claude_cli.default");
|
||||
assert_eq!(provider_alias_for("kimi"), "claude_cli.default");
|
||||
assert_eq!(provider_alias_for("groq"), "groq.default");
|
||||
assert_eq!(
|
||||
provider_alias_for("llama-3.3-70b-versatile"),
|
||||
"groq.default"
|
||||
);
|
||||
assert_eq!(provider_alias_for("claude"), "anthropic.default");
|
||||
assert_eq!(provider_alias_for("claude-sonnet-5"), "anthropic.default");
|
||||
assert_eq!(provider_alias_for("claude-opus-4-8"), "anthropic.default");
|
||||
assert_eq!(provider_alias_for("anything-else"), "anthropic.default");
|
||||
// Claude models spawn the real CLI against the subscription.
|
||||
assert_eq!(provider_alias_for("claude"), "claude_cli.default");
|
||||
assert_eq!(provider_alias_for("claude-sonnet-5"), "claude_cli.default");
|
||||
assert_eq!(provider_alias_for("claude-opus-4-8"), "claude_cli.default");
|
||||
assert_eq!(provider_alias_for("anything-else"), "claude_cli.default");
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -0,0 +1,262 @@
|
||||
//! Run a whole mission as ONE headless agent session.
|
||||
//!
|
||||
//! The alternative to `phase_runner`. Instead of splitting a mission into
|
||||
//! phases that hand work to each other through a shared checkout, this hands
|
||||
//! the entire task to a single agent session and asks the forge afterwards
|
||||
//! what actually landed.
|
||||
//!
|
||||
//! # Why
|
||||
//!
|
||||
//! The phase machinery moves state between processes through a filesystem, and
|
||||
//! that seam produced most of a week's defects: two uids fighting over
|
||||
//! `.git/objects`, a missing git identity, `reset --hard` deleting the
|
||||
//! previous phase's work, a capture base overloaded with two meanings. None of
|
||||
//! those failures are *possible* inside one session, because there is no
|
||||
//! handoff to get wrong — step two knows what step one did because it is the
|
||||
//! same context.
|
||||
//!
|
||||
//! Measured against the same task (create a file, read it back, extend it,
|
||||
//! push it): the phase path took nine production runs and five distinct bug
|
||||
//! fixes to do reliably; a single session did it in 23 seconds, 19 times out
|
||||
//! of 20, first try.
|
||||
//!
|
||||
//! # What this deliberately does NOT trust
|
||||
//!
|
||||
//! The agent's own account of what it did. In the same 60-run experiment one
|
||||
//! session exited 0, ran for 18 seconds, and pushed nothing — a clean exit
|
||||
//! status with no work delivered, about 5% of the time. That is the same
|
||||
//! "reported success while doing nothing" shape as every scaffolding bug, and
|
||||
//! it is why [`verify_landed`] asks the forge rather than reading the summary.
|
||||
//!
|
||||
//! Deleting the phase machinery is justified by the evidence. Deleting the
|
||||
//! verification is not — the evidence points the other way.
|
||||
|
||||
use std::time::Duration;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::container_exec;
|
||||
|
||||
/// Ceiling for one mission session. Long, because a real coding task with a
|
||||
/// test suite legitimately takes minutes; bounded, because a wedged session
|
||||
/// must not hold a container forever.
|
||||
const SESSION_TIMEOUT: Duration = Duration::from_secs(3600);
|
||||
|
||||
/// Tools the session may use without prompting.
|
||||
///
|
||||
/// `--dangerously-skip-permissions` is refused by the CLI when running as
|
||||
/// root, which mission containers do, and blanket bypass is the wrong default
|
||||
/// for something driving a real repository anyway. An explicit allow-list is
|
||||
/// both accepted as root and easier to defend.
|
||||
const ALLOWED_TOOLS: &[&str] = &["Read", "Edit", "Write", "Bash"];
|
||||
|
||||
/// What one session did, as observed from outside it.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SessionOutcome {
|
||||
/// The agent's closing summary. Diagnostic only — never evidence.
|
||||
pub summary: String,
|
||||
pub exit_code: Option<i64>,
|
||||
/// Whether the expected branch actually appeared on the forge.
|
||||
pub landed: bool,
|
||||
/// Head sha of the branch, when it landed.
|
||||
pub head_sha: Option<String>,
|
||||
}
|
||||
|
||||
impl SessionOutcome {
|
||||
/// The session both finished cleanly *and* delivered.
|
||||
///
|
||||
/// Both halves are required. `exit_code == Some(0)` alone is what the
|
||||
/// 5% silent-nothing case looks like from the inside.
|
||||
pub fn delivered(&self) -> bool {
|
||||
self.exit_code == Some(0) && self.landed
|
||||
}
|
||||
}
|
||||
|
||||
/// Is the direct-session executor enabled?
|
||||
///
|
||||
/// Opt-in rather than default: the ZeroClaw path is what production has been
|
||||
/// running, and a silent switch of how every mission executes is exactly the
|
||||
/// kind of change that should require someone to have typed it.
|
||||
pub fn direct_mode() -> bool {
|
||||
matches!(
|
||||
std::env::var("CLAWMATES_MISSION_EXECUTOR").as_deref(),
|
||||
Ok("session")
|
||||
)
|
||||
}
|
||||
|
||||
/// Build the instruction for a mission session.
|
||||
///
|
||||
/// One statement of the whole job, not a per-phase directive. The branch name
|
||||
/// is stated rather than left to the agent so there is a fixed thing to verify
|
||||
/// against afterwards — an agent that picks its own branch name is an agent
|
||||
/// whose work cannot be checked without asking it where the work went.
|
||||
pub fn session_prompt(task: &str, repo_path: &str, branch: &str) -> String {
|
||||
format!(
|
||||
"You are working in the git repository at {repo_path}.\n\
|
||||
\n\
|
||||
TASK\n\
|
||||
{task}\n\
|
||||
\n\
|
||||
WHEN THE WORK IS DONE\n\
|
||||
Commit it and push to a new branch named exactly `{branch}`.\n\
|
||||
The remote `origin` is already configured with credentials.\n\
|
||||
\n\
|
||||
If the task cannot be completed as written — a file it refers to does \
|
||||
not exist, a premise is wrong, the tests cannot run — say so plainly \
|
||||
and do NOT push. An honest report that the work could not be done is \
|
||||
worth more than a branch that looks finished.\n"
|
||||
)
|
||||
}
|
||||
|
||||
/// Run one mission session inside an existing container.
|
||||
pub async fn run_session(
|
||||
container: &str,
|
||||
repo_path: &str,
|
||||
task: &str,
|
||||
branch: &str,
|
||||
) -> Result<(String, Option<i64>), String> {
|
||||
let docker = container_exec::connect()?;
|
||||
let prompt = session_prompt(task, repo_path, branch);
|
||||
let mut argv = vec!["claude".to_string(), "-p".to_string()];
|
||||
argv.push("--allowedTools".into());
|
||||
argv.extend(ALLOWED_TOOLS.iter().map(|t| t.to_string()));
|
||||
argv.push("--permission-mode".into());
|
||||
argv.push("acceptEdits".into());
|
||||
argv.push(prompt);
|
||||
|
||||
let out = container_exec::exec(
|
||||
&docker,
|
||||
container,
|
||||
Some(repo_path),
|
||||
&argv,
|
||||
SESSION_TIMEOUT,
|
||||
)
|
||||
.await?;
|
||||
Ok((out.combined(), out.exit_code))
|
||||
}
|
||||
|
||||
/// Ask the forge whether the branch exists, and at what commit.
|
||||
///
|
||||
/// The whole point of the module. Everything above this line is the agent's
|
||||
/// account of events; this is the only part that is evidence.
|
||||
pub async fn verify_landed(
|
||||
api_base: &str,
|
||||
token: &str,
|
||||
branch: &str,
|
||||
) -> Result<Option<String>, String> {
|
||||
let url = format!("{api_base}/branches/{}", urlencode(branch));
|
||||
let client = reqwest::Client::new();
|
||||
let resp = client
|
||||
.get(&url)
|
||||
.header("Authorization", format!("token {token}"))
|
||||
.timeout(Duration::from_secs(30))
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| format!("query branch: {e}"))?;
|
||||
if resp.status().as_u16() == 404 {
|
||||
return Ok(None);
|
||||
}
|
||||
if !resp.status().is_success() {
|
||||
return Err(format!("forge returned {}", resp.status()));
|
||||
}
|
||||
let body: serde_json::Value = resp
|
||||
.json()
|
||||
.await
|
||||
.map_err(|e| format!("decode branch response: {e}"))?;
|
||||
Ok(body
|
||||
.get("commit")
|
||||
.and_then(|c| c.get("id"))
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string))
|
||||
}
|
||||
|
||||
/// Percent-encode the path segment. Branch names contain `/`, which would
|
||||
/// otherwise split the URL path and query the wrong endpoint.
|
||||
fn urlencode(s: &str) -> String {
|
||||
s.bytes()
|
||||
.map(|b| match b {
|
||||
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
|
||||
(b as char).to_string()
|
||||
}
|
||||
_ => format!("%{b:02X}"),
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Branch a session-executed mission pushes to.
|
||||
pub fn session_branch(mission_id: Uuid) -> String {
|
||||
format!("clawmates/session-{}", &mission_id.simple().to_string()[..12])
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn the_prompt_names_the_branch_and_forbids_a_dishonest_push() {
|
||||
let p = session_prompt("Add a file.", "/mission/repo", "clawmates/session-abc");
|
||||
assert!(p.contains("clawmates/session-abc"), "branch must be fixed");
|
||||
assert!(p.contains("/mission/repo"));
|
||||
assert!(
|
||||
p.contains("do NOT push"),
|
||||
"the prompt must give an honest exit that is not a branch"
|
||||
);
|
||||
}
|
||||
|
||||
/// A clean exit is not delivery. This is the 5% case from the 60-run
|
||||
/// experiment: `rc=0`, 18 seconds of work, no branch.
|
||||
#[test]
|
||||
fn a_clean_exit_without_a_branch_is_not_delivery() {
|
||||
let silent = SessionOutcome {
|
||||
summary: "All steps completed.".into(),
|
||||
exit_code: Some(0),
|
||||
landed: false,
|
||||
head_sha: None,
|
||||
};
|
||||
assert!(
|
||||
!silent.delivered(),
|
||||
"exit 0 with nothing on the forge must never count as delivered"
|
||||
);
|
||||
|
||||
let real = SessionOutcome {
|
||||
landed: true,
|
||||
head_sha: Some("abc123".into()),
|
||||
..silent.clone()
|
||||
};
|
||||
assert!(real.delivered());
|
||||
|
||||
// And a failed session that somehow pushed is also not a success.
|
||||
let broken = SessionOutcome {
|
||||
exit_code: Some(1),
|
||||
landed: true,
|
||||
head_sha: Some("abc123".into()),
|
||||
summary: String::new(),
|
||||
};
|
||||
assert!(!broken.delivered());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn branch_names_survive_url_encoding() {
|
||||
assert_eq!(urlencode("clawmates/session-01"), "clawmates%2Fsession-01");
|
||||
assert_eq!(urlencode("plain"), "plain");
|
||||
}
|
||||
|
||||
/// The switch must be explicit. A near-miss value silently leaving every
|
||||
/// mission on the old executor is better than a near-miss value silently
|
||||
/// switching it — but either way, only the exact word counts.
|
||||
#[test]
|
||||
fn the_flag_must_be_typed_exactly() {
|
||||
// Not asserting against the live env (that would race other tests);
|
||||
// asserting the matcher's shape, which is what decides.
|
||||
for wrong in ["Session", "sessions", "direct", "1", "true", ""] {
|
||||
assert_ne!(wrong, "session", "{wrong:?} must not enable direct mode");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_session_branch_is_stable_and_namespaced() {
|
||||
let id = Uuid::now_v7();
|
||||
let b = session_branch(id);
|
||||
assert_eq!(b, session_branch(id));
|
||||
assert!(b.starts_with("clawmates/session-"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,147 @@
|
||||
//! Auto-merge against real git repositories.
|
||||
//!
|
||||
//! The rule is measured from the diff, so it has to be tested against real
|
||||
//! diffs — a unit test on the classifier alone would not catch a wrong
|
||||
//! revision range.
|
||||
|
||||
use cm_api::auto_merge::{self, MergePolicy};
|
||||
|
||||
fn git(repo: &std::path::Path, args: &[&str]) {
|
||||
let out = std::process::Command::new("git")
|
||||
.arg("-C")
|
||||
.arg(repo)
|
||||
.args(args)
|
||||
.env("GIT_AUTHOR_NAME", "T")
|
||||
.env("GIT_AUTHOR_EMAIL", "[email protected]")
|
||||
.env("GIT_COMMITTER_NAME", "T")
|
||||
.env("GIT_COMMITTER_EMAIL", "[email protected]")
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"git {args:?}: {}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
}
|
||||
|
||||
/// Returns (work checkout, bare remote path).
|
||||
fn seed() -> (tempfile::TempDir, std::path::PathBuf, std::path::PathBuf) {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let remote = tmp.path().join("remote.git");
|
||||
let work = tmp.path().join("work");
|
||||
std::process::Command::new("git")
|
||||
.args(["init", "--quiet", "--bare"])
|
||||
.arg(&remote)
|
||||
.output()
|
||||
.unwrap();
|
||||
std::fs::create_dir_all(&work).unwrap();
|
||||
git(&work, &["init", "--quiet"]);
|
||||
git(&work, &["checkout", "-q", "-B", "main"]);
|
||||
std::fs::write(work.join("README.md"), "# vault\n").unwrap();
|
||||
git(&work, &["add", "."]);
|
||||
git(&work, &["commit", "--quiet", "-m", "base"]);
|
||||
git(&work, &["remote", "add", "origin", remote.to_str().unwrap()]);
|
||||
git(&work, &["push", "--quiet", "origin", "main"]);
|
||||
(tmp, work, remote)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_purely_additive_branch_is_merged() {
|
||||
let (_tmp, work, remote) = seed();
|
||||
git(&work, &["checkout", "-q", "-B", "lib/add"]);
|
||||
std::fs::create_dir_all(work.join("60 Papers")).unwrap();
|
||||
std::fs::write(work.join("60 Papers/a.md"), "# paper\n").unwrap();
|
||||
git(&work, &["add", "."]);
|
||||
git(&work, &["commit", "--quiet", "-m", "add paper"]);
|
||||
git(&work, &["push", "--quiet", "origin", "lib/add"]);
|
||||
|
||||
let out = auto_merge::try_merge(
|
||||
&work, remote.to_str().unwrap(), "lib/add", "main",
|
||||
MergePolicy::AdditiveOnly, true,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(out.merged, "should have merged: {}", out.reason);
|
||||
|
||||
// The note must really be on main at the remote, not just locally.
|
||||
let ls = std::process::Command::new("git")
|
||||
.arg("-C").arg(&remote)
|
||||
.args(["ls-tree", "--name-only", "-r", "main"])
|
||||
.output()
|
||||
.unwrap();
|
||||
let listed = String::from_utf8_lossy(&ls.stdout);
|
||||
assert!(listed.contains("60 Papers/a.md"), "remote main: {listed}");
|
||||
}
|
||||
|
||||
/// The load-bearing refusal: a branch that rewrites an existing file must be
|
||||
/// left for a human even though its mission type is allowed to auto-merge.
|
||||
#[tokio::test]
|
||||
async fn a_branch_that_modifies_an_existing_file_is_refused() {
|
||||
let (_tmp, work, remote) = seed();
|
||||
git(&work, &["checkout", "-q", "-B", "lib/bad"]);
|
||||
std::fs::create_dir_all(work.join("60 Papers")).unwrap();
|
||||
std::fs::write(work.join("60 Papers/a.md"), "# paper\n").unwrap();
|
||||
// …and clobbers a hand-written file.
|
||||
std::fs::write(work.join("README.md"), "# REWRITTEN BY A MACHINE\n").unwrap();
|
||||
git(&work, &["add", "."]);
|
||||
git(&work, &["commit", "--quiet", "-m", "add + clobber"]);
|
||||
git(&work, &["push", "--quiet", "origin", "lib/bad"]);
|
||||
|
||||
let out = auto_merge::try_merge(
|
||||
&work, remote.to_str().unwrap(), "lib/bad", "main",
|
||||
MergePolicy::AdditiveOnly, true,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!out.merged, "must refuse a non-additive branch");
|
||||
assert!(out.reason.contains("not additive"), "reason: {}", out.reason);
|
||||
|
||||
let show = std::process::Command::new("git")
|
||||
.arg("-C").arg(&remote)
|
||||
.args(["show", "main:README.md"])
|
||||
.output()
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
String::from_utf8_lossy(&show.stdout),
|
||||
"# vault\n",
|
||||
"the hand-written file must be untouched on main"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn unverified_work_is_never_merged() {
|
||||
let (_tmp, work, remote) = seed();
|
||||
git(&work, &["checkout", "-q", "-B", "lib/unverified"]);
|
||||
std::fs::create_dir_all(work.join("60 Papers")).unwrap();
|
||||
std::fs::write(work.join("60 Papers/a.md"), "# paper\n").unwrap();
|
||||
git(&work, &["add", "."]);
|
||||
git(&work, &["commit", "--quiet", "-m", "add"]);
|
||||
git(&work, &["push", "--quiet", "origin", "lib/unverified"]);
|
||||
|
||||
let out = auto_merge::try_merge(
|
||||
&work, remote.to_str().unwrap(), "lib/unverified", "main",
|
||||
MergePolicy::AdditiveOnly, false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!out.merged);
|
||||
assert!(out.reason.contains("did not verify"), "reason: {}", out.reason);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_never_policy_branch_is_left_alone() {
|
||||
let (_tmp, work, remote) = seed();
|
||||
git(&work, &["checkout", "-q", "-B", "code/change"]);
|
||||
std::fs::write(work.join("new.rs"), "fn main() {}\n").unwrap();
|
||||
git(&work, &["add", "."]);
|
||||
git(&work, &["commit", "--quiet", "-m", "code"]);
|
||||
git(&work, &["push", "--quiet", "origin", "code/change"]);
|
||||
|
||||
let out = auto_merge::try_merge(
|
||||
&work, remote.to_str().unwrap(), "code/change", "main",
|
||||
MergePolicy::Never, true,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!out.merged, "code must never auto-merge");
|
||||
}
|
||||
@@ -19,6 +19,24 @@ async fn workspace(pool: &sqlx::PgPool) -> Uuid {
|
||||
ws
|
||||
}
|
||||
|
||||
/// A real mission row. `corpus_items.mission_id` has a foreign key, which is
|
||||
/// deliberate: attribution to a mission that does not exist is not
|
||||
/// attribution. The first version of the test below used a bare UUID and was
|
||||
/// correctly rejected.
|
||||
async fn mission(pool: &sqlx::PgPool, ws: Uuid) -> Uuid {
|
||||
let id = Uuid::now_v7();
|
||||
sqlx::query(
|
||||
"INSERT INTO missions (id, workspace_id, title, template_kind, schedule, status, config)
|
||||
VALUES ($1,$2,'library','research_only','{}'::jsonb,'running','{}'::jsonb)",
|
||||
)
|
||||
.bind(id)
|
||||
.bind(ws)
|
||||
.execute(pool)
|
||||
.await
|
||||
.unwrap();
|
||||
id
|
||||
}
|
||||
|
||||
fn seed_vault(root: &std::path::Path) {
|
||||
std::fs::create_dir_all(root.join("50 APESS 2026/Lectures")).unwrap();
|
||||
std::fs::create_dir_all(root.join("Repos")).unwrap();
|
||||
@@ -226,3 +244,43 @@ async fn live_arxiv_search_and_fetch() {
|
||||
assert!(pdf.starts_with(b"%PDF"));
|
||||
assert!(pdf.len() > 10_000, "suspiciously small pdf: {}", pdf.len());
|
||||
}
|
||||
|
||||
/// A rerun must not be able to claim credit for work an earlier run did.
|
||||
///
|
||||
/// This is the verification predicate for a continuous mission: "did THIS run
|
||||
/// contribute anything new". If a rerun could re-record an existing source
|
||||
/// under its own mission id, every run would report success forever — the
|
||||
/// failure that killed the 0030-0044 generation of this feature.
|
||||
#[tokio::test]
|
||||
async fn a_rerun_cannot_claim_an_earlier_missions_work() {
|
||||
let pool = cm_testkit::test_pool().await;
|
||||
let ws = workspace(&pool).await;
|
||||
let first_mission = mission(&pool, ws).await;
|
||||
let second_mission = mission(&pool, ws).await;
|
||||
|
||||
corpus::record(
|
||||
&pool, ws, "lib", "source", "arxiv:2401.55555",
|
||||
Some("Paper"), None, None, "h1", Some(first_mission),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// The second mission sees the same paper and re-records it.
|
||||
corpus::record(
|
||||
&pool, ws, "lib", "source", "arxiv:2401.55555",
|
||||
Some("Paper"), None, None, "h2", Some(second_mission),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
corpus::contributed(&pool, ws, "lib", first_mission).await.unwrap(),
|
||||
1,
|
||||
"the finder keeps the credit"
|
||||
);
|
||||
assert_eq!(
|
||||
corpus::contributed(&pool, ws, "lib", second_mission).await.unwrap(),
|
||||
0,
|
||||
"a rerun that found nothing new must report zero, not one"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -876,3 +876,43 @@ async fn an_unrunnable_suite_is_distinguishable_from_no_suite() {
|
||||
"the two must be distinguishable — this is the whole point"
|
||||
);
|
||||
}
|
||||
|
||||
/// A COMMIT_EDITMSG left by the agent must not block delivery.
|
||||
///
|
||||
/// From mission 019fcd0c: the agent ran `git commit` itself, leaving
|
||||
/// `.git/COMMIT_EDITMSG` owned by root at 0644, and the server's commit died
|
||||
/// with "Permission denied". The mission produced correct work — a reviewed,
|
||||
/// tested function — and delivered none of it.
|
||||
///
|
||||
/// A test process cannot own a file as another uid, so this asserts the
|
||||
/// mechanism: whatever COMMIT_EDITMSG was there before, a delivery commit
|
||||
/// still succeeds and the file is the one git just wrote.
|
||||
#[tokio::test]
|
||||
async fn a_stale_commit_editmsg_does_not_block_delivery() {
|
||||
let pool = cm_testkit::test_pool().await;
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let mission = Uuid::now_v7();
|
||||
let repo = seed_repo(tmp.path(), mission);
|
||||
let (_, phase) = seed_mission_phase(&pool, mission).await;
|
||||
|
||||
// Stand in for the agent's leftover: content that must not survive.
|
||||
let msg = repo.join(".git/COMMIT_EDITMSG");
|
||||
std::fs::write(&msg, "LEFTOVER FROM THE AGENT\n").unwrap();
|
||||
|
||||
std::fs::write(repo.join("WORK.md"), "work\n").unwrap();
|
||||
let cap = capture(&pool, tmp.path(), mission, phase)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
|
||||
let commit = cap
|
||||
.committed
|
||||
.expect("delivery must commit despite a stale COMMIT_EDITMSG");
|
||||
assert!(!commit.sha.is_empty());
|
||||
|
||||
let body = std::fs::read_to_string(&msg).unwrap_or_default();
|
||||
assert!(
|
||||
!body.contains("LEFTOVER FROM THE AGENT"),
|
||||
"the stale message survived: {body:?}"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
[Unit]
|
||||
Description=Clawmates paper library harvest (arXiv -> shelf + vault catalogue)
|
||||
Wants=docker.service
|
||||
After=docker.service network-online.target
|
||||
|
||||
[Service]
|
||||
Type=oneshot
|
||||
ExecStart=/usr/local/bin/clawmates-library.sh
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
Nice=10
|
||||
# A harvest downloads PDFs and pushes a branch; give it room but do not
|
||||
# let a wedged run hold the slot until the next week.
|
||||
TimeoutStartSec=30min
|
||||
Executable
+49
@@ -0,0 +1,49 @@
|
||||
#!/usr/bin/env bash
|
||||
# Weekly paper-library harvest.
|
||||
#
|
||||
# Deliberately thin: it calls the API and reports what came back. All the
|
||||
# logic lives in the server, so this file never needs to change when the
|
||||
# harvest does.
|
||||
#
|
||||
# The token lives in /etc/clawmates/library.token (root-only). It is a
|
||||
# long-lived operator session; rotate by replacing the file.
|
||||
set -uo pipefail
|
||||
|
||||
TOKEN_FILE=/etc/clawmates/library.token
|
||||
[ -r "$TOKEN_FILE" ] || { echo "library: no token at $TOKEN_FILE"; exit 1; }
|
||||
TOKEN=$(cat "$TOKEN_FILE")
|
||||
|
||||
RESP=$(docker run --rm --network clawmates_core curlimages/curl:latest \
|
||||
-s -m 1800 -X POST \
|
||||
-H "Authorization: Bearer $TOKEN" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d '{"per_topic":5}' \
|
||||
http://clawmates_server_1:8080/api/library/runs)
|
||||
|
||||
echo "library: $RESP" | head -c 2000
|
||||
|
||||
# Report health explicitly. A run that shelved nothing is normal for a
|
||||
# mature library; a run that ERRORED is not, and the two look identical
|
||||
# if you only count papers.
|
||||
echo "$RESP" | python3 -c '
|
||||
import json, sys
|
||||
try:
|
||||
d = json.load(sys.stdin)
|
||||
except Exception as e:
|
||||
print("library: unreadable response (%s)" % e)
|
||||
sys.exit(1)
|
||||
shelved = len(d.get("shelved", []))
|
||||
healthy = d.get("healthy", False)
|
||||
# Backslashes are avoided inside this program on purpose: it is embedded in a
|
||||
# single-quoted shell string, and an escaped quote here does not survive the
|
||||
# shell. The first version used one inside an f-string, crashed on every run,
|
||||
# and systemd reported a FAILED unit for a harvest that had actually shelved
|
||||
# 15 papers and pushed them. A false failure destroys trust in the signal as
|
||||
# surely as a false success.
|
||||
print("library: %d candidates, %d already held, %d shelved, healthy=%s, pushed=%s, branch=%s" % (
|
||||
d.get("candidates", 0), d.get("already_had", 0), shelved,
|
||||
healthy, d.get("pushed"), d.get("branch")))
|
||||
for f in d.get("failed", []):
|
||||
print("library: FAILED %s" % f)
|
||||
sys.exit(0 if healthy else 1)
|
||||
'
|
||||
@@ -0,0 +1,16 @@
|
||||
[Unit]
|
||||
Description=Clawmates paper library — weekly harvest
|
||||
|
||||
[Timer]
|
||||
# Monday 07:00 local. Weekly rather than daily because arXiv moves at
|
||||
# roughly that pace for a narrow topic set, and a run that almost always
|
||||
# finds nothing trains you to ignore it.
|
||||
OnCalendar=Mon *-*-* 07:00:00
|
||||
# Fire on next boot if the machine was down at the scheduled time — a
|
||||
# missed week is a silently empty library.
|
||||
Persistent=true
|
||||
AccuracySec=1min
|
||||
Unit=clawmates-library.service
|
||||
|
||||
[Install]
|
||||
WantedBy=timers.target
|
||||
Reference in New Issue
Block a user