Drive presence over USER_SYNC instead of a per-sync set_state relay

Previously a worker proxied a ReplicationPresenceSetState HTTP call to the
presence writer on every sync request, which made the writer's load scale with
the client sync rate and dominated its CPU on large deployments.

Instead, carry presence over the existing USER_SYNC replication command:

* UserSyncCommand gains a presence_state field. Workers edge-trigger it: a
  USER_SYNC is sent only when a device first starts syncing or changes the
  presence state it is syncing with, so a steady sync loop generates no presence
  traffic. The reconnect resync and an explicit REST set_state are handled so
  the writer's view stays correct.

* The writer applies the state from USER_SYNC (factored into _apply_sync_state).
  Combined with the previous commit, a device is treated as active for as long
  as it stays in the syncing set, so it neither idles nor goes offline while
  syncing without the writer being told about every sync.

* A new USER_SYNC_KEEPALIVE command lets a worker periodically remind the writer
  that its syncing users are still live, so a quiet worker isn't expired (which
  would otherwise mark all its users offline).

* The sync-driven half of the relay throttle (#19941) is superseded: syncs no
  longer relay a set_state at all, so the throttle now only covers the
  per-user-action activity bumps. Its entries are reset whenever the device's
  state on the writer may have changed (an explicit set_state, a change of
  sync state, or the device stopping syncing), so a bump that may un-idle the
  device is always relayed immediately.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PRR4mYoQDuv2k1FwEyreL6
This commit is contained in:
Erik Johnston
2026-07-10 14:13:07 +00:00
co-authored by Claude Opus 4.8
parent 55c7431ede
commit 2684275e1a
4 changed files with 542 additions and 215 deletions
+279 -125
View File
@@ -205,6 +205,12 @@ EXTERNAL_PROCESS_EXPIRY = 5 * 60 * 1000
# syncing.
UPDATE_SYNCING_USERS = Duration(seconds=10)
# How often a worker reminds the presence writer that its set of syncing users
# is still live. Must be comfortably shorter than EXTERNAL_PROCESS_EXPIRY,
# otherwise a worker whose users are all quietly syncing (and so sends no
# USER_SYNC commands) would be expired and its users marked offline.
SYNC_PRESENCE_KEEPALIVE_INTERVAL = Duration(minutes=1)
class BasePresenceHandler(abc.ABC):
"""Parts of the PresenceHandler that are shared between workers and presence
@@ -275,15 +281,16 @@ class BasePresenceHandler(abc.ABC):
@abc.abstractmethod
def get_currently_syncing_users_for_replication(
self,
) -> Iterable[tuple[str, str | None]]:
) -> Iterable[tuple[str, str | None, str]]:
"""Get an iterable of syncing users and devices on this worker, to send to the presence handler
This is called when a replication connection is established. It should return
a list of tuples of user ID & device ID, which are then sent as USER_SYNC commands
to inform the process handling presence about those users/devices.
a list of tuples of user ID, device ID & the presence state the device is
syncing with, which are then sent as USER_SYNC commands to inform the
process handling presence about those users/devices.
Returns:
An iterable of tuples of user ID and device ID.
An iterable of tuples of user ID, device ID and presence state.
"""
async def get_state(self, target_user: UserID) -> UserPresenceState:
@@ -376,6 +383,7 @@ class BasePresenceHandler(abc.ABC):
user_id: str,
device_id: str | None,
is_syncing: bool,
presence_state: str,
sync_time_msec: int,
) -> None:
"""Update the syncing users for an external process as a delta.
@@ -389,9 +397,20 @@ class BasePresenceHandler(abc.ABC):
user_id: The user who has started or stopped syncing
device_id: The user's device that has started or stopped syncing
is_syncing: Whether or not the user is now syncing
presence_state: The presence state the device is syncing with. Only
meaningful when the device is starting (or changing) its sync.
sync_time_msec: Time in ms when the user was last syncing
"""
async def update_external_syncs_keepalive( # noqa: B027 (no-op by design)
self, process_id: str
) -> None:
"""Refresh the "last updated" time for an external process, so that its
syncing users aren't expired while it is quiet.
This is a no-op when presence is handled by a different worker.
"""
async def update_external_syncs_clear( # noqa: B027 (no-op by design)
self, process_id: str
) -> None:
@@ -530,25 +549,29 @@ class WorkerPresenceHandler(BasePresenceHandler):
# syncing but we haven't notified the presence writer of that yet
self._user_devices_going_offline: dict[tuple[str, str | None], int] = {}
# How often to relay an unchanged sync-driven presence state to the
# presence writer. The relayed updates are what feed the writer's device
# last_sync_ts/last_active_ts timers, so this must sit comfortably below
# the timers it feeds — the (configurable) sync online timeout and
# last-active granularity — or users would flap offline / lose
# "currently active" between relays. We use 5/6 of the tighter of the
# two, i.e. the historic 25s at the default 30s sync online timeout.
self._sync_presence_relay_interval = (
# How often to relay repeated activity bumps for a device to the
# presence writer. The relayed bumps are what feed the writer's device
# last_active_ts timer, so this must sit comfortably below the
# (configurable) last-active granularity — or users would drop out of
# "currently active" between relays. We use 5/6 of the tighter of it
# and the sync online timeout, i.e. the historic 25s at the defaults.
self._bump_relay_interval = (
min(self._sync_online_timeout, self._last_active_granularity) * 5 // 6
)
# (user_id, device_id) -> (state, last_sent_ms) of the most recent
# sync-driven presence update we proxied to the presence writer. Used
# to suppress the per-sync-request set_state/bump calls, which are
# no-ops on the writer at finer granularity than its timers: while
# the state is unchanged there is no point relaying more than one
# update per relay interval. Entries older than the window are swept by
# `_sweep_last_sent_presence`.
self._last_sent_presence: dict[tuple[str, str | None], tuple[str, int]] = {}
# (user_id, device_id) -> the time we last relayed an activity bump for
# the device to the presence writer. Used to suppress the per-user-action
# bump calls, which are no-ops on the writer at finer granularity than
# its timers: one bump per relay interval is enough to keep them fed.
# Entries older than the window are swept by `_sweep_last_relayed_bumps`.
self._last_relayed_bump_ms: dict[tuple[str, str | None], int] = {}
# (user_id, device_id) -> presence state. The presence state we last told
# the presence writer each syncing device is in. Used to avoid re-sending
# a USER_SYNC when a sync repeats the same state, and to work out what to
# send on reconnect. An entry exists iff the writer believes the device is
# syncing.
self._user_device_last_sync_state: dict[tuple[str, str | None], str] = {}
self._bump_active_client = ReplicationBumpPresenceActiveTime.make_client(hs)
self._set_state_client = ReplicationPresenceSetState.make_client(hs)
@@ -556,7 +579,10 @@ class WorkerPresenceHandler(BasePresenceHandler):
if self._track_presence:
self.clock.looping_call(self.send_stop_syncing, UPDATE_SYNCING_USERS)
self.clock.looping_call(
self._sweep_last_sent_presence, Duration(minutes=30)
self.send_sync_keepalive, SYNC_PRESENCE_KEEPALIVE_INTERVAL
)
self.clock.looping_call(
self._sweep_last_relayed_bumps, Duration(minutes=30)
)
hs.register_async_shutdown_handler(
@@ -577,22 +603,57 @@ class WorkerPresenceHandler(BasePresenceHandler):
user_id: str,
device_id: str | None,
is_syncing: bool,
presence_state: str,
last_sync_ms: int,
) -> None:
if self._track_presence:
self.hs.get_replication_command_handler().send_user_sync(
self.instance_id, user_id, device_id, is_syncing, last_sync_ms
self.instance_id,
user_id,
device_id,
is_syncing,
presence_state,
last_sync_ms,
)
def mark_as_coming_online(self, user_id: str, device_id: str | None) -> None:
"""A user has started syncing. Send a UserSync to the presence writer,
unless they had recently stopped syncing.
def send_sync_keepalive(self) -> None:
"""Periodically remind the presence writer that our syncing users are
still live, so it doesn't expire us during quiet periods (when no device
is starting, stopping or changing its sync).
"""
going_offline = self._user_devices_going_offline.pop((user_id, device_id), None)
if not going_offline:
# Safe to skip because we haven't yet told the presence writer they
# were offline
self.send_user_sync(user_id, device_id, True, self.clock.time_msec())
if self._track_presence and any(
self._user_device_to_num_current_syncs.values()
):
self.hs.get_replication_command_handler().send_user_sync_keepalive(
self.instance_id
)
def mark_as_coming_online(
self, user_id: str, device_id: str | None, presence_state: str
) -> None:
"""A user's device has started syncing, or changed the presence state it
is syncing with. Tell the presence writer, unless it already believes the
device is syncing in this exact state (e.g. the device only briefly
stopped syncing, or the state is unchanged since the last sync).
"""
key = (user_id, device_id)
# Cancel any pending "stopped syncing" notification: the device is (still)
# syncing.
self._user_devices_going_offline.pop(key, None)
if self._user_device_last_sync_state.get(key) == presence_state:
# The writer already knows this device is syncing in this state, so
# there's nothing to send.
return
self._user_device_last_sync_state[key] = presence_state
# The device's state on the writer is about to change; drop its
# bump-throttle entry so a following bump (which may un-idle the new
# state) is relayed immediately.
self._last_relayed_bump_ms.pop(key, None)
self.send_user_sync(
user_id, device_id, True, presence_state, self.clock.time_msec()
)
def mark_as_going_offline(self, user_id: str, device_id: str | None) -> None:
"""A user has stopped syncing. We wait before notifying the presence
@@ -613,24 +674,29 @@ class WorkerPresenceHandler(BasePresenceHandler):
self._user_devices_going_offline.items()
):
if now - last_sync_ms > UPDATE_SYNCING_USERS.as_millis():
self._user_devices_going_offline.pop((user_id, device_id), None)
self.send_user_sync(user_id, device_id, False, last_sync_ms)
# Once the writer knows the device stopped syncing it may time
# the user out, so if the device comes back we must relay its
# state again rather than suppress it as a repeat.
self._last_sent_presence.pop((user_id, device_id), None)
key = (user_id, device_id)
self._user_devices_going_offline.pop(key, None)
# The writer no longer believes this device is syncing; forget the
# state we last sent so a fresh sync re-sends it, and drop the
# device's bump-throttle entry.
last_state = self._user_device_last_sync_state.pop(
key, PresenceState.OFFLINE
)
self._last_relayed_bump_ms.pop(key, None)
self.send_user_sync(user_id, device_id, False, last_state, last_sync_ms)
def _sweep_last_sent_presence(self) -> None:
"""Drop expired presence-throttling entries.
def _sweep_last_relayed_bumps(self) -> None:
"""Drop expired bump-throttling entries.
Entries should be dropped in `send_stop_syncing`, but we add a safety
net here to ensure that the dict deesn't grow unbounded.
Entries are dropped in `send_stop_syncing` when a device stops syncing,
but we add a safety net here to ensure that the dict doesn't grow
unbounded.
"""
now = self.clock.time_msec()
for key, (_, last_sent_ms) in list(self._last_sent_presence.items()):
if now - last_sent_ms >= self._sync_presence_relay_interval:
self._last_sent_presence.pop(key, None)
for key, last_relayed_ms in list(self._last_relayed_bump_ms.items()):
if now - last_relayed_ms >= self._bump_relay_interval:
self._last_relayed_bump_ms.pop(key, None)
async def user_syncing(
self,
@@ -647,23 +713,19 @@ class WorkerPresenceHandler(BasePresenceHandler):
if not affect_presence or not self._track_presence:
return _NullContextManager()
# Note that this causes last_active_ts to be incremented which is not
# what the spec wants. (This call is throttled in `set_state`: while
# the state is unchanged, only one update per relay interval is relayed
# to the presence writer.)
await self.set_state(
UserID.from_string(user_id),
device_id,
state={"presence": presence_state},
is_sync=True,
)
curr_sync = self._user_device_to_num_current_syncs.get((user_id, device_id), 0)
self._user_device_to_num_current_syncs[(user_id, device_id)] = curr_sync + 1
# If this is the first in-flight sync, notify replication
if self._user_device_to_num_current_syncs[(user_id, device_id)] == 1:
self.mark_as_coming_online(user_id, device_id)
# Tell the presence writer that this device is syncing (and with what
# presence state). This is a no-op if the writer already believes the
# device is syncing in this state, so an ongoing sync loop that keeps
# requesting the same state generates no replication traffic; only the
# first sync and any change of state are sent.
#
# Note that, unlike the old behaviour, we no longer proxy a set_state to
# the writer on every sync: the writer treats a syncing device as active
# for as long as it is syncing (see PresenceHandler.handle_timeout).
self.mark_as_coming_online(user_id, device_id, presence_state)
def _end() -> None:
# We check that the user_id is in user_to_num_current_syncs because
@@ -749,10 +811,18 @@ class WorkerPresenceHandler(BasePresenceHandler):
def get_currently_syncing_users_for_replication(
self,
) -> Iterable[tuple[str, str | None]]:
) -> Iterable[tuple[str, str | None, str]]:
return [
user_id_device_id
for user_id_device_id, count in self._user_device_to_num_current_syncs.items()
(
user_id,
device_id,
self._user_device_last_sync_state.get(
(user_id, device_id), PresenceState.ONLINE
),
)
for (user_id, device_id), count in (
self._user_device_to_num_current_syncs.items()
)
if count > 0
]
@@ -787,27 +857,16 @@ class WorkerPresenceHandler(BasePresenceHandler):
if not self._track_presence:
return
now = self.clock.time_msec()
if is_sync and not force_notify:
# Sync-driven updates arrive on every /sync request, which is far
# finer-grained than any of the writer's presence timers need:
# while the state is unchanged, relaying one update per relay
# interval is enough to keep them fed. State changes always go
# through immediately.
last_sent = self._last_sent_presence.get((user_id, device_id))
if last_sent is not None:
last_presence, last_sent_ms = last_sent
if (
presence == last_presence
and now - last_sent_ms < self._sync_presence_relay_interval
):
return
self._last_sent_presence[(user_id, device_id)] = (presence, now)
else:
# An explicit (non-sync) update doesn't refresh the writer's
# last_sync_ts, so it must not count as a recent relay: drop any
# entry so the next sync-driven update goes through.
self._last_sent_presence.pop((user_id, device_id), None)
if not is_sync:
# An explicit state change (e.g. a client PUT /presence) is applied
# on the writer directly below. Forget the state we last told the
# writer this device was syncing with, so that the device's next sync
# re-asserts its sync state (which may differ from the state just set).
self._user_device_last_sync_state.pop((user_id, device_id), None)
# Also forget the last relayed bump: the state just set may be one
# that a bump would un-idle, so the next bump must be relayed
# immediately rather than suppressed.
self._last_relayed_bump_ms.pop((user_id, device_id), None)
# Proxy request to instance that writes presence
await self._set_state_client(
@@ -832,20 +891,20 @@ class WorkerPresenceHandler(BasePresenceHandler):
user_id = user.to_string()
# A bump's only effects on the writer are updating last_active_ts and
# flipping an idle device back online. Going idle takes far longer
# than the relay window, so if we relayed an *online* update within
# the window the user cannot have gone idle since, and this bump is a
# no-op: skip it. Bumps after any other state (or an unknown one) go
# through immediately, as they may un-idle the device.
# flipping an idle device back online. If we relayed a bump within the
# relay window the writer's last_active_ts is already fresh and this
# bump is a no-op: skip it. (Entries are dropped whenever the device's
# state may have changed on the writer — an explicit set_state, a change
# of sync state, or the device stopping syncing — so a bump that might
# un-idle the device always goes through immediately.)
now = self.clock.time_msec()
last_sent = self._last_sent_presence.get((user_id, device_id))
last_relayed_ms = self._last_relayed_bump_ms.get((user_id, device_id))
if (
last_sent is not None
and last_sent[0] == PresenceState.ONLINE
and now - last_sent[1] < self._sync_presence_relay_interval
last_relayed_ms is not None
and now - last_relayed_ms < self._bump_relay_interval
):
return
self._last_sent_presence[(user_id, device_id)] = (PresenceState.ONLINE, now)
self._last_relayed_bump_ms[(user_id, device_id)] = now
# Proxy request to instance that writes presence
await self._bump_active_client(
@@ -932,6 +991,12 @@ class PresenceHandler(BasePresenceHandler):
# this is non zero a user will never go offline.
self._user_device_to_num_current_syncs: dict[tuple[str, str | None], int] = {}
# (user_id, device_id) -> presence state we last applied for a local sync.
# Used to edge-trigger: we only re-apply a device's sync state when it
# first starts syncing or changes the state it is syncing with, rather
# than on every sync request.
self._user_device_last_sync_state: dict[tuple[str, str | None], str] = {}
# Keeps track of the number of *ongoing* syncs on other processes.
#
# While any sync is ongoing on another process the user's device will never
@@ -1280,30 +1345,44 @@ class PresenceHandler(BasePresenceHandler):
if not affect_presence or not self._track_presence:
return _NullContextManager()
curr_sync = self._user_device_to_num_current_syncs.get((user_id, device_id), 0)
self._user_device_to_num_current_syncs[(user_id, device_id)] = curr_sync + 1
key = (user_id, device_id)
curr_sync = self._user_device_to_num_current_syncs.get(key, 0)
self._user_device_to_num_current_syncs[key] = curr_sync + 1
# Note that this causes last_active_ts to be incremented which is not
# what the spec wants.
await self.set_state(
UserID.from_string(user_id),
device_id,
state={"presence": presence_state},
is_sync=True,
)
# Edge-trigger the presence state: apply it when the device first starts
# syncing or when the requested state changes, rather than on every sync.
# While the device keeps syncing it is treated as active (see
# handle_timeout), so repeating the same state is a no-op.
if (
curr_sync == 0
or self._user_device_last_sync_state.get(key) != presence_state
):
self._user_device_last_sync_state[key] = presence_state
await self._apply_sync_state(user_id, device_id, presence_state)
async def _end() -> None:
try:
self._user_device_to_num_current_syncs[(user_id, device_id)] -= 1
# `_user_device_to_num_current_syncs` may have been cleared if we
# are shutting down.
if key not in self._user_device_to_num_current_syncs:
return
prev_state = await self.current_state_for_user(user_id)
await self._update_states(
[
prev_state.copy_and_replace(
last_user_sync_ts=self.clock.time_msec()
)
]
)
count = self._user_device_to_num_current_syncs[key] - 1
self._user_device_to_num_current_syncs[key] = count
# Only when the last in-flight sync for the device ends do we
# record it as no longer syncing (bumping last_user_sync_ts so it
# can time out) and forget the state we were tracking.
if count == 0:
self._user_device_last_sync_state.pop(key, None)
prev_state = await self.current_state_for_user(user_id)
await self._update_states(
[
prev_state.copy_and_replace(
last_user_sync_ts=self.clock.time_msec()
)
]
)
except Exception:
logger.exception("Error updating presence after sync")
@@ -1318,7 +1397,7 @@ class PresenceHandler(BasePresenceHandler):
def get_currently_syncing_users_for_replication(
self,
) -> Iterable[tuple[str, str | None]]:
) -> Iterable[tuple[str, str | None, str]]:
# since we are the process handling presence, there is nothing to do here.
return []
@@ -1328,6 +1407,7 @@ class PresenceHandler(BasePresenceHandler):
user_id: str,
device_id: str | None,
is_syncing: bool,
presence_state: str,
sync_time_msec: int,
) -> None:
"""Update the syncing users for an external process as a delta.
@@ -1339,30 +1419,32 @@ class PresenceHandler(BasePresenceHandler):
user_id: The user who has started or stopped syncing
device_id: The user's device that has started or stopped syncing
is_syncing: Whether or not the user is now syncing
presence_state: The presence state the device is syncing with. Only
used when the device is starting (or changing) its sync.
sync_time_msec: Time in ms when the user was last syncing
"""
async with self.external_sync_linearizer.queue(process_id):
prev_state = await self.current_state_for_user(user_id)
process_presence = self.external_process_to_current_syncs.setdefault(
process_id, set()
)
# USER_SYNC is sent when a user's device starts or stops syncing on
# a remote # process. (But only for the initial and last sync for that
# device.)
# A worker sends USER_SYNC "start" when a device first starts syncing
# and whenever the presence state it is syncing with changes, and
# "end" when the device stops syncing (see WorkerPresenceHandler).
#
# When a device *starts* syncing it also calls set_state(...) which
# will update the state, last_active_ts, and last_user_sync_ts.
# Simply ensure the user & device is tracked as syncing in this case.
# On start (or a state change) apply the presence state and mark the
# device as syncing. The device is then treated as active for as long
# as it remains in `process_presence` (see handle_timeout), so — unlike
# the old per-sync set_state relay — no further updates are needed
# until the state changes or the device stops.
#
# When a device *stops* syncing, update the last_user_sync_ts and mark
# them as no longer syncing. Note this doesn't quite match the
# monolith behaviour, which updates last_user_sync_ts at the end of
# every sync, not just the last in-flight sync.
if is_syncing and (user_id, device_id) not in process_presence:
# On stop, update last_user_sync_ts and mark the device as no longer
# syncing, so it can time out.
if is_syncing:
process_presence.add((user_id, device_id))
elif not is_syncing and (user_id, device_id) in process_presence:
await self._apply_sync_state(user_id, device_id, presence_state)
elif (user_id, device_id) in process_presence:
prev_state = await self.current_state_for_user(user_id)
devices = self._user_to_device_to_current_state.setdefault(user_id, {})
device_state = devices.setdefault(
device_id, UserDevicePresenceState.default(user_id, device_id)
@@ -1378,6 +1460,22 @@ class PresenceHandler(BasePresenceHandler):
self.external_process_last_updated_ms[process_id] = self.clock.time_msec()
async def update_external_syncs_keepalive(self, process_id: str) -> None:
"""Refresh the "last updated" time for an external process so that its
syncing users aren't expired during quiet periods.
Sent periodically by workers that have syncing users but haven't had any
starts/stops/state-changes to report (see
WorkerPresenceHandler.send_sync_keepalive).
"""
async with self.external_sync_linearizer.queue(process_id):
# Only refresh a process we're actually tracking; a keepalive for an
# unknown process (e.g. one we've already expired) is ignored.
if process_id in self.external_process_last_updated_ms:
self.external_process_last_updated_ms[process_id] = (
self.clock.time_msec()
)
async def update_external_syncs_clear(self, process_id: str) -> None:
"""Marks all users that had been marked as syncing by a given process
as offline.
@@ -1487,6 +1585,58 @@ class PresenceHandler(BasePresenceHandler):
).inc(len(updates))
await self._update_states(updates)
async def _apply_sync_state(
self, user_id: str, device_id: str | None, presence_state: str
) -> None:
"""Mark a device as actively syncing in the given presence state.
This is the equivalent of ``set_state(..., is_sync=True)`` but is driven
directly by USER_SYNC (rather than a per-sync ``set_state`` relay). It is
applied when a device first starts syncing and whenever the presence
state it is syncing with changes; the device is then treated as active
for as long as it keeps syncing (see ``handle_timeout``).
"""
presence = presence_state
if presence not in self.VALID_PRESENCE:
logger.warning(
"Ignoring invalid presence state %r from a sync for %s",
presence_state,
user_id,
)
presence = PresenceState.ONLINE
now = self.clock.time_msec()
prev_state = await self.current_state_for_user(user_id)
# Syncs do not override a previous presence of busy.
#
# TODO: This is a hack for lack of multi-device support. Unfortunately
# removing this requires coordination with clients.
if prev_state.state == PresenceState.BUSY:
presence = PresenceState.BUSY
# Update the device specific information.
devices = self._user_to_device_to_current_state.setdefault(user_id, {})
device_state = devices.setdefault(
device_id,
UserDevicePresenceState.default(user_id, device_id),
)
device_state.state = presence
device_state.last_active_ts = now
device_state.last_sync_ts = now
# Based on the state of each user's device calculate the new presence state.
presence = _combine_device_states(devices.values())
new_fields: JsonDict = {"state": presence, "last_user_sync_ts": now}
if presence == PresenceState.ONLINE or presence == PresenceState.BUSY:
new_fields["last_active_ts"] = now
await self._update_states(
[prev_state.copy_and_replace(**new_fields)],
syncing_user_ids={user_id},
)
async def set_state(
self,
target_user: UserID,
@@ -1520,6 +1670,10 @@ class PresenceHandler(BasePresenceHandler):
user_id = target_user.to_string()
now = self.clock.time_msec()
# Forget any state we were tracking for a sync on this device: an
# explicit state change means the next sync should re-assert its state.
self._user_device_last_sync_state.pop((user_id, device_id), None)
prev_state = await self.current_state_for_user(user_id)
# Syncs do not override a previous presence of busy.
+61 -6
View File
@@ -271,16 +271,26 @@ class UserSyncCommand(Command):
This is used by the process handling presence (typically the master) to
calculate who is online and who is not.
Includes a timestamp of when the last user sync was.
Includes the presence state the sync requested and a timestamp of when the
last user sync was.
Format::
USER_SYNC <instance_id> <user_id> <state> <last_sync_ms>
USER_SYNC <instance_id> <user_id> <device_id> <state> <presence_state> <last_sync_ms>
Where <state> is either "start" or "end"
Where <state> is either "start" or "end", and <presence_state> is the
presence state the sync requested (e.g. "online" or "unavailable"). The
presence state is only meaningful when <state> is "start".
"""
__slots__ = ["instance_id", "user_id", "device_id", "is_syncing", "last_sync_ms"]
__slots__ = [
"instance_id",
"user_id",
"device_id",
"is_syncing",
"presence_state",
"last_sync_ms",
]
NAME = "USER_SYNC"
@@ -290,18 +300,22 @@ class UserSyncCommand(Command):
user_id: str,
device_id: str | None,
is_syncing: bool,
presence_state: str,
last_sync_ms: int,
):
self.instance_id = instance_id
self.user_id = user_id
self.device_id = device_id
self.is_syncing = is_syncing
self.presence_state = presence_state
self.last_sync_ms = last_sync_ms
@classmethod
def from_line(cls: type["UserSyncCommand"], line: str) -> "UserSyncCommand":
device_id: str | None
instance_id, user_id, device_id, state, last_sync_ms = line.split(" ", 4)
instance_id, user_id, device_id, state, presence_state, last_sync_ms = (
line.split(" ", 5)
)
if device_id == "None":
device_id = None
@@ -309,7 +323,14 @@ class UserSyncCommand(Command):
if state not in ("start", "end"):
raise Exception("Invalid USER_SYNC state %r" % (state,))
return cls(instance_id, user_id, device_id, state == "start", int(last_sync_ms))
return cls(
instance_id,
user_id,
device_id,
state == "start",
presence_state,
int(last_sync_ms),
)
def to_line(self) -> str:
return " ".join(
@@ -318,6 +339,7 @@ class UserSyncCommand(Command):
self.user_id,
str(self.device_id),
"start" if self.is_syncing else "end",
self.presence_state,
str(self.last_sync_ms),
)
)
@@ -351,6 +373,37 @@ class ClearUserSyncsCommand(Command):
return self.instance_id
class UserSyncKeepaliveCommand(Command):
"""Sent periodically by a worker to let the presence writer know that its
set of syncing users is still live, even when no user has started or
stopped syncing recently.
Without this the presence writer would expire the process (see
``EXTERNAL_PROCESS_EXPIRY``) after a period of quiet and mark all of its
syncing users as offline.
Format::
USER_SYNC_KEEPALIVE <instance_id>
"""
__slots__ = ["instance_id"]
NAME = "USER_SYNC_KEEPALIVE"
def __init__(self, instance_id: str):
self.instance_id = instance_id
@classmethod
def from_line(
cls: type["UserSyncKeepaliveCommand"], line: str
) -> "UserSyncKeepaliveCommand":
return cls(line)
def to_line(self) -> str:
return self.instance_id
class FederationAckCommand(Command):
"""Sent by the client when it has processed up to a given point in the
federation stream. This allows the master to drop in-memory caches of the
@@ -530,6 +583,7 @@ _COMMANDS: tuple[type[Command], ...] = (
UserIpCommand,
RemoteServerUpCommand,
ClearUserSyncsCommand,
UserSyncKeepaliveCommand,
LockReleasedCommand,
NewActiveTaskCommand,
CancelTaskCommand,
@@ -556,6 +610,7 @@ VALID_CLIENT_COMMANDS = (
PingCommand.NAME,
UserSyncCommand.NAME,
ClearUserSyncsCommand.NAME,
UserSyncKeepaliveCommand.NAME,
FederationAckCommand.NAME,
UserIpCommand.NAME,
ErrorCommand.NAME,
+34 -3
View File
@@ -47,6 +47,7 @@ from synapse.replication.tcp.commands import (
ReplicateCommand,
UserIpCommand,
UserSyncCommand,
UserSyncKeepaliveCommand,
)
from synapse.replication.tcp.context import ClientContextFactory
from synapse.replication.tcp.protocol import IReplicationConnection
@@ -472,6 +473,7 @@ class ReplicationCommandHandler:
cmd.user_id,
cmd.device_id,
cmd.is_syncing,
cmd.presence_state,
cmd.last_sync_ms,
)
else:
@@ -485,6 +487,16 @@ class ReplicationCommandHandler:
else:
return None
def on_USER_SYNC_KEEPALIVE(
self, conn: IReplicationConnection, cmd: UserSyncKeepaliveCommand
) -> Awaitable[None] | None:
if self._is_presence_writer:
return self._presence_handler.update_external_syncs_keepalive(
cmd.instance_id
)
else:
return None
def on_FEDERATION_ACK(
self, conn: IReplicationConnection, cmd: FederationAckCommand
) -> None:
@@ -787,9 +799,16 @@ class ReplicationCommandHandler:
)
now = self._clock.time_msec()
for user_id, device_id in currently_syncing:
for user_id, device_id, presence_state in currently_syncing:
connection.send_command(
UserSyncCommand(self._instance_id, user_id, device_id, True, now)
UserSyncCommand(
self._instance_id,
user_id,
device_id,
True,
presence_state,
now,
)
)
def lost_connection(self, connection: IReplicationConnection) -> None:
@@ -846,13 +865,25 @@ class ReplicationCommandHandler:
user_id: str,
device_id: str | None,
is_syncing: bool,
presence_state: str,
last_sync_ms: int,
) -> None:
"""Poke the master that a user has started/stopped syncing."""
self.send_command(
UserSyncCommand(instance_id, user_id, device_id, is_syncing, last_sync_ms)
UserSyncCommand(
instance_id,
user_id,
device_id,
is_syncing,
presence_state,
last_sync_ms,
)
)
def send_user_sync_keepalive(self, instance_id: str) -> None:
"""Poke the master that this worker's syncing users are still live."""
self.send_command(UserSyncKeepaliveCommand(instance_id))
def send_user_ip(
self,
user_id: str,
+168 -81
View File
@@ -1255,6 +1255,48 @@ class PresenceHandlerTestCase(BaseMultiWorkerStreamTestCase):
state = self.get_success(self.presence_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.OFFLINE)
def test_external_process_keepalive(self) -> None:
"""A process that keeps sending keep-alives should not have its syncing
users expired, even if no user starts/stops/changes syncing.
This is what stops a worker whose users are all quietly syncing (and so
which sends no USER_SYNC commands under the pure-USER_SYNC model) from
being expired and having all of its users marked offline.
"""
user_id = f"@test:{self.hs.config.server.server_name}"
user_id_obj = UserID.from_string(user_id)
process_id = "proc-1"
# A user starts syncing on the external process.
self.get_success(
self.presence_handler.update_external_syncs_row(
process_id,
user_id,
"dev-1",
True,
PresenceState.ONLINE,
self.clock.time_msec(),
)
)
state = self.get_success(self.presence_handler.get_state(user_id_obj))
self.assertEqual(state.state, PresenceState.ONLINE)
# Advance past the expiry twice, sending a keep-alive each time in
# between. The user should stay online throughout.
for _ in range(3):
self.reactor.advance(EXTERNAL_PROCESS_EXPIRY * 0.75 / 1000)
self.get_success(
self.presence_handler.update_external_syncs_keepalive(process_id)
)
state = self.get_success(self.presence_handler.get_state(user_id_obj))
self.assertEqual(state.state, PresenceState.ONLINE)
# Once the keep-alives stop, the process is expired and the user times out.
self.reactor.advance(EXTERNAL_PROCESS_EXPIRY * 2 / 1000)
state = self.get_success(self.presence_handler.get_state(user_id_obj))
self.assertEqual(state.state, PresenceState.OFFLINE)
def test_user_goes_offline_by_timeout_status_msg_remain(self) -> None:
"""Test that if a user doesn't update the records for a while
users presence goes `OFFLINE` because of timeout and `status_msg` remains.
@@ -1585,12 +1627,15 @@ class PresenceHandlerTestCase(BaseMultiWorkerStreamTestCase):
Test the behaviour of multiple devices syncing at the same time.
Roughly the user's presence state should be set to the "highest" priority
of all the devices. When a device then goes offline its state should be
discarded and the next highest should win.
of all the devices.
Note that these tests use the idle timer (and don't close the syncs), it
is unlikely that a *single* sync would last this long, but is close enough
to continually syncing with that current state.
The syncs are never closed, i.e. each device keeps syncing with a fixed
state. A device that is actively syncing is treated as active and is *not*
auto-idled (a client that wants to appear idle says so via the sync's
presence state), so the combined state does not change as time passes.
This matches a client that continually re-syncs with the same state; it is
only once a device stops syncing that its state is discarded (see
test_set_presence_from_non_syncing_multi_device).
"""
user_id = f"@test:{self.hs.config.server.server_name}"
@@ -1629,6 +1674,12 @@ class PresenceHandlerTestCase(BaseMultiWorkerStreamTestCase):
)
)
# On a worker the sync notifies the presence writer over replication
# (a USER_SYNC command); pump the reactor so it is delivered and applied
# before we assert.
if test_with_workers:
self.reactor.pump([0.1])
# 4. Assert the expected presence state.
state = self.get_success(
self.presence_handler.get_state(UserID.from_string(user_id))
@@ -1859,6 +1910,12 @@ class PresenceHandlerTestCase(BaseMultiWorkerStreamTestCase):
)
)
# On a worker the syncs notify the presence writer over replication
# (USER_SYNC commands); pump the reactor so they are delivered and applied
# before we assert.
if test_with_workers:
self.reactor.pump([0.1])
# 3. Assert the expected presence state.
state = self.get_success(
self.presence_handler.get_state(UserID.from_string(user_id))
@@ -2481,9 +2538,10 @@ class PresenceJoinTestCase(unittest.HomeserverTestCase):
class WorkerPresenceThrottleTestCase(BaseMultiWorkerStreamTestCase):
"""Tests that sync workers suppress the per-sync-request presence updates
that the presence writer would discard anyway, while relaying genuine
state changes immediately."""
"""Tests that sync workers relay no per-sync-request presence updates to
the presence writer (presence is carried by edge-triggered USER_SYNC
commands instead), and that the per-user-action activity bumps are
throttled."""
def prepare(self, reactor: MemoryReactor, clock: Clock, hs: HomeServer) -> None:
self.user_id = f"@throttled:{hs.config.server.server_name}"
@@ -2525,133 +2583,162 @@ class WorkerPresenceThrottleTestCase(BaseMultiWorkerStreamTestCase):
presence.user_syncing(self.user_id, self.device_id, True, state),
)
def test_repeated_syncs_are_throttled(self) -> None:
def test_syncs_relay_no_set_state(self) -> None:
"""Syncs no longer proxy a set_state to the writer at all: presence is
carried by USER_SYNC, so the writer still sees the user come online."""
presence, set_state_calls, _ = self._make_sync_worker()
# Several syncs in quick succession only relay one set_state.
for _ in range(3):
self._sync(presence)
self.assertEqual(len(set_state_calls), 1)
self.assertEqual(len(set_state_calls), 0)
# The user did come online on the writer.
# The user did come online on the writer (via USER_SYNC). Pump the
# reactor so the replicated command is delivered and applied before we
# assert.
self.reactor.pump([0.1])
state = self.get_success(self.writer_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.ONLINE)
# Once the relay window has passed, the next sync relays again.
self.reactor.advance(presence._sync_presence_relay_interval / 1000 + 1)
self._sync(presence)
self.assertEqual(len(set_state_calls), 2)
def test_state_changes_are_relayed_immediately(self) -> None:
def test_sync_state_changes_reach_writer(self) -> None:
"""A change in the state a device is syncing with reaches the writer
immediately (via a fresh USER_SYNC), still without any set_state."""
presence, set_state_calls, _ = self._make_sync_worker()
# After each sync, pump the reactor so the replicated USER_SYNC command
# is delivered and applied before we assert on the writer's state.
self._sync(presence, PresenceState.ONLINE)
self._sync(presence, PresenceState.UNAVAILABLE)
self._sync(presence, PresenceState.ONLINE)
self.assertEqual(len(set_state_calls), 3)
self.reactor.pump([0.1])
state = self.get_success(self.writer_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.ONLINE)
self._sync(presence, PresenceState.UNAVAILABLE)
self.reactor.pump([0.1])
state = self.get_success(self.writer_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.UNAVAILABLE)
# A repeat of the current state within the window is suppressed.
self._sync(presence, PresenceState.ONLINE)
self.assertEqual(len(set_state_calls), 3)
self.reactor.pump([0.1])
state = self.get_success(self.writer_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.ONLINE)
self.assertEqual(len(set_state_calls), 0)
def test_resends_after_device_stops_syncing(self) -> None:
"""After a USER_SYNC stop is sent the writer may time the user out, so
a device that reconnects within the window must be relayed afresh."""
a device that reconnects must be told to the writer afresh."""
presence, set_state_calls, _ = self._make_sync_worker()
with self._sync(presence):
pass
self.assertEqual(len(set_state_calls), 1)
# Wait for the going-offline grace period to elapse: USER_SYNC stop is
# sent and the throttle entry evicted. The writer then times the user
# out to offline. Advance in steps (rather than one jump) so the
# replicated stop command is delivered before the writer's timeout
# loop fires (which only starts 30s after startup).
# sent. The writer then times the user out to offline. Advance in steps
# (rather than one jump) so the replicated stop command is delivered
# before the writer's timeout loop fires (which only starts 30s after
# startup).
for _ in range(4):
self.reactor.advance(12)
state = self.get_success(self.writer_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.OFFLINE)
# Reconnecting relays the state immediately, even though the presence
# value is unchanged from the last relayed one.
# Reconnecting sends a fresh USER_SYNC, even though the presence value
# is unchanged from the last one sent, and brings the user back online.
self._sync(presence)
self.assertEqual(len(set_state_calls), 2)
self.reactor.pump([0.1])
state = self.get_success(self.writer_handler.get_state(self.user_id_obj))
self.assertEqual(state.state, PresenceState.ONLINE)
self.assertEqual(len(set_state_calls), 0)
def test_bumps_are_throttled(self) -> None:
presence, set_state_calls, bump_calls = self._make_sync_worker()
presence, _, bump_calls = self._make_sync_worker()
# While we recently relayed an online state, bumps are suppressed.
self._sync(presence, PresenceState.ONLINE)
for _ in range(3):
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 0)
# After the window passes, a bump goes through (and then suppresses
# further bumps).
self.reactor.advance(presence._sync_presence_relay_interval / 1000 + 1)
for _ in range(2):
# The first bump is relayed; repeats within the window are suppressed.
for _ in range(3):
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 1)
# After the window passes, the next bump goes through (and then
# suppresses further bumps).
self.reactor.advance(presence._bump_relay_interval / 1000 + 1)
for _ in range(2):
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 2)
def test_explicit_set_state_always_relayed_and_resets(self) -> None:
"""An explicit (non-sync) set_state is always relayed, and resets the
throttle so the next sync-driven update is relayed afresh."""
presence, set_state_calls, _ = self._make_sync_worker()
self._sync(presence, PresenceState.ONLINE)
self.assertEqual(len(set_state_calls), 1)
# An explicit update of the same state within the window still goes
# through (it isn't sync-driven)...
self.get_success(
presence.set_state(
self.user_id_obj,
self.device_id,
{"presence": PresenceState.ONLINE},
),
)
self.assertEqual(len(set_state_calls), 2)
# ...and the following sync-driven update is relayed rather than
# suppressed, re-establishing the writer's sync timestamps.
self._sync(presence, PresenceState.ONLINE)
self.assertEqual(len(set_state_calls), 3)
def test_bump_after_non_online_state_goes_through(self) -> None:
bump throttle: the state just set may be one a bump would un-idle, so
the next bump must be relayed afresh."""
presence, set_state_calls, bump_calls = self._make_sync_worker()
# The user is unavailable; a bump may un-idle them so it must not be
# suppressed.
self._sync(presence, PresenceState.UNAVAILABLE)
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 1)
self.get_success(
presence.set_state(
self.user_id_obj,
self.device_id,
{"presence": PresenceState.UNAVAILABLE},
),
)
self.assertEqual(len(set_state_calls), 1)
# A bump within the window is relayed rather than suppressed, as it
# un-idles the state just set.
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 2)
def test_sync_state_change_resets_bump_throttle(self) -> None:
"""A change in the state a device is syncing with resets the bump
throttle, so a bump that may un-idle the new state goes through."""
presence, _, bump_calls = self._make_sync_worker()
self._sync(presence, PresenceState.ONLINE)
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 1)
# The device switches to syncing as unavailable; a following bump may
# un-idle it so it must not be suppressed.
self._sync(presence, PresenceState.UNAVAILABLE)
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 2)
@override_config({"presence": {"sync_online_timeout": "12s"}})
def test_relay_interval_scales_with_config(self) -> None:
"""The throttle window is derived from the configurable presence timers,
so it stays comfortably below a lowered sync online timeout rather than
being a hardcoded 25s (which would make users flap)."""
presence, set_state_calls, _ = self._make_sync_worker()
so it stays comfortably below them rather than being a hardcoded 25s."""
presence, _, bump_calls = self._make_sync_worker()
# 5/6 of min(12s sync online timeout, 60s default last-active
# granularity).
self.assertEqual(presence._sync_presence_relay_interval, 10 * 1000)
self.assertEqual(presence._bump_relay_interval, 10 * 1000)
# A repeat within the (now shorter) window is still suppressed...
self._sync(presence)
self._sync(presence)
self.assertEqual(len(set_state_calls), 1)
# A repeated bump within the (now shorter) window is still suppressed...
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 1)
# ...and once it passes, the next sync relays again.
self.reactor.advance(presence._sync_presence_relay_interval / 1000 + 1)
self._sync(presence)
self.assertEqual(len(set_state_calls), 2)
# ...and once it passes, the next bump relays again.
self.reactor.advance(presence._bump_relay_interval / 1000 + 1)
self.get_success(
presence.bump_presence_active_time(self.user_id_obj, self.device_id),
)
self.assertEqual(len(bump_calls), 2)