From 2c958a39b801b0ee446d551ccf8ef983fb022a62 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Thu, 24 Sep 2026 18:03:22 +0000 Subject: [PATCH 1/6] MIR-858: Surface sandbox health in app status and doctor --- .../entityserver_v1alpha/rpc.gen.go | 122 +++++++++++ api/entityserver/rpc.yml | 5 + cli/commands/app_status.go | 21 +- cli/commands/app_status_health.go | 206 ++++++++++++++++++ cli/commands/app_status_health_test.go | 159 ++++++++++++++ cli/commands/doctor_check.go | 15 ++ cli/commands/doctor_resources.go | 191 ++++++++++++++++ servers/entityserver/entityserver.go | 13 ++ servers/entityserver/entityserver_test.go | 25 +++ 9 files changed, 755 insertions(+), 2 deletions(-) create mode 100644 cli/commands/app_status_health.go create mode 100644 cli/commands/app_status_health_test.go create mode 100644 cli/commands/doctor_resources.go diff --git a/api/entityserver/entityserver_v1alpha/rpc.gen.go b/api/entityserver/entityserver_v1alpha/rpc.gen.go index 70b4768d8..9e96f4d25 100644 --- a/api/entityserver/entityserver_v1alpha/rpc.gen.go +++ b/api/entityserver/entityserver_v1alpha/rpc.gen.go @@ -2190,6 +2190,58 @@ func (v *EntityAccessReindexResults) UnmarshalJSON(data []byte) error { return json.Unmarshal(data, &v.data) } +type entityAccessCheckIndexHealthArgsData struct{} + +type EntityAccessCheckIndexHealthArgs struct { + call rpc.Call + data entityAccessCheckIndexHealthArgsData +} + +func (v *EntityAccessCheckIndexHealthArgs) MarshalCBOR() ([]byte, error) { + return cbor.Marshal(v.data) +} + +func (v *EntityAccessCheckIndexHealthArgs) UnmarshalCBOR(data []byte) error { + return cbor.Unmarshal(data, &v.data) +} + +func (v *EntityAccessCheckIndexHealthArgs) MarshalJSON() ([]byte, error) { + return json.Marshal(v.data) +} + +func (v *EntityAccessCheckIndexHealthArgs) UnmarshalJSON(data []byte) error { + return json.Unmarshal(data, &v.data) +} + +type entityAccessCheckIndexHealthResultsData struct { + OrphanedEntries *int64 `cbor:"0,keyasint,omitempty" json:"orphaned_entries,omitempty"` +} + +type EntityAccessCheckIndexHealthResults struct { + call rpc.Call + data entityAccessCheckIndexHealthResultsData +} + +func (v *EntityAccessCheckIndexHealthResults) SetOrphanedEntries(orphaned_entries int64) { + v.data.OrphanedEntries = &orphaned_entries +} + +func (v *EntityAccessCheckIndexHealthResults) MarshalCBOR() ([]byte, error) { + return cbor.Marshal(v.data) +} + +func (v *EntityAccessCheckIndexHealthResults) UnmarshalCBOR(data []byte) error { + return cbor.Unmarshal(data, &v.data) +} + +func (v *EntityAccessCheckIndexHealthResults) MarshalJSON() ([]byte, error) { + return json.Marshal(v.data) +} + +func (v *EntityAccessCheckIndexHealthResults) UnmarshalJSON(data []byte) error { + return json.Unmarshal(data, &v.data) +} + type entityAccessGetAttributesByTagArgsData struct { Tag *string `cbor:"0,keyasint,omitempty" json:"tag,omitempty"` } @@ -2802,6 +2854,32 @@ func (t *EntityAccessReindex) Results() *EntityAccessReindexResults { return results } +type EntityAccessCheckIndexHealth struct { + rpc.Call + args EntityAccessCheckIndexHealthArgs + results EntityAccessCheckIndexHealthResults +} + +func (t *EntityAccessCheckIndexHealth) Args() *EntityAccessCheckIndexHealthArgs { + args := &t.args + if args.call != nil { + return args + } + args.call = t.Call + t.Call.Args(args) + return args +} + +func (t *EntityAccessCheckIndexHealth) Results() *EntityAccessCheckIndexHealthResults { + results := &t.results + if results.call != nil { + return results + } + results.call = t.Call + t.Call.Results(results) + return results +} + type EntityAccessGetAttributesByTag struct { rpc.Call args EntityAccessGetAttributesByTagArgs @@ -2850,6 +2928,7 @@ type EntityAccess interface { RevokeSession(ctx context.Context, state *EntityAccessRevokeSession) error PingSession(ctx context.Context, state *EntityAccessPingSession) error Reindex(ctx context.Context, state *EntityAccessReindex) error + CheckIndexHealth(ctx context.Context, state *EntityAccessCheckIndexHealth) error GetAttributesByTag(ctx context.Context, state *EntityAccessGetAttributesByTag) error } @@ -2941,6 +3020,10 @@ func (reexportEntityAccess) Reindex(ctx context.Context, state *EntityAccessRein panic("not implemented") } +func (reexportEntityAccess) CheckIndexHealth(ctx context.Context, state *EntityAccessCheckIndexHealth) error { + panic("not implemented") +} + func (reexportEntityAccess) GetAttributesByTag(ctx context.Context, state *EntityAccessGetAttributesByTag) error { panic("not implemented") } @@ -3161,6 +3244,16 @@ func AdaptEntityAccess(t EntityAccess) *rpc.Interface { return t.Reindex(ctx, &EntityAccessReindex{Call: call}) }, }, + { + Name: "check_index_health", + InterfaceName: "EntityAccess", + Index: 0, + Public: false, + Params: []string{}, + Handler: func(ctx context.Context, call rpc.Call) error { + return t.CheckIndexHealth(ctx, &EntityAccessCheckIndexHealth{Call: call}) + }, + }, { Name: "get_attributes_by_tag", InterfaceName: "EntityAccess", @@ -3933,6 +4026,35 @@ func (v EntityAccessClient) Reindex(ctx context.Context, dry_run bool) (*EntityA return &EntityAccessClientReindexResults{client: v.Client, data: ret}, nil } +type EntityAccessClientCheckIndexHealthResults struct { + client rpc.Client + data entityAccessCheckIndexHealthResultsData +} + +func (v *EntityAccessClientCheckIndexHealthResults) HasOrphanedEntries() bool { + return v.data.OrphanedEntries != nil +} + +func (v *EntityAccessClientCheckIndexHealthResults) OrphanedEntries() int64 { + if v.data.OrphanedEntries == nil { + return 0 + } + return *v.data.OrphanedEntries +} + +func (v EntityAccessClient) CheckIndexHealth(ctx context.Context) (*EntityAccessClientCheckIndexHealthResults, error) { + args := EntityAccessCheckIndexHealthArgs{} + + var ret entityAccessCheckIndexHealthResultsData + + err := v.Call(ctx, "check_index_health", &args, &ret) + if err != nil { + return nil, err + } + + return &EntityAccessClientCheckIndexHealthResults{client: v.Client, data: ret}, nil +} + type EntityAccessClientGetAttributesByTagResults struct { client rpc.Client data entityAccessGetAttributesByTagResultsData diff --git a/api/entityserver/rpc.yml b/api/entityserver/rpc.yml index a5f0a3f43..3dada1db4 100644 --- a/api/entityserver/rpc.yml +++ b/api/entityserver/rpc.yml @@ -373,6 +373,11 @@ interfaces: type: list element: ReindexStat + - name: check_index_health + results: + - name: orphaned_entries + type: int64 + - name: get_attributes_by_tag parameters: - name: tag diff --git a/cli/commands/app_status.go b/cli/commands/app_status.go index ec2cecd51..cb73b4aa4 100644 --- a/cli/commands/app_status.go +++ b/cli/commands/app_status.go @@ -62,10 +62,11 @@ func AppStatus(ctx *Context, opts struct { if err != nil { recentResult = nil } + health, healthErr := fetchServiceHealth(ctx, opts.App) // JSON output if opts.IsJSON() { - return printAppStatusJSON(appResult, appConfig, activeDeployment, recentResult, opts.App, clusterId) + return printAppStatusJSON(appResult, appConfig, activeDeployment, recentResult, opts.App, clusterId, health, healthErr) } // Define styles @@ -231,6 +232,16 @@ func AppStatus(ctx *Context, opts struct { ctx.Printf("\n%s\n", yellowStyle.Render("No active deployment found")) } + ctx.Printf("\n%s\n", labelStyle.Render("Service Sandboxes:")) + switch { + case healthErr != nil: + ctx.Printf(" Health unavailable: %v\n", healthErr) + case len(health) == 0: + ctx.Printf(" No service pools found\n") + default: + ctx.Printf("%s", renderServiceHealth(health)) + } + // Recent deployments summary ctx.Printf("\n%s\n", labelStyle.Render("Recent Activity:")) @@ -288,7 +299,7 @@ func printAppStatusJSON( appConfig *app_v1alpha.Configuration, activeDeployment *deployment_v1alpha.DeploymentClientGetActiveDeploymentResults, recentResult *deployment_v1alpha.DeploymentClientListDeploymentsResults, - app, cluster string, + app, cluster string, health []serviceHealth, healthErr error, ) error { type gitInfo struct { Sha string `json:"sha,omitempty"` @@ -396,11 +407,17 @@ func printAppStatusJSON( Configuration *configuration `json:"configuration,omitempty"` ActiveDeployment *deploymentJSON `json:"active_deployment,omitempty"` RecentDeployments []deploymentJSON `json:"recent_deployments,omitempty"` + Services []serviceHealth `json:"services,omitempty"` + HealthError string `json:"health_error,omitempty"` }{ App: app, Cluster: cluster, WorkloadRole: appResult.WorkloadRole(), MaintenanceRoutes: appResult.MaintenanceRoutes(), + Services: health, + } + if healthErr != nil { + output.HealthError = healthErr.Error() } if appResult.HasVersionId() && appResult.VersionId() != "" { diff --git a/cli/commands/app_status_health.go b/cli/commands/app_status_health.go new file mode 100644 index 000000000..3cd6ef59f --- /dev/null +++ b/cli/commands/app_status_health.go @@ -0,0 +1,206 @@ +package commands + +import ( + "fmt" + "slices" + "sort" + "strings" + "time" + + "miren.dev/runtime/api/app/app_v1alpha" + "miren.dev/runtime/api/compute/compute_v1alpha" + "miren.dev/runtime/api/core/core_v1alpha" + "miren.dev/runtime/api/entityserver" + "miren.dev/runtime/api/entityserver/entityserver_v1alpha" + "miren.dev/runtime/pkg/entity" + "miren.dev/runtime/pkg/rpc/standard" +) + +const failureWindow = 10 * time.Minute + +type serviceHealth struct { + Service string `json:"service"` + Running int `json:"running"` + Dead int `json:"dead"` + CrashStreak int `json:"crash_streak"` + CrashLooping bool `json:"crash_looping"` + LastExitCode *int64 `json:"last_exit_code,omitempty"` + LastFailureLog string `json:"last_failure_log,omitempty"` + + lastSandbox string + lastExit time.Time + lastFailed bool +} + +type sandboxHealthRecord struct { + Sandbox compute_v1alpha.Sandbox + Pool string + Updated time.Time +} + +func summarizeServiceHealth(pools []compute_v1alpha.SandboxPool, sandboxes []sandboxHealthRecord, now time.Time) []serviceHealth { + byPool := make(map[string]*serviceHealth) + byService := make(map[string]*serviceHealth) + for _, pool := range pools { + h := byService[pool.Service] + if h == nil { + h = &serviceHealth{Service: pool.Service} + byService[pool.Service] = h + } + byPool[pool.ID.String()] = h + if pool.DesiredInstances > pool.ReadyInstances && !pool.LastCrashTime.After(now) && pool.LastCrashTime.After(now.Add(-failureWindow)) { + h.CrashStreak += int(pool.ConsecutiveCrashCount) + if pool.CooldownUntil.After(now) && pool.ConsecutiveCrashCount >= 2 { + h.CrashLooping = true + } + } + } + for _, record := range sandboxes { + sb := record.Sandbox + h := byPool[record.Pool] + if h == nil { + continue + } + switch sb.Status { + case compute_v1alpha.RUNNING: + h.Running++ + case compute_v1alpha.DEAD: + h.Dead++ + when := sb.Exit.At + if when.IsZero() { + when = record.Updated + } + failed := !sb.Exit.At.IsZero() && sb.Exit.Code != 0 + if failed && !h.lastFailed || failed == h.lastFailed && when.After(h.lastExit) { + h.lastFailed = failed + if !sb.Exit.At.IsZero() { + code := sb.Exit.Code + h.LastExitCode = &code + } else { + h.LastExitCode = nil + } + h.lastExit = when + h.lastSandbox = sb.ID.String() + } + case compute_v1alpha.PENDING, compute_v1alpha.NOT_READY, compute_v1alpha.STOPPED: + } + } + result := make([]serviceHealth, 0, len(byService)) + for _, h := range byService { + result = append(result, *h) + } + sort.Slice(result, func(i, j int) bool { return result[i].Service < result[j].Service }) + return result +} + +func renderServiceHealth(health []serviceHealth) string { + var b strings.Builder + for _, svc := range health { + state := fmt.Sprintf("%d running, %d dead; %d crashes in current streak (last crash in 10 minutes)", svc.Running, svc.Dead, svc.CrashStreak) + if svc.CrashLooping { + state = infoRed.Render("crash-looping") + ", " + state + } + fmt.Fprintf(&b, " %s: %s\n", svc.Service, state) + if svc.LastExitCode != nil { + fmt.Fprintf(&b, " Last exit code: %d\n", *svc.LastExitCode) + } + if svc.LastFailureLog != "" { + b.WriteString(" Last failure logs:\n") + for _, line := range strings.Split(svc.LastFailureLog, "\n") { + fmt.Fprintf(&b, " %s\n", line) + } + } + } + return b.String() +} + +func activeServicePools(pools []compute_v1alpha.SandboxPool, version entity.Id) []compute_v1alpha.SandboxPool { + var active []compute_v1alpha.SandboxPool + for _, pool := range pools { + if version == "" || slices.Contains(pool.ReferencedByVersions, version) || len(pool.ReferencedByVersions) == 0 && pool.SandboxSpec.Version == version { + active = append(active, pool) + } + } + return active +} + +func fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { + cl, err := ctx.RPCClient("entities") + if err != nil { + return nil, err + } + defer cl.Close() + eac := entityserver_v1alpha.NewEntityAccessClient(cl) + var appEntity core_v1alpha.App + if err := entityserver.NewClient(ctx.Log, eac).Get(ctx, app, &appEntity); err != nil { + return nil, fmt.Errorf("get app %q: %w", app, err) + } + poolRes, err := eac.List(ctx, entity.Ref(compute_v1alpha.SandboxPoolAppId, appEntity.ID)) + if err != nil { + return nil, err + } + var pools []compute_v1alpha.SandboxPool + for _, entry := range poolRes.Values() { + var pool compute_v1alpha.SandboxPool + pool.Decode(entry.Entity()) + pools = append(pools, pool) + } + pools = activeServicePools(pools, appEntity.ActiveVersion) + if len(pools) == 0 { + return []serviceHealth{}, nil + } + kind, err := eac.LookupKind(ctx, "sandbox") + if err != nil { + return nil, err + } + sbRes, err := eac.List(ctx, kind.Attr()) + if err != nil { + return nil, err + } + var sandboxes []sandboxHealthRecord + for _, entry := range sbRes.Values() { + var meta core_v1alpha.Metadata + meta.Decode(entry.Entity()) + pool, _ := meta.Labels.Get("pool") + var sb compute_v1alpha.Sandbox + sb.Decode(entry.Entity()) + sandboxes = append(sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) + } + health := summarizeServiceHealth(pools, sandboxes, time.Now()) + for i := range health { + if health[i].lastSandbox != "" && health[i].lastExit.After(time.Now().Add(-failureWindow)) && health[i].CrashStreak > 0 { + health[i].LastFailureLog = recentSandboxFailureLog(ctx, health[i].lastSandbox, health[i].lastExit) + } + } + return health, nil +} + +func recentSandboxFailureLog(ctx *Context, sandbox string, exit time.Time) string { + cl, err := ctx.RPCClient(rpcLogs) + if err != nil { + return "" + } + defer cl.Close() + logs, err := app_v1alpha.NewLogsClient(cl).SandboxLogs(ctx, sandbox, standard.ToTimestamp(exit.Add(-time.Minute)), false) + if err != nil { + return "" + } + var lines []string + for _, log := range logs.Logs() { + if log.HasTimestamp() && standard.FromTimestamp(log.Timestamp()).After(exit.Add(time.Second)) { + continue + } + if line := strings.TrimSpace(log.Line()); line != "" { + lines = append(lines, line) + } + } + if len(lines) > 3 { + lines = lines[len(lines)-3:] + } + for i, line := range lines { + if len(line) > 200 { + lines[i] = line[:200] + "..." + } + } + return strings.Join(lines, "\n") +} diff --git a/cli/commands/app_status_health_test.go b/cli/commands/app_status_health_test.go new file mode 100644 index 000000000..5b25dd2f0 --- /dev/null +++ b/cli/commands/app_status_health_test.go @@ -0,0 +1,159 @@ +package commands + +import ( + "context" + "encoding/json" + "strings" + "testing" + "time" + + "miren.dev/runtime/api/compute/compute_v1alpha" + "miren.dev/runtime/api/storage/storage_v1alpha" + "miren.dev/runtime/pkg/entity" +) + +func TestSummarizeServiceHealth(t *testing.T) { + now := time.Date(2026, 9, 23, 12, 0, 0, 0, time.UTC) + pools := []compute_v1alpha.SandboxPool{ + {ID: "pool-db", Service: "db", DesiredInstances: 1, ConsecutiveCrashCount: 2, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}, + {ID: "pool-web", Service: "web"}, + } + sandboxes := []sandboxHealthRecord{ + {Pool: "pool-web", Sandbox: compute_v1alpha.Sandbox{ID: "web-1", Status: compute_v1alpha.RUNNING}}, + {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-1", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-9 * time.Minute)}}}, + {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-2", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-2 * time.Minute)}}}, + {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-3", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 5, At: now.Add(-11 * time.Minute)}}}, + {Pool: "pool-other", Sandbox: compute_v1alpha.Sandbox{ID: "foreign", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 99, At: now}}}, + } + got := summarizeServiceHealth(pools, sandboxes, now) + if len(got) != 2 || got[0].Service != "db" || got[0].Dead != 3 || got[0].Running != 0 || got[0].CrashStreak != 2 || !got[0].CrashLooping || got[0].LastExitCode == nil || *got[0].LastExitCode != 17 || got[0].lastSandbox != "db-2" { + t.Fatalf("db summary = %+v", got) + } + if got[1].Service != "web" || got[1].Running != 1 || got[1].CrashLooping { + t.Fatalf("web summary = %+v", got[1]) + } + got[0].LastFailureLog = "database unavailable\nconnection refused" + output := renderServiceHealth(got) + for _, fragment := range []string{"db: crash-looping, 0 running, 3 dead; 2 crashes in current streak", "Last exit code: 17", " connection refused", "web: 1 running, 0 dead"} { + if !strings.Contains(output, fragment) { + t.Errorf("status output %q missing %q", output, fragment) + } + } + encoded, err := json.Marshal(got[0]) + if err != nil || !strings.Contains(string(encoded), `"last_exit_code":17`) || !strings.Contains(string(encoded), `"crash_streak":2`) { + t.Fatalf("JSON status = %s, err = %v", encoded, err) + } + zero := summarizeServiceHealth(pools, []sandboxHealthRecord{{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-zero", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-time.Minute)}}}}, now)[0] + if zero.LastExitCode == nil || *zero.LastExitCode != 0 { + t.Fatalf("zero exit code must be present: %+v", zero) + } + early := summarizeServiceHealth(pools, []sandboxHealthRecord{{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-early", Status: compute_v1alpha.DEAD}, Updated: now.Add(-time.Minute)}}, now)[0] + if early.CrashStreak != 2 || early.LastExitCode != nil || early.lastSandbox != "db-early" { + t.Fatalf("early failure without an exit record = %+v", early) + } + if idle := summarizeServiceHealth([]compute_v1alpha.SandboxPool{{ID: "pool-db", Service: "db", DesiredInstances: 0, ConsecutiveCrashCount: 4, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}}, sandboxes, now)[0]; idle.CrashStreak != 0 || idle.CrashLooping { + t.Fatalf("scaled-to-zero service should not appear to be failing: %+v", idle) + } +} + +func TestScaleDownIsNotACrashLoop(t *testing.T) { + now := time.Date(2026, 9, 23, 12, 0, 0, 0, time.UTC) + retired := make([]sandboxHealthRecord, 3) + for i := range retired { + retired[i] = sandboxHealthRecord{Pool: "pool-web", Sandbox: compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD}, Updated: now.Add(-time.Minute)} + } + for _, desired := range []int64{1, 0, 1} { // 4→1, scaled to zero, then woken + pool := compute_v1alpha.SandboxPool{ID: "pool-web", Service: "web", DesiredInstances: desired} + got := summarizeServiceHealth([]compute_v1alpha.SandboxPool{pool}, retired, now)[0] + if got.Dead != 3 || got.CrashStreak != 0 || got.CrashLooping { + t.Fatalf("desired %d: %+v", desired, got) + } + } + env := verdictEnv(remoteHost, nil, probeOpen, probeOpen) + env.resources = &doctorResources{pools: []compute_v1alpha.SandboxPool{{ID: "pool-web", Service: "web", DesiredInstances: 1}}, sandboxes: retired} + if got := checkSandboxes(env); got.Status != checkOK { + t.Fatalf("doctor treated scale-down as a failure: %+v", got) + } + // A later intentional retirement must not replace the last known failed exit. + pool := compute_v1alpha.SandboxPool{ID: "pool-web", Service: "web", DesiredInstances: 1, ConsecutiveCrashCount: 1, LastCrashTime: now.Add(-2 * time.Minute)} + retired = append(retired, sandboxHealthRecord{Pool: "pool-web", Sandbox: compute_v1alpha.Sandbox{ID: "failed", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-2 * time.Minute)}}}) + got := summarizeServiceHealth([]compute_v1alpha.SandboxPool{pool}, retired, now)[0] + if got.LastExitCode == nil || *got.LastExitCode != 17 || got.lastSandbox != "failed" { + t.Fatalf("known failure lost to retirement: %+v", got) + } +} + +func TestActiveServicePoolsUsesVersionMembership(t *testing.T) { + pools := []compute_v1alpha.SandboxPool{ + {ID: "current", SandboxSpec: compute_v1alpha.SandboxSpec{Version: "app-version-456"}, ReferencedByVersions: []entity.Id{"app-version-456", "app-version-123"}}, + {ID: "old", SandboxSpec: compute_v1alpha.SandboxSpec{Version: "app-version-123"}, ReferencedByVersions: []entity.Id{"app-version-456"}}, + {ID: "legacy", SandboxSpec: compute_v1alpha.SandboxSpec{Version: "app-version-789"}}, + } + got := activeServicePools(pools, "app-version-123") + if len(got) != 1 || got[0].ID != "current" { + t.Fatalf("pools for active version = %+v", got) + } + if got := activeServicePools(pools, "app-version-789"); len(got) != 1 || got[0].ID != "legacy" { + t.Fatalf("legacy pool membership = %+v", got) + } + if got := activeServicePools(pools, ""); len(got) != 3 { + t.Fatalf("pools before first deployment = %+v", got) + } +} + +func TestDoctorResourceChecks(t *testing.T) { + now := time.Now() + env := verdictEnv(remoteHost, nil, probeOpen, probeOpen) + env.resources = &doctorResources{ + pools: []compute_v1alpha.SandboxPool{{ID: "pool-db", Service: "db", DesiredInstances: 1, ReadyInstances: 0, ConsecutiveCrashCount: 3, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}}, + disks: []storage_v1alpha.Disk{{ID: entity.Id("disk-1"), Status: storage_v1alpha.ERROR}}, + volumes: []storage_v1alpha.DiskVolume{{ID: entity.Id("volume-1"), ActualState: storage_v1alpha.DV_ERROR, ErrorMessage: "attach failed"}}, + } + env.orphans = 2 + for i := 0; i < 3; i++ { + env.resources.sandboxes = append(env.resources.sandboxes, sandboxHealthRecord{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{At: now.Add(-time.Duration(i+1) * time.Minute)}}}) + } + for _, tc := range []struct { + name string + check func(*doctorEnv) checkResult + want checkStatus + fragment string + }{ + {"indexes", checkEntityIndexes, checkWarn, "2 orphaned"}, + {"sandboxes", checkSandboxes, checkFail, "3 consecutive crashes"}, + {"pools", checkPools, checkFail, "db"}, + {"disks", checkDisks, checkFail, "attach failed"}, + } { + t.Run(tc.name, func(t *testing.T) { + got := tc.check(env) + if got.Status != tc.want || !strings.Contains(got.Summary, tc.fragment) || got.Problem == nil || len(got.Problem.Actions) == 0 { + t.Fatalf("check = %+v, want %s with action and %q", got, tc.want, tc.fragment) + } + }) + } + env.indexErr = context.DeadlineExceeded + if got := checkEntityIndexes(env); got.Status != checkWarn { + t.Fatalf("timed-out index check = %+v, want warning", got) + } + resources := env.resources + env.resources = nil + env.resourcesErr = context.DeadlineExceeded + env.indexErr = nil + if got := checkEntityIndexes(env); got.Status != checkWarn { + t.Fatalf("independent index check = %+v, want warning for two orphans", got) + } + env.resourcesErr = nil + env.resources = resources + if got := checkPools(env); got.Status != checkFail { + t.Fatalf("pool check = %+v, want independent failure", got) + } + // Two apps with a service called db must not combine into a false alert. + env.resources.pools[0].App = "app-a" + env.resources.pools[0].ConsecutiveCrashCount = 2 + env.resources.pools = append(env.resources.pools, compute_v1alpha.SandboxPool{ID: "pool-other", App: "app-b", Service: "db"}) + env.resources.sandboxes = env.resources.sandboxes[:2] + env.resources.sandboxes = append(env.resources.sandboxes, sandboxHealthRecord{Pool: "pool-other", Sandbox: compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{At: now.Add(-time.Minute)}}}) + if got := checkSandboxes(env); got.Status != checkOK { + t.Fatalf("distinct app failures = %+v, want no aggregate alert", got) + } +} diff --git a/cli/commands/doctor_check.go b/cli/commands/doctor_check.go index 9805dcf1a..af32779d8 100644 --- a/cli/commands/doctor_check.go +++ b/cli/commands/doctor_check.go @@ -115,6 +115,11 @@ type doctorEnv struct { // auth is only attempted when the cluster names an identity. auth authResult + + resources *doctorResources + resourcesErr error + orphans int64 + indexErr error } // local reports whether the active cluster runs on this machine, which decides @@ -189,6 +194,12 @@ func gatherCluster(ctx *Context, opts ConfigCentric, env *doctorEnv) { wg.Go(func() { env.serverVersion, env.serverVersionErr = fetchServerVersion(ctx) }) + wg.Go(func() { + env.resources, env.resourcesErr = gatherDoctorResources(ctx) + }) + wg.Go(func() { + env.orphans, env.indexErr = gatherDoctorIndex(ctx) + }) wg.Go(func() { env.tcp = probeTCP(env.cluster.Hostname) @@ -214,5 +225,9 @@ func doctorChecks() []check { {Name: "Server", Run: checkServer}, {Name: "Version", Run: checkVersion}, {Name: "Authentication", Run: checkAuthentication}, + {Name: "Entity indexes", Run: checkEntityIndexes}, + {Name: "Sandboxes", Run: checkSandboxes}, + {Name: "Pools", Run: checkPools}, + {Name: "Disks and volumes", Run: checkDisks}, } } diff --git a/cli/commands/doctor_resources.go b/cli/commands/doctor_resources.go new file mode 100644 index 000000000..7c1ea12a5 --- /dev/null +++ b/cli/commands/doctor_resources.go @@ -0,0 +1,191 @@ +package commands + +import ( + "context" + "fmt" + "sort" + "strings" + "time" + + "miren.dev/runtime/api/compute/compute_v1alpha" + "miren.dev/runtime/api/core/core_v1alpha" + "miren.dev/runtime/api/entityserver/entityserver_v1alpha" + "miren.dev/runtime/api/storage/storage_v1alpha" + "miren.dev/runtime/pkg/ui" +) + +type doctorResources struct { + pools []compute_v1alpha.SandboxPool + sandboxes []sandboxHealthRecord + disks []storage_v1alpha.Disk + volumes []storage_v1alpha.DiskVolume +} + +func gatherDoctorResources(ctx *Context) (*doctorResources, error) { + cl, err := ctx.RPCClient("entities") + if err != nil { + return nil, err + } + defer cl.Close() + eac := entityserver_v1alpha.NewEntityAccessClient(cl) + r := &doctorResources{} + for _, kindName := range []string{"sandbox_pool", "sandbox", "disk", "disk_volume"} { + kind, err := eac.LookupKind(ctx, kindName) + if err != nil { + return nil, fmt.Errorf("lookup %s: %w", kindName, err) + } + list, err := eac.List(ctx, kind.Attr()) + if err != nil { + return nil, fmt.Errorf("list %s: %w", kindName, err) + } + for _, entry := range list.Values() { + switch kindName { + case "sandbox_pool": + var pool compute_v1alpha.SandboxPool + pool.Decode(entry.Entity()) + r.pools = append(r.pools, pool) + case "sandbox": + var sb compute_v1alpha.Sandbox + sb.Decode(entry.Entity()) + var meta core_v1alpha.Metadata + meta.Decode(entry.Entity()) + pool, _ := meta.Labels.Get("pool") + r.sandboxes = append(r.sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) + case "disk": + var disk storage_v1alpha.Disk + disk.Decode(entry.Entity()) + r.disks = append(r.disks, disk) + case "disk_volume": + var volume storage_v1alpha.DiskVolume + volume.Decode(entry.Entity()) + r.volumes = append(r.volumes, volume) + } + } + } + return r, nil +} + +func gatherDoctorIndex(ctx *Context) (int64, error) { + probeCtx, cancel := context.WithTimeout(ctx, probeTimeout) + defer cancel() + cl, err := ctx.rpcClient(probeCtx, "entities") + if err != nil { + return 0, err + } + defer cl.Close() + result, err := entityserver_v1alpha.NewEntityAccessClient(cl).CheckIndexHealth(probeCtx) + if err != nil { + return 0, err + } + return result.OrphanedEntries(), nil +} + +func resourceUnavailable(env *doctorEnv) (checkResult, bool) { + if !env.configured() { + return checkResult{Status: checkSkip, Summary: "(no cluster configured)"}, true + } + if env.connErr != nil { + return checkResult{Status: checkSkip, Summary: "(server unreachable)"}, true + } + if env.resourcesErr != nil || env.resources == nil { + return checkResult{Status: checkSkip, Summary: fmt.Sprintf("(resource scan unavailable: %v)", env.resourcesErr)}, true + } + return checkResult{}, false +} + +func checkEntityIndexes(env *doctorEnv) checkResult { + if !env.configured() { + return checkResult{Status: checkSkip, Summary: "(no cluster configured)"} + } + if env.connErr != nil { + return checkResult{Status: checkSkip, Summary: "(server unreachable)"} + } + if env.indexErr != nil { + return checkResult{Status: checkWarn, Summary: "index scan unavailable", Problem: &ui.Diagnostic{ + Summary: "could not check for orphaned index entries", Detail: env.indexErr.Error(), + Actions: []ui.Action{{Command: "miren version", Note: "check that the server supports index diagnostics"}, {Command: "miren doctor", Note: "retry the bounded scan"}}, + }} + } + n := env.orphans + if n == 0 { + return checkResult{Status: checkOK, Summary: "no orphaned index entries"} + } + return checkResult{Status: checkWarn, Summary: fmt.Sprintf("%d orphaned index entries", n), Problem: &ui.Diagnostic{ + Summary: fmt.Sprintf("%d index entries point to missing entities", n), + Actions: []ui.Action{{Command: "miren debug reindex --dry-run", Note: "inspect index drift"}, {Command: "miren debug reindex", Note: "repair indexes after review"}}, + }} +} + +func checkSandboxes(env *doctorEnv) checkResult { + if skip, unavailable := resourceUnavailable(env); unavailable { + return skip + } + // The same service name may be used by several apps. Keep their counts separate. + pools := append([]compute_v1alpha.SandboxPool(nil), env.resources.pools...) + for i := range pools { + if pools[i].App != "" { + pools[i].Service = pools[i].App.String() + "/" + pools[i].Service + } + } + services := summarizeServiceHealth(pools, env.resources.sandboxes, time.Now()) + var failing []string + for _, svc := range services { + if svc.CrashStreak >= 3 { + failing = append(failing, fmt.Sprintf("%s (%d consecutive crashes, last in 10m)", svc.Service, svc.CrashStreak)) + } + } + if len(failing) == 0 { + return checkResult{Status: checkOK, Summary: "no elevated recent sandbox crash streaks"} + } + return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ + Summary: "services have repeated sandbox failures", + Detail: strings.Join(failing, ", "), + Actions: []ui.Action{{Command: "miren sandbox list --all", Note: "inspect dead sandboxes"}, {Command: "miren app status -a ", Note: "view exit codes and failure logs"}}, + }} +} + +func checkPools(env *doctorEnv) checkResult { + if skip, unavailable := resourceUnavailable(env); unavailable { + return skip + } + var failing []string + now := time.Now() + for _, pool := range env.resources.pools { + if pool.DesiredInstances > 0 && pool.ReadyInstances == 0 && pool.CooldownUntil.After(now) && pool.ConsecutiveCrashCount >= 2 { + failing = append(failing, fmt.Sprintf("%s/%s (%d crashes, retry in %s)", pool.App, pool.Service, pool.ConsecutiveCrashCount, formatDuration(pool.CooldownUntil.Sub(now)))) + } + } + if len(failing) == 0 { + return checkResult{Status: checkOK, Summary: "no pools in crash cooldown"} + } + sort.Strings(failing) + return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ + Summary: "service pools are crash-looping", Detail: strings.Join(failing, ", "), + Actions: []ui.Action{{Command: "miren sandbox pool list", Note: "inspect failing pools"}, {Command: "miren sandbox list --all", Note: "find failed sandboxes"}}, + }} +} + +func checkDisks(env *doctorEnv) checkResult { + if skip, unavailable := resourceUnavailable(env); unavailable { + return skip + } + var failing []string + for _, disk := range env.resources.disks { + if disk.Status == storage_v1alpha.ERROR { + failing = append(failing, fmt.Sprintf("disk %s", disk.ID)) + } + } + for _, volume := range env.resources.volumes { + if volume.ActualState == storage_v1alpha.DV_ERROR { + failing = append(failing, fmt.Sprintf("volume %s: %s", volume.ID, volume.ErrorMessage)) + } + } + if len(failing) == 0 { + return checkResult{Status: checkOK, Summary: "no disk or volume errors"} + } + sort.Strings(failing) + return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ + Summary: "disk or volume errors detected", Detail: strings.Join(failing, ", "), + Actions: []ui.Action{{Command: "miren debug disk list", Note: "inspect disk states"}}, + }} +} diff --git a/servers/entityserver/entityserver.go b/servers/entityserver/entityserver.go index 9d74a8f79..6c43ce206 100644 --- a/servers/entityserver/entityserver.go +++ b/servers/entityserver/entityserver.go @@ -1303,6 +1303,19 @@ func (e *EntityServer) Reindex(ctx context.Context, req *entityserver_v1alpha.En return nil } +func (e *EntityServer) CheckIndexHealth(ctx context.Context, req *entityserver_v1alpha.EntityAccessCheckIndexHealth) error { + store, ok := e.Store.(*entity.EtcdStore) + if !ok { + return fmt.Errorf("index health requires EtcdStore") + } + stats, err := store.CleanupStaleCollectionEntries(ctx, e.Log, entity.CleanupOptions{DryRun: true}) + if err != nil { + return fmt.Errorf("scan index health: %w", err) + } + req.Results().SetOrphanedEntries(stats.OrphanedEntriesFound) + return nil +} + func (e *EntityServer) GetAttributesByTag(ctx context.Context, req *entityserver_v1alpha.EntityAccessGetAttributesByTag) error { args := req.Args() tag := args.Tag() diff --git a/servers/entityserver/entityserver_test.go b/servers/entityserver/entityserver_test.go index 067cbb6ed..603121279 100644 --- a/servers/entityserver/entityserver_test.go +++ b/servers/entityserver/entityserver_test.go @@ -10,6 +10,7 @@ import ( "testing" "time" + "github.com/mr-tron/base58" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "go.etcd.io/etcd/api/v3/etcdserverpb" @@ -29,6 +30,30 @@ func setupTestEtcd(t *testing.T) (*clientv3.Client, string) { return etcdtest.TestEtcdClient(t) } +func TestCheckIndexHealthIsReadOnly(t *testing.T) { + client, prefix := setupTestEtcd(t) + store, err := entity.NewEtcdStore(t.Context(), slog.Default(), client, prefix) + require.NoError(t, err) + server, err := NewEntityServer(slog.Default(), store) + require.NoError(t, err) + id := entity.Id("missing/sandbox") + key := prefix + "/collections/test_orphan/" + base58.Encode([]byte(id)) + _, err = client.Put(t.Context(), key, id.String()) + require.NoError(t, err) + + eac := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} + result, err := eac.CheckIndexHealth(t.Context()) + require.NoError(t, err) + assert.EqualValues(t, 1, result.OrphanedEntries()) + resp, err := client.Get(t.Context(), key) + require.NoError(t, err) + assert.Len(t, resp.Kvs, 1, "diagnostic must not repair the index") + canceled, cancel := context.WithCancel(t.Context()) + cancel() + _, err = eac.CheckIndexHealth(canceled) + require.Error(t, err, "an interrupted scan must not report a clean result") +} + func TestEntityServer_Get(t *testing.T) { store := entity.NewMockStore() server := &EntityServer{ From f2cd185b7498327c3e0efae38b60557982d95e4d Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Thu, 24 Sep 2026 18:23:35 +0000 Subject: [PATCH 2/6] Address MIR-858 status and doctor review feedback --- cli/commands/app_status_health.go | 67 +++++++++++++---------- cli/commands/app_status_health_test.go | 21 +++++++- cli/commands/doctor_resources.go | 73 +++++++++++++++----------- 3 files changed, 102 insertions(+), 59 deletions(-) diff --git a/cli/commands/app_status_health.go b/cli/commands/app_status_health.go index 3cd6ef59f..30d14e216 100644 --- a/cli/commands/app_status_health.go +++ b/cli/commands/app_status_health.go @@ -29,7 +29,7 @@ type serviceHealth struct { lastSandbox string lastExit time.Time - lastFailed bool + lastFailure time.Time } type sandboxHealthRecord struct { @@ -66,20 +66,13 @@ func summarizeServiceHealth(pools []compute_v1alpha.SandboxPool, sandboxes []san h.Running++ case compute_v1alpha.DEAD: h.Dead++ - when := sb.Exit.At - if when.IsZero() { - when = record.Updated + if !sb.Exit.At.IsZero() && (h.LastExitCode == nil || sb.Exit.At.After(h.lastExit)) { + code := sb.Exit.Code + h.LastExitCode = &code + h.lastExit = sb.Exit.At } - failed := !sb.Exit.At.IsZero() && sb.Exit.Code != 0 - if failed && !h.lastFailed || failed == h.lastFailed && when.After(h.lastExit) { - h.lastFailed = failed - if !sb.Exit.At.IsZero() { - code := sb.Exit.Code - h.LastExitCode = &code - } else { - h.LastExitCode = nil - } - h.lastExit = when + if sb.Exit.Code != 0 && !sb.Exit.At.IsZero() && (h.lastSandbox == "" || sb.Exit.At.After(h.lastFailure)) { + h.lastFailure = sb.Exit.At h.lastSandbox = sb.ID.String() } case compute_v1alpha.PENDING, compute_v1alpha.NOT_READY, compute_v1alpha.STOPPED: @@ -96,7 +89,10 @@ func summarizeServiceHealth(pools []compute_v1alpha.SandboxPool, sandboxes []san func renderServiceHealth(health []serviceHealth) string { var b strings.Builder for _, svc := range health { - state := fmt.Sprintf("%d running, %d dead; %d crashes in current streak (last crash in 10 minutes)", svc.Running, svc.Dead, svc.CrashStreak) + state := fmt.Sprintf("%d running, %d dead", svc.Running, svc.Dead) + if svc.CrashStreak > 0 { + state += fmt.Sprintf("; %d crashes in current streak (last within 10m)", svc.CrashStreak) + } if svc.CrashLooping { state = infoRed.Render("crash-looping") + ", " + state } @@ -153,23 +149,38 @@ func fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { if err != nil { return nil, err } - sbRes, err := eac.List(ctx, kind.Attr()) - if err != nil { - return nil, err - } var sandboxes []sandboxHealthRecord - for _, entry := range sbRes.Values() { - var meta core_v1alpha.Metadata - meta.Decode(entry.Entity()) - pool, _ := meta.Labels.Get("pool") - var sb compute_v1alpha.Sandbox - sb.Decode(entry.Entity()) - sandboxes = append(sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) + selected := make(map[string]bool, len(pools)) + for _, pool := range pools { + selected[pool.ID.String()] = true + } + // Pool is a metadata label, not an indexed attribute. Page the kind index + // rather than materializing every sandbox in the cluster in one RPC. + for cursor := ""; ; { + page, err := eac.ListPage(ctx, kind.Attr(), cursor, 200) + if err != nil { + return nil, err + } + for _, entry := range page.Values() { + var meta core_v1alpha.Metadata + meta.Decode(entry.Entity()) + pool, _ := meta.Labels.Get("pool") + if !selected[pool] { + continue + } + var sb compute_v1alpha.Sandbox + sb.Decode(entry.Entity()) + sandboxes = append(sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) + } + cursor = page.Cursor() + if cursor == "" { + break + } } health := summarizeServiceHealth(pools, sandboxes, time.Now()) for i := range health { - if health[i].lastSandbox != "" && health[i].lastExit.After(time.Now().Add(-failureWindow)) && health[i].CrashStreak > 0 { - health[i].LastFailureLog = recentSandboxFailureLog(ctx, health[i].lastSandbox, health[i].lastExit) + if health[i].lastSandbox != "" && health[i].lastFailure.After(time.Now().Add(-failureWindow)) && health[i].CrashStreak > 0 { + health[i].LastFailureLog = recentSandboxFailureLog(ctx, health[i].lastSandbox, health[i].lastFailure) } } return health, nil diff --git a/cli/commands/app_status_health_test.go b/cli/commands/app_status_health_test.go index 5b25dd2f0..ffba4126a 100644 --- a/cli/commands/app_status_health_test.go +++ b/cli/commands/app_status_health_test.go @@ -39,6 +39,9 @@ func TestSummarizeServiceHealth(t *testing.T) { t.Errorf("status output %q missing %q", output, fragment) } } + if strings.Contains(output[strings.Index(output, "web:"):], "crash") { + t.Fatalf("healthy web service reports a crash: %q", output) + } encoded, err := json.Marshal(got[0]) if err != nil || !strings.Contains(string(encoded), `"last_exit_code":17`) || !strings.Contains(string(encoded), `"crash_streak":2`) { t.Fatalf("JSON status = %s, err = %v", encoded, err) @@ -48,9 +51,17 @@ func TestSummarizeServiceHealth(t *testing.T) { t.Fatalf("zero exit code must be present: %+v", zero) } early := summarizeServiceHealth(pools, []sandboxHealthRecord{{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-early", Status: compute_v1alpha.DEAD}, Updated: now.Add(-time.Minute)}}, now)[0] - if early.CrashStreak != 2 || early.LastExitCode != nil || early.lastSandbox != "db-early" { + if early.CrashStreak != 2 || early.LastExitCode != nil || early.lastSandbox != "" { t.Fatalf("early failure without an exit record = %+v", early) } + latest := summarizeServiceHealth(pools, []sandboxHealthRecord{ + {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "success", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-time.Minute)}}}, + {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "failure", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-2 * time.Minute)}}}, + {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "retired", Status: compute_v1alpha.DEAD}, Updated: now}, + }, now)[0] + if latest.LastExitCode == nil || *latest.LastExitCode != 0 || latest.lastSandbox != "failure" || !latest.lastFailure.Equal(now.Add(-2*time.Minute)) { + t.Fatalf("newer successful exit hid earlier failure or was replaced by retirement: %+v", latest) + } if idle := summarizeServiceHealth([]compute_v1alpha.SandboxPool{{ID: "pool-db", Service: "db", DesiredInstances: 0, ConsecutiveCrashCount: 4, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}}, sandboxes, now)[0]; idle.CrashStreak != 0 || idle.CrashLooping { t.Fatalf("scaled-to-zero service should not appear to be failing: %+v", idle) } @@ -131,6 +142,9 @@ func TestDoctorResourceChecks(t *testing.T) { } }) } + if got := checkPools(env); got.Problem == nil || got.Problem.Actions[0].Command != "miren sandbox-pool list" { + t.Fatalf("pool diagnostic suggests invalid command: %+v", got) + } env.indexErr = context.DeadlineExceeded if got := checkEntityIndexes(env); got.Status != checkWarn { t.Fatalf("timed-out index check = %+v, want warning", got) @@ -139,6 +153,11 @@ func TestDoctorResourceChecks(t *testing.T) { env.resources = nil env.resourcesErr = context.DeadlineExceeded env.indexErr = nil + for _, check := range []func(*doctorEnv) checkResult{checkSandboxes, checkPools, checkDisks} { + if got := check(env); got.Status != checkSkip { + t.Fatalf("incomplete resource scan reported health: %+v", got) + } + } if got := checkEntityIndexes(env); got.Status != checkWarn { t.Fatalf("independent index check = %+v, want warning for two orphans", got) } diff --git a/cli/commands/doctor_resources.go b/cli/commands/doctor_resources.go index 7c1ea12a5..7d564aaf4 100644 --- a/cli/commands/doctor_resources.go +++ b/cli/commands/doctor_resources.go @@ -22,7 +22,9 @@ type doctorResources struct { } func gatherDoctorResources(ctx *Context) (*doctorResources, error) { - cl, err := ctx.RPCClient("entities") + probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + cl, err := ctx.rpcClient(probeCtx, "entities") if err != nil { return nil, err } @@ -30,43 +32,54 @@ func gatherDoctorResources(ctx *Context) (*doctorResources, error) { eac := entityserver_v1alpha.NewEntityAccessClient(cl) r := &doctorResources{} for _, kindName := range []string{"sandbox_pool", "sandbox", "disk", "disk_volume"} { - kind, err := eac.LookupKind(ctx, kindName) + kind, err := eac.LookupKind(probeCtx, kindName) if err != nil { return nil, fmt.Errorf("lookup %s: %w", kindName, err) } - list, err := eac.List(ctx, kind.Attr()) - if err != nil { - return nil, fmt.Errorf("list %s: %w", kindName, err) - } - for _, entry := range list.Values() { - switch kindName { - case "sandbox_pool": - var pool compute_v1alpha.SandboxPool - pool.Decode(entry.Entity()) - r.pools = append(r.pools, pool) - case "sandbox": - var sb compute_v1alpha.Sandbox - sb.Decode(entry.Entity()) - var meta core_v1alpha.Metadata - meta.Decode(entry.Entity()) - pool, _ := meta.Labels.Get("pool") - r.sandboxes = append(r.sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) - case "disk": - var disk storage_v1alpha.Disk - disk.Decode(entry.Entity()) - r.disks = append(r.disks, disk) - case "disk_volume": - var volume storage_v1alpha.DiskVolume - volume.Decode(entry.Entity()) - r.volumes = append(r.volumes, volume) + for cursor := ""; ; { + page, err := eac.ListPage(probeCtx, kind.Attr(), cursor, 200) + if err != nil { + return nil, fmt.Errorf("list %s: %w", kindName, err) + } + for _, entry := range page.Values() { + switch kindName { + case "sandbox_pool": + var pool compute_v1alpha.SandboxPool + pool.Decode(entry.Entity()) + r.pools = append(r.pools, pool) + case "sandbox": + var sb compute_v1alpha.Sandbox + sb.Decode(entry.Entity()) + var meta core_v1alpha.Metadata + meta.Decode(entry.Entity()) + pool, _ := meta.Labels.Get("pool") + r.sandboxes = append(r.sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) + case "disk": + var disk storage_v1alpha.Disk + disk.Decode(entry.Entity()) + r.disks = append(r.disks, disk) + case "disk_volume": + var volume storage_v1alpha.DiskVolume + volume.Decode(entry.Entity()) + r.volumes = append(r.volumes, volume) + } + } + cursor = page.Cursor() + if cursor == "" { + break } } } + if err := probeCtx.Err(); err != nil { + return nil, err + } return r, nil } func gatherDoctorIndex(ctx *Context) (int64, error) { - probeCtx, cancel := context.WithTimeout(ctx, probeTimeout) + // This is a full-store scan, not a connectivity probe. Keep it bounded, + // but allow a larger cluster enough time to finish. + probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() cl, err := ctx.rpcClient(probeCtx, "entities") if err != nil { @@ -103,7 +116,7 @@ func checkEntityIndexes(env *doctorEnv) checkResult { if env.indexErr != nil { return checkResult{Status: checkWarn, Summary: "index scan unavailable", Problem: &ui.Diagnostic{ Summary: "could not check for orphaned index entries", Detail: env.indexErr.Error(), - Actions: []ui.Action{{Command: "miren version", Note: "check that the server supports index diagnostics"}, {Command: "miren doctor", Note: "retry the bounded scan"}}, + Actions: []ui.Action{{Command: "miren version", Note: "check that the server supports index diagnostics"}, {Command: "miren debug reindex --dry-run", Note: "run a complete scan without the doctor deadline"}}, }} } n := env.orphans @@ -161,7 +174,7 @@ func checkPools(env *doctorEnv) checkResult { sort.Strings(failing) return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ Summary: "service pools are crash-looping", Detail: strings.Join(failing, ", "), - Actions: []ui.Action{{Command: "miren sandbox pool list", Note: "inspect failing pools"}, {Command: "miren sandbox list --all", Note: "find failed sandboxes"}}, + Actions: []ui.Action{{Command: "miren sandbox-pool list", Note: "inspect failing pools"}, {Command: "miren sandbox list --all", Note: "find failed sandboxes"}}, }} } From f204eeb870629d09ebeb9bc854715998ba76a6b5 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Thu, 24 Sep 2026 19:30:41 +0000 Subject: [PATCH 3/6] Page app status pool lookup --- cli/commands/app_status_health.go | 22 ++++++++++++++-------- 1 file changed, 14 insertions(+), 8 deletions(-) diff --git a/cli/commands/app_status_health.go b/cli/commands/app_status_health.go index 30d14e216..8ae310544 100644 --- a/cli/commands/app_status_health.go +++ b/cli/commands/app_status_health.go @@ -131,15 +131,21 @@ func fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { if err := entityserver.NewClient(ctx.Log, eac).Get(ctx, app, &appEntity); err != nil { return nil, fmt.Errorf("get app %q: %w", app, err) } - poolRes, err := eac.List(ctx, entity.Ref(compute_v1alpha.SandboxPoolAppId, appEntity.ID)) - if err != nil { - return nil, err - } var pools []compute_v1alpha.SandboxPool - for _, entry := range poolRes.Values() { - var pool compute_v1alpha.SandboxPool - pool.Decode(entry.Entity()) - pools = append(pools, pool) + for cursor := ""; ; { + page, err := eac.ListPage(ctx, entity.Ref(compute_v1alpha.SandboxPoolAppId, appEntity.ID), cursor, 200) + if err != nil { + return nil, err + } + for _, entry := range page.Values() { + var pool compute_v1alpha.SandboxPool + pool.Decode(entry.Entity()) + pools = append(pools, pool) + } + cursor = page.Cursor() + if cursor == "" { + break + } } pools = activeServicePools(pools, appEntity.ActiveVersion) if len(pools) == 0 { From c71f47583dd96153f4022e056fe3af99c4c48a90 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Fri, 25 Sep 2026 17:49:27 +0000 Subject: [PATCH 4/6] MIR-858: Share app health classification with status and doctor --- api/app/app_v1alpha/rpc.gen.go | 181 ++++++++++++++++++ api/app/rpc.yml | 35 ++++ .../entityserver_v1alpha/rpc.gen.go | 122 ------------ api/entityserver/rpc.yml | 5 - cli/commands/app_status_health.go | 166 +++------------- cli/commands/app_status_health_test.go | 181 ++++-------------- cli/commands/doctor_check.go | 11 +- cli/commands/doctor_resources.go | 100 ++-------- servers/app/health.go | 106 ++++++++++ servers/app/health_test.go | 48 +++++ servers/app/runtime.go | 7 + servers/entityserver/entityserver.go | 13 -- servers/entityserver/entityserver_test.go | 25 --- 13 files changed, 459 insertions(+), 541 deletions(-) diff --git a/api/app/app_v1alpha/rpc.gen.go b/api/app/app_v1alpha/rpc.gen.go index f51dd10fb..b250750fb 100644 --- a/api/app/app_v1alpha/rpc.gen.go +++ b/api/app/app_v1alpha/rpc.gen.go @@ -1235,6 +1235,170 @@ func (v *PoolStatus) UnmarshalJSON(data []byte) error { return json.Unmarshal(data, &v.data) } +type serviceHealthData struct { + Service *string `cbor:"0,keyasint,omitempty" json:"service,omitempty"` + Health *string `cbor:"1,keyasint,omitempty" json:"health,omitempty"` + Running *int32 `cbor:"2,keyasint,omitempty" json:"running,omitempty"` + Dead *int32 `cbor:"3,keyasint,omitempty" json:"dead,omitempty"` + CrashCount *int64 `cbor:"4,keyasint,omitempty" json:"crash_count,omitempty"` + CooldownSeconds *int32 `cbor:"5,keyasint,omitempty" json:"cooldown_seconds,omitempty"` + LastExitCode *int64 `cbor:"6,keyasint,omitempty" json:"last_exit_code,omitempty"` + LastFailureSandbox *string `cbor:"7,keyasint,omitempty" json:"last_failure_sandbox,omitempty"` + LastFailureAt *standard.Timestamp `cbor:"8,keyasint,omitempty" json:"last_failure_at,omitempty"` +} + +type ServiceHealth struct { + data serviceHealthData +} + +func (v *ServiceHealth) HasService() bool { + return v.data.Service != nil +} + +func (v *ServiceHealth) Service() string { + if v.data.Service == nil { + return "" + } + return *v.data.Service +} + +func (v *ServiceHealth) SetService(service string) { + v.data.Service = &service +} + +func (v *ServiceHealth) HasHealth() bool { + return v.data.Health != nil +} + +func (v *ServiceHealth) Health() string { + if v.data.Health == nil { + return "" + } + return *v.data.Health +} + +func (v *ServiceHealth) SetHealth(health string) { + v.data.Health = &health +} + +func (v *ServiceHealth) HasRunning() bool { + return v.data.Running != nil +} + +func (v *ServiceHealth) Running() int32 { + if v.data.Running == nil { + return 0 + } + return *v.data.Running +} + +func (v *ServiceHealth) SetRunning(running int32) { + v.data.Running = &running +} + +func (v *ServiceHealth) HasDead() bool { + return v.data.Dead != nil +} + +func (v *ServiceHealth) Dead() int32 { + if v.data.Dead == nil { + return 0 + } + return *v.data.Dead +} + +func (v *ServiceHealth) SetDead(dead int32) { + v.data.Dead = &dead +} + +func (v *ServiceHealth) HasCrashCount() bool { + return v.data.CrashCount != nil +} + +func (v *ServiceHealth) CrashCount() int64 { + if v.data.CrashCount == nil { + return 0 + } + return *v.data.CrashCount +} + +func (v *ServiceHealth) SetCrashCount(crashCount int64) { + v.data.CrashCount = &crashCount +} + +func (v *ServiceHealth) HasCooldownSeconds() bool { + return v.data.CooldownSeconds != nil +} + +func (v *ServiceHealth) CooldownSeconds() int32 { + if v.data.CooldownSeconds == nil { + return 0 + } + return *v.data.CooldownSeconds +} + +func (v *ServiceHealth) SetCooldownSeconds(cooldownSeconds int32) { + v.data.CooldownSeconds = &cooldownSeconds +} + +func (v *ServiceHealth) HasLastExitCode() bool { + return v.data.LastExitCode != nil +} + +func (v *ServiceHealth) LastExitCode() int64 { + if v.data.LastExitCode == nil { + return 0 + } + return *v.data.LastExitCode +} + +func (v *ServiceHealth) SetLastExitCode(lastExitCode int64) { + v.data.LastExitCode = &lastExitCode +} + +func (v *ServiceHealth) HasLastFailureSandbox() bool { + return v.data.LastFailureSandbox != nil +} + +func (v *ServiceHealth) LastFailureSandbox() string { + if v.data.LastFailureSandbox == nil { + return "" + } + return *v.data.LastFailureSandbox +} + +func (v *ServiceHealth) SetLastFailureSandbox(lastFailureSandbox string) { + v.data.LastFailureSandbox = &lastFailureSandbox +} + +func (v *ServiceHealth) HasLastFailureAt() bool { + return v.data.LastFailureAt != nil +} + +func (v *ServiceHealth) LastFailureAt() *standard.Timestamp { + return v.data.LastFailureAt +} + +func (v *ServiceHealth) SetLastFailureAt(lastFailureAt *standard.Timestamp) { + v.data.LastFailureAt = lastFailureAt +} + +func (v *ServiceHealth) MarshalCBOR() ([]byte, error) { + return cbor.Marshal(v.data) +} + +func (v *ServiceHealth) UnmarshalCBOR(data []byte) error { + return cbor.Unmarshal(data, &v.data) +} + +func (v *ServiceHealth) MarshalJSON() ([]byte, error) { + return json.Marshal(v.data) +} + +func (v *ServiceHealth) UnmarshalJSON(data []byte) error { + return json.Unmarshal(data, &v.data) +} + type windowStatusData struct { Version *string `cbor:"0,keyasint,omitempty" json:"version,omitempty"` Leases *int32 `cbor:"1,keyasint,omitempty" json:"leases,omitempty"` @@ -1329,6 +1493,7 @@ type applicationStatusData struct { BoundPorts *[]*BoundPort `cbor:"19,keyasint,omitempty" json:"bound_ports,omitempty"` WorkloadRole *string `cbor:"20,keyasint,omitempty" json:"workload_role,omitempty"` MaintenanceRoutes *[]string `cbor:"21,keyasint,omitempty" json:"maintenance_routes,omitempty"` + Services *[]*ServiceHealth `cbor:"22,keyasint,omitempty" json:"services,omitempty"` } type ApplicationStatus struct { @@ -1671,6 +1836,22 @@ func (v *ApplicationStatus) SetMaintenanceRoutes(maintenanceRoutes []string) { v.data.MaintenanceRoutes = &x } +func (v *ApplicationStatus) HasServices() bool { + return v.data.Services != nil +} + +func (v *ApplicationStatus) Services() []*ServiceHealth { + if v.data.Services == nil { + return nil + } + return *v.data.Services +} + +func (v *ApplicationStatus) SetServices(services []*ServiceHealth) { + x := slices.Clone(services) + v.data.Services = &x +} + func (v *ApplicationStatus) MarshalCBOR() ([]byte, error) { return cbor.Marshal(v.data) } diff --git a/api/app/rpc.yml b/api/app/rpc.yml index 9fc6379f2..ce007df55 100644 --- a/api/app/rpc.yml +++ b/api/app/rpc.yml @@ -258,6 +258,37 @@ types: type: int64 index: 3 + - type: ServiceHealth + fields: + - name: service + type: string + index: 0 + - name: health + type: string + index: 1 + - name: running + type: int32 + index: 2 + - name: dead + type: int32 + index: 3 + - name: crashCount + type: int64 + index: 4 + - name: cooldownSeconds + type: int32 + index: 5 + - name: lastExitCode + type: int64 + index: 6 + optional: true + - name: lastFailureSandbox + type: string + index: 7 + - name: lastFailureAt + type: standard.Timestamp + index: 8 + - type: WindowStatus fields: - name: version @@ -359,6 +390,10 @@ types: element: string index: 21 doc: "Hostnames of this app's routes currently in maintenance, serving a holding page instead of the app. The default route appears as an empty string." + - name: services + type: list + element: '*ServiceHealth' + index: 22 - type: LogEntry fields: diff --git a/api/entityserver/entityserver_v1alpha/rpc.gen.go b/api/entityserver/entityserver_v1alpha/rpc.gen.go index 9e96f4d25..70b4768d8 100644 --- a/api/entityserver/entityserver_v1alpha/rpc.gen.go +++ b/api/entityserver/entityserver_v1alpha/rpc.gen.go @@ -2190,58 +2190,6 @@ func (v *EntityAccessReindexResults) UnmarshalJSON(data []byte) error { return json.Unmarshal(data, &v.data) } -type entityAccessCheckIndexHealthArgsData struct{} - -type EntityAccessCheckIndexHealthArgs struct { - call rpc.Call - data entityAccessCheckIndexHealthArgsData -} - -func (v *EntityAccessCheckIndexHealthArgs) MarshalCBOR() ([]byte, error) { - return cbor.Marshal(v.data) -} - -func (v *EntityAccessCheckIndexHealthArgs) UnmarshalCBOR(data []byte) error { - return cbor.Unmarshal(data, &v.data) -} - -func (v *EntityAccessCheckIndexHealthArgs) MarshalJSON() ([]byte, error) { - return json.Marshal(v.data) -} - -func (v *EntityAccessCheckIndexHealthArgs) UnmarshalJSON(data []byte) error { - return json.Unmarshal(data, &v.data) -} - -type entityAccessCheckIndexHealthResultsData struct { - OrphanedEntries *int64 `cbor:"0,keyasint,omitempty" json:"orphaned_entries,omitempty"` -} - -type EntityAccessCheckIndexHealthResults struct { - call rpc.Call - data entityAccessCheckIndexHealthResultsData -} - -func (v *EntityAccessCheckIndexHealthResults) SetOrphanedEntries(orphaned_entries int64) { - v.data.OrphanedEntries = &orphaned_entries -} - -func (v *EntityAccessCheckIndexHealthResults) MarshalCBOR() ([]byte, error) { - return cbor.Marshal(v.data) -} - -func (v *EntityAccessCheckIndexHealthResults) UnmarshalCBOR(data []byte) error { - return cbor.Unmarshal(data, &v.data) -} - -func (v *EntityAccessCheckIndexHealthResults) MarshalJSON() ([]byte, error) { - return json.Marshal(v.data) -} - -func (v *EntityAccessCheckIndexHealthResults) UnmarshalJSON(data []byte) error { - return json.Unmarshal(data, &v.data) -} - type entityAccessGetAttributesByTagArgsData struct { Tag *string `cbor:"0,keyasint,omitempty" json:"tag,omitempty"` } @@ -2854,32 +2802,6 @@ func (t *EntityAccessReindex) Results() *EntityAccessReindexResults { return results } -type EntityAccessCheckIndexHealth struct { - rpc.Call - args EntityAccessCheckIndexHealthArgs - results EntityAccessCheckIndexHealthResults -} - -func (t *EntityAccessCheckIndexHealth) Args() *EntityAccessCheckIndexHealthArgs { - args := &t.args - if args.call != nil { - return args - } - args.call = t.Call - t.Call.Args(args) - return args -} - -func (t *EntityAccessCheckIndexHealth) Results() *EntityAccessCheckIndexHealthResults { - results := &t.results - if results.call != nil { - return results - } - results.call = t.Call - t.Call.Results(results) - return results -} - type EntityAccessGetAttributesByTag struct { rpc.Call args EntityAccessGetAttributesByTagArgs @@ -2928,7 +2850,6 @@ type EntityAccess interface { RevokeSession(ctx context.Context, state *EntityAccessRevokeSession) error PingSession(ctx context.Context, state *EntityAccessPingSession) error Reindex(ctx context.Context, state *EntityAccessReindex) error - CheckIndexHealth(ctx context.Context, state *EntityAccessCheckIndexHealth) error GetAttributesByTag(ctx context.Context, state *EntityAccessGetAttributesByTag) error } @@ -3020,10 +2941,6 @@ func (reexportEntityAccess) Reindex(ctx context.Context, state *EntityAccessRein panic("not implemented") } -func (reexportEntityAccess) CheckIndexHealth(ctx context.Context, state *EntityAccessCheckIndexHealth) error { - panic("not implemented") -} - func (reexportEntityAccess) GetAttributesByTag(ctx context.Context, state *EntityAccessGetAttributesByTag) error { panic("not implemented") } @@ -3244,16 +3161,6 @@ func AdaptEntityAccess(t EntityAccess) *rpc.Interface { return t.Reindex(ctx, &EntityAccessReindex{Call: call}) }, }, - { - Name: "check_index_health", - InterfaceName: "EntityAccess", - Index: 0, - Public: false, - Params: []string{}, - Handler: func(ctx context.Context, call rpc.Call) error { - return t.CheckIndexHealth(ctx, &EntityAccessCheckIndexHealth{Call: call}) - }, - }, { Name: "get_attributes_by_tag", InterfaceName: "EntityAccess", @@ -4026,35 +3933,6 @@ func (v EntityAccessClient) Reindex(ctx context.Context, dry_run bool) (*EntityA return &EntityAccessClientReindexResults{client: v.Client, data: ret}, nil } -type EntityAccessClientCheckIndexHealthResults struct { - client rpc.Client - data entityAccessCheckIndexHealthResultsData -} - -func (v *EntityAccessClientCheckIndexHealthResults) HasOrphanedEntries() bool { - return v.data.OrphanedEntries != nil -} - -func (v *EntityAccessClientCheckIndexHealthResults) OrphanedEntries() int64 { - if v.data.OrphanedEntries == nil { - return 0 - } - return *v.data.OrphanedEntries -} - -func (v EntityAccessClient) CheckIndexHealth(ctx context.Context) (*EntityAccessClientCheckIndexHealthResults, error) { - args := EntityAccessCheckIndexHealthArgs{} - - var ret entityAccessCheckIndexHealthResultsData - - err := v.Call(ctx, "check_index_health", &args, &ret) - if err != nil { - return nil, err - } - - return &EntityAccessClientCheckIndexHealthResults{client: v.Client, data: ret}, nil -} - type EntityAccessClientGetAttributesByTagResults struct { client rpc.Client data entityAccessGetAttributesByTagResultsData diff --git a/api/entityserver/rpc.yml b/api/entityserver/rpc.yml index 3dada1db4..a5f0a3f43 100644 --- a/api/entityserver/rpc.yml +++ b/api/entityserver/rpc.yml @@ -373,11 +373,6 @@ interfaces: type: list element: ReindexStat - - name: check_index_health - results: - - name: orphaned_entries - type: int64 - - name: get_attributes_by_tag parameters: - name: tag diff --git a/cli/commands/app_status_health.go b/cli/commands/app_status_health.go index 8ae310544..2a7ad50f8 100644 --- a/cli/commands/app_status_health.go +++ b/cli/commands/app_status_health.go @@ -2,98 +2,30 @@ package commands import ( "fmt" - "slices" - "sort" "strings" "time" "miren.dev/runtime/api/app/app_v1alpha" - "miren.dev/runtime/api/compute/compute_v1alpha" - "miren.dev/runtime/api/core/core_v1alpha" - "miren.dev/runtime/api/entityserver" - "miren.dev/runtime/api/entityserver/entityserver_v1alpha" - "miren.dev/runtime/pkg/entity" "miren.dev/runtime/pkg/rpc/standard" ) -const failureWindow = 10 * time.Minute - type serviceHealth struct { - Service string `json:"service"` - Running int `json:"running"` - Dead int `json:"dead"` - CrashStreak int `json:"crash_streak"` - CrashLooping bool `json:"crash_looping"` - LastExitCode *int64 `json:"last_exit_code,omitempty"` - LastFailureLog string `json:"last_failure_log,omitempty"` - - lastSandbox string - lastExit time.Time - lastFailure time.Time -} - -type sandboxHealthRecord struct { - Sandbox compute_v1alpha.Sandbox - Pool string - Updated time.Time -} - -func summarizeServiceHealth(pools []compute_v1alpha.SandboxPool, sandboxes []sandboxHealthRecord, now time.Time) []serviceHealth { - byPool := make(map[string]*serviceHealth) - byService := make(map[string]*serviceHealth) - for _, pool := range pools { - h := byService[pool.Service] - if h == nil { - h = &serviceHealth{Service: pool.Service} - byService[pool.Service] = h - } - byPool[pool.ID.String()] = h - if pool.DesiredInstances > pool.ReadyInstances && !pool.LastCrashTime.After(now) && pool.LastCrashTime.After(now.Add(-failureWindow)) { - h.CrashStreak += int(pool.ConsecutiveCrashCount) - if pool.CooldownUntil.After(now) && pool.ConsecutiveCrashCount >= 2 { - h.CrashLooping = true - } - } - } - for _, record := range sandboxes { - sb := record.Sandbox - h := byPool[record.Pool] - if h == nil { - continue - } - switch sb.Status { - case compute_v1alpha.RUNNING: - h.Running++ - case compute_v1alpha.DEAD: - h.Dead++ - if !sb.Exit.At.IsZero() && (h.LastExitCode == nil || sb.Exit.At.After(h.lastExit)) { - code := sb.Exit.Code - h.LastExitCode = &code - h.lastExit = sb.Exit.At - } - if sb.Exit.Code != 0 && !sb.Exit.At.IsZero() && (h.lastSandbox == "" || sb.Exit.At.After(h.lastFailure)) { - h.lastFailure = sb.Exit.At - h.lastSandbox = sb.ID.String() - } - case compute_v1alpha.PENDING, compute_v1alpha.NOT_READY, compute_v1alpha.STOPPED: - } - } - result := make([]serviceHealth, 0, len(byService)) - for _, h := range byService { - result = append(result, *h) - } - sort.Slice(result, func(i, j int) bool { return result[i].Service < result[j].Service }) - return result + Service string `json:"service"` + Running int `json:"running"` + Dead int `json:"dead"` + CrashStreak int `json:"crash_streak"` + CrashLooping bool `json:"crash_looping"` + CooldownSeconds int32 `json:"cooldown_seconds,omitempty"` + LastExitCode *int64 `json:"last_exit_code,omitempty"` + LastFailureLog string `json:"last_failure_log,omitempty"` } func renderServiceHealth(health []serviceHealth) string { var b strings.Builder for _, svc := range health { state := fmt.Sprintf("%d running, %d dead", svc.Running, svc.Dead) - if svc.CrashStreak > 0 { - state += fmt.Sprintf("; %d crashes in current streak (last within 10m)", svc.CrashStreak) - } if svc.CrashLooping { + state += fmt.Sprintf("; %d crashes in current streak, retry in %s", svc.CrashStreak, formatDuration(time.Duration(svc.CooldownSeconds)*time.Second)) state = infoRed.Render("crash-looping") + ", " + state } fmt.Fprintf(&b, " %s: %s\n", svc.Service, state) @@ -110,84 +42,30 @@ func renderServiceHealth(health []serviceHealth) string { return b.String() } -func activeServicePools(pools []compute_v1alpha.SandboxPool, version entity.Id) []compute_v1alpha.SandboxPool { - var active []compute_v1alpha.SandboxPool - for _, pool := range pools { - if version == "" || slices.Contains(pool.ReferencedByVersions, version) || len(pool.ReferencedByVersions) == 0 && pool.SandboxSpec.Version == version { - active = append(active, pool) - } - } - return active -} - func fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { - cl, err := ctx.RPCClient("entities") + cl, err := ctx.RPCClient("dev.miren.runtime/app") if err != nil { return nil, err } defer cl.Close() - eac := entityserver_v1alpha.NewEntityAccessClient(cl) - var appEntity core_v1alpha.App - if err := entityserver.NewClient(ctx.Log, eac).Get(ctx, app, &appEntity); err != nil { - return nil, fmt.Errorf("get app %q: %w", app, err) - } - var pools []compute_v1alpha.SandboxPool - for cursor := ""; ; { - page, err := eac.ListPage(ctx, entity.Ref(compute_v1alpha.SandboxPoolAppId, appEntity.ID), cursor, 200) - if err != nil { - return nil, err - } - for _, entry := range page.Values() { - var pool compute_v1alpha.SandboxPool - pool.Decode(entry.Entity()) - pools = append(pools, pool) - } - cursor = page.Cursor() - if cursor == "" { - break - } - } - pools = activeServicePools(pools, appEntity.ActiveVersion) - if len(pools) == 0 { - return []serviceHealth{}, nil - } - kind, err := eac.LookupKind(ctx, "sandbox") + res, err := app_v1alpha.NewAppStatusClient(cl).AppInfo(ctx, app) if err != nil { return nil, err } - var sandboxes []sandboxHealthRecord - selected := make(map[string]bool, len(pools)) - for _, pool := range pools { - selected[pool.ID.String()] = true + if res.Status().Health() == "unknown" && res.Status().ActiveVersion() != "" { + return nil, fmt.Errorf("service health unavailable for %s", app) } - // Pool is a metadata label, not an indexed attribute. Page the kind index - // rather than materializing every sandbox in the cluster in one RPC. - for cursor := ""; ; { - page, err := eac.ListPage(ctx, kind.Attr(), cursor, 200) - if err != nil { - return nil, err + var health []serviceHealth + for _, svc := range res.Status().Services() { + h := serviceHealth{Service: svc.Service(), Running: int(svc.Running()), Dead: int(svc.Dead()), CrashStreak: int(svc.CrashCount()), CrashLooping: svc.Health() == "crashed", CooldownSeconds: svc.CooldownSeconds()} + if svc.HasLastExitCode() { + code := svc.LastExitCode() + h.LastExitCode = &code } - for _, entry := range page.Values() { - var meta core_v1alpha.Metadata - meta.Decode(entry.Entity()) - pool, _ := meta.Labels.Get("pool") - if !selected[pool] { - continue - } - var sb compute_v1alpha.Sandbox - sb.Decode(entry.Entity()) - sandboxes = append(sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) - } - cursor = page.Cursor() - if cursor == "" { - break - } - } - health := summarizeServiceHealth(pools, sandboxes, time.Now()) - for i := range health { - if health[i].lastSandbox != "" && health[i].lastFailure.After(time.Now().Add(-failureWindow)) && health[i].CrashStreak > 0 { - health[i].LastFailureLog = recentSandboxFailureLog(ctx, health[i].lastSandbox, health[i].lastFailure) + if svc.HasLastFailureAt() && svc.LastFailureSandbox() != "" && svc.Health() == "crashed" { + h.LastFailureLog = recentSandboxFailureLog(ctx, svc.LastFailureSandbox(), standard.FromTimestamp(svc.LastFailureAt())) } + health = append(health, h) } return health, nil } diff --git a/cli/commands/app_status_health_test.go b/cli/commands/app_status_health_test.go index ffba4126a..defe9e31f 100644 --- a/cli/commands/app_status_health_test.go +++ b/cli/commands/app_status_health_test.go @@ -5,174 +5,65 @@ import ( "encoding/json" "strings" "testing" - "time" - "miren.dev/runtime/api/compute/compute_v1alpha" + "miren.dev/runtime/api/app/app_v1alpha" "miren.dev/runtime/api/storage/storage_v1alpha" "miren.dev/runtime/pkg/entity" ) -func TestSummarizeServiceHealth(t *testing.T) { - now := time.Date(2026, 9, 23, 12, 0, 0, 0, time.UTC) - pools := []compute_v1alpha.SandboxPool{ - {ID: "pool-db", Service: "db", DesiredInstances: 1, ConsecutiveCrashCount: 2, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}, - {ID: "pool-web", Service: "web"}, - } - sandboxes := []sandboxHealthRecord{ - {Pool: "pool-web", Sandbox: compute_v1alpha.Sandbox{ID: "web-1", Status: compute_v1alpha.RUNNING}}, - {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-1", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-9 * time.Minute)}}}, - {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-2", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-2 * time.Minute)}}}, - {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-3", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 5, At: now.Add(-11 * time.Minute)}}}, - {Pool: "pool-other", Sandbox: compute_v1alpha.Sandbox{ID: "foreign", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 99, At: now}}}, - } - got := summarizeServiceHealth(pools, sandboxes, now) - if len(got) != 2 || got[0].Service != "db" || got[0].Dead != 3 || got[0].Running != 0 || got[0].CrashStreak != 2 || !got[0].CrashLooping || got[0].LastExitCode == nil || *got[0].LastExitCode != 17 || got[0].lastSandbox != "db-2" { - t.Fatalf("db summary = %+v", got) - } - if got[1].Service != "web" || got[1].Running != 1 || got[1].CrashLooping { - t.Fatalf("web summary = %+v", got[1]) - } - got[0].LastFailureLog = "database unavailable\nconnection refused" - output := renderServiceHealth(got) - for _, fragment := range []string{"db: crash-looping, 0 running, 3 dead; 2 crashes in current streak", "Last exit code: 17", " connection refused", "web: 1 running, 0 dead"} { +func TestRenderServiceHealth(t *testing.T) { + code := int64(0) + health := []serviceHealth{ + {Service: "db", Running: 0, Dead: 3, CrashStreak: 7, CrashLooping: true, CooldownSeconds: 240, LastExitCode: &code, LastFailureLog: "connection refused"}, + {Service: "web", Running: 1}, + } + output := renderServiceHealth(health) + for _, fragment := range []string{"db: crash-looping, 0 running, 3 dead; 7 crashes in current streak, retry in", "Last exit code: 0", "connection refused", "web: 1 running, 0 dead"} { if !strings.Contains(output, fragment) { t.Errorf("status output %q missing %q", output, fragment) } } if strings.Contains(output[strings.Index(output, "web:"):], "crash") { - t.Fatalf("healthy web service reports a crash: %q", output) + t.Fatalf("healthy service reports a crash: %q", output) } - encoded, err := json.Marshal(got[0]) - if err != nil || !strings.Contains(string(encoded), `"last_exit_code":17`) || !strings.Contains(string(encoded), `"crash_streak":2`) { + encoded, err := json.Marshal(health[0]) + if err != nil || !strings.Contains(string(encoded), `"last_exit_code":0`) || !strings.Contains(string(encoded), `"crash_streak":7`) { t.Fatalf("JSON status = %s, err = %v", encoded, err) } - zero := summarizeServiceHealth(pools, []sandboxHealthRecord{{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-zero", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-time.Minute)}}}}, now)[0] - if zero.LastExitCode == nil || *zero.LastExitCode != 0 { - t.Fatalf("zero exit code must be present: %+v", zero) - } - early := summarizeServiceHealth(pools, []sandboxHealthRecord{{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "db-early", Status: compute_v1alpha.DEAD}, Updated: now.Add(-time.Minute)}}, now)[0] - if early.CrashStreak != 2 || early.LastExitCode != nil || early.lastSandbox != "" { - t.Fatalf("early failure without an exit record = %+v", early) - } - latest := summarizeServiceHealth(pools, []sandboxHealthRecord{ - {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "success", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-time.Minute)}}}, - {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "failure", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-2 * time.Minute)}}}, - {Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{ID: "retired", Status: compute_v1alpha.DEAD}, Updated: now}, - }, now)[0] - if latest.LastExitCode == nil || *latest.LastExitCode != 0 || latest.lastSandbox != "failure" || !latest.lastFailure.Equal(now.Add(-2*time.Minute)) { - t.Fatalf("newer successful exit hid earlier failure or was replaced by retirement: %+v", latest) - } - if idle := summarizeServiceHealth([]compute_v1alpha.SandboxPool{{ID: "pool-db", Service: "db", DesiredInstances: 0, ConsecutiveCrashCount: 4, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}}, sandboxes, now)[0]; idle.CrashStreak != 0 || idle.CrashLooping { - t.Fatalf("scaled-to-zero service should not appear to be failing: %+v", idle) - } -} - -func TestScaleDownIsNotACrashLoop(t *testing.T) { - now := time.Date(2026, 9, 23, 12, 0, 0, 0, time.UTC) - retired := make([]sandboxHealthRecord, 3) - for i := range retired { - retired[i] = sandboxHealthRecord{Pool: "pool-web", Sandbox: compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD}, Updated: now.Add(-time.Minute)} - } - for _, desired := range []int64{1, 0, 1} { // 4→1, scaled to zero, then woken - pool := compute_v1alpha.SandboxPool{ID: "pool-web", Service: "web", DesiredInstances: desired} - got := summarizeServiceHealth([]compute_v1alpha.SandboxPool{pool}, retired, now)[0] - if got.Dead != 3 || got.CrashStreak != 0 || got.CrashLooping { - t.Fatalf("desired %d: %+v", desired, got) - } - } - env := verdictEnv(remoteHost, nil, probeOpen, probeOpen) - env.resources = &doctorResources{pools: []compute_v1alpha.SandboxPool{{ID: "pool-web", Service: "web", DesiredInstances: 1}}, sandboxes: retired} - if got := checkSandboxes(env); got.Status != checkOK { - t.Fatalf("doctor treated scale-down as a failure: %+v", got) - } - // A later intentional retirement must not replace the last known failed exit. - pool := compute_v1alpha.SandboxPool{ID: "pool-web", Service: "web", DesiredInstances: 1, ConsecutiveCrashCount: 1, LastCrashTime: now.Add(-2 * time.Minute)} - retired = append(retired, sandboxHealthRecord{Pool: "pool-web", Sandbox: compute_v1alpha.Sandbox{ID: "failed", Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-2 * time.Minute)}}}) - got := summarizeServiceHealth([]compute_v1alpha.SandboxPool{pool}, retired, now)[0] - if got.LastExitCode == nil || *got.LastExitCode != 17 || got.lastSandbox != "failed" { - t.Fatalf("known failure lost to retirement: %+v", got) - } -} - -func TestActiveServicePoolsUsesVersionMembership(t *testing.T) { - pools := []compute_v1alpha.SandboxPool{ - {ID: "current", SandboxSpec: compute_v1alpha.SandboxSpec{Version: "app-version-456"}, ReferencedByVersions: []entity.Id{"app-version-456", "app-version-123"}}, - {ID: "old", SandboxSpec: compute_v1alpha.SandboxSpec{Version: "app-version-123"}, ReferencedByVersions: []entity.Id{"app-version-456"}}, - {ID: "legacy", SandboxSpec: compute_v1alpha.SandboxSpec{Version: "app-version-789"}}, - } - got := activeServicePools(pools, "app-version-123") - if len(got) != 1 || got[0].ID != "current" { - t.Fatalf("pools for active version = %+v", got) - } - if got := activeServicePools(pools, "app-version-789"); len(got) != 1 || got[0].ID != "legacy" { - t.Fatalf("legacy pool membership = %+v", got) - } - if got := activeServicePools(pools, ""); len(got) != 3 { - t.Fatalf("pools before first deployment = %+v", got) - } } func TestDoctorResourceChecks(t *testing.T) { - now := time.Now() env := verdictEnv(remoteHost, nil, probeOpen, probeOpen) env.resources = &doctorResources{ - pools: []compute_v1alpha.SandboxPool{{ID: "pool-db", Service: "db", DesiredInstances: 1, ReadyInstances: 0, ConsecutiveCrashCount: 3, LastCrashTime: now.Add(-time.Minute), CooldownUntil: now.Add(time.Minute)}}, disks: []storage_v1alpha.Disk{{ID: entity.Id("disk-1"), Status: storage_v1alpha.ERROR}}, volumes: []storage_v1alpha.DiskVolume{{ID: entity.Id("volume-1"), ActualState: storage_v1alpha.DV_ERROR, ErrorMessage: "attach failed"}}, } - env.orphans = 2 - for i := 0; i < 3; i++ { - env.resources.sandboxes = append(env.resources.sandboxes, sandboxHealthRecord{Pool: "pool-db", Sandbox: compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{At: now.Add(-time.Duration(i+1) * time.Minute)}}}) + crashed := &app_v1alpha.AppInfo{} + crashed.SetName("demo") + crashed.SetHealth("crashed") + crashed.SetCrashCount(7) + crashed.SetCooldownSeconds(240) + healthy := &app_v1alpha.AppInfo{} + healthy.SetName("web") + healthy.SetHealth("healthy") + env.apps = []*app_v1alpha.AppInfo{healthy, crashed} + if got := checkApps(env); got.Status != checkFail || !strings.Contains(got.Summary, "demo (7 crashes") || got.Problem == nil { + t.Fatalf("apps check = %+v", got) + } + if got := checkDisks(env); got.Status != checkFail || !strings.Contains(got.Summary, "attach failed") { + t.Fatalf("disk check = %+v", got) + } + env.apps = []*app_v1alpha.AppInfo{healthy} + if got := checkApps(env); got.Status != checkOK { + t.Fatalf("healthy apps check = %+v", got) + } + env.appsErr = context.DeadlineExceeded + if got := checkApps(env); got.Status != checkSkip { + t.Fatalf("incomplete app scan reported health: %+v", got) } - for _, tc := range []struct { - name string - check func(*doctorEnv) checkResult - want checkStatus - fragment string - }{ - {"indexes", checkEntityIndexes, checkWarn, "2 orphaned"}, - {"sandboxes", checkSandboxes, checkFail, "3 consecutive crashes"}, - {"pools", checkPools, checkFail, "db"}, - {"disks", checkDisks, checkFail, "attach failed"}, - } { - t.Run(tc.name, func(t *testing.T) { - got := tc.check(env) - if got.Status != tc.want || !strings.Contains(got.Summary, tc.fragment) || got.Problem == nil || len(got.Problem.Actions) == 0 { - t.Fatalf("check = %+v, want %s with action and %q", got, tc.want, tc.fragment) - } - }) - } - if got := checkPools(env); got.Problem == nil || got.Problem.Actions[0].Command != "miren sandbox-pool list" { - t.Fatalf("pool diagnostic suggests invalid command: %+v", got) - } - env.indexErr = context.DeadlineExceeded - if got := checkEntityIndexes(env); got.Status != checkWarn { - t.Fatalf("timed-out index check = %+v, want warning", got) - } - resources := env.resources env.resources = nil env.resourcesErr = context.DeadlineExceeded - env.indexErr = nil - for _, check := range []func(*doctorEnv) checkResult{checkSandboxes, checkPools, checkDisks} { - if got := check(env); got.Status != checkSkip { - t.Fatalf("incomplete resource scan reported health: %+v", got) - } - } - if got := checkEntityIndexes(env); got.Status != checkWarn { - t.Fatalf("independent index check = %+v, want warning for two orphans", got) - } - env.resourcesErr = nil - env.resources = resources - if got := checkPools(env); got.Status != checkFail { - t.Fatalf("pool check = %+v, want independent failure", got) - } - // Two apps with a service called db must not combine into a false alert. - env.resources.pools[0].App = "app-a" - env.resources.pools[0].ConsecutiveCrashCount = 2 - env.resources.pools = append(env.resources.pools, compute_v1alpha.SandboxPool{ID: "pool-other", App: "app-b", Service: "db"}) - env.resources.sandboxes = env.resources.sandboxes[:2] - env.resources.sandboxes = append(env.resources.sandboxes, sandboxHealthRecord{Pool: "pool-other", Sandbox: compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{At: now.Add(-time.Minute)}}}) - if got := checkSandboxes(env); got.Status != checkOK { - t.Fatalf("distinct app failures = %+v, want no aggregate alert", got) + if got := checkDisks(env); got.Status != checkSkip { + t.Fatalf("incomplete disk scan reported health: %+v", got) } } diff --git a/cli/commands/doctor_check.go b/cli/commands/doctor_check.go index af32779d8..3a114edbc 100644 --- a/cli/commands/doctor_check.go +++ b/cli/commands/doctor_check.go @@ -5,6 +5,7 @@ import ( "errors" "sync" + "miren.dev/runtime/api/app/app_v1alpha" "miren.dev/runtime/clientconfig" "miren.dev/runtime/pkg/release" "miren.dev/runtime/pkg/ui" @@ -118,8 +119,8 @@ type doctorEnv struct { resources *doctorResources resourcesErr error - orphans int64 - indexErr error + apps []*app_v1alpha.AppInfo + appsErr error } // local reports whether the active cluster runs on this machine, which decides @@ -198,7 +199,7 @@ func gatherCluster(ctx *Context, opts ConfigCentric, env *doctorEnv) { env.resources, env.resourcesErr = gatherDoctorResources(ctx) }) wg.Go(func() { - env.orphans, env.indexErr = gatherDoctorIndex(ctx) + env.apps, env.appsErr = gatherDoctorApps(ctx) }) wg.Go(func() { @@ -225,9 +226,7 @@ func doctorChecks() []check { {Name: "Server", Run: checkServer}, {Name: "Version", Run: checkVersion}, {Name: "Authentication", Run: checkAuthentication}, - {Name: "Entity indexes", Run: checkEntityIndexes}, - {Name: "Sandboxes", Run: checkSandboxes}, - {Name: "Pools", Run: checkPools}, + {Name: "Apps", Run: checkApps}, {Name: "Disks and volumes", Run: checkDisks}, } } diff --git a/cli/commands/doctor_resources.go b/cli/commands/doctor_resources.go index 7d564aaf4..085ad40a8 100644 --- a/cli/commands/doctor_resources.go +++ b/cli/commands/doctor_resources.go @@ -7,18 +7,15 @@ import ( "strings" "time" - "miren.dev/runtime/api/compute/compute_v1alpha" - "miren.dev/runtime/api/core/core_v1alpha" + "miren.dev/runtime/api/app/app_v1alpha" "miren.dev/runtime/api/entityserver/entityserver_v1alpha" "miren.dev/runtime/api/storage/storage_v1alpha" "miren.dev/runtime/pkg/ui" ) type doctorResources struct { - pools []compute_v1alpha.SandboxPool - sandboxes []sandboxHealthRecord - disks []storage_v1alpha.Disk - volumes []storage_v1alpha.DiskVolume + disks []storage_v1alpha.Disk + volumes []storage_v1alpha.DiskVolume } func gatherDoctorResources(ctx *Context) (*doctorResources, error) { @@ -31,7 +28,7 @@ func gatherDoctorResources(ctx *Context) (*doctorResources, error) { defer cl.Close() eac := entityserver_v1alpha.NewEntityAccessClient(cl) r := &doctorResources{} - for _, kindName := range []string{"sandbox_pool", "sandbox", "disk", "disk_volume"} { + for _, kindName := range []string{"disk", "disk_volume"} { kind, err := eac.LookupKind(probeCtx, kindName) if err != nil { return nil, fmt.Errorf("lookup %s: %w", kindName, err) @@ -43,17 +40,6 @@ func gatherDoctorResources(ctx *Context) (*doctorResources, error) { } for _, entry := range page.Values() { switch kindName { - case "sandbox_pool": - var pool compute_v1alpha.SandboxPool - pool.Decode(entry.Entity()) - r.pools = append(r.pools, pool) - case "sandbox": - var sb compute_v1alpha.Sandbox - sb.Decode(entry.Entity()) - var meta core_v1alpha.Metadata - meta.Decode(entry.Entity()) - pool, _ := meta.Labels.Get("pool") - r.sandboxes = append(r.sandboxes, sandboxHealthRecord{Sandbox: sb, Pool: pool, Updated: time.UnixMilli(entry.UpdatedAt())}) case "disk": var disk storage_v1alpha.Disk disk.Decode(entry.Entity()) @@ -76,21 +62,19 @@ func gatherDoctorResources(ctx *Context) (*doctorResources, error) { return r, nil } -func gatherDoctorIndex(ctx *Context) (int64, error) { - // This is a full-store scan, not a connectivity probe. Keep it bounded, - // but allow a larger cluster enough time to finish. +func gatherDoctorApps(ctx *Context) ([]*app_v1alpha.AppInfo, error) { probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() - cl, err := ctx.rpcClient(probeCtx, "entities") + cl, err := ctx.rpcClient(probeCtx, "dev.miren.runtime/app") if err != nil { - return 0, err + return nil, err } defer cl.Close() - result, err := entityserver_v1alpha.NewEntityAccessClient(cl).CheckIndexHealth(probeCtx) + res, err := app_v1alpha.NewCrudClient(cl).List(probeCtx) if err != nil { - return 0, err + return nil, err } - return result.OrphanedEntries(), nil + return res.Apps(), nil } func resourceUnavailable(env *doctorEnv) (checkResult, bool) { @@ -106,75 +90,29 @@ func resourceUnavailable(env *doctorEnv) (checkResult, bool) { return checkResult{}, false } -func checkEntityIndexes(env *doctorEnv) checkResult { +func checkApps(env *doctorEnv) checkResult { if !env.configured() { return checkResult{Status: checkSkip, Summary: "(no cluster configured)"} } if env.connErr != nil { return checkResult{Status: checkSkip, Summary: "(server unreachable)"} } - if env.indexErr != nil { - return checkResult{Status: checkWarn, Summary: "index scan unavailable", Problem: &ui.Diagnostic{ - Summary: "could not check for orphaned index entries", Detail: env.indexErr.Error(), - Actions: []ui.Action{{Command: "miren version", Note: "check that the server supports index diagnostics"}, {Command: "miren debug reindex --dry-run", Note: "run a complete scan without the doctor deadline"}}, - }} - } - n := env.orphans - if n == 0 { - return checkResult{Status: checkOK, Summary: "no orphaned index entries"} - } - return checkResult{Status: checkWarn, Summary: fmt.Sprintf("%d orphaned index entries", n), Problem: &ui.Diagnostic{ - Summary: fmt.Sprintf("%d index entries point to missing entities", n), - Actions: []ui.Action{{Command: "miren debug reindex --dry-run", Note: "inspect index drift"}, {Command: "miren debug reindex", Note: "repair indexes after review"}}, - }} -} - -func checkSandboxes(env *doctorEnv) checkResult { - if skip, unavailable := resourceUnavailable(env); unavailable { - return skip - } - // The same service name may be used by several apps. Keep their counts separate. - pools := append([]compute_v1alpha.SandboxPool(nil), env.resources.pools...) - for i := range pools { - if pools[i].App != "" { - pools[i].Service = pools[i].App.String() + "/" + pools[i].Service - } - } - services := summarizeServiceHealth(pools, env.resources.sandboxes, time.Now()) - var failing []string - for _, svc := range services { - if svc.CrashStreak >= 3 { - failing = append(failing, fmt.Sprintf("%s (%d consecutive crashes, last in 10m)", svc.Service, svc.CrashStreak)) - } - } - if len(failing) == 0 { - return checkResult{Status: checkOK, Summary: "no elevated recent sandbox crash streaks"} - } - return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ - Summary: "services have repeated sandbox failures", - Detail: strings.Join(failing, ", "), - Actions: []ui.Action{{Command: "miren sandbox list --all", Note: "inspect dead sandboxes"}, {Command: "miren app status -a ", Note: "view exit codes and failure logs"}}, - }} -} - -func checkPools(env *doctorEnv) checkResult { - if skip, unavailable := resourceUnavailable(env); unavailable { - return skip + if env.appsErr != nil { + return checkResult{Status: checkSkip, Summary: fmt.Sprintf("(app health unavailable: %v)", env.appsErr)} } var failing []string - now := time.Now() - for _, pool := range env.resources.pools { - if pool.DesiredInstances > 0 && pool.ReadyInstances == 0 && pool.CooldownUntil.After(now) && pool.ConsecutiveCrashCount >= 2 { - failing = append(failing, fmt.Sprintf("%s/%s (%d crashes, retry in %s)", pool.App, pool.Service, pool.ConsecutiveCrashCount, formatDuration(pool.CooldownUntil.Sub(now)))) + for _, app := range env.apps { + if app.Health() == "crashed" { + failing = append(failing, fmt.Sprintf("%s (%d crashes, retry in %s)", app.Name(), app.CrashCount(), formatDuration(time.Duration(app.CooldownSeconds())*time.Second))) } } if len(failing) == 0 { - return checkResult{Status: checkOK, Summary: "no pools in crash cooldown"} + return checkResult{Status: checkOK, Summary: "no apps in crash cooldown"} } sort.Strings(failing) return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ - Summary: "service pools are crash-looping", Detail: strings.Join(failing, ", "), - Actions: []ui.Action{{Command: "miren sandbox-pool list", Note: "inspect failing pools"}, {Command: "miren sandbox list --all", Note: "find failed sandboxes"}}, + Summary: "apps are crash-looping", Detail: strings.Join(failing, ", "), + Actions: []ui.Action{{Command: "miren app status -a ", Note: "inspect service failures and exit codes"}, {Command: "miren sandbox list --all", Note: "inspect dead sandboxes"}}, }} } diff --git a/servers/app/health.go b/servers/app/health.go index eedd667b0..d61bf21b9 100644 --- a/servers/app/health.go +++ b/servers/app/health.go @@ -2,6 +2,7 @@ package app import ( "context" + "sort" "time" "miren.dev/runtime/api/app/app_v1alpha" @@ -9,6 +10,7 @@ import ( core_v1alpha "miren.dev/runtime/api/core/core_v1alpha" "miren.dev/runtime/pkg/apphealth" "miren.dev/runtime/pkg/entity" + "miren.dev/runtime/pkg/rpc/standard" ) // specAllowsScaleToZero reports whether an app's resolved config lets it sit at @@ -73,6 +75,110 @@ func (h *poolHealth) accumulate(pool *compute_v1alpha.SandboxPool, now time.Time } } +// serviceSandboxHealth uses the same pool classifier as app list, while +// retaining the sandbox details needed to explain a failure in app status. +type serviceSandboxHealth struct { + pool poolHealth + running int32 + dead int32 + lastExit time.Time + lastCode int64 + hasExit bool + lastFailed time.Time + failedID string +} + +func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1alpha.SandboxPool, spec *core_v1alpha.ConfigSpec, now time.Time) ([]*app_v1alpha.ServiceHealth, error) { + byService := make(map[string]*serviceSandboxHealth) + byPool := make(map[string]*serviceSandboxHealth) + for i := range pools { + pool := &pools[i] + h := byService[pool.Service] + if h == nil { + h = &serviceSandboxHealth{pool: poolHealth{isAutoscale: true}} + if spec != nil { + for _, svc := range spec.Services { + if svc.Name == pool.Service && svc.Concurrency.Mode == "fixed" { + h.pool.isAutoscale = false + } + } + } + byService[pool.Service] = h + } + h.pool.accumulate(pool, now) + byPool[pool.ID.String()] = h + } + if len(pools) == 0 { + return nil, nil + } + + list, err := r.EC.List(ctx, entity.Ref(entity.EntityKind, compute_v1alpha.KindSandbox)) + if err != nil { + return nil, err + } + for list.Next() { + md := list.Metadata() + if md == nil { + continue + } + poolID, _ := md.Labels.Get("pool") + h := byPool[poolID] + if h == nil { + continue + } + var sb compute_v1alpha.Sandbox + if err := list.Read(&sb); err != nil { + continue + } + switch sb.Status { + case compute_v1alpha.RUNNING: + h.running++ + case compute_v1alpha.DEAD: + h.dead++ + if !sb.Exit.At.IsZero() { + if !h.hasExit || sb.Exit.At.After(h.lastExit) { + h.lastCode = sb.Exit.Code + h.lastExit = sb.Exit.At + h.hasExit = true + } + if sb.Exit.Code != 0 && (h.failedID == "" || sb.Exit.At.After(h.lastFailed)) { + h.lastFailed = sb.Exit.At + h.failedID = sb.ID.String() + } + } + case compute_v1alpha.PENDING, compute_v1alpha.NOT_READY, compute_v1alpha.STOPPED: + } + } + + names := make([]string, 0, len(byService)) + for name := range byService { + names = append(names, name) + } + sort.Strings(names) + result := make([]*app_v1alpha.ServiceHealth, 0, len(names)) + for _, name := range names { + h := byService[name] + var svc app_v1alpha.ServiceHealth + svc.SetService(name) + svc.SetHealth(h.pool.classify()) + svc.SetRunning(h.running) + svc.SetDead(h.dead) + if h.pool.inCooldown { + svc.SetCrashCount(h.pool.crashCount) + svc.SetCooldownSeconds(int32(h.pool.cooldownLeft.Seconds())) + } + if h.hasExit { + svc.SetLastExitCode(h.lastCode) + } + if h.failedID != "" { + svc.SetLastFailureSandbox(h.failedID) + svc.SetLastFailureAt(standard.ToTimestamp(h.lastFailed)) + } + result = append(result, &svc) + } + return result, nil +} + // collectBoundPortDivergence scans the sandboxes belonging to the given pools // and returns the ports they actually bound that diverge from the configured // port. The sandbox controller records a bound_port component only on diff --git a/servers/app/health_test.go b/servers/app/health_test.go index b0127423e..aef89ceff 100644 --- a/servers/app/health_test.go +++ b/servers/app/health_test.go @@ -1,14 +1,20 @@ package app import ( + "context" + "log/slog" "testing" "time" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" "miren.dev/runtime/api/compute/compute_v1alpha" "miren.dev/runtime/api/core/core_v1alpha" + "miren.dev/runtime/api/entityserver" "miren.dev/runtime/pkg/apphealth" + "miren.dev/runtime/pkg/entity/testutils" + "miren.dev/runtime/pkg/entity/types" ) func TestPoolHealthClassify(t *testing.T) { @@ -77,6 +83,48 @@ func TestPoolHealthAccumulate_ExpiredCooldownIgnored(t *testing.T) { assert.Equal(t, apphealth.Healthy, h.classify()) } +func TestCollectServiceHealthUsesSharedCooldownAndLatestExit(t *testing.T) { + ctx := context.Background() + inmem, cleanup := testutils.NewInMemEntityServer(t) + t.Cleanup(cleanup) + ec := entityserver.NewClient(slog.Default(), inmem.EAC) + r := &AppInfo{Log: slog.Default(), EC: ec} + now := time.Now().Truncate(time.Second) + pools := []compute_v1alpha.SandboxPool{ + {ID: "pool-db", Service: "db", DesiredInstances: 1, ConsecutiveCrashCount: 7, CooldownUntil: now.Add(3 * time.Minute), LastCrashTime: now.Add(-12 * time.Minute)}, + {ID: "pool-web", Service: "web", DesiredInstances: 1, ReadyInstances: 1}, + } + failedID, err := ec.Create(ctx, "failed", &compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 17, At: now.Add(-12 * time.Minute)}}, entityserver.WithLabels(types.LabelSet("pool", "pool-db"))) + require.NoError(t, err) + _, err = ec.Create(ctx, "success", &compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 0, At: now.Add(-11 * time.Minute)}}, entityserver.WithLabels(types.LabelSet("pool", "pool-db"))) + require.NoError(t, err) + _, err = ec.Create(ctx, "retired", &compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD}, entityserver.WithLabels(types.LabelSet("pool", "pool-db"))) + require.NoError(t, err) + _, err = ec.Create(ctx, "web", &compute_v1alpha.Sandbox{Status: compute_v1alpha.RUNNING}, entityserver.WithLabels(types.LabelSet("pool", "pool-web"))) + require.NoError(t, err) + _, err = ec.Create(ctx, "foreign", &compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 99, At: now}}, entityserver.WithLabels(types.LabelSet("pool", "other"))) + require.NoError(t, err) + + got, err := r.collectServiceHealth(ctx, pools, nil, now) + require.NoError(t, err) + require.Len(t, got, 2) + assert.Equal(t, "db", got[0].Service()) + assert.Equal(t, apphealth.Crashed, got[0].Health(), "cooldown still active past the 10-minute failure window") + assert.EqualValues(t, 7, got[0].CrashCount()) + assert.EqualValues(t, 3, got[0].Dead()) + require.True(t, got[0].HasLastExitCode()) + assert.Zero(t, got[0].LastExitCode(), "latest successful exit replaces the older code") + assert.Equal(t, failedID.String(), got[0].LastFailureSandbox()) + assert.Equal(t, apphealth.Healthy, got[1].Health()) + assert.EqualValues(t, 1, got[1].Running()) + fixed, err := r.collectServiceHealth(ctx, + []compute_v1alpha.SandboxPool{{ID: "pool-fixed", Service: "fixed"}}, + &core_v1alpha.ConfigSpec{Services: []core_v1alpha.ConfigSpecServices{{Name: "fixed", Concurrency: core_v1alpha.ConfigSpecServicesConcurrency{Mode: "fixed"}}}}, now) + require.NoError(t, err) + require.Len(t, fixed, 1) + assert.Equal(t, apphealth.Starting, fixed[0].Health(), "fixed service at zero must not be classified as idle") +} + func TestSpecNeedsNoService(t *testing.T) { assert.False(t, specNeedsNoService(nil)) assert.False(t, specNeedsNoService(&core_v1alpha.ConfigSpec{}), "an empty app has no valid workload") diff --git a/servers/app/runtime.go b/servers/app/runtime.go index 483b8b2e6..e0d037e3b 100644 --- a/servers/app/runtime.go +++ b/servers/app/runtime.go @@ -183,6 +183,7 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn rai.SetHealth(apphealth.Unknown) } else { var pools []*app_v1alpha.PoolStatus + var servicePools []compute_v1alpha.SandboxPool poolIDs := make(map[string]bool) hasInstances := false now := time.Now() @@ -203,6 +204,7 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn for poolsResp.Next() { var pool compute_v1alpha.SandboxPool poolsResp.Read(&pool) + servicePools = append(servicePools, pool) poolIDs[pool.ID.String()] = true if pool.CurrentInstances > 0 { @@ -229,6 +231,11 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn } rai.SetPools(pools) + services, err := a.collectServiceHealth(ctx, servicePools, spec, now) + if err != nil { + return err + } + rai.SetServices(services) rai.SetHealth(health.classify()) rai.SetReadyInstances(int32(health.ready)) rai.SetDesiredInstances(int32(health.desired)) diff --git a/servers/entityserver/entityserver.go b/servers/entityserver/entityserver.go index 6c43ce206..9d74a8f79 100644 --- a/servers/entityserver/entityserver.go +++ b/servers/entityserver/entityserver.go @@ -1303,19 +1303,6 @@ func (e *EntityServer) Reindex(ctx context.Context, req *entityserver_v1alpha.En return nil } -func (e *EntityServer) CheckIndexHealth(ctx context.Context, req *entityserver_v1alpha.EntityAccessCheckIndexHealth) error { - store, ok := e.Store.(*entity.EtcdStore) - if !ok { - return fmt.Errorf("index health requires EtcdStore") - } - stats, err := store.CleanupStaleCollectionEntries(ctx, e.Log, entity.CleanupOptions{DryRun: true}) - if err != nil { - return fmt.Errorf("scan index health: %w", err) - } - req.Results().SetOrphanedEntries(stats.OrphanedEntriesFound) - return nil -} - func (e *EntityServer) GetAttributesByTag(ctx context.Context, req *entityserver_v1alpha.EntityAccessGetAttributesByTag) error { args := req.Args() tag := args.Tag() diff --git a/servers/entityserver/entityserver_test.go b/servers/entityserver/entityserver_test.go index 603121279..067cbb6ed 100644 --- a/servers/entityserver/entityserver_test.go +++ b/servers/entityserver/entityserver_test.go @@ -10,7 +10,6 @@ import ( "testing" "time" - "github.com/mr-tron/base58" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "go.etcd.io/etcd/api/v3/etcdserverpb" @@ -30,30 +29,6 @@ func setupTestEtcd(t *testing.T) (*clientv3.Client, string) { return etcdtest.TestEtcdClient(t) } -func TestCheckIndexHealthIsReadOnly(t *testing.T) { - client, prefix := setupTestEtcd(t) - store, err := entity.NewEtcdStore(t.Context(), slog.Default(), client, prefix) - require.NoError(t, err) - server, err := NewEntityServer(slog.Default(), store) - require.NoError(t, err) - id := entity.Id("missing/sandbox") - key := prefix + "/collections/test_orphan/" + base58.Encode([]byte(id)) - _, err = client.Put(t.Context(), key, id.String()) - require.NoError(t, err) - - eac := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} - result, err := eac.CheckIndexHealth(t.Context()) - require.NoError(t, err) - assert.EqualValues(t, 1, result.OrphanedEntries()) - resp, err := client.Get(t.Context(), key) - require.NoError(t, err) - assert.Len(t, resp.Kvs, 1, "diagnostic must not repair the index") - canceled, cancel := context.WithCancel(t.Context()) - cancel() - _, err = eac.CheckIndexHealth(canceled) - require.Error(t, err, "an interrupted scan must not report a clean result") -} - func TestEntityServer_Get(t *testing.T) { store := entity.NewMockStore() server := &EntityServer{ From a4d5a39abb34d31b4ecdbe91e0353d8010d42125 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Fri, 25 Sep 2026 18:10:47 +0000 Subject: [PATCH 5/6] MIR-858: Keep app health available when sandbox detail scan fails --- cli/commands/app_status_health.go | 4 +- servers/app/health.go | 79 ++++++++----------------------- servers/app/health_test.go | 27 +++++++++-- servers/app/runtime.go | 22 +++------ 4 files changed, 52 insertions(+), 80 deletions(-) diff --git a/cli/commands/app_status_health.go b/cli/commands/app_status_health.go index 2a7ad50f8..7e92c7524 100644 --- a/cli/commands/app_status_health.go +++ b/cli/commands/app_status_health.go @@ -43,7 +43,7 @@ func renderServiceHealth(health []serviceHealth) string { } func fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { - cl, err := ctx.RPCClient("dev.miren.runtime/app") + cl, err := ctx.RPCClient(rpcAppStatus) if err != nil { return nil, err } @@ -52,7 +52,7 @@ func fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { if err != nil { return nil, err } - if res.Status().Health() == "unknown" && res.Status().ActiveVersion() != "" { + if res.Status().ActiveVersion() != "" && (res.Status().Health() == "unknown" || !res.Status().HasServices()) { return nil, fmt.Errorf("service health unavailable for %s", app) } var health []serviceHealth diff --git a/servers/app/health.go b/servers/app/health.go index d61bf21b9..91e5da481 100644 --- a/servers/app/health.go +++ b/servers/app/health.go @@ -88,7 +88,7 @@ type serviceSandboxHealth struct { failedID string } -func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1alpha.SandboxPool, spec *core_v1alpha.ConfigSpec, now time.Time) ([]*app_v1alpha.ServiceHealth, error) { +func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1alpha.SandboxPool, spec *core_v1alpha.ConfigSpec, now time.Time, hasInstances bool) ([]*app_v1alpha.ServiceHealth, []*app_v1alpha.BoundPort, error) { byService := make(map[string]*serviceSandboxHealth) byPool := make(map[string]*serviceSandboxHealth) for i := range pools { @@ -109,13 +109,15 @@ func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1al byPool[pool.ID.String()] = h } if len(pools) == 0 { - return nil, nil + return nil, nil, nil } list, err := r.EC.List(ctx, entity.Ref(entity.EntityKind, compute_v1alpha.KindSandbox)) if err != nil { - return nil, err + return nil, nil, err } + seenPorts := make(map[int64]bool) + var boundPorts []*app_v1alpha.BoundPort for list.Next() { md := list.Metadata() if md == nil { @@ -148,6 +150,20 @@ func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1al } case compute_v1alpha.PENDING, compute_v1alpha.NOT_READY, compute_v1alpha.STOPPED: } + // The controller only records bound_port on divergence. Preserve the + // existing behavior of reporting it for running or booting instances. + if hasInstances && (sb.Status == compute_v1alpha.RUNNING || sb.Status == compute_v1alpha.PENDING) { + for _, bp := range sb.BoundPort { + if bp.Port == 0 || seenPorts[bp.Port] { + continue + } + seenPorts[bp.Port] = true + var port app_v1alpha.BoundPort + port.SetPort(bp.Port) + port.SetAddress(bp.Address) + boundPorts = append(boundPorts, &port) + } + } } names := make([]string, 0, len(byService)) @@ -176,62 +192,7 @@ func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1al } result = append(result, &svc) } - return result, nil -} - -// collectBoundPortDivergence scans the sandboxes belonging to the given pools -// and returns the ports they actually bound that diverge from the configured -// port. The sandbox controller records a bound_port component only on -// divergence (MIR-1246), so any bound_port present is by definition a port the -// app chose for itself. Best effort: a listing error just yields no divergence. -func (r *AppInfo) collectBoundPortDivergence(ctx context.Context, poolIDs map[string]bool) []*app_v1alpha.BoundPort { - if len(poolIDs) == 0 { - return nil - } - - sbList, err := r.EC.List(ctx, entity.Ref(entity.EntityKind, compute_v1alpha.KindSandbox)) - if err != nil { - r.Log.Warn("failed to list sandboxes for bound-port check", "error", err) - return nil - } - - seen := make(map[int64]bool) - var result []*app_v1alpha.BoundPort - - for sbList.Next() { - var sb compute_v1alpha.Sandbox - if err := sbList.Read(&sb); err != nil { - continue - } - - md := sbList.Metadata() - if md == nil { - continue - } - poolLabel, _ := md.Labels.Get("pool") - if !poolIDs[poolLabel] { - continue - } - - // Only living sandboxes describe where the app is serving now. - if sb.Status != compute_v1alpha.RUNNING && sb.Status != compute_v1alpha.PENDING { - continue - } - - for _, bp := range sb.BoundPort { - if bp.Port == 0 || seen[bp.Port] { - continue - } - seen[bp.Port] = true - - var rbp app_v1alpha.BoundPort - rbp.SetPort(bp.Port) - rbp.SetAddress(bp.Address) - result = append(result, &rbp) - } - } - - return result + return result, boundPorts, nil } // classify maps the aggregate to a health string. A pool in cooldown is diff --git a/servers/app/health_test.go b/servers/app/health_test.go index aef89ceff..fe24180d0 100644 --- a/servers/app/health_test.go +++ b/servers/app/health_test.go @@ -2,6 +2,7 @@ package app import ( "context" + "errors" "log/slog" "testing" "time" @@ -13,6 +14,7 @@ import ( "miren.dev/runtime/api/core/core_v1alpha" "miren.dev/runtime/api/entityserver" "miren.dev/runtime/pkg/apphealth" + "miren.dev/runtime/pkg/entity" "miren.dev/runtime/pkg/entity/testutils" "miren.dev/runtime/pkg/entity/types" ) @@ -100,13 +102,15 @@ func TestCollectServiceHealthUsesSharedCooldownAndLatestExit(t *testing.T) { require.NoError(t, err) _, err = ec.Create(ctx, "retired", &compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD}, entityserver.WithLabels(types.LabelSet("pool", "pool-db"))) require.NoError(t, err) - _, err = ec.Create(ctx, "web", &compute_v1alpha.Sandbox{Status: compute_v1alpha.RUNNING}, entityserver.WithLabels(types.LabelSet("pool", "pool-web"))) + _, err = ec.Create(ctx, "web", &compute_v1alpha.Sandbox{Status: compute_v1alpha.RUNNING, BoundPort: []compute_v1alpha.BoundPort{{Port: 8888, Address: "0.0.0.0"}}}, entityserver.WithLabels(types.LabelSet("pool", "pool-web"))) require.NoError(t, err) _, err = ec.Create(ctx, "foreign", &compute_v1alpha.Sandbox{Status: compute_v1alpha.DEAD, Exit: compute_v1alpha.Exit{Code: 99, At: now}}, entityserver.WithLabels(types.LabelSet("pool", "other"))) require.NoError(t, err) - got, err := r.collectServiceHealth(ctx, pools, nil, now) + got, ports, err := r.collectServiceHealth(ctx, pools, nil, now, true) require.NoError(t, err) + require.Len(t, ports, 1) + assert.EqualValues(t, 8888, ports[0].Port()) require.Len(t, got, 2) assert.Equal(t, "db", got[0].Service()) assert.Equal(t, apphealth.Crashed, got[0].Health(), "cooldown still active past the 10-minute failure window") @@ -117,14 +121,29 @@ func TestCollectServiceHealthUsesSharedCooldownAndLatestExit(t *testing.T) { assert.Equal(t, failedID.String(), got[0].LastFailureSandbox()) assert.Equal(t, apphealth.Healthy, got[1].Health()) assert.EqualValues(t, 1, got[1].Running()) - fixed, err := r.collectServiceHealth(ctx, + fixed, ports, err := r.collectServiceHealth(ctx, []compute_v1alpha.SandboxPool{{ID: "pool-fixed", Service: "fixed"}}, - &core_v1alpha.ConfigSpec{Services: []core_v1alpha.ConfigSpecServices{{Name: "fixed", Concurrency: core_v1alpha.ConfigSpecServicesConcurrency{Mode: "fixed"}}}}, now) + &core_v1alpha.ConfigSpec{Services: []core_v1alpha.ConfigSpecServices{{Name: "fixed", Concurrency: core_v1alpha.ConfigSpecServicesConcurrency{Mode: "fixed"}}}}, now, false) require.NoError(t, err) + assert.Empty(t, ports) require.Len(t, fixed, 1) assert.Equal(t, apphealth.Starting, fixed[0].Health(), "fixed service at zero must not be classified as idle") } +func TestCollectServiceHealthListFailure(t *testing.T) { + inmem, cleanup := testutils.NewInMemEntityServer(t) + t.Cleanup(cleanup) + inmem.Store.OnListIndex = func(context.Context, entity.Attr) ([]entity.Id, error) { + return nil, errors.New("sandbox index unavailable") + } + r := &AppInfo{Log: slog.Default(), EC: entityserver.NewClient(slog.Default(), inmem.EAC)} + services, ports, err := r.collectServiceHealth(context.Background(), + []compute_v1alpha.SandboxPool{{ID: "pool-web", Service: "web", DesiredInstances: 1}}, nil, time.Now(), true) + require.ErrorContains(t, err, "sandbox index unavailable") + assert.Empty(t, services) + assert.Empty(t, ports) +} + func TestSpecNeedsNoService(t *testing.T) { assert.False(t, specNeedsNoService(nil)) assert.False(t, specNeedsNoService(&core_v1alpha.ConfigSpec{}), "an empty app has no valid workload") diff --git a/servers/app/runtime.go b/servers/app/runtime.go index e0d037e3b..435505815 100644 --- a/servers/app/runtime.go +++ b/servers/app/runtime.go @@ -184,7 +184,6 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn } else { var pools []*app_v1alpha.PoolStatus var servicePools []compute_v1alpha.SandboxPool - poolIDs := make(map[string]bool) hasInstances := false now := time.Now() @@ -206,7 +205,6 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn poolsResp.Read(&pool) servicePools = append(servicePools, pool) - poolIDs[pool.ID.String()] = true if pool.CurrentInstances > 0 { hasInstances = true } @@ -231,11 +229,15 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn } rai.SetPools(pools) - services, err := a.collectServiceHealth(ctx, servicePools, spec, now) + services, boundPorts, err := a.collectServiceHealth(ctx, servicePools, spec, now, hasInstances) if err != nil { - return err + a.Log.Warn("failed to collect service sandbox details", "error", err) + } else { + rai.SetServices(services) + if len(boundPorts) > 0 { + rai.SetBoundPorts(boundPorts) + } } - rai.SetServices(services) rai.SetHealth(health.classify()) rai.SetReadyInstances(int32(health.ready)) rai.SetDesiredInstances(int32(health.desired)) @@ -244,16 +246,6 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn rai.SetCooldownSeconds(int32(health.cooldownLeft.Seconds())) } - // Scan for port divergence once there's an instance (running or - // booting), not only once it's ready: a wrong bound port can be the - // reason ready stays 0, and that's exactly when we want to surface - // it. Skipping when there are no instances keeps the global sandbox - // scan out of the scaled-to-zero path. - if hasInstances { - if bp := a.collectBoundPortDivergence(ctx, poolIDs); len(bp) > 0 { - rai.SetBoundPorts(bp) - } - } } } else { rai.SetHealth(apphealth.Unknown) From f856605abbfe4cae0e2717e0805832a8bf608448 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Fri, 25 Sep 2026 21:41:18 +0000 Subject: [PATCH 6/6] MIR-858: Preserve empty service health through RPC encoding --- servers/app/health.go | 2 +- servers/app/health_test.go | 17 +++++++++++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/servers/app/health.go b/servers/app/health.go index 91e5da481..19e145810 100644 --- a/servers/app/health.go +++ b/servers/app/health.go @@ -109,7 +109,7 @@ func (r *AppInfo) collectServiceHealth(ctx context.Context, pools []compute_v1al byPool[pool.ID.String()] = h } if len(pools) == 0 { - return nil, nil, nil + return []*app_v1alpha.ServiceHealth{}, nil, nil } list, err := r.EC.List(ctx, entity.Ref(entity.EntityKind, compute_v1alpha.KindSandbox)) diff --git a/servers/app/health_test.go b/servers/app/health_test.go index fe24180d0..16dbeb14f 100644 --- a/servers/app/health_test.go +++ b/servers/app/health_test.go @@ -7,9 +7,11 @@ import ( "testing" "time" + "github.com/fxamacker/cbor/v2" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "miren.dev/runtime/api/app/app_v1alpha" "miren.dev/runtime/api/compute/compute_v1alpha" "miren.dev/runtime/api/core/core_v1alpha" "miren.dev/runtime/api/entityserver" @@ -144,6 +146,21 @@ func TestCollectServiceHealthListFailure(t *testing.T) { assert.Empty(t, ports) } +func TestEmptyServiceHealthSurvivesRPCEncoding(t *testing.T) { + r := &AppInfo{} + services, ports, err := r.collectServiceHealth(context.Background(), nil, nil, time.Now(), false) + require.NoError(t, err) + assert.Empty(t, ports) + var status app_v1alpha.ApplicationStatus + status.SetServices(services) + data, err := cbor.Marshal(&status) + require.NoError(t, err) + var decoded app_v1alpha.ApplicationStatus + require.NoError(t, cbor.Unmarshal(data, &decoded)) + assert.True(t, decoded.HasServices(), "empty services must differ from an unavailable sandbox scan") + assert.Empty(t, decoded.Services()) +} + func TestSpecNeedsNoService(t *testing.T) { assert.False(t, specNeedsNoService(nil)) assert.False(t, specNeedsNoService(&core_v1alpha.ConfigSpec{}), "an empty app has no valid workload")