go.temporal.io/server/chasm/visibility_manager.go

180 LOC · 3 covered · 177 uncovered · 1 ranges · 66 concepts · 1 introducers · 13 tests

File neighbourhood

The centred file is linked to every concept that introduces one of its ranges, every test that runs code from the file, and the gray connector concepts standing between those tests and the file's own introducer concepts. Undirected links join concepts to every file where they introduce source and concepts to the tests they introduce; arrows show specialization between the displayed concepts and bridge only concepts omitted from this view. Concept colors match the source ranges below; connector concepts have no source color and are shown in gray.

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the file

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the related-file, concept, and source links on this page.

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the filego.temporal.io/server/chasm/interceptors.go · 67 LOCchasm/interceptors.goworkflow_handler.go ×11 · 117 introduced LOCworkflow_handler.go ×11common.go ×1 · 11 introduced LOCcommon.go ×1request_response.pb.go ×6 · 246 introduced LOCrequest_response.pb.go ×…visibility_store.go ×17 · 171 introduced LOCvisibility_store.go ×17telemetry.go ×2 · 15 introduced LOCtelemetry.go ×2fx.go ×1 · 4 introduced LOCfx.go ×1collector.go ×7 · 33 introduced LOCcollector.go ×7data_store_factory.go ×29 · 703 introduced LOCdata_store_factory.go ×2…TestNewServer · 0 introduced LOCTestNewServerpri_matcher.go ×1 · 2 introduced LOCpri_matcher.go ×1timer_queue_active_task_executor.go ×1 · 3 introduced LOCtimer_queue_active_task_…TestNewServer · 0 introduced LOCTestNewServermetric_client.go ×3 · 18 introduced LOCmetric_client.go ×3request_response.pb.go ×12 · 137 introduced LOCrequest_response.pb.go ×…workflow_task_completed_handler.go ×9 · 75 introduced LOCworkflow_task_completed_…pri_forwarder.go ×2 · 4 introduced LOCpri_forwarder.go ×2request_response.pb.go ×1 · 7 introduced LOCrequest_response.pb.go ×…logger.go ×1 · 5 introduced LOClogger.go ×1connections.go ×1 · 9 introduced LOCconnections.go ×1persistence_rate_limited_clients.go ×2 · 17 introduced LOCpersistence_rate_limited…logger.go ×1 · 2 introduced LOClogger.go ×1metric_client.go ×2 · 7 introduced LOCmetric_client.go ×2pri_matcher.go ×1 · 4 introduced LOCpri_matcher.go ×1task_queue_partition_manager.go ×2 · 10 introduced LOCtask_queue_partition_man…connections.go ×2 · 12 introduced LOCconnections.go ×2workflow_handler.go ×8 · 144 introduced LOCworkflow_handler.go ×8pri_matcher.go ×1 · 2 introduced LOCpri_matcher.go ×1handler.go ×1 · 20 introduced LOChandler.go ×1pri_matcher.go ×8 · 55 introduced LOCpri_matcher.go ×8matching_engine.go ×1 · 8 introduced LOCmatching_engine.go ×1server.go ×1 · 3 introduced LOCserver.go ×1logger.go ×2 · 21 introduced LOClogger.go ×2queue_scheduled.go ×1 · 2 introduced LOCqueue_scheduled.go ×1handler.go ×25 · 728 introduced LOChandler.go ×25onebox.go ×75 · 1256 introduced LOConebox.go ×75metric_client_gen.go ×4 · 50 introduced LOCmetric_client_gen.go ×4pri_forwarder.go ×1 · 6 introduced LOCpri_forwarder.go ×1matching_service_server_gen.go ×1 · 2 introduced LOCmatching_service_server_…http_api_server.go ×23 · 195 introduced LOChttp_api_server.go ×23db.go ×1 · 2 introduced LOCdb.go ×1pri_task_writer.go ×3 · 14 introduced LOCpri_task_writer.go ×3namespace_handover.go ×3 · 8 introduced LOCnamespace_handover.go ×3task_queue_partition_manager.go ×2 · 4 introduced LOCtask_queue_partition_man…matching_engine.go ×3 · 5 introduced LOCmatching_engine.go ×3service_grpc.pb.go ×19 · 357 introduced LOCservice_grpc.pb.go ×19request_response.pb.go ×6 · 114 introduced LOCrequest_response.pb.go ×…reader.go ×2 · 8 introduced LOCreader.go ×2service_grpc.pb.go ×20 · 740 introduced LOCservice_grpc.pb.go ×20endpoint_registry.go ×2 · 26 introduced LOCendpoint_registry.go ×2server.go ×3 · 13 introduced LOCserver.go ×3scanner.go ×1 · 2 introduced LOCscanner.go ×1mask_internal_error.go ×1 · 2 introduced LOCmask_internal_error.go ×…matching_engine.go ×2 · 5 introduced LOCmatching_engine.go ×2service_resolver.go ×4 · 49 introduced LOCservice_resolver.go ×4lite_server.go ×25 · 303 introduced LOClite_server.go ×25fx.go ×44 · 705 introduced LOCfx.go ×44adaptive_pool.go ×1 · 3 introduced LOCadaptive_pool.go ×1queue_immediate.go ×1 · 1 introduced LOCqueue_immediate.go ×1fx.go ×1 · 1 introduced LOCfx.go ×1service.go ×8 · 413 introduced LOCservice.go ×8rpc.go ×1 · 13 introduced LOCrpc.go ×1queue_scheduled.go ×1 · 1 introduced LOCqueue_scheduled.go ×1fx.go ×1 · 2 introduced LOCfx.go ×1fx.go ×44 · 4693 introduced LOCfx.go ×44engine.go ×1 · 7 introduced LOCengine.go ×1interceptors.go ×2 · 12 introduced LOCinterceptors.go ×2TestChasmVisibilityInterceptor_ShouldRespond · introduced test · go.temporal.io/server/chasm/TestChasmVisibilityInterceptor_ShouldRespondTestChasmVisibilityInter…TestNewServer · introduced test · go.temporal.io/server/temporal/TestNewServerTestNewServerTestNewServerWithJSONEncoding · introduced test · go.temporal.io/server/temporal/TestNewServerWithJSONEncodingTestNewServerWithJSONEnc…with_OTEL_Collector_running · introduced test · go.temporal.io/server/temporal/TestNewServerWithOTEL/with_OTEL_Collector_runningwith_OTEL_Collector_runn…without_OTEL_Collector_running · introduced test · go.temporal.io/server/temporal/TestNewServerWithOTEL/without_OTEL_Collector_runningwithout_OTEL_Collector_r…ExampleNewServer · introduced test · go.temporal.io/server/temporaltest/ExampleNewServerExampleNewServerTestBaseServerOptions · introduced test · go.temporal.io/server/temporaltest/TestBaseServerOptionsTestBaseServerOptionsTestClientWithCustomInterceptor · introduced test · go.temporal.io/server/temporaltest/TestClientWithCustomInterceptorTestClientWithCustomInte…TestDefaultWorkerOptions · introduced test · go.temporal.io/server/temporaltest/TestDefaultWorkerOptionsTestDefaultWorkerOptionsTestNewServer · introduced test · go.temporal.io/server/temporaltest/TestNewServerTestNewServerTestNewWorkerWithOptions · introduced test · go.temporal.io/server/temporaltest/TestNewWorkerWithOptionsTestNewWorkerWithOptionsTestSearchAttributeRegistration · introduced test · go.temporal.io/server/temporaltest/TestSearchAttributeRegistrationTestSearchAttributeRegis…TestWorkerServiceHealthCheck · introduced test · go.temporal.io/server/tests/testcore/TestFunctionalTestBaseSuite/TestWorkerServiceHealthCheckTestWorkerServiceHealthC…Focused file · go.temporal.io/server/chasm/visibility_manager.go · 180 LOCchasm/visibility_manager…

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the related-file, concept, and source links on this page while the interactive map is unavailable.

1 //go:generate mockgen -package $GOPACKAGE -source $GOFILE -destination visibility_manager_mock.go
2
3 package chasm
4
5 import (
6 "context"
7 "reflect"
8 "time"
9
10 commonpb "go.temporal.io/api/common/v1"
11 "go.temporal.io/api/serviceerror"
12 "go.temporal.io/server/api/visibilityservice/v1"
13 "go.temporal.io/server/common/payload"
14 "google.golang.org/protobuf/proto"
15 )
16
17 type VisibilityManager interface {
18 ListExecutions(
19 context.Context,
20 reflect.Type,
21 *ListExecutionsRequest,
22 ) (*visibilityservice.ListChasmExecutionsResponse, error)
23
24 CountExecutions(
25 context.Context,
26 reflect.Type,
27 *CountExecutionsRequest,
28 ) (*visibilityservice.CountChasmExecutionsResponse, error)
29 }
30
31 type VisibilityExecutionInfo[M proto.Message] struct {
32 BusinessID string
33 RunID string
34 StartTime time.Time
35 CloseTime time.Time
36 HistoryLength int64
37 HistorySizeBytes int64
38 StateTransitionCount int64
39 ChasmSearchAttributes SearchAttributesMap
40 CustomSearchAttributes map[string]*commonpb.Payload
41 Memo *commonpb.Memo
42 ChasmMemo M
43 }
44
45 type ListExecutionsRequest struct {
46 NamespaceName string
47 Query string
48 PageSize int
49 NextPageToken []byte
50 }
51
52 type ListExecutionsResponse[M proto.Message] struct {
53 Executions []*VisibilityExecutionInfo[M]
54 NextPageToken []byte
55 }
56
57 type CountExecutionsRequest struct {
58 NamespaceName string
59 Query string
60 }
61
62 type CountExecutionsResponse struct {
63 Count int64
64 Groups []Group
65 }
66
67 type Group struct {
68 Values []*commonpb.Payload
69 Count int64
70 }
71
72 // ListExecutions lists the executions of a CHASM archetype given an initial query.
73 // The query string can specify any combination of CHASM, custom, and predefined/system search attributes.
74 // The generic parameter C is the CHASM component type used for executions and search attribute filtering.
75 // The generic parameter M is the type of the memo payload to be unmarshaled from the execution.
76 // PageSize is required, must be greater than 0.
77 // NextPageToken is optional, set on subsequent requests to continue listing the next page of executions.
78 // Note: For CHASM executions, TemporalNamespaceDivision is the predefined search attribute
79 // that is used to identify the archetype of the execution.
80 // If the query string does not specify TemporalNamespaceDivision, the archetype C of the request will be used to filter the executions.
81 // If the initial query already specifies TemporalNamespaceDivision, the archetype C of the request will
82 // only be used to get the registered SearchAttributes.
83 func ListExecutions[C Component, M proto.Message](
84 ctx context.Context,
85 request *ListExecutionsRequest,
86 ) (*ListExecutionsResponse[M], error) {
87 archetypeType := reflect.TypeFor[C]()
88 response, err := visibilityManagerFromContext(ctx).ListExecutions(ctx, archetypeType, request)
89 if err != nil {
90 return nil, err
91 }
92
93 // Convert response: decode ChasmSearchAttributes and ChasmMemo to type M
94 executions := make([]*VisibilityExecutionInfo[M], len(response.Executions))
95 for i, execution := range response.Executions {
96 chasmSAs, err := newSearchAttributesMapFromProto(execution.ChasmSearchAttributes)
97 if err != nil {
98 return nil, err
99 }
100
101 chasmMemoInterface := reflect.New(reflect.TypeFor[M]().Elem()).Interface()
102 chasmMemo, ok := chasmMemoInterface.(M)
103 if !ok {
104 return nil, serviceerror.NewInternalf("failed to cast chasm memo to type %s", reflect.TypeFor[M]().String())
105 }
106 if err := payload.Decode(execution.ChasmMemo, chasmMemo); err != nil {
107 return nil, serviceerror.NewInternalf("failed to decode chasm memo: %v", err)
108 }
109 executions[i] = &VisibilityExecutionInfo[M]{
110 BusinessID: execution.BusinessId,
111 RunID: execution.RunId,
112 StartTime: execution.StartTime.AsTime(),
113 CloseTime: execution.CloseTime.AsTime(),
114 HistoryLength: execution.HistoryLength,
115 HistorySizeBytes: execution.HistorySizeBytes,
116 StateTransitionCount: execution.StateTransitionCount,
117 ChasmSearchAttributes: chasmSAs,
118 CustomSearchAttributes: execution.CustomSearchAttributes.GetIndexedFields(),
119 Memo: execution.Memo,
120 ChasmMemo: chasmMemo,
121 }
122 }
123
124 return &ListExecutionsResponse[M]{
125 Executions: executions,
126 NextPageToken: response.NextPageToken,
127 }, nil
128 }
129
130 // CountExecutions counts the executions of a CHASM archetype given an initial query.
131 // The generic parameter C is the CHASM component type used for executions and search attribute filtering.
132 // The query string can specify any combination of CHASM, custom, and predefined/system search attributes.
133 // Note: For CHASM executions, TemporalNamespaceDivision is the predefined search attribute
134 // that is used to identify the archetype of the execution.
135 // If the query string does not specify TemporalNamespaceDivision, the archetype C of the request will be used to count the executions.
136 // If the initial query already specifies TemporalNamespaceDivision, the archetype C of the request will
137 // only be used to get the registered SearchAttributes.
138 func CountExecutions[C Component](
139 ctx context.Context,
140 request *CountExecutionsRequest,
141 ) (*CountExecutionsResponse, error) {
142 archetypeType := reflect.TypeFor[C]()
143 visResponse, err := visibilityManagerFromContext(ctx).CountExecutions(ctx, archetypeType, request)
144 if err != nil {
145 return nil, err
146 }
147
148 response := &CountExecutionsResponse{
149 Count: visResponse.Count,
150 Groups: make([]Group, len(visResponse.Groups)),
151 }
152 for k, group := range visResponse.Groups {
153 response.Groups[k] = Group{
154 Values: group.GroupValues,
155 Count: group.Count,
156 }
157 }
158 return response, nil
159 }
160
161 type visibilityManagerCtxKeyType string
162
163 const visibilityManagerCtxKey visibilityManagerCtxKeyType = "chasmVisibilityManager"
164
165 func NewVisibilityManagerContext(
166 ctx context.Context,
167 engine VisibilityManager,
168 > ) context.Context { interceptors.go ×2
169 > return context.WithValue(ctx, visibilityManagerCtxKey, engine)
170 > }
171
172 func visibilityManagerFromContext(
173 ctx context.Context,
174 ) VisibilityManager {
175 e, ok := ctx.Value(visibilityManagerCtxKey).(VisibilityManager)
176 if !ok {
177 return nil
178 }
179 return e
180 }