sync all repos
This commit is contained in:
@@ -68,8 +68,10 @@ pub fn announcements_by(public_key: PublicKey) -> Filter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// All repository announcements (for global discovery).
|
/// All repository announcements (for global discovery).
|
||||||
pub fn all_announcements(limit: usize) -> Filter {
|
///
|
||||||
Filter::new()
|
/// Unbounded: intended for negentropy sync, which reconciles sets
|
||||||
.kind(Kind::GitRepoAnnouncement)
|
/// efficiently regardless of size. Local database queries with this
|
||||||
.limit(limit)
|
/// filter are served by LMDB, so they stay fast as the database grows.
|
||||||
|
pub fn all_announcements() -> Filter {
|
||||||
|
Filter::new().kind(Kind::GitRepoAnnouncement)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -45,6 +45,14 @@ pub enum BackendEvent {
|
|||||||
/// so stores should re-query (no [`BackendEvent::NostrUpdate`] is fired
|
/// so stores should re-query (no [`BackendEvent::NostrUpdate`] is fired
|
||||||
/// for synced events).
|
/// for synced events).
|
||||||
Synced,
|
Synced,
|
||||||
|
/// A negentropy sync is in flight. Stores may re-query to render
|
||||||
|
/// incrementally; UI can show `current`/`total` progress.
|
||||||
|
SyncProgress {
|
||||||
|
/// Total events to process.
|
||||||
|
total: u64,
|
||||||
|
/// Events processed so far.
|
||||||
|
current: u64,
|
||||||
|
},
|
||||||
/// An event built locally was signed, broadcast and stored.
|
/// An event built locally was signed, broadcast and stored.
|
||||||
Published(Box<Event>),
|
Published(Box<Event>),
|
||||||
/// An error occurred.
|
/// An error occurred.
|
||||||
@@ -67,6 +75,7 @@ pub struct Backend {
|
|||||||
inner: NostrBackend,
|
inner: NostrBackend,
|
||||||
current_user: Option<PublicKey>,
|
current_user: Option<PublicKey>,
|
||||||
connected: bool,
|
connected: bool,
|
||||||
|
sync_progress: Option<(u64, u64)>,
|
||||||
tasks: Vec<Task<Result<(), Error>>>,
|
tasks: Vec<Task<Result<(), Error>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -114,6 +123,7 @@ impl Backend {
|
|||||||
inner,
|
inner,
|
||||||
current_user: None,
|
current_user: None,
|
||||||
connected: false,
|
connected: false,
|
||||||
|
sync_progress: None,
|
||||||
tasks: vec![pump],
|
tasks: vec![pump],
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -388,6 +398,11 @@ impl Backend {
|
|||||||
self.connected
|
self.connected
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Progress of the in-flight negentropy sync, if any: `(total, current)`.
|
||||||
|
pub fn sync_progress(&self) -> Option<(u64, u64)> {
|
||||||
|
self.sync_progress
|
||||||
|
}
|
||||||
|
|
||||||
/// Update the signer (any type implementing the async signer traits,
|
/// Update the signer (any type implementing the async signer traits,
|
||||||
/// e.g. `Keys`, `NostrConnect`, a browser extension proxy).
|
/// e.g. `Keys`, `NostrConnect`, a browser extension proxy).
|
||||||
pub fn set_signer<T>(&mut self, new_signer: T, cx: &mut Context<Self>)
|
pub fn set_signer<T>(&mut self, new_signer: T, cx: &mut Context<Self>)
|
||||||
@@ -506,12 +521,47 @@ impl Backend {
|
|||||||
|
|
||||||
/// Negentropy-sync the given filter against the bootstrap relays:
|
/// Negentropy-sync the given filter against the bootstrap relays:
|
||||||
/// reconciles the local database with the relays in both directions.
|
/// reconciles the local database with the relays in both directions.
|
||||||
/// Emits [`BackendEvent::Synced`] on completion.
|
/// Emits [`BackendEvent::SyncProgress`] while running (throttled to
|
||||||
|
/// whole-percent changes) and [`BackendEvent::Synced`] on completion.
|
||||||
pub fn sync_bootstrap(&mut self, filter: Filter, cx: &mut Context<Self>) {
|
pub fn sync_bootstrap(&mut self, filter: Filter, cx: &mut Context<Self>) {
|
||||||
let backend = self.inner.clone();
|
let backend = self.inner.clone();
|
||||||
|
|
||||||
|
self.sync_progress = Some((0, 0));
|
||||||
|
cx.notify();
|
||||||
|
|
||||||
|
let (tx, mut rx) = SyncProgress::channel();
|
||||||
|
|
||||||
|
self.tasks.push(cx.spawn(async move |this, cx| {
|
||||||
|
let mut last_percent: u64 = 0;
|
||||||
|
|
||||||
|
while rx.changed().await.is_ok() {
|
||||||
|
let progress = *rx.borrow_and_update();
|
||||||
|
let percent = (progress.percentage() * 100.0) as u64;
|
||||||
|
|
||||||
|
if progress.current > 0 && percent != last_percent {
|
||||||
|
last_percent = percent;
|
||||||
|
|
||||||
|
let alive = this.update(cx, |this, cx| {
|
||||||
|
this.sync_progress = Some((progress.total, progress.current));
|
||||||
|
cx.emit(BackendEvent::SyncProgress {
|
||||||
|
total: progress.total,
|
||||||
|
current: progress.current,
|
||||||
|
});
|
||||||
|
cx.notify();
|
||||||
|
});
|
||||||
|
|
||||||
|
if alive.is_err() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}));
|
||||||
|
|
||||||
let task = cx.background_spawn(async move {
|
let task = cx.background_spawn(async move {
|
||||||
sync_bootstrap_only(&backend.client(), filter).await
|
let opts = SyncOptions::default().progress(tx);
|
||||||
|
sync_bootstrap_only(&backend.client(), filter, opts).await
|
||||||
});
|
});
|
||||||
|
|
||||||
self.tasks.push(cx.spawn(async move |this, cx| {
|
self.tasks.push(cx.spawn(async move |this, cx| {
|
||||||
@@ -522,13 +572,17 @@ impl Backend {
|
|||||||
summary.received.len(),
|
summary.received.len(),
|
||||||
summary.sent.len()
|
summary.sent.len()
|
||||||
);
|
);
|
||||||
this.update(cx, |_this, cx| {
|
this.update(cx, |this, cx| {
|
||||||
|
this.sync_progress = None;
|
||||||
cx.emit(BackendEvent::Synced);
|
cx.emit(BackendEvent::Synced);
|
||||||
cx.notify();
|
cx.notify();
|
||||||
})?;
|
})?;
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
this.update(cx, |_this, cx| cx.emit(BackendEvent::error(e.to_string())))?;
|
this.update(cx, |this, cx| {
|
||||||
|
this.sync_progress = None;
|
||||||
|
cx.emit(BackendEvent::error(e.to_string()))
|
||||||
|
})?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
@@ -596,7 +650,8 @@ pub(crate) async fn subscribe_bootstrap_only(client: &Client, filters: Vec<Filte
|
|||||||
pub(crate) async fn sync_bootstrap_only(
|
pub(crate) async fn sync_bootstrap_only(
|
||||||
client: &Client,
|
client: &Client,
|
||||||
filter: Filter,
|
filter: Filter,
|
||||||
|
opts: SyncOptions,
|
||||||
) -> Result<SyncSummary, Error> {
|
) -> Result<SyncSummary, Error> {
|
||||||
let output = client.sync(filter).with(BOOTSTRAP_RELAYS).await?;
|
let output = client.sync(filter).with(BOOTSTRAP_RELAYS).opts(opts).await?;
|
||||||
Ok(output.value)
|
Ok(output.value)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -257,7 +257,7 @@ impl ProfileStore {
|
|||||||
// Negentropy-sync with the bootstrap relays. Synced events
|
// Negentropy-sync with the bootstrap relays. Synced events
|
||||||
// are written to the database directly (no NostrUpdate), so
|
// are written to the database directly (no NostrUpdate), so
|
||||||
// re-apply from the database afterwards.
|
// re-apply from the database afterwards.
|
||||||
match sync_bootstrap_only(&client, filter).await {
|
match sync_bootstrap_only(&client, filter, SyncOptions::default()).await {
|
||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
this.update(cx, |this, cx| this.apply_seen(cx))?;
|
this.update(cx, |this, cx| this.apply_seen(cx))?;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -31,7 +31,7 @@ impl RepoListStore {
|
|||||||
event.kind == Kind::GitRepoAnnouncement
|
event.kind == Kind::GitRepoAnnouncement
|
||||||
&& this.author.is_none_or(|a| a == event.pubkey)
|
&& this.author.is_none_or(|a| a == event.pubkey)
|
||||||
}
|
}
|
||||||
BackendEvent::Synced => true,
|
BackendEvent::Synced | BackendEvent::SyncProgress { .. } => true,
|
||||||
_ => false,
|
_ => false,
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -69,7 +69,7 @@ impl RepoListStore {
|
|||||||
backend.update(cx, |backend, cx| {
|
backend.update(cx, |backend, cx| {
|
||||||
let filter = match author {
|
let filter = match author {
|
||||||
Some(a) => filters::announcements_by(a),
|
Some(a) => filters::announcements_by(a),
|
||||||
None => filters::all_announcements(500),
|
None => filters::all_announcements(),
|
||||||
};
|
};
|
||||||
backend.sync_bootstrap(filter, cx);
|
backend.sync_bootstrap(filter, cx);
|
||||||
});
|
});
|
||||||
@@ -93,7 +93,7 @@ impl RepoListStore {
|
|||||||
loop {
|
loop {
|
||||||
let filter = match author {
|
let filter = match author {
|
||||||
Some(a) => filters::announcements_by(a),
|
Some(a) => filters::announcements_by(a),
|
||||||
None => filters::all_announcements(500),
|
None => filters::all_announcements(),
|
||||||
};
|
};
|
||||||
|
|
||||||
let events = match client.database().query(filter).await {
|
let events = match client.database().query(filter).await {
|
||||||
|
|||||||
@@ -18,10 +18,15 @@ impl Workspace {
|
|||||||
let repo_list = cx.new(|cx| RepoListView::new(window, cx));
|
let repo_list = cx.new(|cx| RepoListView::new(window, cx));
|
||||||
|
|
||||||
let connected = backend.read(cx).is_connected();
|
let connected = backend.read(cx).is_connected();
|
||||||
|
let sync_progress = backend.read(cx).sync_progress();
|
||||||
|
|
||||||
let subscription = cx.subscribe(&backend, |this, _backend, event, cx| {
|
let subscription = cx.subscribe(&backend, |this, _backend, event, cx| {
|
||||||
match event {
|
match event {
|
||||||
BackendEvent::Connected => this.status = "Connected".into(),
|
BackendEvent::Connected => this.status = "Connected".into(),
|
||||||
|
BackendEvent::SyncProgress { total, current } => {
|
||||||
|
this.status = format!("Syncing repositories... {current}/{total}").into()
|
||||||
|
}
|
||||||
|
BackendEvent::Synced => this.status = "Connected".into(),
|
||||||
BackendEvent::Error(error) => this.status = error.clone().into(),
|
BackendEvent::Error(error) => this.status = error.clone().into(),
|
||||||
_ => return,
|
_ => return,
|
||||||
}
|
}
|
||||||
@@ -30,7 +35,9 @@ impl Workspace {
|
|||||||
|
|
||||||
Self {
|
Self {
|
||||||
active_screen: repo_list.into(),
|
active_screen: repo_list.into(),
|
||||||
status: if connected {
|
status: if let Some((total, current)) = sync_progress {
|
||||||
|
format!("Syncing repositories... {current}/{total}").into()
|
||||||
|
} else if connected {
|
||||||
"Connected".into()
|
"Connected".into()
|
||||||
} else {
|
} else {
|
||||||
"Connecting...".into()
|
"Connecting...".into()
|
||||||
|
|||||||
Reference in New Issue
Block a user