From 37441464f2ca0c3b6d1c690a2e4c34c7044d0d2b Mon Sep 17 00:00:00 2001 From: Erik Johnston Date: Thu, 9 Jan 2025 14:10:25 +0000 Subject: [PATCH] More move --- synapse/storage/databases/main/state.py | 63 ---------------------- synapse/storage/databases/state/store.py | 68 ++++++++++++++++++++++++ 2 files changed, 68 insertions(+), 63 deletions(-) diff --git a/synapse/storage/databases/main/state.py b/synapse/storage/databases/main/state.py index d785cea9a3..788f7d1e32 100644 --- a/synapse/storage/databases/main/state.py +++ b/synapse/storage/databases/main/state.py @@ -48,7 +48,6 @@ from synapse.api.room_versions import KNOWN_ROOM_VERSIONS, RoomVersion from synapse.events import EventBase from synapse.events.snapshot import EventContext from synapse.logging.opentracing import trace -from synapse.metrics.background_process_metrics import wrap_as_background_process from synapse.replication.tcp.streams import UnPartialStatedEventStream from synapse.replication.tcp.streams.partial_state import UnPartialStatedEventStreamRow from synapse.storage._base import SQLBaseStore @@ -119,68 +118,6 @@ class StateGroupWorkerStore(EventsWorkerStore, SQLBaseStore): super().__init__(database, db_conn, hs) self._instance_name: str = hs.get_instance_name() - if hs.config.worker.run_background_tasks: - self._clock.looping_call_now(self._advance_state_epoch, 2 * 60 * 1000) - - @wrap_as_background_process("_advance_state_epoch") - async def _advance_state_epoch(self) -> None: - """Advances the state epoch, checking that we haven't advanced it too - recently. - """ - - now = self._clock.time_msec() - update_if_before_ts = now - 10 * 60 * 1000 - - def advance_state_epoch_txn(txn: LoggingTransaction) -> None: - sql = """ - UPDATE state_epoch - SET state_epoch = state_epoch + 1, updated_ts = ? - WHERE updated_ts <= ? - """ - txn.execute( - sql, - ( - now, - update_if_before_ts, - ), - ) - - await self.db_pool.runInteraction( - "_advance_state_epoch", advance_state_epoch_txn, db_autocommit=True - ) - - @cached() - async def is_state_group_pending_deletion_before( - self, state_epoch: int, state_group: int - ) -> bool: - """Check if a state group is marked as pending deletion in a previous - epoch, but does not check the current epoch.""" - - def is_state_group_pending_deletion_before_txn(txn: LoggingTransaction) -> bool: - sql = """ - SELECT 1 FROM state_groups_pending_deletion - WHERE state_epoch < ? AND state_group = ? - """ - txn.execute(sql, (state_epoch, state_group)) - - return txn.fetchone() is not None - - return await self.db_pool.runInteraction( - "is_state_group_pending_deletion_before", - is_state_group_pending_deletion_before_txn, - ) - - async def mark_state_group_as_used(self, state_group: int) -> None: - """Mark that a given state group is used""" - - # TODO: Also assert that the state group hasn't advanced too much - - await self.db_pool.simple_delete( - table="state_groups_pending_deletion", - keyvalues={"state_group": state_group}, - desc="mark_state_group_as_used", - ) - def process_replication_rows( self, stream_name: str, diff --git a/synapse/storage/databases/state/store.py b/synapse/storage/databases/state/store.py index 9944f90015..bfb03dede6 100644 --- a/synapse/storage/databases/state/store.py +++ b/synapse/storage/databases/state/store.py @@ -38,6 +38,7 @@ from synapse.api.constants import EventTypes from synapse.events import EventBase from synapse.events.snapshot import UnpersistedEventContext, UnpersistedEventContextBase from synapse.logging.opentracing import tag_args, trace +from synapse.metrics.background_process_metrics import wrap_as_background_process from synapse.storage._base import SQLBaseStore from synapse.storage.database import ( DatabasePool, @@ -140,6 +141,68 @@ class StateGroupDataStore(StateBackgroundUpdateStore, SQLBaseStore): id_column="id", ) + if hs.config.worker.run_background_tasks: + self._clock.looping_call_now(self._advance_state_epoch, 2 * 60 * 1000) + + @wrap_as_background_process("_advance_state_epoch") + async def _advance_state_epoch(self) -> None: + """Advances the state epoch, checking that we haven't advanced it too + recently. + """ + + now = self._clock.time_msec() + update_if_before_ts = now - 10 * 60 * 1000 + + def advance_state_epoch_txn(txn: LoggingTransaction) -> None: + sql = """ + UPDATE state_epoch + SET state_epoch = state_epoch + 1, updated_ts = ? + WHERE updated_ts <= ? + """ + txn.execute( + sql, + ( + now, + update_if_before_ts, + ), + ) + + await self.db_pool.runInteraction( + "_advance_state_epoch", advance_state_epoch_txn, db_autocommit=True + ) + + @cached() + async def is_state_group_pending_deletion_before( + self, state_epoch: int, state_group: int + ) -> bool: + """Check if a state group is marked as pending deletion in a previous + epoch, but does not check the current epoch.""" + + def is_state_group_pending_deletion_before_txn(txn: LoggingTransaction) -> bool: + sql = """ + SELECT 1 FROM state_groups_pending_deletion + WHERE state_epoch < ? AND state_group = ? + """ + txn.execute(sql, (state_epoch, state_group)) + + return txn.fetchone() is not None + + return await self.db_pool.runInteraction( + "is_state_group_pending_deletion_before", + is_state_group_pending_deletion_before_txn, + ) + + async def mark_state_group_as_used(self, state_group: int) -> None: + """Mark that a given state group is used""" + + # TODO: Also assert that the state group hasn't advanced too much + + await self.db_pool.simple_delete( + table="state_groups_pending_deletion", + keyvalues={"state_group": state_group}, + desc="mark_state_group_as_used", + ) + @cached(max_entries=10000, iterable=True) async def get_state_group_delta(self, state_group: int) -> _GetStateGroupDelta: """Given a state group try to return a previous group and a delta between @@ -467,6 +530,9 @@ class StateGroupDataStore(StateBackgroundUpdateStore, SQLBaseStore): Returns: A list of state groups """ + + # TODO: Check state epochs + is_in_db = self.db_pool.simple_select_one_onecol_txn( txn, table="state_groups", @@ -601,6 +667,8 @@ class StateGroupDataStore(StateBackgroundUpdateStore, SQLBaseStore): The state group if successfully created, or None if the state needs to be persisted as a full state. """ + + # TODO: Check state epoch. is_in_db = self.db_pool.simple_select_one_onecol_txn( txn, table="state_groups",