From: Joshua Watt <jpewhacker@gmail.com>
To: bitbake-devel@lists.openembedded.org
Cc: Michal Sieron <michal.sieron@nokia.com>,
Joshua Watt <JPEWhacker@gmail.com>
Subject: [bitbake-devel][PATCH 6/6] hashserv: tests: Add test for upstream pipelining
Date: Fri, 24 Jul 2026 15:25:14 -0600 [thread overview]
Message-ID: <20260724212822.1165552-7-JPEWhacker@gmail.com> (raw)
In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com>
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 <michal.sieron@nokia.com>
Signed-off-by: Joshua Watt <JPEWhacker@gmail.com>
---
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
next prev parent reply other threads:[~2026-07-24 21:28 UTC|newest]
Thread overview: 20+ messages / expand[flat|nested] mbox.gz Atom feed top
2026-07-24 21:25 [bitbake-devel][PATCH 0/6] hashserv: Pipeline Upstream Queries Joshua Watt
2026-07-24 21:25 ` [bitbake-devel][PATCH 1/6] hashserv: client: Add asynchronous streaming API Joshua Watt
2026-07-24 21:25 ` [bitbake-devel][PATCH 2/6] hashserv: server: Fix formatting Joshua Watt
2026-07-24 21:25 ` [bitbake-devel][PATCH 3/6] hashserv: server: Add queued streaming API Joshua Watt
2026-07-24 21:25 ` [bitbake-devel][PATCH 4/6] hashserv: server: Use streaming and queue API for upstream unihash queries Joshua Watt
2026-07-24 21:25 ` [bitbake-devel][PATCH 5/6] hashserv: server: Use streaming and queue API for upstream exist queries Joshua Watt
2026-07-24 21:25 ` Joshua Watt [this message]
2026-07-26 13:50 ` [bitbake-devel][PATCH 0/6] hashserv: Pipeline Upstream Queries Richard Purdie
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 00/10] " Joshua Watt
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 01/10] asyncrpc: Add Task Group Joshua Watt
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 02/10] asyncrpc: serv: Use " Joshua Watt
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 03/10] asyncrpc: serv: Cancel all clients on server stop Joshua Watt
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 04/10] hashserv: tests: Improve test logging Joshua Watt
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 05/10] hashserv: client: Add asynchronous streaming API Joshua Watt
2026-07-30 18:30 ` [bitbake-devel][PATCH v2 06/10] hashserv: server: Add queued " Joshua Watt
2026-07-30 18:31 ` [bitbake-devel][PATCH v2 07/10] hashserv: server: Use streaming and queue API for upstream unihash queries Joshua Watt
2026-07-30 18:31 ` [bitbake-devel][PATCH v2 08/10] hashserv: server: Use streaming and queue API for upstream exist queries Joshua Watt
2026-07-30 18:31 ` [bitbake-devel][PATCH v2 09/10] hashserv: tests: Add more upstream tests Joshua Watt
2026-07-30 18:31 ` [bitbake-devel][PATCH v2 10/10] hashserv: tests: Add test for upstream pipelining Joshua Watt
2026-07-31 13:05 ` [PATCH v2 00/10] hashserv: Pipeline Upstream Queries Michal Sieron
Reply instructions:
You may reply publicly to this message via plain-text email
using any one of the following methods:
* Save the following mbox file, import it into your mail client,
and reply-to-all from there: mbox
Avoid top-posting and favor interleaved quoting:
https://en.wikipedia.org/wiki/Posting_style#Interleaved_style
* Reply using the --to, --cc, and --in-reply-to
switches of git-send-email(1):
git send-email \
--in-reply-to=20260724212822.1165552-7-JPEWhacker@gmail.com \
--to=jpewhacker@gmail.com \
--cc=bitbake-devel@lists.openembedded.org \
--cc=michal.sieron@nokia.com \
/path/to/YOUR_REPLY
https://kernel.org/pub/software/scm/git/docs/git-send-email.html
* If your mail client supports setting the In-Reply-To header
via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line
before the message body.
This is a public inbox, see mirroring instructions
for how to clone and mirror all data and code used for this inbox