fix partial server metric aggregation
This commit is contained in:
@@ -35,12 +35,12 @@ func (svc *CoreService) IngestMetricBatch(batch domain.MetricBatchIngest) (domai
|
|||||||
return domain.MetricBatchIngestResult{}, validationError("metric sample server must belong to runEndpointId")
|
return domain.MetricBatchIngestResult{}, validationError("metric sample server must belong to runEndpointId")
|
||||||
}
|
}
|
||||||
sample.RunEndpointID = batch.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() {
|
if sample.CollectedAt.IsZero() {
|
||||||
sample.CollectedAt = svc.now()
|
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 {
|
if err := validator.ValidateMetricSample(sample); err != nil {
|
||||||
return domain.MetricBatchIngestResult{}, err
|
return domain.MetricBatchIngestResult{}, err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2200,19 +2200,14 @@ func (svc *CoreService) latestMetricsForServer(instance domain.ServerInstance) d
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if found {
|
if found {
|
||||||
return domain.ServerMetrics{
|
metrics := domain.ServerMetrics{
|
||||||
ServerInstanceID: latest.ServerInstanceID,
|
ServerInstanceID: latest.ServerInstanceID,
|
||||||
Online: latest.Online,
|
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,
|
Source: latest.Source,
|
||||||
CollectedAt: latest.CollectedAt,
|
CollectedAt: latest.CollectedAt,
|
||||||
}
|
}
|
||||||
|
mergeRecentMetricFields(&metrics, samples, latest.CollectedAt, instance.RunEndpointID)
|
||||||
|
return metrics
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return domain.ServerMetrics{
|
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 {
|
func buildLogicalServerConfig(instance domain.ServerInstance) string {
|
||||||
lines := []string{
|
lines := []string{
|
||||||
"# platform logical server config",
|
"# platform logical server config",
|
||||||
|
|||||||
@@ -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) {
|
func TestCoreServiceConfigWriteAndFileDispatchAreScoped(t *testing.T) {
|
||||||
svc := newTestCoreService()
|
svc := newTestCoreService()
|
||||||
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
|
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
|
||||||
|
|||||||
Reference in New Issue
Block a user