Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 61 additions & 39 deletions cloud/collector.go
Original file line number Diff line number Diff line change
Expand Up @@ -639,18 +639,57 @@ func (c *Collector) captureAndSend(ctx context.Context, conn *wsConn, serviceKey

// Run is the reconnect loop. It blocks until ctx is cancelled. Each iteration
// dials, performs the hello handshake, and (on accept) serves until the
// connection drops, then backs off — unless the rejection is terminal.
// connection drops, then backs off. A rejection decides the pace (see
// helloRejection); a terminal one parks the collector in stopped, which keeps
// repeating the reason until the container restarts.
func (c *Collector) Run(ctx context.Context) {
bo := newBackoff()
var rareSince time.Time // start of the current run of consecutive rejectRetryRare rejections
var hinted string // reason whose full hint was last logged, at hintedAt
var hintedAt time.Time
for ctx.Err() == nil {
if c.session(ctx, bo) {
return // terminal rejection
}
accepted, rej := c.session(ctx, bo)
if ctx.Err() != nil {
return
}
d := bo.next()
c.log.Info().Dur("backoff", d).Msg("cloud: reconnecting after backoff")
if accepted || (rej != nil && rej.action != rejectRetryRare) {
rareSince = time.Time{}
}
if accepted {
hinted = ""
}
var d time.Duration
switch {
case rej == nil:
d = bo.next()
case rej.action == rejectStop:
c.stopped(ctx, *rej)
return
case rej.action == rejectRetryRare:
if rareSince.IsZero() {
rareSince = time.Now()
} else if time.Since(rareSince) >= rareRetryWindow {
c.stopped(ctx, rej.gaveUp())
return
}
d = bo.around(rareRetryInterval)
case rej.action == rejectRetrySlow:
bo.slow()
d = bo.next()
default:
d = bo.next()
}
switch {
case rej == nil:
c.log.Info().Dur("backoff", d).Msg("cloud: reconnecting after backoff")
case rej.reason != hinted || time.Since(hintedAt) >= reminderInterval:
// The full hint on the first rejection and every reminderInterval;
// the 30–60 s retries in between log one short line.
hinted, hintedAt = rej.reason, time.Now()
c.log.Warn().Str("reason", rej.reason).Dur("retry_in", d.Round(time.Second)).Msg("cloud: connection rejected. " + rej.hint)
default:
c.log.Warn().Str("reason", rej.reason).Dur("retry_in", d.Round(time.Second)).Msg("cloud: connection rejected again")
}
select {
case <-ctx.Done():
return
Expand All @@ -659,19 +698,22 @@ func (c *Collector) Run(ctx context.Context) {
}
}

// session runs one connection lifetime; stop=true means give up entirely.
func (c *Collector) session(ctx context.Context, bo *backoff) (stop bool) {
// session runs one connection lifetime. accepted reports that the cloud took the
// hello; rej is set when the cloud refused the connection (a 401/403 upgrade or a
// rejecting hello_ack). Both zero means an ordinary failure or disconnect.
func (c *Collector) session(ctx context.Context, bo *backoff) (accepted bool, rej *rejection) {
dialCtx, dialCancel := context.WithTimeout(ctx, 20*time.Second)
conn, err := dial(dialCtx, c.cfg.URL, c.cfg.Key, c.log)
dialCancel()
if err != nil {
var de *dialError
if asDialError(err, &de) && (de.statusCode == 401 || de.statusCode == 403) {
c.log.Error().Int("status", de.statusCode).Msg("cloud: connection rejected (auth) — stopping")
return true
if asDialError(err, &de) {
if r := httpRejection(de.statusCode); r != nil {
return false, r
}
}
c.log.Warn().Err(err).Msg("cloud: dial failed")
return false
return false, nil
}

connCtx, connCancel := context.WithCancel(ctx)
Expand Down Expand Up @@ -699,42 +741,31 @@ func (c *Collector) session(ctx context.Context, bo *backoff) (stop bool) {
if !c.sendHello(connCtx, conn) {
connCancel()
<-runDone
return false
return false, nil
}

select {
case <-connCtx.Done():
<-runDone
return false
return false, nil
case err := <-runDone:
if err != nil {
c.log.Warn().Err(err).Msg("cloud: connection closed before hello_ack")
}
return false
return false, nil
case ack := <-ackCh:
if !ack.Accepted {
if terminalReject(ack.Reason) {
c.log.Error().Str("reason", string(ack.Reason)).Msg("cloud: hello rejected (terminal) — stopping")
connCancel()
<-runDone
return true
}
if ack.Reason == proto.RejectEnrollmentClosed {
// An operator can reopen the key's enrollment window, so keep
// retrying automatically without flooding the logs while closed.
bo.slow()
}
c.log.Warn().Str("reason", string(ack.Reason)).Msg("cloud: hello rejected — will retry")
r := helloRejection(ack.Reason)
connCancel()
<-runDone
return false
return false, &r
}
c.log.Info().Str("host_id", ack.HostID).Int("config_version", ack.ConfigVersion).Msg("cloud: connected and accepted")
case <-time.After(15 * time.Second):
c.log.Warn().Msg("cloud: timed out waiting for hello_ack")
connCancel()
<-runDone
return false
return false, nil
}

// Accepted. Publish the connection, send an immediate snapshot from the last
Expand Down Expand Up @@ -767,7 +798,7 @@ func (c *Collector) session(ctx context.Context, bo *backoff) (stop bool) {
if err != nil {
c.log.Warn().Err(err).Msg("cloud: connection closed")
}
return false
return true, nil
}

func (c *Collector) heartbeatLoop(ctx context.Context, conn *wsConn) {
Expand Down Expand Up @@ -1218,15 +1249,6 @@ func (c *Collector) logModeFor(serviceKey string) string {
return c.logMode
}

func terminalReject(reason proto.RejectCode) bool {
switch reason {
case proto.RejectInvalidKey, proto.RejectBlocked, proto.RejectProtocolMismatch:
return true
default:
return false
}
}

func asDialError(err error, target **dialError) bool {
for err != nil {
if de, ok := err.(*dialError); ok {
Expand Down
153 changes: 153 additions & 0 deletions cloud/reject.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
package cloud

import (
"context"
"fmt"
"net/http"
"time"

"github.com/marvinvr/docktail/cloud/proto"
)

// dashboardURL is the DockTail Cloud dashboard that rejection hints point at.
// A DOCKTAIL_CLOUD_URL override (local development) does not change it.
const dashboardURL = "https://cloud.docktail.org"

const (
// rareRetryInterval is how often the collector asks again after a rejection
// only an operator (or a fixed cloud) can lift — a blocked host, a refused
// protocol version, a 403 from something in front of the cloud. None needs the
// agent restarted to clear, so it keeps knocking, but rarely.
rareRetryInterval = 15 * time.Minute
// rareRetryWindow bounds that knocking: after this long the collector stops
// like any terminal rejection, and a restart resumes it.
rareRetryWindow = 24 * time.Hour
// reminderInterval is how often a stopped collector repeats why, and how often
// a retrying one repeats the full hint (the retries in between log one short
// line), so the cause stays near the tail of the container log.
reminderInterval = 30 * time.Minute
)

// rejectAction is what the collector does after the cloud refuses a connection.
type rejectAction int

const (
rejectRetry rejectAction = iota // ordinary exponential backoff
rejectRetrySlow // operator-actionable: retry every 30–60 s
rejectRetryRare // retry about every rareRetryInterval, for at most rareRetryWindow
rejectStop // nothing changes until the container restarts
)

// rejection is a refused connection translated for the operator: the machine
// reason (a RejectCode, or "http_<status>" for a refused upgrade), what the
// collector does next, and one or two sentences saying what happened and how to
// fix it. A rejectRetryRare rejection also carries stopHint, the hint for when
// rareRetryWindow runs out.
type rejection struct {
reason string
action rejectAction
hint string
stopHint string
}

var (
agentKeysURL = dashboardURL + "/settings/agent-keys"
hostsURL = dashboardURL + "/hosts"
billingURL = dashboardURL + "/settings/billing"
)

// invalidKeyHint covers every "this key no longer authenticates" outcome. The
// key comes from the environment, so any fix needs a restart.
var invalidKeyHint = fmt.Sprintf("DockTail Cloud does not accept this workspace key: it was revoked, its workspace was deleted, or it is mistyped. "+
"Create a new agent key at %s, set it as %s, then recreate this container.", agentKeysURL, EnvKey)

// httpRejection classifies a refused WSS upgrade. Only 401/403 are rejections;
// any other status is an ordinary retryable dial failure (nil). DockTail Cloud
// answers an unknown or revoked key with 401; it has no 403 for an agent, so a
// 403 comes from something in between (a proxy, firewall or WAF) and is not the
// key's fault.
func httpRejection(status int) *rejection {
switch status {
case http.StatusUnauthorized:
return &rejection{reason: "http_401", action: rejectStop, hint: invalidKeyHint}
case http.StatusForbidden:
lead := "Something between this host and DockTail Cloud (a proxy, firewall or WAF) refused the connection with HTTP 403. " +
"Allow outbound WebSocket connections to the DockTail Cloud endpoint"
return &rejection{
reason: "http_403",
action: rejectRetryRare,
hint: lead + "; " + rareRetryNote + ".",
stopHint: lead + ", then restart this container.",
}
default:
return nil
}
}

// rareRetryNote tells the operator how a rejectRetryRare rejection proceeds.
var rareRetryNote = fmt.Sprintf("the agent checks again about every %d minutes for up to %s, and after that needs a restart",
int(rareRetryInterval/time.Minute), formatHours(rareRetryWindow))

// helloRejection classifies a non-accepted hello_ack.
func helloRejection(code proto.RejectCode) rejection {
r := rejection{reason: string(code)}
switch code {
case proto.RejectInvalidKey:
r.action, r.hint = rejectStop, invalidKeyHint
case proto.RejectBlocked:
lead := fmt.Sprintf("This host is blocked in DockTail Cloud. Unblock it at %s", hostsURL)
r.action = rejectRetryRare
r.hint = lead + "; " + rareRetryNote + "."
r.stopHint = lead + ", then restart this container."
case proto.RejectProtocolMismatch:
// Stopping would be right for an outdated image, but a cloud-side mistake
// would then silence every agent until each is restarted; a rare retry
// heals that on its own.
lead := fmt.Sprintf("DockTail Cloud does not accept this agent's wire protocol (DockTail %s, protocol v%d). "+
"Pull the latest DockTail image and recreate this container", agentVersion, proto.ProtocolVersion)
r.action = rejectRetryRare
r.hint = lead + "; until then " + rareRetryNote + "."
r.stopHint = lead + "."
case proto.RejectEnrollmentClosed:
r.action = rejectRetrySlow
r.hint = fmt.Sprintf("This workspace key's enrollment window has closed, so it cannot add a new host. "+
"Reopen enrollment for the key at %s (or create a new key and recreate this container with it); the agent retries automatically.", agentKeysURL)
case proto.RejectOverCap:
r.action = rejectRetrySlow
r.hint = fmt.Sprintf("This workspace has reached its host limit. Upgrade the plan at %s or remove an offline host at %s; the agent retries automatically.",
billingURL, hostsURL)
default:
// duplicate_identity (legacy), an empty reason, or a code newer than this
// agent: none is known to be permanent, so keep the ordinary backoff.
r.action = rejectRetry
r.hint = fmt.Sprintf("DockTail Cloud rejected the connection (reason %q); the agent retries automatically.", code)
}
return r
}

// gaveUp turns a rejectRetryRare rejection into the terminal one it ends in
// once rareRetryWindow has passed.
func (r rejection) gaveUp() rejection {
return rejection{reason: r.reason, action: rejectStop, hint: r.stopHint}
}

func formatHours(d time.Duration) string {
return fmt.Sprintf("%d hours", int(d/time.Hour))
}

// stopped logs a terminal rejection and then repeats it every
// reminderInterval until ctx ends, so an operator reading the latest
// container logs still finds the reason the host went quiet.
func (c *Collector) stopped(ctx context.Context, r rejection) {
c.log.Error().Str("reason", r.reason).Msg("cloud: stopped reporting until this container restarts. " + r.hint)
t := time.NewTicker(reminderInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
c.log.Error().Str("reason", r.reason).Msg("cloud: still not reporting. " + r.hint)
}
}
}
67 changes: 67 additions & 0 deletions cloud/reject_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
package cloud

import (
"strings"
"testing"
"time"

"github.com/marvinvr/docktail/cloud/proto"
)

func TestHelloRejectionActions(t *testing.T) {
cases := []struct {
code proto.RejectCode
action rejectAction
want string // a fragment the operator hint must carry
}{
{proto.RejectInvalidKey, rejectStop, "/settings/agent-keys"},
{proto.RejectBlocked, rejectRetryRare, "/hosts"},
{proto.RejectProtocolMismatch, rejectRetryRare, "protocol v"},
{proto.RejectEnrollmentClosed, rejectRetrySlow, "/settings/agent-keys"},
{proto.RejectOverCap, rejectRetrySlow, "/settings/billing"},
{proto.RejectDuplicate, rejectRetry, "retries automatically"},
{"", rejectRetry, "retries automatically"},
{"some_future_code", rejectRetry, "some_future_code"},
}
for _, tc := range cases {
r := helloRejection(tc.code)
if r.action != tc.action {
t.Errorf("%q: action = %d, want %d", tc.code, r.action, tc.action)
}
if r.reason != string(tc.code) {
t.Errorf("%q: reason = %q", tc.code, r.reason)
}
if !strings.Contains(r.hint, tc.want) {
t.Errorf("%q: hint %q does not mention %q", tc.code, r.hint, tc.want)
}
if r.action == rejectRetryRare {
if g := r.gaveUp(); g.action != rejectStop || g.reason != r.reason || !strings.Contains(g.hint, tc.want) {
t.Errorf("%q: gaveUp() = %+v", tc.code, g)
}
}
}
}

func TestHTTPRejection(t *testing.T) {
if r := httpRejection(401); r == nil || r.action != rejectStop || !strings.Contains(r.hint, EnvKey) {
t.Errorf("401: got %+v, want a terminal key rejection", r)
}
if r := httpRejection(403); r == nil || r.action != rejectRetryRare || strings.Contains(r.hint, EnvKey) || r.stopHint == "" {
t.Errorf("403: got %+v, want a rare retry that does not blame the key", r)
}
for _, status := range []int{400, 429, 500, 502, 503} {
if r := httpRejection(status); r != nil {
t.Errorf("status %d: got %+v, want a retryable dial failure", status, r)
}
}
}

func TestBackoffAround(t *testing.T) {
b := newBackoff()
for i := 0; i < 100; i++ {
d := b.around(rareRetryInterval)
if d < 12*time.Minute || d > 18*time.Minute {
t.Fatalf("around(%s) = %s, want within ±20%%", rareRetryInterval, d)
}
}
}
Loading
Loading