[bitbake-devel][PATCH v2 08/10] hashserv: server: Use streaming and queue API for upstream exist queries
Joshua Watt <[email protected]> Thu, 30 Jul 2026 12:31:01 -0600
| Newsgroups | org.openembedded.lists.bitbake-devel |
|---|---|
| Message-ID | <[email protected]> |
Reworks the "unihash-exists" handler to use the new server queue API and client streaming API to efficiently stream requests to the upstream server instead of having to wait for a roundtrip on the requests. Signed-off-by: Joshua Watt <[email protected]> --- lib/hashserv/server.py | 31 ++++++++++++++++++++++++++----- lib/hashserv/tests.py | 5 +++++ 2 files changed, 31 insertions(+), 5 deletions(-) diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py index d0e6f23fc..0153730fa 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -522,17 +522,38 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): @permissions(READ_PERM) async def handle_exists_stream(self, request): - async def handler(l): + async def exists_handler(l): if await self.db.unihash_exists(l): return "true" + return "false" + + if not self.upstream_client: + return await self._stream_handler(exists_handler) - if self.upstream_client is not None: - if await self.upstream_client.unihash_exists(l): + async with self.upstream_client.unihash_exists_stream() as stream: + + async def get_local_result(m): + if await self.db.unihash_exists(m): return "true" + return None - return "false" + async def get_upstream_result(m): + exists = await stream.get_result() + return "true" if exists else "false" - return await self._stream_handler(handler) + queue = asyncio.Queue() + upstream = UpstreamQueue( + queue, + get_local_result, + stream.send_query, + get_upstream_result, + ) + + await bb.asyncrpc.TaskGroup.run( + self._stream_queue_handler(upstream.handler, queue), + upstream.process_results(), + ) + return self.NO_RESPONSE async def report_readonly(self, data): method = data["method"] diff --git a/lib/hashserv/tests.py b/lib/hashserv/tests.py index e24bdcacb..a7ce7425e 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -374,20 +374,25 @@ class HashEquivalenceCommonTests(object): nonlocal side_client # check upstream server + self.assertTrue(self.client.unihash_exists(unihash)) self.assertClientGetHash(self.client, taskhash, unihash) # Hash should *not* be present on the side server + if old_sidehash and unihash != old_sidehash: + self.assertFalse(side_client.unihash_exists(unihash)) self.assertClientGetHash(side_client, taskhash, old_sidehash) # Hash should be present on the downstream server, since it # will defer to the upstream server. This will trigger # the backfill in the downstream server + self.assertTrue(down_client.unihash_exists(unihash)) self.assertClientGetHash(down_client, taskhash, unihash) # After waiting for the downstream client to finish backfilling the # task from the upstream server, it should appear in the side server # since the database is populated down_client.backfill_wait() + self.assertTrue(side_client.unihash_exists(unihash)) self.assertClientGetHash(side_client, taskhash, unihash) # Basic report -- 2.54.0