diff --git a/internal/grpc/pb/vantage.pb.go b/internal/grpc/pb/vantage.pb.go index d96cd29..94cec48 100644 --- a/internal/grpc/pb/vantage.pb.go +++ b/internal/grpc/pb/vantage.pb.go @@ -166,8 +166,14 @@ type ServerCommand struct { RunStep *RunStepCmd `json:"run_step,omitempty"` CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"` OpenProxy *OpenProxyCmd `json:"open_proxy,omitempty"` + Ping *PingCmd `json:"ping,omitempty"` } +// PingCmd is a server-originated liveness beat. It carries nothing and expects +// no reply: its arrival is the entire message. See the .proto for why gRPC +// keepalive is not sufficient on its own. +type PingCmd struct{} + type CleanupWorkspaceCmd struct { diff --git a/internal/sync/sync.go b/internal/sync/sync.go index de06199..9ba6001 100644 --- a/internal/sync/sync.go +++ b/internal/sync/sync.go @@ -123,6 +123,15 @@ func poll(client *grpcclient.Client, cfg *config.Config, version string) error { // continuation of a run of failures. const streamHealthyAfter = time.Minute +// Stream staleness. The server beats every 20s, so 70s tolerates three missed +// beats before the stream is written off — high enough that a slow network or a +// briefly busy server does not cost a reconnect, low enough that an agent is +// not uncommandable for minutes after a control-plane restart. +const ( + streamStaleAfter = 70 * time.Second + streamStaleCheck = 10 * time.Second +) + func runCommandStream(ctx context.Context, cfg *config.Config) { backoff := time.Second @@ -185,7 +194,13 @@ func connectAndHandleStream(ctx context.Context, cfg *config.Config) error { } defer client.Close() - stream, err := client.CommandStream(ctx) + // Cancelling this context is what unblocks Recv when the stream has gone + // quiet. Without it the watchdog below would have no way to interrupt a + // read that is never going to return. + streamCtx, abandon := context.WithCancel(ctx) + defer abandon() + + stream, err := client.CommandStream(streamCtx) if err != nil { return fmt.Errorf("open stream: %w", err) } @@ -207,11 +222,64 @@ func connectAndHandleStream(ctx context.Context, cfg *config.Config) error { return stream.Send(msg) } + // Stream liveness, tracked here rather than left to gRPC keepalive. + // + // Keepalive operates on the transport, and behind an L7 proxy the transport + // ends at the proxy: it answers pings whether or not the server behind it + // is still running. A control-plane pod that dies therefore leaves this + // agent blocked in Recv on a stream that will never deliver another message + // and never error, while the control plane dispatches commands into it and + // the operator watches nothing happen. + // + // The watchdog only arms once a ping has actually been seen. A server too + // old to send them must not be treated as dead — that would put the agent + // in a reconnect loop against a control plane that is working perfectly. + var ( + lastMu sync.Mutex + lastRecv = time.Now() + pinged bool + ) + markRecv := func(isPing bool) { + lastMu.Lock() + lastRecv = time.Now() + if isPing { + pinged = true + } + lastMu.Unlock() + } + + go func() { + t := time.NewTicker(streamStaleCheck) + defer t.Stop() + for { + select { + case <-streamCtx.Done(): + return + case <-t.C: + lastMu.Lock() + idle, armed := time.Since(lastRecv), pinged + lastMu.Unlock() + if armed && idle > streamStaleAfter { + log.Printf("command stream silent for %s, assuming it is dead", idle.Truncate(time.Second)) + abandon() + return + } + } + } + }() + for { cmd, err := stream.Recv() if err != nil { return fmt.Errorf("recv: %w", err) } + markRecv(cmd.Ping != nil) + + // Pings carry nothing and are not acknowledged; being received is their + // whole purpose. + if cmd.Ping != nil { + continue + } if cmd.GenerateKey != nil { go handleGenerateKey(cfg, cmd)