forgetest: events.stream

The acceptance test for forgectrl's GET /events. The daemon counts event
streams per peer address, so the test is several peers at once: it binds its
client sockets to 127.0.0.2 through 127.0.0.5, which are all this host. Three
streams must each get 200, text/event-stream, and the hello event; a fourth
address must get 503 with the reason while GET /settings still answers; a
second stream from the first address must be served, and the older one must
get bye and the end of its response; POST /motion/release and
/motion/energize must arrive on a stream that was open all along as
motors.released and motors.energized, in order, with ids that count up; and
with every stream closed a place must come back within 25 s, since a closed
client is only noticed at the daemon's next write to it (its keep-alive).

The SSE client was run against a host-built daemon before any bench time was
spent, which found two defects in it: its reads blocked for the socket
timeout, not the window the caller asked for, and a replaced stream ends
with the last chunk of the response, not with the socket closing.

Passes on the bench reference: the fourth stream refused in words, the
replacement told and ended, alarm, motors.released, and motors.energized as
ids 1 to 3, and a place back 3 s after every stream closed. It covers
forgectrl's src/events.*, src/main.c, src/status.*, and src/grblport.*.
This commit is contained in:
ScottW514
2026-09-20 07:14:04 -04:00
parent 8856fb0099
commit 25cc42b879
+139
View File
@@ -8,6 +8,7 @@ import json
import os import os
import socket import socket
import time import time
import urllib.parse
from ..catalog import test from ..catalog import test
from .. import hw from .. import hw
@@ -417,3 +418,141 @@ def panel_serves(ctx):
st, text = fc.get("/curve/ladder.gcode", raw=True) st, text = fc.get("/curve/ladder.gcode", raw=True)
ctx.check(st == 200 and b"S1000" in (text or b"") and b"M5" in (text or b""), ctx.check(st == 200 and b"S1000" in (text or b"") and b"M5" in (text or b""),
"GET /curve/ladder.gcode -> %s without the ladder", st) "GET /curve/ladder.gcode -> %s without the ladder", st)
class _EventStream:
"""GET /events over a raw socket bound to a chosen loopback source
address. Every address in 127.0.0.0/8 is this host, so the test can be
several peers at once, and the daemon counts streams per peer address."""
def __init__(self, base, source, timeout=5.0):
u = urllib.parse.urlsplit(base)
self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.sock.bind((source, 0))
self.sock.settimeout(timeout)
self.sock.connect((u.hostname or "127.0.0.1", u.port or 80))
self.sock.sendall(b"GET /events HTTP/1.1\r\nHost: %s\r\nAccept: text/event-stream\r\n"
b"Connection: close\r\n\r\n" % (u.netloc or "127.0.0.1").encode())
self.buf = b""
self.eof = False # the response is over: the last chunk, or the socket closed
head = self._until(b"\r\n\r\n")
self.status = int(head.split(b" ", 2)[1]) if head.startswith(b"HTTP/") else 0
self.head = head.decode("latin-1", "replace")
def _fill(self):
try:
chunk = self.sock.recv(4096)
except socket.timeout:
return False
if not chunk:
self.eof = True
return False
self.buf += chunk
return True
def _until(self, mark):
while mark not in self.buf:
if not self._fill():
break
i = self.buf.find(mark)
if i < 0:
out, self.buf = self.buf, b""
return out
out, self.buf = self.buf[:i], self.buf[i + len(mark):]
return out
def text(self, seconds):
"""Everything the stream says in the next `seconds` (chunk framing
and all: the checks look for event names inside it)."""
end = time.time() + seconds
self.sock.settimeout(0.2) # the window is the caller's, not one long read
while time.time() < end and not self.eof:
self._fill()
if self.buf.endswith(b"\r\n0\r\n\r\n"):
self.eof = True
out, self.buf = self.buf, b""
return out.decode("utf-8", "replace")
def close(self):
try:
self.sock.close()
except OSError:
pass
@test("events.stream", title="The event stream: edges arrive, and the cap holds",
subsystem="forgectrl", kind="auto", mode="grbl", est_min=2,
covers=[("forgectrl", "src/events.*"), ("forgectrl", "src/main.c"), ("forgectrl", "src/status.*"),
("forgectrl", "src/grblport.*")],
description="GET /events from three loopback source addresses: each gets 200, "
"text/event-stream, and the hello event. A fourth address gets 503 with the "
"reason, and GET /settings still answers while the three are held. A second "
"stream from the first address is served, and the older one gets bye and the "
"end of its response. An edge arrives on an open stream: POST /motion/release "
"and /motion/energize show up as motors.released and motors.energized, with "
"ids that count up (X and Y lose their reference, as after any release). With "
"every stream closed, the places come back within two keep-alive intervals.")
def events_stream(ctx):
ev = ctx.evidence
fc = hw.Forgectrl()
streams = []
try:
for src in ("127.0.0.2", "127.0.0.3", "127.0.0.4"):
s = _EventStream(fc.base, src)
streams.append(s)
hello = s.text(1.0)
ctx.check(s.status == 200 and "text/event-stream" in s.head and "event: hello" in hello,
"GET /events from %s -> %s, %r", src, s.status, hello[:120])
fourth = _EventStream(fc.base, "127.0.0.5")
body = fourth.text(1.0)
fourth.close()
ev["fourth"] = [fourth.status, body.strip()[-120:]]
ctx.log("the fourth stream -> %s %s", fourth.status, ev["fourth"][1])
ctx.check(fourth.status == 503 and "every event stream is taken" in body,
"the fourth stream -> %s %r, expected 503 and the reason", fourth.status, body[:160])
st, _body = fc.get("/settings")
ctx.check(st == 200, "GET /settings -> %s with three event streams held", st)
# One per address, by replacement.
newer = _EventStream(fc.base, "127.0.0.2")
hello = newer.text(1.0)
ctx.check(newer.status == 200 and "event: hello" in hello,
"a second stream from one address -> %s %r", newer.status, hello[:120])
bye = streams[0].text(2.0)
ev["replaced"] = {"bye": "event: bye" in bye, "ended": streams[0].eof}
ctx.check("event: bye" in bye and "replaced" in bye, "the older stream was not told: %r", bye[-160:])
ctx.check(streams[0].eof, "the older stream was not ended")
streams[0].close()
streams[0] = newer
# An edge, on a stream that has been open all along.
st, body = fc.post("/motion/release")
ctx.check(st == 200, "POST /motion/release -> %s %s", st, body)
ctx.sleep(1.0)
st, body = fc.post("/motion/energize")
ctx.check(st == 200, "POST /motion/energize -> %s %s", st, body)
text = streams[1].text(2.0)
ev["edges"] = [l for l in text.splitlines() if l.startswith(("id:", "event:"))]
ctx.log("the stream carried: %s", ev["edges"])
i_rel, i_en = text.find("event: motors.released"), text.find("event: motors.energized")
ctx.check(0 <= i_rel < i_en, "the release and the energize did not arrive in order: %r", text[-300:])
ids = [int(l.split(":")[1]) for l in text.splitlines() if l.startswith("id:")]
ctx.check(ids and ids == sorted(ids) and len(set(ids)) == len(ids), "the ids do not count up: %s", ids)
finally:
for s in streams:
s.close()
if fc.status().get("motors_released"):
fc.post("/motion/energize")
# A closed client is noticed at the daemon's next write to it.
def reopened():
s = _EventStream(fc.base, "127.0.0.5")
try:
return s.status == 200
finally:
s.close()
took = ctx.wait_for(reopened, 25, poll=1.0)
ev["place_back_s"] = took
ctx.check(took is not None, "no stream could be opened 25 s after every stream was closed")
ctx.log("PASS: three streams, the fourth refused in words with the settings route answering, a "
"replacement told and ended, two edges in order, and a place back after %.0f s", took)