Repository navigation
Conversation
yycptt
left a comment
There was a problem hiding this comment.
I'd consider introduce a NotExist condition to cassandra's executionCASCondition so we can be more explicit about what we are checking (and the error message return).
(I haven't check today's SQL implementation)
When a WorkflowConditionFailedError is generated due to ^ condition, it can also be higher priority than CurrentWorkflowConditionFailedError, which allow us to dedup the request without attempting a second write.
Maybe we need a new condition failed error type here, since we need other information like status as well.
| } | ||
|
|
||
| // uuidSpaceRunID is a randomly generated UUIDv5 namespace for derived run IDs | ||
| var uuidSpaceRunID = uuid.MustParse("e96ecb9a-482c-4f69-9b66-a40bc6b50553") |
There was a problem hiding this comment.
we probably want to move the following implementation to a common place. we need it for chasm executions creation as well.
| s.request.StartRequest.GetWorkflowId(), | ||
| chasm.WorkflowArchetypeID, | ||
| s.request.StartRequest.GetRequestId(), | ||
| s.shardContext.GetClusterMetadata().GetClusterID(), |
There was a problem hiding this comment.
I feel we need some check in conflict solution logic to detect the same runID but different start event case and error out.
It only happens when a run got deleted from a source cluster and a new run is started with the same runID and the previous got replicated back from a remote cluster, so very rare but at least deletion + replicating back is some thing happened before. Can be done as a follow up.
| //go:generate mockgen -package $GOPACKAGE -source $GOFILE -destination transaction_mock.go | ||
| type ( | ||
| Transaction interface { | ||
| CreateWorkflowExecution( |
There was a problem hiding this comment.
shall we have the same flag here? Looks like this is used by resetWorkflowExecution and we can have the same stronger dedup there.
Also related to: #12042
| workflow.MutableStateFailoverVersion(resetWorkflow.GetMutableState()), | ||
| resetWorkflowSnapshot, | ||
| resetWorkflowEventsSeq, | ||
| // verifyRunIDUniqueness: reset must keep a random run ID; a derived one would collide with the base run. |
There was a problem hiding this comment.
hmm why? Reset operation has it's own requestID and can derive based on that?
| reactivationSignaler, | ||
| uws.workflowLeaseCallback(ctx), | ||
| // allowDerivedRunID: false. The held lease would self-deadlock on terminate-existing with a derived run ID. | ||
| // TODO: support run-ID dedup for Update-with-Start and drop this parameter. |
| return resp, outcome, conflictErr | ||
| } | ||
| return nil, StartErr, err | ||
| return s.runIDDedupResponse(ctx, creationParams.runID, err) |
There was a problem hiding this comment.
when will we get to here for the dup runID case?
There was a problem hiding this comment.
oh current run (run2) get deleted and start request for run1 get retried?
| var currentWorkflowConditionFailedError *persistence.CurrentWorkflowConditionFailedError | ||
| if errors.As(err, ¤tWorkflowConditionFailedError) && len(currentWorkflowConditionFailedError.RunID) > 0 { | ||
| // The history and mutable state generated above will be deleted by a background process. | ||
| resp, outcome, conflictErr := s.handleConflict(ctx, creationParams, currentWorkflowConditionFailedError) |
There was a problem hiding this comment.
CurrentWorkflowConditionFailedError has higher priority than WorkflowConditionFailedError, does that mean we will need two write operations to fail the dup request?
| return &historyservice.StartWorkflowExecutionResponse{ | ||
| RunId: attemptedRunID, | ||
| FirstExecutionRunId: attemptedRunID, | ||
| Started: true, |
There was a problem hiding this comment.
false, because no run actually get started?
| return nil, StartErr, err | ||
| } | ||
|
|
||
| info, loadErr := s.getMutableStateInfo(ctx, attemptedRunID) |
There was a problem hiding this comment.
this is another read operation? can we return the status as part of the error from persistence?
|
|
||
| return &historyservice.StartWorkflowExecutionResponse{ | ||
| RunId: attemptedRunID, | ||
| FirstExecutionRunId: attemptedRunID, |
There was a problem hiding this comment.
nit: I assume this will always equal to info.firstExecutionRunID?
What changed?
StartWorkflowExecution derives the new run ID from the request instead of generating a random one: UUIDv5 of (namespace ID, workflow ID, archetype ID, request ID, cluster ID). A retried start with the same request ID then produces the same run ID; the store rejects it as a duplicate, and the server returns the existing run as a successful start.
Adds:
DeriveRunIDand the dedup response in the Starter; all three create paths (brand new, create as current, terminate existing) map a conflict on the attempted run ID to a dedup.VerifyRunIDUniquenesson the create/update persistence requests, for stores that do not enforce run ID uniqueness through their primary key.WorkflowConditionFailedError.RunID, populated by SQL and Cassandra, so a conflict on the new run can be told apart from a version conflict on the current run.UpdateWorkflowExecutionnow surfaces asWorkflowConditionFailedError, matching SQL.allowDerivedRunID=false).history.enableCrossRunRequestIDDedupdynamic config (default off),Out of scope (will be addressed in follow-ups)
allowDerivedRunID=false).Why?
Today a start is deduplicated only against the current run. If another writer starts a new run in-between, a retry of the original request creates a duplicate. Deriving the run ID from the request ID solves the gap, and relies on the pre-existing protections in store for duplicate RunID writes.
How did you test it?
Potential risks