From patchwork Fri Jul 24 21:25:09 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93475 Return-Path: X-Spam-Checker-Version: SpamAssassin 3.4.0 (2014-02-07) on aws-us-west-2-korg-lkml-1.web.codeaurora.org Received: from aws-us-west-2-korg-lkml-1.web.codeaurora.org (localhost.localdomain [127.0.0.1]) by smtp.lore.kernel.org (Postfix) with ESMTP id 1AD8DC53209 for ; Fri, 24 Jul 2026 21:28:31 +0000 (UTC) Received: from mail-ot1-f50.google.com (mail-ot1-f50.google.com [209.85.210.50]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.28918.1784928507495863074 for ; Fri, 24 Jul 2026 14:28:27 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=DuQGag9S; spf=pass (domain: gmail.com, ip: 209.85.210.50, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f50.google.com with SMTP id 46e09a7af769-7e6b5737bb2so798349a34.1 for ; Fri, 24 Jul 2026 14:28:27 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1784928507; x=1785533307; darn=lists.openembedded.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=KkvlTppzlY9z98B+XfnNU0JHQBJqdVgCIejfzQAzX74=; b=DuQGag9SEF+rPkro4cYmfu/kGrcIo8uOdGyn4Q0x85AviYhINkNbnXXJ7cI90itdqC uGKLG0BWksZX3OptSTzWuSSqpRpgpzgDDOsdMCb6G1U6Nuz7kHa/T7zzsjT3mTKb5Jcy lRiJC2S8zdIuCu7KnjwK4/id4xzsHXWXDQUlXrmhyAcf47ZLQt/YSSh1aUkOsm6UZQ8q e8mGAfEESJmRUJH1CckeDPgFp+GbcPKOgKocC9q5pvqeTQ2ehVVhPZ4T1FNlNtZxT4cT urGA3rcc+ijTczeRNoNl+WLYMdTg+XSslQOEBhthgtVias8uK1nnM57MOYQngscBADww lxWw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1784928507; x=1785533307; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=KkvlTppzlY9z98B+XfnNU0JHQBJqdVgCIejfzQAzX74=; b=h4/gRlBvaKbqPvhcSLlpUOkb0BlowWPY7OitoPqD2r7ajges/a7vq1Hiqg0B/VM7Gz ijjLN5V+Kn3twAQcOmIuQsKOYHOI/QUjb8Ud2gmt5vwJjZ2dcBXLgVfSYANwHDxehpfC GPWeDrNta5eSzsX40+FSPwUdcJAh7+YFdjrn2KVN0Qr2hiXH1QMTiS68SqdQGG/iRbnK EVWcHxprj3eeAsFIIv4u86dHRjQOfSDi6+TH397jfr0vMoaJhVKHgmoEwswH0t+RhhNW KX7o0blBAPbMUK9BImi1+H8TXrwfuV8gKWQHVTCxNOylqzber1XZmJsVuGmYe4EUy7j1 yUCg== X-Gm-Message-State: AOJu0YyVgUUkd6Fv/NvmkC4C+jr7EfnxJFbq5b24NO/Yv3l4tRZGlyyO 9DhKRgCNRvoLTSfoFUyAq8/D0raE2kiGhcd650V+UrTteSmgjP5mbIk2dCDVZg== X-Gm-Gg: AR+sD12DUFBw+QIrE5pGHlwYrRRkAwqZAniGGZzjAWqkbUHk/XWZCxSqARovBHf9VYo pcsSZbXA5DZ4W8p6bfXlm7EnMkL3892vfL9Jf9Q1c1dTPT1mYLSyMG7nbyxveu9aLvvuwv2h2Q/ iQe1+M+BGehmy8OZZruJWhnnTdVVtLBec+/JkPu8wQoS52bIRu1osS4VGEK1i9a02PeRU2aUpUC eEJDN/6egGyMxxv4WEHTqZZqeK9RWwPPyPOb75cTUeSanRMwYbBl4KzJXMW2z8OLdAR+8eQjUYu vHjlxdGcI9pCJ0b9fBOAstDOVw+43NCSaR9FGKRaKX12Z2dv5gjIzrQ18A4PirCsmhBdvrKuBWc 9XI79TEnFsBs3C1XKQ3pWjHLw2pa5PrFJDyrzO48Io5fv1l173ki/8KG1fI9UuPZnSyJdYKFBXm YVk9wGL64QRA== X-Received: by 2002:a05:6808:14c4:b0:4a1:296:adec with SMTP id 5614622812f47-4ab6a0dec55mr186123b6e.15.1784928506580; Fri, 24 Jul 2026 14:28:26 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::5d97]) by smtp.gmail.com with ESMTPSA id 586e51a60fabf-457673d2dc1sm8183883fac.10.2026.07.24.14.28.25 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 24 Jul 2026 14:28:26 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Michal Sieron , Joshua Watt Subject: [bitbake-devel][PATCH 1/6] hashserv: client: Add asynchronous streaming API Date: Fri, 24 Jul 2026 15:25:09 -0600 Message-ID: <20260724212822.1165552-2-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> MIME-Version: 1.0 List-Id: X-Webhook-Received: from 45-33-107-173.ip.linodeusercontent.com [45.33.107.173] by aws-us-west-2-korg-lkml-1.web.codeaurora.org with HTTPS for ; Fri, 24 Jul 2026 21:28:31 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19851 Implments a true asynchronous streaming API for getting a unihash (get_unihash_stream()) and checking if a unihash exists (unihash_exists_stream()). These APIs allow a client to send a query the to the server, then wait for the reply later. Using this API, it is possible to interleave queries and replies. In addition, the gc_mark_stream() API is renamed gc_mark_batch() to match the pattern of the other APIs, and an actual gc_mark_stream() API is implemented that allows asynchronous streaming like the others. Note that the new stream APIs are not accessible in the "synchronous" client API since async constructs must be used for them to make sense. Signed-off-by: Joshua Watt --- bin/bitbake-hashclient | 2 +- lib/hashserv/client.py | 306 ++++++++++++++++++++++++++++++++--------- lib/hashserv/tests.py | 4 +- 3 files changed, 243 insertions(+), 69 deletions(-) diff --git a/bin/bitbake-hashclient b/bin/bitbake-hashclient index 3a2bf5c0d..ee5c32f5f 100755 --- a/bin/bitbake-hashclient +++ b/bin/bitbake-hashclient @@ -240,7 +240,7 @@ def main(): marked_hashes = 0 try: - result = client.gc_mark_stream(args.mark, stdin) + result = client.gc_mark_batch(args.mark, stdin) marked_hashes = result["count"] except ConnectionError: logger.warning( diff --git a/lib/hashserv/client.py b/lib/hashserv/client.py index 8cb18050a..8e9eaaca3 100644 --- a/lib/hashserv/client.py +++ b/lib/hashserv/client.py @@ -8,44 +8,212 @@ import socket import asyncio import bb.asyncrpc import json +from abc import abstractmethod +from collections.abc import AsyncIterable +from contextlib import asynccontextmanager +from dataclasses import dataclass from . import create_async_client - logger = logging.getLogger("hashserv.client") +class AsyncQueue(AsyncIterable): + class Shutdown(Exception): + pass + + SHUTDOWN_SENTINEL = object() + + def __init__(self, *args, **kwargs): + self.__queue = asyncio.Queue() + self.__shutdown = False + self.__is_done = False + + async def done(self): + if self.__is_done: + return + self.__is_done = True + await self.__queue.put(self.SHUTDOWN_SENTINEL) + + async def put(self, item): + if self.__is_done: + raise self.Shutdown + await self.__queue.put(item) + + async def get(self): + if self.__shutdown: + raise self.Shutdown + + item = await self.__queue.get() + if item is self.SHUTDOWN_SENTINEL: + self.__shutdown = True + raise self.Shutdown + + return item + + def __aiter__(self): + return self + + async def __anext__(self): + try: + return await self.get() + except self.Shutdown: + raise StopAsyncIteration + + +@dataclass(eq=False, frozen=True) +class AsyncPipe: + send_queue: AsyncQueue + recv_queue: AsyncQueue + + +class Stream(AsyncIterable): + def __init__(self, pipe): + self._pipe = pipe + + async def done(self): + await self._pipe.send_queue.done() + + def __aiter__(self): + return self + + async def __anext__(self): + try: + return await self.get_result() + except AsyncQueue.Shutdown: + raise StopAsyncIteration + + @abstractmethod + async def _send_batch_input(self, i): + raise NotImplementedError("Not implemented") + + @abstractmethod + async def get_result(self): + raise NotImplementedError("Not implemented") + + async def batch(self, inputs): + """ + Does a "batch" process of stream messages. This sends the query + messages as fast as possible, and simultaneously attempts to read the + messages back. This helps to mitigate the effects of latency to the + hash equivalence server be allowing multiple queries to be "in-flight" + at once + + The input may be a generator or an async generator + """ + + async def get_inputs(): + if isinstance(inputs, AsyncIterable): + async for i in inputs: + yield i + else: + for i in inputs: + yield i + + async def send(): + try: + async for i in get_inputs(): + await self._send_batch_input(i) + finally: + await self.done() + + results = [] + + async def recv(): + async for item in self: + results.append(item) + + await asyncio.gather(send(), recv()) + return results + + +class GetUnihashStream(Stream): + def __init__(self, pipe): + super().__init__(pipe) + + async def _send_batch_input(self, i): + method, taskhash = i + await self.send_query(method, taskhash) + + async def send_query(self, method, taskhash): + await self._pipe.send_queue.put(f"{method} {taskhash}") + + async def get_result(self): + r = await self._pipe.recv_queue.get() + return r if r else None + + +class UnihashExistsStream(Stream): + def __init__(self, pipe): + super().__init__(pipe) + + async def _send_batch_input(self, i): + await self.send_query(i) + + async def send_query(self, unihash): + await self._pipe.send_queue.put(unihash) + + async def get_result(self): + r = await self._pipe.recv_queue.get() + return r == "true" + + +class GcMarkStream(Stream): + def __init__(self, pipe, mark): + super().__init__(pipe) + self.mark = mark + + async def _send_batch_input(self, i): + def row_to_dict(row): + pairs = row.split() + return dict(zip(pairs[::2], pairs[1::2])) + + await self.send_mark(row_to_dict(i)) + + async def send_mark(self, where): + await self._pipe.send_queue.put(json.dumps({"mark": self.mark, "where": where})) + + async def get_result(self): + r = await self._pipe.recv_queue.get() + return json.loads(r) + + class Batch(object): - def __init__(self): + def __init__(self, send_queue, recv_queue): + self.send_queue = send_queue + self.recv_queue = recv_queue self.done = False self.cond = asyncio.Condition() self.pending = [] - self.results = [] self.sent_count = 0 + self.recv_count = 0 async def recv(self, socket): - while True: - async with self.cond: - await self.cond.wait_for(lambda: self.pending or self.done) + try: + while True: + async with self.cond: + await self.cond.wait_for(lambda: self.pending or self.done) - if not self.pending: - if self.done: - return - continue + if not self.pending: + if self.done: + return + continue - r = await socket.recv() - self.results.append(r) + await self.recv_queue.put(await socket.recv()) - async with self.cond: - self.pending.pop(0) + async with self.cond: + self.recv_count += 1 + self.pending.pop(0) + finally: + await self.recv_queue.done() - async def send(self, socket, msgs): + async def send(self, socket): try: # In the event of a restart due to a reconnect, all in-flight # messages need to be resent first to keep to result count in sync for m in self.pending: await socket.send(m) - for m in msgs: + async for m in self.send_queue: # Add the message to the pending list before attempting to send # it so that if the send fails it will be retried async with self.cond: @@ -60,19 +228,14 @@ class Batch(object): self.done = True self.cond.notify() - async def process(self, socket, msgs): - await asyncio.gather( - self.recv(socket), - self.send(socket, msgs), - ) + async def stream(self, socket): + await asyncio.gather(self.send(socket), self.recv(socket)) - if len(self.results) != self.sent_count: - raise ValueError( - f"Expected result count {len(self.results)}. Expected {self.sent_count}" + if self.sent_count != self.recv_count: + raise ConnectionError( + f"Sent {self.sent_count} messages but only received {self.recv_count}" ) - return self.results - class AsyncClient(bb.asyncrpc.AsyncClient): MODE_NORMAL = 0 @@ -98,29 +261,32 @@ class AsyncClient(bb.asyncrpc.AsyncClient): if become: await self.become_user(become) - async def send_stream_batch(self, mode, msgs): - """ - Does a "batch" process of stream messages. This sends the query - messages as fast as possible, and simultaneously attempts to read the - messages back. This helps to mitigate the effects of latency to the - hash equivalence server be allowing multiple queries to be "in-flight" - at once - - The implementation does more complicated tracking using a count of sent - messages so that `msgs` can be a generator function (i.e. its length is - unknown) - - """ - - b = Batch() + @asynccontextmanager + async def send_stream(self, mode): + send_queue = AsyncQueue() + recv_queue = AsyncQueue() + b = Batch(send_queue, recv_queue) async def proc(): - nonlocal b - await self._set_mode(mode) - return await b.process(self.socket, msgs) + await b.stream(self.socket) + + async def process(): + try: + await self._send_wrapper(proc) + finally: + await recv_queue.done() - return await self._send_wrapper(proc) + # Create background process to process messages + task = asyncio.create_task(process()) + + try: + yield AsyncPipe(send_queue, recv_queue) + except AsyncQueue.Shutdown: + pass + finally: + await send_queue.done() + await task async def invoke(self, *args, skip_mode=False, **kwargs): # It's OK if connection errors cause a failure here, because the mode @@ -173,15 +339,18 @@ class AsyncClient(bb.asyncrpc.AsyncClient): self.mode = new_mode async def get_unihash(self, method, taskhash): - r = await self.get_unihash_batch([(method, taskhash)]) - return r[0] + async with self.get_unihash_stream() as stream: + await stream.send_query(method, taskhash) + return await stream.get_result() async def get_unihash_batch(self, args): - result = await self.send_stream_batch( - self.MODE_GET_STREAM, - (f"{method} {taskhash}" for method, taskhash in args), - ) - return [r if r else None for r in result] + async with self.get_unihash_stream() as stream: + return await stream.batch(args) + + @asynccontextmanager + async def get_unihash_stream(self): + async with self.send_stream(self.MODE_GET_STREAM) as pipe: + yield GetUnihashStream(pipe) async def report_unihash(self, taskhash, method, outhash, unihash, extra={}): m = extra.copy() @@ -204,12 +373,18 @@ class AsyncClient(bb.asyncrpc.AsyncClient): ) async def unihash_exists(self, unihash): - r = await self.unihash_exists_batch([unihash]) - return r[0] + async with self.unihash_exists_stream() as stream: + await stream.send_query(unihash) + return await stream.get_result() async def unihash_exists_batch(self, unihashes): - result = await self.send_stream_batch(self.MODE_EXIST_STREAM, unihashes) - return [r == "true" for r in result] + async with self.unihash_exists_stream() as stream: + return await stream.batch(unihashes) + + @asynccontextmanager + async def unihash_exists_stream(self): + async with self.send_stream(self.MODE_EXIST_STREAM) as pipe: + yield UnihashExistsStream(pipe) async def get_outhash(self, method, outhash, taskhash, with_unihash=True): return await self.invoke( @@ -309,23 +484,22 @@ class AsyncClient(bb.asyncrpc.AsyncClient): """ return await self.invoke({"gc-mark": {"mark": mark, "where": where}}) - async def gc_mark_stream(self, mark, rows): + async def gc_mark_batch(self, mark, rows): """ Similar to `gc-mark`, but accepts a list of "where" key-value pair conditions. It utilizes stream mode to mark hashes, which helps reduce the impact of latency when communicating with the hash equivalence server. """ - def row_to_dict(row): - pairs = row.split() - return dict(zip(pairs[::2], pairs[1::2])) + async with self.gc_mark_stream(mark) as stream: + results = await stream.batch(rows) - responses = await self.send_stream_batch( - self.MODE_MARK_STREAM, - (json.dumps({"mark": mark, "where": row_to_dict(row)}) for row in rows), - ) + return {"count": sum(int(r["count"]) for r in results)} - return {"count": sum(int(json.loads(r)["count"]) for r in responses)} + @asynccontextmanager + async def gc_mark_stream(self, mark): + async with self.send_stream(self.MODE_MARK_STREAM) as pipe: + yield GcMarkStream(pipe, mark) async def gc_sweep(self, mark): """ @@ -372,7 +546,7 @@ class Client(bb.asyncrpc.Client): "get_db_query_columns", "gc_status", "gc_mark", - "gc_mark_stream", + "gc_mark_batch", "gc_sweep", ) diff --git a/lib/hashserv/tests.py b/lib/hashserv/tests.py index 7c736d6cc..3acfdcd8c 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -1055,7 +1055,7 @@ class HashEquivalenceCommonTests(object): # First hash is still present self.assertClientGetHash(self.client, taskhash, unihash) - def test_gc_stream(self): + def test_gc_batch(self): taskhash = '53b8dce672cb6d0c73170be43f540460bfc347b4' outhash = '5a9cb1649625f0bf41fc7791b635cd9c2d7118c7f021ba87dcd03f72b67ce7a8' unihash = '46edb5140d2613049332d0bf3745d9fafec9c559dac8cc61813739a28007fcdf' @@ -1078,7 +1078,7 @@ class HashEquivalenceCommonTests(object): self.assertClientGetHash(self.client, taskhash3, unihash3) # Mark the first unihash to be kept - ret = self.client.gc_mark_stream("ABC", (f"unihash {h}" for h in [unihash, unihash2])) + ret = self.client.gc_mark_batch("ABC", (f"unihash {h}" for h in [unihash, unihash2])) self.assertEqual(ret, {"count": 2}) ret = self.client.gc_status() From patchwork Fri Jul 24 21:25:10 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93472 Return-Path: X-Spam-Checker-Version: SpamAssassin 3.4.0 (2014-02-07) on aws-us-west-2-korg-lkml-1.web.codeaurora.org Received: from aws-us-west-2-korg-lkml-1.web.codeaurora.org (localhost.localdomain [127.0.0.1]) by smtp.lore.kernel.org (Postfix) with ESMTP id 0E773C531FC for ; Fri, 24 Jul 2026 21:28:31 +0000 (UTC) Received: from mail-oa1-f53.google.com (mail-oa1-f53.google.com [209.85.160.53]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.28978.1784928508300628099 for ; Fri, 24 Jul 2026 14:28:28 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=n/unC41h; spf=pass (domain: gmail.com, ip: 209.85.160.53, mailfrom: jpewhacker@gmail.com) Received: by mail-oa1-f53.google.com with SMTP id 586e51a60fabf-448479e0eb6so464784fac.2 for ; Fri, 24 Jul 2026 14:28:28 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1784928507; x=1785533307; darn=lists.openembedded.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=m6B/k+wGcWhzzpYrSeOAg49ZYJraXxDiY7Tsj1AzqJ4=; b=n/unC41hhUphPKFuELGTkm/yihlEUbGD/ybW4GS2xLmaadpO5TGzyhyGQgwEB7NQu/ 239vvpa86zQ9a5HPCnmE9Tu1O1o6rPsHqGbxkscQsU0G1VvlvbyAkoNu/ZlWjI0mPfFP BKDXO9N9VgiE5lmSV3NmifiVgw5qJoBIM2ynR5i2PNsn9tUTTXxjfP2HL4KXZrdX/hv8 9SeDGVm0i/MOY7je1sfui7ZmsBjRQ4H4G09ZXt2EO+D90coV4uD2FhBnbbkGXYe2rT3k HqEsJudusERweWHXFmte4xUtL6W7UBdEV1GvYYcDUOfcNPnxSpIoQslmF/QnHAko032w 4GZA== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1784928507; x=1785533307; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=m6B/k+wGcWhzzpYrSeOAg49ZYJraXxDiY7Tsj1AzqJ4=; b=iqXSLlS8vABL65iskws9VpyYsaMs6KrS/ht47yx/kMkZp2Z7anSuRIq/Yg312+oR8k Pwt6GWeizFb1ZkPnGhA8bWZ2r2oHz5NZ03PASnJ57rtZumqN7ENQKjsw9Lgq93BV5DmN vNUxkyz0MUvw/KJ3HF6HIWh5zZXfmQkp6BE7nexzkP7pZjY7omjTRl6qT7gTlzWIpkN9 /t10k08itSVzQxiaMEnqumulELFFCAoZl4X5rTAqkaovKmzis9Kpu7lvHWyTaE6pAT1X j1yAdm/DiZb0zsz2Xzpe21QvZx8MVWq3iGCXceS9jhyJu7mFq6vzd+YoMHr2GJXhk9jK sECg== X-Gm-Message-State: AOJu0YyQmNeH7KOl/caDERuvuI9LnD2flm7cF415mOrYmnz3G7aEcBH7 A/4WkfUKIoYBxxIdaSf1qj/ePz6BkFKnxKJuVj5zXvxo3byiphqf3xSJyqzF5Q== X-Gm-Gg: AR+sD111AfFLcF1V7D6+kaTXPJQWbdiyqKwPuwp8WYQ3Mwmgh6vZBkgmL1Jg122Rgko GPlVlCF2+TpI/J5AX/I35zkjMyq2yMKlklNEUBn5+ZM5SFX/frH/TZPiDijeKcuBspTvg5s/xuZ Zduy5g7AaLIqRNqL8GdR1yjXrNQYQN1lWMeTMoumvLFr9TcsvH2b/1ZqCzVRSQI/Y0u7dnpv7jR 9xcKRpFkGf5tlbJkAh8DJsxgs9smvS4KnLaFKCWtl9qQVkxsGl22f9lDHveGP0w/f1cf/3QFWU6 J6IQj5HFyg8HehqY6hM5EdFNr6b159j8NqxCbVDFy1X+2Asc26i/Vwju7Cn4Q8dOs43NBDFwp1V gJWxr81k56G+oA7lwEWMdeaVdAl5/rfH/Cb69pscWCSkVIWaSo29UzsFXQEg4prX4RSYS0OrDmQ A= X-Received: by 2002:a05:6870:524d:b0:448:9a86:aff9 with SMTP id 586e51a60fabf-457f247a42fmr175175fac.15.1784928507323; Fri, 24 Jul 2026 14:28:27 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::5d97]) by smtp.gmail.com with ESMTPSA id 586e51a60fabf-457673d2dc1sm8183883fac.10.2026.07.24.14.28.26 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 24 Jul 2026 14:28:27 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Michal Sieron , Joshua Watt Subject: [bitbake-devel][PATCH 2/6] hashserv: server: Fix formatting Date: Fri, 24 Jul 2026 15:25:10 -0600 Message-ID: <20260724212822.1165552-3-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> MIME-Version: 1.0 List-Id: X-Webhook-Received: from 45-33-107-173.ip.linodeusercontent.com [45.33.107.173] by aws-us-west-2-korg-lkml-1.web.codeaurora.org with HTTPS for ; Fri, 24 Jul 2026 21:28:31 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19852 Reformats the file with `black` Signed-off-by: Joshua Watt --- lib/hashserv/server.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py index 3ff434785..e7e79196f 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -424,7 +424,7 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): @permissions(READ_PERM) async def handle_get_stream(self, request): async def handler(l): - (method, taskhash) = l.split() + method, taskhash = l.split() # self.logger.debug('Looking up %s %s' % (method, taskhash)) row = await self.db.get_equivalent(method, taskhash) @@ -646,7 +646,7 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): @permissions(DB_ADMIN_PERM) async def handle_gc_status(self, request): - (keep_rows, remove_rows, current_mark) = await self.db.gc_status() + keep_rows, remove_rows, current_mark = await self.db.gc_status() return { "keep": keep_rows, "remove": remove_rows, @@ -903,7 +903,9 @@ class Server(bb.asyncrpc.AsyncServer): d = await client.get_taskhash(method, taskhash) if d is not None: if is_valid_unihash(d.get("unihash")): - await db.insert_unihash(d["method"], d["taskhash"], d["unihash"]) + await db.insert_unihash( + d["method"], d["taskhash"], d["unihash"] + ) else: self.logger.warning("Upstream server returned invalid unihash") self.backfill_queue.task_done() From patchwork Fri Jul 24 21:25:11 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93473 Return-Path: X-Spam-Checker-Version: SpamAssassin 3.4.0 (2014-02-07) on aws-us-west-2-korg-lkml-1.web.codeaurora.org Received: from aws-us-west-2-korg-lkml-1.web.codeaurora.org (localhost.localdomain [127.0.0.1]) by smtp.lore.kernel.org (Postfix) with ESMTP id 0BC22C531F9 for ; Fri, 24 Jul 2026 21:28:31 +0000 (UTC) Received: from mail-ot1-f45.google.com (mail-ot1-f45.google.com [209.85.210.45]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.28919.1784928508896857307 for ; Fri, 24 Jul 2026 14:28:29 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=F5XdI5AD; spf=pass (domain: gmail.com, ip: 209.85.210.45, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f45.google.com with SMTP id 46e09a7af769-7eb9b427da2so1172301a34.0 for ; Fri, 24 Jul 2026 14:28:28 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1784928508; x=1785533308; darn=lists.openembedded.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=cJgDGOKrgOhO4GqIerRQ7ja2rXw7cSOKtYqO0XqDjnU=; b=F5XdI5AD5GF5e+t9fXticu+COuALLT/1pwBAV8DGno2qdwolxDfwpmecE7s8K69oGx N9noCJw3UoG+alXHQna40J+u8UJ/mX8d9/seEF4rD+GV8uynbdERNnzyosNacsxfRfLG uevOcitlca/1p2Vfnk0FzwA2qNZ+8slJlgc5+JYqErGUxhuaISR1douKzJAjMAT0FiT3 LWVv0wwQoZQXfqFbRu+GrgbPq2SOORoAnmKjO7P4xbUVwizSJLtldKJv0rkloqcfm0Eg fET6ZUZRLB9e4lPiVaEe27cjYo6HytoeanscbAq6JLNI0x/Z2Ndg/+HWXFDlQndXOLxR zrQg== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1784928508; x=1785533308; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=cJgDGOKrgOhO4GqIerRQ7ja2rXw7cSOKtYqO0XqDjnU=; b=pJte5n346JBuFFbGYyJHRvfYvAQfyjplhzZ7AXjbC4h8/rJf+e1Llx4FNe0AXCW5xz EMF/0MKp0pEE1Q5t3NLHmtb+W0DNwANMoSefBo2XGDZOzWnqfD6LqFJTcMy33A7AiY3V dBCZvNbRC+Ket0VMi9dDGokzL8JtGYmkmrvSaDKEDGhkfqVH+itmq616fpkTKYADPntY zVYgQfJYM86wtwbcNjB5YPE2yk2k1astkV2uI5+qz9Z4VfiKivVJ4yIGjuXbMHCRRLq3 qB8g6nud1q86djNhtvjxPzkndfivEAKmUrX8cRbSjcAU0xnCFq4ZdVtdZGMPcQ+DctfG rR8g== X-Gm-Message-State: AOJu0YzEG/5OEqzL9uB821R3LVW7LDP9z9kVriZiWz/B00RqkhB2IxGQ rbtjYbGJAQgkODbqa3Vp+wsZn6zCDwG+d3leM8PFxgCV/1oviBhqYeE7sYw4Kg== X-Gm-Gg: AR+sD12B7RZKs47IQcsmIURVpsgxiMFYK/Nd1MpW+yNPjGZavmpdRyuqzROayjNB/8z J7NxXCCzUQmjWm6NMKMrDP8N4Eh0JduL5CDjkMpqRC/qH8RBzTj6xQRQjznKF4o/nQJ8YTVuC9o M8y1tMSzhAC3f5EVd1IYFDS+/q3KVXX4AgTWVNvTzEzwaNu+2EPfDaM+dnG7CIaDOdZmVXyb61l R6dTrU8QGJB0QhkZQjPlyIpENnP1dfC/6C3v4kLrGuKBUZ5uojcVlhQvpuqM1gVSf9KhVH9rk7w khno0UjDqELghkYI9XeUVeYijeA1isjglNQVphzIw6/oxiVynJyPRnAzbTCzVCJuUG7AMZWBbud x0ir/zT9ddhviebMa6RGfM+x2abcVvjnkW0k7l4tijz4BNe4eQz3PUajCiLB2cJgeXgBn3bNiUY U= X-Received: by 2002:a05:6808:c2b8:b0:496:9ee:e538 with SMTP id 5614622812f47-4ab5e33b438mr2336213b6e.5.1784928508048; Fri, 24 Jul 2026 14:28:28 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::5d97]) by smtp.gmail.com with ESMTPSA id 586e51a60fabf-457673d2dc1sm8183883fac.10.2026.07.24.14.28.27 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 24 Jul 2026 14:28:27 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Michal Sieron , Joshua Watt Subject: [bitbake-devel][PATCH 3/6] hashserv: server: Add queued streaming API Date: Fri, 24 Jul 2026 15:25:11 -0600 Message-ID: <20260724212822.1165552-4-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> MIME-Version: 1.0 List-Id: X-Webhook-Received: from 45-33-107-173.ip.linodeusercontent.com [45.33.107.173] by aws-us-west-2-korg-lkml-1.web.codeaurora.org with HTTPS for ; Fri, 24 Jul 2026 21:28:31 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19853 Adds an API that allows a stream handler to more precisely control when a response is sent to the client. The new API does not directly send the response from the handler to the remote client, but instead the handler is expected to put the result in a provided queue when the response is ready. In particular, this allows a stream handler to defer to an upstream server (utilizing the new client streaming API) in an efficient way that does not require waiting on a roundtrip with the upstream server. Instead, several queries to the upstream can be in-flight at once. Signed-off-by: Joshua Watt --- lib/hashserv/server.py | 100 ++++++++++++++++++++++++++++++++--------- 1 file changed, 79 insertions(+), 21 deletions(-) diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py index e7e79196f..991b99b86 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -229,6 +229,58 @@ def permissions(*permissions, allow_anon=True, allow_self_service=False): return wrapper +class UpstreamQueue(object): + UPSTREAM_NONCE = object() + + def __init__(self, queue, get_local_result, send_upstream, get_upstream_result): + self.queue = queue + self.pending = [] + self.cond = asyncio.Condition() + self.done = False + self.get_local_result = get_local_result + self.send_upstream = send_upstream + self.get_upstream_result = get_upstream_result + + async def process_results(self): + try: + while True: + async with self.cond: + await self.cond.wait_for(lambda: self.pending or self.done) + if not self.pending: + if self.done: + return + continue + + value, m = self.pending.pop(0) + + if value is self.UPSTREAM_NONCE: + value = await self.get_upstream_result(m) + + await self.queue.put(value) + finally: + await self.queue.put(None) + + async def handler(self, m): + try: + if m is None: + return + + value = await self.get_local_result(m) + if value is None: + await self.send_upstream(m) + value = self.UPSTREAM_NONCE + + async with self.cond: + self.pending.append((value, m)) + self.cond.notify_all() + + finally: + async with self.cond: + self.done = True + self.cond.notify_all() + # await stream.done() + + class ServerClient(bb.asyncrpc.AsyncServerConnection): def __init__(self, socket, server): super().__init__(socket, "OEHASHEQUIV", server.logger) @@ -390,35 +442,41 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): validate_unihash(unihash) return await self.db.insert_unihash(method, taskhash, unihash) - async def _stream_handler(self, handler): + async def _stream_queue_handler(self, handler, queue): await self.socket.send_message("ok") - while True: - upstream = None + async def recv(): + try: + while True: + m = await self.socket.recv() + if not m or m == "END": + break - l = await self.socket.recv() - if not l: - break + await handler(m) + finally: + await handler(None) - try: - # This inner loop is very sensitive and must be as fast as - # possible (which is why the request sample is handled manually - # instead of using 'with', and also why logging statements are - # commented out. - self.request_sample = self.server.request_stats.start_sample() - request_measure = self.request_sample.measure() - request_measure.start() - - if l == "END": + async def process(): + while True: + m = await queue.get() + if m is None: break - msg = await handler(l) - await self.socket.send(msg) - finally: - request_measure.end() - self.request_sample.end() + await self.socket.send(m) + await asyncio.gather(recv(), process()) await self.socket.send("ok") + + async def _stream_handler(self, handler): + queue = asyncio.Queue(1000) + + async def h(m): + if m is None: + await queue.put(None) + else: + await queue.put(await handler(m)) + + await self._stream_queue_handler(h, queue) return self.NO_RESPONSE @permissions(READ_PERM) From patchwork Fri Jul 24 21:25:12 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93474 Return-Path: X-Spam-Checker-Version: SpamAssassin 3.4.0 (2014-02-07) on aws-us-west-2-korg-lkml-1.web.codeaurora.org Received: from aws-us-west-2-korg-lkml-1.web.codeaurora.org (localhost.localdomain [127.0.0.1]) by smtp.lore.kernel.org (Postfix) with ESMTP id 2DDA2C531C9 for ; Fri, 24 Jul 2026 21:28:31 +0000 (UTC) Received: from mail-oa1-f54.google.com (mail-oa1-f54.google.com [209.85.160.54]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.28980.1784928509707117223 for ; Fri, 24 Jul 2026 14:28:29 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=nNKu7csa; spf=pass (domain: gmail.com, ip: 209.85.160.54, mailfrom: jpewhacker@gmail.com) Received: by mail-oa1-f54.google.com with SMTP id 586e51a60fabf-451d8064238so505670fac.3 for ; Fri, 24 Jul 2026 14:28:29 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1784928509; x=1785533309; darn=lists.openembedded.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=kzMNtYN/R/LBgC0qMvBwsXC6aLcsl21fA5MMzT99uNg=; b=nNKu7csaUiU1QfdD7V6BwcuLfpuhw1xQvbPdhqNbmgBctzXpznwkAurfa4z65Zpv2C EskCM88NJpP2qAW8+R+FrUIO7urSCm8n8yD1stgezaQXT7eKcoKt8oqbKfPeSgHP/y6M 1ezrn+poYt7GOVv0hYrVYRCMSnCn1TOv+LRr9BmCxuHTdo16zTF+pb9krhAOrev8NTlp k8dB/TbJAfvqTviRTAS55sdRmurlhXydweYZU81w+ru4OZ8IbBzOaKGPr4PDatJ8TR34 dsRwPWccCwBuH3eZqEMFPA9DvB7nGAD4uwIkMww4lCv4/9TuGPEplWIeOuRP8diPHWm/ K4lw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1784928509; x=1785533309; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=kzMNtYN/R/LBgC0qMvBwsXC6aLcsl21fA5MMzT99uNg=; b=mEfb9vvOuYO5P/FMdD9A3PUygLlkIlhXlShV3lwHZsOuQ8dpFPfaAwf/OMl1LAa6n0 gbz1zplw6j4DDaySGxNtZbcZUNxqCgEuDSm+4CGNRSOBV1B68FAokBkyCLw8S7hZqNsc 2+J2ua77J/p+1px/HzHodCqHxFZGCcxRbRN/yLSq/T8HPdKyFpk3ho/FjyyQihWLbwNP Mk5g9qE6t0pvqCFfuGeDPbnT0TO6MJiayp7tms31J1dGgWrOmTwAlSuB8N+Jqll55WHv a1r15dB8hacMSJZhB8coSDOS711fT6Y+zo0eEjTP08l3M0O5EynyOFMR8FR0P9rbAM+r RUEg== X-Gm-Message-State: AOJu0YyN5U1JJrwpv9uMlN4UEeP3ATE+SVOQlD39i7nNLWOZ8P3rw4Mx T6+Gq61AiGaug+CqQheDbsv21TYF8k/yvIvu1+TFEaRyCE8h6D9wdnIaBveDJg== X-Gm-Gg: AR+sD11E9M9kG+PKiOry3P8/yD1fNHT9NnrJafaI+pfzXc04G4j6sGY810Sd6fn+P+4 ldYwW9hPvQgOZGUy/2crv3LiUoeM2WXARO/7LtNCu80JRjLIKFogzOqYgWCQRjHbs9CUHgJMFRe 8ip3y/4Aj18524SMgjFFWshb4aWfaeem8UI3W5f6xcKg6z0tXolFVA0tl3vSPo/GDztD7jjskjO 9homjxc0v3DXlwbahXudj6nCsnFE8UKHGiGYH4XT4oASCbPVztuO/vVuXVyID/rSTLvEm9Pr/rm qXOkqeFZ3ZB9qvwEj0PLt02L+xPSTtBJy0Tdfwf2oZE68jZ89gB0AvRHCIofPyqL9Ns1XgJuOs+ 62MOJjXOUQ4XzvUo+ocAVVfm7BkYNT+h+LGRiTyTowFqg3Lf5TBVjl9bCmSo6ZEvYgT2gIpo8LU o= X-Received: by 2002:a05:6870:3195:b0:448:c1f5:90f7 with SMTP id 586e51a60fabf-457f24df7d7mr207095fac.19.1784928508864; Fri, 24 Jul 2026 14:28:28 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::5d97]) by smtp.gmail.com with ESMTPSA id 586e51a60fabf-457673d2dc1sm8183883fac.10.2026.07.24.14.28.28 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 24 Jul 2026 14:28:28 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Michal Sieron , Joshua Watt Subject: [bitbake-devel][PATCH 4/6] hashserv: server: Use streaming and queue API for upstream unihash queries Date: Fri, 24 Jul 2026 15:25:12 -0600 Message-ID: <20260724212822.1165552-5-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> MIME-Version: 1.0 List-Id: X-Webhook-Received: from 45-33-107-173.ip.linodeusercontent.com [45.33.107.173] by aws-us-west-2-korg-lkml-1.web.codeaurora.org with HTTPS for ; Fri, 24 Jul 2026 21:28:31 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19854 Reworks the "get-unihash" 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 --- lib/hashserv/server.py | 52 ++++++++++++++++++++++++++++++------------ 1 file changed, 38 insertions(+), 14 deletions(-) diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py index 991b99b86..d392ac19f 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -481,24 +481,48 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): @permissions(READ_PERM) async def handle_get_stream(self, request): - async def handler(l): - method, taskhash = l.split() - # self.logger.debug('Looking up %s %s' % (method, taskhash)) - row = await self.db.get_equivalent(method, taskhash) - - if row is not None: - # self.logger.debug('Found equivalent task %s -> %s', (row['taskhash'], row['unihash'])) + async def get_unihash(m): + method, taskhash = m.split() + if (row := await self.db.get_equivalent(method, taskhash)) is not None: return row["unihash"] - if self.upstream_client is not None: - upstream = await self.upstream_client.get_unihash(method, taskhash) - if upstream: - await self.server.backfill_queue.put((method, taskhash)) - return upstream - return "" - return await self._stream_handler(handler) + if not self.upstream_client: + return await self._stream_handler(get_unihash) + + async with self.upstream_client.get_unihash_stream() as stream: + + async def get_local_result(m): + method, taskhash = m.split() + if (row := await self.db.get_equivalent(method, taskhash)) is not None: + return row["unihash"] + return None + + async def send_upstream(m): + method, taskhash = m.split() + await stream.send_query(method, taskhash) + + async def get_upstream_result(m): + unihash = await stream.get_result() + if unihash: + method, taskhash = m.split() + await self.server.backfill_queue.put((method, taskhash)) + return unihash + + queue = asyncio.Queue() + upstream = UpstreamQueue( + queue, + get_local_result, + send_upstream, + get_upstream_result, + ) + + await asyncio.gather( + self._stream_queue_handler(upstream.handler, queue), + upstream.process_results(), + ) + return self.NO_RESPONSE @permissions(READ_PERM) async def handle_exists_stream(self, request): From patchwork Fri Jul 24 21:25:13 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93476 Return-Path: X-Spam-Checker-Version: SpamAssassin 3.4.0 (2014-02-07) on aws-us-west-2-korg-lkml-1.web.codeaurora.org Received: from aws-us-west-2-korg-lkml-1.web.codeaurora.org (localhost.localdomain [127.0.0.1]) by smtp.lore.kernel.org (Postfix) with ESMTP id 3CCA0C5321A for ; Fri, 24 Jul 2026 21:28:31 +0000 (UTC) Received: from mail-oa1-f47.google.com (mail-oa1-f47.google.com [209.85.160.47]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.28920.1784928510597662444 for ; Fri, 24 Jul 2026 14:28:30 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=q3MtA4ib; spf=pass (domain: gmail.com, ip: 209.85.160.47, mailfrom: jpewhacker@gmail.com) Received: by mail-oa1-f47.google.com with SMTP id 586e51a60fabf-4513435cdd2so399346fac.2 for ; Fri, 24 Jul 2026 14:28:30 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1784928510; x=1785533310; darn=lists.openembedded.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=IPuFqV+fsZTcQsNuHhGmmdXtXVUrpvyQg/hP39CHzIA=; b=q3MtA4ibSne1DNBv12rNwcIMKAnbhwaQLG0iGcS/Zw/7DdSrCuQExKs+I2wEObXQLa dqYEqzyT/D9EMH7c1wzkeI3P1GXTNoSyYaGACr0vj6UrDQD610Vx5O6ZXG7n93bn/ewx YJG1wirUJxm7aNOJcYaFDrbNi+hg4LJeR15w80+Npbv4jHarra8VmnBqheq0J/dVuN/w ETRoyCHNrAi/eHxfXjKtQkxYuNiHyU4fwlQpr1OtS87valr46m+hPuGhwF5eI2q634/+ X7LrA6kSe1gkvSzvGiIYS8FSOShhvm5/aLLfyMI1ksTXCMd0Rna9sOy/Jz0SX6sBQ8kr 1mdQ== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1784928510; x=1785533310; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=IPuFqV+fsZTcQsNuHhGmmdXtXVUrpvyQg/hP39CHzIA=; b=jF68Ah/RhrEExOFOJTvzU6JQY3kjqF4+Ux331HgFhjz/seEdtxIznupoQN96pWC0go zBFp5FaEsf1tdrdwwFwkQBQolDQTqwWb3KkHz3+ZZhONzdPi/WMZZYNe6Eh+o8J3eUPF XdI1XIG2Mu5Hhl282tqW7udRC5u44FRdu8Pza98BDEFES9xHKvj7vI/CC+XPEljXeAsb h3vmNQUaNSDXSQ5wQH7zTqX9upwxBPf8IM5WAUoLwPro/OR8xE7/eVT4DN1Cm33/r4da bUyk2HBewMX+VTMkzwdRpDgWxCyMCvo2ir3E0cDpDph3yvwL9/77FNfcwVsCzCHfM6Uv eunw== X-Gm-Message-State: AOJu0YzzfoGx7y+LYLoInf3dy/Our/kPh+ZYhQx0zpyKG+IM5llGMBB4 6AoWMiuGiYy+91FPdgWDsGf4DuENAp4mTIWbuasXa9+TkkEZDf4K2KstlC8Z/Q== X-Gm-Gg: AR+sD11JHnUzt8HEq0TlbAO+JTKDNyJXjYqyqCui5d56QVo59Yk7pK92EqkqjA49e+0 TrugHok7rFQybMrmKeBnbB3NhxqR2god3VNeM3OhNVnbBJwlr9DF4uH2qaqeXlfC41CFKOwdNl2 nc8tPQgFARmI1mfRgkDkenzZJhaVx3oKPTwxNDPvr60pEhv1YX1gXh3LFTZrpSm+CunbZxLm5Ys 4X/34Zj/pq9z1U48DCE5HoY9xzlNdbW0fP6A2Cz6cZxXKil0xdhTkIrbKZ9MRtD3Inx8J3vj5cB NRC1YzZW0XvHglakp/YDY+M3fZW2zZlEtp9SwtcK+fd1dsHJ2Rlt+JSYKN2Bnmfajdu66+VXmI4 VNum/sfGzCsafDjxjpX6qKzj1moPHjOdPr6B06MASbi+aAclMJsgKOYhfQa9EzAE3vMeknH/H7v fxqBvftFIQYA== X-Received: by 2002:a05:6870:488:b0:448:a0cf:5afd with SMTP id 586e51a60fabf-457f23900bfmr182899fac.4.1784928509681; Fri, 24 Jul 2026 14:28:29 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::5d97]) by smtp.gmail.com with ESMTPSA id 586e51a60fabf-457673d2dc1sm8183883fac.10.2026.07.24.14.28.29 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 24 Jul 2026 14:28:29 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Michal Sieron , Joshua Watt Subject: [bitbake-devel][PATCH 5/6] hashserv: server: Use streaming and queue API for upstream exist queries Date: Fri, 24 Jul 2026 15:25:13 -0600 Message-ID: <20260724212822.1165552-6-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> MIME-Version: 1.0 List-Id: X-Webhook-Received: from 45-33-107-173.ip.linodeusercontent.com [45.33.107.173] by aws-us-west-2-korg-lkml-1.web.codeaurora.org with HTTPS for ; Fri, 24 Jul 2026 21:28:31 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19855 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 --- 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 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 From patchwork Fri Jul 24 21:25:14 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93477 Return-Path: X-Spam-Checker-Version: SpamAssassin 3.4.0 (2014-02-07) on aws-us-west-2-korg-lkml-1.web.codeaurora.org Received: from aws-us-west-2-korg-lkml-1.web.codeaurora.org (localhost.localdomain [127.0.0.1]) by smtp.lore.kernel.org (Postfix) with ESMTP id 3C13EC531C9 for ; Fri, 24 Jul 2026 21:28:41 +0000 (UTC) Received: from mail-oa1-f47.google.com (mail-oa1-f47.google.com [209.85.160.47]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.28981.1784928511376313270 for ; Fri, 24 Jul 2026 14:28:31 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=i67ktGzV; spf=pass (domain: gmail.com, ip: 209.85.160.47, mailfrom: jpewhacker@gmail.com) Received: by mail-oa1-f47.google.com with SMTP id 586e51a60fabf-4563ac048f4so494093fac.0 for ; Fri, 24 Jul 2026 14:28:31 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1784928510; x=1785533310; darn=lists.openembedded.org; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:from:to:cc:subject:date :message-id:reply-to:content-type; bh=ZpKwBLV5BxsmWBXTu0qLwGCa5/4EAPrutBLij9MoUVY=; b=i67ktGzVrt62taxIEl7VEqYJUhsqiIT0+ChpmXctLPoP4U/tU7s4cENkOHxZUkI/X9 8Lv8wk7+cc3nh2IxFm7AUS/Fs54ZdnkK8Z4ScrYtLHRAdTfDEeRmQslBSn0a5jnlGJLT lTXBPaGu9OvTfy1OfJIS8D2n4ovXusm16h/QaeWOk6BPKSnjUL63u0bLl7B64E8Fm2AD Z7cNn3UD0JhVhsu6eB+bkei3ly0/TO8EUzbQqWmc0mI2E7a0ZbxEdxNtkGS6sZ5ZwIAg xuK//wyjH+Tw6RJ7yQe4YR22FXd51KNT1bqrCvQDV4U3t7uz1a6WdUtnx5bThaPJT88A aLHw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1784928510; x=1785533310; h=content-transfer-encoding:mime-version:references:in-reply-to :message-id:date:subject:cc:to:from:x-gm-gg:x-gm-message-state:from :to:cc:subject:date:message-id:reply-to:content-type; bh=ZpKwBLV5BxsmWBXTu0qLwGCa5/4EAPrutBLij9MoUVY=; b=Ars9nu2njzg4FdriQqQp7Z658wduE8q+MpOIKCTJQN4ooPtlcpcqRLixCsyItgRCq7 2u2Bn35BZAVHeBNTagEfAb9uhBBuPmnEnc5eE7VCSQgMe6LEBnQ8ryeoqGx7cFAQu9gW 32RUBZasT0Pj0jWMXGVU++NI4bKm0btJIRQCW2HX/u4WmBE5c66XnOan1vXfTrGppATK w13uwOZEd9+Tey4JPSxR283F7JY5D/xf9Hrlfwx4sAcCH9CrtocBDJdoeDKib4tX7x2i hA5XjUNzy5NWrXczonNS6nwUBvYgoIECfxejTSy5eGJ1oylv8ZP0dNoOaqZ9d9LnOzq1 q3vQ== X-Gm-Message-State: AOJu0YxdeZlx+wMX9b/F52M0Eb4HsSPfvVyLV7lCM+eAYcyWPYRzEY2l r6Kllq682EaeNV7MlFR41MhFD44C8wwohqfyqkNmvh3Frdjo6Kyd7+f2+GZqHA== X-Gm-Gg: AR+sD11jaUTI7x0iaaNf5bjDE578a9TfJau/H3Dk23Yf/MlYR3YjJdKxSfx8olJA6b0 k1lQzdV441npOXG09p6PzyoxDcP3FyxcaaXesod4uU3h9E5+xEcn5bKUQvWBU/vT+fkIrJ52vJJ 2cny5GXrE51tGlKH+2dffATj69YIlLuSKyBrrkhbibpPyJsaYf4axuM7OeLvjULHuK+2lV93crk vLKtr89eOvlfH+m1YIlpGIz+MCX4/nhSCwGYcOjC+cE+DPTO7tS7ZafRXxIqfpVAw15ns97BLte 65E+r/uQOiVOhzXNXscacwncvD53kt5Ido2gExDwuFDXJKPwgsoqvMIkPUHsPNR/w02CyuiAJYV Z9J/hCLNLrme8opHayeTQiv6a0iZDgQjiOQExoa6Y1Ik6RGIGH2d3AW8T2TC2QmE3MLp2QtKL9Q g= X-Received: by 2002:a05:6870:34d:b0:456:ae13:5583 with SMTP id 586e51a60fabf-457f27758ffmr179363fac.31.1784928510535; Fri, 24 Jul 2026 14:28:30 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::5d97]) by smtp.gmail.com with ESMTPSA id 586e51a60fabf-457673d2dc1sm8183883fac.10.2026.07.24.14.28.29 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 24 Jul 2026 14:28:30 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Michal Sieron , Joshua Watt Subject: [bitbake-devel][PATCH 6/6] hashserv: tests: Add test for upstream pipelining Date: Fri, 24 Jul 2026 15:25:14 -0600 Message-ID: <20260724212822.1165552-7-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260724212822.1165552-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> MIME-Version: 1.0 List-Id: X-Webhook-Received: from 45-33-107-173.ip.linodeusercontent.com [45.33.107.173] by aws-us-west-2-korg-lkml-1.web.codeaurora.org with HTTPS for ; Fri, 24 Jul 2026 21:28:41 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19856 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 Signed-off-by: Joshua Watt --- 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):