feat: add event fetching strategy (#24)
Reviewed-on: #24
This commit was merged in pull request #24.
This commit is contained in:
@@ -23,12 +23,7 @@ pub const USER_KEYRING: &str = "Signed Safe Storage";
|
||||
pub const NOSTR_CONNECT_TIMEOUT: u64 = 60;
|
||||
|
||||
/// Relays connected at startup, before any user-specific relay config is known.
|
||||
pub const BOOTSTRAP_RELAYS: [&str; 4] = [
|
||||
"wss://relay.primal.net",
|
||||
"wss://relay.ditto.pub",
|
||||
"wss://index.ngit.dev",
|
||||
"wss://profiles.nostr1.com",
|
||||
];
|
||||
pub const BOOTSTRAP_RELAYS: [&str; 2] = ["wss://relay.ditto.pub", "wss://index.ngit.dev"];
|
||||
|
||||
/// Relays used to index the user's NIP-65 relay list.
|
||||
pub const INDEXER_RELAYS: [&str; 2] = ["wss://indexer.coracle.social", "wss://user.kindpag.es"];
|
||||
@@ -42,20 +37,15 @@ pub enum BackendEvent {
|
||||
PassphraseRequired,
|
||||
/// The signer changed on login, logout or account switch.
|
||||
SignerChanged,
|
||||
/// New events were received from a relay and stored in the database.
|
||||
///
|
||||
/// Batched: [`Backend`]'s notification pump coalesces everything a
|
||||
/// relay delivers within one debounce window into a single event,
|
||||
/// instead of emitting per-event and making every subscriber debounce
|
||||
/// the same burst independently.
|
||||
NostrUpdate(Vec<Update>),
|
||||
/// Kind-0 metadata arrived for these authors; re-read them from the store.
|
||||
ProfileUpdates(Vec<PublicKey>),
|
||||
/// Repository events arrived: announcements, states, activity and deletions.
|
||||
RepoUpdates(Vec<Update>),
|
||||
Synced,
|
||||
SyncProgress {
|
||||
total: u64,
|
||||
current: u64,
|
||||
},
|
||||
/// An event built locally was signed, broadcast and stored.
|
||||
Published(Box<Event>),
|
||||
Error(String),
|
||||
}
|
||||
|
||||
@@ -100,14 +90,16 @@ impl Backend {
|
||||
|
||||
let pump: Task<Result<(), Error>> = cx.spawn(async move |this, cx| {
|
||||
let mut notifications = pump_client.notifications();
|
||||
let mut pending: Vec<Update> = Vec::new();
|
||||
let mut pending_profiles: HashSet<PublicKey> = HashSet::new();
|
||||
let mut pending_repos: Vec<Update> = Vec::new();
|
||||
let mut seen: HashSet<EventId> = HashSet::new();
|
||||
|
||||
'outer: loop {
|
||||
match notifications.next().await {
|
||||
Some(ClientNotification::Event { event, .. }) => {
|
||||
pending.push(Update::from_event(&event));
|
||||
match next_update(&mut notifications, &mut seen).await {
|
||||
Some(UpdateEvent::Profile(author)) => {
|
||||
pending_profiles.insert(author);
|
||||
}
|
||||
Some(_) => continue,
|
||||
Some(UpdateEvent::Repo(update)) => pending_repos.push(update),
|
||||
None => break,
|
||||
}
|
||||
|
||||
@@ -123,28 +115,33 @@ impl Backend {
|
||||
let timer = cx.background_executor().timer(deadline - now);
|
||||
futures::pin_mut!(timer);
|
||||
|
||||
let next = notifications.next();
|
||||
let next = next_update(&mut notifications, &mut seen);
|
||||
futures::pin_mut!(next);
|
||||
|
||||
match futures::future::select(next, timer).await {
|
||||
futures::future::Either::Left((
|
||||
Some(ClientNotification::Event { event, .. }),
|
||||
_,
|
||||
)) => {
|
||||
pending.push(Update::from_event(&event));
|
||||
futures::future::Either::Left((Some(UpdateEvent::Profile(author)), _)) => {
|
||||
pending_profiles.insert(author);
|
||||
}
|
||||
futures::future::Either::Left((Some(UpdateEvent::Repo(update)), _)) => {
|
||||
pending_repos.push(update);
|
||||
}
|
||||
futures::future::Either::Left((Some(_), _)) => continue,
|
||||
futures::future::Either::Left((None, _)) => break 'outer,
|
||||
futures::future::Either::Right(_) => break,
|
||||
}
|
||||
}
|
||||
|
||||
let batch = std::mem::take(&mut pending);
|
||||
let profiles: Vec<PublicKey> = pending_profiles.drain().collect();
|
||||
let repos = std::mem::take(&mut pending_repos);
|
||||
|
||||
if let Err(e) =
|
||||
this.update(cx, |_this, cx| cx.emit(BackendEvent::NostrUpdate(batch)))
|
||||
{
|
||||
log::warn!("failed to emit nostr update: {e}");
|
||||
if let Err(e) = this.update(cx, |_this, cx| {
|
||||
if !profiles.is_empty() {
|
||||
cx.emit(BackendEvent::ProfileUpdates(profiles));
|
||||
}
|
||||
if !repos.is_empty() {
|
||||
cx.emit(BackendEvent::RepoUpdates(repos));
|
||||
}
|
||||
}) {
|
||||
log::warn!("failed to emit backend update: {e}");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -446,9 +443,6 @@ impl Backend {
|
||||
let client = self.client.clone();
|
||||
|
||||
// Initialize directly at the user's chosen destination.
|
||||
// No mirror is pre-populated: `GitCache::ensure_clone` lazily clones
|
||||
// from the grasp server the first time the repo detail view needs it,
|
||||
// exactly like every other repository.
|
||||
let destination = {
|
||||
let dir_name = signed_git::sanitize_path_component(&name);
|
||||
let dir_name = if dir_name.is_empty() {
|
||||
@@ -513,9 +507,7 @@ impl Backend {
|
||||
let builder = announcement.into_event_builder();
|
||||
let event = builder.finalize_async(&signer).await?;
|
||||
let output = client.send_event(&event).broadcast().await?;
|
||||
let event = require_relay_accepted(output, event)?;
|
||||
this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?;
|
||||
event
|
||||
require_relay_accepted(output, event)?
|
||||
};
|
||||
|
||||
// The state event is the push authorization. Stage it on each
|
||||
@@ -570,14 +562,10 @@ impl Backend {
|
||||
// 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.
|
||||
if let Some(state_event) = &outcome.state_event {
|
||||
if let Err(e) = client.send_event(state_event).broadcast().await {
|
||||
log::warn!("failed to broadcast repository state: {e}");
|
||||
}
|
||||
this.update(cx, |this, cx| {
|
||||
this.announce_published(state_event.clone(), cx)
|
||||
})
|
||||
.ok();
|
||||
if let Some(state_event) = &outcome.state_event
|
||||
&& let Err(e) = client.send_event(state_event).broadcast().await
|
||||
{
|
||||
log::warn!("failed to broadcast repository state: {e}");
|
||||
}
|
||||
|
||||
let announcement = Announcement::from_event(&event)
|
||||
@@ -666,9 +654,7 @@ impl Backend {
|
||||
let builder = announcement.into_event_builder();
|
||||
let event = builder.finalize_async(&signer).await?;
|
||||
let output = client.send_event(&event).broadcast().await?;
|
||||
let event = require_relay_accepted(output, event)?;
|
||||
this.update(cx, |this, cx| this.announce_published(event.clone(), cx))?;
|
||||
event
|
||||
require_relay_accepted(output, event)?
|
||||
};
|
||||
|
||||
let refs = state.refs.clone();
|
||||
@@ -725,12 +711,10 @@ impl Backend {
|
||||
|
||||
// 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.
|
||||
if let Some(state_event) = &outcome.state_event {
|
||||
if let Err(e) = client.send_event(state_event).broadcast().await {
|
||||
log::warn!("failed to broadcast repository state: {e}");
|
||||
}
|
||||
this.update(cx, |this, cx| this.announce_published(state_event.clone(), cx))
|
||||
.ok();
|
||||
if let Some(state_event) = &outcome.state_event
|
||||
&& let Err(e) = client.send_event(state_event).broadcast().await
|
||||
{
|
||||
log::warn!("failed to broadcast repository state: {e}");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -907,14 +891,10 @@ impl Backend {
|
||||
// Fan the state out to the relays once a git server holds the objects.
|
||||
// Staging already stored the event locally, publishing notifies
|
||||
// the repository views and other relays and clients.
|
||||
if let Some(state_event) = &outcome.state_event {
|
||||
if let Err(e) = client.send_event(state_event).broadcast().await {
|
||||
log::warn!("failed to broadcast repository state: {e}");
|
||||
}
|
||||
this.update(cx, |this, cx| {
|
||||
this.announce_published(state_event.clone(), cx)
|
||||
})
|
||||
.ok();
|
||||
if let Some(state_event) = &outcome.state_event
|
||||
&& let Err(e) = client.send_event(state_event).broadcast().await
|
||||
{
|
||||
log::warn!("failed to broadcast repository state: {e}");
|
||||
}
|
||||
|
||||
Ok(outcome)
|
||||
@@ -1086,7 +1066,7 @@ impl Backend {
|
||||
)
|
||||
.await?;
|
||||
|
||||
for url in user_grasp_list_servers(client.clone(), public_key).await? {
|
||||
for url in user_grasp_list_servers(&client, public_key).await? {
|
||||
client.add_relay(url).and_connect().await.ok();
|
||||
}
|
||||
|
||||
@@ -1205,6 +1185,21 @@ impl Backend {
|
||||
.detach();
|
||||
}
|
||||
|
||||
/// Sync filters through the SDK's NIP-65 gossip targeting.
|
||||
pub fn sync_auto(&mut self, filters: Vec<Filter>, cx: &mut Context<Self>) {
|
||||
let client = self.client.clone();
|
||||
|
||||
cx.spawn(async move |_this, _cx| {
|
||||
for filter in filters {
|
||||
if let Err(e) = client.sync(filter).await {
|
||||
log::warn!("gossip relay fetch failed: {e}");
|
||||
}
|
||||
}
|
||||
Ok::<(), Error>(())
|
||||
})
|
||||
.detach();
|
||||
}
|
||||
|
||||
pub fn subscribe_bootstrap(&mut self, filters: Vec<Filter>, cx: &mut Context<Self>) {
|
||||
let client = self.client.clone();
|
||||
|
||||
@@ -1222,7 +1217,8 @@ impl Backend {
|
||||
.detach();
|
||||
}
|
||||
|
||||
pub fn sync_bootstrap(&mut self, filter: Filter, cx: &mut Context<Self>) {
|
||||
/// Sync several bootstrap filters in order, within a single task.
|
||||
pub fn sync_bootstraps(&mut self, filters: Vec<Filter>, cx: &mut Context<Self>) {
|
||||
let client = self.client.clone();
|
||||
let (tx, mut rx) = SyncProgress::channel();
|
||||
|
||||
@@ -1259,8 +1255,19 @@ impl Backend {
|
||||
.detach();
|
||||
|
||||
let sync = cx.background_spawn(async move {
|
||||
let opts = SyncOptions::default().progress(tx);
|
||||
sync_bootstrap_only(&client, filter, opts).await
|
||||
let mut first_error = None;
|
||||
|
||||
for filter in filters {
|
||||
let opts = SyncOptions::default().progress(tx.clone());
|
||||
if let Err(error) = sync_bootstrap_only(&client, filter, opts).await {
|
||||
first_error.get_or_insert(error);
|
||||
}
|
||||
}
|
||||
|
||||
match first_error {
|
||||
Some(error) => Err(error),
|
||||
None => Ok(()),
|
||||
}
|
||||
});
|
||||
|
||||
cx.spawn(async move |this, cx| {
|
||||
@@ -1285,14 +1292,6 @@ impl Backend {
|
||||
.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.
|
||||
///
|
||||
/// Each target gets its own deletion event: a relay rejecting or
|
||||
@@ -1315,6 +1314,42 @@ impl Backend {
|
||||
}
|
||||
}
|
||||
|
||||
/// A relay event the backend routes to a store group.
|
||||
enum UpdateEvent {
|
||||
Profile(PublicKey),
|
||||
Repo(Update),
|
||||
}
|
||||
|
||||
/// 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<UpdateEvent> {
|
||||
loop {
|
||||
match notifications.next().await {
|
||||
Some(ClientNotification::Message { message, .. }) => {
|
||||
let RelayMessage::Event { event, .. } = *message else {
|
||||
continue;
|
||||
};
|
||||
|
||||
let update = match event.kind {
|
||||
Kind::Metadata => UpdateEvent::Profile(event.pubkey),
|
||||
kind if filters::is_repo_kind(kind) => {
|
||||
UpdateEvent::Repo(Update::from_event(&event))
|
||||
}
|
||||
_ => continue,
|
||||
};
|
||||
|
||||
if seen.insert(event.id) {
|
||||
return Some(update);
|
||||
}
|
||||
}
|
||||
Some(_) => continue,
|
||||
None => return None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Sign and send a single NIP-09 deletion request for `event`.
|
||||
async fn retract_event(
|
||||
client: &Client,
|
||||
@@ -1324,16 +1359,14 @@ async fn retract_event(
|
||||
let builder = EventDeletionRequest::new()
|
||||
.id(event.id)
|
||||
.into_event_builder();
|
||||
|
||||
let deletion = builder.finalize_async(signer).await?;
|
||||
client.send_event(&deletion).broadcast().await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// 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(
|
||||
output: SendEventOutput,
|
||||
event: Event,
|
||||
@@ -1419,6 +1452,7 @@ pub(crate) async fn sync_bootstrap_only(
|
||||
.with(BOOTSTRAP_RELAYS)
|
||||
.opts(opts)
|
||||
.await?;
|
||||
|
||||
Ok(output.value)
|
||||
}
|
||||
|
||||
@@ -1491,7 +1525,7 @@ fn latest_grasp_list_servers(events: Vec<Event>) -> Vec<RelayUrl> {
|
||||
}
|
||||
|
||||
pub async fn user_grasp_list_servers(
|
||||
client: Client,
|
||||
client: &Client,
|
||||
user: PublicKey,
|
||||
) -> Result<Vec<RelayUrl>, Error> {
|
||||
let events: Vec<Event> = client
|
||||
|
||||
Reference in New Issue
Block a user