diff mbox series

[bitbake-devel,5/6] hashserv: server: Use streaming and queue API for upstream exist queries

Message ID 20260724212822.1165552-6-JPEWhacker@gmail.com
State New
Headers show
Series hashserv: Pipeline Upstream Queries | expand

Commit Message

Joshua Watt July 24, 2026, 9:25 p.m. UTC
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 <JPEWhacker@gmail.com>
---
 lib/hashserv/server.py | 31 ++++++++++++++++++++++++++-----
 lib/hashserv/tests.py  |  5 +++++
 2 files changed, 31 insertions(+), 5 deletions(-)
diff mbox series

Patch

diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py
index d392ac19f..5017320c1 100644
--- a/lib/hashserv/server.py
+++ b/lib/hashserv/server.py
@@ -526,17 +526,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 asyncio.gather(
+                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 3acfdcd8c..0fbc19c1d 100644
--- a/lib/hashserv/tests.py
+++ b/lib/hashserv/tests.py
@@ -372,20 +372,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