From 43c32ac9f73057c88d738e931b8c88ddbf06fb94 Mon Sep 17 00:00:00 2001 From: Alan Wu Date: Mon, 27 Jul 2026 21:27:11 -0400 Subject: [PATCH] Add dynamic config to flip skip persistence optimization (#11318) ## What changed? Add dynamic config to flip skip persistence optimization. ## Why? Allow backwards compatible rollout of skip persistence flag. ## How did you test it? - [X] built - [X] run locally and tested manually - [X] covered by existing tests - [X] added new unit test(s) - [ ] added new functional test(s) --- chasm/node_backend_mock.go | 8 ++ chasm/skip_persistence_test.go | 82 ++++++++++++++++--- chasm/tree.go | 8 +- chasm/tree_test.go | 12 ++- common/dynamicconfig/constants.go | 8 ++ service/history/chasm_engine_test.go | 2 + service/history/configs/config.go | 2 + .../history/workflow/mutable_state_impl.go | 5 ++ 8 files changed, 114 insertions(+), 13 deletions(-) diff --git a/chasm/node_backend_mock.go b/chasm/node_backend_mock.go index 0fcc7faba9..6b60243224 100644 --- a/chasm/node_backend_mock.go +++ b/chasm/node_backend_mock.go @@ -27,6 +27,7 @@ type MockNodeBackend struct { HandleGetCurrentVersion func() int64 HandleNextTransitionCount func() int64 HandleGetApproximatePersistedSize func() int + HandleChasmSkipPersistenceEnabled func() bool HandleCurrentVersionedTransition func() *persistencespb.VersionedTransition HandleGetWorkflowKey func() definition.WorkflowKey HandleUpdateWorkflowStateStatus func(state enumsspb.WorkflowExecutionState, status enumspb.WorkflowExecutionStatus) (bool, error) @@ -71,6 +72,13 @@ func (m *MockNodeBackend) GetApproximatePersistedSize() int { return 0 } +func (m *MockNodeBackend) ChasmSkipPersistenceEnabled() bool { + if m.HandleChasmSkipPersistenceEnabled != nil { + return m.HandleChasmSkipPersistenceEnabled() + } + return false +} + func (m *MockNodeBackend) GetCurrentVersion() int64 { if m.HandleGetCurrentVersion != nil { return m.HandleGetCurrentVersion() diff --git a/chasm/skip_persistence_test.go b/chasm/skip_persistence_test.go index 8857817d10..b793b32496 100644 --- a/chasm/skip_persistence_test.go +++ b/chasm/skip_persistence_test.go @@ -16,11 +16,40 @@ func (s *nodeSuite) minimalTestComponent() *TestComponent { } } +func (s *nodeSuite) newSkipPersistenceTestTree( + serializedNodes map[string]*persistencespb.ChasmNode, + enabled func() bool, +) *Node { + s.nodeBackend.HandleChasmSkipPersistenceEnabled = enabled + if len(serializedNodes) == 0 { + return NewEmptyTree( + s.registry, + s.timeSource, + s.nodeBackend, + s.nodePathEncoder, + s.logger, + s.metricsHandler, + ) + } + + root, err := NewTreeFromDB( + serializedNodes, + s.registry, + s.timeSource, + s.nodeBackend, + s.nodePathEncoder, + s.logger, + s.metricsHandler, + ) + s.NoError(err) + return root +} + func (s *nodeSuite) TestSkipPersistenceIfClean_NewNode() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 1 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode := NewEmptyTree(s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) + rootNode := s.newSkipPersistenceTestTree(nil, func() bool { return true }) s.NoError(rootNode.SetRootComponent(s.minimalTestComponent())) mutation, err := rootNode.CloseTransaction() @@ -34,7 +63,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_LoadedUnmodified() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 1 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode := NewEmptyTree(s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) + rootNode := s.newSkipPersistenceTestTree(nil, func() bool { return true }) s.NoError(rootNode.SetRootComponent(s.minimalTestComponent())) firstMutation, err := rootNode.CloseTransaction() @@ -49,8 +78,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_LoadedUnmodified() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 2 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode2, err := NewTreeFromDB(persistedNodes, s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) - s.NoError(err) + rootNode2 := s.newSkipPersistenceTestTree(persistedNodes, func() bool { return true }) ctx := NewMutableContext(context.Background(), rootNode2) _, err = rootNode2.Component(ctx, ComponentRef{}) @@ -71,7 +99,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_LoadedModified() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 1 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode := NewEmptyTree(s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) + rootNode := s.newSkipPersistenceTestTree(nil, func() bool { return true }) s.NoError(rootNode.SetRootComponent(s.minimalTestComponent())) firstMutation, err := rootNode.CloseTransaction() @@ -81,8 +109,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_LoadedModified() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 2 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode2, err := NewTreeFromDB(persistedNodes, s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) - s.NoError(err) + rootNode2 := s.newSkipPersistenceTestTree(persistedNodes, func() bool { return true }) ctx := NewMutableContext(context.Background(), rootNode2) component, err := rootNode2.Component(ctx, ComponentRef{}) @@ -105,7 +132,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_WithNewTask() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 1 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode := NewEmptyTree(s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) + rootNode := s.newSkipPersistenceTestTree(nil, func() bool { return true }) s.NoError(rootNode.SetRootComponent(s.minimalTestComponent())) firstMutation, err := rootNode.CloseTransaction() @@ -115,8 +142,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_WithNewTask() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 2 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode2, err := NewTreeFromDB(persistedNodes, s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) - s.NoError(err) + rootNode2 := s.newSkipPersistenceTestTree(persistedNodes, func() bool { return true }) ctx := NewMutableContext(context.Background(), rootNode2) component, err := rootNode2.Component(ctx, ComponentRef{}) @@ -138,7 +164,7 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_DeleteUnpersistedNode() { s.nodeBackend.HandleNextTransitionCount = func() int64 { return 1 } s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } - rootNode := NewEmptyTree(s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler) + rootNode := s.newSkipPersistenceTestTree(nil, func() bool { return true }) s.NoError(rootNode.SetRootComponent(&TestComponent{ ComponentData: &persistencespb.WorkflowExecutionState{RunId: "root"}, SubComponent1: NewComponentField(nil, &TestSubComponent1{ @@ -156,3 +182,37 @@ func (s *nodeSuite) TestSkipPersistenceIfClean_DeleteUnpersistedNode() { s.Contains(mutation.UpdatedNodes, "", "root still needs initial persistence") s.Empty(mutation.DeletedNodes, "node created and deleted before first persistence must not emit a tombstone") } + +func (s *nodeSuite) TestSkipPersistenceIfClean_DynamicConfig() { + s.nodeBackend.HandleNextTransitionCount = func() int64 { return 1 } + s.nodeBackend.HandleGetCurrentVersion = func() int64 { return 1 } + + rootNode := s.newSkipPersistenceTestTree(nil, func() bool { return false }) + s.NoError(rootNode.SetRootComponent(s.minimalTestComponent())) + + firstMutation, err := rootNode.CloseTransaction() + s.NoError(err) + persistedNodes := common.CloneProtoMap(firstMutation.UpdatedNodes) + + enabled := false + nextTransitionCount := int64(2) + s.nodeBackend.HandleNextTransitionCount = func() int64 { return nextTransitionCount } + rootNode2 := s.newSkipPersistenceTestTree(persistedNodes, func() bool { return enabled }) + + ctx := NewMutableContext(context.Background(), rootNode2) + _, err = rootNode2.Component(ctx, ComponentRef{}) + s.NoError(err) + mutation, err := rootNode2.CloseTransaction() + s.NoError(err) + s.Contains(mutation.UpdatedNodes, "", "disabled optimization must preserve pre-optimization persistence behavior") + s.Equal(int64(2), rootNode2.serializedNode.GetMetadata().GetLastUpdateVersionedTransition().GetTransitionCount()) + + enabled = true + nextTransitionCount = 3 + _, err = rootNode2.Component(ctx, ComponentRef{}) + s.NoError(err) + mutation, err = rootNode2.CloseTransaction() + s.NoError(err) + s.NotContains(mutation.UpdatedNodes, "", "enabled optimization must omit an unchanged node") + s.Equal(int64(2), rootNode2.serializedNode.GetMetadata().GetLastUpdateVersionedTransition().GetTransitionCount()) +} diff --git a/chasm/tree.go b/chasm/tree.go index 1e21802c33..35160aaaa3 100644 --- a/chasm/tree.go +++ b/chasm/tree.go @@ -208,6 +208,7 @@ type ( GetExecutionState() *persistencespb.WorkflowExecutionState GetExecutionInfo() *persistencespb.WorkflowExecutionInfo GetApproximatePersistedSize() int + ChasmSkipPersistenceEnabled() bool GetNamespaceEntry() *namespace.Namespace GetCurrentVersion() int64 NextTransitionCount() int64 @@ -439,6 +440,9 @@ func (n *Node) markSubtreeDirty() { } } +// clearAncestorNodeValues invalidates hydrated component ancestors after tree-structure changes. +// Replication must always call this for child-only mutations because the source may omit an +// unchanged parent component when its skip-persistence optimization is enabled. func (n *Node) clearAncestorNodeValues(parent *Node) { for node := parent; node != nil; node = node.parent { if node.serializedNode == nil || !node.isComponent() || node.value == nil { @@ -1933,6 +1937,7 @@ func (n *Node) closeTransactionForceUpdateVisibility( } func (n *Node) closeTransactionSerializeNodes() error { + skipPersistenceIfClean := n.backend.ChasmSkipPersistenceEnabled() for nodePath, node := range n.andAllChildren() { if node.valueState > valueStateNeedSerialize { return serviceerror.NewInternalf("invalid valueState for serializing: %v", node.valueState) @@ -1954,7 +1959,8 @@ func (n *Node) closeTransactionSerializeNodes() error { prevVersionedTransition := common.CloneProto( node.serializedNode.GetMetadata().GetLastUpdateVersionedTransition(), ) - skipIfClean := (node.isComponent() || node.isData() || node.isMap()) && + skipIfClean := skipPersistenceIfClean && + (node.isComponent() || node.isData() || node.isMap()) && prevVersionedTransition != nil && !node.hasNewTransactionSideEffects() var prevData *commonpb.DataBlob diff --git a/chasm/tree_test.go b/chasm/tree_test.go index 3daeefd322..1c80f30ce8 100644 --- a/chasm/tree_test.go +++ b/chasm/tree_test.go @@ -1155,7 +1155,16 @@ func (s *nodeSuite) TestApplyMutation_InvalidatesHydratedMapAncestors() { mutation NodesMutation, expected map[string]string, ) { - target, err := s.newTestTree(common.CloneProtoMap(persistedNodes)) + s.nodeBackend.HandleChasmSkipPersistenceEnabled = func() bool { return false } + target, err := NewTreeFromDB( + common.CloneProtoMap(persistedNodes), + s.registry, + s.timeSource, + s.nodeBackend, + s.nodePathEncoder, + s.logger, + s.metricsHandler, + ) s.NoError(err) component, err := target.Component(NewContext(context.Background(), target), ComponentRef{}) s.NoError(err) @@ -4621,6 +4630,7 @@ func (s *nodeSuite) TestAndAllChildren_PathIndependence() { func (s *nodeSuite) newTestTree( serializedNodes map[string]*persistencespb.ChasmNode, ) (*Node, error) { + s.nodeBackend.HandleChasmSkipPersistenceEnabled = func() bool { return true } if len(serializedNodes) == 0 { return NewEmptyTree(s.registry, s.timeSource, s.nodeBackend, s.nodePathEncoder, s.logger, s.metricsHandler), nil } diff --git a/common/dynamicconfig/constants.go b/common/dynamicconfig/constants.go index aa88c1f11b..a91837b57e 100644 --- a/common/dynamicconfig/constants.go +++ b/common/dynamicconfig/constants.go @@ -3101,6 +3101,14 @@ Requires service restart to take effect.`, true, "Use real chasm tree implementation instead of the noop one", ) + EnableCHASMSkipPersistence = NewNamespaceBoolSetting( + "history.enableCHASMSkipPersistence", + false, + `EnableCHASMSkipPersistence controls whether CHASM CloseTransaction omits nodes whose serialized data is unchanged. +This optimization should only be enabled after every cluster that may receive CHASM replication supports invalidating +hydrated ancestor components when applying child-node mutations.`, + ) + ChasmMaxInMemoryPureTasks = NewGlobalIntSetting( "history.chasmMaxInMemoryPureTasks", 32, diff --git a/service/history/chasm_engine_test.go b/service/history/chasm_engine_test.go index f3467a6b8e..71d61b0b6a 100644 --- a/service/history/chasm_engine_test.go +++ b/service/history/chasm_engine_test.go @@ -1071,6 +1071,8 @@ func (s *chasmEngineSuite) TestPollComponent_SetsContextMetadata() { // TestPollComponent_Success_Wait tests the waiting behavior of PollComponent. func (s *chasmEngineSuite) TestPollComponent_Success_Wait() { + s.config.EnableCHASMSkipPersistence = dynamicconfig.GetBoolPropertyFnFilteredByNamespace(true) + testCases := []struct { name string useEmptyRunID bool diff --git a/service/history/configs/config.go b/service/history/configs/config.go index e04dce32ac..c801f16ef5 100644 --- a/service/history/configs/config.go +++ b/service/history/configs/config.go @@ -75,6 +75,7 @@ type Config struct { MaxCallbacksPerExecution dynamicconfig.IntPropertyFnWithNamespaceFilter MaxCallbacksPerUpdateID dynamicconfig.IntPropertyFnWithNamespaceFilter EnableChasm dynamicconfig.BoolPropertyFnWithNamespaceFilter + EnableCHASMSkipPersistence dynamicconfig.BoolPropertyFnWithNamespaceFilter EnableChasmNexusWorkflowOperations dynamicconfig.BoolPropertyFnWithNamespaceFilter ChasmNexusWorkflowOperationsRolloutPercent dynamicconfig.IntPropertyFnWithNamespaceFilter EnableCHASMCallbacks dynamicconfig.BoolPropertyFnWithNamespaceFilter @@ -505,6 +506,7 @@ func NewConfig( MaxCallbacksPerExecution: callback.MaxPerExecution.Get(dc), MaxCallbacksPerUpdateID: dynamicconfig.MaxCallbacksPerUpdateID.Get(dc), EnableChasm: dynamicconfig.EnableChasm.Get(dc), + EnableCHASMSkipPersistence: dynamicconfig.EnableCHASMSkipPersistence.Get(dc), EnableChasmNexusWorkflowOperations: nexusoperation.EnableChasmWorkflowOperations.Get(dc), ChasmNexusWorkflowOperationsRolloutPercent: nexusoperation.ChasmWorkflowOperationsRolloutPercent.Get(dc), ChasmMaxInMemoryPureTasks: dynamicconfig.ChasmMaxInMemoryPureTasks.Get(dc), diff --git a/service/history/workflow/mutable_state_impl.go b/service/history/workflow/mutable_state_impl.go index 5369067dab..c67458c035 100644 --- a/service/history/workflow/mutable_state_impl.go +++ b/service/history/workflow/mutable_state_impl.go @@ -688,6 +688,11 @@ func (ms *MutableStateImpl) ChasmEnabled() bool { return !isNoop } +func (ms *MutableStateImpl) ChasmSkipPersistenceEnabled() bool { + return ms.config.EnableCHASMSkipPersistence != nil && + ms.config.EnableCHASMSkipPersistence(ms.GetNamespaceEntry().Name().String()) +} + // chasmCallbacksEnabled returns true if CHASM callbacks are enabled for this workflow. func (ms *MutableStateImpl) chasmCallbacksEnabled() bool { if !ms.ChasmEnabled() {