From 6b8d7172a16358b090985b9bbe6daa42c5b491c6 Mon Sep 17 00:00:00 2001 From: yashsinghcodes Date: Tue, 13 Jan 2026 13:04:55 +0530 Subject: [PATCH 1/3] fix: worker execution takeover for cloud actions --- functions/onprem/worker/worker.go | 36 +++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 50910006..939f8bee 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -2652,6 +2652,42 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { } } + if strings.ToLower(actionResult.Action.Environment) == "cloud" { + + log.Printf("[WARNING] Got an action for %s environment forwarding it to the backend", actionResult.Action.Environment) + + streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl) + req, err := http.NewRequest( + "POST", + streamUrl, + bytes.NewBuffer([]byte(body)), + ) + + if err != nil { + log.Printf("[ERROR] Error building subflow (%s) request: %s", workflowExecution.ExecutionId, err) + return + } + + newresp, err := topClient.Do(req) + if err != nil { + log.Printf("[ERROR] Error running subflow (%s) request: %s", workflowExecution.ExecutionId, err) + return + } + + defer newresp.Body.Close() + if newresp.StatusCode != 200 { + body, err := ioutil.ReadAll(newresp.Body) + if err != nil { + log.Printf("[INFO][%s] Failed reading body after subflow request: %s", workflowExecution.ExecutionId, err) + return + } else { + log.Printf("[ERROR][%s] Failed forwarding subflow request of length %d\n: %s", workflowExecution.ExecutionId, len(actionResult.Result), string(body)) + } + } + + return + } + log.Printf("[DEBUG][%s] Action: Received, Label: '%s', Action: '%s', Status: %s, Run status: %s, Extra=Retry:%d", workflowExecution.ExecutionId, actionResult.Action.Label, actionResult.Action.AppName, actionResult.Status, workflowExecution.Status, retries) // results = append(results, actionResult) From ed6417b9be769951b14a53ccc5fe22e48fd0965f Mon Sep 17 00:00:00 2001 From: yashsinghcodes Date: Tue, 13 Jan 2026 17:46:47 +0530 Subject: [PATCH 2/3] added check for environment in handleQueue worker --- functions/onprem/worker/worker.go | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 939f8bee..0bd8ba1c 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -2652,8 +2652,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { } } - if strings.ToLower(actionResult.Action.Environment) == "cloud" { - + if strings.ToLower(actionResult.Action.Environment) != environment { log.Printf("[WARNING] Got an action for %s environment forwarding it to the backend", actionResult.Action.Environment) streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl) From 6921a63cc5ae60ad581867bb51ad6c8eaf7a17d9 Mon Sep 17 00:00:00 2001 From: yashsinghcodes Date: Tue, 13 Jan 2026 19:09:56 +0530 Subject: [PATCH 3/3] fix: env check condition --- functions/onprem/worker/worker.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 0bd8ba1c..f5a85d52 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -2652,7 +2652,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { } } - if strings.ToLower(actionResult.Action.Environment) != environment { + if strings.ToLower(actionResult.Action.Environment) != environment && len(environment) > 0 { log.Printf("[WARNING] Got an action for %s environment forwarding it to the backend", actionResult.Action.Environment) streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl)