From d34db35a5b9c18653f54a8469f85c2b6003f95c4 Mon Sep 17 00:00:00 2001 From: Harsh Thakare Date: Thu, 23 Jul 2026 15:20:30 +0530 Subject: [PATCH 1/3] Fix disk state manager debounced writes --- reflex/istate/manager/disk.py | 15 ++++++++----- tests/units/test_state.py | 41 ++++++++++++++++++++++++++++++++++- 2 files changed, 49 insertions(+), 7 deletions(-) diff --git a/reflex/istate/manager/disk.py b/reflex/istate/manager/disk.py index 15dcc53f2ed..fe9f4e741f4 100644 --- a/reflex/istate/manager/disk.py +++ b/reflex/istate/manager/disk.py @@ -332,14 +332,17 @@ async def set_state( context: The state modification context. """ token = self._coerce_token(token) + is_base_state_token = isinstance(token, BaseStateToken) + if not is_base_state_token: + self.states[token.cache_key] = state if self._write_debounce_seconds > 0: # Deferred write to reduce disk IO overhead. - if token not in self._write_queue: - self._write_queue[token] = QueueItem( - token=token, - state=state, - timestamp=time.time(), - ) + queued_item = self._write_queue.get(token) + self._write_queue[token] = QueueItem[TOKEN_TYPE]( + token=token, + state=state, + timestamp=queued_item.timestamp if queued_item else time.time(), + ) else: # Immediate write to disk. await self.set_state_for_substate(token, state) diff --git a/tests/units/test_state.py b/tests/units/test_state.py index 112aa1e8d43..6af86b99f89 100644 --- a/tests/units/test_state.py +++ b/tests/units/test_state.py @@ -46,7 +46,7 @@ from reflex.istate.manager.disk import StateManagerDisk from reflex.istate.manager.memory import StateManagerMemory from reflex.istate.manager.redis import StateManagerRedis -from reflex.istate.manager.token import BaseStateToken +from reflex.istate.manager.token import BaseStateToken, StateToken from reflex.istate.proxy import StateProxy from reflex.state import ( BaseState, @@ -4432,6 +4432,45 @@ async def test_state_manager_disk_close_resets_write_queue_task(): assert state_manager._write_queue_task is None +@pytest.mark.asyncio +async def test_state_manager_disk_debounced_set_state_updates_non_base_state_cache( + tmp_path, monkeypatch +): + """Test that debounced non-BaseState writes are visible before disk flush.""" + monkeypatch.setattr(prerequisites, "get_states_dir", lambda: tmp_path) + state_manager = StateManagerDisk(_write_debounce_seconds=60) + token = StateToken(ident="client", cls=int) + + await state_manager.set_state(token, 1) + + assert await state_manager.get_state(token) == 1 + + await state_manager.close() + + +@pytest.mark.asyncio +async def test_state_manager_disk_debounced_set_state_flushes_latest_non_base_state( + tmp_path, monkeypatch +): + """Test that debounced non-BaseState writes flush the latest queued value.""" + monkeypatch.setattr(prerequisites, "get_states_dir", lambda: tmp_path) + state_manager = StateManagerDisk(_write_debounce_seconds=60) + token = StateToken(ident="client", cls=int) + + await state_manager.set_state(token, 1) + first_timestamp = state_manager._write_queue[token].timestamp + await state_manager.set_state(token, 2) + + assert state_manager._write_queue[token].timestamp == first_timestamp + + await state_manager.close() + + fresh_state_manager = StateManagerDisk(_write_debounce_seconds=0) + assert await fresh_state_manager.get_state(token) == 2 + + await fresh_state_manager.close() + + class Obj(Base): """A object containing a callable for testing fallback pickle.""" From b81153e92c86e08a093f11329ef16294ec03b42c Mon Sep 17 00:00:00 2001 From: Harsh Thakare Date: Thu, 23 Jul 2026 15:49:50 +0530 Subject: [PATCH 2/3] Avoid caching failed disk writes --- reflex/istate/manager/disk.py | 6 ++++-- tests/units/test_state.py | 24 ++++++++++++++++++++++++ 2 files changed, 28 insertions(+), 2 deletions(-) diff --git a/reflex/istate/manager/disk.py b/reflex/istate/manager/disk.py index fe9f4e741f4..2d45c3df827 100644 --- a/reflex/istate/manager/disk.py +++ b/reflex/istate/manager/disk.py @@ -333,10 +333,10 @@ async def set_state( """ token = self._coerce_token(token) is_base_state_token = isinstance(token, BaseStateToken) - if not is_base_state_token: - self.states[token.cache_key] = state if self._write_debounce_seconds > 0: # Deferred write to reduce disk IO overhead. + if not is_base_state_token: + self.states[token.cache_key] = state queued_item = self._write_queue.get(token) self._write_queue[token] = QueueItem[TOKEN_TYPE]( token=token, @@ -346,6 +346,8 @@ async def set_state( else: # Immediate write to disk. await self.set_state_for_substate(token, state) + if not is_base_state_token: + self.states[token.cache_key] = state # Ensure the processing task is scheduled to handle expirations and any deferred writes. await self._schedule_process_write_queue() diff --git a/tests/units/test_state.py b/tests/units/test_state.py index 6af86b99f89..8f062ca4c63 100644 --- a/tests/units/test_state.py +++ b/tests/units/test_state.py @@ -4471,6 +4471,30 @@ async def test_state_manager_disk_debounced_set_state_flushes_latest_non_base_st await fresh_state_manager.close() +@pytest.mark.asyncio +async def test_state_manager_disk_immediate_set_state_failure_keeps_previous_cache( + tmp_path, monkeypatch, mocker +): + """Test failed immediate writes do not cache an unpersisted value.""" + monkeypatch.setattr(prerequisites, "get_states_dir", lambda: tmp_path) + state_manager = StateManagerDisk(_write_debounce_seconds=0) + token = StateToken(ident="client", cls=int) + + await state_manager.set_state(token, 1) + mocker.patch.object( + state_manager, + "set_state_for_substate", + side_effect=RuntimeError("write failed"), + ) + + with pytest.raises(RuntimeError, match="write failed"): + await state_manager.set_state(token, 2) + + assert await state_manager.get_state(token) == 1 + + await state_manager.close() + + class Obj(Base): """A object containing a callable for testing fallback pickle.""" From 509427838a60beb885d5d918ff44f02009d0796e Mon Sep 17 00:00:00 2001 From: Harsh Thakare Date: Thu, 23 Jul 2026 15:52:36 +0530 Subject: [PATCH 3/3] Add changelog fragment for disk state fix --- news/6807.bugfix.md | 1 + 1 file changed, 1 insertion(+) create mode 100644 news/6807.bugfix.md diff --git a/news/6807.bugfix.md b/news/6807.bugfix.md new file mode 100644 index 00000000000..247a184e645 --- /dev/null +++ b/news/6807.bugfix.md @@ -0,0 +1 @@ +Fixed debounced disk state writes so non-`BaseState` values stay cache-coherent and flush the latest queued value.