-
Notifications
You must be signed in to change notification settings - Fork 28
Expand file tree
/
Copy pathsocket_reactor_test.py
More file actions
165 lines (132 loc) · 6.32 KB
/
Copy pathsocket_reactor_test.py
File metadata and controls
165 lines (132 loc) · 6.32 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
"""SocketReactor Python test.
Exercises the espp.SocketReactor Python binding: lifecycle, UDP receiver
registration + echo round-trip, sender info, callbacks returning None (no
response), multiple receivers, dynamic remove, and input validation. A plain
Python UDP socket is used as the client.
Exit code 0 on full pass, 1 on any failure.
"""
import socket
import sys
import time
from typing import List, Optional, Tuple
import espp
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
results: List[Tuple[str, bool]] = []
def check(test: str, condition: bool, desc: str) -> bool:
ok = bool(condition)
print(f" {'PASS' if ok else 'FAIL'} [{test}]: {desc}")
results.append((test, ok))
return ok
def udp_send_recv(port: int, payload: bytes, timeout: float = 1.0) -> Optional[bytes]:
"""Send `payload` to 127.0.0.1:port and return the reply bytes (or None on timeout)."""
cli = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
cli.settimeout(timeout)
try:
cli.sendto(payload, ("127.0.0.1", port))
reply, _ = cli.recvfrom(4096)
return reply
except socket.timeout:
return None
finally:
cli.close()
def wait_until(predicate, timeout: float = 1.0, interval: float = 0.01) -> bool:
deadline = time.time() + timeout
while time.time() < deadline:
if predicate():
return True
time.sleep(interval)
return predicate()
VERB = espp.Logger.Verbosity.warn
# ---------------------------------------------------------------------------
# 1. Lifecycle: auto_start=False -> start() -> is_running() -> stop()
# ---------------------------------------------------------------------------
name = "lifecycle"
print(f"--- {name} ---")
reactor = espp.SocketReactor(worker_count=2, auto_start=False, log_level=VERB)
check(name, not reactor.is_running(), "not running before start()")
check(name, reactor.start(), "start() returns True")
check(name, reactor.is_running(), "running after start()")
check(name, reactor.num_registered() == 0, "no registrations initially")
# ---------------------------------------------------------------------------
# 2. UDP echo: a Python callback reverses the payload; the reactor sends it back
# ---------------------------------------------------------------------------
name = "udp echo"
print(f"--- {name} ---")
server_a = espp.UdpSocket(espp.UdpSocket.Config(VERB))
def reverse_cb(data, sender):
return bytes(reversed(data))
id_a = reactor.add_udp_receiver(server_a, 5111, 1500, reverse_cb)
check(name, id_a != espp.SocketReactor.INVALID_ID, "add_udp_receiver returns a valid id")
check(name, reactor.num_registered() == 1, "one registration")
payload = bytes(range(40))
reply = udp_send_recv(5111, payload)
check(name, reply == bytes(reversed(payload)), "reversed echo received via the reactor")
# ---------------------------------------------------------------------------
# 3. Sender info + a callback that returns None sends no response
# ---------------------------------------------------------------------------
name = "no response / sender info"
print(f"--- {name} ---")
server_b = espp.UdpSocket(espp.UdpSocket.Config(VERB))
seen = {}
def sink_cb(data, sender):
seen["address"] = sender.address
seen["length"] = len(data)
return None # -> reactor sends nothing back
id_b = reactor.add_udp_receiver(server_b, 5112, 1500, sink_cb)
check(name, id_b != espp.SocketReactor.INVALID_ID, "second receiver registered")
check(name, reactor.num_registered() == 2, "two registrations")
no_reply = udp_send_recv(5112, b"hello", timeout=0.4)
check(name, no_reply is None, "no reply when the callback returns None")
check(name, wait_until(lambda: seen.get("length") == 5), "callback observed the 5-byte payload")
check(name, seen.get("address") == "127.0.0.1", f"sender address seen: {seen.get('address')}")
# ---------------------------------------------------------------------------
# 4. Dynamic remove()
# ---------------------------------------------------------------------------
name = "remove"
print(f"--- {name} ---")
check(name, reactor.remove(id_a) is True, "remove() returns True for a valid id")
check(name, wait_until(lambda: reactor.num_registered() == 1), "count drops to 1 after remove()")
check(name, reactor.remove(999999) is False, "remove() of an unknown id returns False")
check(name, udp_send_recv(5111, b"x", timeout=0.3) is None, "removed receiver no longer echoes")
# ---------------------------------------------------------------------------
# 5. A raising callback must not crash the worker; the reactor keeps serving.
# ---------------------------------------------------------------------------
name = "callback exception"
print(f"--- {name} ---")
raiser = espp.UdpSocket(espp.UdpSocket.Config(VERB))
def raising_cb(data, sender):
raise ValueError("boom")
check(name, reactor.add_udp_receiver(raiser, 5113, 1500, raising_cb) != espp.SocketReactor.INVALID_ID,
"raising receiver registered")
# The datagram triggers the exception (reported via sys.unraisablehook); no reply,
# no crash.
check(name, udp_send_recv(5113, b"boom", timeout=0.4) is None, "raising callback yields no reply")
# The reactor is still alive: the earlier echo receiver on 5112's sibling keeps working.
echo_srv = espp.UdpSocket(espp.UdpSocket.Config(VERB))
reactor.add_udp_receiver(echo_srv, 5114, 1500, reverse_cb)
payload2 = bytes(range(24))
check(name, udp_send_recv(5114, payload2) == bytes(reversed(payload2)),
"reactor still serves other receivers after a callback exception")
# ---------------------------------------------------------------------------
# 6. Stop
# ---------------------------------------------------------------------------
name = "stop"
print(f"--- {name} ---")
reactor.stop()
check(name, not reactor.is_running(), "not running after stop()")
# ---------------------------------------------------------------------------
# Summary
# ---------------------------------------------------------------------------
passed = sum(1 for _, ok in results if ok)
total = len(results)
print()
print(f"======== SocketReactor python test: {passed}/{total} checks passed ========")
if passed != total:
for test_name, ok in results:
if not ok:
print(f" FAILED: {test_name}")
sys.exit(1)
print("ALL CHECKS PASSED")
sys.exit(0)