[bitbake-devel][PATCH 6/6] hashserv: tests: Add test for upstream pipelining
Joshua Watt <[email protected]> Fri, 24 Jul 2026 15:25:14 -0600
| Newsgroups | org.openembedded.lists.bitbake-devel |
|---|---|
| Message-ID | <[email protected]> |
Adds a test that verifies that pipelining to an upstream server works as expected. AI-Generated: Uses Cursor, Claude Sonnet 5 Co-authored-by: Michal Sieron <[email protected]> Signed-off-by: Joshua Watt <[email protected]> --- lib/hashserv/tests.py | 96 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 96 insertions(+) diff --git a/lib/hashserv/tests.py b/lib/hashserv/tests.py index 0fbc19c1d..41c1f249f 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -1561,6 +1561,102 @@ class TestHashEquivalenceTCPServer(HashEquivalenceTestSetup, HashEquivalenceComm # case it is more reliable to resolve the IP address explicitly. return socket.gethostbyname("localhost") + ":0" + def test_get_stream_upstream_pipelined(self): + # Verify that, with an upstream configured, the get-stream handler + # pipelines its upstream queries instead of doing one blocking + # round-trip per task. A latency proxy injects RTT between this server + # and its upstream; a serial (one round-trip per task) implementation + # could not beat N * RTT, so finishing far faster proves the queries + # are pipelined. + import asyncio + + LATENCY = 0.01 # 10ms each direction => ~20ms round trip + upstream_host, upstream_port = self.server_address.rsplit(":", 1) + upstream_port = int(upstream_port) + + ready = threading.Event() + proxy = {} + + def run_proxy(): + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + + async def handle(creader, cwriter): + ureader, uwriter = await asyncio.open_connection(upstream_host, upstream_port) + + async def pipe(r, w): + try: + while True: + data = await r.read(65536) + if not data: + break + await asyncio.sleep(LATENCY) + w.write(data) + await w.drain() + except Exception: + pass + finally: + try: + w.close() + except Exception: + pass + + await asyncio.gather(pipe(creader, uwriter), pipe(ureader, cwriter)) + + stop_event = asyncio.Event() + proxy["stop_event"] = stop_event + proxy["loop"] = loop + + async def main(): + server = await asyncio.start_server(handle, "127.0.0.1", 0) + proxy["port"] = server.sockets[0].getsockname()[1] + ready.set() + async with server: + await stop_event.wait() + + try: + loop.run_until_complete(main()) + except Exception: + pass + finally: + loop.close() + + proxy_thread = threading.Thread(target=run_proxy, daemon=True) + proxy_thread.start() + self.assertTrue(ready.wait(10), "latency proxy did not start") + self.addCleanup(proxy_thread.join, 10) + self.addCleanup(lambda: proxy["loop"].call_soon_threadsafe(proxy["stop_event"].set)) + proxy_addr = "127.0.0.1:%d" % proxy["port"] + + N = 400 + expected = [] + for i in range(N): + taskhash = hashlib.sha256(("task%d" % i).encode()).hexdigest() + outhash = hashlib.sha256(("out%d" % i).encode()).hexdigest() + unihash = hashlib.sha256(("uni%d" % i).encode()).hexdigest() + self.client.report_unihash(taskhash, self.METHOD, outhash, unihash) + expected.append((self.METHOD, taskhash, unihash)) + + # Downstream server with an EMPTY local DB whose upstream is the slow proxy. + down_server = self.start_server(upstream=proxy_addr) + down_client = self.start_client(down_server.address) + + args = [(m, th) for (m, th, _uh) in expected] + start = time.time() + results = down_client.get_unihash_batch(args) + elapsed = time.time() - start + + self.assertEqual(results, [uh for (_m, _th, uh) in expected]) + + # A serial (one round-trip per task) implementation cannot beat N*RTT. + # Allow a generous margin to avoid flakiness on loaded CI machines. + serial_lower_bound = N * 2 * LATENCY + self.assertLess(elapsed, serial_lower_bound / 4, + "Upstream queries are not being pipelined " + "(%.3fs for %d queries at %.0fms injected RTT)" + % (elapsed, N, 2 * LATENCY * 1000)) + + class TestHashEquivalenceWebsocketServer(HashEquivalenceTestSetup, HashEquivalenceCommonTests, unittest.TestCase): def setUp(self): -- 2.54.0