diff --git a/rust/crates/caos-cli/src/lib.rs b/rust/crates/caos-cli/src/lib.rs index f99a2df04..66f871f56 100644 --- a/rust/crates/caos-cli/src/lib.rs +++ b/rust/crates/caos-cli/src/lib.rs @@ -32,9 +32,9 @@ use conversation_protocol::v3::oid::{ensure_genesis, G3}; use conversation_protocol::v3::paths; pub use conversation_protocol::v3::records::TurnStatus; use conversation_protocol::v3::records::{ - Block, Descriptor, Evidence, Identity, IdentityKind, Proposal, PublicationRecord, - PublicationStatus, Role, SourceTreeResolution, ToolResult as ProtocolToolResult, - TranscriptEntry, TurnOutcome as ProtocolTurnOutcome, TurnRecord, + Block, Descriptor, Identity, IdentityKind, Proposal, PublicationRecord, PublicationStatus, + Role, SourceTreeResolution, ToolResult as ProtocolToolResult, TranscriptEntry, + TurnOutcome as ProtocolTurnOutcome, TurnRecord, }; use conversation_protocol::v3::refs; use conversation_protocol::v3::view::Conversation; @@ -2272,52 +2272,7 @@ fn append_publication_pending( Ok(result.expect("publication append always records a result")) } -struct PublicationOutcome { - status: PublicationStatus, - evidence: Evidence, - observed: Option, -} - -impl PublicationOutcome { - fn new( - status: PublicationStatus, - kind: &str, - diagnostic: Option, - observed: Option, - ) -> Self { - Self { - status, - evidence: Evidence { - kind: kind.to_string(), - diagnostic, - }, - observed, - } - } - - fn from_observation( - pending: &PublicationRecord, - observed: Option, - diagnostic: String, - lease_rejected: bool, - ) -> Self { - let (status, kind) = if observed.as_ref() == Some(&pending.planned_head) { - (PublicationStatus::Complete, "ref-converged") - } else if observed == pending.expected_old { - (PublicationStatus::Uncertain, "ambiguous") - } else { - ( - PublicationStatus::Conflict, - if lease_rejected { - "lease-rejected" - } else { - "ref-drift" - }, - ) - }; - Self::new(status, kind, Some(diagnostic), observed) - } -} +use conversation_protocol::v3::publication::Outcome as PublicationOutcome; fn append_publication_terminal( t: &GitTransport, diff --git a/rust/crates/conversation-protocol/src/v3/mod.rs b/rust/crates/conversation-protocol/src/v3/mod.rs index 477bad2ff..0829e49be 100644 --- a/rust/crates/conversation-protocol/src/v3/mod.rs +++ b/rust/crates/conversation-protocol/src/v3/mod.rs @@ -9,6 +9,7 @@ pub mod ids; pub mod kinds; pub mod oid; pub mod paths; +pub mod publication; pub mod reconcile; pub mod records; pub mod refs; diff --git a/rust/crates/conversation-protocol/src/v3/publication.rs b/rust/crates/conversation-protocol/src/v3/publication.rs new file mode 100644 index 000000000..e26a0992d --- /dev/null +++ b/rust/crates/conversation-protocol/src/v3/publication.rs @@ -0,0 +1,49 @@ +//! Publication outcomes shared by the host and agent publishers. +use super::{Evidence, Oid, PublicationRecord, PublicationStatus}; + +pub struct Outcome { + pub status: PublicationStatus, + pub evidence: Evidence, + pub observed: Option, +} + +impl Outcome { + pub fn new( + status: PublicationStatus, + kind: &str, + diagnostic: Option, + observed: Option, + ) -> Self { + Self { + status, + evidence: Evidence { + kind: kind.to_string(), + diagnostic, + }, + observed, + } + } + + pub fn from_observation( + pending: &PublicationRecord, + observed: Option, + diagnostic: String, + lease_rejected: bool, + ) -> Self { + let (status, kind) = if observed.as_ref() == Some(&pending.planned_head) { + (PublicationStatus::Complete, "ref-converged") + } else if observed == pending.expected_old { + (PublicationStatus::Uncertain, "ambiguous") + } else { + ( + PublicationStatus::Conflict, + if lease_rejected { + "lease-rejected" + } else { + "ref-drift" + }, + ) + }; + Self::new(status, kind, Some(diagnostic), observed) + } +} diff --git a/std/llm-step/src/import_source.rs b/std/llm-step/src/import_source.rs index 77ae93061..7af73be75 100644 --- a/std/llm-step/src/import_source.rs +++ b/std/llm-step/src/import_source.rs @@ -35,6 +35,11 @@ fn is_github_remote(source: &str) -> bool { }) } +pub(super) fn token_file(source: &str) -> Option<&'static str> { + (is_github_remote(source) && Path::new("/secret/github-token").exists()) + .then_some("/secret/github-token") +} + fn free_destination(view: &Conversation<'_>, into: &str) -> Result<(), String> { for name in [into.to_string(), format!("{into}.source.json")] { // entry refuses traversal through a gitlink, file or symlink as well. @@ -58,8 +63,7 @@ pub(super) fn execute(state: &mut progress::State, site: &CallSite<'_>) -> Resul into, } = &args; let revision = revision.as_deref(); - let token_file = (is_github_remote(source) && Path::new("/secret/github-token").exists()) - .then_some("/secret/github-token"); + let token_file = token_file(source); if let Err(error) = free_destination(&state.conversation()?, into) { return site.fail(state, &error); } diff --git a/std/llm-step/src/main.rs b/std/llm-step/src/main.rs index 43bcaf677..86a32b680 100644 --- a/std/llm-step/src/main.rs +++ b/std/llm-step/src/main.rs @@ -4,6 +4,7 @@ mod async_work; mod githist; mod import_source; mod progress; +mod publish_source; mod source_trees; mod subagents; mod timing; @@ -1054,6 +1055,10 @@ fn drive_call( .tool(request, round.declaring_round, &call.id)? { if existing.status == CallStatus::Started { + if existing.name == "publish_source" { + publish_source::execute(state, &site)?; + return Ok(true); + } if existing.name == "import_source" { import_source::execute(state, &site)?; return Ok(true); @@ -1078,6 +1083,10 @@ fn drive_call( harvest_agent_call(state, &site)?; return Ok(true); } + if call.name == "publish_source" { + publish_source::execute(state, &site)?; + return Ok(true); + } if call.name == "import_source" { import_source::execute(state, &site)?; return Ok(true); @@ -2948,6 +2957,9 @@ fn source_tree_paths(state: &mut progress::State) -> Result, String> fn registry(cfg: &Config) -> Result, String> { let mut registry = vec![bash_tool()]; registry.extend(tools::declarations()); + registry.push(with_source_tree(tools::tree_tool_declaration( + &tools::builtin_tool("publish_source", publish_source::HELP), + ))); if cfg.run_and_update_ref_image.is_some() { registry.extend(subagents::declarations()); registry.push(async_work::declaration()); @@ -4308,6 +4320,130 @@ mod tests { } } + #[test] + fn publication_pin_keeps_commit_and_lease_across_source_edits_and_restart() { + use conversation_protocol::v3::{Descriptor, PublicationRecord, PublicationStatus}; + for rejected in [false, true] { + let args = json!({"source_tree":"feature/lower","repository":"https://example.com/repo.git","branch":"topic"}); + let golden = golden_with_first("publish_source", args.clone()).unwrap(); + let call = Call { + id: "first".into(), + name: "publish_source".into(), + input: args, + }; + let site = CallSite::at(&golden.request, 0, &call, ASSISTANT_ID); + let store = ImportStore { + objects: golden.store, + head: golden.head.clone(), + race: None, + lost_ack: true, + }; + let mut state = progress::State::from_store( + store, + "refs/conversations/conversation/head".into(), + golden.head, + ) + .unwrap(); + let commit = Oid::parse(&"a".repeat(40), "source").unwrap(); + let old = Oid::parse(&"b".repeat(40), "old").unwrap(); + let descriptor = Descriptor { + source_base: old.clone(), + source_head: commit.clone(), + target_base: old.clone(), + policy: "preserve".into(), + implementation: "caos/server-push".into(), + commit_policy: "preserve".into(), + }; + let conversation = state.conversation().unwrap().identity().unwrap().id; + let key = publish_source::invocation(&conversation, &site).unwrap()[..32].to_string(); + let id = ids::publication_id( + &conversation, + &key, + &ids::projection_id(&descriptor.to_value()).unwrap(), + &commit, + "https://example.com/repo.git", + "refs/heads/topic", + Some(&old), + ) + .unwrap(); + let pending = PublicationRecord { + id, + key, + descriptor, + planned_head: commit.clone(), + expected_old: Some(old.clone()), + repository: "https://example.com/repo.git".into(), + refname: "refs/heads/topic".into(), + source_tree_name: "feature/lower".into(), + status: PublicationStatus::Pending, + evidence: None, + observed: None, + }; + let saved = publish_source::pin(&mut state, &site, &pending) + .unwrap() + .unwrap(); + assert_eq!(saved, pending); + let newer = Oid::parse(&"c".repeat(40), "newer source").unwrap(); + state + .append(Transition::FilesApply { + files: vec![( + "feature/lower".into(), + Some((Mode::Commit, newer.encode_line())), + )], + }) + .unwrap(); + state.reload().unwrap(); + let mut competing = pending.clone(); + competing.planned_head = newer.clone(); + competing.expected_old = None; + let recovered = publish_source::pin(&mut state, &site, &competing) + .unwrap() + .unwrap(); + assert_eq!(recovered, pending); + let outcome = if rejected { + conversation_protocol::v3::publication::Outcome::new( + PublicationStatus::Conflict, + "validation-rejected", + Some("Invalid source".into()), + None, + ) + } else { + // Git may resend after a lost acknowledgement and report a + // rejected old value even though its first push succeeded. + publish_source::reconcile( + &recovered, + conversation_protocol::v3::publication::Outcome::new( + PublicationStatus::Conflict, + "lease-rejected", + None, + None, + ), + Ok(Some(commit.clone())), + ) + }; + publish_source::finish(&mut state, &site, &pending, Some(outcome)).unwrap(); + let view = state.conversation().unwrap(); + assert_eq!( + view.source_tree("feature/lower").unwrap().unwrap().commit, + newer + ); + let receipt = view.publication(&pending.id).unwrap().unwrap(); + assert_eq!( + receipt.status, + if rejected { + PublicationStatus::Conflict + } else { + PublicationStatus::Complete + } + ); + assert_eq!(receipt.planned_head, commit); + assert_eq!(receipt.expected_old, Some(old)); + assert!(publish_source::pin(&mut state, &site, &competing) + .unwrap() + .is_none()); + } + } + #[test] fn import_attachment_is_atomic_replayable_and_does_not_overwrite() { for race_path in [ diff --git a/std/llm-step/src/publish_source.rs b/std/llm-step/src/publish_source.rs new file mode 100644 index 000000000..7766b72a5 --- /dev/null +++ b/std/llm-step/src/publish_source.rs @@ -0,0 +1,336 @@ +//! Pin publication intent in the conversation before any remote mutation. +use super::*; +use conversation_protocol::v3::publication::Outcome; +use conversation_protocol::v3::{Descriptor, PublicationRecord, PublicationStatus}; + +pub(super) const HELP: &str = "Publish the exact selected source commit to an HTTPS Git repository branch, preserving its history. Test and inspect the intended PR diff first. Resolve merge conflicts and clear .caos/conflicts before publishing. The endpoint pushes the commit unchanged; it does not filter files or apply .gitignore. This does not create a PR or change the source gitlink. Only fast-forward updates are supported: import and merge remote changes before retrying a conflict. A receipt names the exact published commit even if the source later changes. On uncertainty, inspect the remote before taking another action. +@param repository HTTPS Git repository URL, without credentials. +@param branch Destination branch name (without refs/heads/)."; + +#[derive(serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct Parameters { + source_tree: String, + repository: String, + branch: String, +} + +fn parameters(call: &Call) -> Result { + let p: Parameters = serde_json::from_value(call.input.clone()) + .map_err(|e| format!("invalid publish_source arguments: {e}"))?; + paths::validate_source_tree_name(&p.source_tree)?; + git_locator::import::remote(&p.repository)?; + conversation_protocol::v3::source_trees::validate_branch(&p.branch)?; + if p.branch.starts_with("refs/") { + return Err("branch must omit refs/heads/".into()); + } + Ok(p) +} + +fn pin_path(site: &CallSite<'_>) -> String { + format!( + "{}/publication.json", + paths::call_payload_dir(site.request.as_str(), site.round, &site.call.id) + ) +} + +fn pinned( + view: &Conversation<'_>, + site: &CallSite<'_>, +) -> Result, String> { + if view + .tool(site.request, site.round, &site.call.id)? + .is_none() + { + return Ok(None); + } + let id: String = serde_json::from_slice(&view.payload(&pin_path(site))?) + .map_err(|_| "invalid publication pin")?; + view.publication(&id)? + .ok_or("pinned publication is missing".into()) + .map(Some) +} + +pub(super) fn execute(state: &mut progress::State, site: &CallSite<'_>) -> Result<(), String> { + let p = match parameters(site.call) { + Ok(p) => p, + Err(error) => return site.fail(state, &error), + }; + let token_file = import_source::token_file(&p.repository); + let token = token_file + .map(fs::read_to_string) + .transpose() + .map_err(|_| "reading GitHub token")?; + let token = token.as_deref().map(str::trim_end); + let observe = || { + git_locator::publish::read_branch(&p.repository, &p.branch, token) + .and_then(|h| h.map(|h| Oid::parse(&h, "remote head")).transpose()) + }; + let pending = match pinned(&state.conversation()?, site)? { + Some(record) => record, + None => { + let view = state.conversation()?; + let head = match view.source_tree(&p.source_tree)? { + Some(source) => source.commit, + None => { + return site.fail(state, "publish_source requires an existing source gitlink") + } + }; + let base = view.reference_start(&p.source_tree)?; + let old = match observe() { + Ok(old) => old, + Err(error) => return site.fail(state, &error), + }; + let id = view.identity()?.id; + let descriptor = Descriptor { + source_base: base.clone(), + source_head: head.clone(), + target_base: base, + policy: "preserve".into(), + implementation: "caos/server-push".into(), + commit_policy: "preserve".into(), + }; + let key = invocation(&id, site)?[..32].to_string(); + let publication = ids::publication_id( + &id, + &key, + &ids::projection_id(&descriptor.to_value())?, + &head, + &p.repository, + &format!("refs/heads/{}", p.branch), + old.as_ref(), + )?; + let record = PublicationRecord { + id: publication, + key, + descriptor, + planned_head: head, + repository: p.repository.clone(), + refname: format!("refs/heads/{}", p.branch), + expected_old: old, + source_tree_name: p.source_tree.clone(), + status: PublicationStatus::Pending, + evidence: None, + observed: None, + }; + let Some(record) = pin(state, site, &record)? else { + return Ok(()); + }; + record + } + }; + if pending.status != PublicationStatus::Pending { + return finish(state, site, &pending, None); + } + // An attempt that pinned this intent may send it, including concurrent + // attempts that joined the identical transition. The exact lease makes + // those pushes converge. Recovery never obtains a fresh source or lease. + let outcome = { + let mut command = std::process::Command::new("caos"); + command + .args([ + "push-git", + &pending.repository, + pending.planned_head.as_str(), + &p.branch, + ]) + .arg(format!( + "--expected={}", + pending + .expected_old + .as_ref() + .map(Oid::as_str) + .unwrap_or("absent") + )); + if let Some(file) = token_file { + command.arg(format!("--github-token-file={file}")); + } + let child = command + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .spawn(); + match child { + Err(_) => Outcome::new( + PublicationStatus::Conflict, + "validation-rejected", + Some("Could not start caos push-git; no push was attempted.".into()), + None, + ), + Ok(child) => match child.wait_with_output() { + Ok(output) if output.status.success() => { + let value: Value = serde_json::from_slice(&output.stdout) + .map_err(|_| "invalid push-git result")?; + let status: PublicationStatus = serde_json::from_value(value["status"].clone()) + .map_err(|_| "invalid push status")?; + let observed: Option = serde_json::from_value(value["observed"].clone()) + .map_err(|_| "invalid remote head")?; + let outcome = Outcome::new( + status, + value["kind"] + .as_str() + .ok_or("missing publication evidence")?, + value["diagnostic"].as_str().map(str::to_owned), + observed, + ); + if status == PublicationStatus::Complete { + outcome + } else { + reconcile(&pending, outcome, observe()) + } + } + Ok(output) if output.status.code() == Some(1) => Outcome::new( + PublicationStatus::Conflict, + "validation-rejected", + Some(String::from_utf8_lossy(&output.stderr).trim().to_string()), + None, + ), + _ => recovered(&pending, observe()), + }, + } + }; + finish(state, site, &pending, Some(outcome)) +} + +pub(super) fn invocation(conversation: &str, site: &CallSite<'_>) -> Result { + ids::protocol_id( + "external-tool", + &json!({ + "conversation":conversation, "request":site.request, "round":site.round, "call":site.call.id + }), + ) +} + +pub(super) fn reconcile( + pending: &PublicationRecord, + outcome: Outcome, + observed: Result, String>, +) -> Outcome { + match observed { + Ok(head) if head.as_ref() == Some(&pending.planned_head) => { + Outcome::new(PublicationStatus::Complete, "ref-converged", None, head) + } + observed if outcome.status == PublicationStatus::Uncertain => recovered(pending, observed), + _ => outcome, + } +} + +fn recovered(pending: &PublicationRecord, observed: Result, String>) -> Outcome { + match observed { + Ok(head) => Outcome::from_observation( + pending, + head, + "No confirmed push result. An unchanged branch may still have an in-flight update." + .into(), + false, + ), + Err(_) => Outcome::new( + PublicationStatus::Uncertain, + "ambiguous", + Some( + "Could not observe the pinned publication; inspect the remote before continuing." + .into(), + ), + None, + ), + } +} + +pub(super) fn pin( + state: &mut progress::State, + site: &CallSite<'_>, + pending: &PublicationRecord, +) -> Result, String> { + for _ in 0..32 { + state.reload()?; + let view = state.conversation()?; + if view + .tool(site.request, site.round, &site.call.id)? + .is_some_and(|r| r.is_terminal()) + { + return Ok(None); + } + if let Some(record) = pinned(&view, site)? { + return Ok(Some(record)); + } + let mut tool = site.stub(None); + tool.status = CallStatus::Started; + let expected = state.head().clone(); + if matches!( + state.try_append_pair_at( + &expected, + Transition::PublicationPending { + record: pending.clone() + }, + Transition::ToolStart { + record: tool, + payloads: vec![( + "publication.json".into(), + canonical_payload_bytes(&json!(pending.id))? + )] + } + )?, + progress::TryAppend::Appended(_) + ) { + return Ok(Some(pending.clone())); + } + } + Err("conversation kept moving while pinning publication".into()) +} + +pub(super) fn finish( + state: &mut progress::State, + site: &CallSite<'_>, + pending: &PublicationRecord, + outcome: Option, +) -> Result<(), String> { + for _ in 0..32 { + state.reload()?; + let view = state.conversation()?; + if view + .tool(site.request, site.round, &site.call.id)? + .is_some_and(|r| r.is_terminal()) + { + return Ok(()); + } + let mut record = view + .publication(&pending.id)? + .ok_or("publication disappeared")?; + let mut transitions = Vec::new(); + if record.status == PublicationStatus::Pending { + let out = outcome.as_ref().ok_or("missing publication outcome")?; + record.status = out.status; + record.evidence = Some(out.evidence.clone()); + record.observed = out.observed.clone(); + transitions.push(Transition::PublicationTerminal { + publication: record.id.clone(), + status: out.status, + evidence: out.evidence.clone(), + observed: out.observed.clone(), + }); + } + let text = serde_json::to_string(&record).map_err(|e| e.to_string())?; + let block = result_block( + &site.call.id, + &text, + record.status != PublicationStatus::Complete, + ); + let stub = site.stub(None); + let tool = completed_record( + &stub, + ToolResult::Complete { + observation: observation_path(&stub), + proposal: None, + }, + None, + ); + transitions.push(tool_complete_transition(tool, &block, Vec::new())?); + let expected = state.head().clone(); + if matches!( + state.try_append_many_at(&expected, transitions)?, + progress::TryAppend::Appended(_) + ) { + return Ok(()); + } + } + Err("conversation kept moving while completing publication".into()) +} diff --git a/std/llm-step/src/source_trees.rs b/std/llm-step/src/source_trees.rs index 4a5730548..acae55c89 100644 --- a/std/llm-step/src/source_trees.rs +++ b/std/llm-step/src/source_trees.rs @@ -13,7 +13,7 @@ Conversation filesystem: - Before integrating a publication base, compare the intended feature change with the full source-versus-destination difference. Merging upstream preserves all existing branch changes; rebasing a whole branch may replay inherited changes too. Transplanting only the requested edit onto a new base is a separate operation. If a small task would publish unrelated inherited work, explain that scope and ask whether to retain the branch or transplant the edit. Do not claim that targeting an older upstream isolates the edit. Preserve existing snapshots when making a new starting point. - Resolve source-tree conflicts by editing affected files and clearing their ledger entries. Saving an edited source tree removes an empty .caos/conflicts ledger and prunes its .caos directory if empty. Unresolved entries and other metadata stay intact. Do not recreate a removed ledger to register resolution. This source-tree ledger is distinct from protected conversation-root .caos. - Preparing, delegating, merging, testing, and organizing a PR stack requires no publishing destination. Preserve the imported starting commit at 00-base. Finish the requested work without asking for repository publication settings. -- Publication is an explicit client operation. The user runs /pr with an explicit source path and base branch; the repository URL is supplied or inferred from unambiguous import provenance, then confirmed in a preview. Do not create publication policy files. Incorporate the chosen destination base and each preceding boundary, resolve conflicts, and run relevant checks before publication. The TUI pushes the exact previewed commits; it does not edit or test code. +- When branch publication is requested, use publish_source with an explicit gitlink, HTTPS repository and branch. Inspect the complete PR diff, incorporate the chosen base and prior stack boundary, resolve conflicts and run checks first. No publication policy file or local branch database is needed. Import provenance can identify the repository and default base; ask only when ambiguous. - If a tool fails, report the observed error and uncertainty; do not invent storage behavior or claim success from a failed check. Use available repository tools for relevant tests; describe a specific missing capability rather than assuming workers cannot run tests or reach a server. - Git tools require the target commit-entry path explicitly. Invoke repository tools with run_tool at a conversation-relative path. UI selection never changes your execution context. diff --git a/std/llm-step/src/tools.rs b/std/llm-step/src/tools.rs index 3f4564e74..7ed757d33 100644 --- a/std/llm-step/src/tools.rs +++ b/std/llm-step/src/tools.rs @@ -28,7 +28,10 @@ const MAX_ENTRIES: usize = 1_000; /// True if `name` is one of the inline tools this module executes. pub fn is_inline(name: &str) -> bool { - matches!(name, "read" | "ls" | "write" | "edit" | "import_source") + matches!( + name, + "read" | "ls" | "write" | "edit" | "import_source" | "publish_source" + ) } /// Help text for the built-in tools, authored exactly like a caos-tools @@ -109,6 +112,7 @@ pub fn grep_declaration() -> Value { /// are standard, not project-defined. const RESERVED_TOOLS: &[&str] = &[ "import_source", + "publish_source", "bash", "grep", "read",