From 953b8f63ac73e57ae1c17cde23f85fa69edc0b2c Mon Sep 17 00:00:00 2001 From: et-nik Date: Wed, 23 Sep 2026 12:16:47 +0200 Subject: [PATCH] windows stats --- internal/processmanager/README.md | 35 +- internal/processmanager/errors.go | 3 + .../process_snapshot_windows.go | 158 +++++++++ .../process_snapshot_windows_test.go | 98 ++++++ internal/processmanager/process_tree.go | 135 ++++++++ .../processmanager/process_tree_metrics.go | 217 ++++++++++++ .../process_tree_metrics_test.go | 326 ++++++++++++++++++ internal/processmanager/process_tree_test.go | 264 ++++++++++++++ internal/processmanager/shawl_windows.go | 24 +- 9 files changed, 1251 insertions(+), 9 deletions(-) create mode 100644 internal/processmanager/process_snapshot_windows.go create mode 100644 internal/processmanager/process_snapshot_windows_test.go create mode 100644 internal/processmanager/process_tree.go create mode 100644 internal/processmanager/process_tree_metrics.go create mode 100644 internal/processmanager/process_tree_metrics_test.go create mode 100644 internal/processmanager/process_tree_test.go diff --git a/internal/processmanager/README.md b/internal/processmanager/README.md index feac8c9..e594b4b 100644 --- a/internal/processmanager/README.md +++ b/internal/processmanager/README.md @@ -91,10 +91,11 @@ forwards to the panel via gRPC. | `docker` | yes | yes | yes | yes | yes (Linux) | yes | | `podman` | yes | yes | yes | yes | yes | yes | | `systemd` | yes | yes | yes | yes | yes | yes | -| `tmux` / `simple` / `winsw` / `shawl` | yes | — | — | — | — | — | +| `shawl` | yes | yes | usage only | — | yes (logical I/O) | yes (threads) | +| `tmux` / `simple` / `winsw` | yes | — | — | — | — | — | Container-backed managers tag their metrics with `{server_id, server_uuid, container}`. -The systemd manager tags its metrics with `{server_id, server_uuid, service}`. +The systemd and shawl managers tag their metrics with `{server_id, server_uuid, service}`. The systemd manager reads metrics from `systemctl show` and relies on the `CPUAccounting=yes`, `MemoryAccounting=yes`, `IOAccounting=yes`, @@ -105,7 +106,7 @@ CPU%) until the next start/restart regenerates the unit. Metrics are also suppressed for the first sample after each restart, since the cumulative CPU counter has no baseline yet. -PID-based stats for `tmux` / `simple` / `winsw` / `shawl` are tracked as a follow-up. +PID-based stats for `tmux` / `simple` / `winsw` are tracked as a follow-up. ## SystemD scopes @@ -208,7 +209,33 @@ and the start command's program is in neither the working directory nor PATH, th so before the log tail, naming both — the difference between an archive that unpacked into a subdirectory and a game that crashed on startup, which are fixed in entirely different places. -Metrics are liveness-only; see the table above. +### Metrics + +The metrics describe the processes shawl started for the server: the game server and anything it +runs through or starts, such as `cmd.exe` for a `.bat` or `.cmd` start command. shawl itself is not +counted, just as systemd and container runtimes stay out of a unit's or a container's accounting. +The shawl process is the one the service control manager reports for the service, and a single +`NtQuerySystemInformation(SystemProcessInformation)` call per metrics tick reads every process on +the host. No process has to be opened, so the account a server runs under does not matter. + +- `gameap_server_cpu_usage_percent` is a percentage of one core, as for docker and systemd. Task + Manager divides by the number of cores instead. +- `gameap_server_memory_usage_bytes` is the private working set, the "Memory" column of Task + Manager. Shared DLL pages are left out, so the sum over several processes is not inflated. A + service has no memory limit, so `gameap_server_memory_limit_bytes` and + `gameap_server_memory_usage_percent` are not reported. +- `gameap_server_block_io_*_bytes_total` count the bytes the processes read and wrote through files + and pipes, including reads served from the file cache: the I/O the game asked for, not what + reached the disk. +- `gameap_server_process_pids` is the number of threads, the unit `pids.current` counts on Linux. +- Windows keeps no per-process network counters, only ETW traces carry them, so no network metrics + are reported. + +A process counts as a child only when it is not older than its parent: Windows keeps the parent ID +of a process whose parent has exited and hands that ID to the next process that starts. CPU time and +I/O are tracked per process, so the counters keep growing when shawl restarts a crashed game, and +CPU is reported from the second sample after a start. What a process used between the last sample +and its exit is not counted. ### Changing the restart policy diff --git a/internal/processmanager/errors.go b/internal/processmanager/errors.go index b6a8d7f..b0fcc40 100644 --- a/internal/processmanager/errors.go +++ b/internal/processmanager/errors.go @@ -18,4 +18,7 @@ var ( ErrUserMismatch = errors.New( "server user does not match daemon user (required for systemctl --user mode)", ) + + ErrProcessSnapshotMalformed = errors.New("malformed process list") + ErrProcessSnapshotTooLarge = errors.New("process list does not fit the buffer") ) diff --git a/internal/processmanager/process_snapshot_windows.go b/internal/processmanager/process_snapshot_windows.go new file mode 100644 index 0000000..93ad33f --- /dev/null +++ b/internal/processmanager/process_snapshot_windows.go @@ -0,0 +1,158 @@ +//go:build windows + +package processmanager + +import ( + "context" + "sync" + "time" + "unsafe" + + "github.com/gameap/daemon/internal/app/config" + "github.com/gameap/daemon/internal/app/domain" + "github.com/gameap/daemon/pkg/logger" + "github.com/pkg/errors" + "golang.org/x/sys/windows" +) + +const ( + processSnapshotInitialSize = 512 * 1024 + processSnapshotMaxSize = 64 * 1024 * 1024 + processSnapshotAttempts = 5 + + // processSnapshotMaxAge lets every server of one metrics tick share a snapshot. It is half the + // shortest collection interval, so two ticks never do. + processSnapshotMaxAge = config.MetricsMinCollectionInterval / 2 +) + +var nativeProcessRecordLayout = processRecordLayout{ + size: int(unsafe.Sizeof(windows.SYSTEM_PROCESS_INFORMATION{})), + threads: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.NumberOfThreads)), + privateWorkingSet: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.WorkingSetPrivateSize)), + createTime: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.CreateTime)), + userTime: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.UserTime)), + kernelTime: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.KernelTime)), + pid: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.UniqueProcessID)), + parentPID: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.InheritedFromUniqueProcessID)), + readBytes: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.ReadTransferCount)), + writeBytes: int(unsafe.Offsetof(windows.SYSTEM_PROCESS_INFORMATION{}.WriteTransferCount)), +} + +// processSnapshotter reads the process list of the host. The list covers every process and thread +// on the host, hundreds of kilobytes, so one read serves every server of a metrics tick and the +// buffer is kept for the next tick. +type processSnapshotter struct { + mu sync.Mutex + buf []uint64 + procs []processInfo + takenAt time.Time +} + +// snapshot returns the process list and the time it was read. The list is shared between callers +// and must not be modified. +func (s *processSnapshotter) snapshot() ([]processInfo, time.Time, error) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.procs != nil && time.Since(s.takenAt) < processSnapshotMaxAge { + return s.procs, s.takenAt, nil + } + + procs, takenAt, err := s.read() + if err != nil { + return nil, time.Time{}, err + } + + s.procs, s.takenAt = procs, takenAt + + return procs, takenAt, nil +} + +// read queries the process list. The buffer is a []uint64 because the records carry 64-bit fields, +// which the kernel writes aligned. +func (s *processSnapshotter) read() ([]processInfo, time.Time, error) { + if len(s.buf) == 0 { + s.buf = make([]uint64, processSnapshotInitialSize/8) + } + + for range processSnapshotAttempts { + size := len(s.buf) * 8 + takenAt := time.Now() + + var written uint32 + + err := windows.NtQuerySystemInformation( + windows.SystemProcessInformation, unsafe.Pointer(&s.buf[0]), uint32(size), &written, + ) + if errors.Is(err, windows.STATUS_INFO_LENGTH_MISMATCH) { + // Processes start between two calls, so the buffer gets more room than the kernel + // asked for a moment ago. + grown := max(2*size, int(written)+int(written)/4) + if grown > processSnapshotMaxSize { + return nil, time.Time{}, errors.WithMessagef( + ErrProcessSnapshotTooLarge, "%d bytes needed, at most %d allowed", written, processSnapshotMaxSize, + ) + } + + s.buf = make([]uint64, (grown+7)/8) + + continue + } + if err != nil { + return nil, time.Time{}, errors.Wrap(err, "failed to query the process list") + } + + buf := unsafe.Slice((*byte)(unsafe.Pointer(&s.buf[0])), min(int(written), size)) + + procs, err := parseProcessSnapshot(buf, nativeProcessRecordLayout) + if err != nil { + return nil, time.Time{}, err + } + + return procs, takenAt, nil + } + + return nil, time.Time{}, errors.WithMessagef( + ErrProcessSnapshotTooLarge, "the process list kept growing over %d attempts", processSnapshotAttempts, + ) +} + +// serviceProcessMetrics reports what the processes a Windows service started are using: the game +// server and whatever it runs through, such as cmd.exe for a script. The service process itself is +// a supervisor, and systemd and container runtimes keep theirs out of a unit's or container's +// accounting as well. +// +// It returns nothing while the service has no process; the caller reports liveness on its own. +func serviceProcessMetrics( + ctx context.Context, serviceName string, snapshots *processSnapshotter, usage *processTreeSampler, +) []domain.Metric { + status, err := queryService(serviceName) + if err != nil || status.ProcessId == 0 { + usage.forget(serviceName) + + if err != nil && !errors.Is(err, ErrServiceNotFound) { + logger.WithError(ctx, err).Debug("Failed to query service " + serviceName + " for metrics") + } + + return nil + } + + procs, takenAt, err := snapshots.snapshot() + if err != nil { + logger.WithError(ctx, err).Debug("Failed to read the process list for metrics") + + return nil + } + + // The snapshot can be up to processSnapshotMaxAge older than the query. A service that started + // in between is not in it yet and is measured on the next tick; in the rare case that its + // process ID was still held by another process then, that process is reported for one tick. + tree, found := processDescendants(procs, status.ProcessId) + if !found { + usage.forget(serviceName) + + return nil + } + + return processTreeMetrics(takenAt, serviceName, usage.observe(serviceName, tree, takenAt)) +} diff --git a/internal/processmanager/process_snapshot_windows_test.go b/internal/processmanager/process_snapshot_windows_test.go new file mode 100644 index 0000000..b018d17 --- /dev/null +++ b/internal/processmanager/process_snapshot_windows_test.go @@ -0,0 +1,98 @@ +//go:build windows + +package processmanager + +import ( + "os" + "os/exec" + "testing" + "time" + "unsafe" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestNativeProcessRecordLayout(t *testing.T) { + want := processRecordLayout64 + if unsafe.Sizeof(uintptr(0)) == 4 { + want = processRecordLayout386 + } + + assert.Equal(t, want, nativeProcessRecordLayout) +} + +func TestProcessSnapshotterRead_ContainsCurrentProcess(t *testing.T) { + procs, takenAt, err := (&processSnapshotter{}).read() + require.NoError(t, err) + assert.WithinDuration(t, time.Now(), takenAt, time.Minute) + + var self *processInfo + + for i := range procs { + if procs[i].PID == uint32(os.Getpid()) { + self = &procs[i] + } + } + + require.NotNil(t, self, "the snapshot must contain the test process") + assert.Equal(t, uint32(os.Getppid()), self.ParentPID) + assert.Positive(t, self.Threads) + assert.Positive(t, self.PrivateWorkingSet) + assert.Positive(t, self.CPUTime) + assert.Less(t, self.CreateTime, filetimeTicks(time.Now())) +} + +func TestProcessSnapshotterRead_GrowsBuffer(t *testing.T) { + snapshots := &processSnapshotter{buf: make([]uint64, 1)} + + procs, _, err := snapshots.read() + + require.NoError(t, err) + assert.NotEmpty(t, procs) + assert.Greater(t, len(snapshots.buf), 1) +} + +func TestProcessSnapshotterSnapshot_SharesRecentRead(t *testing.T) { + snapshots := &processSnapshotter{} + + _, first, err := snapshots.snapshot() + require.NoError(t, err) + + _, second, err := snapshots.snapshot() + require.NoError(t, err) + + assert.Equal(t, first, second) +} + +func TestProcessDescendants_FindsChildAndGrandchild(t *testing.T) { + cmd := exec.Command("cmd.exe", "/c", "ping", "-n", "3", "127.0.0.1") + require.NoError(t, cmd.Start()) + + t.Cleanup(func() { + _ = cmd.Wait() + }) + + cmdPID := uint32(cmd.Process.Pid) + + assert.Eventually(t, func() bool { + procs, _, err := (&processSnapshotter{}).read() + if err != nil { + return false + } + + tree, found := processDescendants(procs, uint32(os.Getpid())) + if !found { + return false + } + + var hasCmd, hasPing bool + + for _, p := range tree { + hasCmd = hasCmd || p.PID == cmdPID + hasPing = hasPing || p.ParentPID == cmdPID + } + + return hasCmd && hasPing + }, 5*time.Second, 100*time.Millisecond, "cmd.exe and the ping it runs must be in the test process's tree") +} diff --git a/internal/processmanager/process_tree.go b/internal/processmanager/process_tree.go new file mode 100644 index 0000000..3c3bcc1 --- /dev/null +++ b/internal/processmanager/process_tree.go @@ -0,0 +1,135 @@ +package processmanager + +import ( + "encoding/binary" + + "github.com/pkg/errors" +) + +// processInfo is one process as a system-wide snapshot reports it. Times are FILETIME ticks of +// 100ns; CreateTime counts them from 1601-01-01 UTC, CPUTime is user and kernel time together. +type processInfo struct { + PID uint32 + ParentPID uint32 + CreateTime int64 + CPUTime int64 + Threads uint32 + PrivateWorkingSet uint64 + ReadBytes uint64 + WriteBytes uint64 +} + +// processRecordLayout is the size of one SYSTEM_PROCESS_INFORMATION record and the offsets of the +// fields the daemon reads. The fields after ImageName move with the pointer size, so the Windows +// build takes every offset from x/sys instead of keeping a table of its own. +// +// Records are read field by field rather than cast to the x/sys struct: the struct has pointer +// fields that the kernel fills with values the Go runtime must never treat as pointers. +type processRecordLayout struct { + size int + threads int + privateWorkingSet int + createTime int + userTime int + kernelTime int + pid int + parentPID int + readBytes int + writeBytes int +} + +// parseProcessSnapshot reads the records NtQuerySystemInformation(SystemProcessInformation) wrote to +// buf. Every record starts with the distance to the next one, zero on the last, and the thread +// records between two process records are skipped with that distance. +func parseProcessSnapshot(buf []byte, layout processRecordLayout) ([]processInfo, error) { + procs := make([]processInfo, 0, 256) + + for offset := 0; ; { + if len(buf)-offset < layout.size { + return nil, errors.WithMessagef( + ErrProcessSnapshotMalformed, + "record at offset %d overruns the %d bytes returned", offset, len(buf), + ) + } + + record := buf[offset : offset+layout.size] + + // The IDs sit in HANDLE-sized fields but are DWORDs, and on a little-endian machine the low + // 32 bits come first whatever the pointer size. + procs = append(procs, processInfo{ + PID: binary.LittleEndian.Uint32(record[layout.pid:]), + ParentPID: binary.LittleEndian.Uint32(record[layout.parentPID:]), + CreateTime: int64(binary.LittleEndian.Uint64(record[layout.createTime:])), + CPUTime: int64(binary.LittleEndian.Uint64(record[layout.userTime:])) + + int64(binary.LittleEndian.Uint64(record[layout.kernelTime:])), + Threads: binary.LittleEndian.Uint32(record[layout.threads:]), + PrivateWorkingSet: binary.LittleEndian.Uint64(record[layout.privateWorkingSet:]), + ReadBytes: binary.LittleEndian.Uint64(record[layout.readBytes:]), + WriteBytes: binary.LittleEndian.Uint64(record[layout.writeBytes:]), + }) + + next := uint64(binary.LittleEndian.Uint32(record)) + if next == 0 { + return procs, nil + } + + if next < uint64(layout.size) || next > uint64(len(buf)-offset) { + return nil, errors.WithMessagef( + ErrProcessSnapshotMalformed, + "record at offset %d puts the next one %d bytes on, in a %d-byte buffer", offset, next, len(buf), + ) + } + + offset += int(next) + } +} + +// processDescendants returns every process that rootPID started, directly or through its children, +// and reports whether rootPID itself is in the snapshot. The root is not part of the result. +// +// A process keeps the ID of the parent that created it for its whole life, while Windows hands a +// freed ID to the next process that starts. A process that names the parent's ID but is older than +// the parent was created by an earlier holder of that ID, so it is skipped with everything it +// started. Equal times are accepted because CreateTime comes from a clock that advances in steps of +// several milliseconds. A wall clock set back by more than a parent's age would make a child started +// afterwards look older than its parent; such a child is skipped too. +func processDescendants(procs []processInfo, rootPID uint32) ([]processInfo, bool) { + root := -1 + children := make(map[uint32][]int, len(procs)) + + for i := range procs { + if procs[i].PID == rootPID { + root = i + } + + children[procs[i].ParentPID] = append(children[procs[i].ParentPID], i) + } + + if root < 0 { + return nil, false + } + + var tree []processInfo + + visited := map[uint32]struct{}{rootPID: {}} + queue := []int{root} + + for len(queue) > 0 { + parent := procs[queue[0]] + queue = queue[1:] + + for _, i := range children[parent.PID] { + child := procs[i] + + if _, seen := visited[child.PID]; seen || child.CreateTime < parent.CreateTime { + continue + } + + visited[child.PID] = struct{}{} + tree = append(tree, child) + queue = append(queue, i) + } + } + + return tree, true +} diff --git a/internal/processmanager/process_tree_metrics.go b/internal/processmanager/process_tree_metrics.go new file mode 100644 index 0000000..c1ea649 --- /dev/null +++ b/internal/processmanager/process_tree_metrics.go @@ -0,0 +1,217 @@ +package processmanager + +import ( + "sync" + "time" + + "github.com/gameap/daemon/internal/app/domain" +) + +const ( + filetimeTick = 100 * time.Nanosecond + + // filetimeUnixEpoch is the Unix epoch in FILETIME ticks, which count from 1601-01-01 UTC. + filetimeUnixEpoch = 116444736000000000 +) + +// filetimeTicks converts t to the scale process creation times use. +func filetimeTicks(t time.Time) int64 { + return t.UnixNano()/int64(filetimeTick) + filetimeUnixEpoch +} + +// processKey identifies one process for its whole life. The ID alone does not: Windows hands a +// freed ID to the next process that starts. +type processKey struct { + pid uint32 + createTime int64 +} + +type processCounters struct { + cpuTime int64 + readBytes uint64 + writeBytes uint64 +} + +// processTreeSample is what one observation of a service's processes leaves for the next one: the +// counters of every process, and the I/O totals reported so far. +type processTreeSample struct { + at time.Time + counters map[processKey]processCounters + readBytes uint64 + writeBytes uint64 +} + +// processTreeUsage is what the processes of a service used, in the units the metrics report. +type processTreeUsage struct { + CPUPercent float64 + HasCPU bool + PrivateWorkingSet uint64 + ReadBytes uint64 + WriteBytes uint64 + Threads uint64 +} + +// computeProcessTreeUsage measures tree against the previous observation of the same service and +// returns the sample the next observation is measured against. +// +// CPU time and I/O are counted per process, so a process that exits or starts between two +// observations neither pulls the totals down nor counts twice: +// - a process seen last time contributes what it used since then; +// - a process started after the last observation contributes everything it used, all of which +// falls into the interval; +// - any other process has no baseline yet and contributes from the next observation on. +// +// The I/O totals therefore only grow while the service keeps running, starting from what its +// processes report when it is first observed. CPU needs an interval, so the first observation +// reports none. What a process used between the last observation and its exit is lost. +func computeProcessTreeUsage( + prior *processTreeSample, tree []processInfo, at time.Time, +) (processTreeUsage, processTreeSample) { + usage := processTreeUsage{} + next := processTreeSample{ + at: at, + counters: make(map[processKey]processCounters, len(tree)), + } + + var priorTicks int64 + if prior != nil { + priorTicks = filetimeTicks(prior.at) + } + + var cpuTicks int64 + var readBytes, writeBytes uint64 + + for _, p := range tree { + key := processKey{pid: p.PID, createTime: p.CreateTime} + current := processCounters{cpuTime: p.CPUTime, readBytes: p.ReadBytes, writeBytes: p.WriteBytes} + next.counters[key] = current + + usage.PrivateWorkingSet += p.PrivateWorkingSet + usage.Threads += uint64(p.Threads) + + if prior == nil { + readBytes += current.readBytes + writeBytes += current.writeBytes + + continue + } + + before, seen := prior.counters[key] + + switch { + case seen: + cpuTicks += max(current.cpuTime-before.cpuTime, 0) + readBytes += counterGrowth(before.readBytes, current.readBytes) + writeBytes += counterGrowth(before.writeBytes, current.writeBytes) + case p.CreateTime >= priorTicks: + cpuTicks += current.cpuTime + readBytes += current.readBytes + writeBytes += current.writeBytes + } + } + + next.readBytes, next.writeBytes = readBytes, writeBytes + + if prior != nil { + next.readBytes += prior.readBytes + next.writeBytes += prior.writeBytes + + if wall := at.Sub(prior.at); wall > 0 { + usage.CPUPercent = float64(time.Duration(cpuTicks)*filetimeTick) / float64(wall) * 100 + usage.HasCPU = true + } + } + + usage.ReadBytes, usage.WriteBytes = next.readBytes, next.writeBytes + + return usage, next +} + +func counterGrowth(before, after uint64) uint64 { + if after < before { + return 0 + } + + return after - before +} + +// processTreeSampler keeps the last observation of every service, so that each observation is +// measured against the one before it. +type processTreeSampler struct { + mu sync.Mutex + samples map[string]processTreeSample +} + +func newProcessTreeSampler() *processTreeSampler { + return &processTreeSampler{samples: make(map[string]processTreeSample)} +} + +// observe measures the processes of a service and keeps them as the baseline of the next observation. +func (s *processTreeSampler) observe(service string, tree []processInfo, at time.Time) processTreeUsage { + s.mu.Lock() + defer s.mu.Unlock() + + var prior *processTreeSample + if sample, ok := s.samples[service]; ok { + prior = &sample + } + + usage, next := computeProcessTreeUsage(prior, tree, at) + s.samples[service] = next + + return usage +} + +// forget drops the baseline of a service whose processes are gone, so the service is measured from +// scratch once it runs again instead of against processes that no longer exist. +func (s *processTreeSampler) forget(service string) { + s.mu.Lock() + defer s.mu.Unlock() + + delete(s.samples, service) +} + +// processTreeMetrics reports the usage of a service's processes as per-server metrics. There is no +// memory limit to compare against, and Windows keeps no per-process network counters. +func processTreeMetrics(ts time.Time, service string, usage processTreeUsage) []domain.Metric { + metric := func( + name string, metricType domain.MetricType, unit domain.MetricUnit, value domain.MetricValue, + ) domain.Metric { + return domain.Metric{ + Name: name, + Type: metricType, + Unit: unit, + Labels: map[string]string{metricLabelService: service}, + Timestamp: ts, + Value: value, + } + } + + out := make([]domain.Metric, 0, 5) + + if usage.HasCPU { + out = append(out, metric( + metricServerCPUUsagePercent, domain.MetricTypeGauge, domain.MetricUnitPercent, + domain.Float64Value(usage.CPUPercent), + )) + } + + return append(out, + metric( + metricServerMemoryUsageBytes, domain.MetricTypeGauge, domain.MetricUnitBytes, + domain.Uint64Value(usage.PrivateWorkingSet), + ), + metric( + metricServerBlockIOReadBytesTotal, domain.MetricTypeCounter, domain.MetricUnitBytes, + domain.Uint64Value(usage.ReadBytes), + ), + metric( + metricServerBlockIOWriteBytesTotal, domain.MetricTypeCounter, domain.MetricUnitBytes, + domain.Uint64Value(usage.WriteBytes), + ), + metric( + metricServerProcessPIDs, domain.MetricTypeGauge, domain.MetricUnitCount, + domain.Uint64Value(usage.Threads), + ), + ) +} diff --git a/internal/processmanager/process_tree_metrics_test.go b/internal/processmanager/process_tree_metrics_test.go new file mode 100644 index 0000000..643c6ce --- /dev/null +++ b/internal/processmanager/process_tree_metrics_test.go @@ -0,0 +1,326 @@ +package processmanager + +import ( + "strconv" + "sync" + "testing" + "time" + + "github.com/gameap/daemon/internal/app/domain" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +var usageBase = time.Date(2026, 9, 23, 12, 0, 0, 0, time.UTC) + +func cpuTicks(d time.Duration) int64 { + return int64(d / filetimeTick) +} + +// startedAt is the creation time of a process started d after usageBase. +func startedAt(d time.Duration) int64 { + return filetimeTicks(usageBase.Add(d)) +} + +func TestFiletimeTicks(t *testing.T) { + assert.Equal(t, int64(116444736000000000), filetimeTicks(time.Unix(0, 0))) + assert.Equal(t, int64(116444736000000000+10_000_000), filetimeTicks(time.Unix(1, 0))) +} + +func TestComputeProcessTreeUsage_FirstObservation(t *testing.T) { + tree := []processInfo{ + { + PID: 3100, CreateTime: startedAt(-time.Hour), CPUTime: cpuTicks(5 * time.Second), Threads: 1, + PrivateWorkingSet: 2 << 20, ReadBytes: 100, WriteBytes: 10, + }, + { + PID: 3180, CreateTime: startedAt(-time.Hour), CPUTime: cpuTicks(90 * time.Second), Threads: 24, + PrivateWorkingSet: 700 << 20, ReadBytes: 5000, WriteBytes: 800, + }, + } + + usage, sample := computeProcessTreeUsage(nil, tree, usageBase) + + assert.Equal(t, processTreeUsage{ + PrivateWorkingSet: 702 << 20, + ReadBytes: 5100, + WriteBytes: 810, + Threads: 25, + }, usage) + assert.Equal(t, usageBase, sample.at) + assert.Equal(t, uint64(5100), sample.readBytes) + assert.Equal(t, uint64(810), sample.writeBytes) + assert.Equal(t, map[processKey]processCounters{ + {pid: 3100, createTime: startedAt(-time.Hour)}: {cpuTime: cpuTicks(5 * time.Second), readBytes: 100, writeBytes: 10}, + {pid: 3180, createTime: startedAt(-time.Hour)}: {cpuTime: cpuTicks(90 * time.Second), readBytes: 5000, writeBytes: 800}, + }, sample.counters) +} + +func TestComputeProcessTreeUsage_AgainstPriorObservation(t *testing.T) { + game := processInfo{ + PID: 3180, CreateTime: startedAt(-time.Hour), CPUTime: cpuTicks(90 * time.Second), + Threads: 24, PrivateWorkingSet: 700 << 20, ReadBytes: 5000, WriteBytes: 800, + } + helper := processInfo{ + PID: 3300, CreateTime: startedAt(-time.Minute), CPUTime: cpuTicks(2 * time.Second), + Threads: 2, PrivateWorkingSet: 10 << 20, ReadBytes: 300, WriteBytes: 30, + } + + grown := func(p processInfo, cpu time.Duration, read, write uint64) processInfo { + p.CPUTime += cpuTicks(cpu) + p.ReadBytes += read + p.WriteBytes += write + + return p + } + + tests := []struct { + name string + prior []processInfo + current []processInfo + wall time.Duration + wantCPU float64 + wantRead uint64 + wantWrite uint64 + }{ + { + name: "idle", + prior: []processInfo{game}, + current: []processInfo{game}, + wall: time.Second, + wantCPU: 0, + wantRead: 5000, + wantWrite: 800, + }, + { + name: "one_core", + prior: []processInfo{game}, + current: []processInfo{grown(game, 5*time.Second, 64, 16)}, + wall: 5 * time.Second, + wantCPU: 100, + wantRead: 5064, + wantWrite: 816, + }, + { + name: "several_cores", + prior: []processInfo{game, helper}, + current: []processInfo{grown(game, 2*time.Second, 0, 0), grown(helper, time.Second, 0, 0)}, + wall: time.Second, + wantCPU: 300, + wantRead: 5300, + wantWrite: 830, + }, + { + name: "process_started_after_the_prior_observation_counts_in_full", + prior: []processInfo{game}, + current: []processInfo{game, { + PID: 3400, CreateTime: startedAt(500 * time.Millisecond), CPUTime: cpuTicks(500 * time.Millisecond), + ReadBytes: 64, WriteBytes: 8, + }}, + wall: time.Second, + wantCPU: 50, + wantRead: 5064, + wantWrite: 808, + }, + { + name: "process_older_than_the_prior_observation_has_no_baseline", + prior: []processInfo{game}, + current: []processInfo{game, helper}, + wall: time.Second, + wantCPU: 0, + wantRead: 5000, + wantWrite: 800, + }, + { + name: "reused_process_id_counts_as_a_new_process", + prior: []processInfo{game}, + current: []processInfo{{ + PID: game.PID, CreateTime: startedAt(750 * time.Millisecond), CPUTime: cpuTicks(250 * time.Millisecond), + ReadBytes: 10, WriteBytes: 1, + }}, + wall: time.Second, + wantCPU: 25, + wantRead: 5010, + wantWrite: 801, + }, + { + name: "counters_going_back_are_clamped", + prior: []processInfo{game}, + current: []processInfo{{ + PID: game.PID, CreateTime: game.CreateTime, CPUTime: game.CPUTime - cpuTicks(time.Second), + ReadBytes: game.ReadBytes - 100, WriteBytes: game.WriteBytes - 10, + }}, + wall: time.Second, + wantCPU: 0, + wantRead: 5000, + wantWrite: 800, + }, + { + name: "exited_process_keeps_its_share_of_the_io_totals", + prior: []processInfo{game, helper}, + current: []processInfo{grown(game, time.Second, 100, 10)}, + wall: 2 * time.Second, + wantCPU: 50, + wantRead: 5400, + wantWrite: 840, + }, + { + name: "service_restarted_between_observations", + prior: []processInfo{game, helper}, + current: []processInfo{{ + PID: 4100, CreateTime: startedAt(2 * time.Second), CPUTime: cpuTicks(time.Second), + ReadBytes: 50, WriteBytes: 5, + }}, + wall: 4 * time.Second, + wantCPU: 25, + wantRead: 5350, + wantWrite: 835, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, prior := computeProcessTreeUsage(nil, tt.prior, usageBase) + + usage, next := computeProcessTreeUsage(&prior, tt.current, usageBase.Add(tt.wall)) + + require.True(t, usage.HasCPU) + assert.InDelta(t, tt.wantCPU, usage.CPUPercent, 0.001) + assert.Equal(t, tt.wantRead, usage.ReadBytes) + assert.Equal(t, tt.wantWrite, usage.WriteBytes) + assert.Equal(t, tt.wantRead, next.readBytes) + assert.Equal(t, tt.wantWrite, next.writeBytes) + assert.Len(t, next.counters, len(tt.current)) + }) + } +} + +func TestComputeProcessTreeUsage_NoIntervalReportsNoCPU(t *testing.T) { + tree := []processInfo{{PID: 3180, CreateTime: startedAt(-time.Hour), CPUTime: cpuTicks(time.Minute)}} + + _, prior := computeProcessTreeUsage(nil, tree, usageBase) + + sameInstant, _ := computeProcessTreeUsage(&prior, tree, usageBase) + assert.False(t, sameInstant.HasCPU) + + clockSetBack, _ := computeProcessTreeUsage(&prior, tree, usageBase.Add(-time.Second)) + assert.False(t, clockSetBack.HasCPU) +} + +func TestProcessTreeSampler_MeasuresEachServiceAgainstItsOwnPast(t *testing.T) { + sampler := newProcessTreeSampler() + + game := processInfo{PID: 3180, CreateTime: startedAt(-time.Hour), CPUTime: cpuTicks(time.Minute)} + other := processInfo{PID: 5120, CreateTime: startedAt(-time.Hour), CPUTime: cpuTicks(time.Hour)} + + assert.False(t, sampler.observe("gameapServer1", []processInfo{game}, usageBase).HasCPU) + assert.False(t, sampler.observe("gameapServer2", []processInfo{other}, usageBase).HasCPU) + + game.CPUTime += cpuTicks(time.Second) + usage := sampler.observe("gameapServer1", []processInfo{game}, usageBase.Add(2*time.Second)) + + require.True(t, usage.HasCPU) + assert.InDelta(t, 50.0, usage.CPUPercent, 0.001) +} + +func TestProcessTreeSampler_ForgetStartsOver(t *testing.T) { + sampler := newProcessTreeSampler() + + game := processInfo{PID: 3180, CreateTime: startedAt(-time.Hour), ReadBytes: 100, WriteBytes: 10} + sampler.observe("gameapServer3", []processInfo{game}, usageBase) + + sampler.forget("gameapServer3") + + restarted := processInfo{PID: 4100, CreateTime: startedAt(time.Second), ReadBytes: 40, WriteBytes: 4} + usage := sampler.observe("gameapServer3", []processInfo{restarted}, usageBase.Add(2*time.Second)) + + assert.False(t, usage.HasCPU) + assert.Equal(t, uint64(40), usage.ReadBytes) + assert.Equal(t, uint64(4), usage.WriteBytes) +} + +func TestProcessTreeSampler_ConcurrentObservations(t *testing.T) { + sampler := newProcessTreeSampler() + + var wg sync.WaitGroup + + for worker := range 8 { + wg.Go(func() { + service := "gameapServer" + strconv.Itoa(worker%3) + + for step := range 50 { + tree := []processInfo{{ + PID: uint32(3000 + worker), + CreateTime: startedAt(-time.Hour), + CPUTime: cpuTicks(time.Duration(step) * time.Millisecond), + }} + sampler.observe(service, tree, usageBase.Add(time.Duration(step)*time.Second)) + + if step%10 == 0 { + sampler.forget(service) + } + } + }) + } + + wg.Wait() + + assert.NotPanics(t, func() { + sampler.observe("gameapServer1", nil, usageBase.Add(time.Hour)) + }) +} + +func TestProcessTreeMetrics(t *testing.T) { + got := processTreeMetrics(usageBase, "gameapServer5", processTreeUsage{ + CPUPercent: 142.5, + HasCPU: true, + PrivateWorkingSet: 700 << 20, + ReadBytes: 5300, + WriteBytes: 830, + Threads: 26, + }) + + want := []struct { + name string + metricType domain.MetricType + unit domain.MetricUnit + value domain.MetricValue + }{ + {metricServerCPUUsagePercent, domain.MetricTypeGauge, domain.MetricUnitPercent, domain.Float64Value(142.5)}, + {metricServerMemoryUsageBytes, domain.MetricTypeGauge, domain.MetricUnitBytes, domain.Uint64Value(700 << 20)}, + {metricServerBlockIOReadBytesTotal, domain.MetricTypeCounter, domain.MetricUnitBytes, domain.Uint64Value(5300)}, + {metricServerBlockIOWriteBytesTotal, domain.MetricTypeCounter, domain.MetricUnitBytes, domain.Uint64Value(830)}, + {metricServerProcessPIDs, domain.MetricTypeGauge, domain.MetricUnitCount, domain.Uint64Value(26)}, + } + + require.Len(t, got, len(want)) + + for i, w := range want { + assert.Equal(t, w.name, got[i].Name) + assert.Equal(t, w.metricType, got[i].Type, w.name) + assert.Equal(t, w.unit, got[i].Unit, w.name) + assert.Equal(t, w.value, got[i].Value, w.name) + assert.Equal(t, usageBase, got[i].Timestamp, w.name) + assert.Equal(t, map[string]string{metricLabelService: "gameapServer5"}, got[i].Labels, w.name) + } + + // The metrics collector adds the server labels to every metric in place. + got[0].Labels["server_id"] = "5" + assert.NotContains(t, got[1].Labels, "server_id") +} + +func TestProcessTreeMetrics_WithoutCPU(t *testing.T) { + got := processTreeMetrics(usageBase.Add(time.Minute), "gameapServer9", processTreeUsage{PrivateWorkingSet: 1 << 20}) + + require.Len(t, got, 4) + assert.Empty(t, collectByName(got, metricServerCPUUsagePercent)) + + for _, name := range []string{ + metricServerMemoryLimitBytes, + metricServerMemoryUsagePercent, + metricServerNetworkReceiveBytesTotal, + metricServerNetworkTransmitBytesTotal, + } { + assert.Empty(t, collectByName(got, name), name) + } +} diff --git a/internal/processmanager/process_tree_test.go b/internal/processmanager/process_tree_test.go new file mode 100644 index 0000000..9146d23 --- /dev/null +++ b/internal/processmanager/process_tree_test.go @@ -0,0 +1,264 @@ +package processmanager + +import ( + "bytes" + "encoding/binary" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The layouts of SYSTEM_PROCESS_INFORMATION that x/sys has on 64-bit Windows and on 386. +var ( + processRecordLayout64 = processRecordLayout{ + size: 256, threads: 4, privateWorkingSet: 8, createTime: 32, userTime: 40, kernelTime: 48, + pid: 80, parentPID: 88, readBytes: 232, writeBytes: 240, + } + processRecordLayout386 = processRecordLayout{ + size: 184, threads: 4, privateWorkingSet: 8, createTime: 32, userTime: 40, kernelTime: 48, + pid: 68, parentPID: 72, readBytes: 160, writeBytes: 168, + } +) + +// systemThreadInformationSize is the size of the thread records the kernel writes after every +// process record. +const systemThreadInformationSize = 80 + +// encodeProcessSnapshot lays procs out the way NtQuerySystemInformation does: every process record +// is followed by one record per thread, and the last one has no next offset. Bytes the parser must +// not read are filled with garbage, and CPU time is split between the user and kernel fields. +func encodeProcessSnapshot(layout processRecordLayout, procs ...processInfo) []byte { + pointerSize := layout.parentPID - layout.pid + + putID := func(b []byte, id uint32) { + if pointerSize == 8 { + binary.LittleEndian.PutUint64(b, uint64(id)) + + return + } + + binary.LittleEndian.PutUint32(b, id) + } + + var buf []byte + + for i, p := range procs { + record := bytes.Repeat([]byte{0xAB}, layout.size+int(p.Threads)*systemThreadInformationSize) + + next := uint32(0) + if i < len(procs)-1 { + next = uint32(len(record)) + } + + userTime := p.CPUTime / 3 + + binary.LittleEndian.PutUint32(record, next) + binary.LittleEndian.PutUint32(record[layout.threads:], p.Threads) + binary.LittleEndian.PutUint64(record[layout.privateWorkingSet:], p.PrivateWorkingSet) + binary.LittleEndian.PutUint64(record[layout.createTime:], uint64(p.CreateTime)) + binary.LittleEndian.PutUint64(record[layout.userTime:], uint64(userTime)) + binary.LittleEndian.PutUint64(record[layout.kernelTime:], uint64(p.CPUTime-userTime)) + putID(record[layout.pid:], p.PID) + putID(record[layout.parentPID:], p.ParentPID) + binary.LittleEndian.PutUint64(record[layout.readBytes:], p.ReadBytes) + binary.LittleEndian.PutUint64(record[layout.writeBytes:], p.WriteBytes) + + buf = append(buf, record...) + } + + return buf +} + +func TestParseProcessSnapshot(t *testing.T) { + system := processInfo{ + PID: 4, CreateTime: 133_700_000_000_000_000, CPUTime: 91_234_567, Threads: 3, + PrivateWorkingSet: 196 << 10, ReadBytes: 1 << 20, WriteBytes: 3 << 20, + } + shawl := processInfo{ + PID: 7112, ParentPID: 640, CreateTime: 133_700_000_100_000_000, CPUTime: 1_500_000, Threads: 5, + PrivateWorkingSet: 2 << 20, ReadBytes: 40_960, WriteBytes: 81_920, + } + game := processInfo{ + PID: 7340, ParentPID: 7112, CreateTime: 133_700_000_100_000_001, CPUTime: 7_200_000_000, Threads: 41, + PrivateWorkingSet: 1 << 30, ReadBytes: 5 << 30, WriteBytes: 12_345, + } + + layouts := []struct { + name string + layout processRecordLayout + }{ + {name: "64_bit", layout: processRecordLayout64}, + {name: "386", layout: processRecordLayout386}, + } + + tests := []struct { + name string + procs []processInfo + }{ + {name: "one_record", procs: []processInfo{system}}, + {name: "several_records", procs: []processInfo{system, shawl, game}}, + } + + for _, l := range layouts { + for _, tt := range tests { + t.Run(l.name+"_"+tt.name, func(t *testing.T) { + got, err := parseProcessSnapshot(encodeProcessSnapshot(l.layout, tt.procs...), l.layout) + + require.NoError(t, err) + assert.Equal(t, tt.procs, got) + }) + } + } +} + +func TestParseProcessSnapshot_Malformed(t *testing.T) { + layout := processRecordLayout64 + + // 416 bytes for the first process and its two threads, 336 for the second and its one thread. + valid := encodeProcessSnapshot(layout, + processInfo{PID: 4, Threads: 2}, + processInfo{PID: 88, ParentPID: 4, Threads: 1}, + ) + require.Len(t, valid, 752) + + withFirstNext := func(next uint32) []byte { + buf := bytes.Clone(valid) + binary.LittleEndian.PutUint32(buf, next) + + return buf + } + + tests := []struct { + name string + buf []byte + wantError string + }{ + { + name: "empty_buffer", + buf: nil, + wantError: "record at offset 0 overruns the 0 bytes returned", + }, + { + name: "first_record_cut_short", + buf: valid[:layout.size-1], + wantError: "record at offset 0 overruns the 255 bytes returned", + }, + { + name: "last_record_cut_short", + buf: valid[:416+layout.size-1], + wantError: "record at offset 416 overruns the 671 bytes returned", + }, + { + name: "next_record_inside_this_one", + buf: withFirstNext(uint32(layout.size - 8)), + wantError: "record at offset 0 puts the next one 248 bytes on, in a 752-byte buffer", + }, + { + name: "next_record_past_the_end", + buf: withFirstNext(760), + wantError: "record at offset 0 puts the next one 760 bytes on, in a 752-byte buffer", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := parseProcessSnapshot(tt.buf, layout) + + require.ErrorIs(t, err, ErrProcessSnapshotMalformed) + assert.Contains(t, err.Error(), tt.wantError) + assert.Nil(t, got) + }) + } +} + +func TestProcessDescendants(t *testing.T) { + hour := int64(time.Hour / filetimeTick) + + services := processInfo{PID: 640, ParentPID: 520, CreateTime: 1 * hour} + shawl := processInfo{PID: 1200, ParentPID: 640, CreateTime: 10 * hour} + cmd := processInfo{PID: 3100, ParentPID: 1200, CreateTime: 10*hour + 5} + game := processInfo{PID: 3180, ParentPID: 3100, CreateTime: 10*hour + 9} + conhost := processInfo{PID: 3192, ParentPID: 3180, CreateTime: 10*hour + 9} + otherShawl := processInfo{PID: 1300, ParentPID: 640, CreateTime: 11 * hour} + otherGame := processInfo{PID: 4400, ParentPID: 1300, CreateTime: 11*hour + 3} + + // Started by an earlier process that held the ID shawl has now. + orphan := processInfo{PID: 2900, ParentPID: 1200, CreateTime: 9 * hour} + orphanChild := processInfo{PID: 2950, ParentPID: 2900, CreateTime: 10*hour + 7} + + idle := processInfo{PID: 0, ParentPID: 0, CreateTime: 0} + system := processInfo{PID: 4, ParentPID: 0, CreateTime: 0} + smss := processInfo{PID: 380, ParentPID: 4, CreateTime: 1} + + tests := []struct { + name string + procs []processInfo + rootPID uint32 + want []processInfo + wantFound bool + }{ + { + name: "root_without_children", + procs: []processInfo{services, shawl, otherShawl}, + rootPID: shawl.PID, + want: nil, + wantFound: true, + }, + { + name: "children_and_grandchildren_but_nothing_else", + procs: []processInfo{services, otherGame, conhost, game, shawl, otherShawl, cmd}, + rootPID: shawl.PID, + want: []processInfo{cmd, game, conhost}, + wantFound: true, + }, + { + name: "root_missing", + procs: []processInfo{services, cmd, game}, + rootPID: shawl.PID, + want: nil, + wantFound: false, + }, + { + name: "orphan_under_a_reused_id_is_skipped_with_its_children", + procs: []processInfo{services, shawl, orphan, orphanChild, cmd, game}, + rootPID: shawl.PID, + want: []processInfo{cmd, game}, + wantFound: true, + }, + { + name: "child_started_in_the_same_clock_step_as_its_parent", + procs: []processInfo{shawl, {PID: 3500, ParentPID: shawl.PID, CreateTime: shawl.CreateTime}}, + rootPID: shawl.PID, + want: []processInfo{{PID: 3500, ParentPID: shawl.PID, CreateTime: shawl.CreateTime}}, + wantFound: true, + }, + { + name: "root_that_is_its_own_parent", + procs: []processInfo{idle, system, smss}, + rootPID: idle.PID, + want: []processInfo{system, smss}, + wantFound: true, + }, + { + name: "parents_that_point_at_each_other", + procs: []processInfo{ + {PID: 10, ParentPID: 20, CreateTime: 5}, + {PID: 20, ParentPID: 10, CreateTime: 5}, + }, + rootPID: 10, + want: []processInfo{{PID: 20, ParentPID: 10, CreateTime: 5}}, + wantFound: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, found := processDescendants(tt.procs, tt.rootPID) + + assert.Equal(t, tt.wantFound, found) + assert.Equal(t, tt.want, got) + }) + } +} diff --git a/internal/processmanager/shawl_windows.go b/internal/processmanager/shawl_windows.go index 1be3fb1..946e279 100644 --- a/internal/processmanager/shawl_windows.go +++ b/internal/processmanager/shawl_windows.go @@ -35,12 +35,19 @@ const ( type Shawl struct { cfg *config.Config + + processes *processSnapshotter + usage *processTreeSampler } // NewShawl builds the Windows process manager. It drives the service control manager through // its API rather than sc.exe, so it needs no executor. func NewShawl(cfg *config.Config, _, _ contracts.Executor) *Shawl { - return &Shawl{cfg: cfg} + return &Shawl{ + cfg: cfg, + processes: &processSnapshotter{}, + usage: newProcessTreeSampler(), + } } func (pm *Shawl) Install(ctx context.Context, server *domain.Server, out io.Writer) (domain.Result, error) { @@ -72,6 +79,8 @@ func (pm *Shawl) Uninstall(ctx context.Context, server *domain.Server, out io.Wr return domain.ErrorResult, errors.WithMessage(err, "failed to delete service") } + pm.usage.forget(serviceName) + configFile := pm.configFile(server) if err := os.Remove(configFile); err != nil && !errors.Is(err, os.ErrNotExist) { logger.WithError(ctx, err).Warn("failed to remove service config file") @@ -822,8 +831,13 @@ func (pm *Shawl) HasOwnInstallation(_ *domain.Server) bool { return false } -// Metrics returns only the cached process-active gauge for now. Resource -// stats via Windows performance counters / WMI are tracked as a follow-up. -func (pm *Shawl) Metrics(_ context.Context, server *domain.Server) ([]domain.Metric, error) { - return []domain.Metric{livenessMetric(server, time.Now())}, nil +// Metrics reports the liveness gauge and what the processes shawl started for the server use. +// shawl itself is not counted; see serviceProcessMetrics. +func (pm *Shawl) Metrics(ctx context.Context, server *domain.Server) ([]domain.Metric, error) { + processMetrics := serviceProcessMetrics(ctx, pm.serviceName(server), pm.processes, pm.usage) + + out := make([]domain.Metric, 0, len(processMetrics)+1) + out = append(out, livenessMetric(server, time.Now())) + + return append(out, processMetrics...), nil }