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
85 changes: 85 additions & 0 deletions crates/agent-controller/src/community.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
//! Owner-signed intent to configure one identity/community pair.
//! Local import and existing setups do not depend on this confirmation.
use crate::config::{canonical_key, canonical_relay};
use crate::Result;
use secp256k1::{schnorr::Signature, Secp256k1, XOnlyPublicKey};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};

#[derive(Clone, Deserialize, Serialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct CommunityResolution {
pub pubkey: String,
pub relay_url: String,
pub owner: String,
pub signature: String,
}
impl CommunityResolution {
pub fn verify(&self, auth: &str) -> Result<()> {
crate::secret::validate_attestation(auth, &self.pubkey)?;
let tag: Vec<String> = serde_json::from_str(auth).map_err(|_| "Invalid source owner")?;
if !canonical_key(&self.pubkey)
|| tag[1] != self.owner
|| canonical_relay(&self.relay_url)? != self.relay_url
{
return Err(
"Community resolution does not match the source owner or destination".into(),
);
}
let owner: XOnlyPublicKey = self.owner.parse().map_err(|_| "Invalid resolution owner")?;
let signature: Signature = self
.signature
.parse()
.map_err(|_| "Invalid resolution signature")?;
let digest = Sha256::digest(format!(
"nostr:agent-community:{}:{}",
self.pubkey, self.relay_url
));
Secp256k1::verification_only()
.verify_schnorr(&signature, &digest, &owner)
.map_err(|_| "Community resolution was not signed by the source owner".into())
}
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
use secp256k1::{Keypair, SecretKey};
const PUB: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798";
pub(crate) fn resolution(relay: &str) -> CommunityResolution {
let secp = Secp256k1::new();
let mut bytes = [0; 32];
bytes[31] = 2;
let pair = Keypair::from_secret_key(&secp, &SecretKey::from_byte_array(bytes).unwrap());
CommunityResolution {
pubkey: PUB.into(),
relay_url: relay.into(),
owner: pair.x_only_public_key().0.to_string(),
signature: secp
.sign_schnorr_no_aux_rand(
&Sha256::digest(format!("nostr:agent-community:{PUB}:{relay}")),
&pair,
)
.to_string(),
}
}
#[test]
fn owner_can_confirm_multiple_pairs_without_import_or_storage() {
let auth = crate::secret::test_attestation(PUB);
resolution("wss://one.example").verify(&auth).unwrap();
resolution("wss://two.example").verify(&auth).unwrap();
}
#[test]
fn unsigned_or_retargeted_resolution_is_rejected() {
let auth = crate::secret::test_attestation(PUB);
for change in ["destination", "identity", "owner", "signature"] {
let mut value = resolution("wss://one.example");
match change {
"destination" => value.relay_url = "wss://two.example".into(),
"identity" => value.pubkey = "ab".repeat(32),
"owner" => value.owner = "cd".repeat(32),
_ => value.signature = "00".repeat(64),
}
assert!(value.verify(&auth).is_err());
}
}
}
5 changes: 5 additions & 0 deletions crates/agent-controller/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ pub struct AgentView {
pub restart_diff: Vec<crate::restart::RestartDiffEntry>,
pub deployed_remote: bool,
pub needs_team_import: bool,
pub configured: bool,
}
#[derive(Clone, Serialize)]
#[serde(rename_all = "camelCase")]
Expand Down Expand Up @@ -126,6 +127,9 @@ pub(crate) struct Agent {
pub extra: BTreeMap<String, Value>,
}
impl Agent {
pub(crate) fn configured(&self) -> bool {
self.extra.get("configured") != Some(&Value::Bool(false))
}
/// `inherited` supplies launch selectors; the saved harness is shown as saved.
pub(crate) fn view(&self, inherited: &crate::agent_defaults::AgentDefaults) -> AgentView {
let defaults = crate::build_defaults();
Expand Down Expand Up @@ -153,6 +157,7 @@ impl Agent {
status: ProcessStatus::Stopped,
error: None,
diagnostics: Vec::new(),
configured: self.configured(),
profile_pending: self.extra.get("profilePending") == Some(&Value::Bool(true)),
start_on_app_launch: self.starts_on_launch(),
respond_to: self.respond_to(defaults.owner_only).ok().map(str::to_owned),
Expand Down
30 changes: 29 additions & 1 deletion crates/agent-controller/src/import.rs
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,7 @@ impl Imports {
if data.digest != pending.digest {
return Err("Source changed after preview; preview it again".into());
}
let reservation = store.reserve_import()?;
let existing = store.agents()?;
let mut agents = Vec::new();
let mut repairs = Vec::new();
Expand Down Expand Up @@ -184,11 +185,15 @@ impl Imports {
));
continue;
}
if existing.iter().any(|a| a.pubkey == candidate.pubkey) {
return Err("Selected identity is already imported".into());
}
let agent = resolve(&data, record, &pending.workspace, &candidate.relay_url)?;
agent.validate()?;
agents.push((agent, string(record, "private_key_nsec").to_owned()));
}
Ok(PreparedImport {
reservation,
agents,
repairs,
source_kind: pending.source_kind,
Expand All @@ -212,13 +217,15 @@ impl Imports {
/// Native-only import plan; never serialized. Credential operations can happen
/// outside the controller mutex. The source snapshot is copied, never mutated.
pub struct PreparedImport {
reservation: crate::store::ImportReservation,
agents: Vec<(Agent, String)>,
repairs: Vec<(String, u64, String)>,
source_kind: LegacySource,
source: PathBuf,
digest: String,
}
pub struct CredentialedImport {
reservation: crate::store::ImportReservation,
agents: Vec<Agent>,
repairs: Vec<(String, u64, String)>,
source: PathBuf,
Expand Down Expand Up @@ -250,6 +257,7 @@ impl PreparedImport {
}
}
Ok(CredentialedImport {
reservation: self.reservation,
agents: agents.into_iter().map(|(a, _)| a).collect(),
repairs: self.repairs,
source: self.source,
Expand All @@ -259,10 +267,30 @@ impl PreparedImport {
}
impl CredentialedImport {
pub fn commit(self, store: &mut Store) -> Result<()> {
if !self.reservation.belongs_to(store) {
return Err("Import belongs to another agent store".into());
}
if read_source(&self.source)?.digest != self.digest {
return Err("Source changed during credential access; preview again".into());
}
store.import(self.agents, self.repairs)
let existing = store.agents()?;
if self
.agents
.iter()
.any(|incoming| existing.iter().any(|saved| saved.pubkey == incoming.pubkey))
{
return Err("Selected identity is already imported".into());
}
let agents = self
.agents
.into_iter()
.map(|mut agent| {
agent.extra.insert("configured".into(), Value::Bool(true));
agent.enabled = false;
agent
})
.collect();
store.import(agents, self.repairs)
}
}
fn string<'a>(value: &'a Value, key: &str) -> &'a str {
Expand Down
85 changes: 85 additions & 0 deletions crates/agent-controller/src/import/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,8 @@ fn preview_is_keyless_commit_resolves_preserves_and_never_enables_or_mutates_sou
assert!(!saved.enabled);
// Source `start_on_app_launch: true` does not auto-start an imported record.
assert_eq!(saved.start_on_app_launch, Some(false));
assert!(saved.configured());
assert_eq!(saved.relay_url, "wss://relay.example");
assert_eq!(saved.system_prompt, "definition-prompt");
assert_eq!(saved.harness.model, "definition-model");
assert_eq!(saved.harness.provider, "global-provider");
Expand Down Expand Up @@ -826,3 +828,86 @@ fn team_snapshot_handles_empty_deleted_and_invalid_teams() {
assert!(team_instructions(&data, &data.records[1]).is_err());
}
}
#[test]
fn import_excludes_overlapping_destinations_before_credentials_through_commit() {
let old = tempfile::tempdir().unwrap();
let dest = tempfile::tempdir().unwrap();
source(old.path());
let keys = Memory::default();
let mut store = Store::open(dest.path().into()).unwrap();
let mut first = Imports::default();
let preview = first
.preview(
LegacySource::Installed,
old.path().into(),
dest.path().into(),
"wss://first.example",
)
.unwrap();
let mut second = Imports::default();
let other = second
.preview(
LegacySource::Installed,
old.path().into(),
dest.path().into(),
"wss://second.example",
)
.unwrap();
let ids = [preview.candidates[0].id.clone()];
let other_ids = [other.candidates[0].id.clone()];
let prepared = first.prepare(&preview.token, &ids, &store).unwrap();
// Separate preview owners still share the same store reservation.
assert!(second
.prepare(&other.token, &other_ids, &store)
.err()
.unwrap()
.contains("import is in progress"));
assert!(keys.keys.lock().unwrap().is_empty());
// Cancellation before credential access releases the reservation.
drop(prepared);
let prepared = first.prepare(&preview.token, &ids, &store).unwrap();
let unavailable = Memory {
fail: true,
..Default::default()
};
assert!(prepared.acquire(&unavailable).is_err());
// Credential failure releases it too, so the same preview can be retried.
let pending = first
.prepare(&preview.token, &ids, &store)
.unwrap()
.acquire(&keys)
.unwrap();
let reads = keys.reads.load(Ordering::SeqCst);
assert!(second
.commit(&other.token, &other_ids, &mut store, &keys)
.unwrap_err()
.contains("import is in progress"));
assert_eq!(keys.reads.load(Ordering::SeqCst), reads);
assert_eq!(
keys.keys
.lock()
.unwrap()
.keys()
.cloned()
.collect::<Vec<_>>(),
ids
);
pending.commit(&mut store).unwrap();
let before = fs::read(dest.path().join("agents.json")).unwrap();
assert!(second
.commit(&other.token, &other_ids, &mut store, &keys)
.unwrap_err()
.contains("already imported"));
assert_eq!(keys.reads.load(Ordering::SeqCst), reads);
assert_eq!(
keys.keys
.lock()
.unwrap()
.keys()
.cloned()
.collect::<Vec<_>>(),
ids
);
assert_eq!(fs::read(dest.path().join("agents.json")).unwrap(), before);
assert_eq!(store.agents().unwrap().len(), 1);
}
2 changes: 2 additions & 0 deletions crates/agent-controller/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@
//! No legacy desktop dependency, implicit identity creation, or credential projection.
mod agent_defaults;
mod bundle;
mod community;
mod config;
pub use community::CommunityResolution;
pub mod connection;
mod create;
mod credentials;
Expand Down
38 changes: 30 additions & 8 deletions crates/agent-controller/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -486,6 +486,14 @@ impl Controller {
|| !agent.imported.is_null(),
)
}
pub fn use_here(
&mut self,
id: &str,
resolution: crate::CommunityResolution,
) -> Result<ControlSnapshot> {
self.store.use_here(id, resolution)?;
self.snapshot()
}
pub fn prepare_import(
&self,
imports: &mut crate::Imports,
Expand Down Expand Up @@ -584,12 +592,18 @@ impl Controller {
.collect())
}
pub fn delete(&mut self, id: &str, revision: u64) -> Result<ControlSnapshot> {
let agent = self
.store
.agents()?
.into_iter()
let agents = self.store.agents()?;
let agent = agents
.iter()
.find(|agent| agent.id == id)
.cloned()
.ok_or("Agent no longer exists")?;
// Use here setups of one identity share its key; keep it for the others.
let shared = agents.iter().any(|other| {
other.id != agent.id
&& other.credential_id == agent.credential_id
&& other.pubkey == agent.pubkey
});
if agent.revision != revision {
return Err("Agent settings changed. Reload before deleting".into());
}
Expand All @@ -601,8 +615,10 @@ impl Controller {
self.store.enabled(id, false)?;
// A failed settings write leaves the card available for an explicit retry.
// Credential deletion is idempotent, so that retry can finish cleanup.
self.credentials
.delete(&agent.credential_id, &agent.pubkey)?;
if !shared {
self.credentials
.delete(&agent.credential_id, &agent.pubkey)?;
}
self.store.remove(id, revision)?;
self.errors.remove(id);
self.snapshot()
Expand Down Expand Up @@ -652,7 +668,7 @@ impl Controller {
.store
.agents()?
.into_iter()
.filter(Agent::starts_on_launch)
.filter(|a| a.starts_on_launch() && a.configured())
{
// Like the host restore, a launch preference is a Start: it enables.
if let Err(error) = self
Expand All @@ -673,6 +689,9 @@ impl Controller {
.find(|a| a.id == id)
.ok_or("Agent no longer exists")?;
let agent = crate::agent_defaults::effective(&agent, &self.store.defaults()?);
if !agent.configured() {
return Err("Choose Use here before opening this identity’s credentials".into());
}
let workspace = effective_databricks(&agent)?.map(|s| s.host);
self.bundle.as_ref().map_err(Clone::clone)?;
Ok((agent.credential_id, agent.pubkey, agent.revision, workspace))
Expand Down Expand Up @@ -714,7 +733,7 @@ impl Controller {
.store
.agents()?
.into_iter()
.filter(Agent::starts_on_launch)
.filter(|a| a.starts_on_launch() && a.configured())
.map(|a| a.id)
.collect())
}
Expand All @@ -739,6 +758,9 @@ impl Controller {
.into_iter()
.find(|a| a.id == id)
.ok_or("Agent no longer exists")?;
if !agent.configured() {
return Err("Choose Use here before starting this imported identity".into());
}
if !agent.enabled {
return Err("Agent is disabled".into());
}
Expand Down
Loading
Loading