go.temporal.io/server/chasm/registry_test.go
559 LOC · 0 covered · 559 uncovered · 0 ranges · 0 concepts · 0 introducers · 0 tests
1
package chasm_test
2
3
import (
4
"reflect"
5
"testing"
6
7
"github.com/nexus-rpc/sdk-go/nexus"
8
"github.com/stretchr/testify/require"
9
"github.com/stretchr/testify/suite"
10
enumspb "go.temporal.io/api/enums/v1"
11
"go.temporal.io/server/chasm"
12
"go.temporal.io/server/common/log"
13
"go.temporal.io/server/common/searchattribute/sadefs"
14
"go.uber.org/mock/gomock"
15
)
16
17
type (
18
RegistryTestSuite struct {
19
suite.Suite
20
logger log.Logger
21
}
22
23
testTask1 struct{}
24
testTask2 struct{}
25
testTaskComponentInterface interface {
26
DoSomething()
27
}
28
29
// testComponentWithVisibility is a test component that has a Visibility field.
30
testComponentWithVisibility struct {
31
chasm.UnimplementedComponent
32
Visibility chasm.Field[*chasm.Visibility]
33
}
34
)
35
36
func (t *testComponentWithVisibility) LifecycleState(_ chasm.Context) chasm.LifecycleState {
37
return chasm.LifecycleStateRunning
38
}
39
40
func TestRegistryTestSuite(t *testing.T) {
41
suite.Run(t, new(RegistryTestSuite))
42
}
43
44
func (s *RegistryTestSuite) SetupTest() {
45
s.logger = log.NewTestLogger()
46
}
47
48
func (s *RegistryTestSuite) TestRegistry_RegisterComponents_Success() {
49
r := chasm.NewRegistry(s.logger)
50
ctrl := gomock.NewController(s.T())
51
lib := chasm.NewMockLibrary(ctrl)
52
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
53
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
54
chasm.NewRegistrableComponent[*chasm.MockComponent]("Component1"),
55
})
56
57
lib.EXPECT().Tasks().Return(nil)
58
lib.EXPECT().NexusServices().Return(nil)
59
lib.EXPECT().NexusServiceProcessors().Return(nil)
60
61
err := r.Register(lib)
62
require.NoError(s.T(), err)
63
64
rc1, ok := r.Component("TestLibrary.Component1")
65
require.True(s.T(), ok)
66
require.Equal(s.T(), "TestLibrary.Component1", rc1.FqType())
67
68
missingRC, ok := r.Component("TestLibrary.Component2")
69
require.False(s.T(), ok)
70
require.Nil(s.T(), missingRC)
71
72
cInstance1 := chasm.NewMockComponent(ctrl)
73
rc2, ok := r.ComponentFor(cInstance1)
74
require.True(s.T(), ok)
75
require.Equal(s.T(), "TestLibrary.Component1", rc2.FqType())
76
77
rc2, ok = r.ComponentOf(reflect.TypeFor[*chasm.MockComponent]())
78
require.True(s.T(), ok)
79
require.Equal(s.T(), "TestLibrary.Component1", rc2.FqType())
80
81
cInstance2 := "invalid component instance"
82
rc3, ok := r.ComponentFor(cInstance2)
83
require.False(s.T(), ok)
84
require.Nil(s.T(), rc3)
85
}
86
87
func (s *RegistryTestSuite) TestRegistry_RegisterComponents_WithDetached() {
88
r := chasm.NewRegistry(s.logger)
89
ctrl := gomock.NewController(s.T())
90
lib := chasm.NewMockLibrary(ctrl)
91
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
92
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
93
chasm.NewRegistrableComponent[*chasm.MockComponent]("DetachedComponent", chasm.WithDetached()),
94
})
95
lib.EXPECT().Tasks().Return(nil)
96
lib.EXPECT().NexusServices().Return(nil)
97
lib.EXPECT().NexusServiceProcessors().Return(nil)
98
99
err := r.Register(lib)
100
s.Require().NoError(err)
101
102
// Detached component should have IsDetached() return true
103
detachedRC, ok := r.Component("TestLibrary.DetachedComponent")
104
s.Require().True(ok)
105
s.Require().True(detachedRC.IsDetached())
106
107
// Verify that a component without WithDetached() has IsDetached() return false
108
normalRC := chasm.NewRegistrableComponent[*chasm.MockComponent]("NormalComponent")
109
s.Require().False(normalRC.IsDetached())
110
}
111
112
func (s *RegistryTestSuite) TestRegistry_RegisterTasks_Success() {
113
r := chasm.NewRegistry(s.logger)
114
ctrl := gomock.NewController(s.T())
115
lib := chasm.NewMockLibrary(ctrl)
116
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
117
lib.EXPECT().Components().Return(nil)
118
lib.EXPECT().NexusServices().Return(nil)
119
lib.EXPECT().NexusServiceProcessors().Return(nil)
120
121
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{
122
chasm.NewRegistrableSideEffectTask(
123
"Task1",
124
chasm.NewMockSideEffectTaskHandler[*chasm.MockComponent, testTask1](ctrl),
125
chasm.WithTaskGroup("test-task-group"),
126
),
127
chasm.NewRegistrablePureTask(
128
"Task2",
129
chasm.NewMockPureTaskHandler[testTaskComponentInterface, testTask2](ctrl),
130
),
131
})
132
133
err := r.Register(lib)
134
require.NoError(s.T(), err)
135
136
rt1, ok := r.Task("TestLibrary.Task1")
137
require.True(s.T(), ok)
138
require.Equal(s.T(), "TestLibrary.Task1", rt1.FqType())
139
s.Require().Equal("test-task-group", rt1.TaskGroup())
140
141
missingRT, ok := r.Task("TestLibrary.TaskMissing")
142
require.False(s.T(), ok)
143
require.Nil(s.T(), missingRT)
144
145
tInstance1 := testTask2{}
146
rt2, ok := r.TaskFor(tInstance1)
147
require.True(s.T(), ok)
148
require.Equal(s.T(), "TestLibrary.Task2", rt2.FqType())
149
s.Require().Equal(rt2.FqType(), rt2.TaskGroup())
150
151
rt2, ok = r.TaskOf(reflect.TypeFor[testTask2]())
152
require.True(s.T(), ok)
153
require.Equal(s.T(), "TestLibrary.Task2", rt2.FqType())
154
155
tInstance2 := "invalid task instance"
156
rt3, ok := r.TaskFor(tInstance2)
157
require.False(s.T(), ok)
158
require.Nil(s.T(), rt3)
159
}
160
161
func (s *RegistryTestSuite) TestRegistry_Register_LibraryError() {
162
ctrl := gomock.NewController(s.T())
163
lib := chasm.NewMockLibrary(ctrl)
164
165
s.T().Run("library name must not be empty", func(t *testing.T) {
166
lib.EXPECT().Name().Return("")
167
r := chasm.NewRegistry(s.logger)
168
err := r.Register(lib)
169
require.Error(t, err)
170
require.Contains(t, err.Error(), "name must not be empty")
171
})
172
173
s.T().Run("library name must follow rules", func(t *testing.T) {
174
lib.EXPECT().Name().Return("bad.lib.name")
175
r := chasm.NewRegistry(s.logger)
176
err := r.Register(lib)
177
require.Error(t, err)
178
require.Contains(t, err.Error(), "name must follow golang identifier rules")
179
})
180
}
181
182
func (s *RegistryTestSuite) TestRegistry_RegisterComponents_Error() {
183
ctrl := gomock.NewController(s.T())
184
lib := chasm.NewMockLibrary(ctrl)
185
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
186
187
s.T().Run("component name must not be empty", func(t *testing.T) {
188
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
189
chasm.NewRegistrableComponent[*chasm.MockComponent](""),
190
})
191
r := chasm.NewRegistry(s.logger)
192
err := r.Register(lib)
193
require.Error(t, err)
194
require.Contains(t, err.Error(), "name must not be empty")
195
})
196
197
s.T().Run("component name must follow rules", func(t *testing.T) {
198
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
199
chasm.NewRegistrableComponent[*chasm.MockComponent]("bad.component.name"),
200
})
201
r := chasm.NewRegistry(s.logger)
202
err := r.Register(lib)
203
require.Error(t, err)
204
require.Contains(t, err.Error(), "name must follow golang identifier rules")
205
})
206
207
s.T().Run("component is already registered by name", func(t *testing.T) {
208
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
209
chasm.NewRegistrableComponent[*chasm.MockComponent]("Component1"),
210
chasm.NewRegistrableComponent[*chasm.MockComponent]("Component1"),
211
})
212
r := chasm.NewRegistry(s.logger)
213
err := r.Register(lib)
214
require.Error(t, err)
215
require.Contains(t, err.Error(), "is already registered")
216
})
217
218
s.T().Run("component is already registered by type", func(t *testing.T) {
219
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
220
chasm.NewRegistrableComponent[*chasm.MockComponent]("Component1"),
221
chasm.NewRegistrableComponent[*chasm.MockComponent]("Component2"),
222
})
223
r := chasm.NewRegistry(s.logger)
224
225
err := r.Register(lib)
226
require.Error(t, err)
227
require.Contains(t, err.Error(), "is already registered")
228
})
229
230
s.T().Run("component is already registered in another library", func(t *testing.T) {
231
lib2 := chasm.NewMockLibrary(ctrl)
232
lib2.EXPECT().Name().Return("TestLibrary2").AnyTimes()
233
234
component := chasm.NewRegistrableComponent[*chasm.MockComponent]("Component1")
235
lib2.EXPECT().Components().Return([]*chasm.RegistrableComponent{
236
component,
237
})
238
lib2.EXPECT().Tasks().Return(nil)
239
lib2.EXPECT().NexusServices().Return(nil)
240
lib2.EXPECT().NexusServiceProcessors().Return(nil)
241
r2 := chasm.NewRegistry(s.logger)
242
err := r2.Register(lib2)
243
require.NoError(t, err)
244
245
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
246
component,
247
})
248
r := chasm.NewRegistry(s.logger)
249
250
err = r.Register(lib)
251
require.Error(t, err)
252
require.Contains(t, err.Error(), "is already registered in library TestLibrary2")
253
})
254
255
s.T().Run("component must be a struct", func(t *testing.T) {
256
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
257
chasm.NewRegistrableComponent[chasm.Component]("Component1"),
258
})
259
r := chasm.NewRegistry(s.logger)
260
261
err := r.Register(lib)
262
require.Error(t, err)
263
require.Contains(t, err.Error(), "must be struct or pointer to struct")
264
})
265
266
s.Run("duplicate search attribute alias panics", func() {
267
s.Require().PanicsWithValue("registrable component validation error: search attribute alias \"MyAlias\" is already defined",
268
func() {
269
chasm.NewRegistrableComponent[*chasm.MockComponent](
270
"Component1",
271
chasm.WithSearchAttributes(
272
chasm.NewSearchAttributeBool("MyAlias", chasm.SearchAttributeFieldBool01),
273
chasm.NewSearchAttributeInt("MyAlias", chasm.SearchAttributeFieldInt01),
274
),
275
)
276
},
277
)
278
})
279
280
s.Run("duplicate search attribute field panics", func() {
281
s.Require().PanicsWithValue("registrable component validation error: search attribute field \"TemporalBool01\" is already defined",
282
func() {
283
chasm.NewRegistrableComponent[*chasm.MockComponent](
284
"Component1",
285
chasm.WithSearchAttributes(
286
chasm.NewSearchAttributeBool("Alias1", chasm.SearchAttributeFieldBool01),
287
chasm.NewSearchAttributeBool("Alias2", chasm.SearchAttributeFieldBool01),
288
),
289
)
290
},
291
)
292
})
293
294
s.Run("valid search attributes do not panic", func() {
295
s.Require().NotPanics(func() {
296
chasm.NewRegistrableComponent[*chasm.MockComponent](
297
"Component1",
298
chasm.WithSearchAttributes(
299
chasm.NewSearchAttributeBool("Completed", chasm.SearchAttributeFieldBool01),
300
chasm.NewSearchAttributeInt("Count", chasm.SearchAttributeFieldInt01),
301
chasm.NewSearchAttributeKeyword("Status", chasm.SearchAttributeFieldKeyword01),
302
),
303
)
304
})
305
})
306
307
s.Run("ExecutionStatus alias is allowed for CHASM components", func() {
308
s.Require().NotPanics(func() {
309
chasm.NewRegistrableComponent[*chasm.MockComponent](
310
"Component1",
311
chasm.WithSearchAttributes(
312
chasm.NewSearchAttributeKeyword("ExecutionStatus", chasm.SearchAttributeFieldLowCardinalityKeyword01),
313
),
314
)
315
})
316
})
317
318
s.Run("TaskQueue preallocated search attribute is allowed", func() {
319
s.Require().NotPanics(func() {
320
chasm.NewRegistrableComponent[*chasm.MockComponent](
321
"Component1",
322
chasm.WithSearchAttributes(
323
chasm.SearchAttributeTaskQueue,
324
),
325
)
326
})
327
})
328
329
s.Run("CHASM system search attribute alias panics", func() {
330
s.Require().PanicsWithValue(
331
"registrable component validation error: CHASM search attribute alias \"WorkflowId\" is a CHASM system search attribute",
332
func() {
333
chasm.NewRegistrableComponent[*chasm.MockComponent](
334
"Component1",
335
chasm.WithSearchAttributes(
336
chasm.NewSearchAttributeKeyword("WorkflowId", chasm.SearchAttributeFieldKeyword01),
337
),
338
)
339
},
340
)
341
})
342
343
s.Run("identity-mapped system search attributes are registered as overrides", func() {
344
var rc *chasm.RegistrableComponent
345
s.Require().NotPanics(func() {
346
rc = chasm.NewRegistrableComponent[*chasm.MockComponent](
347
"Component1",
348
chasm.WithSearchAttributes(
349
chasm.SearchAttributeExecutionTime,
350
chasm.SearchAttributeTaskQueue,
351
),
352
)
353
})
354
mapper := rc.SearchAttributesMapper()
355
s.Require().True(mapper.IsSystemOverride(sadefs.ExecutionTime))
356
s.Require().True(mapper.IsSystemOverride(sadefs.TaskQueue))
357
s.Require().Equal(enumspb.INDEXED_VALUE_TYPE_DATETIME, mapper.OverriddenSystemFields()[sadefs.ExecutionTime])
358
s.Require().Equal(enumspb.INDEXED_VALUE_TYPE_KEYWORD, mapper.OverriddenSystemFields()[sadefs.TaskQueue])
359
360
// Overrides are recorded only in overriddenSystemFields, not the alias/field maps; the
361
// query path resolves them via the system column instead.
362
_, err := mapper.Field(sadefs.ExecutionTime)
363
s.Require().Error(err)
364
})
365
366
s.Run("component with Visibility field must have businessID alias", func() {
367
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
368
chasm.NewRegistrableComponent[*testComponentWithVisibility]("ComponentWithVis"),
369
})
370
r := chasm.NewRegistry(s.logger)
371
err := r.Register(lib)
372
s.Require().Error(err)
373
s.Require().Contains(err.Error(), "has Field[*Visibility] but no businessID alias")
374
})
375
376
s.Run("component with Visibility field and businessID alias succeeds", func() {
377
lib.EXPECT().Components().Return([]*chasm.RegistrableComponent{
378
chasm.NewRegistrableComponent[*testComponentWithVisibility](
379
"ComponentWithVis",
380
chasm.WithBusinessIDAlias("MyBusinessId"),
381
),
382
})
383
lib.EXPECT().Tasks().Return(nil)
384
lib.EXPECT().NexusServices().Return(nil)
385
lib.EXPECT().NexusServiceProcessors().Return(nil)
386
r := chasm.NewRegistry(s.logger)
387
err := r.Register(lib)
388
s.Require().NoError(err)
389
})
390
391
}
392
393
func (s *RegistryTestSuite) TestRegistry_RegisterTasks_Error() {
394
ctrl := gomock.NewController(s.T())
395
lib := chasm.NewMockLibrary(ctrl)
396
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
397
lib.EXPECT().Components().Return(nil).AnyTimes()
398
399
s.T().Run("task name must not be empty", func(t *testing.T) {
400
r := chasm.NewRegistry(s.logger)
401
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{
402
chasm.NewRegistrablePureTask[*chasm.MockComponent, testTask1](
403
"",
404
chasm.NewMockPureTaskHandler[*chasm.MockComponent, testTask1](ctrl),
405
),
406
})
407
err := r.Register(lib)
408
require.Error(t, err)
409
require.Contains(t, err.Error(), "name must not be empty")
410
})
411
412
s.T().Run("task name must follow rules", func(t *testing.T) {
413
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{
414
chasm.NewRegistrablePureTask[*chasm.MockComponent, testTask1](
415
"bad.task.name",
416
chasm.NewMockPureTaskHandler[*chasm.MockComponent, testTask1](ctrl),
417
),
418
})
419
r := chasm.NewRegistry(s.logger)
420
err := r.Register(lib)
421
require.Error(t, err)
422
require.Contains(t, err.Error(), "name must follow golang identifier rules")
423
})
424
425
s.T().Run("task is already registered by name", func(t *testing.T) {
426
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{
427
chasm.NewRegistrablePureTask[*chasm.MockComponent, testTask1](
428
"Task1",
429
chasm.NewMockPureTaskHandler[*chasm.MockComponent, testTask1](ctrl),
430
),
431
chasm.NewRegistrableSideEffectTask[*chasm.MockComponent, testTask1](
432
"Task1",
433
chasm.NewMockSideEffectTaskHandler[*chasm.MockComponent, testTask1](ctrl),
434
),
435
})
436
r := chasm.NewRegistry(s.logger)
437
err := r.Register(lib)
438
require.Error(t, err)
439
require.Contains(t, err.Error(), "is already registered")
440
})
441
442
s.T().Run("task is already registered by type", func(t *testing.T) {
443
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{
444
chasm.NewRegistrablePureTask[*chasm.MockComponent, testTask1](
445
"Task1",
446
chasm.NewMockPureTaskHandler[*chasm.MockComponent, testTask1](ctrl),
447
),
448
chasm.NewRegistrablePureTask[*chasm.MockComponent, testTask1](
449
"Task2",
450
chasm.NewMockPureTaskHandler[*chasm.MockComponent, testTask1](ctrl),
451
),
452
})
453
r := chasm.NewRegistry(s.logger)
454
err := r.Register(lib)
455
require.Error(t, err)
456
require.Contains(t, err.Error(), "is already registered")
457
})
458
459
s.Run("task is already registered in another library", func() {
460
lib2 := chasm.NewMockLibrary(ctrl)
461
lib2.EXPECT().Name().Return("TestLibrary2").AnyTimes()
462
463
lib2.EXPECT().Components().Return(nil)
464
lib2.EXPECT().NexusServices().Return(nil)
465
lib2.EXPECT().NexusServiceProcessors().Return(nil)
466
task := chasm.NewRegistrablePureTask[*chasm.MockComponent, testTask1](
467
"Task1",
468
chasm.NewMockPureTaskHandler[*chasm.MockComponent, testTask1](ctrl),
469
)
470
lib2.EXPECT().Tasks().Return([]*chasm.RegistrableTask{task})
471
r2 := chasm.NewRegistry(s.logger)
472
err := r2.Register(lib2)
473
s.Require().NoError(err)
474
475
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{task})
476
r := chasm.NewRegistry(s.logger)
477
478
err = r.Register(lib)
479
s.ErrorContains(err, "is already registered in library TestLibrary2")
480
})
481
482
s.Run("task must be struct", func() {
483
lib.EXPECT().Tasks().Return([]*chasm.RegistrableTask{
484
chasm.NewRegistrablePureTask[*chasm.MockComponent, string](
485
"Task1",
486
chasm.NewMockPureTaskHandler[*chasm.MockComponent, string](ctrl),
487
),
488
})
489
r := chasm.NewRegistry(s.logger)
490
err := r.Register(lib)
491
s.ErrorContains(err, "must be struct or pointer to struct")
492
})
493
}
494
495
func (s *RegistryTestSuite) TestRegistry_RegisterNexusServices_Success() {
496
r := chasm.NewRegistry(s.logger)
497
ctrl := gomock.NewController(s.T())
498
lib := chasm.NewMockLibrary(ctrl)
499
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
500
lib.EXPECT().Components().Return(nil)
501
lib.EXPECT().Tasks().Return(nil)
502
lib.EXPECT().NexusServiceProcessors().Return(nil)
503
504
svc1 := nexus.NewService("Service1")
505
svc2 := nexus.NewService("Service2")
506
lib.EXPECT().NexusServices().Return([]*nexus.Service{svc1, svc2})
507
508
err := r.Register(lib)
509
s.Require().NoError(err)
510
511
services := r.NexusServices()
512
s.Require().Len(services, 2)
513
s.Require().Contains(services, "Service1")
514
s.Require().Contains(services, "Service2")
515
s.Require().Equal(svc1, services["Service1"])
516
s.Require().Equal(svc2, services["Service2"])
517
}
518
519
func (s *RegistryTestSuite) TestRegistry_RegisterNexusServices_Error() {
520
ctrl := gomock.NewController(s.T())
521
lib := chasm.NewMockLibrary(ctrl)
522
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
523
lib.EXPECT().Components().Return(nil).AnyTimes()
524
lib.EXPECT().Tasks().Return(nil).AnyTimes()
525
lib.EXPECT().NexusServiceProcessors().Return(nil).AnyTimes()
526
527
s.Run("nexus service is already registered", func() {
528
svc := nexus.NewService("Service1")
529
lib.EXPECT().NexusServices().Return([]*nexus.Service{svc, svc})
530
r := chasm.NewRegistry(s.logger)
531
err := r.Register(lib)
532
s.Require().ErrorContains(err, "is already registered")
533
})
534
}
535
536
func (s *RegistryTestSuite) TestRegistry_RegisterNexusServiceProcessors() {
537
r := chasm.NewRegistry(s.logger)
538
ctrl := gomock.NewController(s.T())
539
lib := chasm.NewMockLibrary(ctrl)
540
lib.EXPECT().Name().Return("TestLibrary").AnyTimes()
541
lib.EXPECT().Components().Return(nil)
542
lib.EXPECT().Tasks().Return(nil)
543
lib.EXPECT().NexusServices().Return(nil)
544
545
proc1 := chasm.NewNexusServiceProcessor("ServiceProcessor1")
546
proc2 := chasm.NewNexusServiceProcessor("ServiceProcessor2")
547
lib.EXPECT().NexusServiceProcessors().Return([]*chasm.NexusServiceProcessor{proc1, proc2})
548
549
err := r.Register(lib)
550
s.Require().NoError(err)
551
552
// Verify the processors were registered by attempting to use them
553
// We can verify registration indirectly by trying to register them again which should fail
554
err = r.NexusEndpointProcessor.RegisterServiceProcessor(proc1)
555
s.Require().ErrorContains(err, "already registered")
556
557
err = r.NexusEndpointProcessor.RegisterServiceProcessor(proc2)
558
s.Require().ErrorContains(err, "already registered")
559
}