Skip to content

Commit e1878fa

Browse files
authored
Make batch run --follow actually stream, and terminate (#16)
--follow previously waited for the execution to finish and then printed its output in one block, which for a long job means staring at nothing for its whole duration. The first attempt at fixing that held a server-side follow open instead, and was worse: the server never closes a followed batch log, so the command hung until killed, and because it attached only after the execution left the queue it still showed everything at once. Tailing from a byte offset instead gives both properties. Output appears as it is produced -- verified against a job that prints every two seconds, and the lines arrive two seconds apart -- and the loop ends when the execution reaches a terminal state, which is where the answer to "is there more?" actually lives. The status is read before the log so a final read cannot miss anything written between the two. batch logs --follow uses the same loop, so it no longer hangs on an execution that has already finished.
1 parent f5cdc28 commit e1878fa

2 files changed

Lines changed: 85 additions & 39 deletions

File tree

‎api/client.go‎

Lines changed: 16 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -585,35 +585,27 @@ func (client *Client) CancelBatchExecution(project, service, execution string) e
585585
return client.request(http.MethodPost, path, &resp, nil)
586586
}
587587

588-
// StreamBatchExecutionLogs writes one execution's container output to w.
588+
// ReadBatchExecutionLogs writes an execution's output from byteOffset onward and
589+
// returns how many bytes it wrote.
589590
//
590591
// Batch logs are addressed per execution rather than per service: a batch
591592
// service has no deployment for the service-wide endpoint to derive a label
592593
// from.
593594
//
594-
// A followed stream reconnects and resumes at an exact byte offset, the same
595-
// way service logs do. It always starts from the beginning of the log, since an
596-
// execution's output is finite and usually short -- there is no reason to
597-
// anchor to the end and risk missing the first lines.
598-
func (client *Client) StreamBatchExecutionLogs(project, service, execution string, w io.Writer, follow bool) error {
599-
base := batchPath(project, service) + "/" + execution + "/logs"
600-
logPath := func(params url.Values) string {
601-
if follow {
602-
params.Set("follow", "true")
603-
params.Set("keepalive", keepaliveInterval.String())
604-
}
605-
return base + "?" + params.Encode()
606-
}
607-
start := url.Values{"seek": {"start"}}
608-
if !follow {
609-
return client.streamOnce(logPath(start), w)
610-
}
611-
return client.follow(followOptions{path: func(resumeAfter int) string {
612-
if resumeAfter == 0 {
613-
return logPath(start)
614-
}
615-
return logPath(url.Values{"seek": {"start"}, "offset": {strconv.Itoa(resumeAfter)}})
616-
}}, w)
595+
// This is a plain read rather than a server-side follow. A followed stream is
596+
// never closed by the server when an execution's log completes, so a follower
597+
// has no way to learn it is done and hangs. Reading from an offset lets the
598+
// caller decide when to stop -- which for a batch execution is when the
599+
// execution itself reaches a terminal state.
600+
func (client *Client) ReadBatchExecutionLogs(project, service, execution string, byteOffset int, w io.Writer) (int, error) {
601+
params := url.Values{"seek": {"start"}}
602+
if byteOffset > 0 {
603+
params.Set("offset", strconv.Itoa(byteOffset))
604+
}
605+
path := batchPath(project, service) + "/" + execution + "/logs?" + params.Encode()
606+
counter := &countingWriter{w: w}
607+
err := client.streamOnce(path, counter)
608+
return counter.bytes, err
617609
}
618610

619611
// GetBatchConfig retrieves a batch service's configuration.

‎cmd/batch.go‎

Lines changed: 69 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -57,20 +57,15 @@ latest successful build is used.`,
5757
if !batchRunFollow && !batchRunWait {
5858
return nil
5959
}
60-
final, err := awaitBatchCompletion(client, cfg, exec.ID)
61-
if err != nil {
62-
return err
63-
}
6460
if batchRunFollow {
65-
// Deliberately fetched after completion rather than streamed
66-
// live. Batch containers routinely finish in under a second, so
67-
// a tail started on "running" usually attaches to a log that is
68-
// already closed and the stream ends before anything arrives.
69-
// Waiting first makes the output deterministic.
70-
if err := client.StreamBatchExecutionLogs(cfg.Project(), cfg.Service(), exec.ID, os.Stdout, false); err != nil {
61+
if err := followBatchOutput(client, cfg, exec.ID); err != nil {
7162
fmt.Fprintf(os.Stderr, "could not read execution output: %v\n", err)
7263
}
7364
}
65+
final, err := awaitBatchCompletion(client, cfg, exec.ID)
66+
if err != nil {
67+
return err
68+
}
7469
fmt.Printf("Execution %s %s (exit %d)\n", final.ID, colorize(final.Status), final.ExitCode)
7570
if final.Status != "succeeded" {
7671
if final.Error != "" {
@@ -160,7 +155,11 @@ latest successful build is used.`,
160155
if err != nil {
161156
return err
162157
}
163-
return client.StreamBatchExecutionLogs(cfg.Project(), cfg.Service(), args[0], os.Stdout, batchLogsFollow)
158+
if batchLogsFollow {
159+
return followBatchOutput(client, cfg, args[0])
160+
}
161+
_, err = client.ReadBatchExecutionLogs(cfg.Project(), cfg.Service(), args[0], 0, os.Stdout)
162+
return err
164163
}),
165164
}
166165

@@ -299,28 +298,83 @@ func shortID(id string) string {
299298
return id
300299
}
301300

302-
const batchPollInterval = 2 * time.Second
301+
const (
302+
batchPollInterval = 2 * time.Second
303+
// Log tailing polls faster than status, so output feels live.
304+
batchLogPollInterval = 750 * time.Millisecond
305+
)
303306

304307
func awaitBatchCompletion(client *api.Client, cfg config.ServiceConfig, id string) (*api.BatchExecution, error) {
305308
for {
306309
exec, err := client.InspectBatchExecution(cfg.Project(), cfg.Service(), id)
307310
if err != nil {
308311
return nil, err
309312
}
310-
switch exec.Status {
311-
case "succeeded", "failed", "cancelled":
313+
if isTerminalBatchStatus(exec.Status) {
312314
return exec, nil
313315
}
314316
time.Sleep(batchPollInterval)
315317
}
316318
}
317319

320+
// followBatchOutput shows an execution's output as it is produced, returning
321+
// once the execution has finished and its output has been fully drained.
322+
//
323+
// It polls from a byte offset rather than holding a server-side follow open. The
324+
// server never closes a followed batch log, so a follower cannot tell when the
325+
// execution is done and simply hangs -- which it did. Polling puts the
326+
// termination condition where the answer actually lives: the execution's status.
327+
func followBatchOutput(client *api.Client, cfg config.ServiceConfig, id string) error {
328+
offset := 0
329+
for {
330+
exec, err := client.InspectBatchExecution(cfg.Project(), cfg.Service(), id)
331+
if err != nil {
332+
return err
333+
}
334+
// A queued execution has no log yet, and asking for one is an error
335+
// rather than an empty read.
336+
if exec.Status != "pending" {
337+
n, err := client.ReadBatchExecutionLogs(cfg.Project(), cfg.Service(), id, offset, os.Stdout)
338+
offset += n
339+
if err != nil && !isTerminalBatchStatus(exec.Status) {
340+
// Transient while the container is still coming up.
341+
time.Sleep(batchPollInterval)
342+
continue
343+
} else if err != nil {
344+
return err
345+
}
346+
}
347+
if isTerminalBatchStatus(exec.Status) {
348+
// The status was read before the log, so anything written between
349+
// the two reads is still outstanding.
350+
n, err := client.ReadBatchExecutionLogs(cfg.Project(), cfg.Service(), id, offset, os.Stdout)
351+
offset += n
352+
if err != nil {
353+
return err
354+
}
355+
if n == 0 {
356+
return nil
357+
}
358+
continue
359+
}
360+
time.Sleep(batchLogPollInterval)
361+
}
362+
}
363+
364+
func isTerminalBatchStatus(status string) bool {
365+
switch status {
366+
case "succeeded", "failed", "cancelled":
367+
return true
368+
}
369+
return false
370+
}
371+
318372
func initBatchCmd() {
319373
batchCmd.PersistentFlags().StringVarP(&batchService, "service", "s", "", "The batch service to operate on")
320374

321375
batchRunCmd.Flags().StringVarP(&batchRunBuild, "build", "b", "", "Build to run (default: latest successful)")
322376
batchRunCmd.Flags().StringArrayVarP(&batchRunEnv, "env", "e", nil, "Environment variable as KEY=VALUE, repeatable")
323-
batchRunCmd.Flags().BoolVarP(&batchRunFollow, "follow", "f", false, "Wait for the execution to finish, then print its output")
377+
batchRunCmd.Flags().BoolVarP(&batchRunFollow, "follow", "f", false, "Show the output as the execution runs, and wait for it to finish")
324378
batchRunCmd.Flags().BoolVarP(&batchRunWait, "wait", "w", false, "Wait for the execution to finish and exit non-zero if it failed")
325379
batchLogsCmd.Flags().BoolVarP(&batchLogsFollow, "follow", "f", false, "Follow the output")
326380

0 commit comments

Comments
 (0)