go.temporal.io/server/tools/tdbg/commands.go

1307 LOC · 241 covered · 1066 uncovered · 72 ranges · 37 concepts · 31 introducers · 19 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/tools/tdbg/tdbg_commands.go · 1028 LOCtdbg/tdbg_commands.gogo.temporal.io/server/tools/tdbg/util.go · 322 LOCtdbg/util.gocommands.go ×2 · 7 introduced LOCcommands.go ×2TestMigrateSchedule_FromVisibility_Execute, TestMigrateSchedule_FromVisibility_Workers · 0 introduced LOCTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_DryRun · 0 introduced LOCTestMigrateSchedule_From…TestMigrateSchedule_Stdin_OutputLog · 0 introduced LOCTestMigrateSchedule_Stdi…commands.go ×1 · 1 introduced LOCcommands.go ×1TestMigrateSchedule_Stdin_Workers · 0 introduced LOCTestMigrateSchedule_Stdi…commands.go ×3 · 11 introduced LOCcommands.go ×3commands.go ×2 · 4 introduced LOCcommands.go ×2TestMigrateSchedule_Stdin_DryRun · 0 introduced LOCTestMigrateSchedule_Stdi…TestMigrateSchedule_FromVisibility_CustomQuery · 0 introduced LOCTestMigrateSchedule_From…commands.go ×4 · 13 introduced LOCcommands.go ×4commands.go ×2 · 2 introduced LOCcommands.go ×2commands.go ×11 · 44 introduced LOCcommands.go ×11commands.go ×9 · 45 introduced LOCcommands.go ×9commands.go ×5 · 28 introduced LOCcommands.go ×5commands.go ×2 · 10 introduced LOCcommands.go ×2commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×2 · 4 introduced LOCcommands.go ×2commands.go ×3 · 8 introduced LOCcommands.go ×3commands.go ×1 · 3 introduced LOCcommands.go ×1commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×1 · 1 introduced LOCcommands.go ×1commands.go ×2 · 4 introduced LOCcommands.go ×2commands.go ×2 · 7 introduced LOCcommands.go ×2commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×3 · 18 introduced LOCcommands.go ×3commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×1 · 1 introduced LOCcommands.go ×1commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×1 · 1 introduced LOCcommands.go ×1commands.go ×4 · 15 introduced LOCcommands.go ×4commands.go ×1 · 3 introduced LOCcommands.go ×1commands.go ×1 · 3 introduced LOCcommands.go ×1commands.go ×1 · 1 introduced LOCcommands.go ×1commands.go ×1 · 2 introduced LOCcommands.go ×1commands.go ×1 · 3 introduced LOCcommands.go ×1different_case · introduced test · go.temporal.io/server/tools/tdbg/TestGetCategory/different_casedifferent_casenot_found · introduced test · go.temporal.io/server/tools/tdbg/TestGetCategory/not_foundnot_foundsame_case · introduced test · go.temporal.io/server/tools/tdbg/TestGetCategory/same_casesame_caseTestMigrateSchedule_FromVisibility_CustomQuery · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_CustomQueryTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_DefaultQueryToChasm · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_DefaultQueryToChasmTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_DryRun · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_DryRunTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_Execute · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_ExecuteTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_OutputLog · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_OutputLogTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_RejectsScheduleID · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_RejectsScheduleIDTestMigrateSchedule_From…TestMigrateSchedule_FromVisibility_Workers · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_FromVisibility_WorkersTestMigrateSchedule_From…TestMigrateSchedule_RejectsQueryWithoutFromVisibility · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_RejectsQueryWithoutFromVisibilityTestMigrateSchedule_Reje…TestMigrateSchedule_RejectsWorkersWithScheduleID · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_RejectsWorkersWithScheduleIDTestMigrateSchedule_Reje…TestMigrateSchedule_Stdin_DryRun · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_Stdin_DryRunTestMigrateSchedule_Stdi…TestMigrateSchedule_Stdin_Execute · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_Stdin_ExecuteTestMigrateSchedule_Stdi…TestMigrateSchedule_Stdin_OutputLog · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_Stdin_OutputLogTestMigrateSchedule_Stdi…TestMigrateSchedule_Stdin_Workers · introduced test · go.temporal.io/server/tools/tdbg/TestMigrateSchedule_Stdin_WorkersTestMigrateSchedule_Stdi…TestScheduleStatus_Basic · introduced test · go.temporal.io/server/tools/tdbg/TestScheduleStatus_BasicTestScheduleStatus_BasicTestScheduleStatus_CountError · introduced test · go.temporal.io/server/tools/tdbg/TestScheduleStatus_CountErrorTestScheduleStatus_Count…TestScheduleStatus_DefaultsToDefaultNamespace · introduced test · go.temporal.io/server/tools/tdbg/TestScheduleStatus_DefaultsToDefaultNamespaceTestScheduleStatus_Defau…Focused file · go.temporal.io/server/tools/tdbg/commands.go · 1307 LOCtdbg/commands.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 tdbg
2
3 import (
4 "bufio"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "os"
10 "strings"
11 "sync"
12 "time"
13
14 "github.com/fatih/color"
15 "github.com/google/uuid"
16 "github.com/urfave/cli/v2"
17 commonpb "go.temporal.io/api/common/v1"
18 historypb "go.temporal.io/api/history/v1"
19 "go.temporal.io/api/workflowservice/v1"
20 "go.temporal.io/server/api/adminservice/v1"
21 enumsspb "go.temporal.io/server/api/enums/v1"
22 historyspb "go.temporal.io/server/api/history/v1"
23 persistencespb "go.temporal.io/server/api/persistence/v1"
24 "go.temporal.io/server/chasm"
25 "go.temporal.io/server/common"
26 "go.temporal.io/server/common/codec"
27 "go.temporal.io/server/common/log"
28 "go.temporal.io/server/common/namespace"
29 "go.temporal.io/server/common/persistence/serialization"
30 "go.temporal.io/server/common/persistence/versionhistory"
31 "go.temporal.io/server/common/primitives"
32 "go.temporal.io/server/common/primitives/timestamp"
33 "go.temporal.io/server/service/history/tasks"
34 "go.temporal.io/server/service/worker/scheduler"
35 "google.golang.org/protobuf/encoding/prototext"
36 "google.golang.org/protobuf/types/known/durationpb"
37 "google.golang.org/protobuf/types/known/timestamppb"
38 )
39
40 const (
41 historyImportBlobSize = 16
42 historyImportPageSize = 256 * 1024 // 256K
43 )
44
45 // AdminShowWorkflow shows history
46 func AdminShowWorkflow(c *cli.Context, clientFactory ClientFactory) error {
47 nsName, err := getRequiredOption(c, FlagNamespace)
48 if err != nil {
49 return err
50 }
51 wid, err := getRequiredOption(c, FlagWorkflowID)
52 if err != nil {
53 return err
54 }
55 rid := c.String(FlagRunID)
56 startEventId := c.Int64(FlagMinEventID)
57 endEventId := c.Int64(FlagMaxEventID)
58 startEventVerion := int64(c.Int(FlagMinEventVersion))
59 endEventVersion := int64(c.Int(FlagMaxEventVersion))
60 outputFileName := c.String(FlagOutputFilename)
61 decode := c.Bool(FlagDecode)
62
63 client := clientFactory.AdminClient(c)
64 serializer := serialization.NewSerializer()
65
66 ctx, cancel := newContext(c)
67 defer cancel()
68
69 nsID, err := getNamespaceID(c, clientFactory, namespace.Name(nsName))
70 if err != nil {
71 return err
72 }
73
74 var histories []*commonpb.DataBlob
75 var token []byte
76 for doContinue := true; doContinue; doContinue = len(token) != 0 {
77 resp, err := client.GetWorkflowExecutionRawHistoryV2(ctx, &adminservice.GetWorkflowExecutionRawHistoryV2Request{
78 NamespaceId: nsID.String(),
79 Execution: &commonpb.WorkflowExecution{
80 WorkflowId: wid,
81 RunId: rid,
82 },
83 StartEventId: startEventId,
84 EndEventId: endEventId,
85 StartEventVersion: startEventVerion,
86 EndEventVersion: endEventVersion,
87 MaximumPageSize: 100,
88 NextPageToken: token,
89 })
90 if err != nil {
91 return fmt.Errorf("unable to recv History Branch: %s", err)
92 }
93 histories = append(histories, resp.HistoryBatches...)
94 token = resp.NextPageToken
95 }
96
97 var historyBatches []*historypb.History
98 totalSize := 0
99 var errs []error
100 for idx, b := range histories {
101 totalSize += len(b.Data)
102 // nolint:errcheck // assuming that write will succeed.
103 fmt.Fprintf(c.App.Writer, "======== batch %v, blob len: %v ======\n", idx+1, len(b.Data))
104 historyBatch, err := serializer.DeserializeEvents(b)
105 if err != nil {
106 err := fmt.Errorf("unable to deserialize Events: %s", err)
107 fmt.Fprintln(c.App.Writer, err)
108 errs = append(errs, err)
109 continue
110 }
111 historyBatches = append(historyBatches, &historypb.History{Events: historyBatch})
112 encoder := codec.NewJSONPBEncoder()
113 data, err := encoder.EncodeHistoryEvents(historyBatch)
114 if err != nil {
115 err := fmt.Errorf("unable to encode History Events: %s", err)
116 // nolint:errcheck // assuming that write will succeed.
117 fmt.Fprintln(c.App.Writer, err)
118 text, terr := prototext.Marshal(&historypb.History{Events: historyBatch})
119 if terr == nil {
120 fmt.Fprintln(c.App.Writer, "marshal to text:")
121 fmt.Fprintln(c.App.Writer, string(text))
122 }
123 errs = append(errs, err)
124 continue
125 }
126 if decode {
127 data = decodePayloadsInJSON(data)
128 }
129 // nolint:errcheck // assuming that write will succeed.
130 fmt.Fprintln(c.App.Writer, string(data))
131 }
132 // nolint:errcheck // assuming that write will succeed.
133 fmt.Fprintf(c.App.Writer, "======== total batches %v, total blob len: %v ======\n", len(histories), totalSize)
134
135 // Show info to user about the option to decode payloads.
136 if !decode {
137 // nolint:errcheck // assuming that write will succeed.
138 fmt.Fprintf(c.App.ErrWriter, "(use --%s to decode payloads to JSON)\n", FlagDecode)
139 }
140
141 err = errors.Join(errs...)
142 if err != nil {
143 return err
144 }
145
146 if outputFileName != "" {
147 encoder := codec.NewJSONPBEncoder()
148 data, err := encoder.EncodeHistories(historyBatches)
149 if err != nil {
150 return fmt.Errorf("unable to serialize History data: %s", err)
151 }
152 if decode {
153 data = decodePayloadsInJSON(data)
154 }
155 if err := os.WriteFile(outputFileName, data, 0666); err != nil {
156 return fmt.Errorf("unable to write History data file: %s", err)
157 }
158 }
159 return nil
160 }
161
162 // AdminImportWorkflow imports history
163 func AdminImportWorkflow(c *cli.Context, clientFactory ClientFactory) error {
164 nsName, err := getRequiredOption(c, FlagNamespace)
165 if err != nil {
166 return err
167 }
168 wid, err := getRequiredOption(c, FlagWorkflowID)
169 if err != nil {
170 return err
171 }
172 rid, err := getRequiredOption(c, FlagRunID)
173 if err != nil {
174 return err
175 }
176 inputFileName := c.String(FlagInputFilename)
177
178 client := clientFactory.AdminClient(c)
179 serializer := serialization.NewSerializer()
180
181 ctx, cancel := newContext(c)
182 defer cancel()
183
184 data, err := os.ReadFile(inputFileName)
185 if err != nil {
186 return fmt.Errorf("unable to read History data file: %s", err)
187 }
188 encoder := codec.NewJSONPBEncoder()
189 historyBatches, err := encoder.DecodeHistories(data)
190 if err != nil {
191 return fmt.Errorf("unable to deserialize History data: %s", err)
192 }
193
194 versionHistory := &historyspb.VersionHistory{}
195 for _, historyBatch := range historyBatches {
196 for _, event := range historyBatch.Events {
197 item := versionhistory.NewVersionHistoryItem(event.EventId, event.Version)
198 if err := versionhistory.AddOrUpdateVersionHistoryItem(versionHistory, item); err != nil {
199 return fmt.Errorf("unable to generate version history: %s", err)
200 }
201 }
202 }
203
204 var token []byte
205 var blobs []*commonpb.DataBlob
206 blobSize := 0
207 for i := 0; i < len(historyBatches)+1; i++ {
208 if i < len(historyBatches) {
209 historyBatch := historyBatches[i]
210 blob, err := serializer.SerializeEvents(historyBatch.Events)
211 if err != nil {
212 return fmt.Errorf("unable to deserialize Events: %s", err)
213 }
214 blobSize += len(blob.Data)
215 blobs = append(blobs, blob)
216 }
217 if blobSize >= historyImportBlobSize ||
218 len(blobs) >= historyImportPageSize ||
219 (i == len(historyBatches) && len(blobs) > 0) {
220 resp, err := client.ImportWorkflowExecution(ctx, &adminservice.ImportWorkflowExecutionRequest{
221 Namespace: nsName,
222 Execution: &commonpb.WorkflowExecution{
223 WorkflowId: wid,
224 RunId: rid,
225 },
226 HistoryBatches: blobs,
227 VersionHistory: versionHistory,
228 Token: token,
229 })
230 if err != nil {
231 return fmt.Errorf("unable to send History Branch: %s", err)
232 }
233 token = resp.Token
234
235 blobs = []*commonpb.DataBlob{}
236 blobSize = 0
237 }
238 }
239 if len(blobs) != 0 {
240 return errors.New("unable to import workflow events, some events are not sent")
241 }
242 // call with empty history to commit
243 resp, err := client.ImportWorkflowExecution(ctx, &adminservice.ImportWorkflowExecutionRequest{
244 Namespace: nsName,
245 Execution: &commonpb.WorkflowExecution{
246 WorkflowId: wid,
247 RunId: rid,
248 },
249 HistoryBatches: []*commonpb.DataBlob{},
250 VersionHistory: versionHistory,
251 Token: token,
252 })
253 if err != nil {
254 return fmt.Errorf("unable to import workflow events: %s", err)
255 }
256 if len(resp.Token) != 0 {
257 return errors.New("unable to import workflow events, not committed")
258 }
259 return nil
260 }
261
262 // AdminDescribeExecution describes a Temporal execution (CHASM tree or workflow).
263 func AdminDescribeExecution(c *cli.Context, clientFactory ClientFactory) error {
264 resp, err := describeMutableState(c, clientFactory)
265 if err != nil {
266 return err
267 }
268 if resp == nil {
269 return errors.New("no mutable state returned")
270 }
271
272 // nolint:errcheck // assuming that write will succeed.
273 fmt.Fprintln(c.App.Writer, color.GreenString("Cache mutable state:"))
274 if resp.GetCacheMutableState() != nil {
275 prettyPrintJSONObject(c, resp.GetCacheMutableState())
276 }
277 if resp.GetDatabaseMutableState() != nil {
278 // nolint:errcheck // assuming that write will succeed.
279 fmt.Fprintln(c.App.Writer, color.GreenString("Database mutable state:"))
280 prettyPrintJSONObject(c, resp.GetDatabaseMutableState())
281 }
282
283 // CHASM executions also print their tree.
284 if len(resp.GetDatabaseMutableState().GetChasmNodes()) > 0 {
285 err := dumpChasmTree(resp, c)
286 if err != nil {
287 // nolint:errcheck // assuming that write will succeed.
288 fmt.Fprintln(c.App.Writer, color.RedString("Unable to dump CHASM tree:"), err)
289 }
290 }
291
292 if resp.GetDatabaseMutableState() != nil {
293 // nolint:errcheck // assuming that write will succeed.
294 fmt.Fprintln(c.App.Writer, color.GreenString("Current branch token:"))
295 versionHistories := resp.GetDatabaseMutableState().GetExecutionInfo().GetVersionHistories()
296 // if VersionHistories is set, then all branch infos are stored in VersionHistories
297 currentVersionHistory, err := versionhistory.GetCurrentVersionHistory(versionHistories)
298 if err != nil {
299 // nolint:errcheck // assuming that write will succeed.
300 fmt.Fprintln(c.App.Writer, color.RedString("Unable to get current version history:"), err)
301 } else {
302 currentBranchToken := persistencespb.HistoryBranch{}
303 err := currentBranchToken.Unmarshal(currentVersionHistory.BranchToken)
304 if err != nil {
305 // nolint:errcheck // assuming that write will succeed.
306 fmt.Fprintln(c.App.Writer, color.RedString("Unable to unmarshal current branch token:"), err)
307 } else {
308 prettyPrintJSONObject(c, &currentBranchToken)
309 }
310 }
311 }
312
313 // nolint:errcheck // assuming that write will succeed.
314 fmt.Fprintf(c.App.Writer, "History service address: %s\n", resp.GetHistoryAddr())
315 // nolint:errcheck // assuming that write will succeed.
316 fmt.Fprintf(c.App.Writer, "Shard Id: %s\n", resp.GetShardId())
317
318 return nil
319 }
320
321 func dumpChasmTree(resp *adminservice.DescribeMutableStateResponse, c *cli.Context) error {
322 chasmNodes := resp.GetDatabaseMutableState().GetChasmNodes()
323 if len(chasmNodes) == 0 {
324 return nil
325 }
326
327 logger := log.NewNoopLogger()
328 registry, err := newChasmRegistry(logger)
329 if err != nil {
330 return fmt.Errorf("failed to create CHASM registry: %w", err)
331 }
332
333 decodedNodes, err := decodeChasmNodes(chasmNodes, registry)
334 if err != nil {
335 return fmt.Errorf("failed to decode CHASM nodes: %w", err)
336 }
337
338 fmt.Fprintln(c.App.Writer, color.GreenString("CHASM Tree Nodes:")) // nolint:errcheck // assuming that write will succeed.
339 prettyPrintJSONObject(c, decodedNodes)
340
341 return nil
342 }
343
344 func describeMutableState(c *cli.Context, clientFactory ClientFactory) (*adminservice.DescribeMutableStateResponse, error) {
345 adminClient := clientFactory.AdminClient(c)
346
347 namespace, err := getRequiredOption(c, FlagNamespace)
348 if err != nil {
349 return nil, err
350 }
351 bid, err := getRequiredOption(c, FlagBusinessID)
352 if err != nil {
353 return nil, err
354 }
355 rid := c.String(FlagRunID)
356
357 ctx, cancel := newContext(c)
358 defer cancel()
359
360 resp, err := adminClient.DescribeMutableState(ctx, &adminservice.DescribeMutableStateRequest{
361 Namespace: namespace,
362 Execution: &commonpb.WorkflowExecution{
363 WorkflowId: bid,
364 RunId: rid,
365 },
366 Archetype: getArchetype(c),
367 ArchetypeId: chasm.ArchetypeID(c.Uint(FlagArchetypeID)),
368 })
369 if err != nil {
370 return nil, fmt.Errorf("unable to get Mutable State: %s", err)
371 }
372 return resp, nil
373 }
374
375 // AdminDeleteWorkflow force deletes a workflow's mutable state (both concrete and current), history, and visibility
376 // records as long as it's possible.
377 // It should only be used as a troubleshooting tool since no additional check will be done before the deletion.
378 // (e.g. if a child workflow has recorded its result in the parent workflow)
379 // Please use normal workflow delete command to gracefully delete a workflow execution.
380 func AdminDeleteWorkflow(c *cli.Context, clientFactory ClientFactory, prompter *Prompter) error {
381 adminClient := clientFactory.AdminClient(c)
382
383 namespace, err := getRequiredOption(c, FlagNamespace)
384 if err != nil {
385 return err
386 }
387 wid, err := getRequiredOption(c, FlagWorkflowID)
388 if err != nil {
389 return err
390 }
391 rid := c.String(FlagRunID)
392
393 msg := fmt.Sprintf("Namespace: %s WorkflowID: %s RunID: %s\nForce delete above workflow execution?", namespace, wid, rid)
394 prompter.Prompt(msg)
395
396 ctx, cancel := newContext(c)
397 defer cancel()
398
399 resp, err := adminClient.DeleteWorkflowExecution(ctx, &adminservice.DeleteWorkflowExecutionRequest{
400 Namespace: namespace,
401 Execution: &commonpb.WorkflowExecution{
402 WorkflowId: wid,
403 RunId: rid,
404 },
405 Archetype: getArchetype(c),
406 ArchetypeId: chasm.ArchetypeID(c.Uint(FlagArchetypeID)),
407 })
408 if err != nil {
409 return fmt.Errorf("unable to delete workflow execution: %s", err)
410 }
411
412 if len(resp.Warnings) != 0 {
413 // nolint:errcheck // assuming that write will succeed.
414 fmt.Fprintln(c.App.Writer, "Warnings:")
415 for _, warning := range resp.Warnings {
416 fmt.Fprintf(c.App.Writer, "- %s\n", warning)
417 }
418 // nolint:errcheck // assuming that write will succeed.
419 fmt.Fprintln(c.App.Writer, "")
420 }
421
422 // nolint:errcheck // assuming that write will succeed.
423 fmt.Fprintln(c.App.Writer, "Workflow execution deleted.")
424
425 return nil
426 }
427
428 // AdminGetShardID get shardID
429 func AdminGetShardID(c *cli.Context) error {
430 namespaceID := c.String(FlagNamespaceID)
431 wid, err := getRequiredOption(c, FlagWorkflowID)
432 if err != nil {
433 return err
434 }
435 numberOfShards := int32(c.Int(FlagNumberOfShards))
436
437 if numberOfShards <= 0 {
438 return fmt.Errorf("missing required parameter number of Shards")
439 }
440 shardID := common.WorkflowIDToHistoryShard(namespaceID, wid, numberOfShards)
441 // nolint:errcheck // assuming that write will succeed.
442 fmt.Fprintf(c.App.Writer, "ShardId for namespace, workflowId: %v, %v is %v \n", namespaceID, wid, shardID)
443 return nil
444 }
445
446 // getCategory first searches the registry for the category by the [tasks.Category.Name].
447 > func getCategory(registry tasks.TaskCategoryRegistry, key string) (tasks.Category, error) { commands.go ×1
448 > for _, category := range registry.GetCategories() {
449 > if strings.EqualFold(category.Name(), key) {
450 > return category, nil commands.go ×1
451 > }
452 }
453 > return tasks.Category{}, fmt.Errorf("unknown task category %q", key) commands.go ×1
454 }
455
456 // AdminListShardTasks outputs a list of a tasks for given Shard and Task Category
457 func AdminListShardTasks(c *cli.Context, clientFactory ClientFactory, registry tasks.TaskCategoryRegistry) error {
458 sid := int32(c.Int(FlagShardID))
459 categoryStr := c.String(FlagTaskCategory)
460 category, err := getCategory(registry, categoryStr)
461 if err != nil {
462 return err
463 }
464
465 client := clientFactory.AdminClient(c)
466 pageSize := defaultPageSize
467 if c.IsSet(FlagPageSize) {
468 pageSize = c.Int(FlagPageSize)
469 }
470
471 minFireTime, err := parseTime(c.String(FlagMinVisibilityTimestamp), time.Unix(0, 0), time.Now().UTC())
472 if err != nil {
473 return err
474 }
475 maxFireTime, err := parseTime(c.String(FlagMaxVisibilityTimestamp), time.Unix(0, 0), time.Now().UTC())
476 if err != nil {
477 return err
478 }
479 req := &adminservice.ListHistoryTasksRequest{
480 ShardId: sid,
481 Category: int32(category.ID()),
482 TaskRange: &historyspb.TaskRange{
483 InclusiveMinTaskKey: &historyspb.TaskKey{
484 FireTime: timestamppb.New(minFireTime),
485 TaskId: c.Int64(FlagMinTaskID),
486 },
487 ExclusiveMaxTaskKey: &historyspb.TaskKey{
488 FireTime: timestamppb.New(maxFireTime),
489 TaskId: c.Int64(FlagMaxTaskID),
490 },
491 },
492 BatchSize: int32(pageSize),
493 }
494
495 ctx, cancel := newContext(c)
496 defer cancel()
497 paginationFunc := func(paginationToken []byte) ([]any, []byte, error) {
498 req.NextPageToken = paginationToken
499 response, err := client.ListHistoryTasks(ctx, req)
500 if err != nil {
501 return nil, nil, err
502 }
503 token := response.NextPageToken
504
505 var items []any
506 for _, task := range response.Tasks {
507 items = append(items, task)
508 }
509 return items, token, nil
510 }
511 if err := paginate(c, paginationFunc, pageSize); err != nil {
512 return fmt.Errorf("unable to list History tasks: %s", err)
513 }
514 return nil
515 }
516
517 // AdminRemoveTask describes history host
518 func AdminRemoveTask(
519 c *cli.Context,
520 clientFactory ClientFactory,
521 taskCategoryRegistry tasks.TaskCategoryRegistry,
522 ) error {
523 adminClient := clientFactory.AdminClient(c)
524 shardID := c.Int(FlagShardID)
525 taskID := c.Int64(FlagTaskID)
526 category, err := getCategory(taskCategoryRegistry, c.String(FlagTaskCategory))
527 if err != nil {
528 return err
529 }
530 var visibilityTimestamp int64
531 if category.Type() == tasks.CategoryTypeScheduled {
532 if !c.IsSet(FlagTaskVisibilityTimestamp) {
533 //nolint:errorlint
534 return fmt.Errorf("%s is required to remove %s tasks", FlagTaskVisibilityTimestamp, category.Name())
535 }
536 visibilityTimestamp = c.Int64(FlagTaskVisibilityTimestamp)
537 }
538
539 ctx, cancel := newContext(c)
540 defer cancel()
541
542 req := &adminservice.RemoveTaskRequest{
543 ShardId: int32(shardID),
544 Category: int32(category.ID()),
545 TaskId: taskID,
546 VisibilityTime: timestamppb.New(timestamp.UnixOrZeroTime(visibilityTimestamp)),
547 }
548
549 _, err = adminClient.RemoveTask(ctx, req)
550 if err != nil {
551 return fmt.Errorf("unable to remove Task: %s", err)
552 }
553 return nil
554 }
555
556 // AdminDescribeShard describes shard by shard id
557 func AdminDescribeShard(c *cli.Context, clientFactory ClientFactory) error {
558 sid := c.Int(FlagShardID)
559 adminClient := clientFactory.AdminClient(c)
560 ctx, cancel := newContext(c)
561 defer cancel()
562 response, err := adminClient.GetShard(ctx, &adminservice.GetShardRequest{ShardId: int32(sid)})
563
564 if err != nil {
565 return fmt.Errorf("unable to initialize Shard Manager: %s", err)
566 }
567
568 prettyPrintJSONObject(c, response.ShardInfo)
569 return nil
570 }
571
572 // AdminShardManagement describes history host
573 func AdminShardManagement(c *cli.Context, clientFactory ClientFactory) error {
574 adminClient := clientFactory.AdminClient(c)
575 sid := c.Int(FlagShardID)
576
577 ctx, cancel := newContext(c)
578 defer cancel()
579
580 req := &adminservice.CloseShardRequest{}
581 req.ShardId = int32(sid)
582
583 _, err := adminClient.CloseShard(ctx, req)
584 if err != nil {
585 return fmt.Errorf("unable to close Shard Task: %s", err)
586 }
587 return nil
588 }
589
590 // AdminListGossipMembers outputs a list of gossip members
591 func AdminListGossipMembers(c *cli.Context, clientFactory ClientFactory) error {
592 roleFlag := c.String(FlagClusterMembershipRole)
593
594 adminClient := clientFactory.AdminClient(c)
595 ctx, cancel := newContext(c)
596 defer cancel()
597 response, err := adminClient.DescribeCluster(ctx, &adminservice.DescribeClusterRequest{})
598 if err != nil {
599 return fmt.Errorf("unable to describe Cluster: %s", err)
600 }
601
602 members := response.MembershipInfo.Rings
603 if roleFlag != string(primitives.AllServices) {
604 all := members
605
606 members = members[:0]
607 for _, v := range all {
608 if roleFlag == v.Role {
609 members = append(members, v)
610 }
611 }
612 }
613
614 prettyPrintJSONObject(c, members)
615 return nil
616 }
617
618 // AdminListClusterMembers outputs a list of cluster members
619 func AdminListClusterMembers(c *cli.Context, clientFactory ClientFactory) error {
620 role, _ := StringToEnum(c.String(FlagClusterMembershipRole), enumsspb.ClusterMemberRole_value)
621 // TODO: refactor this: parseTime shouldn't be used for duration.
622 heartbeatFlag, err := parseTime(c.String(FlagFrom), time.Time{}, time.Now().UTC())
623 if err != nil {
624 return fmt.Errorf("unable to parse Heartbeat time: %s", err)
625 }
626 heartbeat := time.Duration(heartbeatFlag.UnixNano())
627
628 adminClient := clientFactory.AdminClient(c)
629 ctx, cancel := newContext(c)
630 defer cancel()
631
632 req := &adminservice.ListClusterMembersRequest{
633 Role: enumsspb.ClusterMemberRole(role),
634 LastHeartbeatWithin: durationpb.New(heartbeat),
635 }
636
637 resp, err := adminClient.ListClusterMembers(ctx, req)
638 if err != nil {
639 return fmt.Errorf("unable to list Cluster Members: %s", err)
640 }
641
642 members := resp.ActiveMembers
643
644 prettyPrintJSONObject(c, members)
645 return nil
646 }
647
648 // AdminDescribeHistoryHost describes history host
649 func AdminDescribeHistoryHost(c *cli.Context, clientFactory ClientFactory) error {
650 adminClient := clientFactory.AdminClient(c)
651
652 namespace := c.String(FlagNamespace)
653 workflowID := c.String(FlagWorkflowID)
654 shardID := c.Int(FlagShardID)
655 historyAddr := c.String(FlagHistoryAddress)
656 printFully := c.Bool(FlagPrintFullyDetail)
657
658 flagsCount := 0
659 if c.IsSet(FlagShardID) {
660 flagsCount++
661 }
662 if c.IsSet(FlagNamespace) && c.IsSet(FlagWorkflowID) {
663 flagsCount++
664 }
665 if c.IsSet(FlagHistoryAddress) {
666 flagsCount++
667 }
668 if flagsCount != 1 {
669 return fmt.Errorf("missing required parameter either Shard Id, Namespace, Workflow Id or Host address")
670 }
671
672 ctx, cancel := newContext(c)
673 defer cancel()
674
675 req := &adminservice.DescribeHistoryHostRequest{}
676 if c.IsSet(FlagShardID) {
677 req.ShardId = int32(shardID)
678 } else if c.IsSet(FlagNamespace) && c.IsSet(FlagWorkflowID) {
679 req.Namespace = namespace
680 req.WorkflowExecution = &commonpb.WorkflowExecution{WorkflowId: workflowID}
681 } else if c.IsSet(FlagHistoryAddress) {
682 req.HostAddress = historyAddr
683 }
684
685 resp, err := adminClient.DescribeHistoryHost(ctx, req)
686 if err != nil {
687 return fmt.Errorf("unable to describe History host: %s", err)
688 }
689
690 if !printFully {
691 resp.ShardIds = nil
692 }
693 prettyPrintJSONObject(c, resp)
694 return nil
695 }
696
697 func adminRefreshWorkflowTasks(c *cli.Context, clientFactory ClientFactory, prompter *Prompter) error {
698 if c.IsSet(FlagVisibilityQuery) && c.IsSet(FlagWorkflowID) && c.IsSet(FlagRunID) {
699 return errors.New("setting parameter visibility query with workflow ID and run ID is not allowed")
700 }
701 if c.IsSet(FlagVisibilityQuery) && !c.IsSet(FlagWorkflowID) && !c.IsSet(FlagRunID) {
702 return AdminBatchRefreshWorkflowTasks(c, clientFactory, prompter)
703 }
704 return AdminRefreshWorkflowTasks(c, clientFactory)
705 }
706
707 // AdminRefreshWorkflowTasks refreshes all the tasks of a workflow
708 func AdminRefreshWorkflowTasks(c *cli.Context, clientFactory ClientFactory) error {
709 adminClient := clientFactory.AdminClient(c)
710
711 nsName, err := getRequiredOption(c, FlagNamespace)
712 if err != nil {
713 return err
714 }
715
716 wid, err := getRequiredOption(c, FlagWorkflowID)
717 if err != nil {
718 return err
719 }
720 rid := c.String(FlagRunID)
721
722 ctx, cancel := newContext(c)
723 defer cancel()
724
725 nsID, err := getNamespaceID(c, clientFactory, namespace.Name(nsName))
726 if err != nil {
727 return err
728 }
729
730 _, err = adminClient.RefreshWorkflowTasks(ctx, &adminservice.RefreshWorkflowTasksRequest{
731 NamespaceId: nsID.String(),
732 Execution: &commonpb.WorkflowExecution{
733 WorkflowId: wid,
734 RunId: rid,
735 },
736 Archetype: getArchetype(c),
737 ArchetypeId: chasm.ArchetypeID(c.Uint(FlagArchetypeID)),
738 })
739 if err != nil {
740 return fmt.Errorf("unable to refresh Workflow Task: %s", err)
741 } else {
742 // nolint:errcheck // assuming that write will succeed.
743 fmt.Fprintln(c.App.Writer, "Refresh workflow task succeeded.")
744 }
745 return nil
746 }
747
748 // AdminBatchRefreshWorkflowTasks starts a batch job to refresh workflow tasks for multiple workflows
749 func AdminBatchRefreshWorkflowTasks(c *cli.Context, clientFactory ClientFactory, prompter *Prompter) error {
750 adminClient := clientFactory.AdminClient(c)
751 workflowClient := clientFactory.WorkflowClient(c)
752
753 nsName, err := getRequiredOption(c, FlagNamespace)
754 if err != nil {
755 return err
756 }
757
758 query, err := getRequiredOption(c, FlagVisibilityQuery)
759 if err != nil {
760 return err
761 }
762
763 reason, err := getRequiredOption(c, FlagReason)
764 if err != nil {
765 return err
766 }
767
768 jobID := c.String(FlagJobID)
769 if jobID == "" {
770 jobID = fmt.Sprintf("batch-refresh-%d", time.Now().UnixNano())
771 }
772
773 ctx, cancel := newContext(c)
774 defer cancel()
775
776 // Count workflows matching the query to confirm with user
777 countResp, err := workflowClient.CountWorkflowExecutions(ctx, &workflowservice.CountWorkflowExecutionsRequest{
778 Namespace: nsName,
779 Query: query,
780 })
781 if err != nil {
782 return fmt.Errorf("unable to count workflow executions: %w", err)
783 }
784
785 msg := fmt.Sprintf("Will refresh tasks for %d execution(s) matching query %q in namespace %q. Continue Y/N?",
786 countResp.GetCount(), query, nsName)
787 prompter.Prompt(msg)
788
789 _, err = adminClient.StartAdminBatchOperation(ctx, &adminservice.StartAdminBatchOperationRequest{
790 Namespace: nsName,
791 VisibilityQuery: query,
792 JobId: jobID,
793 Reason: reason,
794 Identity: getCurrentUserFromEnv(),
795 Operation: &adminservice.StartAdminBatchOperationRequest_RefreshTasksOperation{
796 RefreshTasksOperation: &adminservice.BatchOperationRefreshTasks{},
797 },
798 })
799 if err != nil {
800 return fmt.Errorf("unable to start batch refresh workflow tasks: %w", err)
801 }
802
803 // nolint:errcheck // assuming that write will succeed.
804 fmt.Fprintf(c.App.Writer, "Batch Refresh Workflow Tasks started successfully for Job ID: %s\n", jobID)
805 return nil
806 }
807
808 // AdminRebuildMutableState rebuild a workflow mutable state using persisted history events
809 func AdminRebuildMutableState(c *cli.Context, clientFactory ClientFactory) error {
810 adminClient := clientFactory.AdminClient(c)
811
812 namespace, err := getRequiredOption(c, FlagNamespace)
813 if err != nil {
814 return err
815 }
816 wid, err := getRequiredOption(c, FlagWorkflowID)
817 if err != nil {
818 return err
819 }
820 rid := c.String(FlagRunID)
821
822 ctx, cancel := newContext(c)
823 defer cancel()
824
825 _, err = adminClient.RebuildMutableState(ctx, &adminservice.RebuildMutableStateRequest{
826 Namespace: namespace,
827 Execution: &commonpb.WorkflowExecution{
828 WorkflowId: wid,
829 RunId: rid,
830 },
831 })
832 if err != nil {
833 return fmt.Errorf("rebuild mutable state failed: %s", err)
834 } else {
835 // nolint:errcheck // assuming that write will succeed.
836 fmt.Fprintln(c.App.Writer, "rebuild mutable state succeeded.")
837 }
838 return nil
839 }
840
841 // AdminReplicateWorkflow force replicates a workflow by generating replication tasks
842 func AdminReplicateWorkflow(
843 c *cli.Context,
844 clientFactory ClientFactory,
845 ) error {
846 adminClient := clientFactory.AdminClient(c)
847
848 nsName, err := getRequiredOption(c, FlagNamespace)
849 if err != nil {
850 return err
851 }
852
853 wid, err := getRequiredOption(c, FlagWorkflowID)
854 if err != nil {
855 return err
856 }
857 rid := c.String(FlagRunID)
858
859 ctx, cancel := newContext(c)
860 defer cancel()
861
862 _, err = adminClient.GenerateLastHistoryReplicationTasks(ctx, &adminservice.GenerateLastHistoryReplicationTasksRequest{
863 Namespace: nsName,
864 Execution: &commonpb.WorkflowExecution{
865 WorkflowId: wid,
866 RunId: rid,
867 },
868 Archetype: getArchetype(c),
869 ArchetypeId: chasm.ArchetypeID(c.Uint(FlagArchetypeID)),
870 })
871 if err != nil {
872 return fmt.Errorf("unable to replicate workflow: %w", err)
873 }
874
875 // nolint:errcheck // assuming that write will succeed.
876 fmt.Fprintln(c.App.Writer, "Replication tasks generated successfully.")
877 return nil
878 }
879
880 // AdminMigrateSchedule migrates schedules between V1 (workflow-backed) and V2 (CHASM).
881 //
882 // It supports three mutually-exclusive selection modes, all sharing the required --target flag:
883 // - single: --schedule-id <id> (performs immediately, as before)
884 // - from visibility: --from-visibility [--query <q>] (default query is chosen from --target: the
885 // running V1 schedules when migrating to chasm, the running V2 schedules when migrating to workflow)
886 // - stdin: JSON lines piped on stdin, one {"namespace":..., "schedule_id":...} per line
887 //
888 // The from-visibility and stdin modes default to a dry-run; pass --execute to perform the migration.
889 > func AdminMigrateSchedule(c *cli.Context, clientFactory ClientFactory) error { commands.go ×4
890 > target, targetStr, err := parseMigrateTarget(c)
891 > if err != nil {
892 return err
893 }
894
895 > fromVisibility := c.Bool(FlagFromVisibility) commands.go ×4
896 > scheduleID := c.String(FlagScheduleID)
897 >
898 > // --query only takes effect in --from-visibility mode; reject it elsewhere rather than
899 > // silently ignoring it.
900 > if !fromVisibility && c.IsSet(FlagVisibilityQuery) {
901 > return fmt.Errorf("--%s is only valid with --%s", FlagVisibilityQuery, FlagFromVisibility) commands.go ×1
902 > }
903 // --workers applies to the bulk modes (--from-visibility and stdin); it has no effect when
904 // migrating a single --schedule-id, so reject it there rather than silently ignoring it.
905 > if scheduleID != "" && c.IsSet(FlagWorkers) { commands.go ×1
906 > return fmt.Errorf("--%s is only valid with --%s or when piping JSON lines on stdin", FlagWorkers, FlagFromVisibility) commands.go ×1
907 > }
908
909 > switch { commands.go ×1
910 > case fromVisibility: commands.go ×1
911 > if scheduleID != "" {
912 > return fmt.Errorf("--%s cannot be combined with --%s", FlagFromVisibility, FlagScheduleID) commands.go ×1
913 > }
914 > return migrateSchedulesFromVisibility(c, clientFactory, target, targetStr) commands.go ×9
915 case scheduleID != "":
916 return migrateSingleSchedule(c, clientFactory, target, targetStr, scheduleID)
917 > case isStdinPiped(): commands.go ×11
918 > return migrateSchedulesFromStdin(c, clientFactory, target, targetStr)
919 default:
920 return fmt.Errorf("specify one of: --%s, --%s, or pipe JSON lines on stdin", FlagScheduleID, FlagFromVisibility)
921 }
922 }
923
924 > func parseMigrateTarget(c *cli.Context) (adminservice.MigrateScheduleRequest_SchedulerTarget, string, error) { commands.go ×4
925 > targetStr, err := getRequiredOption(c, FlagTarget)
926 > if err != nil {
927 return 0, "", err
928 }
929 > switch strings.ToLower(targetStr) { commands.go ×4
930 > case "chasm": commands.go ×2
931 > return adminservice.MigrateScheduleRequest_SCHEDULER_TARGET_CHASM, targetStr, nil
932 > case "workflow": commands.go ×1
933 > return adminservice.MigrateScheduleRequest_SCHEDULER_TARGET_WORKFLOW, targetStr, nil
934 default:
935 return 0, "", fmt.Errorf("invalid target %q, valid values are: chasm, workflow", targetStr)
936 }
937 }
938
939 // v1ScheduleVisibilityQuery returns the visibility query selecting running V1 (workflow-backed)
940 // schedules.
941 > func v1ScheduleVisibilityQuery() string { commands.go ×1
942 > return fmt.Sprintf("TemporalNamespaceDivision = '%s' AND ExecutionStatus = 'Running'", scheduler.NamespaceDivision)
943 > }
944
945 // v2ScheduleVisibilityQuery returns the visibility query selecting running V2 (CHASM) schedules.
946 // The explicit TemporalNamespaceDivision filter is required, otherwise the visibility query
947 // converter appends "TemporalNamespaceDivision IS NULL" and excludes CHASM executions.
948 > func v2ScheduleVisibilityQuery() string { commands.go ×1
949 > return fmt.Sprintf("TemporalNamespaceDivision = '%d' AND ExecutionStatus = 'Running'", chasm.SchedulerArchetypeID)
950 > }
951
952 // AdminScheduleStatus reports how many schedules in --namespace are currently V1
953 // (workflow-backed) vs V2 (CHASM), using the same default visibility queries as
954 // `schedule migrate --from-visibility`.
955 > func AdminScheduleStatus(c *cli.Context, clientFactory ClientFactory) error { commands.go ×3
956 > ns, err := getRequiredOption(c, FlagNamespace)
957 > if err != nil {
958 return err
959 }
960
961 > wfClient := clientFactory.WorkflowClient(c) commands.go ×3
962 >
963 > v1Count, err := countScheduleVisibility(c, wfClient, ns, v1ScheduleVisibilityQuery())
964 > if err != nil {
965 > return fmt.Errorf("unable to count V1 schedules: %w", err) commands.go ×2
966 > }
967 > v2Count, err := countScheduleVisibility(c, wfClient, ns, v2ScheduleVisibilityQuery()) commands.go ×3
968 > if err != nil {
969 return fmt.Errorf("unable to count V2 schedules: %w", err)
970 }
971
972 > _, _ = fmt.Fprintf(c.App.Writer, "Namespace: %s\n", ns) commands.go ×3
973 > _, _ = fmt.Fprintf(c.App.Writer, "V1 (workflow-backed): %d\n", v1Count)
974 > _, _ = fmt.Fprintf(c.App.Writer, "V2 (CHASM): %d\n", v2Count)
975 > _, _ = fmt.Fprintf(c.App.Writer, "Total: %d\n", v1Count+v2Count)
976 > return nil
977 }
978
979 > func countScheduleVisibility(c *cli.Context, wfClient workflowservice.WorkflowServiceClient, ns, query string) (int64, error) { commands.go ×3
980 > ctx, cancel := newContext(c)
981 > defer cancel()
982 > resp, err := wfClient.CountWorkflowExecutions(ctx, &workflowservice.CountWorkflowExecutionsRequest{
983 > Namespace: ns,
984 > Query: query,
985 > })
986 > if err != nil {
987 > return 0, err commands.go ×2
988 > }
989 > return resp.GetCount(), nil commands.go ×3
990 }
991
992 // migrateSingleSchedule migrates one schedule and performs the migration immediately.
993 func migrateSingleSchedule(
994 c *cli.Context,
995 clientFactory ClientFactory,
996 target adminservice.MigrateScheduleRequest_SchedulerTarget,
997 targetStr string,
998 scheduleID string,
999 ) error {
1000 ns, err := getRequiredOption(c, FlagNamespace)
1001 if err != nil {
1002 return err
1003 }
1004
1005 adminClient := clientFactory.AdminClient(c)
1006 ctx, cancel := newContext(c)
1007 defer cancel()
1008
1009 if err := migrateScheduleRPC(ctx, adminClient, ns, scheduleID, target); err != nil {
1010 return fmt.Errorf("unable to migrate schedule: %w", err)
1011 }
1012
1013 _, _ = fmt.Fprintf(c.App.Writer, "Successfully initiated migration of schedule %q in namespace %q to %s.\n", scheduleID, ns, targetStr)
1014 return nil
1015 }
1016
1017 // migrateSchedulesFromVisibility selects schedules via a visibility query and migrates each.
1018 // When --query is not supplied the default query is chosen from the --target direction.
1019 func migrateSchedulesFromVisibility(
1020 c *cli.Context,
1021 clientFactory ClientFactory,
1022 target adminservice.MigrateScheduleRequest_SchedulerTarget,
1023 targetStr string,
1024 > ) error { commands.go ×9
1025 > ns, err := getRequiredOption(c, FlagNamespace)
1026 > if err != nil {
1027 return err
1028 }
1029
1030 // When --query is not supplied the default is chosen from the --target direction: migrating
1031 // to CHASM (V2) selects the running V1 (workflow-backed) schedules to move forward, while
1032 // migrating to workflow (V1) selects the running V2 (CHASM) schedules to roll back.
1033 > query := c.String(FlagVisibilityQuery) commands.go ×9
1034 > if query == "" {
1035 > if target == adminservice.MigrateScheduleRequest_SCHEDULER_TARGET_CHASM { commands.go ×2
1036 > // Forward migration V1 -> V2: all running V1 (workflow-backed) schedules. commands.go ×2
1037 > query = v1ScheduleVisibilityQuery()
1038 > } else { commands.go ×2
1039 > // Rollback V2 -> V1: all running V2 (CHASM) schedules. commands.go ×3
1040 > query = v2ScheduleVisibilityQuery()
1041 > }
1042 }
1043
1044 > execute := c.Bool(FlagExecute) commands.go ×9
1045 > workers := max(c.Int(FlagWorkers), 1)
1046 > wfClient := clientFactory.WorkflowClient(c)
1047 > adminClient := clientFactory.AdminClient(c)
1048 >
1049 > // Schedules are listed (paginated) on this goroutine and fed to a pool of workers
1050 > // that migrate them concurrently.
1051 > var summary migrateSummary
1052 > closeLog, err := openMigrateLog(c, &summary)
1053 > if err != nil {
1054 return err
1055 }
1056 > defer closeLog() commands.go ×9
1057 > jobs := make(chan migrateJob)
1058 > var wg sync.WaitGroup
1059 > for range workers {
1060 > wg.Go(func() {
1061 > for job := range jobs {
1062 > migrateOne(c, adminClient, job.namespace, job.scheduleID, target, targetStr, execute, &summary) commands.go ×3
1063 > }
1064 })
1065 }
1066
1067 > var listErr error commands.go ×9
1068 > var nextPageToken []byte
1069 > for {
1070 > ctx, cancel := newContext(c)
1071 > resp, err := wfClient.ListWorkflowExecutions(ctx, &workflowservice.ListWorkflowExecutionsRequest{
1072 > Namespace: ns,
1073 > Query: query,
1074 > NextPageToken: nextPageToken,
1075 > })
1076 > cancel()
1077 > if err != nil {
1078 listErr = fmt.Errorf("unable to list schedules from visibility: %w", err)
1079 break
1080 }
1081
1082 > for _, exec := range resp.GetExecutions() { commands.go ×9
1083 > workflowID := exec.GetExecution().GetWorkflowId() commands.go ×3
1084 > // CHASM scheduler executions store the schedule id directly as the workflow id;
1085 > // TrimPrefix is a no-op for them and handles any V1 records defensively.
1086 > scheduleID := strings.TrimPrefix(workflowID, primitives.ScheduleWorkflowIDPrefix)
1087 > jobs <- migrateJob{namespace: ns, scheduleID: scheduleID}
1088 > }
1089
1090 > nextPageToken = resp.GetNextPageToken() commands.go ×9
1091 > if len(nextPageToken) == 0 {
1092 > break
1093 }
1094 }
1095 > close(jobs) commands.go ×9
1096 > wg.Wait()
1097 >
1098 > // Always report what was migrated before surfacing a listing error: if pagination fails
1099 > // partway through, workers may have already migrated the schedules listed so far, and the
1100 > // user needs to see that partial progress.
1101 > summary.print(c, execute)
1102 > return listErr
1103 }
1104
1105 type migrateJob struct {
1106 namespace string
1107 scheduleID string
1108 }
1109
1110 // openMigrateLog wires summary.logEnc to the --output-log file when the flag is set, returning a
1111 // cleanup func that closes the file (a no-op when the flag is unset).
1112 > func openMigrateLog(c *cli.Context, summary *migrateSummary) (func(), error) { commands.go ×2
1113 > logPath := c.String(FlagOutputLog)
1114 > if logPath == "" {
1115 > return func() {}, nil commands.go ×1
1116 }
1117 > logFile, err := os.Create(logPath) commands.go ×4
1118 > if err != nil {
1119 return nil, fmt.Errorf("unable to open output log %q: %w", logPath, err)
1120 }
1121 > summary.logEnc = json.NewEncoder(logFile) commands.go ×4
1122 > return func() { _ = logFile.Close() }, nil
1123 }
1124
1125 // migrateSchedulesFromStdin reads JSON lines from stdin, one {"namespace","schedule_id"} per line,
1126 // feeding them to a pool of --workers goroutines that migrate them concurrently (mirroring
1127 // --from-visibility mode). With the default of one worker, lines are processed in order.
1128 func migrateSchedulesFromStdin(
1129 c *cli.Context,
1130 clientFactory ClientFactory,
1131 target adminservice.MigrateScheduleRequest_SchedulerTarget,
1132 targetStr string,
1133 > ) error { commands.go ×11
1134 > execute := c.Bool(FlagExecute)
1135 > workers := max(c.Int(FlagWorkers), 1)
1136 > adminClient := clientFactory.AdminClient(c)
1137 >
1138 > var summary migrateSummary
1139 > closeLog, err := openMigrateLog(c, &summary)
1140 > if err != nil {
1141 return err
1142 }
1143 > defer closeLog() commands.go ×11
1144 >
1145 > jobs := make(chan migrateJob)
1146 > var wg sync.WaitGroup
1147 > for range workers {
1148 > wg.Go(func() {
1149 > for job := range jobs {
1150 > migrateOne(c, adminClient, job.namespace, job.scheduleID, target, targetStr, execute, &summary)
1151 > }
1152 })
1153 }
1154
1155 > var readErr error commands.go ×11
1156 > scanner := bufio.NewScanner(os.Stdin)
1157 > for scanner.Scan() {
1158 > line := strings.TrimSpace(scanner.Text())
1159 > if line == "" {
1160 > continue commands.go ×1
1161 }
1162 > var record struct { commands.go ×11
1163 > Namespace string `json:"namespace"`
1164 > ScheduleID string `json:"schedule_id"`
1165 > }
1166 > if err := json.Unmarshal([]byte(line), &record); err != nil {
1167 readErr = fmt.Errorf("invalid JSON line %q: %w", line, err)
1168 break
1169 }
1170 > if record.Namespace == "" || record.ScheduleID == "" { commands.go ×11
1171 readErr = fmt.Errorf("each line must include non-empty \"namespace\" and \"schedule_id\": %q", line)
1172 break
1173 }
1174 > jobs <- migrateJob{namespace: record.Namespace, scheduleID: record.ScheduleID} commands.go ×11
1175 }
1176 > if readErr == nil { commands.go ×11
1177 > if err := scanner.Err(); err != nil {
1178 readErr = fmt.Errorf("error reading stdin: %w", err)
1179 }
1180 }
1181 > close(jobs) commands.go ×11
1182 > wg.Wait()
1183 >
1184 > // Always report what was migrated before surfacing a read error: workers may have already
1185 > // migrated the lines read so far, and the user needs to see that partial progress.
1186 > summary.print(c, execute)
1187 > return readErr
1188 }
1189
1190 // migrateOne migrates a single schedule (or prints the planned action in dry-run), updating
1191 // summary. It is safe to call concurrently from multiple workers: the migration RPC runs
1192 // outside the summary lock, while counter updates and output are serialized.
1193 func migrateOne(
1194 c *cli.Context,
1195 adminClient adminservice.AdminServiceClient,
1196 ns string,
1197 scheduleID string,
1198 target adminservice.MigrateScheduleRequest_SchedulerTarget,
1199 targetStr string,
1200 execute bool,
1201 summary *migrateSummary,
1202 > ) { commands.go ×2
1203 > if !execute {
1204 > summary.recordDryRun(c, ns, scheduleID, targetStr) commands.go ×2
1205 > return
1206 > }
1207
1208 > ctx, cancel := newContext(c) commands.go ×5
1209 > defer cancel()
1210 > err := migrateScheduleRPC(ctx, adminClient, ns, scheduleID, target)
1211 > summary.recordResult(c, ns, scheduleID, targetStr, err)
1212 }
1213
1214 func migrateScheduleRPC(
1215 ctx context.Context,
1216 adminClient adminservice.AdminServiceClient,
1217 ns string,
1218 scheduleID string,
1219 target adminservice.MigrateScheduleRequest_SchedulerTarget,
1220 > ) error { commands.go ×5
1221 > _, err := adminClient.MigrateSchedule(ctx, &adminservice.MigrateScheduleRequest{
1222 > Namespace: ns,
1223 > ScheduleId: scheduleID,
1224 > Target: target,
1225 > Identity: getCurrentUserFromEnv(),
1226 > RequestId: uuid.NewString(),
1227 > })
1228 > return err
1229 > }
1230
1231 // migrateLogRecord is one structured entry written to --output-log per schedule.
1232 type migrateLogRecord struct {
1233 Timestamp string `json:"timestamp"`
1234 Namespace string `json:"namespace"`
1235 ScheduleID string `json:"schedule_id"`
1236 Target string `json:"target"`
1237 Status string `json:"status"` // "migrated", "failed", or "dry-run"
1238 Error string `json:"error,omitempty"`
1239 }
1240
1241 type migrateSummary struct {
1242 mu sync.Mutex
1243 planned int
1244 migrated int
1245 failed int
1246 logEnc *json.Encoder // optional; writes one migrateLogRecord per result
1247 }
1248
1249 > func (s *migrateSummary) recordDryRun(c *cli.Context, ns, scheduleID, targetStr string) { commands.go ×2
1250 > s.mu.Lock()
1251 > defer s.mu.Unlock()
1252 > s.planned++
1253 > _, _ = fmt.Fprintf(c.App.Writer, "[dry-run] would migrate %s/%s -> %s\n", ns, scheduleID, targetStr)
1254 > s.writeLogLocked(ns, scheduleID, targetStr, "dry-run", nil)
1255 > }
1256
1257 > func (s *migrateSummary) recordResult(c *cli.Context, ns, scheduleID, targetStr string, err error) { commands.go ×5
1258 > s.mu.Lock()
1259 > defer s.mu.Unlock()
1260 > s.planned++
1261 > if err != nil {
1262 > s.failed++ commands.go ×2
1263 > _, _ = fmt.Fprintf(c.App.ErrWriter, "failed to migrate %s/%s: %v\n", ns, scheduleID, err)
1264 > s.writeLogLocked(ns, scheduleID, targetStr, "failed", err)
1265 > return
1266 > }
1267 > s.migrated++ commands.go ×5
1268 > _, _ = fmt.Fprintf(c.App.Writer, "migrated %s/%s -> %s\n", ns, scheduleID, targetStr)
1269 > s.writeLogLocked(ns, scheduleID, targetStr, "migrated", nil)
1270 }
1271
1272 // writeLogLocked appends a structured record to the output log. Callers must hold s.mu.
1273 > func (s *migrateSummary) writeLogLocked(ns, scheduleID, targetStr, status string, err error) { commands.go ×2
1274 > if s.logEnc == nil {
1275 > return commands.go ×1
1276 > }
1277 > rec := migrateLogRecord{ commands.go ×4
1278 > Timestamp: time.Now().UTC().Format(time.RFC3339),
1279 > Namespace: ns,
1280 > ScheduleID: scheduleID,
1281 > Target: targetStr,
1282 > Status: status,
1283 > }
1284 > if err != nil {
1285 > rec.Error = err.Error() commands.go ×2
1286 > }
1287 > _ = s.logEnc.Encode(&rec) commands.go ×4
1288 }
1289
1290 > func (s *migrateSummary) print(c *cli.Context, execute bool) { commands.go ×2
1291 > s.mu.Lock()
1292 > defer s.mu.Unlock()
1293 > if !execute {
1294 > _, _ = fmt.Fprintf(c.App.Writer, "Dry-run: %d schedule(s) would be migrated. Re-run with --%s to perform.\n", s.planned, FlagExecute) commands.go ×1
1295 > return
1296 > }
1297 > _, _ = fmt.Fprintf(c.App.Writer, "Done: %d migrated, %d failed (of %d).\n", s.migrated, s.failed, s.planned) commands.go ×5
1298 }
1299
1300 // isStdinPiped reports whether stdin is connected to a pipe or file rather than a terminal.
1301 > func isStdinPiped() bool { commands.go ×11
1302 > fi, err := os.Stdin.Stat()
1303 > if err != nil {
1304 return false
1305 }
1306 > return (fi.Mode() & os.ModeCharDevice) == 0 commands.go ×11
1307 }