Skip to content

Commit 388dd52

Browse files
committed
Compute network byte deltas in ProcessGroupStatsMonitor
`_read_network_bytes()` reads cumulative counters from `/proc/net/dev` which include all traffic since the interface was brought up. The `net_rx_bytes` and `net_tx_bytes` fields in `ProcessGroupResourceUsage` were previously set to these raw cumulative values, making it difficult to interpret how much network I/O occurred during each monitoring interval. This change computes deltas between consecutive snapshots, matching the pattern already used for `cpu_percent`. On the first snapshot, `net_rx_bytes` and `net_tx_bytes` are `None` (no previous value to diff against). The raw cumulative values are threaded through `_collect_pgrp_stats` as internal state so subsequent calls can compute the difference.
1 parent 4dbfa40 commit 388dd52

2 files changed

Lines changed: 68 additions & 19 deletions

File tree

src/spdl/pipeline/_pgrp_stats.py

Lines changed: 49 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -437,16 +437,26 @@ class ProcessGroupResourceUsage:
437437
"""Number of processes in the process group."""
438438

439439
net_rx_bytes: int | None = None
440-
"""Total network bytes received (excluding loopback)."""
440+
"""Network bytes received since the previous snapshot (excluding loopback).
441+
442+
``None`` on the first snapshot (no previous value to diff against).
443+
Note: this is host-wide (per network namespace), not process-group-scoped.
444+
"""
441445

442446
net_tx_bytes: int | None = None
443-
"""Total network bytes transmitted (excluding loopback)."""
447+
"""Network bytes transmitted since the previous snapshot (excluding loopback).
448+
449+
``None`` on the first snapshot (no previous value to diff against).
450+
Note: this is host-wide (per network namespace), not process-group-scoped.
451+
"""
444452

445453

446454
def _collect_pgrp_stats(
447455
prev_cpu_usec: int | None = None,
448456
prev_time_usec: int | None = None,
449-
) -> tuple[ProcessGroupResourceUsage, int | None, int]:
457+
prev_net_rx_bytes: int | None = None,
458+
prev_net_tx_bytes: int | None = None,
459+
) -> tuple[ProcessGroupResourceUsage, int | None, int, int | None, int | None]:
450460
"""Collect all process-group stats.
451461
452462
Reads CPU, memory (RSS, PSS, private), disk IO from
@@ -459,16 +469,24 @@ def _collect_pgrp_stats(
459469
(used to compute ``cpu_percent``). ``None`` on the first call.
460470
prev_time_usec: Wall-clock µs (``time.monotonic()`` × 1e6) of the
461471
previous snapshot.
472+
prev_net_rx_bytes: Cumulative network RX bytes from the previous
473+
snapshot. ``None`` on the first call.
474+
prev_net_tx_bytes: Cumulative network TX bytes from the previous
475+
snapshot. ``None`` on the first call.
462476
463477
Returns:
464-
A tuple of ``(usage, current_cpu_usec, current_time_usec)``.
478+
A tuple of
479+
``(usage, current_cpu_usec, current_time_usec, current_net_rx_bytes, current_net_tx_bytes)``.
465480
``current_cpu_usec`` is ``None`` when /proc could not be read.
481+
``current_net_*_bytes`` are ``None`` when /proc/net/dev could not be read.
466482
"""
467483
import time as _time
468484

469485
now_usec = int(_time.monotonic() * 1_000_000)
470486
result = ProcessGroupResourceUsage(pid=os.getpid(), pgid=os.getpgrp())
471487
current_cpu_usec: int | None = None
488+
current_net_rx_bytes: int | None = None
489+
current_net_tx_bytes: int | None = None
472490

473491
try:
474492
pgrp = _read_pgrp_stats()
@@ -492,12 +510,21 @@ def _collect_pgrp_stats(
492510

493511
try:
494512
net = _read_network_bytes()
495-
result.net_rx_bytes = net.rx_bytes
496-
result.net_tx_bytes = net.tx_bytes
513+
current_net_rx_bytes = net.rx_bytes
514+
current_net_tx_bytes = net.tx_bytes
515+
if prev_net_rx_bytes is not None and prev_net_tx_bytes is not None:
516+
result.net_rx_bytes = net.rx_bytes - prev_net_rx_bytes
517+
result.net_tx_bytes = net.tx_bytes - prev_net_tx_bytes
497518
except RuntimeError as e:
498519
_warn_once("net_dev", "%s", e)
499520

500-
return result, current_cpu_usec, now_usec
521+
return (
522+
result,
523+
current_cpu_usec,
524+
now_usec,
525+
current_net_rx_bytes,
526+
current_net_tx_bytes,
527+
)
501528

502529

503530
# ---------------------------------------------------------------------------
@@ -528,10 +555,23 @@ async def _loop() -> None:
528555
loop = asyncio.get_running_loop()
529556
prev_cpu_usec: int | None = None
530557
prev_time_usec: int | None = None
558+
prev_net_rx_bytes: int | None = None
559+
prev_net_tx_bytes: int | None = None
531560
while running:
532561
try:
533-
usage, prev_cpu_usec, prev_time_usec = await loop.run_in_executor(
534-
None, _collect_pgrp_stats, prev_cpu_usec, prev_time_usec
562+
(
563+
usage,
564+
prev_cpu_usec,
565+
prev_time_usec,
566+
prev_net_rx_bytes,
567+
prev_net_tx_bytes,
568+
) = await loop.run_in_executor(
569+
None,
570+
_collect_pgrp_stats,
571+
prev_cpu_usec,
572+
prev_time_usec,
573+
prev_net_rx_bytes,
574+
prev_net_tx_bytes,
535575
)
536576
await callback(usage)
537577
except Exception:

tests/pipeline/pgrp_stats_test.py

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -57,23 +57,32 @@ def test_read_pgrp_stats_returns_valid_data(self) -> None:
5757

5858
def test_collect_pgrp_stats_returns_complete_snapshot(self) -> None:
5959
"""_collect_pgrp_stats should return a fully populated snapshot."""
60-
result, cpu_usec, time_usec = _collect_pgrp_stats()
60+
result, cpu_usec, time_usec, net_rx, net_tx = _collect_pgrp_stats()
6161
self.assertIsInstance(result, ProcessGroupResourceUsage)
6262
self.assertEqual(result.pid, os.getpid())
6363
self.assertEqual(result.pgid, os.getpgrp())
6464

65-
# First call: cpu_percent should be None (no previous value).
65+
# First call: cpu_percent and net deltas should be None (no previous value).
6666
self.assertIsNone(result.cpu_percent)
6767
self.assertIsNotNone(cpu_usec)
6868
self.assertIsNotNone(result.rss_bytes)
6969
self.assertIsNotNone(result.num_procs)
70-
self.assertIsNotNone(result.net_rx_bytes)
71-
self.assertIsNotNone(result.net_tx_bytes)
72-
73-
# Second call with prev values: cpu_percent should be set.
74-
result2, _, _ = _collect_pgrp_stats(cpu_usec, time_usec)
75-
self.assertIsNotNone(result2.cpu_percent)
76-
self.assertGreaterEqual(result2.cpu_percent, 0.0)
70+
self.assertIsNone(result.net_rx_bytes)
71+
self.assertIsNone(result.net_tx_bytes)
72+
self.assertIsNotNone(net_rx)
73+
self.assertIsNotNone(net_tx)
74+
75+
# Second call with prev values: cpu_percent and net deltas should be set.
76+
result2, _, _, _, _ = _collect_pgrp_stats(cpu_usec, time_usec, net_rx, net_tx)
77+
cpu_pct = result2.cpu_percent
78+
assert cpu_pct is not None
79+
self.assertGreaterEqual(cpu_pct, 0.0)
80+
rx = result2.net_rx_bytes
81+
assert rx is not None
82+
self.assertGreaterEqual(rx, 0)
83+
tx = result2.net_tx_bytes
84+
assert tx is not None
85+
self.assertGreaterEqual(tx, 0)
7786

7887
# Sanity: at least one process (this one) should be counted.
7988
assert result.num_procs is not None
@@ -449,7 +458,7 @@ def test_subprocess_function_collects_and_calls_callback(
449458
net_rx_bytes=100,
450459
net_tx_bytes=200,
451460
)
452-
mock_collect.return_value = (usage, 500000, 1000000)
461+
mock_collect.return_value = (usage, 500000, 1000000, 100, 200)
453462

454463
mock_callback = AsyncMock()
455464

0 commit comments

Comments
 (0)