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/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..7e92c7524 --- /dev/null +++ b/cli/commands/app_status_health.go @@ -0,0 +1,101 @@ +package commands + +import ( + "fmt" + "strings" + "time" + + "miren.dev/runtime/api/app/app_v1alpha" + "miren.dev/runtime/pkg/rpc/standard" +) + +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"` + 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.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) + 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 fetchServiceHealth(ctx *Context, app string) ([]serviceHealth, error) { + cl, err := ctx.RPCClient(rpcAppStatus) + if err != nil { + return nil, err + } + defer cl.Close() + res, err := app_v1alpha.NewAppStatusClient(cl).AppInfo(ctx, app) + if err != nil { + return nil, err + } + if res.Status().ActiveVersion() != "" && (res.Status().Health() == "unknown" || !res.Status().HasServices()) { + return nil, fmt.Errorf("service health unavailable for %s", app) + } + 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 + } + 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 +} + +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..defe9e31f --- /dev/null +++ b/cli/commands/app_status_health_test.go @@ -0,0 +1,69 @@ +package commands + +import ( + "context" + "encoding/json" + "strings" + "testing" + + "miren.dev/runtime/api/app/app_v1alpha" + "miren.dev/runtime/api/storage/storage_v1alpha" + "miren.dev/runtime/pkg/entity" +) + +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 service reports a crash: %q", output) + } + 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) + } +} + +func TestDoctorResourceChecks(t *testing.T) { + env := verdictEnv(remoteHost, nil, probeOpen, probeOpen) + env.resources = &doctorResources{ + 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"}}, + } + 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) + } + env.resources = nil + env.resourcesErr = context.DeadlineExceeded + 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 9805dcf1a..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" @@ -115,6 +116,11 @@ type doctorEnv struct { // auth is only attempted when the cluster names an identity. auth authResult + + resources *doctorResources + resourcesErr error + apps []*app_v1alpha.AppInfo + appsErr error } // local reports whether the active cluster runs on this machine, which decides @@ -189,6 +195,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.apps, env.appsErr = gatherDoctorApps(ctx) + }) wg.Go(func() { env.tcp = probeTCP(env.cluster.Hostname) @@ -214,5 +226,7 @@ func doctorChecks() []check { {Name: "Server", Run: checkServer}, {Name: "Version", Run: checkVersion}, {Name: "Authentication", Run: checkAuthentication}, + {Name: "Apps", Run: checkApps}, + {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..085ad40a8 --- /dev/null +++ b/cli/commands/doctor_resources.go @@ -0,0 +1,142 @@ +package commands + +import ( + "context" + "fmt" + "sort" + "strings" + "time" + + "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 { + disks []storage_v1alpha.Disk + volumes []storage_v1alpha.DiskVolume +} + +func gatherDoctorResources(ctx *Context) (*doctorResources, error) { + probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + cl, err := ctx.rpcClient(probeCtx, "entities") + if err != nil { + return nil, err + } + defer cl.Close() + eac := entityserver_v1alpha.NewEntityAccessClient(cl) + r := &doctorResources{} + 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) + } + 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 "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 gatherDoctorApps(ctx *Context) ([]*app_v1alpha.AppInfo, error) { + probeCtx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + cl, err := ctx.rpcClient(probeCtx, "dev.miren.runtime/app") + if err != nil { + return nil, err + } + defer cl.Close() + res, err := app_v1alpha.NewCrudClient(cl).List(probeCtx) + if err != nil { + return nil, err + } + return res.Apps(), 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 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.appsErr != nil { + return checkResult{Status: checkSkip, Summary: fmt.Sprintf("(app health unavailable: %v)", env.appsErr)} + } + var failing []string + 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 apps in crash cooldown"} + } + sort.Strings(failing) + return checkResult{Status: checkFail, Summary: strings.Join(failing, ", "), Problem: &ui.Diagnostic{ + 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"}}, + }} +} + +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/app/health.go b/servers/app/health.go index eedd667b0..19e145810 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,59 +75,124 @@ func (h *poolHealth) accumulate(pool *compute_v1alpha.SandboxPool, now time.Time } } -// 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 +// 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, hasInstances bool) ([]*app_v1alpha.ServiceHealth, []*app_v1alpha.BoundPort, 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 []*app_v1alpha.ServiceHealth{}, nil, nil } - sbList, err := r.EC.List(ctx, entity.Ref(entity.EntityKind, compute_v1alpha.KindSandbox)) + list, 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 + return nil, nil, err } - - 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() + seenPorts := make(map[int64]bool) + var boundPorts []*app_v1alpha.BoundPort + for list.Next() { + md := list.Metadata() if md == nil { continue } - poolLabel, _ := md.Labels.Get("pool") - if !poolIDs[poolLabel] { + poolID, _ := md.Labels.Get("pool") + h := byPool[poolID] + if h == nil { continue } - - // Only living sandboxes describe where the app is serving now. - if sb.Status != compute_v1alpha.RUNNING && sb.Status != compute_v1alpha.PENDING { + var sb compute_v1alpha.Sandbox + if err := list.Read(&sb); err != nil { continue } - - for _, bp := range sb.BoundPort { - if bp.Port == 0 || seen[bp.Port] { - 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: + } + // 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) } - seen[bp.Port] = true - - var rbp app_v1alpha.BoundPort - rbp.SetPort(bp.Port) - rbp.SetAddress(bp.Address) - result = append(result, &rbp) } } - return result + 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, 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 b0127423e..16dbeb14f 100644 --- a/servers/app/health_test.go +++ b/servers/app/health_test.go @@ -1,14 +1,24 @@ package app import ( + "context" + "errors" + "log/slog" "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" "miren.dev/runtime/pkg/apphealth" + "miren.dev/runtime/pkg/entity" + "miren.dev/runtime/pkg/entity/testutils" + "miren.dev/runtime/pkg/entity/types" ) func TestPoolHealthClassify(t *testing.T) { @@ -77,6 +87,80 @@ 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, 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, 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") + 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, 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, 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 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") diff --git a/servers/app/runtime.go b/servers/app/runtime.go index 483b8b2e6..435505815 100644 --- a/servers/app/runtime.go +++ b/servers/app/runtime.go @@ -183,7 +183,7 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn rai.SetHealth(apphealth.Unknown) } else { var pools []*app_v1alpha.PoolStatus - poolIDs := make(map[string]bool) + var servicePools []compute_v1alpha.SandboxPool hasInstances := false now := time.Now() @@ -203,8 +203,8 @@ 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 { hasInstances = true } @@ -229,6 +229,15 @@ func (a *AppInfo) AppInfo(ctx context.Context, state *app_v1alpha.AppStatusAppIn } rai.SetPools(pools) + services, boundPorts, err := a.collectServiceHealth(ctx, servicePools, spec, now, hasInstances) + if err != nil { + a.Log.Warn("failed to collect service sandbox details", "error", err) + } else { + rai.SetServices(services) + if len(boundPorts) > 0 { + rai.SetBoundPorts(boundPorts) + } + } rai.SetHealth(health.classify()) rai.SetReadyInstances(int32(health.ready)) rai.SetDesiredInstances(int32(health.desired)) @@ -237,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)