apache/beam · error

failed to send staging token

Error message

failed to send staging token

What it means

stageFiles wraps the gRPC Send error when it fails to send the initial ArtifactResponseWrapper carrying the staging token on the ReverseArtifactRetrievalService stream. It indicates the artifact staging stream could not even be established with the token.

Solutions

  1. Check the wrapped gRPC error for connection status (Unavailable, Unimplemented)
  2. If Unimplemented, upgrade the job server to a version supporting reverse artifact retrieval
  3. Verify transport security settings match the server (TLS vs plaintext)
  4. Confirm network connectivity to the artifact endpoint

Example fix

// before
if err := stream.Send(&jobpb.ArtifactResponseWrapper{StagingToken: st}); err != nil {
  return errors.Wrapf(err, "failed to send staging token")
}
// after: diagnose via the wrapped status, e.g. status.Code(err) == codes.Unimplemented -> upgrade job server
Defensive patterns

Strategy: retry

Validate before calling

status := statusFromErr(err) // check codes.Unavailable/Unimplemented after failure
// pre-check: ensure job server supports reverse artifact retrieval

Try / catch

if err := stageFiles(ctx, cc, binary, st); err != nil {
  if status.Code(errors.Unwrap(err)) == codes.Unimplemented {
    // job server too old: upgrade job server
  }
}

Prevention

When it happens

Trigger: stream.Send(&jobpb.ArtifactResponseWrapper{StagingToken: st}) returns a gRPC error — connection dropped, server rejected the stream, or the artifact service does not implement reverse artifact retrieval.

Common situations: Artifact service version doesn't implement ReverseArtifactRetrievalService, TLS/plaintext mismatch on the connection, connection reset because the job server is overloaded or shut down.

Related errors


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/5baff14dae36dade. Report an issue: GitHub.

Appendix: source

Thrown at sdks/go/pkg/beam/runners/universal/runnerlib/stage.go:86

			return errors.Errorf("failed to stage artifacts for token %v in %v attempts: %v", st, attempts, strings.Join(failures, ";\n"))
		}
	}
}

func stageFiles(ctx context.Context, cc *grpc.ClientConn, binary, st string) error {
	client := jobpb.NewArtifactStagingServiceClient(cc)
	stream, err := client.ReverseArtifactRetrievalService(ctx)
	if err != nil {
		return err
	}
	defer func() {
		if err := stream.CloseSend(); err != nil {
			log.Error(ctx, "StageViaPortableApi CloseSend error: ", err)
		}
	}()

	if err := stream.Send(&jobpb.ArtifactResponseWrapper{StagingToken: st}); err != nil {
		return errors.Wrapf(err, "failed to send staging token")
	}

	for {
		in, err := stream.Recv()
		if err == io.EOF {
			return nil
		}
		if err != nil {
			return err
		}

		switch request := in.Request.(type) {
		case *jobpb.ArtifactRequestWrapper_ResolveArtifact:
			err = stream.Send(&jobpb.ArtifactResponseWrapper{
				Response: &jobpb.ArtifactResponseWrapper_ResolveArtifactResponse{
					ResolveArtifactResponse: &jobpb.ResolveArtifactsResponse{
						Replacements: request.ResolveArtifact.Artifacts,
					},

View on GitHub (pinned to 12126d8942)