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 }