mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix: reconcile agent profile on startup when relay publish was missed (#921)
Signed-off-by: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 <dcfd242e557282d7a1e2cf2e6877522682f1e5c6156dc92ca7d90eaedd3b0f95@sprout-oss.stage.blox.sqprod.co> Co-authored-by: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 <dcfd242e557282d7a1e2cf2e6877522682f1e5c6156dc92ca7d90eaedd3b0f95@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
co-authored by
npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
parent
5d927aba0c
commit
762a45969a
@@ -30,6 +30,7 @@ const rules = [
|
||||
// Do not add to this list; split the file instead. Remove each entry as its
|
||||
// file is broken up. Tracked as a follow-up.
|
||||
const overrides = new Map([
|
||||
["src-tauri/src/commands/agents.rs", 1190],
|
||||
["src-tauri/src/managed_agents/nest.rs", 1417],
|
||||
["src-tauri/src/managed_agents/runtime.rs", 1387],
|
||||
["src-tauri/src/huddle/tts.rs", 1364],
|
||||
|
||||
@@ -239,7 +239,10 @@ pub async fn update_managed_agent(
|
||||
.map_err(|e| format!("failed to parse agent keys: {e}"))?;
|
||||
let relay_url = record.relay_url.clone();
|
||||
let display_name = record.name.clone();
|
||||
let avatar_url = managed_agent_avatar_url(&record.agent_command);
|
||||
let avatar_url = record
|
||||
.avatar_url
|
||||
.clone()
|
||||
.or_else(|| managed_agent_avatar_url(&record.agent_command));
|
||||
let auth_tag = record.auth_tag.clone();
|
||||
Some((agent_keys, relay_url, display_name, avatar_url, auth_tag))
|
||||
} else {
|
||||
|
||||
@@ -377,7 +377,7 @@ pub async fn create_managed_agent(
|
||||
};
|
||||
|
||||
// ── Phase 3: save record (sync lock) ───────────────────────────────────────
|
||||
let agent = {
|
||||
let (agent, resolved_avatar_url) = {
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
@@ -453,6 +453,18 @@ pub async fn create_managed_agent(
|
||||
Some((pack_path, slug.to_owned()))
|
||||
});
|
||||
|
||||
// Resolve the avatar URL once at creation and persist it on the record.
|
||||
// This is the same logic the original publish used (user input, else
|
||||
// command-based fallback) — storing it lets reconciliation compare
|
||||
// against what was actually published instead of re-deriving it.
|
||||
let resolved_avatar_url = input
|
||||
.avatar_url
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.map(str::to_string)
|
||||
.or_else(|| managed_agent_avatar_url(&agent_command));
|
||||
|
||||
let record = crate::managed_agents::ManagedAgentRecord {
|
||||
pubkey: pubkey.clone(),
|
||||
name: name.clone(),
|
||||
@@ -460,6 +472,7 @@ pub async fn create_managed_agent(
|
||||
private_key_nsec: private_key_nsec.clone(),
|
||||
auth_tag: auth_tag.clone(),
|
||||
relay_url: resolved_relay_url.clone(),
|
||||
avatar_url: resolved_avatar_url.clone(),
|
||||
acp_command: input
|
||||
.acp_command
|
||||
.as_deref()
|
||||
@@ -534,7 +547,10 @@ pub async fn create_managed_agent(
|
||||
.iter()
|
||||
.find(|record| record.pubkey == pubkey)
|
||||
.ok_or_else(|| "created agent disappeared unexpectedly".to_string())?;
|
||||
build_managed_agent_summary(&app, record, &runtimes)?
|
||||
(
|
||||
build_managed_agent_summary(&app, record, &runtimes)?,
|
||||
resolved_avatar_url,
|
||||
)
|
||||
};
|
||||
|
||||
// ── Phase 3b: local spawn (async preflight outside store lock) ───────────
|
||||
@@ -571,19 +587,14 @@ pub async fn create_managed_agent(
|
||||
try_regenerate_nest(&app);
|
||||
|
||||
// ── Phase 4: sync agent profile on relay (async, outside lock) ───────────
|
||||
let avatar_url = input
|
||||
.avatar_url
|
||||
.as_deref()
|
||||
.map(str::trim)
|
||||
.filter(|value| !value.is_empty())
|
||||
.map(str::to_string)
|
||||
.or_else(|| managed_agent_avatar_url(agent.agent_command.as_str()));
|
||||
// Use the avatar persisted on the record so the published profile and any
|
||||
// later reconciliation agree on the same value.
|
||||
let profile_sync_error = (sync_managed_agent_profile(
|
||||
&state,
|
||||
&resolved_relay_url,
|
||||
&agent_keys,
|
||||
&name,
|
||||
avatar_url.as_deref(),
|
||||
resolved_avatar_url.as_deref(),
|
||||
auth_tag.as_deref(),
|
||||
)
|
||||
.await)
|
||||
@@ -665,6 +676,20 @@ pub async fn create_managed_agent(
|
||||
})
|
||||
}
|
||||
|
||||
/// Data needed for background profile reconciliation after agent start.
|
||||
struct ProfileReconcileData {
|
||||
private_key_nsec: String,
|
||||
name: String,
|
||||
relay_url: String,
|
||||
/// Expected avatar URL for the published profile. Resolved at start from the
|
||||
/// record's persisted `avatar_url` (the exact URL published at creation),
|
||||
/// falling back to persona/command derivation only for pre-existing records
|
||||
/// that have no stored value — so old records still self-heal without
|
||||
/// regressing a user-overridden avatar.
|
||||
avatar_url: Option<String>,
|
||||
auth_tag: Option<String>,
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn start_managed_agent(
|
||||
pubkey: String,
|
||||
@@ -684,7 +709,8 @@ pub async fn start_managed_agent(
|
||||
}
|
||||
|
||||
// Collect backend info under lock; async preflight/spawn happens below.
|
||||
let target = {
|
||||
// Also snapshot profile reconciliation data for the background task.
|
||||
let (target, reconcile_data) = {
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
@@ -701,7 +727,17 @@ pub async fn start_managed_agent(
|
||||
|
||||
let record = find_managed_agent_mut(&mut records, &pubkey)?;
|
||||
|
||||
if record.backend == BackendKind::Local {
|
||||
let expected_avatar = reconcile_avatar(record.avatar_url.as_deref(), &record.agent_command);
|
||||
|
||||
let reconcile = ProfileReconcileData {
|
||||
private_key_nsec: record.private_key_nsec.clone(),
|
||||
name: record.name.clone(),
|
||||
relay_url: record.relay_url.clone(),
|
||||
avatar_url: expected_avatar,
|
||||
auth_tag: record.auth_tag.clone(),
|
||||
};
|
||||
|
||||
let target = if record.backend == BackendKind::Local {
|
||||
StartTarget::Local
|
||||
} else {
|
||||
StartTarget::Provider {
|
||||
@@ -709,10 +745,12 @@ pub async fn start_managed_agent(
|
||||
cached_binary_path: record.provider_binary_path.clone(),
|
||||
agent_json: build_deploy_payload(&app, record)?,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
(target, reconcile)
|
||||
};
|
||||
|
||||
match target {
|
||||
let result = match target {
|
||||
StartTarget::Local => {
|
||||
start_local_agent_with_preflight(&app, &state, &pubkey, &owner_hex, false).await
|
||||
}
|
||||
@@ -751,6 +789,98 @@ pub async fn start_managed_agent(
|
||||
StartTarget::Provider { backend, .. } => Err(format!(
|
||||
"agent {pubkey} has unsupported backend kind: {backend:?}"
|
||||
)),
|
||||
};
|
||||
|
||||
// ── Profile reconciliation (fire-and-forget) ────────────────────────────
|
||||
// On successful start, spawn a background task to ensure the agent's kind:0
|
||||
// profile is published on the relay. This self-heals cases where the initial
|
||||
// profile sync at creation time failed silently.
|
||||
if result.is_ok() {
|
||||
let reconcile_pubkey = pubkey.clone();
|
||||
let reconcile_app = app.clone();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
use tauri::Manager;
|
||||
let state = reconcile_app.state::<AppState>();
|
||||
if let Err(e) =
|
||||
reconcile_agent_profile(&state, &reconcile_pubkey, &reconcile_data).await
|
||||
{
|
||||
eprintln!(
|
||||
"sprout-desktop: profile reconciliation failed for agent {reconcile_pubkey}: {e}"
|
||||
);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
result
|
||||
}
|
||||
|
||||
/// Reconcile an agent's kind:0 profile on the relay.
|
||||
///
|
||||
/// Queries the relay for the agent's existing profile and re-publishes if missing
|
||||
/// or stale (display_name or picture mismatch). This is fire-and-forget — errors
|
||||
/// are returned to the caller for logging but never block agent startup.
|
||||
///
|
||||
/// Query and publish both target the agent's stored `relay_url` so that, under
|
||||
/// an active workspace relay override, reconciliation reads and writes the same
|
||||
/// relay the agent's profile actually lives on.
|
||||
async fn reconcile_agent_profile(
|
||||
state: &AppState,
|
||||
agent_pubkey: &str,
|
||||
data: &ProfileReconcileData,
|
||||
) -> Result<(), String> {
|
||||
use crate::relay::{query_agent_profile, sync_managed_agent_profile};
|
||||
|
||||
// Compare against the avatar persisted at creation time — never re-derive it.
|
||||
let expected_avatar = data.avatar_url.as_deref();
|
||||
|
||||
// Query the same relay the profile is published to (the stored relay_url).
|
||||
let existing = query_agent_profile(state, &data.relay_url, agent_pubkey).await?;
|
||||
|
||||
if !profile_needs_sync(existing.as_ref(), &data.name, expected_avatar) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let agent_keys = Keys::parse(&data.private_key_nsec)
|
||||
.map_err(|e| format!("failed to parse agent keys: {e}"))?;
|
||||
|
||||
sync_managed_agent_profile(
|
||||
state,
|
||||
&data.relay_url,
|
||||
&agent_keys,
|
||||
&data.name,
|
||||
expected_avatar,
|
||||
data.auth_tag.as_deref(),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Decide whether a published profile is missing or stale relative to the
|
||||
/// expected name and avatar. A missing profile always needs sync; a present
|
||||
/// one is stale when either the display name or picture diverges.
|
||||
fn profile_needs_sync(
|
||||
existing: Option<&crate::relay::AgentProfileInfo>,
|
||||
expected_name: &str,
|
||||
expected_avatar: Option<&str>,
|
||||
) -> bool {
|
||||
match existing {
|
||||
None => true,
|
||||
Some(info) => {
|
||||
let name_matches = info.display_name.as_deref() == Some(expected_name);
|
||||
let picture_matches = info.picture.as_deref() == expected_avatar;
|
||||
!name_matches || !picture_matches
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Resolve the avatar a managed agent's profile should reconcile against.
|
||||
/// Stored value (persisted at creation) wins; legacy records that predate the
|
||||
/// field (`stored == None`) fall back to the command-based derivation — the
|
||||
/// same source the create path used. Persona config is never consulted: doing
|
||||
/// so diverges from what was published and overwrites user intent on restart.
|
||||
fn reconcile_avatar(stored: Option<&str>, agent_command: &str) -> Option<String> {
|
||||
match stored {
|
||||
Some(url) => Some(url.to_string()),
|
||||
None => managed_agent_avatar_url(agent_command),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -969,4 +1099,91 @@ mod tests {
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
fn profile(name: Option<&str>, picture: Option<&str>) -> crate::relay::AgentProfileInfo {
|
||||
crate::relay::AgentProfileInfo {
|
||||
display_name: name.map(str::to_string),
|
||||
picture: picture.map(str::to_string),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_needs_sync_when_missing() {
|
||||
assert!(profile_needs_sync(None, "Duncan", Some("https://x/a.png")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_needs_sync_when_name_diverges() {
|
||||
let existing = profile(Some("Stilgar"), Some("https://x/a.png"));
|
||||
assert!(profile_needs_sync(
|
||||
Some(&existing),
|
||||
"Duncan",
|
||||
Some("https://x/a.png")
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_needs_sync_when_picture_diverges() {
|
||||
let existing = profile(Some("Duncan"), Some("https://x/old.png"));
|
||||
assert!(profile_needs_sync(
|
||||
Some(&existing),
|
||||
"Duncan",
|
||||
Some("https://x/new.png")
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_in_sync_when_name_and_picture_match() {
|
||||
let existing = profile(Some("Duncan"), Some("https://x/a.png"));
|
||||
assert!(!profile_needs_sync(
|
||||
Some(&existing),
|
||||
"Duncan",
|
||||
Some("https://x/a.png")
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_in_sync_when_both_avatars_absent() {
|
||||
let existing = profile(Some("Duncan"), None);
|
||||
assert!(!profile_needs_sync(Some(&existing), "Duncan", None));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_needs_sync_when_existing_name_is_none() {
|
||||
let existing = profile(None, Some("https://x/a.png"));
|
||||
assert!(profile_needs_sync(
|
||||
Some(&existing),
|
||||
"Duncan",
|
||||
Some("https://x/a.png"),
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn profile_needs_sync_when_expected_avatar_absent_but_published() {
|
||||
let existing = profile(Some("Duncan"), Some("https://x/a.png"));
|
||||
assert!(profile_needs_sync(Some(&existing), "Duncan", None));
|
||||
}
|
||||
|
||||
/// Legacy records (`avatar_url: None`) must reconcile against
|
||||
/// `managed_agent_avatar_url(agent_command)` — never persona config —
|
||||
/// matching what the original create path published.
|
||||
#[test]
|
||||
fn reconcile_avatar_legacy_record_uses_command_not_persona() {
|
||||
let resolved = reconcile_avatar(None, "goose");
|
||||
|
||||
assert_eq!(resolved, managed_agent_avatar_url("goose"));
|
||||
assert!(
|
||||
resolved.is_some(),
|
||||
"goose command should have a known avatar"
|
||||
);
|
||||
}
|
||||
|
||||
/// New records persist their avatar at creation; the stored value is used
|
||||
/// verbatim, never falling back to command derivation.
|
||||
#[test]
|
||||
fn reconcile_avatar_stored_value_wins() {
|
||||
let resolved = reconcile_avatar(Some("https://custom/avatar.png"), "goose");
|
||||
|
||||
assert_eq!(resolved.as_deref(), Some("https://custom/avatar.png"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -962,6 +962,7 @@ mod tests {
|
||||
private_key_nsec: String::new(),
|
||||
auth_tag: None,
|
||||
relay_url: String::new(),
|
||||
avatar_url: None,
|
||||
acp_command: String::new(),
|
||||
agent_command: String::new(),
|
||||
agent_args: vec![],
|
||||
|
||||
@@ -68,6 +68,7 @@ mod tests {
|
||||
private_key_nsec: "nsec1fake".into(),
|
||||
auth_tag: Some("tag".into()),
|
||||
relay_url: "ws://localhost:3000".into(),
|
||||
avatar_url: None,
|
||||
acp_command: "sprout-acp".into(),
|
||||
agent_command: "goose".into(),
|
||||
agent_args: vec![],
|
||||
|
||||
@@ -1255,6 +1255,7 @@ mod tests {
|
||||
private_key_nsec: "nsec1fake".into(),
|
||||
auth_tag,
|
||||
relay_url: "ws://localhost:3000".into(),
|
||||
avatar_url: None,
|
||||
acp_command: "sprout-acp".into(),
|
||||
agent_command: "goose".into(),
|
||||
agent_args: vec![],
|
||||
|
||||
@@ -82,6 +82,13 @@ pub struct ManagedAgentRecord {
|
||||
#[serde(default)]
|
||||
pub auth_tag: Option<String>,
|
||||
pub relay_url: String,
|
||||
/// Avatar URL resolved at creation time (user-supplied input, else the
|
||||
/// command-based fallback). Persisted so startup reconciliation compares
|
||||
/// against what was actually published rather than re-deriving it from
|
||||
/// persona config — which would silently overwrite user intent on restart.
|
||||
/// `#[serde(default)]` so pre-existing records deserialize as `None`.
|
||||
#[serde(default)]
|
||||
pub avatar_url: Option<String>,
|
||||
pub acp_command: String,
|
||||
pub agent_command: String,
|
||||
pub agent_args: Vec<String>,
|
||||
@@ -614,6 +621,7 @@ mod tests {
|
||||
.expect("legacy agent record without auth_tag should deserialize");
|
||||
|
||||
assert_eq!(record.auth_tag, None);
|
||||
assert_eq!(record.avatar_url, None);
|
||||
assert_eq!(record.pubkey, "abcd1234");
|
||||
}
|
||||
|
||||
|
||||
@@ -151,7 +151,18 @@ pub async fn query_relay(
|
||||
state: &AppState,
|
||||
filters: &[serde_json::Value],
|
||||
) -> Result<Vec<nostr::Event>, String> {
|
||||
let url = format!("{}/query", relay_api_base_url_with_override(state));
|
||||
query_relay_at(state, &relay_api_base_url_with_override(state), filters).await
|
||||
}
|
||||
|
||||
/// Like [`query_relay`] but targets an explicit HTTP API base URL instead of
|
||||
/// the workspace override. Used when a query must hit a specific relay (e.g.
|
||||
/// reconciling an agent's profile on the relay where it was published).
|
||||
pub async fn query_relay_at(
|
||||
state: &AppState,
|
||||
api_base_url: &str,
|
||||
filters: &[serde_json::Value],
|
||||
) -> Result<Vec<nostr::Event>, String> {
|
||||
let url = format!("{}/query", api_base_url);
|
||||
let body_bytes =
|
||||
serde_json::to_vec(filters).map_err(|e| format!("filter serialization failed: {e}"))?;
|
||||
let auth = build_nip98_auth_header(&Method::POST, &url, &body_bytes, state)?;
|
||||
@@ -282,6 +293,56 @@ pub async fn sync_managed_agent_profile(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ── Agent profile query ─────────────────────────────────────────────────────
|
||||
|
||||
/// Query the relay for an agent's kind:0 profile event.
|
||||
///
|
||||
/// Queries the relay identified by `relay_url` (typically the agent's stored
|
||||
/// `relay_url`) so the query targets the same host the profile is published to,
|
||||
/// even when a workspace relay override is active.
|
||||
///
|
||||
/// Returns the parsed profile content (display_name, picture) if a kind:0 event
|
||||
/// exists for the given pubkey, or `None` if no profile is published.
|
||||
pub async fn query_agent_profile(
|
||||
state: &AppState,
|
||||
relay_url: &str,
|
||||
agent_pubkey: &str,
|
||||
) -> Result<Option<AgentProfileInfo>, String> {
|
||||
let filter = serde_json::json!({
|
||||
"authors": [agent_pubkey],
|
||||
"kinds": [0],
|
||||
"limit": 1
|
||||
});
|
||||
|
||||
let events = query_relay_at(state, &relay_http_base_url(relay_url), &[filter]).await?;
|
||||
|
||||
let Some(event) = events.first() else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let Ok(content) = serde_json::from_str::<serde_json::Value>(&event.content) else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
Ok(Some(AgentProfileInfo {
|
||||
display_name: content
|
||||
.get("display_name")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
picture: content
|
||||
.get("picture")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
}))
|
||||
}
|
||||
|
||||
/// Parsed fields from a kind:0 profile event.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct AgentProfileInfo {
|
||||
pub display_name: Option<String>,
|
||||
pub picture: Option<String>,
|
||||
}
|
||||
|
||||
// ── Signed-event submission ─────────────────────────────────────────────────
|
||||
|
||||
/// Response from `POST /events`.
|
||||
|
||||
Reference in New Issue
Block a user