Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
104 changes: 104 additions & 0 deletions experiments/e52-latency-sweep/FINDINGS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
# E52: latency between env client and bridge changes the time, not the weights - up to 25 ms each way, SB3 ends on its in-process weights

2026-10-02 · guangzhao (bridges-venv: torch 2.14.1+cpu, SB3 2.9.0,
gymnasium 1.3.0; plugrl-bridges 79ec044) · protocol: [`PROTOCOL.md`](PROTOCOL.md)
(`ce807c4`, after the pilot in `results/pilot.txt`, before the registered
runs) · no amendment

---

## The result

SB3's PPO trained 16 Pendulum environments through the bridge, with the env
client connecting either directly or through `delay_proxy.py`. The proxy
holds every byte for D ms in each direction.

| seed | weights, every arm | direct | 0 ms | 1 ms | 5 ms | 25 ms |
| --- | --- | --- | --- | --- | --- | --- |
| 0 | `ec5b8226...` | 16.3 s | 16.8 s | 35.7 s | 92.1 s | 355.6 s |
| 1 | `062e92e3...` | 16.1 s | 17.2 s | 35.4 s | 92.7 s | 355.9 s |
| 2 | `3a8d968f...` | 16.2 s | 16.9 s | 35.4 s | 92.4 s | 356.1 s |

Per vector step (6,400 a run), beyond the direct run, and how far that is
above 2D (two one-way delays):

| D | per step | beyond 2D |
| --- | --- | --- |
| 0 ms | +0.07 to +0.17 ms | +0.07 to +0.17 |
| 1 ms | +3.00 to +3.04 ms | +1.00 to +1.04 |
| 5 ms | +11.85 to +11.96 ms | +1.85 to +1.96 |
| 25 ms | +53.02 to +53.10 ms | **+3.02 to +3.10** |

Every check and prediction holds except one:
- **V1:** all 15 runs are complete, and each bridge run served one
connection.
- **P1 holds.** Every run, at every delay, ends on E50's in-process weights
for its seed. A run with 25 ms each way between environment and trainer
takes 22 times as long and ends on the same bytes as a run with the
environments in SB3's own process.
- **P2 is falsified, at 25 ms only.** The band was 2D to 2D + 3 ms. At 1 and
5 ms every seed is inside it. At 25 ms all three seeds are above it by
0.02-0.10 ms.

---

## Where the excess comes from

The excess over 2D grows with D, from 0.1 ms to 3.1 ms. The pilot gave only
one point, 1.87 ms at 5 ms, and P2's band assumed the excess was a constant.
It is not.

After the runs, `proxy_rtt.py` measured the proxy alone. It sent 4 KiB
messages through `delay_proxy.py` to a TCP echo server, 300 round trips per
delay (`results/proxy_rtt.txt`). Beyond 2D, the median was:

| D | one message | two messages 0.2 ms apart |
| --- | --- | --- |
| 0 ms | 0.03 ms | 0.28 ms |
| 1 ms | 0.48 ms | 0.77 ms |
| 5 ms | 1.00 ms | 1.04 ms |
| 25 ms | 0.96 ms | 1.29 ms |

So the proxy itself accounts for about 1 ms of each step's excess. That is
its two `asyncio.sleep` calls, each rounded up in `epoll_wait`, as E47 found
for the server's scheduler. The rest grows with D:
- 0.2-0.5 ms at 1 ms;
- about 0.9 ms at 5 ms;
- about 1.8 ms at 25 ms.

It is **not explained** here. It lies between the proxy and the trainer: in
the WebSocket client, the bridge's event loop, or the order in which they
wake. A tighter claim would need a delay that does not sleep in Python, such
as `tc netem`, which needs root on guangzhao.

---

## What it means

- **The claim holds as stated: latency changes time, not data.** The client
blocks on every action, and the bridge is lockstep. So when an action
arrives decides when a step happens, never what it produces.
- **Its price is two one-way delays per step, plus a little.** On a 25 ms
link that is 53 ms a step, against 2.5 ms of computation. A lockstep
trainer on a long link spends nearly all its time waiting. That is E48's
6.4x slowdown across Wi-Fi, and the reason an asynchronous server (E43:
2x) or a deeper pipeline is the design for distant environments.
- **Measure latency without sleeping where you can.** The tool that adds
latency adds its own (E47).

---

## What E52 does not show

* **Real links.** The proxy adds latency but not bandwidth limits, jitter
or loss. E43 and E48 measured a real Wi-Fi link.
* **An environment that keeps moving.** The client waits for each action.
A robot does not, and latency then changes the data. That is out of
scope (the paper's limitations).
* **The excess's cause beyond the proxy's own share.** See above.

Kept here:
- each run's `result.json`, `monitor.csv`, `episodes.csv`, logs and proxy
logs;
- `registered.out`, `load.txt`, `verdicts.txt` (`summarise.py`);
- `proxy_rtt.txt` (`proxy_rtt.py`, run after the registered runs).
95 changes: 95 additions & 0 deletions experiments/e52-latency-sweep/PROTOCOL.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
# E52 measurement protocol (pre-registered)

**Written 2026-10-02, after the pilot in `results/pilot.txt`, before the
registered runs.**

This file must not be edited after the first registered data point. Anything
learned afterwards goes in `AMENDMENT.md`, dated.

---

## The question

The paper claims that latency between environment and trainer slows a run
without changing what it collects. That rests on a design argument: the env
client waits for every action. E43 adds indirect evidence: two arms on
different machines learned alike, but their physics differ in the last bit,
so their data could not be compared. So:

**On one machine, with latency added between the env client and the bridge,
does training end on the same weights? And does the time grow as the
protocol says it should, by two one-way delays per step?**

---

## Declared in advance: what was already known

1. **The code.**
- `train.py` is E50's, with one added option: `--client-port`, the port
the env client connects to.
- The client is E50's `gym_client.py`. The bridge is plugrl-bridges as
synced to guangzhao for E49 (`79ec044`).
- SB3 2.9.0 PPO with rl-zoo's Pendulum-v1 settings, 16 envs, 102,400
steps, one torch thread, CPU.
2. **`delay_proxy.py`** forwards each chunk D ms after it arrived, in each
direction. It does not limit bandwidth. It sleeps with `asyncio.sleep`,
which on Linux waits in `epoll_wait` and rounds up to the millisecond
(E47). A step is one action out and one feedback-and-infer back, so it
should take 2D longer.
3. **E50, on this machine.** The in-process arm ends on these weights:
- seed 0: `ec5b8226...`;
- seed 1: `062e92e3...`;
- seed 2: `3a8d968f...`.

E50's bridge arm ends on the same weights.
4. **The pilot, seed 9** (`results/pilot.txt`).
- Direct and 5 ms both ended on `3721266c...`, E50's seed-9 hash on
guangzhao.
- Direct took 16.1 s and 5 ms took 92.1 s. That is 11.87 ms more per
vector step (6,400 of them), against 2D = 10 ms.

---

## Design

On guangzhao, through `run.sh`, one run at a time:
- seeds 0, 1 and 2;
- for each, the arms direct (no proxy, as E50), then 0, 1, 5 and 25 ms one
way.

That is 15 runs.

---

## Checks

* **V1 - complete.** Every run has `steps` 102,400 and `client_rc` 0, and its
log has one `plugrl-bridges: slots 0-15` line.

---

## Predictions, and what falsifies each

**P1 - latency leaves the weights alone.** All 15 runs end on E50's
in-process weights for their seed.

> Grounds: known items 3 and 4. The client blocks on every action, so the
> data do not depend on when an action arrives. The bridge is lockstep, so
> neither does the order of the data.

**P2 - and costs two one-way delays per step.** For D = 1, 5 and 25 ms,
the time per vector step beyond the direct run's lies between 2D and
2D + 3 ms, on every seed.

> Grounds: known items 2 and 4. The pilot's excess over 2D was 1.87 ms at
> 5 ms: the proxy's own forwarding, plus the sleep rounding up on each leg.

**Reported, not predicted:**
- the 0 ms arm's time beyond the direct run's, the proxy's own cost;
- each arm's excess over 2D.

---

## Reading order

V1, P1, P2, then the reported figures.
66 changes: 66 additions & 0 deletions experiments/e52-latency-sweep/delay_proxy.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
"""A TCP proxy that holds every byte for D ms in each direction.

python delay_proxy.py --listen 8821 --upstream 8820 --delay-ms 5

A client that connects to --listen reaches 127.0.0.1:--upstream, with each
chunk it sends, and each chunk it receives, forwarded D ms after it arrived,
in order. A request and its answer therefore take 2 D longer, which is what a
link with a one-way latency of D does to a round trip. Bandwidth is not
limited. With --delay-ms 0 it only forwards, which measures its own cost.
"""

from __future__ import annotations

import argparse
import asyncio
import time


async def pump(reader, writer, delay: float) -> None:
queue: asyncio.Queue = asyncio.Queue()

async def read() -> None:
while data := await reader.read(1 << 20):
queue.put_nowait((time.monotonic() + delay, data))
queue.put_nowait(None)

async def write() -> None:
while (item := await queue.get()) is not None:
due, data = item
wait = due - time.monotonic()
if wait > 0:
await asyncio.sleep(wait)
writer.write(data)
await writer.drain()
writer.close()

await asyncio.gather(read(), write())


async def main() -> None:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--listen", type=int, required=True)
p.add_argument("--upstream", type=int, required=True)
p.add_argument("--delay-ms", type=float, required=True)
a = p.parse_args()
delay = a.delay_ms / 1000.0

async def handle(client_reader, client_writer) -> None:
up_reader, up_writer = await asyncio.open_connection("127.0.0.1", a.upstream)
await asyncio.gather(
pump(client_reader, up_writer, delay),
pump(up_reader, client_writer, delay),
return_exceptions=True,
)

server = await asyncio.start_server(handle, "127.0.0.1", a.listen)
print(
f"delaying 127.0.0.1:{a.listen} -> :{a.upstream} by {a.delay_ms} ms each way",
flush=True,
)
async with server:
await server.serve_forever()


if __name__ == "__main__":
asyncio.run(main())
97 changes: 97 additions & 0 deletions experiments/e52-latency-sweep/proxy_rtt.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
"""E52, after the registered runs: what delay_proxy.py itself adds to a round trip.

python proxy_rtt.py [--delays 0 1 5 25] [--rounds 300]

For each delay D it starts delay_proxy.py in front of a TCP echo server, sends
a 4 KiB message through it and waits for the echo, `--rounds` times, and
reports the median and 90th percentile round trip minus 2D. No bridge, no
trainer, no WebSocket: only the proxy and the loopback.

With --burst 2 each round sends two messages 0.2 ms apart, as E52's client
sends a feedback and then the next infer, and waits for both echoes.
"""

from __future__ import annotations

import argparse
import pathlib
import socket
import statistics
import subprocess
import sys
import threading
import time

HERE = pathlib.Path(__file__).resolve().parent
ECHO, PROXY = 8841, 8842
PAYLOAD = b"x" * 4096


def echo_server(stop: threading.Event) -> None:
with socket.create_server(("127.0.0.1", ECHO)) as srv:
srv.settimeout(0.5)
while not stop.is_set():
try:
conn, _ = srv.accept()
except TimeoutError:
continue
with conn:
while data := conn.recv(1 << 16):
conn.sendall(data)


def recv_exactly(sock: socket.socket, n: int) -> None:
got = 0
while got < n:
chunk = sock.recv(n - got)
if not chunk:
raise ConnectionError("closed")
got += len(chunk)


def main() -> int:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--delays", type=float, nargs="+", default=[0, 1, 5, 25])
p.add_argument("--rounds", type=int, default=300)
p.add_argument("--burst", type=int, default=1)
a = p.parse_args()
stop = threading.Event()
threading.Thread(target=echo_server, args=(stop,), daemon=True).start()
time.sleep(0.3)
for d in a.delays:
proxy = subprocess.Popen(
[sys.executable, str(HERE / "delay_proxy.py"), "--listen", str(PROXY),
"--upstream", str(ECHO), "--delay-ms", str(d)],
stdout=subprocess.DEVNULL,
) # fmt: skip
time.sleep(0.5)
try:
with socket.create_connection(("127.0.0.1", PROXY)) as s:
s.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
rtts = []
for _ in range(a.rounds):
start = time.perf_counter()
for k in range(a.burst):
if k:
spin = time.perf_counter() + 0.0002
while time.perf_counter() < spin:
pass
s.sendall(PAYLOAD)
recv_exactly(s, a.burst * len(PAYLOAD))
rtts.append((time.perf_counter() - start) * 1000)
finally:
proxy.terminate()
proxy.wait()
excess = sorted(r - 2 * d for r in rtts)
print(
f"burst {a.burst}, D {d:5.1f} ms: round trip median {statistics.median(rtts):.3f} ms, "
f"beyond 2D median {statistics.median(excess):.3f}, "
f"p90 {excess[int(0.9 * len(excess))]:.3f} ms",
flush=True,
)
stop.set()
return 0


if __name__ == "__main__":
raise SystemExit(main())
2 changes: 2 additions & 0 deletions experiments/e52-latency-sweep/results/load.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
09:21:59 up 27 days, 22:23, 1 user, load average: 0.11, 0.18, 0.14
09:48:23 up 27 days, 22:50, 1 user, load average: 0.48, 0.65, 0.62
2 changes: 2 additions & 0 deletions experiments/e52-latency-sweep/results/pendulum-0-seed0.log
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
plugrl-bridges: slots 0-15 <- ('127.0.0.1', 38370)
{"env": "pendulum", "arm": "bridge", "seed": 0, "weights_sha256": "ec5b8226b227f3fa58bd771b411e3673b751ec2c062b010cf0dcef60baed42bc", "learn_s": 16.76, "steps": 102400, "steps_per_s": 6109.6, "torch": "2.14.1+cpu", "client_rc": 0}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
delaying 127.0.0.1:8821 -> :8820 by 0.0 ms each way
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
closed after 6400 steps: 'plugrl-server-stop'
Loading
Loading