From patchwork Thu Jul 30 18:30:54 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93953 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 0968CC55179 for ; Thu, 30 Jul 2026 18:33:07 +0000 (UTC) Received: from mail-oi1-f177.google.com (mail-oi1-f177.google.com [209.85.167.177]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.18453.1785436377875696463 for ; Thu, 30 Jul 2026 11:32:57 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=V0ar03f8; spf=pass (domain: gmail.com, ip: 209.85.167.177, mailfrom: jpewhacker@gmail.com) Received: by mail-oi1-f177.google.com with SMTP id 5614622812f47-4a456e44e01so109744b6e.1 for ; Thu, 30 Jul 2026 11:32:57 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436377; x=1786041177; 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=dIIcAB25BFKFt5slf26rijE6SlRoyuK/mAAVE1tb+Wk=; b=V0ar03f8HmMNeJMChXEmrO2ozy02iPptm8sv4v70ykH+WaeozAl6gzio4oYdQ1DfDn O4k+WhPoGLZwpWa/zc4TQQQXvVNZFNBPRcUNNQV1IzVDzmoY9PcheoKLiMNxhXbu0qEB Tkb2uopAwesuzauOVLsCMSIlZLouFwhokBep+5Ra42e1oZN2Rou2CLUghjIm8EZmD4Ki d3+0Gak6DWxkybUlnU65f/8fxPSEL8dM+MOM8aAS2wW68uAGF60LNVr60Z0Mg2NFQEQI uaDBlgJ+ASwLhmL0nBhy9AA4XTwLzj1oNpKJSNlGkm5sXm5aUq3aW0rYikwTNem+yQQs 25yw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436377; x=1786041177; 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=dIIcAB25BFKFt5slf26rijE6SlRoyuK/mAAVE1tb+Wk=; b=cqFLSd41R/FGwztmSoUmjfVa60Loo6UIx0K73K6CxyGdEdkLqlmklYVOVccRx2JPwY 3JwsQlJWEp1o4mQuXdjaJkt+QOM9yhNeO9PKS8rQEA65GTJzYkDOe46yC5Ig94ZRRHC+ goItt0+kCgnCu3SBW0aelyfenbHFjmJF2863W81q/y54+e6lImkkcrPIda6Isq49Zuvg w4GPNCfldmrgkkAoLvi2RCmduQ6tR89QqQuaJv8nYoNbwTOAgx/3FEJ63di3A6Mq/oZd iju4OP/kHnOarRPS1Cc+W4ptYpg9tmEL16y5wgaRigUng60cArw1dOeD4nHk3mqdo1II Ihdg== X-Gm-Message-State: AOJu0Yw7Arf8GFbCEJXdKU1TqBXcHtOujBlSVmzPQDoOCz+ytkdbd6tV 0FW0WxnSRvYT0vw5EX4n08iB7suHVtExCwaVuZmehNR7PF+xJl07dR8Gfpg85w== X-Gm-Gg: AR+sD13QMO+ONL8yUYpCBT5fUN6H5DCBAWd8h6FIotBx4jLi67AAkYhlc19yKvIr30x CLuMTISZ5lidHHySWONI4hjPYNsapnquUpt2A4HJfWEpawDIBp+7u8EM8SaRY9bJniKDaEPrqQu AGRTmXXEUwz3sfTFEEvIrllhxXrkI0D0xTEs18wJ5pknnEv6JL97mFSf46yIBi7QwUaefzkveoM urkkI5f+hvI2bZxQbb8hZ/c/iPqt3CRwyKqCqakUfBdiv3BfJ1crBh4vDVkgY9dCZKoQAPKmLF+ sq8Qc3KkGPE1vYuwxk80r8Ht3qnqHeKH5OzUMf6J+iSLNsk9afx32S3TvEMlKyXu/CEWlPi62JB hnN9IyEJFnaIW/0HrSek90W5n4IbsDv2QgDLZETGpb/x4B6Xhsw+loj4GixMiOCzQTInGgkliBp 9IGDumWsrDpAP8dL+KA+03RyIKHvEedm4tOD3HKrVYm8ZLaflJK6sdVpCZiQ== X-Received: by 2002:a05:6808:1687:b0:4a4:ed9a:e790 with SMTP id 5614622812f47-4ad9533b9b2mr1399100b6e.29.1785436376875; Thu, 30 Jul 2026 11:32:56 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.32.56 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:32:56 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 01/10] asyncrpc: Add Task Group Date: Thu, 30 Jul 2026 12:30:54 -0600 Message-ID: <20260730183254.793698-2-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:07 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19881 Adds a Task Group object which can be used to manage multiple asynchronous tasks and correctly cancel them if an exception is raised Signed-off-by: Joshua Watt --- lib/bb/asyncrpc/__init__.py | 1 + lib/bb/asyncrpc/taskgroup.py | 52 ++++++++++++++++++++++++++++++++++++ 2 files changed, 53 insertions(+) create mode 100644 lib/bb/asyncrpc/taskgroup.py diff --git a/lib/bb/asyncrpc/__init__.py b/lib/bb/asyncrpc/__init__.py index a4371643d..2f0956dfa 100644 --- a/lib/bb/asyncrpc/__init__.py +++ b/lib/bb/asyncrpc/__init__.py @@ -14,3 +14,4 @@ from .exceptions import ( ConnectionClosedError, InvokeError, ) +from .taskgroup import TaskGroup diff --git a/lib/bb/asyncrpc/taskgroup.py b/lib/bb/asyncrpc/taskgroup.py new file mode 100644 index 000000000..b61f40385 --- /dev/null +++ b/lib/bb/asyncrpc/taskgroup.py @@ -0,0 +1,52 @@ +# +# Copyright BitBake Contributors +# +# SPDX-License-Identifier: GPL-2.0-only +# +import asyncio +import logging + +logger = logging.getLogger("asyncio.TaskGroup") + + +class TaskGroup(object): + def __init__(self): + self._tasks = [] + + def create_task(self, coro, **kwargs): + self._tasks.append(asyncio.create_task(coro, **kwargs)) + + async def __aenter__(self): + return self + + async def __aexit__(self, exc_type, exc, tb): + try: + if exc is None: + while self._tasks: + done, pending = await asyncio.wait( + self._tasks, return_when=asyncio.FIRST_COMPLETED + ) + self._tasks = pending + + for t in done: + try: + await t + except asyncio.CancelledError: + pass + + finally: + for t in self._tasks: + t.cancel() + try: + await t + except: + # Ignore exceptions + pass + + return False + + @classmethod + async def run(cls, *coros): + async with cls() as group: + for c in coros: + group.create_task(c) From patchwork Thu Jul 30 18:30:55 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93946 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 82B01C5516F for ; Thu, 30 Jul 2026 18:33:06 +0000 (UTC) Received: from mail-oa1-f51.google.com (mail-oa1-f51.google.com [209.85.160.51]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.18816.1785436378754527288 for ; Thu, 30 Jul 2026 11:32:58 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=Dp/RyNzj; spf=pass (domain: gmail.com, ip: 209.85.160.51, mailfrom: jpewhacker@gmail.com) Received: by mail-oa1-f51.google.com with SMTP id 586e51a60fabf-45171f2f608so135675fac.1 for ; Thu, 30 Jul 2026 11:32:58 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436378; x=1786041178; 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=PI3rOpzqM993OnncotOYYrbT5CTP6U+3eCypg+8LKJM=; b=Dp/RyNzjADQZph9vSnG+/XlGWis3+6i32YPSxmf7havYs/5O6LEHMLmV1+SCFc2JgQ pE6Vglr1JhGDVBc82tz4CJtAb4FDpO7yQQE0w0mFnj8zJOaotC/cQHVc5F18BfaVUTjQ cbYJb6ZreHHnFzgr/9DCOukJy5paVByREfqpAyJ5pSfP0FUQkmsf8IXc1vgM1Q9h4hBU IRW+teFSPp/jqK53pCYK5A7XRWWwVsxBJcRBBXS+rArOcBFWgXDS3hPI75OnFLfKzo4c uwEhOKxTemqvnZ9Aws4gBpa2dOnOWlElSWzHzqZKjpef8qnK+GXWPHJpROs946114rxr dVeQ== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436378; x=1786041178; 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=PI3rOpzqM993OnncotOYYrbT5CTP6U+3eCypg+8LKJM=; b=ZNPxjGuww6ya0XhXGJhSTvdtyZDA5Cm+0Lzk7jjTzHxZRIDmHiZjmH+yyNGZYalngn XEwrEEd1MppifdZGB8C0wEKEJ5yEt69qyA6t56VosH5nXxNPXWBZvoUx/WBXVK9wYcQO CzNJPg39nYkMlMGD0465SHKFEoq+VDxvSdgDUI7yCiJV+RqGVvwUTu4mZo7e3ROY8bUn xjiGMPIqSxE96SgMUfV1tfPwLb9xNbERLzlEQvTVKi0K2dFCD/9H/i+scUczKaJBv56R W0H7p8YycEz7XtLyyPtP4v3GpnVtvtC15U7WpPNuqvi1M0RBpWx3O8GudiEs/wTeXh7h qDSA== X-Gm-Message-State: AOJu0Yyec44y8EZhx9G1SeE9jPHhOD8l3a5EiyMEmeNX5tkSYOhV9Ax1 4T9wrURbo/QImBl4/JWetQiEgrM+RsGQLHPIWH7C8VHX7rVLZn9HDXt6HsrjoA== X-Gm-Gg: AR+sD13mMuSZxNgGerrR287GwTt/rDBCLv1GWqambHVqdm3PobF2N77e7b3sWXd4PN1 v8XIxEKdib8v+svBV8dR5FMcQEpD0wDPYL3ayF55F3gbCQi8NbqA9Gkeunc23UdKeuu6mGd6rE6 EfmSLldzGmgVUt0ye2uRgeckmnXyI/2vpt+YG0+BkCro0Oev7W+yTG108cnhJFMZLUdqD/Tc2zk 8BX+/SQwPnjqnFjMaXzSSi2lyH8ZxtfKUq2qtKQBu3WIyFdwQ1Y1ONqIvj1hTQVVwQ0ivfvEjxF vjXWf4GxjRWfdjwKA9zsDN9ij4kaHUBUZPhY1sQO0f7LowMRlextIP8sneiJpQD7X/+TLZU74OE KfebbkCsTceGeDzrYh20eQ5HRu8jJv+Sfj5G7Nt8tJYSjDeZ+iY3Td9Asi2xRDODCeiWXF+VBqd iw1WuOoZNMFnEx879m64kZ590fI4olvMGkDHnK5cxeKJLljK/6iCLcGeifeQ== X-Received: by 2002:a05:6808:179b:b0:4a4:934:b316 with SMTP id 5614622812f47-4ad878cae8cmr2587508b6e.21.1785436377790; Thu, 30 Jul 2026 11:32:57 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.32.56 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:32:57 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 02/10] asyncrpc: serv: Use Task Group Date: Thu, 30 Jul 2026 12:30:55 -0600 Message-ID: <20260730183254.793698-3-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:06 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19882 Uses a task group instead of asyncio.gather(). The task group ensures that all tasks are canceled if an exception occurs, which asyncio.gather() does not. Signed-off-by: Joshua Watt --- lib/bb/asyncrpc/serv.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/bb/asyncrpc/serv.py b/lib/bb/asyncrpc/serv.py index bd1aded8d..d3b1c6c35 100644 --- a/lib/bb/asyncrpc/serv.py +++ b/lib/bb/asyncrpc/serv.py @@ -334,7 +334,7 @@ class AsyncServer(object): self.loop.add_signal_handler(signal.SIGQUIT, self.signal_handler) signal.pthread_sigmask(signal.SIG_UNBLOCK, [signal.SIGTERM]) - self.loop.run_until_complete(asyncio.gather(*tasks)) + self.loop.run_until_complete(bb.asyncrpc.TaskGroup.run(*tasks)) self.logger.debug("Server shutting down") finally: From patchwork Thu Jul 30 18:30:56 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93949 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 D1A7FC55174 for ; Thu, 30 Jul 2026 18:33:06 +0000 (UTC) Received: from mail-oi1-f170.google.com (mail-oi1-f170.google.com [209.85.167.170]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.18818.1785436379523788690 for ; Thu, 30 Jul 2026 11:32:59 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=XDp2zmsc; spf=pass (domain: gmail.com, ip: 209.85.167.170, mailfrom: jpewhacker@gmail.com) Received: by mail-oi1-f170.google.com with SMTP id 5614622812f47-4a427e628a9so91907b6e.0 for ; Thu, 30 Jul 2026 11:32:59 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436379; x=1786041179; 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=NDFCp80fnVzRsGRXZ7ucLzC6zLtSq5ASKZf+80e9RvA=; b=XDp2zmscYJ3UaV1LfM5Wh7JPygsH78G3i2hZ6R9f1oK1Mfwn3lzs1AFWEAB5o8fXDB pPl7SJxZ5RKQp1y7X3igzS2LFIqTIpoTxRCAAMO/cTNMsbTrRX+/zKD4xc5rehdw/62P xaldjOAZUOLcITmx7fExEt2yPhYyG4zZRVFRRcBfh4a4J6DrOKBYL8eu63wPFZ4QUDyw 0ltMVzxR3e3k7YtEaSUsjFZpNyLnFY3Z+8IfxROV10i3jUBTuSvBMW9FHGsbkrC+6DIj Vy24zm8Qg7lHScWGlDeJJ/BhfZ3OmTxJUq+f4SQcDyE1EClh1QF6G1Srid7GTo0DhDHt OR3Q== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436379; x=1786041179; 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=NDFCp80fnVzRsGRXZ7ucLzC6zLtSq5ASKZf+80e9RvA=; b=gv3YLuX0gNl6cP8osk95dexUl0V3dD4vq3h+0Wb1mjR6UmmcHhNO2cNiRYB3PUU+9W SzoN+y8FNiThHgolB9r0r367Ayk7yo6yTTg0+xqDBzVR5o9uQFvQrnIDQlEzcZSojTa+ 13Qi7g3lWN+KEH/hJ/G4rpr45JXDxpfig8b1zvRVsD/Lo6R6J41zdeMHwCcDje39U1dO GXP+fiMGrXdwB723GVVtxPR54ZNFPbr0a9e1GX1SbZzXEVjruSRg0r7mGDNOXgZPKFrD c3GbLWZLvxuP2IHViHL3kqc//QxaI+XxLCCtpryURYARKaDcxrZ0hRnZXrJo5QnmpxbT nFDA== X-Gm-Message-State: AOJu0YxDrEaIDdasyl4/g+fFGBEAMfnkhwbq+jSvcLQP8tD5DXMGyBr7 5zIh+7QMNLBXaUjm82KK5czqs7BrKYV6DPFJPSMBcIXTvVH/CruEv2vYI+y97w== X-Gm-Gg: AR+sD13MAQOtP9SF6RWGBCy4EuoQXyRi9iXKboecnBQGQHMFkcQNT0Ck8zo6r8rGqt/ ItO3yX6Q6YUio4bPOTpfLtDZZylKkwxqqJWdwI3cSaGUDrlRY1hk+03mvnnezGktzM/iSw6EOc1 bxkRbjJ8kYDK4H0Fk5Y5GrqETgslTcXqPvP2swkb78+VsMzb90xgFa04thDuDey/UDKuLkz3x1A Z20rbAOIvXHjiJipcXJIaFv0kH2wnuwqn/p08NRR4TosNGgtUdmLLxOSwzQsB7bbDd07XYFcce/ LZgC0qXziq6oAZEmi8RMg0o8rQTA+mAnvd7E/8HUO/VBjAkQTvd/sG98ZpKdyhoVAx1LPXWKgZz 4unkxS2kGaJ5nWsktFmxVrGqGwlevkOx+FyEKW9MtqzyTPU+hrDZeBfp6opEXsd0ZB6ly6udYCR osn9LnQ22W4Bxv1sqzjBfeRpsVq9t77UCuFoid4YWi5JdvbRETwspUrM1lFQ== X-Received: by 2002:a05:6808:f14:b0:496:2b3:ae71 with SMTP id 5614622812f47-4ad8775789dmr2861803b6e.18.1785436378563; Thu, 30 Jul 2026 11:32:58 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.32.57 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:32:58 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 03/10] asyncrpc: serv: Cancel all clients on server stop Date: Thu, 30 Jul 2026 12:30:56 -0600 Message-ID: <20260730183254.793698-4-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:06 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19883 When the server is stopped, all clients should be stopped, otherwise the server will hang waiting for a disconnect. Signed-off-by: Joshua Watt --- lib/bb/asyncrpc/serv.py | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/lib/bb/asyncrpc/serv.py b/lib/bb/asyncrpc/serv.py index d3b1c6c35..dd77ed0a4 100644 --- a/lib/bb/asyncrpc/serv.py +++ b/lib/bb/asyncrpc/serv.py @@ -56,7 +56,7 @@ class AsyncServerConnection(object): if not client_protocol: return - (client_proto_name, client_proto_version) = client_protocol.split() + client_proto_name, client_proto_version = client_protocol.split() if client_proto_name != self.proto_name: self.logger.debug("Rejecting invalid protocol %s" % (self.proto_name)) return @@ -123,6 +123,7 @@ class StreamServer(object): self.handler = handler self.logger = logger self.closed = False + self.clients = [] async def handle_stream_client(self, reader, writer): # writer.transport.set_write_buffer_limits(0) @@ -131,10 +132,16 @@ class StreamServer(object): await socket.close() return - await self.handler(socket) + self.clients.append(socket) + try: + await self.handler(socket) + finally: + self.clients.remove(socket) async def stop(self): self.closed = True + for socket in self.clients: + await socket.close() class TCPStreamServer(StreamServer): From patchwork Thu Jul 30 18:30:57 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93951 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 03381C55177 for ; Thu, 30 Jul 2026 18:33:07 +0000 (UTC) Received: from mail-oi1-f170.google.com (mail-oi1-f170.google.com [209.85.167.170]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.18455.1785436380078265277 for ; Thu, 30 Jul 2026 11:33:00 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=IouSGcwu; spf=pass (domain: gmail.com, ip: 209.85.167.170, mailfrom: jpewhacker@gmail.com) Received: by mail-oi1-f170.google.com with SMTP id 5614622812f47-4a4c6081f9fso63628b6e.3 for ; Thu, 30 Jul 2026 11:33:00 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436379; x=1786041179; 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=EKnVVvQLGYQVBqqlwMzNNMoh+34r7ffbwuqLdEkaRWk=; b=IouSGcwuhLqgFvpexVInoxdn44L8oVC8OVE92NSic7Jc4TI4i/LoSl6hiHENlaCiZu xo+6uXcU4Jai98Kkl6+y1OMttc8Ez0b+46VvW3wK2eU+e9jr1hGK3bez0t2h7Qiy+xqB dcC6y/V3IQk5pDabgHouEU+JFm8C8u68L0dj5WZ7tTcLAyFbpz+iniieM34EMCZI4Orr NZ5it7BDUqZ47nIBkzy933+JEgVvTPIQ6T5VJIOPXfInBDZeTJdoxIDlwu+JwoaRQE0H 0sekj5jDwvtHQKPRqoAc7AYeOMmF/mgvjKfxaoCVAEjEC+f3cGiwajwo/PWlAiIJGZ7d Dklg== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436379; x=1786041179; 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=EKnVVvQLGYQVBqqlwMzNNMoh+34r7ffbwuqLdEkaRWk=; b=ZfIl6p40+8HDpipu80bS9TcdNKQXY91p+CpTRJQaYginqw9JvWzOB8hNQ+gzYm4Yhu 1s0hezqZyAZJ0DLYHD7xVzAtZZjBGdPKLXNteVpufceYHQhvw8f7B026/gQz3ySPwjg1 uxaZp51Oi2PDm8/HPwWVakJF/1/HZr+i5ZHLFFNnWoAks6tM6Av2e51Ov/2vhRALygjp 45trVl/lZ6lmsE662GlKTW4U0Ua8+WtpdBFgI11TOD0RXbxjeGqawc+ULKwwbGWD+3ds 5fHT/c/XJQ1MTpMT7ngJT3EWXMvUUfmQz66qO0mAuri1paekvCW9awgiMfAn8Yc1rQhT R6HQ== X-Gm-Message-State: AOJu0Yw7uzXBVgOiWDqL8qrxdJZkXQkgRCaSXiVqsIpeSML4eGr/KE5C NpgALVKuxbOH5eZHHp48C0Gry4AET/nkt4CmsMSPO5Ie7MiVaKsG3MTmgDzWOw== X-Gm-Gg: AR+sD1261zt9+0fmONZufmbPyrG3Y3/GeCTHctNwnqn/fIdVGRk9/E5pPPovsrtyYhV 84jdJ+uyIYqlt+MSc6U6R9AfwFQ5SNqF/y0gpKrhnOK6V1CoEN1nrVKN9h6ySB2l/ZfP0NZdFJi OSoYuUs1HfwRgLbLrtdvi2cT37i1lH1AU+6/wY5z5wZt1gAq69sCvsRNfGcmtOEEVwZvFlMohnR rcXeop6y7VR/iDRcBk01VGvw6w9h1weJVVakgvI1Y1jTlI3VifNlJZRofCIasUTRK7XH3nW7Me1 Q9GKKg4g7uyIB58wvpkia6N1P633FBi3OVs6lOybACsUYWtrh1TSaaNVLzEVhxwEWPsH/T9/2ld wkMcGrA0HxySA5bYqw/OWpbbPb9Z1l3QmR+IV/nzepBlLcwS40gI+1M5Q7gFkOAWPriylrXky6k 5bESJdm0XHxeYLozAgyStlXKHw7vNKUiaqey18qjzHCSM7VeCJkk6hUwlQKg== X-Received: by 2002:a05:6808:f13:b0:490:9f92:d00b with SMTP id 5614622812f47-4ad879701ddmr2746549b6e.29.1785436379213; Thu, 30 Jul 2026 11:32:59 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.32.58 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:32:58 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 04/10] hashserv: tests: Improve test logging Date: Thu, 30 Jul 2026 12:30:57 -0600 Message-ID: <20260730183254.793698-5-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:07 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19884 Log client messages to a file to aid in test debugging Signed-off-by: Joshua Watt --- lib/hashserv/tests.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/lib/hashserv/tests.py b/lib/hashserv/tests.py index 7c736d6cc..551b9e298 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -30,7 +30,7 @@ BIN_DIR = THIS_DIR.parent.parent / "bin" def server_prefunc(server, idx): logging.basicConfig(level=logging.DEBUG, filename='bbhashserv-%d.log' % idx, filemode='w', - format='%(levelname)s %(filename)s:%(lineno)d %(message)s') + format='%(levelname)s %(filename)s:%(lineno)d %(message)s', force=True) server.logger.debug("Running server %d" % idx) sys.stdout = open('bbhashserv-stdout-%d.log' % idx, 'w') sys.stderr = sys.stdout @@ -93,6 +93,8 @@ class HashEquivalenceTestSetup(object): return self.start_client(self.auth_server_address, user["username"], user["token"]) def setUp(self): + logging.basicConfig(level=logging.DEBUG, filename='bbhashtest.log', filemode='w', + format='%(levelname)s %(filename)s:%(lineno)d %(message)s') self.temp_dir = tempfile.TemporaryDirectory(prefix='bb-hashserv') self.addCleanup(self.temp_dir.cleanup) From patchwork Thu Jul 30 18:30:58 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93950 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 AB83AC55173 for ; Thu, 30 Jul 2026 18:33:06 +0000 (UTC) Received: from mail-oi1-f170.google.com (mail-oi1-f170.google.com [209.85.167.170]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.18457.1785436380971840477 for ; Thu, 30 Jul 2026 11:33:01 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=j6W2530q; spf=pass (domain: gmail.com, ip: 209.85.167.170, mailfrom: jpewhacker@gmail.com) Received: by mail-oi1-f170.google.com with SMTP id 5614622812f47-497deab2d66so18580b6e.0 for ; Thu, 30 Jul 2026 11:33:00 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436380; x=1786041180; 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=02F4y5xS9K3M1HHYzG3Z8NSjYZKcvzoEnWK/6Xuh8NQ=; b=j6W2530q40d5YXuFuaQWOm0kHpcucE+3EedBtxlx1AX3Z9z4HbQsQqH0Sme1m3rG0l vTW/d+nKVE877ne4HCu9jTcTXmX+5jiI+vk+xUhla9zVFctkcw27jzEJwv0aSkqssKOE ehztv5tluw4XKVjJfvTiemfq9bs9IvxCCKtWvwvAWBdp/ravRpu5Wv2OwG98SuQFzU9J 9zsrSO1yN9d4vL0dXh9myXivEymVC5gStz+cEHWDcW5hLfdELNObdRmTA9JXnFAiZjwp nyfJXbWPkQaUWs2lX7Abm/om7IDeaoU8Zgzg2JojmPxx4QvKBuTk6HfqjIf99uoUF3h7 R7mA== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436380; x=1786041180; 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=02F4y5xS9K3M1HHYzG3Z8NSjYZKcvzoEnWK/6Xuh8NQ=; b=VSfUAt1amm7Jk4z2v6zbJXAOGojbKXP8xuTvDuQjpfpZ1hFpzGclRaHK+Suzh2TKGj o/gRj7kWrXKPRvFKQs2EfCZMudguNVqkIr+E9R7DGBPlJXnjNryCO/xvbvpztpXnkGhA FEnM8YqUpbu2OJy5xs1d7l9gI9HGlo81fMRuyylkFVqMVEstAstx2AqYe/MnwDi222zH IcBd1gEQfnz53QjO41lfkGLvRmvLyGCbSvUEr7jrrhAWbSVUYxhyJS4qjrZUhqcLnMyi iAD/2kbL8QHDefU6duAjNdCdDSc1qXIps7iBBpYsYSYY/zTRF0taN++ZM7MoVk6gQxwv hFHg== X-Gm-Message-State: AOJu0YxpU05LyAXO45nEisCNxcUQkV2YSndqcDhBsa+TQcoaJvhANPVx oBis77XxsXe+FjjGuD0GVR1IVm4Zd7c99n4TJkeXCycNA5fDNXHNQLwAop3K1A== X-Gm-Gg: AR+sD10s9OuMLBlvzZoK+SQUlW18BNvc4/jtymVs6D/aqGEM5jEKkPhNPazvtA+uiL4 hhZQiRBTshGpZHf9wIg2OgblxSATJsYxREcbVUOXCQjTI61rQSXcPjTPjsyMS2osXKZ2LYIMKnE 8ZwExEukaHhe8Wlixva1hCrTMhnO9mnHIBB9izYqjsa+C5NmqPlaKm45qsaA2xtVTWgB6V2oKd/ e+fOd+NpscDblgDfXXFs0p43XbYE/96ZJqmDIrtw8+cZHPSbmR5pG30q/cZDg3Gq5bzTPUITy9b uZNKKShXFND9mq4Lyqn7DN9UaHKPBR/nHaqToaZ3uHvC4oYduyufSWrD9Zl6LTk0gXsNRfd6WLK 3aFk+91xEl5gUEUPuU0Uj3L+/7qDYPou022ZwKeSTl6qtMxpjeh59IExh+B6Irmd0vMOogDD0Hy 5SLTiOnhtEr2foaC2vigFahzNhhkcU9Vrvkm8C/W6eRe/yh43okJwPB4vFFg== X-Received: by 2002:a05:6808:c14b:b0:4a4:66f8:1279 with SMTP id 5614622812f47-4ad958fdd3cmr1033193b6e.1.1785436380055; Thu, 30 Jul 2026 11:33:00 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.32.59 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:32:59 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 05/10] hashserv: client: Add asynchronous streaming API Date: Thu, 30 Jul 2026 12:30:58 -0600 Message-ID: <20260730183254.793698-6-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:06 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19885 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 | 349 ++++++++++++++++++++++++++++++++--------- lib/hashserv/tests.py | 4 +- 3 files changed, 275 insertions(+), 80 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..b447fb259 100644 --- a/lib/hashserv/client.py +++ b/lib/hashserv/client.py @@ -8,70 +8,252 @@ 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 bb.asyncrpc.TaskGroup.run(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): - self.done = False + def __init__(self, send_queue, recv_queue): + self.send_queue = send_queue + self.recv_queue = recv_queue + self.fill_done = False + self.send_done = False self.cond = asyncio.Condition() self.pending = [] - self.results = [] self.sent_count = 0 + self.recv_count = 0 + self.item = None async def recv(self, socket): while True: async with self.cond: - await self.cond.wait_for(lambda: self.pending or self.done) - + await self.cond.wait_for(lambda: self.pending or self.send_done) if not self.pending: - if self.done: + if self.send_done: return continue - r = await socket.recv() - self.results.append(r) + m = await socket.recv() + await self.recv_queue.put(m) async with self.cond: + self.recv_count += 1 self.pending.pop(0) - async def send(self, socket, msgs): - 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 + async def fill(self): + async for m in self.send_queue: + async with self.cond: + # Wait for item to be consumed + await self.cond.wait_for(lambda: self.item is None) + self.item = m + self.cond.notify_all() + + async with self.cond: + self.fill_done = True + self.cond.notify_all() + + async def send(self, socket): + # 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 + async with self.cond: for m in self.pending: await socket.send(m) - for m in msgs: - # 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: - self.pending.append(m) - self.cond.notify() - self.sent_count += 1 - - await socket.send(m) - - finally: + while True: async with self.cond: - self.done = True - self.cond.notify() + await self.cond.wait_for( + lambda: self.item is not None or self.fill_done + ) + if self.item is None: + if self.fill_done: + self.send_done = True + self.cond.notify_all() + return + continue - async def process(self, socket, msgs): - await asyncio.gather( - self.recv(socket), - self.send(socket, msgs), - ) + m = self.item - if len(self.results) != self.sent_count: - raise ValueError( - f"Expected result count {len(self.results)}. Expected {self.sent_count}" - ) + await socket.send(m) - return self.results + async with self.cond: + self.item = None + self.pending.append(m) + self.sent_count += 1 + self.cond.notify_all() + + async def stream(self, socket): + await bb.asyncrpc.TaskGroup.run(self.send(socket), self.recv(socket)) + + def check(self): + if self.sent_count != self.recv_count: + raise ConnectionError( + f"Sent {self.sent_count} messages but only received {self.recv_count}" + ) class AsyncClient(bb.asyncrpc.AsyncClient): @@ -98,29 +280,34 @@ 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) - - return await self._send_wrapper(proc) + await b.stream(self.socket) + + async def process(): + try: + await self._send_wrapper(proc) + finally: + await recv_queue.done() + + # Create background process to process messages + async with bb.asyncrpc.TaskGroup() as group: + group.create_task(process()) + group.create_task(b.fill()) + try: + yield AsyncPipe(send_queue, recv_queue) + b.check() + except AsyncQueue.Shutdown as e: + pass + finally: + await send_queue.done() + await recv_queue.done() async def invoke(self, *args, skip_mode=False, **kwargs): # It's OK if connection errors cause a failure here, because the mode @@ -173,15 +360,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 +394,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 +505,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 +567,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 551b9e298..e24bdcacb 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -1057,7 +1057,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' @@ -1080,7 +1080,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 Thu Jul 30 18:30:59 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93948 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 9160DC55172 for ; Thu, 30 Jul 2026 18:33:06 +0000 (UTC) Received: from mail-oi1-f172.google.com (mail-oi1-f172.google.com [209.85.167.172]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.18458.1785436381830164031 for ; Thu, 30 Jul 2026 11:33:01 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=PoCJYjwc; spf=pass (domain: gmail.com, ip: 209.85.167.172, mailfrom: jpewhacker@gmail.com) Received: by mail-oi1-f172.google.com with SMTP id 5614622812f47-495b98b4f6aso98261b6e.2 for ; Thu, 30 Jul 2026 11:33:01 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436381; x=1786041181; 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=GJ9d+Yy8ysHYMM9q9YAN46OY943/I3hmMwBY4dTdLzs=; b=PoCJYjwcnwZkt4XV0ymEzG5hH9h2gTZA2mauaMqCeXlv2eApGOAeuB/onMqBOT/0u+ 0tSvqKljRc67iOU1NLp0UguaYIPqKmW9XBxDomuL0kjVnh050kdRsTfCIglpAek9yK2X VaVvQvnkHrK3QF1WZVnaA+WOqiRwGY+LmkVWngHtGDgS5eo8xQ5a9wJ8nqvMw/TZbV6v dA0P7YIag1/B+Opm7gxZizRqR/kGUVNRJgCAPGL23foaA1Xf4L0R8Wh/VT4QM6u/gzD8 JuRph/kZ5YILdYgpc2wKwzUhzU310WQPgvgbbLiRvZDImOZOiShYSefPUyvyCCcRGXC0 RiDw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436381; x=1786041181; 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=GJ9d+Yy8ysHYMM9q9YAN46OY943/I3hmMwBY4dTdLzs=; b=QQLmqL5Bqev80n9Q1bX3eT/lL65ezMX/Zs0OEG/ZCi9ffTuI9asD3onyAE+71413tt cIOJY9s1We2T0BNjRJU98m3Y1GDYgWYXK/xKvb1boOTPFp7SmxQa/GYL6VIIGHRxpQpE URxFR7Z54+XKCbvtC2YJ7PZXOg3ugZ4qfOJRTc6EKfIBGrpX/Ku73E4Xwo43tcAn3Soz BDrd1LW2guIz6bvOrjQvO84KB4nWG2gFDR3+oAiffZ9ZudPyM/8o6JM6qeXMn+6kTzNX r4NWT/LQ+Y5s6/UfW18tcNNI5HhXe8kcAlbY/81S6K6RT8jUuGDpj2dmvpk8h1zsQ7EI 3Q8g== X-Gm-Message-State: AOJu0YzaYWJQ6m22KgFHdGEPxv9ciefUjBGsLnNyz7oO0aS6KdSf9C/Q lZxiVBEqAslo5znpsWNRdUNhoAmnrRLDhzP5i0elXFLIAC+UUwbsInYMmpoUbw== X-Gm-Gg: AR+sD12d/c8qXFsEl5cO3M4xOwy5i0N11o/KGssRoFbz8lZfRf0vTa9KIHTgFUUQEBu baDtQKDu2LG/JBo3rrby+EMvfML6ylpUHc4TjteczbdgZdlOJ3ErrWo54H5cyKBy3QvGYnoRdbW BbL2+j/1aydp2wr8OcsafnRUt/QwFci3SxSOYd9vCDHhEXFbTRHHGZUtfg713lLT1SzcK+9ttJ4 3HnjOyGjuElgIinWegUiB5gAXxe+DpM1+XJBef8LRqfy1i8qOdbpUx8wm/MigrvwgHSv9PW7+E8 mr6NOq4krkzmQP0JId/ikpY1Nc7m99+o6+qa2dHzVPpz6SgnGp21Pcw5Z3dkS0i6sStWcCXJV6Y g+tzIx4dmA5ge+Wl88yoo9kDzNzZnFI2lzUPiKawoq44gFDoGRHYDAuiBvBcutzECFWBySdn+3E CiSL6JG9NZOhY+eb67MvqCXV8DhSw7+z8b90bjRA6dBaXsQtm86f5hfkJang== X-Received: by 2002:a05:6808:309b:b0:495:f79d:b08e with SMTP id 5614622812f47-4ad878c4c59mr3125464b6e.21.1785436380881; Thu, 30 Jul 2026 11:33:00 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.33.00 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:33:00 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 06/10] hashserv: server: Add queued streaming API Date: Thu, 30 Jul 2026 12:30:59 -0600 Message-ID: <20260730183254.793698-7-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:06 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19886 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 | 96 +++++++++++++++++++++++++++++++++--------- 1 file changed, 75 insertions(+), 21 deletions(-) diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py index e7e79196f..0fa81e85d 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -229,6 +229,54 @@ 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): + if m is None: + async with self.cond: + self.done = True + self.cond.notify_all() + 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() + + class ServerClient(bb.asyncrpc.AsyncServerConnection): def __init__(self, socket, server): super().__init__(socket, "OEHASHEQUIV", server.logger) @@ -390,35 +438,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 bb.asyncrpc.TaskGroup.run(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 Thu Jul 30 18:31:01 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93947 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 9E50CC55171 for ; Thu, 30 Jul 2026 18:33:06 +0000 (UTC) Received: from mail-oa1-f49.google.com (mail-oa1-f49.google.com [209.85.160.49]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.18460.1785436383590465317 for ; Thu, 30 Jul 2026 11:33:03 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=Zfsf+Zgb; spf=pass (domain: gmail.com, ip: 209.85.160.49, mailfrom: jpewhacker@gmail.com) Received: by mail-oa1-f49.google.com with SMTP id 586e51a60fabf-45133d2974fso145932fac.1 for ; Thu, 30 Jul 2026 11:33:03 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436383; x=1786041183; 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=qqrE/iXHsvexr2oIZF73UCyKJJ2N2jIRTQf2CjxnnlA=; b=Zfsf+ZgbiQ+ejP+HGYUdWhJIQBWaK6gMWbIJVE57LrdJea78VM7RMb4RwQnCNZtiW7 XY109JUJpLWeeJoO41FqAGyRTWOK/lV7jZgThozoh5q14OYAoIzApc6W2Z06dbUPDyW1 LdJX02OD7Ezo4WJZ7c8/rZP8XeuAWbtBEqO8X/xlHL3vJCfSYXeFTceT3eVYNH1t7p1U QA1O/LPRZwYM/olnqnplCb5EM/4oDbhmBxDUgfEsJHsMn3v0Nt5a30qpY9oQuVFbK7fz GWa5Zua0jF2EhLK0584XM0rTn2shuU88sToWXti3WnH/AFSXPHQSgzW7QQSDdGLhTbOf oFoA== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436383; x=1786041183; 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=qqrE/iXHsvexr2oIZF73UCyKJJ2N2jIRTQf2CjxnnlA=; b=rYLHWAQcp4WXqalWGcHJ++pt5AOSa0g1GtVoZdOMTuPrchAdufOPOmziJYwHihv2fk roSsAFbQ6KTopFI0/4zY8fTAucOK1+OVIG3okYT5ic//9OIApEyUWGgR3WN9K52SuW+/ wro/TkD+YDYWPzCCZVJj6s979YSZ7nvgZ827KAdWDOB7n+fojCKHnuFQJXeP4ujI6i7T oo13kauVrmfKXhZvK/eHW5/ZYXGgZ0lutLlnps4l2nYEoI/GB4o6Ojj9KrSmwAfh44r0 SncP7qveCHhhqjfn32v/8W8FfRHgyDuPyXe556u7+c85j55AYuckH+itRtrlakgVlwFJ rWvQ== X-Gm-Message-State: AOJu0YwwVbEsgZfOk5KeOSJusgkat08Af8a3+zpvjDQf40Q+nNvSKuc7 hI7U+eo2ctSVv/yjbTCZxPB59fMFByJHKdokIxMmPZvILDDdLtcIA7MHzL2bKQ== X-Gm-Gg: AR+sD13/rXFCnR92cjoMQz7oVGKqcAJg3RgrYY+tuCZytBOVkVwPls5GSCTYwabXKo7 801BQ00tsz7DSGsE3Lgit4wfhCpwRoK/8vJcouhiY4LqXUY/jdNLARJCdbwePyqnpWvpMlpB11p fAGRPyb8a1l5rIKGNlPicTjU1B4RUyrIOsgSjE3LC1JToWfykT99MxoyOeIaAEkyDzHP2tOAhY+ kaDBPDXHO/IsgSOk5Cx7lKigmbsPoM4rkTFs0luOgalIcrbTXrlrNjVSqxeObRL0t8uC9hBRIJV /B36z5b7dcDtxV06Jj8tkLamCZlPAeAdgoRd6mMm8xiR4F8b0oJAs49dko217/9X4ElvKeM5xtK crMtSruZLy9kDzuGrQBNIlZZeExXciedpd4OkRaDpSYriFF+efmrGMRx0qN630NpQgRXLFOy6L6 n0nTGtqG0H9AFRLFrBc0y2NL0BzlxgAfEIA4A/CP/tbVTafBa8FvwOoUSo8Q== X-Received: by 2002:a05:6808:1246:b0:497:8f1:df07 with SMTP id 5614622812f47-4ad87782704mr2593976b6e.7.1785436382710; Thu, 30 Jul 2026 11:33:02 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.33.02 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:33:02 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v2 08/10] hashserv: server: Use streaming and queue API for upstream exist queries Date: Thu, 30 Jul 2026 12:31:01 -0600 Message-ID: <20260730183254.793698-9-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:06 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19888 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 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 From patchwork Thu Jul 30 18:31:03 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 93952 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 E00DAC55176 for ; Thu, 30 Jul 2026 18:33:06 +0000 (UTC) Received: from mail-oi1-f169.google.com (mail-oi1-f169.google.com [209.85.167.169]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.18823.1785436385113950930 for ; Thu, 30 Jul 2026 11:33:05 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=GBrmDLJN; spf=pass (domain: gmail.com, ip: 209.85.167.169, mailfrom: jpewhacker@gmail.com) Received: by mail-oi1-f169.google.com with SMTP id 5614622812f47-4864ebb6268so95057b6e.3 for ; Thu, 30 Jul 2026 11:33:05 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1785436384; x=1786041184; 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=BZcxZ3UaMHTn38iuVdLYUP8pryE3wcTDVvswqZpNd+k=; b=GBrmDLJNhxstrvmnq+4Dg+jtczLo3jtpXZ4Mpk7Jvj+aDXYwSHw/06u0LfkpakSd7R QbMZuG0vnCfOfwipzZ6fpa5YSinDSku4Rj6del/DtE5GIY9l5Pxm9AerwCMAfZCStWWd 9jVyY3nPea2eE95vhAKEAe6euRXI35ZQABymfjEhTI459KEde+AwhgZQFUKawKhaSjYH CQUSM2hGQs/B69LpFULmVgKsjnV7O3KU4+uOUPyXSNeDXhdMQ0InCdzc4oGCvBl0gChG 1XyfvujChKPGaBJJkfwiJgPU6hE4zskkiSvcTelyQgSCDxelV6FOOHfX3xi4mlPWbRP3 tIgA== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1785436384; x=1786041184; 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=BZcxZ3UaMHTn38iuVdLYUP8pryE3wcTDVvswqZpNd+k=; b=pV13DGI7m0jMNdsjHcdC54Vh4+iJVy2wUZpgW1FvUz7eSBzBeNiOn7egzjEz64pYxi 2NjEp37fvoczWkT5owOtgYWiBv/hT/gfM2NyRSJEYOv5cjb9n1Lmc/VH/tBRhjb4ITER B/2/o3ihrdTUK9w6tkL1A3DzMpUD6U+JLx7md0c1xFaZhv+n4q0/u5ozdp+iPSVPuJlF voBKeregCBCwGq+ZjO6cbGffhr347JRYzbwF/k423HvatkhcAQVYK7/+AURyrTNeFqIa K8ISnWt2Hksp9M+/NK6s8x5dtb8/eh6ANxwP8W/wB9Z/cAUl/FAUHbEPpLn61x5y6+xH 18sw== X-Gm-Message-State: AOJu0Yxt/PyLdKG29xofEfSxOFji+RZ1FmYaT1B5oGzajMPYKBH4OI9L ZlH1yDRL541YXfcAipEPQjs5hDn3MvVkFp0qgFusW8CSs8H90ZRvPEOmkFZ1IA== X-Gm-Gg: AR+sD13dnSjWSDeEzwtC9lVsi9y9cA4YePKYMKaHL/5k3fVJBX2wYfxGOzdem/O4Unl puJGrIXz09Hi88eJhFZuXscl/qXl/A0WpKgKmjhSFjun7oMcexviKl1vkH0TDdUVgWT3Qn+5MpF kGFHJmGi8df4Ze7XtNU3o1nB63tv9YN/P1y7EJ4l5cNeAcUvLQV8hvxe5n64M8G3K0rRB0tNZAr dJsdJdDteencZ9it1YICvix+A7usBT/FydFgsSm50s/lyCi0hcyp5RvTaOSVBGVi3QVZ/QMSMkr NxcWj/hFxYsRgba4wceGks8qmbdkYh10ZLFe0dfTdh/9/iFJ2PNxKK8HAL1ftLTc8Yk9ohA9+e3 TL5m1m6zPgjYIt0bFOtcqlc2CweeYmu7h10/4Uv02BFtBwpKpBtxhoEH7zOcjgOMEqycTBvgMbN nwt26YFea2+Aa8iWlX+rpTw5baoY8osCvNKTnMrvCMU0X3SLy3I9TgZa9vNQ== X-Received: by 2002:a05:6808:30a3:b0:4a3:cea7:4ab4 with SMTP id 5614622812f47-4ad877fc162mr3161720b6e.14.1785436384204; Thu, 30 Jul 2026 11:33:04 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b02:22d0::10c9]) by smtp.gmail.com with ESMTPSA id 5614622812f47-4ad6efc5c3fsm4593047b6e.13.2026.07.30.11.33.03 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Thu, 30 Jul 2026 11:33:03 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt , Michal Sieron Subject: [bitbake-devel][PATCH v2 10/10] hashserv: tests: Add test for upstream pipelining Date: Thu, 30 Jul 2026 12:31:03 -0600 Message-ID: <20260730183254.793698-11-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.54.0 In-Reply-To: <20260730183254.793698-1-JPEWhacker@gmail.com> References: <20260724212822.1165552-1-JPEWhacker@gmail.com> <20260730183254.793698-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 ; Thu, 30 Jul 2026 18:33:06 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/19890 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 bb227c161..6201ce3bd 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -1693,6 +1693,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 bb.asyncrpc.TaskGroup.run(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):