This commit is contained in:
2026-09-27 08:03:56 +07:00
parent 3a70bb7228
commit 4d2d058396
6 changed files with 59 additions and 134 deletions
+40 -53
View File
@@ -37,15 +37,14 @@ pub enum BackendEvent {
PassphraseRequired, PassphraseRequired,
/// The signer changed on login, logout or account switch. /// The signer changed on login, logout or account switch.
SignerChanged, SignerChanged,
/// New events were received from a relay and stored in the database. /// Events received from a relay, including this client's own events echoed
/// back by a relay after publishing.
NostrUpdate(Vec<Update>), NostrUpdate(Vec<Update>),
Synced, Synced,
SyncProgress { SyncProgress {
total: u64, total: u64,
current: u64, current: u64,
}, },
/// An event built locally was signed, broadcast and stored.
Published(Box<Event>),
Error(String), Error(String),
} }
@@ -91,13 +90,11 @@ impl Backend {
let pump: Task<Result<(), Error>> = cx.spawn(async move |this, cx| { let pump: Task<Result<(), Error>> = cx.spawn(async move |this, cx| {
let mut notifications = pump_client.notifications(); let mut notifications = pump_client.notifications();
let mut pending: Vec<Update> = Vec::new(); let mut pending: Vec<Update> = Vec::new();
let mut seen: HashSet<EventId> = HashSet::new();
'outer: loop { 'outer: loop {
match notifications.next().await { match next_update(&mut notifications, &mut seen).await {
Some(ClientNotification::Event { event, .. }) => { Some(update) => pending.push(update),
pending.push(Update::from_event(&event));
}
Some(_) => continue,
None => break, None => break,
} }
@@ -113,17 +110,11 @@ impl Backend {
let timer = cx.background_executor().timer(deadline - now); let timer = cx.background_executor().timer(deadline - now);
futures::pin_mut!(timer); futures::pin_mut!(timer);
let next = notifications.next(); let next = next_update(&mut notifications, &mut seen);
futures::pin_mut!(next); futures::pin_mut!(next);
match futures::future::select(next, timer).await { match futures::future::select(next, timer).await {
futures::future::Either::Left(( futures::future::Either::Left((Some(update), _)) => pending.push(update),
Some(ClientNotification::Event { event, .. }),
_,
)) => {
pending.push(Update::from_event(&event));
}
futures::future::Either::Left((Some(_), _)) => continue,
futures::future::Either::Left((None, _)) => break 'outer, futures::future::Either::Left((None, _)) => break 'outer,
futures::future::Either::Right(_) => break, futures::future::Either::Right(_) => break,
} }
@@ -500,9 +491,7 @@ impl Backend {
let builder = announcement.into_event_builder(); let builder = announcement.into_event_builder();
let event = builder.finalize_async(&signer).await?; let event = builder.finalize_async(&signer).await?;
let output = client.send_event(&event).broadcast().await?; let output = client.send_event(&event).broadcast().await?;
let event = require_relay_accepted(output, event)?; require_relay_accepted(output, event)?
this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?;
event
}; };
// The state event is the push authorization. Stage it on each // The state event is the push authorization. Stage it on each
@@ -557,15 +546,11 @@ impl Backend {
// Fan the state out to the relays once a git server holds the objects. // Fan the state out to the relays once a git server holds the objects.
// Staging already stored the event locally, publishing makes it // Staging already stored the event locally, publishing makes it
// visible to the other relays and clients. // visible to the other relays and clients.
if let Some(state_event) = &outcome.state_event { if let Some(state_event) = &outcome.state_event
if let Err(e) = client.send_event(state_event).broadcast().await { && let Err(e) = client.send_event(state_event).broadcast().await
{
log::warn!("failed to broadcast repository state: {e}"); log::warn!("failed to broadcast repository state: {e}");
} }
this.update(cx, |this, cx| {
this.announce_published(state_event.clone(), cx)
})
.ok();
}
let announcement = Announcement::from_event(&event) let announcement = Announcement::from_event(&event)
.ok_or_else(|| anyhow!("failed to parse announcement"))?; .ok_or_else(|| anyhow!("failed to parse announcement"))?;
@@ -653,9 +638,7 @@ impl Backend {
let builder = announcement.into_event_builder(); let builder = announcement.into_event_builder();
let event = builder.finalize_async(&signer).await?; let event = builder.finalize_async(&signer).await?;
let output = client.send_event(&event).broadcast().await?; let output = client.send_event(&event).broadcast().await?;
let event = require_relay_accepted(output, event)?; require_relay_accepted(output, event)?
this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?;
event
}; };
let refs = state.refs.clone(); let refs = state.refs.clone();
@@ -712,13 +695,11 @@ impl Backend {
// Fan the state out to the relays once a git server holds the objects. // Fan the state out to the relays once a git server holds the objects.
// Staging already stored the event locally, publishing makes it visible to the other relays and clients. // Staging already stored the event locally, publishing makes it visible to the other relays and clients.
if let Some(state_event) = &outcome.state_event { if let Some(state_event) = &outcome.state_event
if let Err(e) = client.send_event(state_event).broadcast().await { && let Err(e) = client.send_event(state_event).broadcast().await
{
log::warn!("failed to broadcast repository state: {e}"); log::warn!("failed to broadcast repository state: {e}");
} }
this.update(cx, |this, cx| this.announce_published(state_event.clone(), cx))
.ok();
}
} }
// Point `origin` at the first grasp server so later pushes have a target. // Point `origin` at the first grasp server so later pushes have a target.
@@ -894,15 +875,11 @@ impl Backend {
// Fan the state out to the relays once a git server holds the objects. // Fan the state out to the relays once a git server holds the objects.
// Staging already stored the event locally, publishing notifies // Staging already stored the event locally, publishing notifies
// the repository views and other relays and clients. // the repository views and other relays and clients.
if let Some(state_event) = &outcome.state_event { if let Some(state_event) = &outcome.state_event
if let Err(e) = client.send_event(state_event).broadcast().await { && let Err(e) = client.send_event(state_event).broadcast().await
{
log::warn!("failed to broadcast repository state: {e}"); log::warn!("failed to broadcast repository state: {e}");
} }
this.update(cx, |this, cx| {
this.announce_published(state_event.clone(), cx)
})
.ok();
}
Ok(outcome) Ok(outcome)
}) })
@@ -1299,14 +1276,6 @@ impl Backend {
.detach(); .detach();
} }
/// Emit [`BackendEvent::Published`] for cross-store invalidation.
///
/// Callers publish with `client.send_event(...)` directly, then call this
/// so stores like `RepoListStore` refresh without re-querying the relays.
pub fn announce_published(&self, event: Event, cx: &mut Context<Self>) {
cx.emit(BackendEvent::Published(Box::new(event)));
}
/// Publish a NIP-09 deletion for each of `events`, best-effort. /// Publish a NIP-09 deletion for each of `events`, best-effort.
/// ///
/// Each target gets its own deletion event: a relay rejecting or /// Each target gets its own deletion event: a relay rejecting or
@@ -1329,6 +1298,26 @@ impl Backend {
} }
} }
/// Await the next relay event from the notification stream.
async fn next_update(
notifications: &mut (impl futures::Stream<Item = ClientNotification> + Unpin),
seen: &mut HashSet<EventId>,
) -> Option<Update> {
loop {
match notifications.next().await {
Some(ClientNotification::Message { message, .. }) => {
if let RelayMessage::Event { event, .. } = *message
&& seen.insert(event.id)
{
return Some(Update::from_event(&event));
}
}
Some(_) => continue,
None => return None,
}
}
}
/// Sign and send a single NIP-09 deletion request for `event`. /// Sign and send a single NIP-09 deletion request for `event`.
async fn retract_event( async fn retract_event(
client: &Client, client: &Client,
@@ -1338,16 +1327,14 @@ async fn retract_event(
let builder = EventDeletionRequest::new() let builder = EventDeletionRequest::new()
.id(event.id) .id(event.id)
.into_event_builder(); .into_event_builder();
let deletion = builder.finalize_async(signer).await?; let deletion = builder.finalize_async(signer).await?;
client.send_event(&deletion).broadcast().await?; client.send_event(&deletion).broadcast().await?;
Ok(()) Ok(())
} }
/// The event was accepted by at least one relay, or a descriptive error otherwise. /// The event was accepted by at least one relay, or a descriptive error otherwise.
///
/// The SDK does not treat "accepted by zero relays" as an error on its own:
/// [`SendEventOutput::success`] may be empty while the call still returns `Ok`.
/// This turns that case into an error the caller can surface.
pub(crate) fn require_relay_accepted( pub(crate) fn require_relay_accepted(
output: SendEventOutput, output: SendEventOutput,
event: Event, event: Event,
+2 -9
View File
@@ -93,8 +93,8 @@ impl ProfileStore {
pub(crate) fn new(cx: &mut Context<Self>) -> Self { pub(crate) fn new(cx: &mut Context<Self>) -> Self {
let backend = Backend::global(cx); let backend = Backend::global(cx);
let subscription = cx.subscribe(&backend, |this, _backend, event, cx| match event { let subscription = cx.subscribe(&backend, |this, _backend, event, cx| {
BackendEvent::NostrUpdate(updates) => { if let BackendEvent::NostrUpdate(updates) = event {
for update in updates for update in updates
.iter() .iter()
.filter(|update| update.kind == Kind::Metadata) .filter(|update| update.kind == Kind::Metadata)
@@ -102,13 +102,6 @@ impl ProfileStore {
this.apply_author(update.author, cx); this.apply_author(update.author, cx);
} }
} }
BackendEvent::Published(event) if event.kind == Kind::Metadata => {
let metadata = Metadata::from_json(&event.content).unwrap_or_default();
this.profiles
.insert(event.pubkey, Profile::new(event.pubkey, metadata));
cx.notify();
}
_ => {}
}); });
// Fetch requests are queued on a channel, batched into one sync per debounce window. // Fetch requests are queued on a channel, batched into one sync per debounce window.
+2 -43
View File
@@ -225,20 +225,6 @@ impl RepoStore {
deletion || coordinate || authored || comment || status deletion || coordinate || authored || comment || status
}), }),
BackendEvent::Published(event) => {
let announcement = event.kind == Kind::GitRepoAnnouncement;
let author = event.pubkey == addr.public_key;
let coordinate = event.tags.coordinates().into_iter().any(|c| c == *addr);
let state = event.kind == Kind::RepoState
&& author
&& event.tags.identifier().as_deref() == Some(addr.identifier.as_str());
let deletion =
event.kind == Kind::EventDeletion || event.kind == Kind::RequestToVanish;
coordinate || (announcement && author) || state || deletion
}
_ => false, _ => false,
}; };
@@ -977,11 +963,6 @@ impl RepoStore {
} }
}; };
this.update(cx, |_this, cx| {
Backend::global(cx)
.update(cx, |backend, cx| backend.announce_published(pr_event.clone(), cx))
})?;
// A draft PR carries a kind-1633 status event, NIP-34. // A draft PR carries a kind-1633 status event, NIP-34.
// Publish it right after the PR event so viewers never show it open. // Publish it right after the PR event so viewers never show it open.
if draft { if draft {
@@ -1171,21 +1152,12 @@ impl RepoStore {
} }
.await; .await;
match publish_result { if let Err(e) = publish_result {
Ok(event) => {
this.update(cx, |_this, cx| {
Backend::global(cx).update(cx, |backend, cx| {
backend.announce_published(event.clone(), cx)
})
})?;
}
Err(e) => {
return this.update(cx, |this, cx| { return this.update(cx, |this, cx| {
this.last_error = Some(e.to_string()); this.last_error = Some(e.to_string());
cx.notify(); cx.notify();
}); });
} }
}
Ok(()) Ok(())
}) })
@@ -1622,20 +1594,12 @@ impl RepoStore {
} }
.await; .await;
match publish_result { if let Err(e) = publish_result {
Ok(event) => {
this.update(cx, |_this, cx| {
Backend::global(cx)
.update(cx, |backend, cx| backend.announce_published(event, cx))
})?;
}
Err(e) => {
this.update(cx, |this, cx| { this.update(cx, |this, cx| {
this.last_error = Some(e.to_string()); this.last_error = Some(e.to_string());
cx.notify(); cx.notify();
})?; })?;
} }
}
Ok(()) Ok(())
}); });
@@ -1767,11 +1731,6 @@ async fn publish_patch_series(
let event = builder.finalize_async(&signer).await?; let event = builder.finalize_async(&signer).await?;
let output = client.send_event(&event).broadcast().await?; let output = client.send_event(&event).broadcast().await?;
let event = require_relay_accepted(output, event)?; let event = require_relay_accepted(output, event)?;
this.update(cx, |_this, cx| {
Backend::global(cx).update(cx, |backend, cx| {
backend.announce_published(event.clone(), cx)
})
})?;
if root.is_none() { if root.is_none() {
root = Some(event.clone()); root = Some(event.clone());
-8
View File
@@ -82,14 +82,6 @@ impl RepoListStore {
is_announcement || is_repo_state is_announcement || is_repo_state
} }
}), }),
BackendEvent::Published(event) => {
let announcement = event.kind == Kind::GitRepoAnnouncement;
let state = event.kind == Kind::RepoState;
let deletion =
event.kind == Kind::EventDeletion || event.kind == Kind::RequestToVanish;
announcement || state || deletion
}
BackendEvent::SignerChanged => { BackendEvent::SignerChanged => {
this.state_synced_repos.clear(); this.state_synced_repos.clear();
true true
+1 -1
View File
@@ -171,7 +171,7 @@ impl InboxView {
self.refresh(cx); self.refresh(cx);
} }
} }
BackendEvent::Synced | BackendEvent::Published(_) => self.refresh(cx), BackendEvent::Synced => self.refresh(cx),
_ => {} _ => {}
} }
} }
+3 -9
View File
@@ -79,16 +79,12 @@ impl SidebarPanel {
let signer_changed = matches!(event, BackendEvent::SignerChanged); let signer_changed = matches!(event, BackendEvent::SignerChanged);
let signer_required = matches!(event, BackendEvent::SignerRequired); let signer_required = matches!(event, BackendEvent::SignerRequired);
if !signer_changed && !signer_required {
return;
}
if signer_required { if signer_required {
this.banner = pick_banner(); this.banner = pick_banner();
cx.notify(); cx.notify();
} }
if this.refresh(cx) || signer_required { if this.refresh(cx) || signer_changed {
cx.notify(); cx.notify();
} }
})); }));
@@ -131,16 +127,14 @@ impl SidebarPanel {
let user = backend.read(cx).current_user(); let user = backend.read(cx).current_user();
let (announcements, local_repos, scanning) = { let (announcements, local_repos, scanning) = {
let repo_list = RepoListStore::global(cx); let repo_list = RepoListStore::global(cx).read(cx);
let repo_list = repo_list.read(cx);
let announcements = user let announcements = user
.as_ref() .as_ref()
.map(|user| repo_list.announcements_of(user)) .map(|user| repo_list.announcements_of(user))
.unwrap_or_default(); .unwrap_or_default();
let local = LocalReposStore::global(cx); let local = LocalReposStore::global(cx).read(cx);
let local = local.read(cx);
let local_repos = let local_repos =
resolve_local_repos(&local.repos, &repo_list.announcements, &announcements); resolve_local_repos(&local.repos, &repo_list.announcements, &announcements);