go.temporal.io/server/service/matching/reachability.go

349 LOC · 168 covered · 181 uncovered · 49 ranges · 83 concepts · 16 introducers · 20 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/common/log/tag/tags.go · 1039 LOCtag/tags.gogo.temporal.io/server/common/worker_versioning/worker_versioning.go · 1327 LOCworker_versioning/worker…workflow_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 ×44TestGetReachability_WithVisibility_WithDeletedRules · 0 introduced LOCTestGetReachability_With…reachability.go ×1 · 2 introduced LOCreachability.go ×1reachability.go ×1 · 3 introduced LOCreachability.go ×1TestGetReachability_WithVisibility_WithRules · 0 introduced LOCTestGetReachability_With…reachability.go ×1 · 2 introduced LOCreachability.go ×1reachability.go ×3 · 8 introduced LOCreachability.go ×3reachability.go ×1 · 1 introduced LOCreachability.go ×1reachability.go ×18 · 86 introduced LOCreachability.go ×18reachability.go ×2 · 2 introduced LOCreachability.go ×2reachability.go ×1 · 2 introduced LOCreachability.go ×1reachability.go ×3 · 5 introduced LOCreachability.go ×3reachability.go ×1 · 8 introduced LOCreachability.go ×1TestIsReachableAssignmentRuleTarget · 0 introduced LOCTestIsReachableAssignmen…reachability.go ×1 · 2 introduced LOCreachability.go ×1reachability.go ×2 · 6 introduced LOCreachability.go ×2reachability.go ×2 · 5 introduced LOCreachability.go ×2reachability.go ×3 · 6 introduced LOCreachability.go ×3reachability.go ×3 · 5 introduced LOCreachability.go ×3reachability.go ×6 · 34 introduced LOCreachability.go ×6TestGetBuildIdsOfInterest · introduced test · go.temporal.io/server/service/matching/TestGetBuildIdsOfInterestTestGetBuildIdsOfInteres…TestGetDefaultBuildId · introduced test · go.temporal.io/server/service/matching/TestGetDefaultBuildIdTestGetDefaultBuildIdTestGetReachability_WithVisibility_WithDeletedRules · introduced test · go.temporal.io/server/service/matching/TestGetReachability_WithVisibility_WithDeletedRulesTestGetReachability_With…TestGetReachability_WithVisibility_WithRules · introduced test · go.temporal.io/server/service/matching/TestGetReachability_WithVisibility_WithRulesTestGetReachability_With…TestGetReachability_WithVisibility_WithoutRules · introduced test · go.temporal.io/server/service/matching/TestGetReachability_WithVisibility_WithoutRulesTestGetReachability_With…TestGetReachability_WithoutVisibility_WithRules · introduced test · go.temporal.io/server/service/matching/TestGetReachability_WithoutVisibility_WithRulesTestGetReachability_With…TestIsReachableAssignmentRuleTarget · introduced test · go.temporal.io/server/service/matching/TestIsReachableAssignmentRuleTargetTestIsReachableAssignmen…TestMakeBuildIdQuery · introduced test · go.temporal.io/server/service/matching/TestMakeBuildIdQueryTestMakeBuildIdQueryTestNewServer · 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/service/matching/reachability.go · 349 LOCmatching/reachability.go

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 package matching
2
3 import (
4 "context"
5 "fmt"
6 "slices"
7 "strings"
8 "time"
9
10 "github.com/temporalio/sqlparser"
11 enumspb "go.temporal.io/api/enums/v1"
12 clockspb "go.temporal.io/server/api/clock/v1"
13 persistencespb "go.temporal.io/server/api/persistence/v1"
14 "go.temporal.io/server/common/cache"
15 hlc "go.temporal.io/server/common/clock/hybrid_logical_clock"
16 "go.temporal.io/server/common/log"
17 "go.temporal.io/server/common/log/tag"
18 "go.temporal.io/server/common/metrics"
19 "go.temporal.io/server/common/namespace"
20 "go.temporal.io/server/common/persistence/visibility/manager"
21 "go.temporal.io/server/common/searchattribute/sadefs"
22 "go.temporal.io/server/common/tqid"
23 "go.temporal.io/server/common/util"
24 "go.temporal.io/server/common/worker_versioning"
25 )
26
27 const (
28 reachabilityCacheMaxSize = 10000
29 )
30
31 type reachabilityExitPoint int32
32
33 const (
34 checkedRuleSourcesForInput reachabilityExitPoint = 0
35 checkedRuleTargetsForUpstream reachabilityExitPoint = 1
36 checkedBacklogForUpstream reachabilityExitPoint = 2
37 checkedOpenWorkflowExecutionsForUpstreamHit reachabilityExitPoint = 3
38 checkedOpenWorkflowExecutionsForUpstreamMiss reachabilityExitPoint = 4
39 checkedClosedWorkflowExecutionsForUpstreamHit reachabilityExitPoint = 5
40 checkedClosedWorkflowExecutionsForUpstreamMiss reachabilityExitPoint = 6
41 reachabilityExitPointTagName = "reachability_exit_point"
42 )
43
44 var (
45 reachabilityExitPoint2TagValue = map[reachabilityExitPoint]string{
46 checkedRuleSourcesForInput: "checked_rule_sources_for_input",
47 checkedRuleTargetsForUpstream: "checked_rule_targets_for_upstream",
48 checkedBacklogForUpstream: "checked_backlog_for_upstream",
49 checkedOpenWorkflowExecutionsForUpstreamHit: "checked_open_wf_executions_for_upstream_hit",
50 checkedOpenWorkflowExecutionsForUpstreamMiss: "checked_open_wf_executions_for_upstream_miss",
51 checkedClosedWorkflowExecutionsForUpstreamHit: "checked_closed_wf_executions_for_upstream_hit",
52 checkedClosedWorkflowExecutionsForUpstreamMiss: "checked_closed_wf_executions_for_upstream_miss",
53 }
54 )
55
56 type reachabilityCalculator struct {
57 cache reachabilityCache
58 nsID namespace.ID
59 nsName namespace.Name
60 taskQueue *tqid.TaskQueueFamily
61 assignmentRules []*persistencespb.AssignmentRule
62 redirectRules []*persistencespb.RedirectRule
63 buildIdVisibilityGracePeriod time.Duration
64 tqConfig *taskQueueConfig
65 }
66
67 func newReachabilityCalculator(
68 data *persistencespb.VersioningData,
69 rCache reachabilityCache,
70 nsID, nsName string,
71 taskQueue *tqid.TaskQueueFamily,
72 buildIdVisibilityGracePeriod time.Duration,
73 tqConfig *taskQueueConfig,
74 ) *reachabilityCalculator {
75 return &reachabilityCalculator{
76 cache: rCache,
77 nsID: namespace.ID(nsID),
78 nsName: namespace.Name(nsName),
79 taskQueue: taskQueue,
80 assignmentRules: data.GetAssignmentRules(),
81 redirectRules: data.GetRedirectRules(),
82 buildIdVisibilityGracePeriod: buildIdVisibilityGracePeriod,
83 tqConfig: tqConfig,
84 }
85 }
86
87 func getBuildIdTaskReachability(
88 ctx context.Context,
89 rc *reachabilityCalculator,
90 metricsHandler metrics.Handler,
91 logger log.Logger,
92 buildId string,
93 > ) (enumspb.BuildIdTaskReachability, error) { reachability.go ×18
94 > reachability, exitPoint, err := rc.run(ctx, buildId)
95 > handler := metrics.GetPerTaskQueueFamilyScope(metricsHandler, rc.nsName.String(), rc.taskQueue, rc.tqConfig.BreakdownMetricsByTaskQueue())
96 > metrics.ReachabilityExitPointCounter.With(handler).Record(1,
97 > metrics.WorkerVersionTag(buildId, rc.tqConfig.BreakdownMetricsByBuildID()),
98 > metrics.StringTag(reachabilityExitPointTagName, reachabilityExitPoint2TagValue[exitPoint]))
99 > logger.Info("Calculated reachability for build id",
100 > tag.WorkerVersion(buildId),
101 > tag.BuildIdTaskReachabilityTag(reachability.String()),
102 > tag.ReachabilityExitPointTag(reachabilityExitPoint2TagValue[exitPoint]),
103 > tag.WorkflowNamespace(rc.nsName.String()),
104 > tag.WorkflowTaskQueueName(rc.taskQueue.Name()),
105 > )
106 > return reachability, err
107 > }
108
109 > func (rc *reachabilityCalculator) run(ctx context.Context, buildId string) (enumspb.BuildIdTaskReachability, reachabilityExitPoint, error) { reachability.go ×18
110 > // 1. Easy UNREACHABLE case
111 > if isActiveRedirectRuleSource(buildId, rc.redirectRules) {
112 > return enumspb.BUILD_ID_TASK_REACHABILITY_UNREACHABLE, checkedRuleSourcesForInput, nil reachability.go ×1
113 > }
114
115 // Gather list of all build ids that could point to buildId
116 > buildIdsOfInterest := rc.getBuildIdsOfInterest(buildId, time.Duration(0)) reachability.go ×18
117 >
118 > // 2. Cases for REACHABLE
119 > // 2a. If buildId is assignable to new tasks
120 > if slices.ContainsFunc(buildIdsOfInterest, rc.isReachableActiveAssignmentRuleTargetOrDefault) {
121 > return enumspb.BUILD_ID_TASK_REACHABILITY_REACHABLE, checkedRuleTargetsForUpstream, nil reachability.go ×1
122 > }
123
124 // 2b. If buildId could be reached from the backlog
125 > if existsBacklog, err := rc.existsBackloggedActivityOrWFTaskAssignedToAny(ctx, buildIdsOfInterest); err != nil { reachability.go ×18
126 return enumspb.BUILD_ID_TASK_REACHABILITY_UNSPECIFIED, checkedBacklogForUpstream, err
127 > } else if existsBacklog { reachability.go ×18
128 return enumspb.BUILD_ID_TASK_REACHABILITY_REACHABLE, checkedBacklogForUpstream, nil
129 }
130
131 // Note: The below cases are not applicable to activity-only task queues, since we don't record those in visibility
132
133 // Gather list of all build ids that could point to buildId, now including deleted rules to account for the delay in updating visibility
134 > buildIdsOfInterest = rc.getBuildIdsOfInterest(buildId, rc.buildIdVisibilityGracePeriod) reachability.go ×18
135 >
136 > // 2c. If buildId is assignable to tasks from open workflows
137 > existsOpenWFAssignedToBuildId, hit, err := rc.existsWFAssignedToAny(ctx, buildIdsOfInterest, true)
138 > if err != nil {
139 return enumspb.BUILD_ID_TASK_REACHABILITY_UNSPECIFIED, checkedOpenWorkflowExecutionsForUpstreamMiss, err
140 }
141 > if existsOpenWFAssignedToBuildId { reachability.go ×18
142 > return enumspb.BUILD_ID_TASK_REACHABILITY_REACHABLE, getCacheExitPoint(true, hit), nil reachability.go ×3
143 > }
144
145 // 3. Cases for CLOSED_WORKFLOWS_ONLY
146 > existsClosedWFAssignedToBuildId, hit, err := rc.existsWFAssignedToAny(ctx, buildIdsOfInterest, false) reachability.go ×18
147 > if err != nil {
148 return enumspb.BUILD_ID_TASK_REACHABILITY_UNSPECIFIED, checkedClosedWorkflowExecutionsForUpstreamMiss, err
149 }
150 > if existsClosedWFAssignedToBuildId { reachability.go ×18
151 > return enumspb.BUILD_ID_TASK_REACHABILITY_CLOSED_WORKFLOWS_ONLY, getCacheExitPoint(false, hit), nil reachability.go ×3
152 > }
153
154 // 4. Otherwise, UNREACHABLE
155 > return enumspb.BUILD_ID_TASK_REACHABILITY_UNREACHABLE, getCacheExitPoint(false, hit), nil reachability.go ×1
156 }
157
158 > func getCacheExitPoint(open, hit bool) reachabilityExitPoint { reachability.go ×18
159 > if open {
160 > if hit { reachability.go ×3
161 > return checkedOpenWorkflowExecutionsForUpstreamHit
162 > }
163 > return checkedOpenWorkflowExecutionsForUpstreamMiss
164 }
165 > if hit { reachability.go ×18
166 > return checkedClosedWorkflowExecutionsForUpstreamHit
167 > }
168 > return checkedClosedWorkflowExecutionsForUpstreamMiss
169 }
170
171 // getBuildIdsOfInterest returns a list of build ids that point to the given buildId in the graph
172 // of redirect rules and adds the given build ID to that list.
173 // It considers rules if the deletion time is nil or within the given deletedRuleInclusionPeriod.
174 func (rc *reachabilityCalculator) getBuildIdsOfInterest(
175 buildId string,
176 > deletedRuleInclusionPeriod time.Duration) []string { reachability.go ×3
177 >
178 > withinRuleInclusionPeriod := func(clk *clockspb.HybridLogicalClock) bool {
179 > if clk == nil { reachability.go ×2
180 return true
181 }
182 > return hlc.Since(clk) <= deletedRuleInclusionPeriod reachability.go ×2
183 }
184
185 > includedRules := util.FilterSlice(slices.Clone(rc.redirectRules), func(rr *persistencespb.RedirectRule) bool { reachability.go ×3
186 > return rr.DeleteTimestamp == nil || withinRuleInclusionPeriod(rr.DeleteTimestamp) reachability.go ×1
187 > })
188
189 > return append(getUpstreamBuildIds(buildId, includedRules), buildId) reachability.go ×3
190 }
191
192 > func (rc *reachabilityCalculator) existsBackloggedActivityOrWFTaskAssignedToAny(ctx context.Context, buildIdsOfInterest []string) (bool, error) { reachability.go ×18
193 > // todo backlog
194 > return false, nil
195 > }
196
197 > func (rc *reachabilityCalculator) isReachableActiveAssignmentRuleTargetOrDefault(buildId string) bool { reachability.go ×3
198 > foundFullyRampedRule := false
199 > for _, r := range getActiveAssignmentRules(rc.assignmentRules) {
200 > if r.GetRule().GetTargetBuildId() == buildId { reachability.go ×2
201 > return true reachability.go ×1
202 > }
203 > if isFullyRamped(r.GetRule()) { reachability.go ×2
204 > // rules after a fully-ramped rule will not be reached
205 > foundFullyRampedRule = true
206 > break
207 }
208 }
209 > if !foundFullyRampedRule && buildId == "" { reachability.go ×3
210 > // unversioned is the default, and is reachable reachability.go ×1
211 > return true
212 > }
213 > return false reachability.go ×3
214 }
215
216 func (rc *reachabilityCalculator) existsWFAssignedToAny(
217 ctx context.Context,
218 buildIdsOfInterest []string,
219 open bool,
220 > ) (exists, hit bool, err error) { reachability.go ×18
221 > query := rc.makeBuildIdQuery(buildIdsOfInterest, open)
222 > return rc.cache.Get(ctx, *rc.makeBuildIdCountRequest(query), open)
223 > }
224
225 > func (rc *reachabilityCalculator) makeBuildIdCountRequest(query string) *manager.CountWorkflowExecutionsRequest { reachability.go ×18
226 > return &manager.CountWorkflowExecutionsRequest{
227 > NamespaceID: rc.nsID,
228 > Namespace: rc.nsName,
229 > Query: query,
230 > }
231 > }
232
233 func (rc *reachabilityCalculator) makeBuildIdQuery(
234 buildIdsOfInterest []string,
235 open bool,
236 > ) string { reachability.go ×6
237 > slices.Sort(buildIdsOfInterest)
238 > escapedTaskQueue := sqlparser.String(sqlparser.NewStrVal([]byte(rc.taskQueue.Name())))
239 > var statusFilter string
240 > var escapedBuildIds []string
241 > var includeNull bool
242 > if open {
243 > statusFilter = fmt.Sprintf(` AND %s = "Running"`, sadefs.ExecutionStatus)
244 > // want: currently assigned to that build-id
245 > // (b1, b2) --> (assigned:b1, assigned:b2)
246 > // (b1, b2, "") --> (assigned:b1, assigned:b2, unversioned, null)
247 > // ("") --> (unversioned, null)
248 > for _, bid := range buildIdsOfInterest {
249 > if bid == "" {
250 > escapedBuildIds = append(escapedBuildIds, sqlparser.String(sqlparser.NewStrVal([]byte(worker_versioning.UnversionedSearchAttribute)))) reachability.go ×3
251 > includeNull = true
252 > } else { reachability.go ×6
253 > escapedBuildIds = append(escapedBuildIds, sqlparser.String(sqlparser.NewStrVal([]byte(worker_versioning.AssignedBuildIdSearchAttribute(bid)))))
254 > }
255 }
256 > } else { reachability.go ×6
257 > statusFilter = fmt.Sprintf(` AND %s != "Running"`, sadefs.ExecutionStatus)
258 > // want: closed AT that build ID, and once used that build ID
259 > // (b1, b2) --> (versioned:b1, versioned:b2)
260 > // (b1, b2, "") --> (versioned:b1, versioned:b2, unversioned, null)
261 > // ("") --> (unversioned, null)
262 > for _, bid := range buildIdsOfInterest {
263 > if bid == "" {
264 > escapedBuildIds = append(escapedBuildIds, sqlparser.String(sqlparser.NewStrVal([]byte(worker_versioning.UnversionedSearchAttribute)))) reachability.go ×3
265 > includeNull = true
266 > } else { reachability.go ×6
267 > escapedBuildIds = append(escapedBuildIds, sqlparser.String(sqlparser.NewStrVal([]byte(worker_versioning.VersionedBuildIdSearchAttribute(bid)))))
268 > }
269 }
270 }
271 > buildIdsFilter := fmt.Sprintf("%s IN (%s)", sadefs.BuildIds, strings.Join(escapedBuildIds, ",")) reachability.go ×6
272 > if includeNull {
273 > buildIdsFilter = fmt.Sprintf("(%s IS NULL OR %s)", sadefs.BuildIds, buildIdsFilter) reachability.go ×3
274 > }
275 > return fmt.Sprintf("%s = %s AND %s%s", sadefs.TaskQueue, escapedTaskQueue, buildIdsFilter, statusFilter) reachability.go ×6
276 }
277
278 // getDefaultBuildId gets the build ID mentioned in the first fully-ramped Assignment Rule.
279 // If there is no default Build ID, the result for the unversioned queue will be returned.
280 // This should only be called on the root.
281 > func getDefaultBuildId(assignmentRules []*persistencespb.AssignmentRule) string { reachability.go ×2
282 > for _, ar := range getActiveAssignmentRules(assignmentRules) {
283 > if isFullyRamped(ar.GetRule()) {
284 > return ar.GetRule().GetTargetBuildId()
285 > }
286 }
287 > return "" reachability.go ×2
288 }
289
290 /*
291 In-memory Reachability Cache of Visibility Queries and Results
292 */
293
294 type reachabilityCache struct {
295 openWFCache cache.Cache
296 closedWFCache cache.Cache // these are separate due to allow for different TTL
297 metricsHandler metrics.Handler
298 visibilityMgr manager.VisibilityManager
299 }
300
301 func newReachabilityCache(
302 handler metrics.Handler,
303 visibilityMgr manager.VisibilityManager,
304 reachabilityCacheOpenWFExecutionTTL,
305 reachabilityCacheClosedWFExecutionTTL time.Duration,
306 > ) reachabilityCache { reachability.go ×1
307 > return reachabilityCache{
308 > openWFCache: cache.New(reachabilityCacheMaxSize, &cache.Options{TTL: reachabilityCacheOpenWFExecutionTTL}),
309 > closedWFCache: cache.New(reachabilityCacheMaxSize, &cache.Options{TTL: reachabilityCacheClosedWFExecutionTTL}),
310 > metricsHandler: handler,
311 > visibilityMgr: visibilityMgr,
312 > }
313 > }
314
315 // Get retrieves the Workflow Count existence value based on the query-string key.
316 > func (c *reachabilityCache) Get(ctx context.Context, countRequest manager.CountWorkflowExecutionsRequest, open bool) (exists, hit bool, err error) { reachability.go ×18
317 > // try cache
318 > var result any
319 > if open {
320 > result = c.openWFCache.Get(countRequest)
321 > } else {
322 > result = c.closedWFCache.Get(countRequest)
323 > }
324 > if result != nil {
325 > // there's no reason that the cache would ever contain a non-bool, but just in case, treat non-bool as a miss
326 > exists, ok := result.(bool)
327 > if ok {
328 > return exists, true, nil
329 > }
330 }
331
332 // cache was cold, ask visibility and put result in cache
333 > countResponse, err := c.visibilityMgr.CountWorkflowExecutions(ctx, &countRequest) reachability.go ×18
334 > if err != nil {
335 return false, false, err
336 }
337 > exists = countResponse.Count > 0 reachability.go ×18
338 > c.Put(countRequest, exists, open)
339 > return exists, false, nil
340 }
341
342 // Put adds an element to the cache.
343 > func (c *reachabilityCache) Put(countRequest manager.CountWorkflowExecutionsRequest, exists, open bool) { reachability.go ×18
344 > if open {
345 > c.openWFCache.Put(countRequest, exists)
346 > } else {
347 > c.closedWFCache.Put(countRequest, exists)
348 > }
349 }