[Nexus] Support internal callback completion routing to CHASM - #8372
Conversation
| NexusComponentRefHeader = "X-CHASM-Component-Ref" | ||
|
|
||
| // Base URL for Nexus->CHASM callbacks. | ||
| NexusCompletionHandlerURL = "temporal://internal/chasm" |
There was a problem hiding this comment.
| NexusCompletionHandlerURL = "temporal://internal/chasm" | |
| NexusCompletionHandlerURL = "temporal://internal" |
|
|
||
| const ( | ||
| // Header name for the CHASM ComponentRef. | ||
| NexusComponentRefHeader = "X-CHASM-Component-Ref" |
There was a problem hiding this comment.
No need, just put this in the callback token header.
There was a problem hiding this comment.
When I synced up with @chaptersix , we decided to use SystemCallbackURL so it goes through the same handler as other internal callbacks, which I think makes sense. However, if we don't have a unique URL nor a unique CHASM header, how would we distinguish this as a CHASM request in the callback executor?
There was a problem hiding this comment.
We are thinking of deprecating the callback headers for the most part, which is why I would ask that you use a single token field.
There was a problem hiding this comment.
Moved to the commonnexus.CallbackTokenHeader field. Relies on temporal://internal being unique to CHASM.
| if err != nil { | ||
| return err | ||
| } | ||
| case *persistencespb.Callback_Hsm: |
There was a problem hiding this comment.
Can just remove this branch and the rest of the code it references.
There was a problem hiding this comment.
Sure, I'll open a separate PR to remove the Callback_Hsm invocation stuff.
| chasmInvokable := chasmInvocation{} | ||
| chasmInvokable.nexus = variant.Nexus | ||
| chasmInvokable.attempt = callback.Attempt | ||
| invokable = chasmInvokable |
There was a problem hiding this comment.
You're going to want to set the completion on here too...
Would you also please just set all of the fields here inline with the struct definition?
Could just do:
invokable = chasmInvocation{
nexus: variant.Nexus,
...
}
There was a problem hiding this comment.
Oops, yeah, cleaned these up.
| } | ||
|
|
||
| // variant struct is immutable and ok to reference without copying | ||
| nexusInvokable := nexusInvocation{} |
There was a problem hiding this comment.
Same here, would be nice to just inline the whole field assignment (unrelated to your change).
| // this similarly as we would a pure task (holding an exclusive lock), as the | ||
| // assumption is that the accessed component will be recording (or generating a | ||
| // task) based on this result. | ||
| _, err = e.ChasmEngine.UpdateComponent(ctx, ref, func(ctx chasm.MutableContext, component chasm.Component) error { |
There was a problem hiding this comment.
This doesn't do an RPC, you're going to have to add an RPC to the shard that owns the execution of the original caller.
We can then optimize it so if the ref is something that we can resolve locally we would bypass the RPC call.
There was a problem hiding this comment.
The RPC would be historyservice.CompleteNexusOperation, you'll want to update that method to invoke ChasmEngine.UpdateComponent and take a different form of completion.
Feel free to rename the completion field to hsm_completion to disambiguate between HSM and CHASM refs.
There was a problem hiding this comment.
Let's also not change the Nexus RPC stuff for now since we are not using that SDK for delivering the completions.
We can change it later when the SDK supports pluggable transports if we want to but that's probably not critical.
| } | ||
| // Internally routed. | ||
| // TODO - validate requests involving this URL come from a history pod. | ||
| if u.Scheme == "temporal" { |
There was a problem hiding this comment.
@chaptersix will take care of this, he's modifying this code right now.
There was a problem hiding this comment.
We need to validate that the caller is from within the cluster if it tries to attach a callback with an internal URL.
There was a problem hiding this comment.
You need to merge your code with main (or rebase if you prefer :))
There was a problem hiding this comment.
I'll do that now that the stacked PRs are merged.
| // Base URL for Nexus->CHASM callbacks. | ||
| NexusCompletionHandlerURL = "temporal://system/nexus/callback/chasm" |
There was a problem hiding this comment.
I'm still planning to get rid of this and use the constant SystemCallbackURL, but this branch is stacked on top of others, so I'll do that after others have merged in.
There was a problem hiding this comment.
I would consider using temporal://internal to avoid having to unpack the callback tokens.
There was a problem hiding this comment.
A good reason to use temporal://internal callbacks is because they will trip a separate circuit breaker in the outbound queue.
We should also consider running these callbacks on the immediate "transfer" queue instead and avoid using circuit breakers altogether but let's add an issue for this, we are going to refactor callbacks to use CHASM soon and can do this work as part of that.
|
|
||
| const ( | ||
| // Header name for the CHASM ComponentRef. | ||
| NexusComponentRefHeader = "X-CHASM-Component-Ref" |
There was a problem hiding this comment.
We are thinking of deprecating the callback headers for the most part, which is why I would ask that you use a single token field.
| // Base URL for Nexus->CHASM callbacks. | ||
| NexusCompletionHandlerURL = "temporal://system/nexus/callback/chasm" |
There was a problem hiding this comment.
I would consider using temporal://internal to avoid having to unpack the callback tokens.
| func (c chasmInvocation) Invoke(ctx context.Context, ns *namespace.Namespace, e taskExecutor, task InvocationTask) invocationResult { | ||
| // Get back the base64-encoded ComponentRef from the header. | ||
| encodedRef, ok := c.nexus.GetHeader()[chasm.NexusComponentRefHeader] | ||
| if !ok { | ||
| return invocationResultFail{errors.New("callback missing CHASM header")} | ||
| } | ||
|
|
||
| decodedRef, err := base64.RawURLEncoding.DecodeString(encodedRef) | ||
| if err != nil { | ||
| return invocationResultFail{fmt.Errorf("failed to decode CHASM ComponentRef: %v", err)} | ||
| } | ||
|
|
||
| ref := &persistencespb.ChasmComponentRef{} | ||
| err = proto.Unmarshal(decodedRef, ref) | ||
| if err != nil { | ||
| return invocationResultFail{fmt.Errorf("failed to unmarshal CHASM ComponentRef: %v", err)} | ||
| } | ||
|
|
||
| request, err := c.getHistoryRequest(ref) | ||
| if err != nil { | ||
| return invocationResultFail{fmt.Errorf("failed to build history request: %v", err)} | ||
| } | ||
|
|
||
| // RPC to History for cross-shard completion delivery. | ||
| _, err = e.HistoryClient.CompleteNexusOperationChasm(ctx, request) | ||
| if err != nil { | ||
| return invocationResultRetry{fmt.Errorf("failed to complete Nexus operation: %v", err)} | ||
| } | ||
|
|
||
| return invocationResultOK{} | ||
| } |
There was a problem hiding this comment.
Are these internal errors ever exposed to users or is this all internal?
There was a problem hiding this comment.
This is internal-facing, for operators who are log-diving.
|
|
||
| func (c chasmInvocation) WrapError(result invocationResult, err error) error { | ||
| if failure, ok := result.(invocationResultFail); ok { | ||
| return queues.NewUnprocessableTaskError(failure.err.Error()) |
There was a problem hiding this comment.
This will prevent the task from being retried, is that what you are trying to do here? Doesn't seem like it...
There was a problem hiding this comment.
Yes, that's what I'm trying to do; the errors I've marked as 'fail' explicitly are non-retryable (e.g., fails validation checks).
There was a problem hiding this comment.
I see. That's fine. I think you just don't want to wrap in a destination down error though to avoid triggering the circuit breaker.
There was a problem hiding this comment.
Updated to not wrap in DestinationDown.
| return nil, fmt.Errorf("failed to convert Nexus links: %v", err) | ||
| } | ||
|
|
||
| payload, err := io.ReadAll(op.Reader) |
There was a problem hiding this comment.
Use proto.Unmarshal into common.Payload message here. And consider that the payload may be nil.
| // ChasmNexusCompletion includes details about a completed Nexus operation. | ||
| message ChasmNexusCompletionInfo { | ||
| // Operation state - may only be successful / failed / canceled. | ||
| string state = 1; |
| } | ||
|
|
||
| // ChasmNexusCompletion includes details about a completed Nexus operation. | ||
| message ChasmNexusCompletionInfo { |
There was a problem hiding this comment.
It's to disambiguate against the token, ChasmNexusCompletion; I could suffix that instead, maybe?
There was a problem hiding this comment.
Renamed to ChasmNexusCompletion.
| } | ||
| // Internally routed. | ||
| // TODO - validate requests involving this URL come from a history pod. | ||
| if u.Scheme == "temporal" { |
There was a problem hiding this comment.
You need to merge your code with main (or rebase if you prefer :))
| _, err := h.chasmEngine.UpdateComponent(ctx, ref, func(ctx chasm.MutableContext, component chasm.Component) error { | ||
| handler, ok := component.(chasm.NexusCompletionHandler) | ||
| if !ok { | ||
| return serviceerror.NewUnimplementedf("component '%T' does not implement NexusCompletionHandler", component) |
There was a problem hiding this comment.
Will the caller consider this a retryable error? We should verify that it doesn't because the callback executor will just spin indefinitely.
There was a problem hiding this comment.
Good point, yeah the ChasmInvocation shouldn't be retrying on this.
| // Base URL for Nexus->CHASM callbacks. | ||
| NexusCompletionHandlerURL = "temporal://system/nexus/callback/chasm" |
There was a problem hiding this comment.
A good reason to use temporal://internal callbacks is because they will trip a separate circuit breaker in the outbound queue.
We should also consider running these callbacks on the immediate "transfer" queue instead and avoid using circuit breakers altogether but let's add an issue for this, we are going to refactor callbacks to use CHASM soon and can do this work as part of that.
| ) (*historyservice.CompleteNexusOperationChasmRequest, error) { | ||
| var req *historyservice.CompleteNexusOperationChasmRequest | ||
|
|
||
| token := &tokenspb.NexusOperationChasmCompletion{ |
| req = &historyservice.CompleteNexusOperationChasmRequest{ | ||
| Completion: token, | ||
| Outcome: &historyservice.CompleteNexusOperationChasmRequest_Failure{ | ||
| Failure: apiFailure, | ||
| }, | ||
| CloseTime: timestamppb.New(op.CloseTime), |
There was a problem hiding this comment.
nit: you could just set the outcome here and initialize the request outside of the switch statement.
There was a problem hiding this comment.
I thought so too, but CloseTime is independent on both CompletionSuccess/CompletionFailure variants, so you have to type assert for that as well anyways.
| // Complete an async Nexus Operation using a completion token. The completion state could be successful, failed, or | ||
| // canceled. | ||
| // | ||
| // Deprecated. Will be renamed to CompleteNexusOperationHsm in a future release. |
There was a problem hiding this comment.
Not to start using it, but have it defined already. You need two versions to safely point the client to use the new API.
| } | ||
|
|
||
| // ChasmNexusCompletion includes details about a completed Nexus operation. | ||
| message ChasmNexusCompletionInfo { |
| } | ||
|
|
||
| // A completion token for a Nexus operation started from a CHASM component. | ||
| message NexusOperationChasmCompletion { |
There was a problem hiding this comment.
Hmm... can we merge this into the existing NexusOperationCompletion and deprecate the original fields?
You can just add a component_ref field and it would work well.
Not critical but it would simplify the decoding later.
| // this similarly as we would a pure task (holding an exclusive lock), as the | ||
| // assumption is that the accessed component will be recording (or generating a | ||
| // task) based on this result. | ||
| _, err := h.chasmEngine.UpdateComponent(ctx, ref, func(ctx chasm.MutableContext, component chasm.Component) error { |
There was a problem hiding this comment.
Use chasm.UpdateComponent and avoid casting. If you merge the current main changes, you'll see that the chasm engine should already be on the request context.
There was a problem hiding this comment.
Use chasm.UpdateComponent and avoid casting
I don't think I can do that here in a way that keeps this generic. I'm using the NexusCompletionHandler interface which doesn't bundle in Component, therefore, if I try to accept a NexusCompletionHandler here, it won't meet the generic requirements of implementing Component's LifecycleState.
There was a problem hiding this comment.
That sounds like a bug in CHASM. You shouldn't need to inspect the type. There's no reason that chasmEngine.UpdateComponent would work and chasm.UpdateComponent(ctx) would not work.
There was a problem hiding this comment.
What component type would you expect I use here? NexusCompletionHandler won't work since it isn't a Component.
There was a problem hiding this comment.
That's what I expect, that you can use the interface instead of a concrete type.
0b752c9 to
9aee894
Compare
| } | ||
|
|
||
| if retry, ok := result.(invocationResultRetry); ok { | ||
| return queues.NewDestinationDownError(retry.err.Error(), err) |
There was a problem hiding this comment.
| return queues.NewDestinationDownError(retry.err.Error(), err) | |
| return err |
| return err | ||
| } | ||
|
|
||
| func (c chasmInvocation) Invoke(ctx context.Context, ns *namespace.Namespace, e taskExecutor, task InvocationTask) invocationResult { |
There was a problem hiding this comment.
Thinking about this some more, we shouldn't expose these errors. They are user visible. I suggest that instead we follow this pattern.
referenceID := uuid.NewString()
e.Logger.Error("failed to decode CHASM ComponentRef", tag.Error(err), tag.NewStringTag("reference-id", referendID))
return invocationResultFail{fmt.Errorf("internal error, reference-id: %v", referenceID)}Same for invocationResultRetry
| { | ||
| name: "success-with-successful-operation", | ||
| setupHistoryClient: func(t *testing.T, ctrl *gomock.Controller) *historyservicemock.MockHistoryServiceClient { | ||
| client := historyservicemock.NewMockHistoryServiceClient(ctrl) |
There was a problem hiding this comment.
Unrelated to this PR but I couldn't help myself. Really the whole mock and codegen here is redundant, this is so much simpler:
type MockHistoryClient struct {
historyservice.HistoryServiceClient
HandleCompleteNexusOperationChasm func(ctx context.Context, req *historyservice.CompleteNexusOperationChasmRequest, opts ...grpc.CallOption) (*historyservice.CompleteNexusOperationChasmResponse, error)
}| requestID string | ||
| } | ||
|
|
||
| var ErrUnimplementedHandler = serviceerror.NewUnimplemented("component does not implement NexusCompletionHandler") |
| // ChasmNexusCompletion includes details about a completed Nexus operation. | ||
| message ChasmNexusCompletion { |
There was a problem hiding this comment.
| // ChasmNexusCompletion includes details about a completed Nexus operation. | |
| message ChasmNexusCompletion { | |
| // NexusCompletion includes details about a completed Nexus operation. | |
| message NexusCompletion { |
There was a problem hiding this comment.
Every other type in this protobuf is prefixed with Chasm, I'd rather keep that consistent or change all at once.
There was a problem hiding this comment.
It's not CHASM specific, you can move this somewhere else but not blocking the PR.
| // this similarly as we would a pure task (holding an exclusive lock), as the | ||
| // assumption is that the accessed component will be recording (or generating a | ||
| // task) based on this result. | ||
| _, err := h.chasmEngine.UpdateComponent(ctx, ref, func(ctx chasm.MutableContext, component chasm.Component) error { |
There was a problem hiding this comment.
That sounds like a bug in CHASM. You shouldn't need to inspect the type. There's no reason that chasmEngine.UpdateComponent would work and chasm.UpdateComponent(ctx) would not work.
…esponse.proto Co-authored-by: Roey Berman <roey@temporal.io>
| // ChasmNexusCompletion includes details about a completed Nexus operation. | ||
| message ChasmNexusCompletion { |
There was a problem hiding this comment.
It's not CHASM specific, you can move this somewhere else but not blocking the PR.
| // this similarly as we would a pure task (holding an exclusive lock), as the | ||
| // assumption is that the accessed component will be recording (or generating a | ||
| // task) based on this result. | ||
| _, err := h.chasmEngine.UpdateComponent(ctx, ref, func(ctx chasm.MutableContext, component chasm.Component) error { |
There was a problem hiding this comment.
That's what I expect, that you can use the interface instead of a concrete type.
What changed?
temporal://internal/chasm, which are routed to CHASM components through a newChasmInvocation.chasm.NexusCompletionHandler, which components can implement to receive Nexus callbacks.GetNexusCallback, has been added that returns aCallbackprotobuf routing to the given component.Why?
How did you test it?
Potential risks
StartWorkflowExecution's frontend flow). But it would be nice to validate that it was only coming from internal pods.