Add upper-bound for presence_since().

Signed-off-by: Jason Volk <jason@zemos.net>
This commit is contained in:
Jason Volk
2025-07-26 04:59:00 +00:00
parent c6836e51b2
commit 63dfe8f7e3
5 changed files with 18 additions and 10 deletions

View File

@@ -19,6 +19,9 @@ pub(crate) enum PresenceCommand {
PresenceSince {
/// UNIX timestamp since (u64)
since: u64,
/// Upper-bound of since
to: Option<u64>,
},
}
@@ -34,11 +37,11 @@ pub(super) async fn process(subcommand: PresenceCommand, context: &Context<'_>)
write!(context, "Query completed in {query_time:?}:\n\n```rs\n{results:#?}\n```")
},
| PresenceCommand::PresenceSince { since } => {
| PresenceCommand::PresenceSince { since, to } => {
let timer = tokio::time::Instant::now();
let results: Vec<(_, _, _)> = services
.presence
.presence_since(since)
.presence_since(since, to)
.map(|(user_id, count, bytes)| (user_id.to_owned(), count, bytes.to_vec()))
.collect()
.await;

View File

@@ -288,7 +288,7 @@ pub(crate) async fn build_sync_events(
let presence_updates: OptionFuture<_> = services
.config
.allow_local_presence
.then(|| process_presence_updates(services, since, sender_user))
.then(|| process_presence_updates(services, since, next_batch, sender_user))
.into();
let account_data = services
@@ -391,11 +391,12 @@ pub(crate) async fn build_sync_events(
async fn process_presence_updates(
services: &Services,
since: u64,
next_batch: u64,
syncing_user: &UserId,
) -> PresenceUpdates {
services
.presence
.presence_since(since)
.presence_since(since, Some(next_batch))
.filter(|(user_id, ..)| {
services
.rooms

View File

@@ -154,13 +154,15 @@ impl Data {
pub(super) fn presence_since(
&self,
since: u64,
to: Option<u64>,
) -> impl Stream<Item = (&UserId, u64, &[u8])> + Send + '_ {
self.presenceid_presence
.raw_stream()
.ignore_err()
.ready_filter_map(move |(key, presence)| {
let (count, user_id) = presenceid_parse(key).ok()?;
(count > since).then_some((user_id, count, presence))
(count > since && to.is_none_or(|to| count <= to))
.then_some((user_id, count, presence))
})
}
}

View File

@@ -239,8 +239,9 @@ impl Service {
pub fn presence_since(
&self,
since: u64,
to: Option<u64>,
) -> impl Stream<Item = (&UserId, u64, &[u8])> + Send + '_ {
self.db.presence_since(since)
self.db.presence_since(since, to)
}
#[inline]

View File

@@ -605,14 +605,15 @@ impl Service {
since: (u64, u64),
max_edu_count: &AtomicU64,
) -> Option<EduBuf> {
let presence_since = self.services.presence.presence_since(since.0);
let presence_since = self
.services
.presence
.presence_since(since.0, Some(since.1));
pin_mut!(presence_since);
let mut presence_updates = HashMap::<OwnedUserId, PresenceUpdate>::new();
while let Some((user_id, count, presence_bytes)) = presence_since.next().await {
if count > since.1 {
break;
}
debug_assert!(count <= since.1, "exceeded upper-bound");
max_edu_count.fetch_max(count, Ordering::Relaxed);
if !self.services.globals.user_is_local(user_id) {