From 6421cceb21d7e6ff7c7e08a086aa01d978dbfddd Mon Sep 17 00:00:00 2001 From: joshvanl Date: Tue, 22 Sep 2026 15:10:08 -0300 Subject: [PATCH 1/3] Actors: invalidate the default state tracker after a reentrant save With reentrancy enabled, each dispatched method call gets its own state change tracker, but activation, reminders and timers run on the default tracker because no reentrancy id reaches them. A key read during activation stays cached there with change kind none forever, while method calls write the same key through their own trackers. A reminder callback that later reads that key is served the stale activation value. An app that skips its write because the value looks unchanged loses that write silently: nothing is logged anywhere, because no write is ever issued. Drop the default tracker's clean copies of keys written through a reentrancy-scoped tracker, so the next read reloads them from the runtime. Reported in dapr/dapr#10532, where a reminder callback's read-modify-write of an actor state key never persisted while the identical write from an ordinary method call did, and only with reentrancy enabled. Should be backported. Signed-off-by: joshvanl --- dapr/actor/runtime/state_manager.py | 11 +++++++++++ tests/actor/test_state_manager.py | 30 +++++++++++++++++++++++++++++ 2 files changed, 41 insertions(+) diff --git a/dapr/actor/runtime/state_manager.py b/dapr/actor/runtime/state_manager.py index c2882debb..03d10998b 100644 --- a/dapr/actor/runtime/state_manager.py +++ b/dapr/actor/runtime/state_manager.py @@ -267,6 +267,17 @@ async def save_state(self) -> None: ) for state_name in states_to_remove: state_change_tracker.pop(state_name, None) + if state_change_tracker is not self._default_state_change_tracker: + self._invalidate_default_tracker(state_changes) + + def _invalidate_default_tracker(self, state_changes: List[ActorStateChange]) -> None: + # Writes made through a reentrancy-scoped tracker are invisible to the default + # tracker, which activation, reminders and timers read from. Drop its clean copies + # of the written keys so the next read reloads them instead of serving stale data. + for change in state_changes: + metadata = self._default_state_change_tracker.get(change.state_name) + if metadata is not None and metadata.change_kind == StateChangeKind.none: + self._default_state_change_tracker.pop(change.state_name) def is_state_marked_for_remove(self, state_name: str) -> bool: state_change_tracker = self._get_contextual_state_tracker() diff --git a/tests/actor/test_state_manager.py b/tests/actor/test_state_manager.py index 5ed4e24e6..566e80457 100644 --- a/tests/actor/test_state_manager.py +++ b/tests/actor/test_state_manager.py @@ -20,6 +20,7 @@ from dapr.actor.id import ActorId from dapr.actor.runtime._type_information import ActorTypeInformation from dapr.actor.runtime.context import ActorRuntimeContext +from dapr.actor.runtime.reentrancy_context import reentrancy_ctx from dapr.actor.runtime.state_change import StateChangeKind from dapr.actor.runtime.state_manager import ActorStateManager, StateMetadata from dapr.serializers import DefaultJSONSerializer @@ -118,6 +119,35 @@ def test_get_state_for_removed_value(self): self.assertFalse(has_value) self.assertIsNone(val) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.get_state', + new=_async_mock(return_value=b'"value1"'), + ) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() + ) + def test_reentrant_save_invalidates_default_tracker(self): + state_manager = ActorStateManager(self._fake_actor) + + # A read outside any reentrancy context caches the value in the default tracker. + has_value, val = _run(state_manager.try_get_state('state1')) + self.assertTrue(has_value) + self.assertEqual('value1', val) + + # A reentrancy-scoped call writes the same key through its own tracker. + reentrancy_ctx.set('reentrancy-id') + state_manager.set_state_context('ctx1') + _run(state_manager.set_state('state1', 'value2')) + _run(state_manager.save_state()) + state_manager.set_state_context(None) + reentrancy_ctx.set(None) + + # The default tracker must reload the key instead of serving its stale copy. + self._fake_client.get_state.mock.return_value = b'"value2"' + has_value, val = _run(state_manager.try_get_state('state1')) + self.assertTrue(has_value) + self.assertEqual('value2', val) + @mock.patch('tests.actor.fake_client.FakeDaprActorClient.get_state', new=_async_mock()) def test_set_state_for_new_state(self): state_manager = ActorStateManager(self._fake_actor) From 83d0e033421d2abbbbb1bc51645172e45d493c06 Mon Sep 17 00:00:00 2001 From: Casper Nielsen Date: Fri, 25 Sep 2026 12:23:37 +0200 Subject: [PATCH 2/3] fix(actor): refresh the default state tracker in place after a reentrant save Instead of dropping the default tracker's clean copy of a key that a reentrant call saved, replace it with the saved value and ttl so the next read from activation, a reminder or a timer is served from cache rather than costing an extra state store read. Removed keys are still dropped and entries with pending changes are left alone. The cached value is passed through the state serializer first, so it has the same shape a fresh read would return (for example a tuple comes back as a list). This mirrors dapr/dotnet-sdk#1912. Signed-off-by: Casper Nielsen Co-Authored-By: Claude Opus 5.5 --- dapr/actor/runtime/_state_provider.py | 4 ++ dapr/actor/runtime/state_manager.py | 20 ++++-- tests/actor/test_state_manager.py | 90 ++++++++++++++++++++++----- 3 files changed, 95 insertions(+), 19 deletions(-) diff --git a/dapr/actor/runtime/_state_provider.py b/dapr/actor/runtime/_state_provider.py index eeb1e4995..86752d150 100644 --- a/dapr/actor/runtime/_state_provider.py +++ b/dapr/actor/runtime/_state_provider.py @@ -51,6 +51,10 @@ async def try_load_state( result = self._state_serializer.deserialize(raw_state_value, state_type) return (True, result) + def round_trip_state_value(self, value: Any) -> Any: + """Returns value as try_load_state would decode it after a save, without a store read.""" + return self._state_serializer.deserialize(self._state_serializer.serialize(value), object) + async def contains_state(self, actor_type: str, actor_id: str, state_name: str) -> bool: raw_state_value = await self._state_client.get_state(actor_type, actor_id, state_name) return (raw_state_value is not None) and len(raw_state_value) > 0 diff --git a/dapr/actor/runtime/state_manager.py b/dapr/actor/runtime/state_manager.py index 03d10998b..bcb4a0310 100644 --- a/dapr/actor/runtime/state_manager.py +++ b/dapr/actor/runtime/state_manager.py @@ -268,16 +268,26 @@ async def save_state(self) -> None: for state_name in states_to_remove: state_change_tracker.pop(state_name, None) if state_change_tracker is not self._default_state_change_tracker: - self._invalidate_default_tracker(state_changes) + self._refresh_default_tracker(state_changes) - def _invalidate_default_tracker(self, state_changes: List[ActorStateChange]) -> None: + def _refresh_default_tracker(self, state_changes: List[ActorStateChange]) -> None: # Writes made through a reentrancy-scoped tracker are invisible to the default - # tracker, which activation, reminders and timers read from. Drop its clean copies - # of the written keys so the next read reloads them instead of serving stale data. + # tracker, which activation, reminders and timers read from. Refresh its clean + # copies of the written keys in place, in the shape a fresh read would return, + # and drop removed keys. Entries with pending changes are left alone. + state_provider = self._actor.runtime_ctx.state_provider for change in state_changes: metadata = self._default_state_change_tracker.get(change.state_name) - if metadata is not None and metadata.change_kind == StateChangeKind.none: + if metadata is None or metadata.change_kind != StateChangeKind.none: + continue + if change.change_kind == StateChangeKind.remove: self._default_state_change_tracker.pop(change.state_name) + else: + self._default_state_change_tracker[change.state_name] = StateMetadata( + state_provider.round_trip_state_value(change.value), + StateChangeKind.none, + change.ttl_in_seconds, + ) def is_state_marked_for_remove(self, state_name: str) -> bool: state_change_tracker = self._get_contextual_state_tracker() diff --git a/tests/actor/test_state_manager.py b/tests/actor/test_state_manager.py index 566e80457..66e86f651 100644 --- a/tests/actor/test_state_manager.py +++ b/tests/actor/test_state_manager.py @@ -119,6 +119,19 @@ def test_get_state_for_removed_value(self): self.assertFalse(has_value) self.assertIsNone(val) + def _run_reentrant(self, state_manager, coro_fn): + # Runs coro_fn inside a reentrancy-scoped call, then saves its tracker. + token = reentrancy_ctx.set('reentrancy-id') + try: + state_manager.set_state_context('ctx1') + try: + _run(coro_fn()) + _run(state_manager.save_state()) + finally: + state_manager.set_state_context(None) + finally: + reentrancy_ctx.reset(token) + @mock.patch( 'tests.actor.fake_client.FakeDaprActorClient.get_state', new=_async_mock(return_value=b'"value1"'), @@ -126,27 +139,76 @@ def test_get_state_for_removed_value(self): @mock.patch( 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() ) - def test_reentrant_save_invalidates_default_tracker(self): + def test_reentrant_update_refreshes_default_tracker(self): state_manager = ActorStateManager(self._fake_actor) + _run(state_manager.try_get_state('state1')) + + self._run_reentrant(state_manager, lambda: state_manager.set_state_ttl('state1', 'v2', 60)) - # A read outside any reentrancy context caches the value in the default tracker. + # The default read is served from the refreshed entry, not the state store. + calls = self._fake_client.get_state.mock.call_count has_value, val = _run(state_manager.try_get_state('state1')) self.assertTrue(has_value) - self.assertEqual('value1', val) + self.assertEqual('v2', val) + self.assertEqual(calls, self._fake_client.get_state.mock.call_count) + state = state_manager._default_state_change_tracker['state1'] + self.assertEqual(StateChangeKind.none, state.change_kind) + self.assertEqual(60, state.ttl_in_seconds) - # A reentrancy-scoped call writes the same key through its own tracker. - reentrancy_ctx.set('reentrancy-id') - state_manager.set_state_context('ctx1') - _run(state_manager.set_state('state1', 'value2')) - _run(state_manager.save_state()) - state_manager.set_state_context(None) - reentrancy_ctx.set(None) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.get_state', + new=_async_mock(return_value=b'"value1"'), + ) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() + ) + def test_reentrant_remove_evicts_default_tracker(self): + state_manager = ActorStateManager(self._fake_actor) + _run(state_manager.try_get_state('state1')) - # The default tracker must reload the key instead of serving its stale copy. - self._fake_client.get_state.mock.return_value = b'"value2"' + self._run_reentrant(state_manager, lambda: state_manager.remove_state('state1')) + + self.assertNotIn('state1', state_manager._default_state_change_tracker) + self._fake_client.get_state.mock.return_value = None has_value, val = _run(state_manager.try_get_state('state1')) - self.assertTrue(has_value) - self.assertEqual('value2', val) + self.assertFalse(has_value) + self.assertIsNone(val) + + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.get_state', + new=_async_mock(return_value=b'"value1"'), + ) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() + ) + def test_reentrant_save_keeps_dirty_default_entry(self): + state_manager = ActorStateManager(self._fake_actor) + _run(state_manager.set_state('state1', 'pending')) + + self._run_reentrant(state_manager, lambda: state_manager.set_state('state1', 'v2')) + + state = state_manager._default_state_change_tracker['state1'] + self.assertEqual('pending', state.value) + self.assertEqual(StateChangeKind.update, state.change_kind) + + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.get_state', + new=_async_mock(return_value=b'[1, 2]'), + ) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() + ) + def test_reentrant_refresh_matches_fresh_read(self): + state_manager = ActorStateManager(self._fake_actor) + _run(state_manager.try_get_state('state1')) + + self._run_reentrant(state_manager, lambda: state_manager.set_state('state1', (3, 4))) + + _, refreshed = _run(state_manager.try_get_state('state1')) + self._fake_client.get_state.mock.return_value = self._serializer.serialize((3, 4)) + _, fresh = _run(ActorStateManager(self._fake_actor).try_get_state('state1')) + self.assertEqual([3, 4], fresh) + self.assertEqual(fresh, refreshed) @mock.patch('tests.actor.fake_client.FakeDaprActorClient.get_state', new=_async_mock()) def test_set_state_for_new_state(self): From c3b3387af98499ccff489e2aec4c2d82e10f6f99 Mon Sep 17 00:00:00 2001 From: Casper Nielsen Date: Fri, 25 Sep 2026 12:26:09 +0200 Subject: [PATCH 3/3] fix(actor): drop the default entry when a reentrant refresh cannot match a fresh read Two cases fall back to #1227's eviction instead of an in-place refresh: a saved None value, which the state provider leaves out of the write, and a state serializer that fails to decode the value after the save has already committed. The save no longer raises after a successful write. Signed-off-by: Casper Nielsen Co-Authored-By: Claude Opus 5.5 --- dapr/actor/runtime/state_manager.py | 19 ++++++++----- tests/actor/test_state_manager.py | 42 +++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+), 7 deletions(-) diff --git a/dapr/actor/runtime/state_manager.py b/dapr/actor/runtime/state_manager.py index bcb4a0310..37b1a6d3e 100644 --- a/dapr/actor/runtime/state_manager.py +++ b/dapr/actor/runtime/state_manager.py @@ -280,14 +280,19 @@ def _refresh_default_tracker(self, state_changes: List[ActorStateChange]) -> Non metadata = self._default_state_change_tracker.get(change.state_name) if metadata is None or metadata.change_kind != StateChangeKind.none: continue - if change.change_kind == StateChangeKind.remove: + # A None value is not written to the store, so let the next read reload it. + if change.change_kind == StateChangeKind.remove or change.value is None: self._default_state_change_tracker.pop(change.state_name) - else: - self._default_state_change_tracker[change.state_name] = StateMetadata( - state_provider.round_trip_state_value(change.value), - StateChangeKind.none, - change.ttl_in_seconds, - ) + continue + try: + value = state_provider.round_trip_state_value(change.value) + except Exception: + # The save has already committed; fall back to reloading on the next read. + self._default_state_change_tracker.pop(change.state_name) + continue + self._default_state_change_tracker[change.state_name] = StateMetadata( + value, StateChangeKind.none, change.ttl_in_seconds + ) def is_state_marked_for_remove(self, state_name: str) -> bool: state_change_tracker = self._get_contextual_state_tracker() diff --git a/tests/actor/test_state_manager.py b/tests/actor/test_state_manager.py index 66e86f651..658959b6a 100644 --- a/tests/actor/test_state_manager.py +++ b/tests/actor/test_state_manager.py @@ -210,6 +210,48 @@ def test_reentrant_refresh_matches_fresh_read(self): self.assertEqual([3, 4], fresh) self.assertEqual(fresh, refreshed) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.get_state', + new=_async_mock(return_value=b'"value1"'), + ) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() + ) + def test_reentrant_save_of_none_evicts_default_tracker(self): + state_manager = ActorStateManager(self._fake_actor) + _run(state_manager.try_get_state('state1')) + + self._run_reentrant(state_manager, lambda: state_manager.set_state('state1', None)) + + self.assertNotIn('state1', state_manager._default_state_change_tracker) + + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.get_state', + new=_async_mock(return_value=b'"value1"'), + ) + @mock.patch( + 'tests.actor.fake_client.FakeDaprActorClient.save_state_transactionally', new=_async_mock() + ) + def test_reentrant_refresh_failure_evicts_default_tracker(self): + state_manager = ActorStateManager(self._fake_actor) + _run(state_manager.try_get_state('state1')) + _run(state_manager.try_get_state('state2')) + + async def set_both(): + await state_manager.set_state('state1', 'v2') + await state_manager.set_state('state2', 'v2') + + # The save must not raise after the write has committed, and every key is handled. + with mock.patch.object( + self._runtime_ctx.state_provider, + 'round_trip_state_value', + side_effect=ValueError('cannot decode'), + ): + self._run_reentrant(state_manager, set_both) + + self.assertNotIn('state1', state_manager._default_state_change_tracker) + self.assertNotIn('state2', state_manager._default_state_change_tracker) + @mock.patch('tests.actor.fake_client.FakeDaprActorClient.get_state', new=_async_mock()) def test_set_state_for_new_state(self): state_manager = ActorStateManager(self._fake_actor)