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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 4 additions & 49 deletions rust/crates/caos-cli/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Oid>,
}

impl PublicationOutcome {
fn new(
status: PublicationStatus,
kind: &str,
diagnostic: Option<String>,
observed: Option<Oid>,
) -> Self {
Self {
status,
evidence: Evidence {
kind: kind.to_string(),
diagnostic,
},
observed,
}
}

fn from_observation(
pending: &PublicationRecord,
observed: Option<Oid>,
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,
Expand Down
1 change: 1 addition & 0 deletions rust/crates/conversation-protocol/src/v3/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
49 changes: 49 additions & 0 deletions rust/crates/conversation-protocol/src/v3/publication.rs
Original file line number Diff line number Diff line change
@@ -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<Oid>,
}

impl Outcome {
pub fn new(
status: PublicationStatus,
kind: &str,
diagnostic: Option<String>,
observed: Option<Oid>,
) -> Self {
Self {
status,
evidence: Evidence {
kind: kind.to_string(),
diagnostic,
},
observed,
}
}

pub fn from_observation(
pending: &PublicationRecord,
observed: Option<Oid>,
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)
}
}
8 changes: 6 additions & 2 deletions std/llm-step/src/import_source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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);
}
Expand Down
136 changes: 136 additions & 0 deletions std/llm-step/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ mod async_work;
mod githist;
mod import_source;
mod progress;
mod publish_source;
mod source_trees;
mod subagents;
mod timing;
Expand Down Expand Up @@ -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);
Expand All @@ -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);
Expand Down Expand Up @@ -2948,6 +2957,9 @@ fn source_tree_paths(state: &mut progress::State) -> Result<Vec<String>, String>
fn registry(cfg: &Config) -> Result<Vec<Value>, 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());
Expand Down Expand Up @@ -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 [
Expand Down
Loading
Loading