feat: Added ping command
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user