workflow.go ×10

Frontier kind: Joint frontier

unlabeled · c_c641d2b26a47

1 test · 4276 LOC · 187 files · introduces 1 test · 33 LOC · 2 files

Introduces — evidence that enters the hierarchy at this concept

Code
13 ranges33 lines · 2 files
Tests
1 test

Contains — complete concept membership

All code (extent)
639 ranges4276 lines · 187 files · Browse complete extent
All tests (intent)
1 testBrowse 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.

1 test introduced at this concept.

Introduced code

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

2 files ranked by introduced lines: 33 introduced LOC across 13 ranges. Expand a file to inspect source; the > gutter marks introduced lines.

go.temporal.io/server/service/history/ndc/workflow.go 19 introduced LOC · 10 ranges

Open complete file

185 }
186
187 > func (r *WorkflowImpl) FlushBufferedEvents() error { workflow.go
188 >
189 > if !r.mutableState.IsWorkflow() {
190 return nil
191 }
192
193 > if !r.mutableState.IsWorkflowExecutionRunning() { workflow.go
194 return nil
195 }
196
197 > if !r.mutableState.HasBufferedEvents() { workflow.go
198 return nil
199 }
203 // events for state only changes as well.
204 // Transition history is not enabled today so LastWriteVersion == LastEventVersion
205 > lastWriteVersion, err := r.mutableState.GetLastWriteVersion() workflow.go
206 > if err != nil {
207 return err
208 }
209
210 > lastWriteCluster := r.clusterMetadata.ClusterNameForFailoverVersion(true, lastWriteVersion) workflow.go
211 > currentCluster := r.clusterMetadata.GetCurrentClusterName()
212 >
213 > if lastWriteCluster != currentCluster {
214 return serviceerror.NewInternal("Workflow encountered workflow with buffered events but last write not from current cluster")
215 }
216
217 > if err := r.mutableState.UpdateCurrentVersion(lastWriteVersion, true); err != nil { workflow.go
218 return err
219 }
220
221 > if _, err = r.failWorkflowTask(); err != nil { workflow.go
222 return err
223 }
224
225 // Don't schedule a new workflow task if the workflow is paused.
226 > if r.mutableState.IsWorkflowExecutionStatusPaused() { workflow.go
227 return nil
228 }
229
230 > if _, err := r.mutableState.AddWorkflowTaskScheduledEvent( workflow.go
231 > false,
232 > enumsspb.WORKFLOW_TASK_TYPE_NORMAL,
233 > ); err != nil {
234 return err
235 }
236 > return nil workflow.go
237 }
238
go.temporal.io/server/service/history/ndc/buffer_event_flusher.go 14 introduced LOC · 3 ranges

Open complete file

63 }
64
65 > targetWorkflow := NewWorkflow( buffer_event_flusher.go
66 > r.clusterMetadata,
67 > r.wfContext,
68 > r.mutableState,
69 > wcache.NoopReleaseFn,
70 > )
71 > if err := targetWorkflow.FlushBufferedEvents(); err != nil {
72 return nil, nil, err
73 }
74
75 // the workflow must be updated as active, to send out replication tasks
76 > if err := targetWorkflow.context.UpdateWorkflowExecutionAsActive( buffer_event_flusher.go
77 > ctx,
78 > r.shardContext,
79 > ); err != nil {
80 return nil, nil, err
81 }
82
83 > r.wfContext = targetWorkflow.GetContext() buffer_event_flusher.go
84 > r.mutableState = targetWorkflow.GetMutableState()
85 > return r.wfContext, r.mutableState, nil
86 }