From 14b1312cfa6f0474a44d82893f194e16572ecb29 Mon Sep 17 00:00:00 2001 From: CMGS Date: Sun, 27 Sep 2026 17:03:51 +0800 Subject: [PATCH 1/4] feat: drop a NodeInventory older than --inventory-stale-after, and let the applier delete one A dead node's last NodeInventory stayed a claim candidate, kept its sandboxes listed as running, and drew warm-pool PUTs and envd-proxy probes until its Node object went away. NodeInventory gains publishedAt; kubeinventory.Source omits an inventory whose publishedAt is older than StaleAfter (90 s, three publish intervals) from ListNodes and answers NotFound for it from NodeInventory and NodeCapacity, so every consumer of the source drops the node at once and the watch emits Deleted for its entries, as mesh mode's MaxStale does. An inventory without publishedAt, from a vk-sandbox that predates the field, counts as fresh so an upgrade never empties the fleet. InventoryApplier gains Delete (not-found is success), for vk-sandbox to remove its inventory on a graceful exit. StaticInventorySource.Delete replaces Remove. sandbox-apiserver and sandbox-envd-proxy take --inventory-stale-after through kubeinventory.Options.AddFlags. --- api/v1beta1/nodeinventory_types.go | 3 + api/v1beta1/zz_generated.deepcopy.go | 1 + cmd/sandbox-apiserver/main.go | 5 +- cmd/sandbox-envd-proxy/main.go | 4 +- docs/api.md | 1 + docs/configuration.md | 14 ++++ docs/e2b-compat.md | 7 ++ docs/scaling-design.md | 2 + ...andbox.cocoonstack.io_nodeinventories.yaml | 3 + pkg/scale/inventorycache_bench_test.go | 3 +- pkg/scale/kubeinventory/applier.go | 10 +++ pkg/scale/kubeinventory/applier_test.go | 17 +++++ pkg/scale/kubeinventory/source.go | 68 ++++++++++++++--- pkg/scale/kubeinventory/source_test.go | 76 +++++++++++++++++++ pkg/scale/sandboxstore_impl.go | 6 +- pkg/scale/sandboxstore_impl_test.go | 2 +- 16 files changed, 204 insertions(+), 18 deletions(-) create mode 100644 pkg/scale/kubeinventory/source_test.go diff --git a/api/v1beta1/nodeinventory_types.go b/api/v1beta1/nodeinventory_types.go index d192c6d..33f4f1d 100644 --- a/api/v1beta1/nodeinventory_types.go +++ b/api/v1beta1/nodeinventory_types.go @@ -87,6 +87,9 @@ type NodeInventory struct { // already holds a warm microVM for a requested (template, net, size). // +optional Pools []PoolCapacity `json:"pools,omitempty"` + // publishedAt is the publisher's last apply; the aggregated apiserver skips an inventory older than its staleness window. + // +optional + PublishedAt metav1.Time `json:"publishedAt,omitempty"` } // +kubebuilder:object:root=true diff --git a/api/v1beta1/zz_generated.deepcopy.go b/api/v1beta1/zz_generated.deepcopy.go index e5a6822..5858c20 100644 --- a/api/v1beta1/zz_generated.deepcopy.go +++ b/api/v1beta1/zz_generated.deepcopy.go @@ -63,6 +63,7 @@ func (in *NodeInventory) DeepCopyInto(out *NodeInventory) { *out = make([]PoolCapacity, len(*in)) copy(*out, *in) } + in.PublishedAt.DeepCopyInto(&out.PublishedAt) } // DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new NodeInventory. diff --git a/cmd/sandbox-apiserver/main.go b/cmd/sandbox-apiserver/main.go index fd9fcf8..ee19d5c 100644 --- a/cmd/sandbox-apiserver/main.go +++ b/cmd/sandbox-apiserver/main.go @@ -53,6 +53,8 @@ type options struct { E2BAPI bool E2B *e2bcompat.Flags + + Inventory kubeinventory.Options } func newOptions() *options { @@ -88,6 +90,7 @@ func (o *options) addFlags(fs *pflag.FlagSet) { fs.BoolVar(&o.E2BAPI, "enable-e2b-api", o.E2BAPI, "Serve the e2b-compatible REST surface, so an unmodified e2b SDK can claim from the same warm pools (point E2B_API_URL at it).") o.E2B.AddFlags(fs) + o.Inventory.AddFlags(fs) } func (o *options) serverConfig() (*genericapiserver.Config, error) { @@ -139,7 +142,7 @@ func run() error { if err != nil { return err } - invSource := kubeinventory.New(reader) + invSource := kubeinventory.New(reader, o.Inventory) store := scale.NewScatterGatherStore( invSource, scale.WithClaimRouting(token, scale.NewSandboxdClientFactory()), diff --git a/cmd/sandbox-envd-proxy/main.go b/cmd/sandbox-envd-proxy/main.go index c6a1149..d7d418b 100644 --- a/cmd/sandbox-envd-proxy/main.go +++ b/cmd/sandbox-envd-proxy/main.go @@ -32,10 +32,12 @@ type options struct { Domain string Namespace string Proxy envdproxy.Flags + Inventory kubeinventory.Options } func (o *options) addFlags(fs *pflag.FlagSet) { o.Proxy.AddFlags(fs, "") + o.Inventory.AddFlags(fs) fs.StringVar(&o.Domain, "domain", o.Domain, "Base domain sandbox hosts are derived from, as {port}-{sandboxID}.{domain}. Must match the apiserver's --e2b-domain.") fs.StringVar(&o.Namespace, "namespace", o.Namespace, @@ -70,7 +72,7 @@ func run(ctx context.Context, o *options) error { if err != nil { return err } - inv := kubeinventory.New(reader) + inv := kubeinventory.New(reader, o.Inventory) resolver, err := envdproxy.NewResolver(scale.NewScatterGatherStore(inv), inv, o.Namespace) if err != nil { return err diff --git a/docs/api.md b/docs/api.md index 58493aa..3977722 100644 --- a/docs/api.md +++ b/docs/api.md @@ -67,6 +67,7 @@ _Appears in:_ | `entries` _[InventoryEntry](#inventoryentry) array_ | entries summarizes the node's live sandboxes. | | Optional: \{\}
| | `address` _string_ | address is the node's sandboxd advertise address ("host:port"); the
aggregated apiserver routes a claim to this node's sandboxd through it. | | Optional: \{\}
| | `pools` _[PoolCapacity](#poolcapacity) array_ | pools is the node's per-pool warm capacity, used to pick a node that
already holds a warm microVM for a requested (template, net, size). | | Optional: \{\}
| +| `publishedAt` _[Time](https://kubernetes.io/docs/reference/generated/kubernetes-api/v1.37/#time-v1-meta)_ | publishedAt is the publisher's last apply; the aggregated apiserver skips an inventory older than its staleness window. | | Optional: \{\}
| #### NodeInventoryList diff --git a/docs/configuration.md b/docs/configuration.md index 6a1f595..453dd67 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -134,6 +134,19 @@ The driver watches `SandboxWarmPool` and `NodeInventory`, resolves each pool's under leader election (`cocoon-warmpool-driver`), so one replica drives the pools. An empty sandboxd token leaves it fail-closed: it logs and sets no pools. +### Node inventory freshness + +| Flag | Default | Help | +|---|---|---| +| `--inventory-stale-after` | `90s` | Drop a node from claims, lists, lookups, the warm-pool driver and the envd-proxy probe once its NodeInventory `publishedAt` is older than this. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | + +A node that stops publishing, because it died or its vk-sandbox stopped, +leaves every read and claim path once its `publishedAt` passes this age. At +vk-sandbox's default 30 s publish cadence that is three missed publishes. Keep +it above vk-sandbox's `-publish-interval`. A vk-sandbox that predates +`publishedAt` never goes stale, so the window starts working on each node as +its vk-sandbox is upgraded. + ### e2b-compatible surface | Flag | Default | Help | @@ -176,6 +189,7 @@ informer and relays the request into the owning node's guest-port endpoint. | `--tls-cert-file` | — | Wildcard certificate for `*.{domain}`. Omit to serve cleartext h2c behind an edge that terminates TLS. | | `--tls-private-key-file` | — | Private key for `--tls-cert-file`. | | `--guest-http2` | `false` | Forward to the guest over cleartext HTTP/2. Off by default: envd 0.8.0 installs no h2c handler and refuses it. Clients still reach this proxy over HTTP/2. | +| `--inventory-stale-after` | `90s` | Drop a node from claims, lists, lookups, the warm-pool driver and the envd-proxy probe once its NodeInventory `publishedAt` is older than this. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | The two TLS flags must be set together. Routing, authorization and failure mapping are in [envd-proxy](envd-proxy.md). diff --git a/docs/e2b-compat.md b/docs/e2b-compat.md index 297ead5..f8b5227 100644 --- a/docs/e2b-compat.md +++ b/docs/e2b-compat.md @@ -143,6 +143,13 @@ const sandbox = await Sandbox.create('registry.example.com/rt:24.04') once by `files`/`commands` works. The proxy's node lookups share one budget per replica (200 new sandboxes/s, see [envd-proxy](envd-proxy.md)); past it a sandbox its node has not published answers `502` until the node publishes. +- **A node that stops publishing leaves after `--inventory-stale-after`** + (90 s by default). Until then a dead node's sandboxes stay listed and a + claim can still sample it; past it they leave `GET /sandboxes`, and + `GET`, `DELETE` and the verbs on them answer `404`. A vk-sandbox that stops + gracefully deletes its inventory at once, so its restart shows the node's + sandboxes absent for those seconds. A vk-sandbox that predates + `publishedAt` never goes stale. - **`envdVersion`** is reported as `0.4.0` unless `--e2b-envd-version` says otherwise. The SDK version-compares it and *kills the sandbox* if it cannot parse it, so it is always sent. Set it to the version actually installed in diff --git a/docs/scaling-design.md b/docs/scaling-design.md index 566c55e..6bfb82e 100644 --- a/docs/scaling-design.md +++ b/docs/scaling-design.md @@ -161,6 +161,8 @@ v1beta1 the `APIService` hands to the aggregated server, which serves only | Node partitioned from aggregated server | Its sandboxes briefly absent from `List` (eventual consistency, same as an informer lag) | No | | A node `inventory` object lost | Rebuilt from the node's own live state on next publish | No | | A node's `Node` object deleted | Its `NodeInventory` is garbage-collected with it: the node leaves the claim path and its sandboxes leave the read view, though sandboxd keeps serving them and their leases still expire there; restarting vk-sandbox registers the node again | No | +| A node dies, its `Node` object kept | Its `NodeInventory` is no longer republished. Once `publishedAt` is older than `--inventory-stale-after` (90 s), the node leaves the claim path and its sandboxes leave the read view, the watch emits Deleted for each, Get and the verbs answer 404, and the warm-pool driver and envd-proxy stop calling it. Its next publish brings it back | No | +| vk-sandbox restarts on a live node | On SIGTERM vk-sandbox deletes its `NodeInventory`, and it republishes on start. For those seconds the node's sandboxes leave the list, Get and connect answer 404, and the watch emits Deleted then Added for each, while sandboxd keeps serving them | No | | Aggregated server restart | Stateless; rebuilds from node fan-out | No | | Client reads before the owning node republishes inventory | Lookups by name and by claim id (e2b, envd proxy), and a list or watch pinned to one name in a namespace, ask the nodes; other lists and watches show it at the next publish | No — fleet `list` and `watch`, and a deleted sandbox until its node publishes, stay eventually consistent | diff --git a/helm/crds/sandbox.cocoonstack.io_nodeinventories.yaml b/helm/crds/sandbox.cocoonstack.io_nodeinventories.yaml index 0449971..f851098 100644 --- a/helm/crds/sandbox.cocoonstack.io_nodeinventories.yaml +++ b/helm/crds/sandbox.cocoonstack.io_nodeinventories.yaml @@ -80,6 +80,9 @@ spec: - warm type: object type: array + publishedAt: + format: date-time + type: string required: - node type: object diff --git a/pkg/scale/inventorycache_bench_test.go b/pkg/scale/inventorycache_bench_test.go index c6a2fac..61b8834 100644 --- a/pkg/scale/inventorycache_bench_test.go +++ b/pkg/scale/inventorycache_bench_test.go @@ -96,6 +96,7 @@ func benchCachedStore(b *testing.B, nodes, perNode int, noCopy bool) (*scale.Sca invs, pool := scale.BenchInventories(nodes, perNode) list := &unstructured.UnstructuredList{} for _, inv := range invs { + inv.PublishedAt = metav1.Now() raw, err := runtime.DefaultUnstructuredConverter.ToUnstructured(inv) if err != nil { b.Fatalf("encode node inventory: %v", err) @@ -137,5 +138,5 @@ func benchCachedStore(b *testing.B, nodes, perNode int, noCopy bool) (*scale.Sca if !invCache.WaitForCacheSync(syncCtx) { b.Fatal("node inventory cache did not sync") } - return scale.NewScatterGatherStore(kubeinventory.New(invCache)).(*scale.ScatterGatherStore), pool + return scale.NewScatterGatherStore(kubeinventory.New(invCache, kubeinventory.Options{})).(*scale.ScatterGatherStore), pool } diff --git a/pkg/scale/kubeinventory/applier.go b/pkg/scale/kubeinventory/applier.go index 8c05068..74351a3 100644 --- a/pkg/scale/kubeinventory/applier.go +++ b/pkg/scale/kubeinventory/applier.go @@ -39,3 +39,13 @@ func (a *ssaApplier) Apply(ctx context.Context, inv *scale.NodeInventory) error } return nil } + +func (a *ssaApplier) Delete(ctx context.Context, node string) error { + u := &unstructured.Unstructured{} + u.SetGroupVersionKind(scale.NodeInventoryGVK) + u.SetName(node) + if err := client.IgnoreNotFound(a.c.Delete(ctx, u)); err != nil { + return fmt.Errorf("kubeinventory: delete node %q inventory: %w", node, err) + } + return nil +} diff --git a/pkg/scale/kubeinventory/applier_test.go b/pkg/scale/kubeinventory/applier_test.go index d050f5e..69c3a0e 100644 --- a/pkg/scale/kubeinventory/applier_test.go +++ b/pkg/scale/kubeinventory/applier_test.go @@ -5,6 +5,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + k8serrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" @@ -37,3 +38,19 @@ func TestSSAApplier_UpsertsOneObjectPerNode(t *testing.T) { require.Len(t, list.Items, 1) assert.Len(t, list.Items[0].Entries, 2) } + +func TestSSAApplier_DeleteTreatsAMissingInventoryAsDone(t *testing.T) { + ctx := t.Context() + scheme := runtime.NewScheme() + require.NoError(t, cocoonv1beta1.AddToScheme(scheme)) + cli := fake.NewClientBuilder().WithScheme(scheme).Build() + applier := NewSSAApplier(cli, "vk-test") + require.NoError(t, applier.Apply(ctx, &scale.NodeInventory{ + Kind: scale.NodeInventoryGVK.Kind, APIVersion: scale.NodeInventoryGVK.GroupVersion().String(), Name: "n1", Node: "n1", + })) + + require.NoError(t, applier.Delete(ctx, "n1")) + err := cli.Get(ctx, client.ObjectKey{Name: "n1"}, &cocoonv1beta1.NodeInventory{}) + assert.True(t, k8serrors.IsNotFound(err), "inventory after delete: %v", err) + require.NoError(t, applier.Delete(ctx, "n1"), "a missing inventory is success") +} diff --git a/pkg/scale/kubeinventory/source.go b/pkg/scale/kubeinventory/source.go index 15ecddd..c9afc8d 100644 --- a/pkg/scale/kubeinventory/source.go +++ b/pkg/scale/kubeinventory/source.go @@ -2,29 +2,51 @@ package kubeinventory import ( + "cmp" "context" "fmt" "slices" + "time" + "github.com/spf13/pflag" + k8serrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/types" "sigs.k8s.io/controller-runtime/pkg/client" "github.com/cocoonstack/sandbox-operator/pkg/scale" ) +const defaultStaleAfter = 90 * time.Second + +var nodeInventories = schema.GroupResource{Group: scale.NodeInventoryGVK.Group, Resource: "nodeinventories"} + +// Options tunes a Source. +type Options struct { + // StaleAfter is how old publishedAt may get before its node leaves the fleet; zero means 90 s. + StaleAfter time.Duration +} + +// AddFlags registers the source flags on fs. +func (o *Options) AddFlags(fs *pflag.FlagSet) { + fs.DurationVar(&o.StaleAfter, "inventory-stale-after", cmp.Or(o.StaleAfter, defaultStaleAfter), + "Drop a node from claims, lists, lookups, the warm-pool driver and the envd-proxy probe once its NodeInventory publishedAt is older than this. An inventory without publishedAt, from a vk-sandbox that predates the field, always stays.") +} + var _ scale.InventorySource = (*Source)(nil) // Source is the production InventorySource over a cache-fed NodeInventory reader. // It never mutates what it reads, since NewCache hands out its cached objects themselves. type Source struct { - reader client.Reader + reader client.Reader + staleAfter time.Duration } // New builds a Source over reader. -func New(reader client.Reader) *Source { - return &Source{reader: reader} +func New(reader client.Reader, opts Options) *Source { + return &Source{reader: reader, staleAfter: cmp.Or(opts.StaleAfter, defaultStaleAfter)} } func (s *Source) ListNodes(ctx context.Context) ([]string, error) { @@ -34,18 +56,20 @@ func (s *Source) ListNodes(ctx context.Context) ([]string, error) { return nil, fmt.Errorf("kubeinventory: list node inventories: %w", err) } nodes := make([]string, 0, len(ul.Items)) + now := time.Now() for i := range ul.Items { - nodes = append(nodes, ul.Items[i].GetName()) + if !s.stale(ul.Items[i].Object, now) { + nodes = append(nodes, ul.Items[i].GetName()) + } } slices.Sort(nodes) return nodes, nil } func (s *Source) NodeInventory(ctx context.Context, node string) (*scale.NodeInventory, error) { - u := &unstructured.Unstructured{} - u.SetGroupVersionKind(scale.NodeInventoryGVK) - if err := s.reader.Get(ctx, types.NamespacedName{Name: node}, u); err != nil { - return nil, fmt.Errorf("kubeinventory: get node %q inventory: %w", node, err) + u, err := s.get(ctx, node) + if err != nil { + return nil, err } inv := &scale.NodeInventory{} if err := runtime.DefaultUnstructuredConverter.FromUnstructured(u.Object, inv); err != nil { @@ -55,10 +79,9 @@ func (s *Source) NodeInventory(ctx context.Context, node string) (*scale.NodeInv } func (s *Source) NodeCapacity(ctx context.Context, node string) (string, []scale.PoolCapacity, error) { - u := &unstructured.Unstructured{} - u.SetGroupVersionKind(scale.NodeInventoryGVK) - if err := s.reader.Get(ctx, types.NamespacedName{Name: node}, u); err != nil { - return "", nil, fmt.Errorf("kubeinventory: get node %q inventory: %w", node, err) + u, err := s.get(ctx, node) + if err != nil { + return "", nil, err } addr, _, err := unstructured.NestedString(u.Object, "address") if err != nil { @@ -82,3 +105,24 @@ func (s *Source) NodeCapacity(ctx context.Context, node string) (string, []scale } return addr, pools, nil } + +func (s *Source) get(ctx context.Context, node string) (*unstructured.Unstructured, error) { + u := &unstructured.Unstructured{} + u.SetGroupVersionKind(scale.NodeInventoryGVK) + if err := s.reader.Get(ctx, types.NamespacedName{Name: node}, u); err != nil { + return nil, fmt.Errorf("kubeinventory: get node %q inventory: %w", node, err) + } + if s.stale(u.Object, time.Now()) { + return nil, fmt.Errorf("kubeinventory: node %q inventory is stale: %w", node, k8serrors.NewNotFound(nodeInventories, node)) + } + return u, nil +} + +func (s *Source) stale(obj map[string]any, now time.Time) bool { + raw, _ := obj["publishedAt"].(string) + if raw == "" { + return false + } + at, err := time.Parse(time.RFC3339, raw) + return err == nil && now.Sub(at) > s.staleAfter +} diff --git a/pkg/scale/kubeinventory/source_test.go b/pkg/scale/kubeinventory/source_test.go new file mode 100644 index 0000000..2f5e407 --- /dev/null +++ b/pkg/scale/kubeinventory/source_test.go @@ -0,0 +1,76 @@ +package kubeinventory + +import ( + "testing" + "testing/synctest" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + k8serrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + cocoonv1beta1 "github.com/cocoonstack/sandbox-operator/api/v1beta1" + "github.com/cocoonstack/sandbox-operator/pkg/scale" +) + +func TestSourceDropsAnInventoryPastTheWindow(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ctx := t.Context() + src := New(fakeReader(t, inventory("fresh", metav1.Now()), inventory("unstamped", metav1.Time{})), Options{}) + + time.Sleep(90 * time.Second) + nodes, err := src.ListNodes(ctx) + require.NoError(t, err) + assert.Equal(t, []string{"fresh", "unstamped"}, nodes) + _, err = src.NodeInventory(ctx, "fresh") + require.NoError(t, err) + addr, pools, err := src.NodeCapacity(ctx, "fresh") + require.NoError(t, err) + assert.Equal(t, "fresh:7777", addr) + assert.Len(t, pools, 1) + + time.Sleep(500 * time.Millisecond) + nodes, err = src.ListNodes(ctx) + require.NoError(t, err) + assert.Equal(t, []string{"unstamped"}, nodes) + _, err = src.NodeInventory(ctx, "fresh") + assert.True(t, k8serrors.IsNotFound(err), "NodeInventory on a stale node: %v", err) + _, _, err = src.NodeCapacity(ctx, "fresh") + assert.True(t, k8serrors.IsNotFound(err), "NodeCapacity on a stale node: %v", err) + + time.Sleep(time.Hour) + _, _, err = src.NodeCapacity(ctx, "unstamped") + require.NoError(t, err, "an inventory without publishedAt never goes stale") + }) +} + +func TestSourceHonorsStaleAfter(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + src := New(fakeReader(t, inventory("n1", metav1.Now())), Options{StaleAfter: 10 * time.Second}) + time.Sleep(11 * time.Second) + nodes, err := src.ListNodes(t.Context()) + require.NoError(t, err) + assert.Empty(t, nodes) + }) +} + +func fakeReader(t *testing.T, objs ...client.Object) client.Reader { + t.Helper() + scheme := runtime.NewScheme() + require.NoError(t, cocoonv1beta1.AddToScheme(scheme)) + return fake.NewClientBuilder().WithScheme(scheme).WithObjects(objs...).Build() +} + +func inventory(node string, publishedAt metav1.Time) *cocoonv1beta1.NodeInventory { + return &cocoonv1beta1.NodeInventory{ + Name: node, + Node: node, + Address: node + ":7777", + Pools: []scale.PoolCapacity{{Template: "rt", Warm: 1, Target: 1}}, + PublishedAt: publishedAt, + } +} diff --git a/pkg/scale/sandboxstore_impl.go b/pkg/scale/sandboxstore_impl.go index 4a0d1b6..cae2f67 100644 --- a/pkg/scale/sandboxstore_impl.go +++ b/pkg/scale/sandboxstore_impl.go @@ -539,6 +539,8 @@ type NodeLiveSource interface { // InventoryApplier server-side-applies a NodeInventory object. type InventoryApplier interface { Apply(ctx context.Context, inv *NodeInventory) error + // Delete removes node's inventory, and a missing one is success. + Delete(ctx context.Context, node string) error } var ( @@ -584,12 +586,12 @@ func (s *StaticInventorySource) Partition(node string) { s.partition[node] = struct{}{} } -// Remove drops a node's inventory object entirely (lost inventory). -func (s *StaticInventorySource) Remove(node string) { +func (s *StaticInventorySource) Delete(_ context.Context, node string) error { s.mu.Lock() defer s.mu.Unlock() delete(s.inv, node) delete(s.partition, node) + return nil } func (s *StaticInventorySource) ListNodes(_ context.Context) ([]string, error) { diff --git a/pkg/scale/sandboxstore_impl_test.go b/pkg/scale/sandboxstore_impl_test.go index 485d94f..f9cdeb0 100644 --- a/pkg/scale/sandboxstore_impl_test.go +++ b/pkg/scale/sandboxstore_impl_test.go @@ -286,7 +286,7 @@ func TestPublisher_RebuildsFromLiveAfterLoss(t *testing.T) { require.NoError(t, err) require.Equal(t, 1, src.ObjectCount()) - src.Remove("n1") + require.NoError(t, src.Delete(ctx, "n1")) require.Equal(t, 0, src.ObjectCount()) live.entries = []InventoryEntry{entry("ns/a", "Running"), entry("ns/b", "Running")} From 25d0e9409a8ba21630d529fb19a17da4ff7be607 Mon Sep 17 00:00:00 2001 From: CMGS Date: Sun, 27 Sep 2026 17:20:59 +0800 Subject: [PATCH 2/4] review: document the CRD step, the clock requirement, and the restart costs; flag help names its own process Helm never upgrades crds/, and an older CRD prunes publishedAt so no node goes stale; the age uses the reading host's clock; a vk-sandbox restart also drops an e2b kill (404), answers 502 on the data plane, and moves the node's warm-pool share. The shared flag help now scopes itself to the process that parses it. --- docs/configuration.md | 14 +++++++++----- docs/e2b-compat.md | 9 +++++---- docs/scaling-design.md | 2 +- helm/README.md | 7 ++++--- pkg/scale/kubeinventory/source.go | 2 +- 5 files changed, 20 insertions(+), 14 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index 453dd67..45ab14d 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -138,14 +138,18 @@ pools. An empty sandboxd token leaves it fail-closed: it logs and sets no pools. | Flag | Default | Help | |---|---|---| -| `--inventory-stale-after` | `90s` | Drop a node from claims, lists, lookups, the warm-pool driver and the envd-proxy probe once its NodeInventory `publishedAt` is older than this. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | +| `--inventory-stale-after` | `90s` | Drop a node from this process's inventory reads once its NodeInventory `publishedAt` is older than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | A node that stops publishing, because it died or its vk-sandbox stopped, leaves every read and claim path once its `publishedAt` passes this age. At vk-sandbox's default 30 s publish cadence that is three missed publishes. Keep -it above vk-sandbox's `-publish-interval`. A vk-sandbox that predates -`publishedAt` never goes stale, so the window starts working on each node as -its vk-sandbox is upgraded. +it above vk-sandbox's `--publish-interval`. The age is measured against this +host's clock, so the nodes and the control plane need synchronized clocks. + +The window works on a node once the `nodeinventories` CRD from `helm/crds` is +applied and that node runs a vk-sandbox that stamps `publishedAt`. An older CRD +prunes the field and an older vk-sandbox never sets it, and either way that +node never goes stale, so an upgrade never empties the fleet. ### e2b-compatible surface @@ -189,7 +193,7 @@ informer and relays the request into the owning node's guest-port endpoint. | `--tls-cert-file` | — | Wildcard certificate for `*.{domain}`. Omit to serve cleartext h2c behind an edge that terminates TLS. | | `--tls-private-key-file` | — | Private key for `--tls-cert-file`. | | `--guest-http2` | `false` | Forward to the guest over cleartext HTTP/2. Off by default: envd 0.8.0 installs no h2c handler and refuses it. Clients still reach this proxy over HTTP/2. | -| `--inventory-stale-after` | `90s` | Drop a node from claims, lists, lookups, the warm-pool driver and the envd-proxy probe once its NodeInventory `publishedAt` is older than this. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | +| `--inventory-stale-after` | `90s` | Drop a node from this process's inventory reads once its NodeInventory `publishedAt` is older than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | The two TLS flags must be set together. Routing, authorization and failure mapping are in [envd-proxy](envd-proxy.md). diff --git a/docs/e2b-compat.md b/docs/e2b-compat.md index f8b5227..63604ec 100644 --- a/docs/e2b-compat.md +++ b/docs/e2b-compat.md @@ -146,10 +146,11 @@ const sandbox = await Sandbox.create('registry.example.com/rt:24.04') - **A node that stops publishing leaves after `--inventory-stale-after`** (90 s by default). Until then a dead node's sandboxes stay listed and a claim can still sample it; past it they leave `GET /sandboxes`, and - `GET`, `DELETE` and the verbs on them answer `404`. A vk-sandbox that stops - gracefully deletes its inventory at once, so its restart shows the node's - sandboxes absent for those seconds. A vk-sandbox that predates - `publishedAt` never goes stale. + `GET`, `DELETE` and the verbs on them answer `404`. A vk-sandbox that exits + deletes its inventory at once, so while it restarts the node's live + sandboxes are absent: `DELETE` answers `404` and the kill is dropped, and + the data plane answers `502`. A vk-sandbox that predates `publishedAt` never + goes stale. - **`envdVersion`** is reported as `0.4.0` unless `--e2b-envd-version` says otherwise. The SDK version-compares it and *kills the sandbox* if it cannot parse it, so it is always sent. Set it to the version actually installed in diff --git a/docs/scaling-design.md b/docs/scaling-design.md index 6bfb82e..73c5dc1 100644 --- a/docs/scaling-design.md +++ b/docs/scaling-design.md @@ -162,7 +162,7 @@ v1beta1 the `APIService` hands to the aggregated server, which serves only | A node `inventory` object lost | Rebuilt from the node's own live state on next publish | No | | A node's `Node` object deleted | Its `NodeInventory` is garbage-collected with it: the node leaves the claim path and its sandboxes leave the read view, though sandboxd keeps serving them and their leases still expire there; restarting vk-sandbox registers the node again | No | | A node dies, its `Node` object kept | Its `NodeInventory` is no longer republished. Once `publishedAt` is older than `--inventory-stale-after` (90 s), the node leaves the claim path and its sandboxes leave the read view, the watch emits Deleted for each, Get and the verbs answer 404, and the warm-pool driver and envd-proxy stop calling it. Its next publish brings it back | No | -| vk-sandbox restarts on a live node | On SIGTERM vk-sandbox deletes its `NodeInventory`, and it republishes on start. For those seconds the node's sandboxes leave the list, Get and connect answer 404, and the watch emits Deleted then Added for each, while sandboxd keeps serving them | No | +| vk-sandbox restarts on a live node | On exit vk-sandbox deletes its `NodeInventory`, and it republishes on start. For those seconds sandboxd keeps serving the node's sandboxes, but they leave the list, Get, connect and an e2b kill answer 404 (the kill is dropped and the sandbox runs to its lease), envd-proxy answers 502, the watch emits Deleted then Added for each, and the warm-pool driver moves the node's share of each pool to the other nodes and back | No | | Aggregated server restart | Stateless; rebuilds from node fan-out | No | | Client reads before the owning node republishes inventory | Lookups by name and by claim id (e2b, envd proxy), and a list or watch pinned to one name in a namespace, ask the nodes; other lists and watches show it at the next publish | No — fleet `list` and `watch`, and a deleted sandbox until its node publishes, stay eventually consistent | diff --git a/helm/README.md b/helm/README.md index 59525ac..150972a 100644 --- a/helm/README.md +++ b/helm/README.md @@ -51,9 +51,10 @@ kubectl apply -f helm/crds/ helm upgrade sandbox-operator ./helm --namespace sandbox-system --reuse-values ``` -An older `nodeinventories` CRD silently prunes entry fields it does not know -(`claimedAt` is one), so until the CRD is applied every node publishes like an -older node and the read view falls back accordingly. +An older `nodeinventories` CRD silently prunes fields it does not know +(`claimedAt` and `publishedAt` are two), so until the CRD is applied every node +publishes like an older node: the read view falls back accordingly, and no node +ever goes stale under `--inventory-stale-after`. `NodeInventory` moved from `extensions.agents.x-k8s.io` to `sandbox.cocoonstack.io`. A fleet coming from the old group runs a vk-sandbox diff --git a/pkg/scale/kubeinventory/source.go b/pkg/scale/kubeinventory/source.go index c9afc8d..5b766aa 100644 --- a/pkg/scale/kubeinventory/source.go +++ b/pkg/scale/kubeinventory/source.go @@ -32,7 +32,7 @@ type Options struct { // AddFlags registers the source flags on fs. func (o *Options) AddFlags(fs *pflag.FlagSet) { fs.DurationVar(&o.StaleAfter, "inventory-stale-after", cmp.Or(o.StaleAfter, defaultStaleAfter), - "Drop a node from claims, lists, lookups, the warm-pool driver and the envd-proxy probe once its NodeInventory publishedAt is older than this. An inventory without publishedAt, from a vk-sandbox that predates the field, always stays.") + "Drop a node from this process's inventory reads once its NodeInventory publishedAt is older than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without publishedAt, from a vk-sandbox that predates the field, always stays.") } var _ scale.InventorySource = (*Source)(nil) From 835f250c728380d72395dc8915e6625af00c5b8f Mon Sep 17 00:00:00 2001 From: CMGS Date: Sun, 27 Sep 2026 18:05:21 +0800 Subject: [PATCH 3/4] feat: measure NodeInventory staleness against the newest publish, and drop the applier delete CMGS's K1 decision. The reference is min(newest publishedAt ListNodes saw, now), kept on the Source as unix nanos; an inventory whose stamp trails it by more than StaleAfter is dropped from ListNodes, and NodeInventory and NodeCapacity compare against the last stored reference (zero before the first list, which leaves every node fresh). A control-plane write outage freezes every stamp together, so it never empties the fleet; a publisher clock ahead of the reader's is capped by now and cannot evict the others; the reader's clock no longer decides a node's age. Missing or unparsable stamps stay fresh. Option 3 is dropped: InventoryApplier is Apply-only again and StaticInventorySource.Remove is restored, so a vk-sandbox restart inside the window changes nothing. Docs drop the restart flap and gain the outage row. --- docs/configuration.md | 18 ++++-- docs/e2b-compat.md | 9 ++- docs/scaling-design.md | 4 +- pkg/scale/kubeinventory/applier.go | 10 --- pkg/scale/kubeinventory/applier_test.go | 17 ----- pkg/scale/kubeinventory/source.go | 39 ++++++++---- pkg/scale/kubeinventory/source_test.go | 82 +++++++++++++++++++------ pkg/scale/sandboxstore_impl.go | 6 +- pkg/scale/sandboxstore_impl_test.go | 2 +- 9 files changed, 113 insertions(+), 74 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index 45ab14d..22bd893 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -138,13 +138,19 @@ pools. An empty sandboxd token leaves it fail-closed: it logs and sets no pools. | Flag | Default | Help | |---|---|---| -| `--inventory-stale-after` | `90s` | Drop a node from this process's inventory reads once its NodeInventory `publishedAt` is older than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | +| `--inventory-stale-after` | `90s` | Drop a node from this process's inventory reads once its NodeInventory `publishedAt` trails the newest publish in the fleet by more than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | A node that stops publishing, because it died or its vk-sandbox stopped, -leaves every read and claim path once its `publishedAt` passes this age. At -vk-sandbox's default 30 s publish cadence that is three missed publishes. Keep -it above vk-sandbox's `--publish-interval`. The age is measured against this -host's clock, so the nodes and the control plane need synchronized clocks. +leaves every read and claim path once its `publishedAt` trails the newest +publish by more than this. At vk-sandbox's default 30 s publish cadence that is +three missed publishes. Keep it above vk-sandbox's `--publish-interval`. + +The age is measured against the newest publish the process has seen, capped by +its own clock. A control-plane write outage freezes every stamp together, so it +never empties the fleet, and a node's clock matters only relative to the other +nodes: one running ahead cannot evict the rest. A fleet whose nodes all stop +publishing at once, or a one-node fleet, keeps its last inventory until the +nodes publish again. The window works on a node once the `nodeinventories` CRD from `helm/crds` is applied and that node runs a vk-sandbox that stamps `publishedAt`. An older CRD @@ -193,7 +199,7 @@ informer and relays the request into the owning node's guest-port endpoint. | `--tls-cert-file` | — | Wildcard certificate for `*.{domain}`. Omit to serve cleartext h2c behind an edge that terminates TLS. | | `--tls-private-key-file` | — | Private key for `--tls-cert-file`. | | `--guest-http2` | `false` | Forward to the guest over cleartext HTTP/2. Off by default: envd 0.8.0 installs no h2c handler and refuses it. Clients still reach this proxy over HTTP/2. | -| `--inventory-stale-after` | `90s` | Drop a node from this process's inventory reads once its NodeInventory `publishedAt` is older than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | +| `--inventory-stale-after` | `90s` | Drop a node from this process's inventory reads once its NodeInventory `publishedAt` trails the newest publish in the fleet by more than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without `publishedAt`, from a vk-sandbox that predates the field, always stays. | The two TLS flags must be set together. Routing, authorization and failure mapping are in [envd-proxy](envd-proxy.md). diff --git a/docs/e2b-compat.md b/docs/e2b-compat.md index 63604ec..b0aeb89 100644 --- a/docs/e2b-compat.md +++ b/docs/e2b-compat.md @@ -146,11 +146,10 @@ const sandbox = await Sandbox.create('registry.example.com/rt:24.04') - **A node that stops publishing leaves after `--inventory-stale-after`** (90 s by default). Until then a dead node's sandboxes stay listed and a claim can still sample it; past it they leave `GET /sandboxes`, and - `GET`, `DELETE` and the verbs on them answer `404`. A vk-sandbox that exits - deletes its inventory at once, so while it restarts the node's live - sandboxes are absent: `DELETE` answers `404` and the kill is dropped, and - the data plane answers `502`. A vk-sandbox that predates `publishedAt` never - goes stale. + `GET`, `DELETE` and the verbs on them answer `404`. The window is measured + against the fleet's newest publish, so a control-plane outage never drops + a node, and a vk-sandbox restart inside the window is invisible. A + vk-sandbox that predates `publishedAt` never goes stale. - **`envdVersion`** is reported as `0.4.0` unless `--e2b-envd-version` says otherwise. The SDK version-compares it and *kills the sandbox* if it cannot parse it, so it is always sent. Set it to the version actually installed in diff --git a/docs/scaling-design.md b/docs/scaling-design.md index 73c5dc1..b8f98b9 100644 --- a/docs/scaling-design.md +++ b/docs/scaling-design.md @@ -161,8 +161,8 @@ v1beta1 the `APIService` hands to the aggregated server, which serves only | Node partitioned from aggregated server | Its sandboxes briefly absent from `List` (eventual consistency, same as an informer lag) | No | | A node `inventory` object lost | Rebuilt from the node's own live state on next publish | No | | A node's `Node` object deleted | Its `NodeInventory` is garbage-collected with it: the node leaves the claim path and its sandboxes leave the read view, though sandboxd keeps serving them and their leases still expire there; restarting vk-sandbox registers the node again | No | -| A node dies, its `Node` object kept | Its `NodeInventory` is no longer republished. Once `publishedAt` is older than `--inventory-stale-after` (90 s), the node leaves the claim path and its sandboxes leave the read view, the watch emits Deleted for each, Get and the verbs answer 404, and the warm-pool driver and envd-proxy stop calling it. Its next publish brings it back | No | -| vk-sandbox restarts on a live node | On exit vk-sandbox deletes its `NodeInventory`, and it republishes on start. For those seconds sandboxd keeps serving the node's sandboxes, but they leave the list, Get, connect and an e2b kill answer 404 (the kill is dropped and the sandbox runs to its lease), envd-proxy answers 502, the watch emits Deleted then Added for each, and the warm-pool driver moves the node's share of each pool to the other nodes and back | No | +| A node dies, its `Node` object kept | Its `NodeInventory` is no longer republished. Once its `publishedAt` trails the newest publish in the fleet by more than `--inventory-stale-after` (90 s), the node leaves the claim path and its sandboxes leave the read view, the watch emits Deleted for each, Get and the verbs answer 404, and the warm-pool driver and envd-proxy stop calling it. Its next publish brings it back. A vk-sandbox restart shorter than the window changes nothing | No | +| The control plane stops accepting writes | Every publish fails and every stamp freezes together, so no node goes stale: claims, lists and lookups keep running from the informer cache, which shows the inventory as of the outage | No | | Aggregated server restart | Stateless; rebuilds from node fan-out | No | | Client reads before the owning node republishes inventory | Lookups by name and by claim id (e2b, envd proxy), and a list or watch pinned to one name in a namespace, ask the nodes; other lists and watches show it at the next publish | No — fleet `list` and `watch`, and a deleted sandbox until its node publishes, stay eventually consistent | diff --git a/pkg/scale/kubeinventory/applier.go b/pkg/scale/kubeinventory/applier.go index 74351a3..8c05068 100644 --- a/pkg/scale/kubeinventory/applier.go +++ b/pkg/scale/kubeinventory/applier.go @@ -39,13 +39,3 @@ func (a *ssaApplier) Apply(ctx context.Context, inv *scale.NodeInventory) error } return nil } - -func (a *ssaApplier) Delete(ctx context.Context, node string) error { - u := &unstructured.Unstructured{} - u.SetGroupVersionKind(scale.NodeInventoryGVK) - u.SetName(node) - if err := client.IgnoreNotFound(a.c.Delete(ctx, u)); err != nil { - return fmt.Errorf("kubeinventory: delete node %q inventory: %w", node, err) - } - return nil -} diff --git a/pkg/scale/kubeinventory/applier_test.go b/pkg/scale/kubeinventory/applier_test.go index 69c3a0e..d050f5e 100644 --- a/pkg/scale/kubeinventory/applier_test.go +++ b/pkg/scale/kubeinventory/applier_test.go @@ -5,7 +5,6 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" - k8serrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" @@ -38,19 +37,3 @@ func TestSSAApplier_UpsertsOneObjectPerNode(t *testing.T) { require.Len(t, list.Items, 1) assert.Len(t, list.Items[0].Entries, 2) } - -func TestSSAApplier_DeleteTreatsAMissingInventoryAsDone(t *testing.T) { - ctx := t.Context() - scheme := runtime.NewScheme() - require.NoError(t, cocoonv1beta1.AddToScheme(scheme)) - cli := fake.NewClientBuilder().WithScheme(scheme).Build() - applier := NewSSAApplier(cli, "vk-test") - require.NoError(t, applier.Apply(ctx, &scale.NodeInventory{ - Kind: scale.NodeInventoryGVK.Kind, APIVersion: scale.NodeInventoryGVK.GroupVersion().String(), Name: "n1", Node: "n1", - })) - - require.NoError(t, applier.Delete(ctx, "n1")) - err := cli.Get(ctx, client.ObjectKey{Name: "n1"}, &cocoonv1beta1.NodeInventory{}) - assert.True(t, k8serrors.IsNotFound(err), "inventory after delete: %v", err) - require.NoError(t, applier.Delete(ctx, "n1"), "a missing inventory is success") -} diff --git a/pkg/scale/kubeinventory/source.go b/pkg/scale/kubeinventory/source.go index 5b766aa..b46beca 100644 --- a/pkg/scale/kubeinventory/source.go +++ b/pkg/scale/kubeinventory/source.go @@ -6,6 +6,7 @@ import ( "context" "fmt" "slices" + "sync/atomic" "time" "github.com/spf13/pflag" @@ -32,7 +33,7 @@ type Options struct { // AddFlags registers the source flags on fs. func (o *Options) AddFlags(fs *pflag.FlagSet) { fs.DurationVar(&o.StaleAfter, "inventory-stale-after", cmp.Or(o.StaleAfter, defaultStaleAfter), - "Drop a node from this process's inventory reads once its NodeInventory publishedAt is older than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without publishedAt, from a vk-sandbox that predates the field, always stays.") + "Drop a node from this process's inventory reads once its NodeInventory publishedAt trails the newest publish in the fleet by more than this; set the same value on sandbox-apiserver and sandbox-envd-proxy. An inventory without publishedAt, from a vk-sandbox that predates the field, always stays.") } var _ scale.InventorySource = (*Source)(nil) @@ -42,6 +43,7 @@ var _ scale.InventorySource = (*Source)(nil) type Source struct { reader client.Reader staleAfter time.Duration + reference atomic.Int64 } // New builds a Source over reader. @@ -55,15 +57,23 @@ func (s *Source) ListNodes(ctx context.Context) ([]string, error) { if err := s.reader.List(ctx, ul); err != nil { return nil, fmt.Errorf("kubeinventory: list node inventories: %w", err) } - nodes := make([]string, 0, len(ul.Items)) - now := time.Now() + nodes := make([]string, len(ul.Items)) + stamps := make([]int64, len(ul.Items)) + var newest int64 for i := range ul.Items { - if !s.stale(ul.Items[i].Object, now) { - nodes = append(nodes, ul.Items[i].GetName()) + nodes[i], stamps[i] = ul.Items[i].GetName(), publishedAt(ul.Items[i].Object) + newest = max(newest, stamps[i]) + } + ref := min(newest, time.Now().UnixNano()) + s.reference.Store(ref) + kept := nodes[:0] + for i, node := range nodes { + if !s.stale(stamps[i], ref) { + kept = append(kept, node) } } - slices.Sort(nodes) - return nodes, nil + slices.Sort(kept) + return kept, nil } func (s *Source) NodeInventory(ctx context.Context, node string) (*scale.NodeInventory, error) { @@ -112,17 +122,24 @@ func (s *Source) get(ctx context.Context, node string) (*unstructured.Unstructur if err := s.reader.Get(ctx, types.NamespacedName{Name: node}, u); err != nil { return nil, fmt.Errorf("kubeinventory: get node %q inventory: %w", node, err) } - if s.stale(u.Object, time.Now()) { + if s.stale(publishedAt(u.Object), s.reference.Load()) { return nil, fmt.Errorf("kubeinventory: node %q inventory is stale: %w", node, k8serrors.NewNotFound(nodeInventories, node)) } return u, nil } -func (s *Source) stale(obj map[string]any, now time.Time) bool { +func (s *Source) stale(stamp, ref int64) bool { + return stamp != 0 && ref-stamp > int64(s.staleAfter) +} + +func publishedAt(obj map[string]any) int64 { raw, _ := obj["publishedAt"].(string) if raw == "" { - return false + return 0 } at, err := time.Parse(time.RFC3339, raw) - return err == nil && now.Sub(at) > s.staleAfter + if err != nil { + return 0 + } + return at.UnixNano() } diff --git a/pkg/scale/kubeinventory/source_test.go b/pkg/scale/kubeinventory/source_test.go index 2f5e407..380b602 100644 --- a/pkg/scale/kubeinventory/source_test.go +++ b/pkg/scale/kubeinventory/source_test.go @@ -17,54 +17,100 @@ import ( "github.com/cocoonstack/sandbox-operator/pkg/scale" ) -func TestSourceDropsAnInventoryPastTheWindow(t *testing.T) { +func TestSourceDropsADeadNodeAgainstTheNewestPublish(t *testing.T) { synctest.Test(t, func(t *testing.T) { ctx := t.Context() - src := New(fakeReader(t, inventory("fresh", metav1.Now()), inventory("unstamped", metav1.Time{})), Options{}) + c := fakeClient(t, inventory("live", metav1.Now()), inventory("dead", metav1.Now()), inventory("unstamped", metav1.Time{})) + src := New(c, Options{}) time.Sleep(90 * time.Second) + republish(t, c, "live") nodes, err := src.ListNodes(ctx) require.NoError(t, err) - assert.Equal(t, []string{"fresh", "unstamped"}, nodes) - _, err = src.NodeInventory(ctx, "fresh") - require.NoError(t, err) - addr, pools, err := src.NodeCapacity(ctx, "fresh") - require.NoError(t, err) - assert.Equal(t, "fresh:7777", addr) - assert.Len(t, pools, 1) + assert.Equal(t, []string{"dead", "live", "unstamped"}, nodes) - time.Sleep(500 * time.Millisecond) + time.Sleep(time.Second) + republish(t, c, "live") nodes, err = src.ListNodes(ctx) require.NoError(t, err) - assert.Equal(t, []string{"unstamped"}, nodes) - _, err = src.NodeInventory(ctx, "fresh") + assert.Equal(t, []string{"live", "unstamped"}, nodes) + _, err = src.NodeInventory(ctx, "dead") assert.True(t, k8serrors.IsNotFound(err), "NodeInventory on a stale node: %v", err) - _, _, err = src.NodeCapacity(ctx, "fresh") + _, _, err = src.NodeCapacity(ctx, "dead") assert.True(t, k8serrors.IsNotFound(err), "NodeCapacity on a stale node: %v", err) + addr, pools, err := src.NodeCapacity(ctx, "live") + require.NoError(t, err) + assert.Equal(t, "live:7777", addr) + assert.Len(t, pools, 1) + }) +} + +func TestSourceKeepsEveryNodeThroughAPublishOutage(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ctx := t.Context() + src := New(fakeClient(t, inventory("a", metav1.Now()), inventory("b", metav1.Now())), Options{}) + time.Sleep(10 * time.Minute) + nodes, err := src.ListNodes(ctx) + require.NoError(t, err) + assert.Equal(t, []string{"a", "b"}, nodes) + _, err = src.NodeInventory(ctx, "a") + require.NoError(t, err, "a frozen fleet still answers lookups") + }) +} +func TestSourceCapsTheReferenceAtTheClock(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ahead := metav1.NewTime(time.Now().Add(time.Hour)) + src := New(fakeClient(t, inventory("rogue", ahead), inventory("a", metav1.Now()), inventory("b", metav1.Now())), Options{}) + nodes, err := src.ListNodes(t.Context()) + require.NoError(t, err) + assert.Equal(t, []string{"a", "b", "rogue"}, nodes, "a stamp ahead of the clock must not evict the others") + }) +} + +func TestSourceLookupBeforeAnyListCountsFresh(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + src := New(fakeClient(t, inventory("a", metav1.Now())), Options{}) time.Sleep(time.Hour) - _, _, err = src.NodeCapacity(ctx, "unstamped") - require.NoError(t, err, "an inventory without publishedAt never goes stale") + _, _, err := src.NodeCapacity(t.Context(), "a") + require.NoError(t, err) }) } func TestSourceHonorsStaleAfter(t *testing.T) { synctest.Test(t, func(t *testing.T) { - src := New(fakeReader(t, inventory("n1", metav1.Now())), Options{StaleAfter: 10 * time.Second}) + c := fakeClient(t, inventory("live", metav1.Now()), inventory("dead", metav1.Now())) + src := New(c, Options{StaleAfter: 10 * time.Second}) time.Sleep(11 * time.Second) + republish(t, c, "live") nodes, err := src.ListNodes(t.Context()) require.NoError(t, err) - assert.Empty(t, nodes) + assert.Equal(t, []string{"live"}, nodes) }) } -func fakeReader(t *testing.T, objs ...client.Object) client.Reader { +func TestPublishedAtTreatsAMissingOrUnparsableStampAsUnset(t *testing.T) { + for _, obj := range []map[string]any{{}, {"publishedAt": nil}, {"publishedAt": "not a time"}} { + assert.Zero(t, publishedAt(obj), "%v", obj) + } + assert.Equal(t, time.Date(2026, 9, 27, 9, 0, 0, 0, time.UTC).UnixNano(), publishedAt(map[string]any{"publishedAt": "2026-09-27T09:00:00Z"})) +} + +func fakeClient(t *testing.T, objs ...client.Object) client.Client { t.Helper() scheme := runtime.NewScheme() require.NoError(t, cocoonv1beta1.AddToScheme(scheme)) return fake.NewClientBuilder().WithScheme(scheme).WithObjects(objs...).Build() } +func republish(t *testing.T, c client.Client, node string) { + t.Helper() + inv := &cocoonv1beta1.NodeInventory{} + require.NoError(t, c.Get(t.Context(), client.ObjectKey{Name: node}, inv)) + inv.PublishedAt = metav1.Now() + require.NoError(t, c.Update(t.Context(), inv)) +} + func inventory(node string, publishedAt metav1.Time) *cocoonv1beta1.NodeInventory { return &cocoonv1beta1.NodeInventory{ Name: node, diff --git a/pkg/scale/sandboxstore_impl.go b/pkg/scale/sandboxstore_impl.go index cae2f67..4a0d1b6 100644 --- a/pkg/scale/sandboxstore_impl.go +++ b/pkg/scale/sandboxstore_impl.go @@ -539,8 +539,6 @@ type NodeLiveSource interface { // InventoryApplier server-side-applies a NodeInventory object. type InventoryApplier interface { Apply(ctx context.Context, inv *NodeInventory) error - // Delete removes node's inventory, and a missing one is success. - Delete(ctx context.Context, node string) error } var ( @@ -586,12 +584,12 @@ func (s *StaticInventorySource) Partition(node string) { s.partition[node] = struct{}{} } -func (s *StaticInventorySource) Delete(_ context.Context, node string) error { +// Remove drops a node's inventory object entirely (lost inventory). +func (s *StaticInventorySource) Remove(node string) { s.mu.Lock() defer s.mu.Unlock() delete(s.inv, node) delete(s.partition, node) - return nil } func (s *StaticInventorySource) ListNodes(_ context.Context) ([]string, error) { diff --git a/pkg/scale/sandboxstore_impl_test.go b/pkg/scale/sandboxstore_impl_test.go index f9cdeb0..485d94f 100644 --- a/pkg/scale/sandboxstore_impl_test.go +++ b/pkg/scale/sandboxstore_impl_test.go @@ -286,7 +286,7 @@ func TestPublisher_RebuildsFromLiveAfterLoss(t *testing.T) { require.NoError(t, err) require.Equal(t, 1, src.ObjectCount()) - require.NoError(t, src.Delete(ctx, "n1")) + src.Remove("n1") require.Equal(t, 0, src.ObjectCount()) live.entries = []InventoryEntry{entry("ns/a", "Running"), entry("ns/b", "Running")} From 6b9cb2e8a1d1e2ec3888c0aada593371d671f582 Mon Sep 17 00:00:00 2001 From: CMGS Date: Sun, 27 Sep 2026 18:23:44 +0800 Subject: [PATCH 4/4] docs: a dead node leaves within one publish interval after the window, and a control-plane outage stops only the driver's lease holder Measured on the kube kit: the reference advances only when a live node publishes, so a dead node left 120.8 s after its last stamp; with kube-apiserver down 150 s the replica holding the warm-pool lease exited on the lost lease (pre-existing), while the second replica and envd-proxy kept every node listed and a claim landed. --- docs/configuration.md | 4 +++- docs/scaling-design.md | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index 22bd893..b3a5888 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -143,7 +143,9 @@ pools. An empty sandboxd token leaves it fail-closed: it logs and sets no pools. A node that stops publishing, because it died or its vk-sandbox stopped, leaves every read and claim path once its `publishedAt` trails the newest publish by more than this. At vk-sandbox's default 30 s publish cadence that is -three missed publishes. Keep it above vk-sandbox's `--publish-interval`. +three missed publishes, and since the newest publish only advances when a live +node publishes, a dead node leaves between this age and one publish interval +later. Keep it above vk-sandbox's `--publish-interval`. The age is measured against the newest publish the process has seen, capped by its own clock. A control-plane write outage freezes every stamp together, so it diff --git a/docs/scaling-design.md b/docs/scaling-design.md index b8f98b9..9a86345 100644 --- a/docs/scaling-design.md +++ b/docs/scaling-design.md @@ -162,7 +162,7 @@ v1beta1 the `APIService` hands to the aggregated server, which serves only | A node `inventory` object lost | Rebuilt from the node's own live state on next publish | No | | A node's `Node` object deleted | Its `NodeInventory` is garbage-collected with it: the node leaves the claim path and its sandboxes leave the read view, though sandboxd keeps serving them and their leases still expire there; restarting vk-sandbox registers the node again | No | | A node dies, its `Node` object kept | Its `NodeInventory` is no longer republished. Once its `publishedAt` trails the newest publish in the fleet by more than `--inventory-stale-after` (90 s), the node leaves the claim path and its sandboxes leave the read view, the watch emits Deleted for each, Get and the verbs answer 404, and the warm-pool driver and envd-proxy stop calling it. Its next publish brings it back. A vk-sandbox restart shorter than the window changes nothing | No | -| The control plane stops accepting writes | Every publish fails and every stamp freezes together, so no node goes stale: claims, lists and lookups keep running from the informer cache, which shows the inventory as of the outage | No | +| The control plane stops accepting writes | Every publish fails and every stamp freezes together, so no node goes stale. envd-proxy and every apiserver replica but one keep claims, lists and lookups running from the informer cache, which shows the inventory as of the outage. The replica holding the warm-pool driver's lease exits once it cannot renew it (about 10 s) and starts again when the control plane returns | No | | Aggregated server restart | Stateless; rebuilds from node fan-out | No | | Client reads before the owning node republishes inventory | Lookups by name and by claim id (e2b, envd proxy), and a list or watch pinned to one name in a namespace, ask the nodes; other lists and watches show it at the next publish | No — fleet `list` and `watch`, and a deleted sandbox until its node publishes, stay eventually consistent |