1131
target adminservice.MigrateScheduleRequest_SchedulerTarget,
1132
targetStr string,
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
}
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
1156
>
scanner := bufio.NewScanner(os.Stdin)
1157
>
for scanner.Scan() {
1158
>
line := strings.TrimSpace(scanner.Text())
1159
>
if line == "" {
1160
continue
1161
}
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
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
1175
}
1177
>
if err := scanner.Err(); err != nil {
1178
readErr = fmt.Errorf("error reading stdin: %w", err)
1179
}
1180
}
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