Files
temporal/chasm/ref.go
tekkaya 3a425f7a94 Add run fallback to CHASM Nexus operation completion (#11028)
## What changed?

CHASM async Nexus operation completion now retries against the current workflow run when the original completion ref points at a run that was replaced by reset.

`CompleteNexusOperationChasm` now:
- first resolves the component at `RefConsistencyLevelComponentCreation`
- retries at `RefConsistencyLevelCurrentRun` if the first lookup returns `NotFound` before reaching the operation handler
- uses the completion request ID to re-establish operation identity after the current-run lookup

This matches the existing HSM completion behavior. The fallback only applies to workflow-backed operations, since standalone Nexus operations do not have a run chain.

## Why?

After reset, a workflow can continue the same scheduled Nexus operation on a new run with the same request ID. HSM completions already handle this by falling back to the current run.

The CHASM path previously used only the run ID from the completion ref, so a valid completion could be rejected after reset.

## 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)

## Potential risks
Low, chasm disabled for nexus operations in workflows currently and this PR adds the same functionality offered in hsm rails to chasm rails.
2026-07-23 01:01:38 +03:00

196 lines
7.3 KiB
Go

package chasm
import (
"reflect"
"go.temporal.io/api/serviceerror"
persistencespb "go.temporal.io/server/api/persistence/v1"
)
// ErrMalformedComponentRef is returned when component ref bytes cannot be deserialized.
var ErrMalformedComponentRef = serviceerror.NewInvalidArgument("malformed component ref")
// ErrInvalidComponentRef is returned when component ref bytes deserialize to an invalid component ref.
var ErrInvalidComponentRef = serviceerror.NewInvalidArgument("invalid component ref")
// ErrInvalidRefConsistencyLevel is returned when a component ref cannot be used at the requested RefConsistencyLevel.
var ErrInvalidRefConsistencyLevel = serviceerror.NewInvalidArgument("invalid consistency level for component ref")
// ExecutionKey uniquely identifies a CHASM execution in the system.
type ExecutionKey struct {
NamespaceID string
BusinessID string
RunID string
}
type ComponentRef struct {
ExecutionKey
// archetypeID is CHASM framework's internal ID for the type of the root component of the CHASM execution.
//
// It is used to find and validate the loaded execution has the right archetype, especially when runID
// is not specified in the ExecutionKey.
archetypeID ArchetypeID
// executionGoType is used for determining the ComponetRef's archetype.
// When CHASM deverloper needs to create a ComponentRef, they will only provide the component type,
// and leave the work of determining archetypeID to the CHASM framework.
executionGoType reflect.Type
// executionLastUpdateVT is the consistency token for the entire execution.
executionLastUpdateVT *persistencespb.VersionedTransition
// componentType is the fully qualified component type name.
// It is for performing partial loading more efficiently in future versions of CHASM.
//
// From the componentType, we can find the registered component struct definition,
// then use reflection to find sub-components and understand if those sub-components
// need to be loaded or not.
// We only need to do this for sub-components, path for parent/ancenstor components
// can be inferred from the current component path and they always needs to be loaded.
//
// componentType string
// componentPath and componentInitialVT are used to identify a component.
componentPath []string
componentInitialVT *persistencespb.VersionedTransition
validationFn func(NodeBackend, Context, Component, *Registry) error
}
// NewComponentRef creates a new ComponentRef with a registered root component go type.
//
// In V1, if you don't have a ref,
// then you can only interact with the (top level) execution.
func NewComponentRef[C Component](
executionKey ExecutionKey,
) ComponentRef {
return ComponentRef{
ExecutionKey: executionKey,
executionGoType: reflect.TypeFor[C](),
}
}
// NewComponentRefByArchetypeID creates a new ComponentRef with a known archetype ID.
// This should only be used by CHASM framework internals.
// CHASM library developers should use [NewComponentRef] instead.
func NewComponentRefByArchetypeID(
executionKey ExecutionKey,
archetypeID ArchetypeID,
) ComponentRef {
return ComponentRef{
ExecutionKey: executionKey,
archetypeID: archetypeID,
}
}
// forConsistencyLevel returns a copy of the ref adjusted for the requested [RefConsistencyLevel] by
// dropping the consistency tokens (and, at the weakest level, the run ID) that the level does not enforce.
// See [RefConsistencyLevel] for the ladder.
func (r *ComponentRef) forConsistencyLevel(level RefConsistencyLevel) (ComponentRef, error) {
ref := *r
switch level {
case RefConsistencyLevelExecutionLastUpdate:
// Strongest (default): enforce the execution-level token; nothing relaxed.
return ref, nil
case RefConsistencyLevelComponentCreation:
// Tolerate a stale execution transition without losing the staleness guarantee: run the
// execution staleness check ([Node.IsStale]) against the component's creation transition
// instead of the execution's last update. This still reloads a mutable state that predates
// the component (so the handler never operates on a state that doesn't yet know about the
// component), while no longer requiring the ref to match the latest execution transition.
// The creation transition is also matched in [Node.Component]; identity is otherwise the
// caller's responsibility (e.g. request ID).
ref.executionLastUpdateVT = ref.componentInitialVT
return ref, nil
case RefConsistencyLevelCurrentRun:
// Only workflow executions have a run chain to resolve a "current run" against, so this level
// is invalid for other archetypes (see ErrInvalidRefConsistencyLevel).
if r.archetypeID != WorkflowArchetypeID {
return ComponentRef{}, ErrInvalidRefConsistencyLevel
}
// Weakest: resolve by component path on the current run. Drop the run ID and every versioned
// transition; identity must be re-established by caller logic (e.g. request ID).
ref.executionLastUpdateVT = nil
ref.componentInitialVT = nil
ref.RunID = ""
return ref, nil
default:
return ref, serviceerror.NewInternalf("unknown ref consistency level: %d", level)
}
}
func (r *ComponentRef) ArchetypeID(
registry *Registry,
) (ArchetypeID, error) {
if r.archetypeID != UnspecifiedArchetypeID {
return r.archetypeID, nil
}
rc, ok := registry.componentOf(r.executionGoType)
if !ok {
return 0, serviceerror.NewInternal("unknown chasm component type: " + r.executionGoType.String())
}
r.archetypeID = rc.componentID
return r.archetypeID, nil
}
func (r *ComponentRef) Serialize(
registry *Registry,
) ([]byte, error) {
if r == nil {
return nil, nil
}
archetypeID, err := r.ArchetypeID(registry)
if err != nil {
return nil, err
}
pRef := persistencespb.ChasmComponentRef{
NamespaceId: r.NamespaceID,
BusinessId: r.BusinessID,
RunId: r.RunID,
ArchetypeId: archetypeID,
ExecutionVersionedTransition: r.executionLastUpdateVT,
ComponentPath: r.componentPath,
ComponentInitialVersionedTransition: r.componentInitialVT,
}
return pRef.Marshal()
}
// DeserializeComponentRef deserializes a byte slice into a ComponentRef.
// Provides caller the access to information including ExecutionKey, Archetype, and ShardingKey.
func DeserializeComponentRef(data []byte) (ComponentRef, error) {
if len(data) == 0 {
return ComponentRef{}, ErrInvalidComponentRef
}
var pRef persistencespb.ChasmComponentRef
if err := pRef.Unmarshal(data); err != nil {
return ComponentRef{}, ErrMalformedComponentRef
}
ref := ProtoRefToComponentRef(&pRef)
if ref.BusinessID == "" || ref.NamespaceID == "" {
return ComponentRef{}, ErrInvalidComponentRef
}
return ref, nil
}
// ProtoRefToComponentRef converts a persistence ChasmComponentRef reference to a
// ComponentRef. This is useful for situations where the protobuf ComponentRef has
// already been deserialized as part of an enclosing message.
func ProtoRefToComponentRef(pRef *persistencespb.ChasmComponentRef) ComponentRef {
return ComponentRef{
ExecutionKey: ExecutionKey{
NamespaceID: pRef.NamespaceId,
BusinessID: pRef.BusinessId,
RunID: pRef.RunId,
},
archetypeID: pRef.ArchetypeId,
executionLastUpdateVT: pRef.ExecutionVersionedTransition,
componentPath: pRef.ComponentPath,
componentInitialVT: pRef.ComponentInitialVersionedTransition,
}
}