Skip to content
Draft
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
13 changes: 13 additions & 0 deletions packages/api/internal/orchestrator/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"go.opentelemetry.io/otel"
"go.uber.org/zap"

"github.com/e2b-dev/infra/packages/api/internal/orchestrator/deadnode"
"github.com/e2b-dev/infra/packages/api/internal/orchestrator/nodemanager"
"github.com/e2b-dev/infra/packages/api/internal/sandbox"
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
Expand Down Expand Up @@ -84,6 +85,9 @@ func (o *Orchestrator) syncNodes(ctx context.Context, store *sandbox.Store, skip
// Wait for nodes discovery to finish
wg.Wait()

// Discovery succeeded (a listNomadNodes failure returns early above)
o.lastDiscoverySync.Store(new(time.Now()))

// Sync state of all nodes currently in the pool
ctx, syncNodesSpan := tracer.Start(ctx, "keep-in-sync-existing")
defer syncNodesSpan.End()
Expand Down Expand Up @@ -114,6 +118,15 @@ func (o *Orchestrator) syncNodes(ctx context.Context, store *sandbox.Store, skip
}

o.deregisterNode(n)

return
}

if _, unreachable := n.UnreachableSince(); !unreachable {
ref := deadnode.NodeRef{ClusterID: n.ClusterID.String(), NodeID: n.ID}
if err := deadnode.RecordNodeSeen(ctx, o.redisClient, ref, time.Now()); err != nil {
logger.L().Warn(ctx, "Failed to record node liveness", zap.Error(err), logger.WithNodeID(n.ID))
}
}
})
}
Expand Down
150 changes: 150 additions & 0 deletions packages/api/internal/orchestrator/deadnode/liveness.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
package deadnode

import (
"context"
"errors"
"fmt"
"strconv"
"time"

"github.com/redis/go-redis/v9"
"go.uber.org/zap"

"github.com/e2b-dev/infra/packages/shared/pkg/logger"
)

// Node liveness is a cross-replica signal stored in Redis: every API replica
// refreshes a per-node "last seen" key after each successful sync cycle with
// that node. The dead-node sweep only considers a node dead when NO replica
// has seen it for the grace period — a single replica with a broken network
// view can never purge sandboxes of a node other replicas still reach.
//
// For nodes without a last-seen key (fresh feature rollout, or the node died
// before any replica ever synced it), a first-missing marker is written
// once (SET NX) so the grace period ages from the first observation and is
// shared across replicas and their restarts.
const (
nodeLastSeenKeyPrefix = "node:last-seen:"
nodeFirstMissingKeyPrefix = "node:first-missing:"

// nodeLivenessKeyTTL self-cleans keys of nodes that were removed from the
// fleet permanently.
nodeLivenessKeyTTL = 24 * time.Hour
)

// NodeRef identifies a node across clusters. Comparable — used as a map key.
type NodeRef struct {
ClusterID string
NodeID string
}

// nodeLiveness is the cross-replica view of one node.
type nodeLiveness struct {
// lastSeen is the last time any replica completed a successful sync with
// the node. Zero when no replica has ever recorded seeing it.
lastSeen time.Time
// firstMissing is when a replica first observed the node having no
// last-seen key. Only meaningful when lastSeen is zero.
firstMissing time.Time
}

func nodeLastSeenKey(ref NodeRef) string {
return nodeLastSeenKeyPrefix + ref.ClusterID + ":" + ref.NodeID
}

func nodeFirstMissingKey(ref NodeRef) string {
return nodeFirstMissingKeyPrefix + ref.ClusterID + ":" + ref.NodeID
}

// RecordNodeSeen refreshes the node's last-seen key and clears any
// first-missing marker. Called by every replica after a successful sync
// cycle unconditionally
func RecordNodeSeen(ctx context.Context, client redis.UniversalClient, ref NodeRef, now time.Time) error {
pipe := client.Pipeline()
pipe.Set(ctx, nodeLastSeenKey(ref), strconv.FormatInt(now.Unix(), 10), nodeLivenessKeyTTL)
pipe.Del(ctx, nodeFirstMissingKey(ref))
_, err := pipe.Exec(ctx)
if err != nil {
return fmt.Errorf("failed to record node seen: %w", err)
}

return nil
}

// fetchNodeLiveness returns the liveness for the given nodes. For nodes
// without a last-seen key it plants (SET NX) and reads back the shared
// first-missing marker so the missing-grace ages consistently across replicas
func fetchNodeLiveness(ctx context.Context, client redis.UniversalClient, refs []NodeRef, now time.Time) (map[NodeRef]nodeLiveness, error) {
if len(refs) == 0 {
return map[NodeRef]nodeLiveness{}, nil
}

pipe := client.Pipeline()
lastSeenCmds := make([]*redis.StringCmd, len(refs))
for i, ref := range refs {
lastSeenCmds[i] = pipe.Get(ctx, nodeLastSeenKey(ref))
}
if _, err := pipe.Exec(ctx); err != nil && !errors.Is(err, redis.Nil) {
return nil, fmt.Errorf("failed to fetch node last-seen keys: %w", err)
}

out := make(map[NodeRef]nodeLiveness, len(refs))
var missing []NodeRef
for i, ref := range refs {
ts, err := parseUnixSeconds(lastSeenCmds[i])
if err != nil {
missing = append(missing, ref)

continue
}

out[ref] = nodeLiveness{lastSeen: ts}
}

if len(missing) == 0 {
return out, nil
}

// Plant the shared first-missing marker (first writer wins) and read the
// authoritative value back in one pipeline.
markerPipe := client.Pipeline()
markerCmds := make([]*redis.StringCmd, len(missing))
for i, ref := range missing {
markerPipe.SetNX(ctx, nodeFirstMissingKey(ref), strconv.FormatInt(now.Unix(), 10), nodeLivenessKeyTTL)
markerCmds[i] = markerPipe.Get(ctx, nodeFirstMissingKey(ref))
}
if _, err := markerPipe.Exec(ctx); err != nil && !errors.Is(err, redis.Nil) {
return nil, fmt.Errorf("failed to plant node first-missing markers: %w", err)
}

for i, ref := range missing {
ts, err := parseUnixSeconds(markerCmds[i])
if err != nil {
// Leave the node without data; the sweep skips nodes it has no evidence about
logger.L().Warn(ctx, "Failed to read node first-missing marker",
zap.Error(err),
logger.WithNodeID(ref.NodeID),
)

continue
}

out[ref] = nodeLiveness{firstMissing: ts}
}

return out, nil
}

func parseUnixSeconds(cmd *redis.StringCmd) (time.Time, error) {
raw, err := cmd.Result()
if err != nil {
return time.Time{}, err
}

sec, err := strconv.ParseInt(raw, 10, 64)
if err != nil {
return time.Time{}, fmt.Errorf("invalid unix timestamp %q: %w", raw, err)
}

return time.Unix(sec, 0), nil
}
100 changes: 100 additions & 0 deletions packages/api/internal/orchestrator/deadnode/liveness_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
package deadnode

import (
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

redis_utils "github.com/e2b-dev/infra/packages/shared/pkg/redis"
)

func TestNodeLiveness_RecordAndFetch(t *testing.T) {
t.Parallel()

client := redis_utils.SetupInstance(t)
ref := NodeRef{ClusterID: "cluster-a", NodeID: "node-1"}
now := time.Now()

require.NoError(t, RecordNodeSeen(t.Context(), client, ref, now))

liveness, err := fetchNodeLiveness(t.Context(), client, []NodeRef{ref}, now)
require.NoError(t, err)

nl, ok := liveness[ref]
require.True(t, ok)
assert.WithinDuration(t, now, nl.lastSeen, time.Second)
assert.True(t, nl.firstMissing.IsZero())
}

func TestNodeLiveness_MissingMarkerIsFirstWriteWins(t *testing.T) {
t.Parallel()

client := redis_utils.SetupInstance(t)
ref := NodeRef{ClusterID: "cluster-a", NodeID: "node-1"}
first := time.Now().Add(-time.Minute)

// First observer (e.g. replica A) plants the marker.
livenessA, err := fetchNodeLiveness(t.Context(), client, []NodeRef{ref}, first)
require.NoError(t, err)
require.WithinDuration(t, first, livenessA[ref].firstMissing, time.Second)

// A later observer (replica B) must read A's marker, not plant its own —
// the grace period ages from the FIRST observation across all replicas.
livenessB, err := fetchNodeLiveness(t.Context(), client, []NodeRef{ref}, time.Now())
require.NoError(t, err)
assert.Equal(t, livenessA[ref].firstMissing.Unix(), livenessB[ref].firstMissing.Unix())
}

func TestNodeLiveness_RecordSeenClearsMissingMarker(t *testing.T) {
t.Parallel()

client := redis_utils.SetupInstance(t)
ref := NodeRef{ClusterID: "cluster-a", NodeID: "node-1"}
past := time.Now().Add(-time.Hour)

// Marker planted long ago.
_, err := fetchNodeLiveness(t.Context(), client, []NodeRef{ref}, past)
require.NoError(t, err)

// Node seen: last-seen written, marker cleared.
require.NoError(t, RecordNodeSeen(t.Context(), client, ref, time.Now()))

// Simulate the last-seen key later expiring (node gone again long after):
// the next missing observation must start a FRESH grace period, not
// resurrect the old marker.
require.NoError(t, client.Del(t.Context(), nodeLastSeenKey(ref)).Err())

now := time.Now()
liveness, err := fetchNodeLiveness(t.Context(), client, []NodeRef{ref}, now)
require.NoError(t, err)
assert.WithinDuration(t, now, liveness[ref].firstMissing, time.Second, "old marker must not survive a successful sync")
}

func TestNodeLiveness_MixedBatch(t *testing.T) {
t.Parallel()

client := redis_utils.SetupInstance(t)
seen := NodeRef{ClusterID: "cluster-a", NodeID: "node-seen"}
missing := NodeRef{ClusterID: "cluster-a", NodeID: "node-missing"}
now := time.Now()

require.NoError(t, RecordNodeSeen(t.Context(), client, seen, now))

liveness, err := fetchNodeLiveness(t.Context(), client, []NodeRef{seen, missing}, now)
require.NoError(t, err)
require.Len(t, liveness, 2)
assert.False(t, liveness[seen].lastSeen.IsZero())
assert.False(t, liveness[missing].firstMissing.IsZero())
}

func TestNodeLiveness_EmptyRefs(t *testing.T) {
t.Parallel()

client := redis_utils.SetupInstance(t)

liveness, err := fetchNodeLiveness(t.Context(), client, nil, time.Now())
require.NoError(t, err)
assert.Empty(t, liveness)
}
Loading
Loading