-
Notifications
You must be signed in to change notification settings - Fork 1.8k
Expand file tree
/
Copy pathconflict_resolver.go
More file actions
182 lines (157 loc) · 5.87 KB
/
Copy pathconflict_resolver.go
File metadata and controls
182 lines (157 loc) · 5.87 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
//go:generate mockgen -package $GOPACKAGE -source $GOFILE -destination conflict_resolver_mock.go
package ndc
import (
"context"
"go.temporal.io/api/serviceerror"
"go.temporal.io/server/common/definition"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/persistence/versionhistory"
"go.temporal.io/server/common/primitives/timestamp"
historyi "go.temporal.io/server/service/history/interfaces"
)
type (
ConflictResolver interface {
GetOrRebuildCurrentMutableState(
ctx context.Context,
branchIndex int32,
incomingVersion int64,
) (historyi.MutableState, bool, error)
GetOrRebuildMutableState(
ctx context.Context,
branchIndex int32,
) (historyi.MutableState, bool, error)
}
ConflictResolverImpl struct {
shard historyi.ShardContext
stateRebuilder StateRebuilder
context historyi.WorkflowContext
mutableState historyi.MutableState
logger log.Logger
}
)
var _ ConflictResolver = (*ConflictResolverImpl)(nil)
func NewConflictResolver(
shard historyi.ShardContext,
wfContext historyi.WorkflowContext,
mutableState historyi.MutableState,
logger log.Logger,
) *ConflictResolverImpl {
return &ConflictResolverImpl{
shard: shard,
stateRebuilder: NewStateRebuilder(shard, logger),
context: wfContext,
mutableState: mutableState,
logger: logger,
}
}
func (r *ConflictResolverImpl) GetOrRebuildCurrentMutableState(
ctx context.Context,
branchIndex int32,
incomingVersion int64,
) (historyi.MutableState, bool, error) {
versionHistories := r.mutableState.GetExecutionInfo().GetVersionHistories()
currentVersionHistoryIndex := versionHistories.GetCurrentVersionHistoryIndex()
currentVersionHistory, err := versionhistory.GetVersionHistory(versionHistories, currentVersionHistoryIndex)
if err != nil {
return nil, false, err
}
currentLastItem, err := versionhistory.GetLastVersionHistoryItem(currentVersionHistory)
if err != nil {
return nil, false, err
}
// mutable state does not need Rebuild
if incomingVersion < currentLastItem.GetVersion() {
return r.mutableState, false, nil
}
if incomingVersion == currentLastItem.GetVersion() && branchIndex != currentVersionHistoryIndex {
return nil, false, serviceerror.NewInvalidArgument("ConflictResolver encountered replication task version == current branch last write version")
}
// incomingVersion > currentLastItem.GetVersion()
return r.getOrRebuildMutableStateByIndex(ctx, branchIndex)
}
func (r *ConflictResolverImpl) GetOrRebuildMutableState(
ctx context.Context,
branchIndex int32,
) (historyi.MutableState, bool, error) {
return r.getOrRebuildMutableStateByIndex(ctx, branchIndex)
}
func (r *ConflictResolverImpl) getOrRebuildMutableStateByIndex(
ctx context.Context,
branchIndex int32,
) (historyi.MutableState, bool, error) {
versionHistories := r.mutableState.GetExecutionInfo().GetVersionHistories()
currentVersionHistoryIndex := versionHistories.GetCurrentVersionHistoryIndex()
// replication task to be applied to current branch
if branchIndex == currentVersionHistoryIndex {
return r.mutableState, false, nil
}
// task.getVersion() > currentLastItem
// incoming replication task, after application, will become the current branch
// (because higher version wins), we need to Rebuild the mutable state for that
rebuiltMutableState, err := r.rebuild(ctx, branchIndex)
if err != nil {
return nil, false, err
}
return rebuiltMutableState, true, nil
}
func (r *ConflictResolverImpl) rebuild(
ctx context.Context,
branchIndex int32,
) (historyi.MutableState, error) {
versionHistories := r.mutableState.GetExecutionInfo().GetVersionHistories()
replayVersionHistory, err := versionhistory.GetVersionHistory(versionHistories, branchIndex)
if err != nil {
return nil, err
}
lastItem, err := versionhistory.GetLastVersionHistoryItem(replayVersionHistory)
if err != nil {
return nil, err
}
executionInfo := r.mutableState.GetExecutionInfo()
executionState := r.mutableState.GetExecutionState()
workflowKey := definition.NewWorkflowKey(
executionInfo.NamespaceId,
executionInfo.WorkflowId,
executionState.RunId,
)
historySize := r.mutableState.GetHistorySize()
externalPayloadSize := r.mutableState.GetExternalPayloadSize()
externalPayloadCount := r.mutableState.GetExternalPayloadCount()
rebuildMutableState, _, err := r.stateRebuilder.Rebuild(
ctx,
timestamp.TimeValue(executionState.StartTime),
workflowKey,
replayVersionHistory.GetBranchToken(),
lastItem.GetEventId(),
new(lastItem.GetVersion()),
workflowKey,
replayVersionHistory.GetBranchToken(),
findStartRequestID(executionState),
)
if err != nil {
return nil, err
}
// after rebuilt verification
rebuildVersionHistories := rebuildMutableState.GetExecutionInfo().GetVersionHistories()
rebuildVersionHistory, err := versionhistory.GetCurrentVersionHistory(rebuildVersionHistories)
rebuildMutableState.GetExecutionInfo().PreviousTransitionHistory = r.mutableState.GetExecutionInfo().PreviousTransitionHistory
rebuildMutableState.GetExecutionInfo().LastTransitionHistoryBreakPoint = r.mutableState.GetExecutionInfo().LastTransitionHistoryBreakPoint
if err != nil {
return nil, err
}
if !rebuildVersionHistory.Equal(replayVersionHistory) {
return nil, serviceerror.NewInternal("ConflictResolver encounter mismatch version history after Rebuild")
}
// set the current branch index to target branch index
if err := versionhistory.SetCurrentVersionHistoryIndex(versionHistories, branchIndex); err != nil {
return nil, err
}
rebuildMutableState.GetExecutionInfo().VersionHistories = versionHistories
rebuildMutableState.AddHistorySize(historySize)
rebuildMutableState.AddExternalPayloadSize(externalPayloadSize)
rebuildMutableState.AddExternalPayloadCount(externalPayloadCount)
// set the update condition from original mutable state
rebuildMutableState.SetUpdateCondition(r.mutableState.GetUpdateCondition())
r.context.Clear()
return rebuildMutableState, nil
}