execution_state_map.go ×18

Frontier kind: Code frontier

unlabeled · c_36305cceeb64

250 tests · 4305 LOC · 170 files · introduces 0 tests · 172 LOC · 3 files

Introduces — evidence that enters the hierarchy at this concept

Code
34 ranges172 lines · 3 files
Tests
0 tests

Contains — complete concept membership

All code (extent)
797 ranges4305 lines · 170 files · Browse complete extent
All tests (intent)
250 testsBrowse complete intent

Neighbourhood graph

The orange circle is the focus. Violet and green circles are every ancestor and descendant, broader and narrower, at any distance; blue squares and pink diamonds are the introduced files and exact introduced tests of every visible concept, not only the focus's. Arrows point from broader to narrower concepts and bridge only concepts omitted from this view. Undirected links show source or test introduction. Concept and file size follows LOC; exact test nodes use test-count units.

Introduced files, introduced tests, and structurally relevant concept specialization

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 native relationship evidence on this page.

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the native relationship evidence on this page while the interactive map is unavailable.

Native relationship evidence

Every exact file and test below is linked only from the concept that introduces it.

Introduced tests

Every collected test enters the hierarchy at exactly one concept.

No tests are introduced at this concept. Its intent tests are introduced by other concepts.

Introduced code

Every collected source range enters the hierarchy at exactly one concept.

3 files ranked by introduced lines: 172 introduced LOC across 34 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/common/persistence/sql/execution.go 77 introduced LOC · 10 ranges

Open complete file

222 })
223 switch err {
224 > case nil: execution.go
225 // noop
226 case sql.ErrNoRows:
230 }
231
232 > state := &p.InternalWorkflowMutableState{ execution.go
233 > ExecutionInfo: p.NewDataBlob(executionsRow.Data, executionsRow.DataEncoding),
234 > ExecutionState: p.NewDataBlob(executionsRow.State, executionsRow.StateEncoding),
235 > NextEventID: executionsRow.NextEventID,
236 >
237 > DBRecordVersion: executionsRow.DBRecordVersion,
238 > }
239 >
240 > state.ActivityInfos, err = getActivityInfoMap(ctx,
241 > m.DB,
242 > request.ShardID,
243 > namespaceID,
244 > workflowID,
245 > runID,
246 > )
247 > if err != nil {
248 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get activity info. Error: %v", err)
249 }
250
251 > state.TimerInfos, err = getTimerInfoMap(ctx, execution.go
252 > m.DB,
253 > request.ShardID,
254 > namespaceID,
255 > workflowID,
256 > runID,
257 > )
258 > if err != nil {
259 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get timer info. Error: %v", err)
260 }
261
262 > state.ChildExecutionInfos, err = getChildExecutionInfoMap(ctx, execution.go
263 > m.DB,
264 > request.ShardID,
265 > namespaceID,
266 > workflowID,
267 > runID,
268 > )
269 > if err != nil {
270 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get child executionsRow info. Error: %v", err)
271 }
272
273 > state.RequestCancelInfos, err = getRequestCancelInfoMap(ctx, execution.go
274 > m.DB,
275 > request.ShardID,
276 > namespaceID,
277 > workflowID,
278 > runID,
279 > )
280 > if err != nil {
281 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get request cancel info. Error: %v", err)
282 }
283
284 > state.SignalInfos, err = getSignalInfoMap(ctx, execution.go
285 > m.DB,
286 > request.ShardID,
287 > namespaceID,
288 > workflowID,
289 > runID,
290 > )
291 > if err != nil {
292 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get signal info. Error: %v", err)
293 }
294
295 > state.BufferedEvents, err = getBufferedEvents(ctx, execution.go
296 > m.DB,
297 > request.ShardID,
298 > namespaceID,
299 > workflowID,
300 > runID,
301 > )
302 > if err != nil {
303 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get buffered events. Error: %v", err)
304 }
305
306 > state.ChasmNodes, err = getChasmNodeMap(ctx, execution.go
307 > m.DB,
308 > request.ShardID,
309 > namespaceID,
310 > workflowID,
311 > runID,
312 > )
313 > if err != nil {
314 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get CHASM nodes. Error: %v", err)
315 }
316
317 > state.SignalRequestedIDs, err = getSignalsRequested(ctx, execution.go
318 > m.DB,
319 > request.ShardID,
320 > namespaceID,
321 > workflowID,
322 > runID,
323 > )
324 > if err != nil {
325 return nil, serviceerror.NewUnavailablef("GetWorkflowExecution: failed to get signals requested. Error: %v", err)
326 }
327
328 > return &p.InternalGetWorkflowExecutionResponse{ execution.go
329 > State: state,
330 > DBRecordVersion: executionsRow.DBRecordVersion,
331 > }, nil
332 }
333
go.temporal.io/server/common/persistence/sql/execution_state_map.go 71 introduced LOC · 18 ranges

Open complete file

65 workflowID string,
66 runID primitives.UUID,
67 > ) (map[int64]*commonpb.DataBlob, error) { execution_state_map.go
68 >
69 > rows, err := db.SelectAllFromActivityInfoMaps(ctx, sqlplugin.ActivityInfoMapsAllFilter{
70 > ShardID: shardID,
71 > NamespaceID: namespaceID,
72 > WorkflowID: workflowID,
73 > RunID: runID,
74 > })
75 > if err != nil && err != sql.ErrNoRows {
76 return nil, serviceerror.NewUnavailablef("Failed to get activity info. Error: %v", err)
77 }
78
79 > ret := make(map[int64]*commonpb.DataBlob) execution_state_map.go
80 > for _, row := range rows {
81 ret[row.ScheduleID] = persistence.NewDataBlob(row.Data, row.DataEncoding)
82 }
83
84 > return ret, nil execution_state_map.go
85 }
86
155 workflowID string,
156 runID primitives.UUID,
157 > ) (map[string]*commonpb.DataBlob, error) { execution_state_map.go
158 >
159 > rows, err := db.SelectAllFromTimerInfoMaps(ctx, sqlplugin.TimerInfoMapsAllFilter{
160 > ShardID: shardID,
161 > NamespaceID: namespaceID,
162 > WorkflowID: workflowID,
163 > RunID: runID,
164 > })
165 > if err != nil && err != sql.ErrNoRows {
166 return nil, serviceerror.NewUnavailablef("Failed to get timer info. Error: %v", err)
167 }
168 > ret := make(map[string]*commonpb.DataBlob) execution_state_map.go
169 > for _, row := range rows {
170 ret[row.TimerID] = persistence.NewDataBlob(row.Data, row.DataEncoding)
171 }
172
173 > return ret, nil execution_state_map.go
174 }
175
244 workflowID string,
245 runID primitives.UUID,
246 > ) (map[int64]*commonpb.DataBlob, error) { execution_state_map.go
247 >
248 > rows, err := db.SelectAllFromChildExecutionInfoMaps(ctx, sqlplugin.ChildExecutionInfoMapsAllFilter{
249 > ShardID: shardID,
250 > NamespaceID: namespaceID,
251 > WorkflowID: workflowID,
252 > RunID: runID,
253 > })
254 > if err != nil && err != sql.ErrNoRows {
255 return nil, serviceerror.NewUnavailablef("Failed to get timer info. Error: %v", err)
256 }
257
258 > ret := make(map[int64]*commonpb.DataBlob) execution_state_map.go
259 > for _, row := range rows {
260 ret[row.InitiatedID] = persistence.NewDataBlob(row.Data, row.DataEncoding)
261 }
262
263 > return ret, nil execution_state_map.go
264 }
265
335 workflowID string,
336 runID primitives.UUID,
337 > ) (map[int64]*commonpb.DataBlob, error) { execution_state_map.go
338 >
339 > rows, err := db.SelectAllFromRequestCancelInfoMaps(ctx, sqlplugin.RequestCancelInfoMapsAllFilter{
340 > ShardID: shardID,
341 > NamespaceID: namespaceID,
342 > WorkflowID: workflowID,
343 > RunID: runID,
344 > })
345 > if err != nil && err != sql.ErrNoRows {
346 return nil, serviceerror.NewUnavailablef("Failed to get request cancel info. Error: %v", err)
347 }
348
349 > ret := make(map[int64]*commonpb.DataBlob) execution_state_map.go
350 > for _, row := range rows {
351 ret[row.InitiatedID] = persistence.NewDataBlob(row.Data, row.DataEncoding)
352 }
353
354 > return ret, nil execution_state_map.go
355 }
356
426 workflowID string,
427 runID primitives.UUID,
428 > ) (map[int64]*commonpb.DataBlob, error) { execution_state_map.go
429 >
430 > rows, err := db.SelectAllFromSignalInfoMaps(ctx, sqlplugin.SignalInfoMapsAllFilter{
431 > ShardID: shardID,
432 > NamespaceID: namespaceID,
433 > WorkflowID: workflowID,
434 > RunID: runID,
435 > })
436 > if err != nil && err != sql.ErrNoRows {
437 return nil, serviceerror.NewUnavailablef("Failed to get signal info. Error: %v", err)
438 }
439
440 > ret := make(map[int64]*commonpb.DataBlob) execution_state_map.go
441 > for _, row := range rows {
442 ret[row.InitiatedID] = persistence.NewDataBlob(row.Data, row.DataEncoding)
443 }
444
445 > return ret, nil execution_state_map.go
446 }
447
521 workflowID string,
522 runID primitives.UUID,
523 > ) (map[string]persistence.InternalChasmNode, error) { execution_state_map.go
524 > rows, err := db.SelectAllFromChasmNodeMaps(ctx, sqlplugin.ChasmNodeMapsAllFilter{
525 > ShardID: shardID,
526 > NamespaceID: namespaceID,
527 > WorkflowID: workflowID,
528 > RunID: runID,
529 > })
530 > if err != nil && err != sql.ErrNoRows {
531 return nil, serviceerror.NewUnavailablef("Failed to get CHASM nodes. Error: %v", err)
532 }
533
534 > ret := make(map[string]persistence.InternalChasmNode) execution_state_map.go
535 > for _, row := range rows {
536 ret[row.ChasmPath] = persistence.InternalChasmNode{
537 Metadata: persistence.NewDataBlob(row.Metadata, row.MetadataEncoding),
go.temporal.io/server/common/persistence/sql/execution_state_non_map.go 24 introduced LOC · 6 ranges

Open complete file

61 workflowID string,
62 runID primitives.UUID,
63 > ) ([]string, error) { execution_state_non_map.go
64 >
65 > rows, err := db.SelectAllFromSignalsRequestedSets(ctx, sqlplugin.SignalsRequestedSetsAllFilter{
66 > ShardID: shardID,
67 > NamespaceID: namespaceID,
68 > WorkflowID: workflowID,
69 > RunID: runID,
70 > })
71 > if err != nil && err != sql.ErrNoRows {
72 return nil, serviceerror.NewUnavailablef("Failed to get signals requested. Error: %v", err)
73 }
74 > var ret = make([]string, len(rows)) execution_state_non_map.go
75 > for i, s := range rows {
76 ret[i] = s.SignalID
77 }
78 > return ret, nil execution_state_non_map.go
79 }
80
134 workflowID string,
135 runID primitives.UUID,
136 > ) ([]*commonpb.DataBlob, error) { execution_state_non_map.go
137 >
138 > rows, err := db.SelectFromBufferedEvents(ctx, sqlplugin.BufferedEventsFilter{
139 > ShardID: shardID,
140 > NamespaceID: namespaceID,
141 > WorkflowID: workflowID,
142 > RunID: runID,
143 > })
144 > if err != nil && err != sql.ErrNoRows {
145 return nil, serviceerror.NewUnavailablef("getBufferedEvents operation failed. Select failed: %v", err)
146 }
147 > var result []*commonpb.DataBlob execution_state_non_map.go
148 > for _, row := range rows {
149 result = append(result, p.NewDataBlob(row.Data, row.DataEncoding))
150 }
151 > return result, nil execution_state_non_map.go
152 }
153