diff --git a/platform/service/observability.go b/platform/service/observability.go index d03a3b7..fc280a2 100644 --- a/platform/service/observability.go +++ b/platform/service/observability.go @@ -35,12 +35,12 @@ func (svc *CoreService) IngestMetricBatch(batch domain.MetricBatchIngest) (domai return domain.MetricBatchIngestResult{}, validationError("metric sample server must belong to runEndpointId") } sample.RunEndpointID = batch.RunEndpointID - if sample.ID == "" { - sample.ID = fmt.Sprintf("metric:%s:%d:%d", sample.ServerInstanceID, sample.CollectedAt.UnixNano(), index) - } if sample.CollectedAt.IsZero() { sample.CollectedAt = svc.now() } + if sample.ID == "" { + sample.ID = fmt.Sprintf("metric:%s:%d:%d", sample.ServerInstanceID, sample.CollectedAt.UnixNano(), index) + } if err := validator.ValidateMetricSample(sample); err != nil { return domain.MetricBatchIngestResult{}, err } diff --git a/platform/service/resources.go b/platform/service/resources.go index e56fefc..e794a22 100644 --- a/platform/service/resources.go +++ b/platform/service/resources.go @@ -2200,19 +2200,14 @@ func (svc *CoreService) latestMetricsForServer(instance domain.ServerInstance) d } } if found { - return domain.ServerMetrics{ + metrics := domain.ServerMetrics{ ServerInstanceID: latest.ServerInstanceID, Online: latest.Online, - PlayerCount: latest.PlayerCount, - MaxPlayers: latest.MaxPlayers, - TPS: latest.TPS, - LatencyMS: latest.LatencyMS, - CPUPercent: latest.CPUPercent, - MemoryPercent: latest.MemoryPercent, - DiskPercent: latest.DiskPercent, Source: latest.Source, CollectedAt: latest.CollectedAt, } + mergeRecentMetricFields(&metrics, samples, latest.CollectedAt, instance.RunEndpointID) + return metrics } } return domain.ServerMetrics{ @@ -2223,6 +2218,49 @@ func (svc *CoreService) latestMetricsForServer(instance domain.ServerInstance) d } } +const recentMetricFieldWindow = 2 * time.Minute + +// mergeRecentMetricFields keeps independently reported fields visible while +// refusing to carry values forward indefinitely when a collector stops. +func mergeRecentMetricFields(metrics *domain.ServerMetrics, samples []domain.MetricSample, latestAt time.Time, runEndpointID string) { + cutoff := latestAt.Add(-recentMetricFieldWindow) + var playerAt, maxPlayersAt, tpsAt, latencyAt, cpuAt, memoryAt, diskAt time.Time + for i := range samples { + sample := samples[i] + if (sample.RunEndpointID != "" && sample.RunEndpointID != runEndpointID) || sample.CollectedAt.Before(cutoff) || sample.CollectedAt.After(latestAt) { + continue + } + if sample.PlayerCount != nil && (metrics.PlayerCount == nil || sample.CollectedAt.After(playerAt)) { + metrics.PlayerCount = sample.PlayerCount + playerAt = sample.CollectedAt + } + if sample.MaxPlayers != nil && (metrics.MaxPlayers == nil || sample.CollectedAt.After(maxPlayersAt)) { + metrics.MaxPlayers = sample.MaxPlayers + maxPlayersAt = sample.CollectedAt + } + if sample.TPS != nil && (metrics.TPS == nil || sample.CollectedAt.After(tpsAt)) { + metrics.TPS = sample.TPS + tpsAt = sample.CollectedAt + } + if sample.LatencyMS != nil && (metrics.LatencyMS == nil || sample.CollectedAt.After(latencyAt)) { + metrics.LatencyMS = sample.LatencyMS + latencyAt = sample.CollectedAt + } + if sample.CPUPercent != nil && (metrics.CPUPercent == nil || sample.CollectedAt.After(cpuAt)) { + metrics.CPUPercent = sample.CPUPercent + cpuAt = sample.CollectedAt + } + if sample.MemoryPercent != nil && (metrics.MemoryPercent == nil || sample.CollectedAt.After(memoryAt)) { + metrics.MemoryPercent = sample.MemoryPercent + memoryAt = sample.CollectedAt + } + if sample.DiskPercent != nil && (metrics.DiskPercent == nil || sample.CollectedAt.After(diskAt)) { + metrics.DiskPercent = sample.DiskPercent + diskAt = sample.CollectedAt + } + } +} + func buildLogicalServerConfig(instance domain.ServerInstance) string { lines := []string{ "# platform logical server config", diff --git a/platform/service/resources_test.go b/platform/service/resources_test.go index b3bd400..1e4c1d5 100644 --- a/platform/service/resources_test.go +++ b/platform/service/resources_test.go @@ -741,6 +741,41 @@ func TestCoreServiceMetricsAndConfigReadAreRoleScoped(t *testing.T) { } } +func TestCoreServiceMergesRecentPartialMetricSamples(t *testing.T) { + svc := newTestCoreService() + plugin, endpoint := createPluginAndRunEndpoint(t, svc) + ownerSession := createServiceUserAndLogin(t, svc, domain.User{ID: "user-partial-metrics", DisplayName: "Partial Metrics Owner", Email: "partial-metrics@example.test", Roles: []string{"server-owner"}, PasswordHash: "secret-password"}) + instance, err := svc.CreateServerInstanceForSession(ownerSession, domain.ServerInstance{ID: "server-partial-metrics", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "Partial Metrics", State: domain.ServerInstanceStateRunning}) + if err != nil { + t.Fatalf("create server: %v", err) + } + runHello := validRunControlHello() + runHello.RunEndpointID = endpoint.ID + registered, err := svc.RegisterRunHello(runHello) + if err != nil { + t.Fatalf("register run: %v", err) + } + cpu := 37.0 + memory := 49.0 + players := 8 + disk := 78.0 + staleMemory := 12.0 + if _, err := svc.IngestMetricBatch(domain.MetricBatchIngest{RunEndpointID: endpoint.ID, SessionToken: registered.SessionToken, Samples: []domain.MetricSample{ + {ServerInstanceID: instance.ID, Online: true, CPUPercent: &cpu, MemoryPercent: &memory, PlayerCount: &players, Source: "run", CollectedAt: fixedTime.Add(-15 * time.Second)}, + {ServerInstanceID: instance.ID, Online: true, MemoryPercent: &staleMemory, Source: "run", CollectedAt: fixedTime.Add(-3 * time.Minute)}, + {ServerInstanceID: instance.ID, Online: true, DiskPercent: &disk, Source: "run", CollectedAt: fixedTime}, + }}); err != nil { + t.Fatalf("ingest partial metrics: %v", err) + } + metrics, err := svc.ListServerMetricsForSession(ownerSession) + if err != nil { + t.Fatalf("list merged metrics: %v", err) + } + if len(metrics) != 1 || metrics[0].CPUPercent == nil || *metrics[0].CPUPercent != cpu || metrics[0].MemoryPercent == nil || *metrics[0].MemoryPercent != memory || metrics[0].PlayerCount == nil || *metrics[0].PlayerCount != players || metrics[0].DiskPercent == nil || *metrics[0].DiskPercent != disk || metrics[0].CollectedAt != fixedTime { + t.Fatalf("expected recent partial metrics to merge, got %+v", metrics) + } +} + func TestCoreServiceConfigWriteAndFileDispatchAreScoped(t *testing.T) { svc := newTestCoreService() plugin, endpoint := createPluginAndRunEndpoint(t, svc)