Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .github/workflows/codeql-analysis.yml
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ jobs:

# Initializes the CodeQL tools for scanning.
- name: Initialize CodeQL
uses: github/codeql-action/init@54f647b7e1bb85c95cddabcd46b0c578ec92bc1a # v4.36.3
uses: github/codeql-action/init@e4fba868fa4b1b91e1fdab776edc8cfbe6e9fb81 # v4.37.3
with:
languages: ${{ matrix.language }}
# If you wish to specify custom queries, you can do so here or in a config file.
Expand All @@ -56,7 +56,7 @@ jobs:
# Autobuild attempts to build any compiled languages (C/C++, C#, or Java).
# If this step fails, then you should remove it and run the build manually (see below)
- name: Autobuild
uses: github/codeql-action/autobuild@8aad20d150bbac5944a9f9d289da16a4b0d87c1e # v4.36.2
uses: github/codeql-action/autobuild@e4fba868fa4b1b91e1fdab776edc8cfbe6e9fb81 # v4.37.3
# 📚 https://git.io/JvXDl

# ✏️ If the Autobuild fails above, remove it and uncomment the following three lines
Expand All @@ -68,4 +68,4 @@ jobs:
# make release

- name: Perform CodeQL Analysis
uses: github/codeql-action/analyze@8aad20d150bbac5944a9f9d289da16a4b0d87c1e # v4.36.2
uses: github/codeql-action/analyze@e4fba868fa4b1b91e1fdab776edc8cfbe6e9fb81 # v4.37.3
64 changes: 45 additions & 19 deletions aggregatedpool/workers.go
Original file line number Diff line number Diff line change
Expand Up @@ -125,59 +125,85 @@ func TemporalWorkers(wDef *Workflow, actDef *Activity, wi []*internal.WorkerInfo
workers := make([]worker.Worker, 0, len(wi))

for i := range wi {
log.Debug("worker info", zap.Any("worker_info", wi[i]))
workerInfo := wi[i]
log.Debug("worker info", zap.Any("worker_info", workerInfo))

// Override to 0: RoadRunner manages worker lifecycle independently
wi[i].Options.WorkerStopTimeout = 0
workerInfo.Options.WorkerStopTimeout = 0

if wi[i].TaskQueue == "" {
wi[i].TaskQueue = temporalClient.DefaultNamespace
if workerInfo.TaskQueue == "" {
workerInfo.TaskQueue = temporalClient.DefaultNamespace
}

if wi[i].Options.Identity == "" {
wi[i].Options.Identity = fmt.Sprintf(
if workerInfo.Options.Identity == "" {
workerInfo.Options.Identity = fmt.Sprintf(
"roadrunner:%s:%s",
wi[i].TaskQueue,
workerInfo.TaskQueue,
uuid.NewString(),
)
}

wi[i].Options.Interceptors = append(wi[i].Options.Interceptors, resolved...)

wrk := worker.New(tc, wi[i].TaskQueue, wi[i].Options)

for j := 0; j < len(wi[i].Workflows); j++ {
wf := wi[i].Workflows[j]
workerInfo.Options.Interceptors = append(workerInfo.Options.Interceptors, resolved...)

wrk := worker.New(tc, workerInfo.TaskQueue, workerInfo.Options)
dynamicWorkflowRegistered := false

for _, wf := range workerInfo.Workflows {
// A dynamic workflow is the catch-all: register it via
// RegisterDynamicWorkflow (not by name) so it handles any workflow
// type that has no statically registered handler. The shared proxy
// (wDef) is a WorkflowDefinitionFactory, which RegisterDynamicWorkflow
// accepts; it forwards the real workflow type name to PHP.
if wf.Dynamic {
if dynamicWorkflowRegistered {
return nil, errors.E(
errors.Op("temporal_workers"),
errors.Errorf("multiple dynamic workflows configured on task queue %q", workerInfo.TaskQueue),
)
}

err := registerWorkflow(func() {
wrk.RegisterDynamicWorkflow(wDef, workflow.DynamicRegisterOptions{})
}, wf.Name, workerInfo.TaskQueue)
if err != nil {
return nil, err
}
dynamicWorkflowRegistered = true

log.Debug("dynamic workflow registered", zap.String(tq, workerInfo.TaskQueue), zap.Any("workflow name", wf.Name))

continue
}

err := registerWorkflow(func() {
wrk.RegisterWorkflowWithOptions(wDef, workflow.RegisterOptions{
Name: wf.Name,
VersioningBehavior: wf.VersioningBehavior,
DisableAlreadyRegisteredCheck: false,
})
}, wf.Name, wi[i].TaskQueue)
}, wf.Name, workerInfo.TaskQueue)
if err != nil {
return nil, err
}

log.Debug("workflow registered", zap.String(tq, wi[i].TaskQueue), zap.Any("workflow name", wf.Name), zap.Int("versioning_behavior", int(wf.VersioningBehavior)))
log.Debug("workflow registered", zap.String(tq, workerInfo.TaskQueue), zap.Any("workflow name", wf.Name), zap.Int("versioning_behavior", int(wf.VersioningBehavior)))
}

if actDef.disableActivityWorkers {
log.Debug("activity workers disabled", zap.String(tq, wi[i].TaskQueue))
log.Debug("activity workers disabled", zap.String(tq, workerInfo.TaskQueue))
// add worker to the pool without activities
workers = append(workers, wrk)
continue
}

for j := 0; j < len(wi[i].Activities); j++ {
for _, activity := range workerInfo.Activities {
wrk.RegisterActivityWithOptions(actDef.execute, tActivity.RegisterOptions{
Name: wi[i].Activities[j].Name,
Name: activity.Name,
DisableAlreadyRegisteredCheck: false,
SkipInvalidStructFunctions: false,
})

log.Debug("activity registered", zap.String(tq, wi[i].TaskQueue), zap.Any("workflow name", wi[i].Activities[j].Name))
log.Debug("activity registered", zap.String(tq, workerInfo.TaskQueue), zap.Any("workflow name", activity.Name))
}
// add worker to the pool
workers = append(workers, wrk)
Expand Down
22 changes: 22 additions & 0 deletions aggregatedpool/workers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,12 @@ import (
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/temporalio/roadrunner-temporal/v5/api"
"github.com/temporalio/roadrunner-temporal/v5/internal"
commonpb "go.temporal.io/api/common/v1"
"go.temporal.io/sdk/client"
"go.temporal.io/sdk/converter"
sdkinterceptor "go.temporal.io/sdk/interceptor"
"go.uber.org/zap"
)

// mockPayloadConverter implements converter.PayloadConverter for testing.
Expand Down Expand Up @@ -278,6 +281,25 @@ func TestRegisterWorkflow_NonStringPanic_Handled(t *testing.T) {
assert.Contains(t, err.Error(), "42", "should preserve a non-string panic value")
}

func TestTemporalWorkers_MultipleDynamicWorkflows_ReturnsError(t *testing.T) {
temporalClient, err := client.NewLazyClient(client.Options{})
require.NoError(t, err)
t.Cleanup(temporalClient.Close)

workers := []*internal.WorkerInfo{{
TaskQueue: "default",
Workflows: []internal.WorkflowInfo{
{Name: "DynamicOne", Dynamic: true},
{Name: "DynamicTwo", Dynamic: true},
},
}}

_, err = TemporalWorkers(nil, nil, workers, zap.NewNop(), temporalClient, nil, nil)
require.Error(t, err)
assert.Contains(t, err.Error(), "multiple dynamic workflows")
assert.Contains(t, err.Error(), "default")
}

func TestResolveDataConverters_EmptyMap_WithConfig(t *testing.T) {
_, err := ResolveDataConverters(map[string]converter.PayloadConverter{}, []string{"encoding/a"})
require.Error(t, err)
Expand Down
8 changes: 4 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module github.com/temporalio/roadrunner-temporal/v5

go 1.26.4
go 1.26.5

require (
github.com/goccy/go-json v0.10.6
Expand All @@ -13,8 +13,8 @@ require (
github.com/roadrunner-server/pool v1.1.3
github.com/stretchr/testify v1.11.1
github.com/uber-go/tally/v4 v4.1.17
go.temporal.io/api v1.62.14
go.temporal.io/sdk v1.44.1
go.temporal.io/api v1.63.4
go.temporal.io/sdk v1.47.0
go.temporal.io/sdk/contrib/tally v0.2.0
go.temporal.io/server v1.31.1
go.uber.org/zap v1.28.0
Expand Down Expand Up @@ -57,7 +57,7 @@ require (
golang.org/x/time v0.15.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect
google.golang.org/grpc v1.81.1
google.golang.org/grpc v1.82.1
gopkg.in/yaml.v3 v3.0.1 // indirect
)

Expand Down
12 changes: 6 additions & 6 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -241,11 +241,11 @@ go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/
go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE=
go.opentelemetry.io/proto/otlp v0.7.0/go.mod h1:PqfVotwruBrMGOCsRd/89rSnXhoiJIqeYNgFYFoEGnI=
go.temporal.io/api v1.5.0/go.mod h1:BqKxEJJYdxb5dqf0ODfzfMxh8UEQ5L3zKS51FiIYYkA=
go.temporal.io/api v1.62.14 h1:Tree3eqoKRt5Vv+nvYHMPp/ROiGmSOTtFUD9d0w8yJE=
go.temporal.io/api v1.62.14/go.mod h1:0k75tRljEuELWGeXjEZZO7zYqBln4+1FrG6+IMOMy7Q=
go.temporal.io/api v1.63.4 h1:p4dVIAP3dJop0MfcyH9QSzjU7+V/ttLDhxFhSRUar58=
go.temporal.io/api v1.63.4/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ=
go.temporal.io/sdk v1.12.0/go.mod h1:lSp3lH1lI0TyOsus0arnO3FYvjVXBZGi/G7DjnAnm6o=
go.temporal.io/sdk v1.44.1 h1:Mt2OZLZpqkzDIdg9YyQzO0Rb/HqCDnnqHlIAGAJ5gqM=
go.temporal.io/sdk v1.44.1/go.mod h1:vkApR12F9/Y8OR+hkxe7WyXQFuCX6clhzqnAk6rzDAM=
go.temporal.io/sdk v1.47.0 h1:lZ39w1+uWSjHTL0F3mSc0t4XUnKX8CCcWxqSiLbaHnc=
go.temporal.io/sdk v1.47.0/go.mod h1:ilKs0twgP4JpP8pfhIgZumnOEyBiYn6ZO/ta//NnKMU=
go.temporal.io/sdk/contrib/tally v0.2.0 h1:XnTJIQcjOv+WuCJ1u8Ve2nq+s2H4i/fys34MnWDRrOo=
go.temporal.io/sdk/contrib/tally v0.2.0/go.mod h1:1kpSuCms/tHeJQDPuuKkaBsMqfHnIIRnCtUYlPNXxuE=
go.temporal.io/server v1.31.1 h1:3rxA0Ls21hLOseKsyzm7e+IrVLuaDsuri0DGX+chIqU=
Expand Down Expand Up @@ -388,8 +388,8 @@ google.golang.org/grpc v1.29.1/go.mod h1:itym6AZVZYACWQqET3MqgPpjcuV5QH3BxFS3Iji
google.golang.org/grpc v1.33.1/go.mod h1:fr5YgcSWrqhRRxogOsw7RzIpsmvOZ6IcH4kBYTpR3n0=
google.golang.org/grpc v1.36.0/go.mod h1:qjiiYl8FncCW8feJPdyg3v6XW24KsRHe+dy9BAGRRjU=
google.golang.org/grpc v1.40.0/go.mod h1:ogyxbiOoUXAkP+4+xa6PZSE9DZgIHtSpzjDTB9KAK34=
google.golang.org/grpc v1.81.1 h1:VnnIIZ88UzOOKLukQi+ImGz8O1Wdp8nAGGnvOfEIWQQ=
google.golang.org/grpc v1.81.1/go.mod h1:xGH9GfzOyMTGIOXBJmXt+BX/V0kcdQbdcuwQ/zNw42I=
google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE=
google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA=
google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8=
google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0=
google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM=
Expand Down
2 changes: 1 addition & 1 deletion go.work
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
go 1.26.3
go 1.26.5

use (
.
Expand Down
3 changes: 3 additions & 0 deletions internal/worker_info.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,9 @@ type WorkflowInfo struct {
Signals []string `json:"signals"`
// VersioningBehavior for the workflow.
VersioningBehavior workflow.VersioningBehavior `json:"versioning_behavior,omitempty"`
// Dynamic marks this as the catch-all workflow, registered via
// RegisterDynamicWorkflow so it handles any otherwise-unregistered type.
Dynamic bool `json:"dynamic,omitempty"`
}

// ActivityInfo describes single worker activity.
Expand Down
8 changes: 4 additions & 4 deletions tests/go.mod
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module tests

go 1.26.4
go 1.26.5

require (
github.com/fatih/color v1.19.0
Expand All @@ -19,8 +19,8 @@ require (
github.com/stretchr/testify v1.11.1
github.com/temporalio/roadrunner-temporal/v5 v5.11.0
go.opentelemetry.io/otel/sdk v1.44.0
go.temporal.io/api v1.62.14
go.temporal.io/sdk v1.44.1
go.temporal.io/api v1.63.4
go.temporal.io/sdk v1.47.0
go.temporal.io/sdk/contrib/opentelemetry v0.7.0
go.uber.org/zap v1.28.0
)
Expand Down Expand Up @@ -94,7 +94,7 @@ require (
golang.org/x/time v0.15.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect
google.golang.org/grpc v1.81.1 // indirect
google.golang.org/grpc v1.82.1 // indirect
google.golang.org/protobuf v1.36.11 // indirect
gopkg.in/natefinch/lumberjack.v2 v2.2.1 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
Expand Down
12 changes: 6 additions & 6 deletions tests/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -290,11 +290,11 @@ go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/
go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE=
go.opentelemetry.io/proto/otlp v0.7.0/go.mod h1:PqfVotwruBrMGOCsRd/89rSnXhoiJIqeYNgFYFoEGnI=
go.temporal.io/api v1.5.0/go.mod h1:BqKxEJJYdxb5dqf0ODfzfMxh8UEQ5L3zKS51FiIYYkA=
go.temporal.io/api v1.62.14 h1:Tree3eqoKRt5Vv+nvYHMPp/ROiGmSOTtFUD9d0w8yJE=
go.temporal.io/api v1.62.14/go.mod h1:0k75tRljEuELWGeXjEZZO7zYqBln4+1FrG6+IMOMy7Q=
go.temporal.io/api v1.63.4 h1:p4dVIAP3dJop0MfcyH9QSzjU7+V/ttLDhxFhSRUar58=
go.temporal.io/api v1.63.4/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ=
go.temporal.io/sdk v1.12.0/go.mod h1:lSp3lH1lI0TyOsus0arnO3FYvjVXBZGi/G7DjnAnm6o=
go.temporal.io/sdk v1.44.1 h1:Mt2OZLZpqkzDIdg9YyQzO0Rb/HqCDnnqHlIAGAJ5gqM=
go.temporal.io/sdk v1.44.1/go.mod h1:vkApR12F9/Y8OR+hkxe7WyXQFuCX6clhzqnAk6rzDAM=
go.temporal.io/sdk v1.47.0 h1:lZ39w1+uWSjHTL0F3mSc0t4XUnKX8CCcWxqSiLbaHnc=
go.temporal.io/sdk v1.47.0/go.mod h1:ilKs0twgP4JpP8pfhIgZumnOEyBiYn6ZO/ta//NnKMU=
go.temporal.io/sdk/contrib/opentelemetry v0.7.0 h1:GSna1HP+1ibNXZ9xlVdQU2zFVqdt5VcdF0dzpeaYccQ=
go.temporal.io/sdk/contrib/opentelemetry v0.7.0/go.mod h1:oQJC6UIl3FbSYh4f2MlUAIYSE6FPw02X1Tw8/bOvfxg=
go.temporal.io/sdk/contrib/tally v0.2.0 h1:XnTJIQcjOv+WuCJ1u8Ve2nq+s2H4i/fys34MnWDRrOo=
Expand Down Expand Up @@ -440,8 +440,8 @@ google.golang.org/grpc v1.29.1/go.mod h1:itym6AZVZYACWQqET3MqgPpjcuV5QH3BxFS3Iji
google.golang.org/grpc v1.33.1/go.mod h1:fr5YgcSWrqhRRxogOsw7RzIpsmvOZ6IcH4kBYTpR3n0=
google.golang.org/grpc v1.36.0/go.mod h1:qjiiYl8FncCW8feJPdyg3v6XW24KsRHe+dy9BAGRRjU=
google.golang.org/grpc v1.40.0/go.mod h1:ogyxbiOoUXAkP+4+xa6PZSE9DZgIHtSpzjDTB9KAK34=
google.golang.org/grpc v1.81.1 h1:VnnIIZ88UzOOKLukQi+ImGz8O1Wdp8nAGGnvOfEIWQQ=
google.golang.org/grpc v1.81.1/go.mod h1:xGH9GfzOyMTGIOXBJmXt+BX/V0kcdQbdcuwQ/zNw42I=
google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE=
google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA=
google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8=
google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0=
google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM=
Expand Down