From patchwork Fri Aug 28 15:58:08 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 96674 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 0360EC61DE0 for ; Fri, 28 Aug 2026 15:59:57 +0000 (UTC) Received: from mail-ot1-f49.google.com (mail-ot1-f49.google.com [209.85.210.49]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.4134.1787932787026584300 for ; Fri, 28 Aug 2026 08:59:47 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=QuXcRA1Q; spf=pass (domain: gmail.com, ip: 209.85.210.49, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f49.google.com with SMTP id 46e09a7af769-7f3ece23165so1346233a34.0 for ; Fri, 28 Aug 2026 08:59:46 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932786; x=1788537586; 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=3fYbe3oErlV0YggmLkEPdbvl0hPmAFAJegxZ6Vj/yjo=; b=QuXcRA1QDifAPRqqrXQa/KnGQsw5jSKlXYRS/4Rn3yVVnpDCYL1eCYfKbKm1SEuFow kOHpcavb3oLUTsZpRTgNVx/tGlhpYkIvU4pQ/GDovm+dyyAldUj82OHfK+56FZjBGLSn 69FI3laFS54e5MqyDeQW5x++LnyPT8USy9XX+Dr6oIrQnfTZNvXGDWG9Fzxi1xJFCAfM fhpl0a9bQB0aBILaXikS76b74ojDEuhBa+WMMUOzdFcmcL/pocSiNArOczNDtktVGGKG BvIPuSXwXW5ovQY/inyObrF3hAVA5V8DDgSd3cJWEvRFMSOoT7KuJRF3zG8Rqm0vAOrU 9nJQ== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932786; x=1788537586; 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=3fYbe3oErlV0YggmLkEPdbvl0hPmAFAJegxZ6Vj/yjo=; b=Gy4y9yHpjrO/jStT74AGl64TCWZMtgdlaR8GPqx457Z9yxXbB04ouRzN7uqQOSy5G+ EccPGL4pIsdu2dMJGnVHtam/iRTQfR4v3dCfns/WFXDjurwaN7V9Of/H7CEYqnFnsD4M zbFquNvjwYa1jI9zwxPka1IKggjQWq5MMxgBD4FhiKJ+Am6vh98My1F5ibkDyJkzSjsg 2YYcja3Fc8czm26sbXxh8AXn7cvsW+PDreDftyNn/4im+Cwp3fjGQwz8SqZSNIBcm4Uo wMzjbLzggwEVhL7yUJMeOZO98+BQCrD0nq0ScGjM3OUbWlbtHqFjjNpSAH2kwScwEOsM J6TQ== X-Gm-Message-State: AFuF++lAqxnj+sKNnG9HwFAi2vzUihGXRQyJXZCnlcVUrv9AG/gWbnVC ViOsc4HL4SFFKtCjqX8i+k6ZXiAPehZ2jS3cKQLZC9WYf0oHKoDB+q5hF/D1LQ== X-Gm-Gg: AR+sD13Kh1HqRQEUExx51mJgwi5i2/ak8VSb6UjQtlOTm9/IishoTvfI7uO7yvmvWOq XV3981EYoiioeeXT/d6vj+tKJ6Od3Aq3eZ3La93mx5p+W6gPagklEGX0w/aB0jXvL13L3YHtWbX ZBxQa/6KSq5YsUilXcm8x83x/wLnXaoYQbCj2gkSyMrL/R24IYEJ+nzu/qtiXDzHNsnLvgRjfyf XFWlqFo9EfJS9zmLMUuJzGkeMJVGEenSsNS38jQQwtqlEw4HIQXmxbnhY5DOmxc2VxMpDewJe8F dPtDHXqyM7VFjcNjI/8/fUWmATwxYsHKEwFzs8eM7fE0zW1mMkJx/9IYj1etyXRVMzL4Ym8bNgQ tDOxaiWUDHOERSyZ0m/XSsI8I5Xx0rtMOYcA/iJt2DGU3lQgrIpWG08VyZgWuoBL8MefeoZXjUT 1wFZIggXVFf7v0HSQCzaHelF5ZF/cU/6BoT+Xk/OYJOkk/79EUhbvNWBVGioD6kic/AZkI X-Received: by 2002:a05:6830:230c:b0:7e6:f2dd:93c3 with SMTP id 46e09a7af769-7f4f24e3adbmr8422055a34.15.1787932785886; Fri, 28 Aug 2026 08:59:45 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.45 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:45 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 01/11] asyncrpc: Add Task Group Date: Fri, 28 Aug 2026 09:58:08 -0600 Message-ID: <20260828155942.1219468-2-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:57 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20107 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 Fri Aug 28 15:58: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: 96669 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 70154C61DBD for ; Fri, 28 Aug 2026 15:59:55 +0000 (UTC) Received: from mail-ot1-f43.google.com (mail-ot1-f43.google.com [209.85.210.43]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.4135.1787932787636007080 for ; Fri, 28 Aug 2026 08:59:47 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=YQ8/6+J4; spf=pass (domain: gmail.com, ip: 209.85.210.43, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f43.google.com with SMTP id 46e09a7af769-7ec1e9d3359so1262566a34.0 for ; Fri, 28 Aug 2026 08:59:47 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932787; x=1788537587; 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=jNbp8ctbFzHX6aLs56BQ70pPObANvXVauMmYbwoHVoY=; b=YQ8/6+J41uEzw0KignzcymTyQBD3R1sdh3aCrT0lh4eu0PXIB45xx7iPC9X8MYctIK DhtB+3iq3MrMHl7oS6tnwEPaph35/gfS2D+EP8ycelpGumKVTYXPkAIoIQqv5/HTp5Nm k47HdqgX2x/f/GKGY97IuEwdEz5QY/Moe6ILnNaO4ehcu401OuPen78VBPxt2nWDpigi cMgo9nyKNhW2/BuRFlgHg9aVhHY3fL8RByi5Z7WgdUTzrIoEn2mbCOKvJy2rEQS4E9qG x1/+8hdfkt7n8ZxTgtLNlDvxzds1YwaoCFQnUcmcm5Of934kd5xGF+Y4UuIQ2WAYK722 rSkg== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932787; x=1788537587; 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=jNbp8ctbFzHX6aLs56BQ70pPObANvXVauMmYbwoHVoY=; b=CvWajm1+Jv/SA19XBb+45sqAbIrYWuYL3L8JfwcmqHLVuL1SIUVnl4zQ4VHlMDbgK2 swdhbN6g/LeOcopI8VY3xbeDNSYz1/OtmQ2/CHcJbGsU7jhys66aAdZmTptoCEAJDeTo AG7nW4IUodQXGrNsdhSDpLXztcD0sjvae6vUj6h6EPhUk+lrj/zYDpH+IOFk8iPoCzIg vsv9V85/zLx0eBF02MWkYzCsH0mi5FdikheKHOq5Vphwct0Evs9vYf6v1+o+SQibBbU1 PUFeREmmI9xbYhtPxRpoad7YKth9CMbuMsJeoLCkoMykKm4HNOoGbUTczLTBDBLEtNkD 0HEQ== X-Gm-Message-State: AFuF++lIqWPTrUraSysfkxAa8E4G0gFXoPEkyat6xTEJtQR9jw2fzqN3 J33RGIxHyxk2vnfEvaFcpPhDSmve4p3UZn3d0c9BDIvAqozpiByTbv9lKhQsng== X-Gm-Gg: AR+sD12YQBMDSeP3kD2b8sAj7ypG09Rx1UPEldbP97U0Pw7l+qgWYTHedlu84VPBCAO jwbmBJgG40USGVGWrfo2melydN2KS57UHJ+ljpzyfA68g26mcl5IrY2GpQQyc1a4Pk0s3Jg5ufh N9Va2erojKpHbQoBfOcqsEAbVAJ6tpICCxp1P8PQM/lyNTZBHwpYzJdAm96EIpKGjRoRHeAwXKS cAxr7hXiOfImq7icF6A4q5en+8OGGkMYZ9JWPqP0pDW0gF3RbW8uOgnmFGW22OuujaTZlkWbek+ uea881B+kjTwVZpQtC1r+T73hV0nYmCA/sM/pSMadWYwElWDZnLWMNWMqk9NvmVdVsaqEp7X8h6 chcmHCVhbus+/R85SHXVj15FGn5Y04X2ZvJFBeTp8uYf5HXqwS0wn0tEVZutNkQW/KEFc57bDZG Ae3DxYVWh8gmd33XuTbtOMbCZByOnq5FO+ILYEMtNnOv0F7+vOq1zOrm0HWw== X-Received: by 2002:a05:6830:439e:b0:7f4:d16e:45f9 with SMTP id 46e09a7af769-7f4f253ea8fmr8507966a34.7.1787932786694; Fri, 28 Aug 2026 08:59:46 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.45 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:46 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 02/11] asyncrpc: serv: Use Task Group Date: Fri, 28 Aug 2026 09:58:09 -0600 Message-ID: <20260828155942.1219468-3-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:55 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20108 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 Fri Aug 28 15:58: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: 96672 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 779A7C61DDA for ; Fri, 28 Aug 2026 15:59:56 +0000 (UTC) Received: from mail-ot1-f52.google.com (mail-ot1-f52.google.com [209.85.210.52]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4200.1787932788448865750 for ; Fri, 28 Aug 2026 08:59:48 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=FyXRdZXr; spf=pass (domain: gmail.com, ip: 209.85.210.52, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f52.google.com with SMTP id 46e09a7af769-7f434cc1c5eso1163834a34.0 for ; Fri, 28 Aug 2026 08:59:48 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932787; x=1788537587; 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=24w7ijzTwTc0o1o+lkilFFejitrP11DCTbSdBpCLG8s=; b=FyXRdZXrGxJ/Itajp6aVk+JG7b5ZDMvNivJu0xUYYVVnnhtoFwK685tZz4w9/zLZ7G JA2EY04vDmvIvKIH0rPR50KWE+4LAI8CE8na55UvaBYww6LHxJ9OqkERBGuj9gmorv1H 5SeqXxX0PWUDZGUOXhR+Ot3U192exRRY668/DYlTQiHdhlZdNQeX+X5Jk7Ctfy0qUzA8 75ko9AFdv3kjTOZgoqwFGJXVpkB2dO86SluxoNEmbLs3BlbcW+zTT1qz0t4UKhOqHGnI iIEuwi8GTrHymZnClv3ST/7RGvdraERqN012uvQeicLKdwKK0tUFm3T02rb2/NDZzS4E s+xA== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932787; x=1788537587; 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=24w7ijzTwTc0o1o+lkilFFejitrP11DCTbSdBpCLG8s=; b=B640rYLsjx+dbMCDrnrDRfTiITP4Mflo8scDWoblskL2/4vSxg7AIn5qsfX1epwwfO QLDzJYPqo8Vxpv9fr+llz3654rddHBkH5xZ23wQ0hS/Sm7dHLPpybk7MXQr9itzanyrN /3KJevzqHe86fD+rfXvlhO8CyhcbXfVC0fZEqnXDTfYODkAhW0gB/4WH81DOKv7ff5bg 7NVvzNBuw1Bbx8dfAu1cCrnY3CdziJ0+EHG+7kN5YgwFy9l4iJEHSe+PC/R9guztdN2i 3Kk9yGq1YWeh1IBBi65OQX4SPGKDRXd4m6g26b98yjlPa8MGtmsNGh5BAGq+f9rxZ+63 OEZQ== X-Gm-Message-State: AFuF++mDpiUtAwXhSPZcmG+nHbL1b1BoRrXrSeHt7TpQcfuFvjYU22th EnH1V/bSnWj2nI2zvwbSn4jcuDxa1gWXv77VfqF95FO0DJScEwhIeohMWU8hKg== X-Gm-Gg: AR+sD11Uf6eAso+izmavrVXNf0kC2kxnHePbBwMX/p63q+dQpRxSULPhBO1NG3qsQTf Z3fkrfqA6Iv7EIp1Lih+n0nZHRuTafCPzCNO4AGOuQRNLJyxNK/Q+3bJk/PY1fueCm5TbFDz7eX pwyGwz7REuygjK0iAcaS+xqvbfELNNdM2rDwq08iEdm22iTFxxYIp6uV+snTdGPo+pXOq2nvddc eE0gayXm0QYHttcAyiIJiXhL1k7SApL/Pe1tn7lRrZ1HSF1d+YGwrmmqfpCuX3+FCcIhdbLncGp b56CEjby01bqT2XfCZ5/1XirfWx6+5rLSd0uUJZVAcaQJIdSFlHSXdHAx+0S8V4Tx6vArnn75yn 8osKiUovQKPEHYWM65yzBodPyjOu4dMPR42PzA8EgU4DDABjizKwQpH8oXTQtEj9tD/KsCxWEAI p30Y9xjBxy1XLfuKh4uLEctt6YFrGPMDtNwW/4mm1cK5oA95grwx+t0J8Wjs4yQPAqKgLn X-Received: by 2002:a05:6830:2106:b0:7e7:8dc0:3951 with SMTP id 46e09a7af769-7f4f228d83emr10529994a34.8.1787932787539; Fri, 28 Aug 2026 08:59:47 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.46 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:47 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 03/11] asyncrpc: serv: Cancel all clients on server stop Date: Fri, 28 Aug 2026 09:58:10 -0600 Message-ID: <20260828155942.1219468-4-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:56 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20109 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 Fri Aug 28 15:58: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: 96673 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 A0B0CC61DDE for ; Fri, 28 Aug 2026 15:59:56 +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.msgproc01-g2.4136.1787932789290362855 for ; Fri, 28 Aug 2026 08:59:49 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=fW50jDk/; 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-7ee37dc91f5so632753a34.3 for ; Fri, 28 Aug 2026 08:59:49 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932788; x=1788537588; 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=DoHa3ZUIiIR//1vZWiyeQap8laLQkHQfURKVV6sd7mY=; b=fW50jDk/ONJ1qrX//y81g3K+Y5c8l1Zsbpp5RD5UQLo53PLsr6ixULv0lHhTDkYMA4 xwU6tfQUod9DAiN8uCCQq30fCn/7HWaOwnZqVab1F698w2Q7tuTumBKvAYkK6qU56BHE 4tEL32TgvNu01jFHs6nedsjPzDxUijdIk8OEbeNWDvSSWYS3TfHhcnHQhPFmuc3V7fsW yMy7Mz3F5apSHr2K6mjJxdu9rRu7IFcf2zDbw5RMgsX3ZUWZanFUZiTEzBTKWwVysnUF 7rEwYLyTE0M3kjBucTtItpF+xM+OhzSodHCl6sE0FY0/NDEaDVJyKpiK8umnYvE2YcMp 7AxQ== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932788; x=1788537588; 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=DoHa3ZUIiIR//1vZWiyeQap8laLQkHQfURKVV6sd7mY=; b=XUa+DYvBgbMvx0KJ4QmiJ0tv3VK1QI1FE6e0M+vnnAxB3BlIW66lU2xRpqOYB8hxMX TrSTQLtntd7slh2zFuzAjWq6AUq/FazxG+6eTA5fZ9S1yMeMFG31Dfa11yCqW5rS4OVT N3BqFxQpXLofVWxYs1u0hlUM0mOk83AmhTg57P8pZb1fTefek0u8BiR+ziHYqr7TDnsY 9xZlSIxQ8VwVWSxq4bs706YduTyVuekkRF0ZYBwiSy7zvdZfzLJG5dWq/90AB7rR/iUT fJ/WXlI07+4fDi8+0m9R6ms6fSQVng7S2lB9v3901hG/b8wvIhcEOAUc2xZd8ga0pPTV /9AQ== X-Gm-Message-State: AFuF++nhl1S89iytZ8lfzC7gJUAHj8XPFf9Yhp/2xNfJIZV3lrjujYDY OsjG/0xes/1zwVJwgVa9v1gKa1+GNHuNZaBiyaGcCMpNKHTa1MpxuL7jNz25SA== X-Gm-Gg: AR+sD10iqobxrRVHk5/+GvwcIJtKx2le8l5+0fFIExhgwSQxJ8BkckAizJkcYjr4th9 WLNl0AcXXAE6hT9a2/HxuyxDmfU3xTfic6t2MMjSosGgsTqQ3NmgkdLSymNCbJLwDJJqsS+XDoO eJXVUmjHuHNo4iLB8ChD0H3TsHZmZGOGZlTlUEaRUeZWEU6OXthKeWTmCW2awhbsxeAobDTPrkb e+1ZghS6WZyyXpXsBesxQMSUPfbCwbhN1l0SK3roTS3Z3iAEUhnMICo2Rq5NoZ4Y822ShdYURi3 hEJOgozKdj3+819SLU9AehxYLosYxTf2cQVMcPrI8LvYGwMaiYkWkcu3iGQtVF9EdRPiONbC5ON 2cidK5m2gsYFnXAM0p2L/iRcs0m/pVr2oImBcNXzUe3rMWCZR+GlddiLPvpRNTkj+pGrPZ4kbAE Yr4TwyipDDUlNcincWUJaTE7JS4WA2EN1LYy0l8pBDVdVsFkPw74KgCiynbjWC49pfAzmV X-Received: by 2002:a05:6830:4118:b0:7f3:a8ac:45f7 with SMTP id 46e09a7af769-7f4f224178emr9021485a34.3.1787932788460; Fri, 28 Aug 2026 08:59:48 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.47 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:47 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 04/11] hashserv: tests: Improve test logging Date: Fri, 28 Aug 2026 09:58:11 -0600 Message-ID: <20260828155942.1219468-5-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:56 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20110 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 Fri Aug 28 15:58: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: 96671 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 DB75AC61DDD for ; Fri, 28 Aug 2026 15:59:55 +0000 (UTC) Received: from mail-ot1-f46.google.com (mail-ot1-f46.google.com [209.85.210.46]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4202.1787932790216193821 for ; Fri, 28 Aug 2026 08:59:50 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=svbxoLm6; spf=pass (domain: gmail.com, ip: 209.85.210.46, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f46.google.com with SMTP id 46e09a7af769-7f18c0e03e3so670741a34.2 for ; Fri, 28 Aug 2026 08:59:50 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932789; x=1788537589; 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=elGM45IJFjC0vRZKWZOIxkCy/9N9jYnBJzioh548U48=; b=svbxoLm6ZT68eyrwyzfXkmBFw60srH9MYTfKP9Wf7KHqyneoA9mLxiMQebmMSogHsR m1f5KmEtHPDe/mj+jSOZclnh3sHtLDSmPsdVaoEDYGktcI4RuFtr4Ws9knVUinhrp5FX ikOE8MtUJ18eFTCC2U/RLx4RxcClmNwgbDFmSG9bOp5IfS8W5FRPbYM3W5cLXIbTfYps 69s+XlH66ErPE2L6la/eY/zsvBZQYeTcufs1DTvT80VuUImwbeXhuIhev599dk4daLMc SDDw5jBACPEe3ePRftGcrrSylsxSlf/CkJSsqMo2/A+5bx2B2sw8U7IyZYQNn5LPu8WU yMNA== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932789; x=1788537589; 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=elGM45IJFjC0vRZKWZOIxkCy/9N9jYnBJzioh548U48=; b=pU9N3qURtatDMENHh72DytXMJLu0zzmBuxmFT2yVj1WhVnHzVtAWUoFleW8vOB9rgB wPBTzR8+bUOONaveaE3Ctb+jTq3SWZdfpvJzFTQ2W5JUIUX4xi14ISlHD+ZHcQcSuVo8 Y1xNf6Wj0pjLNWM8Tj8KGWycA7+8I0qBLUIA02leQx6iJjs/20OW+0gMGlIxICgaIIKX /XDT96St37pE7nQuQl/jsL0LcXmYokyW94VwbvIuJ9LnFkFGBbIYgUieVtjcF/HaVYee 7Cp8o7Be8bquvwUuUGNaVQYhLHVssLv5kYoflVZR7SUxClJjQck4zwEUbmFVGDvcqY5l oJDg== X-Gm-Message-State: AFuF++k6sh+e/sQMkMj1/VaofgW65N0gnd4uKbplnaUHMjFjaKMMbHGg oEn29KYbKYwzJt7itpmDWUnzaFO06asabqMr68DuK42A2IiZVX/1WQ0S1QfPGw== X-Gm-Gg: AR+sD12/R3BVi3YjrZYgvXCrmFnXsfFo0+/bQr3v50kfCozKJNPu35gyVkvhjdDBBBn 41zDOZLWuh+g96rCuUJNKB8TjtOgk2+CSFJ6lCVcxPsCZT5Y2vD5VXZLSfhHibuGWmqGuhDNHRD kEV0GCSTwGYx5VvllnDuyfcxXXO4a3nzUkQHgEquaxx3v7kEsy65kXMEgBdzAb51notEYDhjy5F aGtD8WE/yvfxJ15r1c0HJlRkbiGCsfOpBKd5s6yXZiuA4B8cy47JSjxGdce7FDiBiFkOMfUgRHV uSfcc+liIxduDzVprVdRr84fteTfUZxLCprgPuUzexpysDyjf6U93DtQMU96Z2z9yVeqS7rcWkW //wwtONMgj0MNkLV4hadeZTHhsga7zk8TtzBYx45UQHUglYhNxZo+2SsBkYGY2NoTdi8w+nPTqP BSTf1oL6lHHB8tM0yY7b+GJ+p+YTLBARszaJQyuADuFYPdjCrCsiZGai9gxA== X-Received: by 2002:a05:6830:4408:b0:7f4:f1f3:da3b with SMTP id 46e09a7af769-7f4f227a85emr9168115a34.4.1787932789321; Fri, 28 Aug 2026 08:59:49 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.48 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:48 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 05/11] hashserv: client: Add asynchronous streaming API Date: Fri, 28 Aug 2026 09:58:12 -0600 Message-ID: <20260828155942.1219468-6-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:55 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20111 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 Fri Aug 28 15:58: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: 96675 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 EA224C61DE1 for ; Fri, 28 Aug 2026 15:59:56 +0000 (UTC) Received: from mail-ot1-f48.google.com (mail-ot1-f48.google.com [209.85.210.48]) by mx.groups.io with SMTP id smtpd.msgproc01-g2.4138.1787932791392230838 for ; Fri, 28 Aug 2026 08:59:51 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=svLmxYSB; spf=pass (domain: gmail.com, ip: 209.85.210.48, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f48.google.com with SMTP id 46e09a7af769-7f0167e59a3so741631a34.0 for ; Fri, 28 Aug 2026 08:59:51 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932790; x=1788537590; 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=Jb1Dzfauv3VZlZPygDmNfwFlNe3kwbtyFp5vbhuxCh0=; b=svLmxYSBy2UE7G2czAGTmJkeO+3kB2jG1QdvUA4tmkdb8WgTVrvP+Bo5zYijifrZ/4 gDqvq4UdjcJjqKFpPu9AyRZez8IK5klu1ioO7n+uABenERLgmRCaIcwu6sOuNrXwh1f0 uzTRsTGjI9BwtqYatAbI/7IigGB+pjcTKUvR8A6bnARFElJM6/WsNv0nLN/wuDCu3gP4 cW+YPSc+5U66Zfim98RsRkM8wdqeeQes45rIuI7Ypy/BrLJSNed9uSfqqYppevOicqJr b/uctVi/oPgVmsDUu81ksAXb8ghI82WUFP/abIBEzzK2zYL47ABqIVVSjmFTVFoiSTXp 9vpg== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932790; x=1788537590; 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=Jb1Dzfauv3VZlZPygDmNfwFlNe3kwbtyFp5vbhuxCh0=; b=DNZi66NjdaYkbfCOtpJINAcoa9l3InNButOX3koelMU90Gc/71ROFqDQ52jC17B5ER N28yhEhchUSvLhvujb/dKggNMpJ0QymJY2vBuCqnpTUKvW16HZ7DACM7WexNK9rdZiAK zGR9Xr/Ol6rmjQlJoBXU9BmwUFrTpI3OCbaDFidsU3lsslrZ+Sl2JbIasoW7JQwv0adt BSwZLBLp6ZCvwktqyCfvY6QhB9XOPc8Gt8YH9oJEFqCG2fsgHfpAN1FPpzvOJBneffgD Potm7aqR2x+4u6nZAS5Wjsn0mF3nnARMusFHAqziFH9Mc9orxCKi3haqB2hFWn2PFBY8 KjqQ== X-Gm-Message-State: AFuF++nLHcYHWGenHxBh+SaYVSJQOgrLwdo1TlTjNEdB4vuIce0AyzbS dCSoAqrD6CrKNbn2LuajHfHUK9F1j4syyxat5COHw/VDugVkXVabdbgAHlsZIg== X-Gm-Gg: AR+sD12kU15nCwdOrm0C6dFM2fLVYnAvUgu/n8BR5hGPv3pM7KZffgpAPK1zKaQF5fn WCga13ho8zY9ZPwmdP1se6YaK6TJT7ka2K1LACPAiGxnE3XecZ8FcKDKguKx6b1flpW7KIvyRoU mPG0M3mYdKLpjzBOH851kZLh/NyD7XgpTZMbf2bS7cLM3bEnZzeCowL6UUvBjLzii9yXvqbuYtC 1EJipKjal7wLAF8C+7DBx1woYck1pl78wH5dBj1FwbZ9vZtxV7rIwA3Y+QGYyzwgZ9DWqP9ibi8 c+LmTZLgSm6XDMICVK89s5L+/JkhnCHk4AGSt07J63gTjZxYEVaOHGHAHYpMtiL1mViYHJ/jq3L CUALm+0RDx6h4KN1O6AfcwF2spqfuc81T1q4RTPLGTmG/b+nXu6GTjMmzQd8TSS8LergpXvWFPE SqXkkOae8k5pdJB1s7u2bXSK2WX+0vSx2u5Z852cAh+XLvv30lj4ACNrMSeQ== X-Received: by 2002:a05:6830:4406:b0:7e6:cfd0:42de with SMTP id 46e09a7af769-7f4f2624dc3mr9607954a34.15.1787932790452; Fri, 28 Aug 2026 08:59:50 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.49 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:49 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 06/11] hashserv: server: Add queued streaming API Date: Fri, 28 Aug 2026 09:58:13 -0600 Message-ID: <20260828155942.1219468-7-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:56 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20112 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 Fri Aug 28 15:58: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: 96668 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 82EC8C61DD9 for ; Fri, 28 Aug 2026 15:59:55 +0000 (UTC) Received: from mail-ot1-f44.google.com (mail-ot1-f44.google.com [209.85.210.44]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4203.1787932792309588629 for ; Fri, 28 Aug 2026 08:59:52 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=ZtAczjl0; spf=pass (domain: gmail.com, ip: 209.85.210.44, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f44.google.com with SMTP id 46e09a7af769-7e9eaf04bfaso513870a34.1 for ; Fri, 28 Aug 2026 08:59:52 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932791; x=1788537591; 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=xvj+eveOAyyIuwPOF9snYwI44PrD8WXGB8NyChgMC+4=; b=ZtAczjl0tYCIXNbHcEsRWnlmg1csOjRDpAQRpI/VV0pgJbwnmaRJVGCXhzXyyg9xoU L/HBT5miD9qJ26OuVP66JxoA31zhhbdg7ceXA09e/NyZj7yNxhwDw2nZljHyNGJyYztD CymXRSedxyx+a+J3j5ZpH7mE7D3rYC/KWSH6mRuEzX5BfA9xUiCNvTiVxQ/z9FrpQR0g cv74mb1dI7iMI1Pw33GnjUh8BFUiLaHx++XECAQNFUbrhYfGraMkVfxfM53H/nVoZOW3 +yHrOq7r+w9h5Wy74kSe+eXeu7HO2+5wiaNyMJwKZpv+UJAW6y3OGGsGd0Q9I1/QLbrE Q2mw== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932791; x=1788537591; 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=xvj+eveOAyyIuwPOF9snYwI44PrD8WXGB8NyChgMC+4=; b=XtBp/XKqmwqgRMvq6L33VEQo5bIQwAeVWiJcNwn3WZvqwuNN1lXtTlE5TCXdlG2bxW ZViC7UJ8I/kgLHK0/+AD5jvZvdOMcAemU4nz4DM8ipz0drL4a17dYoKl0aw6noUB4HUG vyfpmCtoWDovfAeXQP+vqdaVm3zTpYq//JTDxNS8Vdn+SCZRe21EigBMDz9lU6HnB5k8 k8/xI83SL41mOz0qPiiSt1+hQP317IEXpRTtrpqQuHHzpsRL5XMj1jjcERH5GSbWokP0 VJMCrOasEbZG7XsFsKzlX7+Yw/Fw1pHg6dxYHpy1Utqa8fX7d4Fh9+n8rjQ12IIoCseL xaNA== X-Gm-Message-State: AFuF++mDmct7IyPkSXJ2bvHtHzBP9ezNZ4uxcaGFdozBoTuwa1JKVosu NQMfnVUyID0CzkiRhwJhBSgRyTtO9OgR7W3MTnujaM0079eBIfanJyO3ZjlBlA== X-Gm-Gg: AR+sD10nHuKqKZs44ueWuWzin274vv6LixUbCBeI8ymV9zyTSoBQTYujsODDhIwxprk 22wM9PqQCWpA5Jx3azhWHJ+eWVsFa1GpqvdsXdpFgpS1qcSAwEmbpEE+M46F/ujWW8ydr82uPec JZeY4evEgox128GlanOzG8VEj+CueazsFHzcDRtQ4iSUT1e/MFHLWtXLe8x30pOf46kLEU3wYaA 4zdOPX0OvW+SaVk/ptW/D7Pyn1YHOxlMx8nzle1LU2xC06wuQ7M7YCxzqehVtYAyvplpgyse2Kx y2b21zwImb9OSiR+5sj2cNx0O+7GhJ75KpBg9l/lvCwPdRarv+wGo+HpeUrrBYv6HcNv2vdP5y3 O1RPa3u8jVrAJmQwuLeTZN9l3zcpN3+2EhPAGqY89RSn3DWwtLqirq1eGORBdwnFHxXQsgCgz94 ANDnUHJo2cKVt4kwVqTiKgJUdud2oIfol/Pi6n60JZm7Fus4iedAC7OnaiNA== X-Received: by 2002:a05:6830:4984:b0:7f4:eb76:8820 with SMTP id 46e09a7af769-7f4f24e6456mr9454246a34.13.1787932791294; Fri, 28 Aug 2026 08:59:51 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.50 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:50 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 07/11] hashserv: server: Use streaming and queue API for upstream unihash queries Date: Fri, 28 Aug 2026 09:58:14 -0600 Message-ID: <20260828155942.1219468-8-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:55 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20113 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 0fa81e85d..d0e6f23fc 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -477,24 +477,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 bb.asyncrpc.TaskGroup.run( + 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 Aug 28 15:58:15 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 96670 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 CD2D3C61DDC for ; Fri, 28 Aug 2026 15:59:55 +0000 (UTC) Received: from mail-ot1-f49.google.com (mail-ot1-f49.google.com [209.85.210.49]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4204.1787932793070333061 for ; Fri, 28 Aug 2026 08:59:53 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=L+uWBUS8; spf=pass (domain: gmail.com, ip: 209.85.210.49, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f49.google.com with SMTP id 46e09a7af769-7f4f53975e6so1100018a34.3 for ; Fri, 28 Aug 2026 08:59:52 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932792; x=1788537592; 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=WH/3ShpehcTLl333LvMRdyNhjlZdFl/NOwoXJHGc7fo=; b=L+uWBUS8EiwtxKCdygMi3gTZtHB+CEa0nPEf8f74Gbz37ERiQgBGNz1mgNnql4fqCp TPlLh5XvqzghjNgjrBv3SdLpiTHvbUA5yasVFRAA52efGjto6LJl8LYZm13yndmcxnpc f3p+MRdAI2qmgCahW4/VgUi5gz/0UxU4/rOq9Jym32/4CHLAJHDCcEk9+k+7TuR2D3E5 n+5wAIPk/4DgHMMqCzzuO7IrCjo+QFGrP2SkQ9PZdS2A/9V3J9dkPoRMlA8Gm9y56KT6 VTSlhhzUWdMxjV7xUa+k+MYN7+2nd1NxDbQg9ZE8JKS7uuU4Tr/9sU2rE1aeNhrd5L5T HlEg== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932792; x=1788537592; 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=WH/3ShpehcTLl333LvMRdyNhjlZdFl/NOwoXJHGc7fo=; b=inmpk79pwXAgQaVxOhopdhZr1caXkcvC46rm+2MZj5ijNy1mQ0VYcOjYKBBDQeYDTk zQI7zGJazZlE5qdO/UJbL0ibOY4y0721yZP8Bwrckbcg7iW2taiWe7hrAvBENtQd/6g4 WuT50vLIfscBPGWAXRM7otHp6yYaG6ToaNrc4J2OZ4Q+1M6J/2Mb2t0WC7hIQxslR2QM rhriz4YO0FHvSIExxXZJIoJUrd+WOzFNBSzZbqw48gMRf+6VVQCJkjo+LmAgiWarfwzJ 7XUWjdDX38w1LM1V64o97kJJBnyoFwvYFBfR5Ac9ozrZ8M08eln9/wakWM2Ewv21Lt9w NPrA== X-Gm-Message-State: AFuF++n4WZRP/ALvp00qRuQ+ZvMcpDFLUYpI1UUJfnls+QhmjSrEIbLB HTtJRAoUcGivGC2IjObYuNA9EnWsH2fihxVAHBjZXn7bgwHmymPToRKBvgeZzA== X-Gm-Gg: AR+sD11pzr8PpJfA17aO0jIlHIHxXrmw5ofLZudGpVBTMp/CAiEq15i7RiE6U/rRAUV uxe77EjMp/rBoAJzsgcypXRCahciyMG84HSIFjfar9wEvfU6Pav+VSNN6o8xswKmKYlJgEodSxQ mroyGnOeuzaDUedvi0BN3O+iYm0mzBYJgurgLjWcRL78VBTsDoVbw7yM4EH8ctdZm8r1czEZxfR ilik1yb/8h4F/LcGCYG1ZY0H/WPSBPDxJM+ibnv22L1a+UMGd79vApZI5fDkr7LlQiV78CBgUVZ WR+o5AUXn5XTXhyPb6HPijnIbz6CTTcVlyfE/EdKGUnK+pnbQIKkgbTWpk3juitTuvvMk4x78Dt C3SYCDrVL+lfcu1UzacH17dttdv/o8yKOMm6aOhq9ORwsarD0GqtIOriBHXqVN50bwnF3+gWKQ1 Y986KK5mbYkpRmTz3uxNduyFQMVsGtv6IJSODSIBLqPKfBrl80MifqQO4akw== X-Received: by 2002:a05:6830:700e:b0:7eb:d848:c856 with SMTP id 46e09a7af769-7f4f227a851mr10195792a34.4.1787932792024; Fri, 28 Aug 2026 08:59:52 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.51 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:51 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 08/11] hashserv: server: Use streaming and queue API for upstream exist queries Date: Fri, 28 Aug 2026 09:58:15 -0600 Message-ID: <20260828155942.1219468-9-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:55 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20114 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 Fri Aug 28 15:58:16 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 96667 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 61D26C61DD3 for ; Fri, 28 Aug 2026 15:59:55 +0000 (UTC) Received: from mail-ot1-f41.google.com (mail-ot1-f41.google.com [209.85.210.41]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4206.1787932793978993685 for ; Fri, 28 Aug 2026 08:59:54 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=nC9Y/khI; spf=pass (domain: gmail.com, ip: 209.85.210.41, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f41.google.com with SMTP id 46e09a7af769-7f4dedd67b8so1391226a34.1 for ; Fri, 28 Aug 2026 08:59:53 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932793; x=1788537593; 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=trJxSlvOYmX9uexDym3Vg3/mLaOnCCusDYcJNFLRf1A=; b=nC9Y/khIo7K6DkJRhLEhYv1xdUhgXgaSkp+bZjbcnnK8Wz7Y8Aq/6bNG61MQtzeaIN FQOr+5m8wp4/1vfikrklCyj2ylKyPe2vsn91KPWG8+/biSGhMnf9/qRuDC+1vFSunSLd 6ctyKKzzcubZWYUzLsE+fh/oXmB2ma+OaF1SZXZkdWBcjV1TUr0vUCUN8cPEuLdmr2hm P/10TDTkc22bONvut3I7QD9QSUBUbAj6QNVVJ6GpY7xOHQGNPL0hgZxw+wTMUN3srcf5 0B/mQsPfjQYcQV7X8vOqttI/NX9kusN/hmJ1z4unQp9pn6+lS/PbdbYXjuIxmBDA0qns vH2w== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932793; x=1788537593; 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=trJxSlvOYmX9uexDym3Vg3/mLaOnCCusDYcJNFLRf1A=; b=oHT4Mp6P1EyW8hOmzwO+FIp6LRKD3Zo7rS88XuIvNhiELPR7MQzZTnS8CH1Z6hnYlC JIpXcKS16YElIF+rU0dxfFIvDt+2YwF6O2XZu4Hgjd7O1UAS3lriq6Mo8s8uRbId0Jf3 Ol6v+pPF2snfE/eVYD3/oduAuCbFKxQ9KzFS0Bc1InF85MBHdFiSFVHx+DF4CDxuzhbR Qk0z8Cvamb4QZuMVM/S59Wxquv1olSLP5HIS9VR3JJaksMSUKd7HYpCbrSlh+IEbX42x jQ172xGAj4u69dO9KrjErAOKPIKEhcCpj3mGMc7zkkkQonaIwovtKhlEpSbBUIr4l00j uLmQ== X-Gm-Message-State: AFuF++lWk+LtEPqh5L4PYkfxPc4wUlkkAGN/VUC8tqkQ6/HBS1zJpJYs aoyYYHYnZse2Qgzw/sLsHh0eygPIa+fRuqmiT8dWFrTFhtANgMaVovQ6yZCKYw== X-Gm-Gg: AR+sD11SXtMv4a+F1OlFS9Lxi52SMxioiRPygrWUKzyPZzTT3xuWS2EO+qLenjnJMpo UI2kRdG7rheO1fQUfRN1/4qD5cafW6B3zlaZ62HRkSvcjBWj2esnf6pND4wW0/Mu/V4fITZlaZF KS0aF5zM718afZybI6ay/5mAlD3JiQ+jsX82sgWttIDrTZG+IydVTSz0iJ8XGzh1Z+X9/v6jLPI UETMgYr2a47WDeCg0B/zAHFOT8PnzM5Qhf/ckmL4B6Zo6iu++r9RprpL4iAVrIVJwTNvmq8nzt7 zGk/KW7X/MGjTlnLsdv8GgXf3mf3AeOebQGBl+1JieS9h3c6ZM/3moxOvnuwSQxAIPzBz2Z5SG6 GiU0O4IcD47BT8SH/5bGuuNNN26VDcm03O0c5efujNrPRlGkOxKkM45s6hMK2h/6TaPYdy/1IBK 0mfwWsq9UPuqtZEY989xQSjnR6AQaoaqfl3mOyhTKGBow2iUFlxttW5Lbl5Q== X-Received: by 2002:a05:6830:618a:b0:7e6:f4a3:1df5 with SMTP id 46e09a7af769-7f4f23b7ac2mr10332399a34.1.1787932793087; Fri, 28 Aug 2026 08:59:53 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.52 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:52 -0700 (PDT) From: Joshua Watt X-Google-Original-From: Joshua Watt To: bitbake-devel@lists.openembedded.org Cc: Joshua Watt Subject: [bitbake-devel][PATCH v3 09/11] hashserv: tests: Add more upstream tests Date: Fri, 28 Aug 2026 09:58:16 -0600 Message-ID: <20260828155942.1219468-10-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:55 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20115 Adds a test to verify that the batch API works properly with an upstream server and also a test to verify that if the upstream server restarts or is temporarily disconnected that the client recovers properly. Signed-off-by: Joshua Watt --- lib/hashserv/tests.py | 152 +++++++++++++++++++++++++++++++++++++++--- 1 file changed, 141 insertions(+), 11 deletions(-) diff --git a/lib/hashserv/tests.py b/lib/hashserv/tests.py index a7ce7425e..bb227c161 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -5,12 +5,13 @@ # SPDX-License-Identifier: GPL-2.0-only # -from . import create_server, create_client +from . import create_server, create_client, create_async_client from .server import DEFAULT_ANON_PERMS, ALL_PERMISSIONS from bb.asyncrpc import InvokeError import hashlib import logging from bb import multiprocessing +import asyncio import os import sys import tempfile @@ -41,19 +42,15 @@ class HashEquivalenceTestSetup(object): server_index = 0 client_index = 0 - def start_server(self, dbpath=None, upstream=None, read_only=False, prefunc=server_prefunc, anon_perms=DEFAULT_ANON_PERMS, admin_username=None, admin_password=None): + def start_server(self, dbpath=None, upstream=None, read_only=False, prefunc=server_prefunc, anon_perms=DEFAULT_ANON_PERMS, admin_username=None, admin_password=None, addr=None): self.server_index += 1 + if addr is None: + addr = self.get_server_addr(self.server_index) + if dbpath is None: dbpath = self.make_dbpath() - def cleanup_server(server): - if server.process.exitcode is not None: - return - - server.process.terminate() - server.process.join() - - server = create_server(self.get_server_addr(self.server_index), + server = create_server(addr, dbpath, upstream=upstream, read_only=read_only, @@ -63,7 +60,7 @@ class HashEquivalenceTestSetup(object): server.dbpath = dbpath server.serve_as_process(prefunc=prefunc, args=(self.server_index,)) - self.addCleanup(cleanup_server, server) + self.addCleanup(self.stop_server, server) return server @@ -83,6 +80,13 @@ class HashEquivalenceTestSetup(object): self.server = self.start_server() return self.server.address + def stop_server(self, server): + if not server or server.process.exitcode is not None: + return + + server.process.terminate() + server.process.join() + def start_auth_server(self): auth_server = self.start_server(self.server.dbpath, anon_perms=[], admin_username="admin", admin_password="password") self.auth_server_address = auth_server.address @@ -476,6 +480,132 @@ class HashEquivalenceCommonTests(object): self.assertEqual(result['taskhash'], taskhash9, 'Server failed to copy unihash from upstream') self.assertEqual(result['method'], self.METHOD) + def test_upstream_batch(self): + down_server = self.start_server(upstream=self.server.address) + down_client = self.start_client(down_server.address) + + taskhash1 = '8aa96fcffb5831b3c2c0cb75f0431e3f8b20554a' + outhash1 = 'afe240a439959ce86f5e322f8c208e1fedefea9e813f2140c81af866cc9edf7e' + unihash1 = '5b521d8a12683086cc08bc2c6d94a7a2dcff17eba53b9911e145d51164689380' + self.client.report_unihash(taskhash1, self.METHOD, outhash1, unihash1) + + taskhash2 = "e3da00593d6a7fb435c7e2114976c59c5fd6d561" + outhash2 = "1cf8713e645f491eb9c959d20b5cae1c47133a292626dda9b10709857cbe688a" + unihash2 = "7aebef07d66a8c0f92d0c4f65ec8b1fbb850a3693c53827b8774b64fa9a8a9fe" + self.client.report_unihash(taskhash2, self.METHOD, outhash2, unihash2) + + taskhash3 = '35788efcb8dfb0a02659d81cf2bfd695fb30faf9' + outhash3 = '2765d4a5884be49b28601445c2760c5f21e7e5c0ee2b7e3fce98fd7e5970796f' + unihash3 = 'a69ec97f5af2e21e1a1f9cc8896965515d5559425666f734e245a3d40cee33d9' + self.client.report_unihash(taskhash3, self.METHOD, outhash3, unihash3) + + def query_generator(): + yield unihash1 + yield unihash2 + yield unihash3 + yield "cc74784b2c0ad5b378a6b783c74c518d2c46b8b52fba29cb39a8430d742440d7" + + results = down_client.unihash_exists_batch(query_generator()) + self.assertEqual(results, [True, True, True, False]) + + def test_upstream_interrupted(self): + up_server = self.start_server() + down_server = self.start_server(upstream=up_server.address) + + def restart_upstream(): + nonlocal up_server + + self.stop_server(up_server) + up_server = self.start_server(addr=up_server.address, dbpath=up_server.dbpath) + + # Report some hashes + with self.start_client(up_server.address) as up_client: + taskhash1 = '8aa96fcffb5831b3c2c0cb75f0431e3f8b20554a' + outhash1 = 'afe240a439959ce86f5e322f8c208e1fedefea9e813f2140c81af866cc9edf7e' + unihash1 = '5b521d8a12683086cc08bc2c6d94a7a2dcff17eba53b9911e145d51164689380' + up_client.report_unihash(taskhash1, self.METHOD, outhash1, unihash1) + + taskhash2 = "e3da00593d6a7fb435c7e2114976c59c5fd6d561" + outhash2 = "1cf8713e645f491eb9c959d20b5cae1c47133a292626dda9b10709857cbe688a" + unihash2 = "7aebef07d66a8c0f92d0c4f65ec8b1fbb850a3693c53827b8774b64fa9a8a9fe" + up_client.report_unihash(taskhash2, self.METHOD, outhash2, unihash2) + + taskhash3 = '35788efcb8dfb0a02659d81cf2bfd695fb30faf9' + outhash3 = '2765d4a5884be49b28601445c2760c5f21e7e5c0ee2b7e3fce98fd7e5970796f' + unihash3 = 'a69ec97f5af2e21e1a1f9cc8896965515d5559425666f734e245a3d40cee33d9' + up_client.report_unihash(taskhash3, self.METHOD, outhash3, unihash3) + + restart_upstream() + + with self.start_client(up_server.address) as up_client: + # Verify that reported hashes are correct after restaring server + self.assertTrue(up_client.unihash_exists(unihash1)) + self.assertClientGetHash(up_client, taskhash1, unihash1) + + self.assertTrue(up_client.unihash_exists(unihash2)) + self.assertClientGetHash(up_client, taskhash2, unihash2) + + async def check_unihashes(): + async with await create_async_client(down_server.address) as down_client: + async with down_client.unihash_exists_stream() as stream: + await stream.send_query(unihash1) + r = await stream.get_result() + self.assertTrue(r) + + restart_upstream() + + await stream.send_query(unihash2) + r = await stream.get_result() + self.assertTrue(r) + + await stream.send_query(unihash3) + r = await stream.get_result() + self.assertTrue(r) + + asyncio.run(check_unihashes()) + + def test_upstream_lost(self): + up_server = self.start_server() + down_server = self.start_server(upstream=up_server.address) + + def restart_upstream(): + nonlocal up_server + + self.stop_server(up_server) + up_server = self.start_server(addr=up_server.address, dbpath=up_server.dbpath) + + # Report some hashes + with self.start_client(up_server.address) as up_client: + taskhash1 = '8aa96fcffb5831b3c2c0cb75f0431e3f8b20554a' + outhash1 = 'afe240a439959ce86f5e322f8c208e1fedefea9e813f2140c81af866cc9edf7e' + unihash1 = '5b521d8a12683086cc08bc2c6d94a7a2dcff17eba53b9911e145d51164689380' + up_client.report_unihash(taskhash1, self.METHOD, outhash1, unihash1) + + taskhash2 = "e3da00593d6a7fb435c7e2114976c59c5fd6d561" + outhash2 = "1cf8713e645f491eb9c959d20b5cae1c47133a292626dda9b10709857cbe688a" + unihash2 = "7aebef07d66a8c0f92d0c4f65ec8b1fbb850a3693c53827b8774b64fa9a8a9fe" + up_client.report_unihash(taskhash2, self.METHOD, outhash2, unihash2) + + taskhash3 = '35788efcb8dfb0a02659d81cf2bfd695fb30faf9' + outhash3 = '2765d4a5884be49b28601445c2760c5f21e7e5c0ee2b7e3fce98fd7e5970796f' + unihash3 = 'a69ec97f5af2e21e1a1f9cc8896965515d5559425666f734e245a3d40cee33d9' + up_client.report_unihash(taskhash3, self.METHOD, outhash3, unihash3) + + async def check_unihashes(): + async with await create_async_client(down_server.address) as down_client: + with self.assertRaises(ConnectionError): + async with down_client.unihash_exists_stream() as stream: + await stream.send_query(unihash1) + r = await stream.get_result() + self.assertTrue(r) + + self.stop_server(up_server) + + await stream.send_query(unihash2) + r = await stream.get_result() + + asyncio.run(check_unihashes()) + def test_unihash_exsits(self): taskhash, outhash, unihash = self.create_test_hash(self.client) self.assertTrue(self.client.unihash_exists(unihash)) From patchwork Fri Aug 28 15:58:17 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 96676 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 24445C61DE2 for ; Fri, 28 Aug 2026 15:59:57 +0000 (UTC) Received: from mail-ot1-f43.google.com (mail-ot1-f43.google.com [209.85.210.43]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4207.1787932794864436840 for ; Fri, 28 Aug 2026 08:59:54 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=NEgj/vAY; spf=pass (domain: gmail.com, ip: 209.85.210.43, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f43.google.com with SMTP id 46e09a7af769-7f4e596e393so1334972a34.1 for ; Fri, 28 Aug 2026 08:59:54 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932794; x=1788537594; 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=krMJFZtT11VxKkbCU4bf/lkNGA2769SAPOC0O85UI64=; b=NEgj/vAYDCJzJb6luHXJo8htlgZkIIgqq3r9/tYxLlWc7a/x+z4PWnAnIB2L1VJYnG csVSnYHkTLX4V6xcFKz2sJUc3EMOZs1WAtYTuSTDzrF/BYhxt9hMMdV20AC24q/wXmDN ZNhNG0AlZYnfthjAS4OOyKZFLlrLdD9btLUpeqHnD5bU/79qu6daRjSa4iy2QIAqs8wJ xXNJQ8bmYPk+aVrofbdNh573l9Ya7kJgkJjQIODZWK1Q9ktTTjslyeMbfZMsKKwXhoCP ua3pGwZ1d9M5HWYLiJ1Sz6Z8hFrLMGg1HylsZUwNLmV731nUJX6gdHuIhMrWjZrYstiI hG/Q== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932794; x=1788537594; 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=krMJFZtT11VxKkbCU4bf/lkNGA2769SAPOC0O85UI64=; b=HaWVdOipUQhZDOKtIiMAsHt2kFxZk2bqnLfAhxuIPYjGIH0cORUQhEmRAoM6qrdwHR oihGWWB1a8uqSR2CH2H77xFP3qmYs8/2t1Cv1IqVL2x0Qct+TAUeRpG9dVWGqiKIBEo2 C6lV36ts68aklGUPGPzfA/aVPR7QAYGHEWzSem8gQf0YaFF4PrNs4n54hcZ3CkiEd3K6 CCB8gEkKbtIX+NF5eGY+VSbppnd4jL4GMUYxKpXw3kbH6HZGyi862ZcWpCS1J6LZNFTG xMQ1yGUp8TsrO8RvAmNU1I4uPJCjYYbe8ay6hPWKwMGEZad241weY2FLCvX/aNgJ9J30 Gq0Q== X-Gm-Message-State: AFuF++krlnR9Zq1cd/AIjwQLbabK1Xn07xAUgaMhvXm7L/C28CpAF4Os u9uHxGHXJn2RH145eLe4meKAANWP5u4SEudYX/QeWMEBThqvIPVZKeS4AdRjBg== X-Gm-Gg: AR+sD13T0hGeZcrcn/RXLDjn6RgB221mnen1nCIW6hRq8iTbCO+rr01VurHNR4+Bxji TnKlD/6LY3Bj8VQvAZoJhqsYoQQOJ8rJzBYwDES1rYrzsp8jJZiECBUR9OQ6KW0l4ZgoWLPnlYI zD2DDNdnxC0z/dBdj86dZhHDlI0cxPq0miH/QoBcrtYlxK6AEdhJYvbqmwJP+4/Meec69y3E7CN mZ10MYU9SMBsURsAKVAWvbeYXBqp3G5xjEXkJdoMttyqb0I+kmlKaF2mmu0gBXo+ovit/OO4dUy Z9MBIub3fuoPmoEyipUtWk2z+bAa7Wn8SC7OtloTvIWFbMrRnzAkJjaRW+g1RSZ4/L47SzZdkqW 9XS1719m+BGhQo9MNT6S4H9J0rzomxjRFVx7akASWOTmCh8jIVFOVJz/knMYBD73ZTSh/w6lyBe EypYzWkvxUxqe42ztJhBOyy2Ym1zeYYH82l+9T5KxJqwCI0pYafJlKpIFCJg== X-Received: by 2002:a05:6830:6519:b0:7e9:e8ae:7049 with SMTP id 46e09a7af769-7f4f22c9710mr10610725a34.7.1787932793944; Fri, 28 Aug 2026 08:59:53 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.53 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:53 -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 v3 10/11] hashserv: tests: Add test for upstream pipelining Date: Fri, 28 Aug 2026 09:58:17 -0600 Message-ID: <20260828155942.1219468-11-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:57 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20116 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): From patchwork Fri Aug 28 15:58:18 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Joshua Watt X-Patchwork-Id: 96677 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 3404BC61DE4 for ; Fri, 28 Aug 2026 15:59:57 +0000 (UTC) Received: from mail-ot1-f49.google.com (mail-ot1-f49.google.com [209.85.210.49]) by mx.groups.io with SMTP id smtpd.msgproc02-g2.4209.1787932796284044940 for ; Fri, 28 Aug 2026 08:59:56 -0700 Authentication-Results: mx.groups.io; dkim=pass header.i=@gmail.com header.s=20251104 header.b=tFz2rjNp; spf=pass (domain: gmail.com, ip: 209.85.210.49, mailfrom: jpewhacker@gmail.com) Received: by mail-ot1-f49.google.com with SMTP id 46e09a7af769-7f3ff92cf4aso1472470a34.1 for ; Fri, 28 Aug 2026 08:59:56 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20251104; t=1787932795; x=1788537595; 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=PfBBnipMtCzeMdVMDwm1MJXt0ux8RleKhlKvKgsAJxc=; b=tFz2rjNpNaQkn5O31sEuEKhEXTxceJUd+hVj4ymcIwZTiL0PUOhpOSSqPCiMpB58rQ 73iKk8wdZyEEWHfnFaD2HYgR3aAkzxKQsEf6rELKre9pcMvSM2OjGe7FaHOE5eazXG3W JW0GUigrjEpcXzfHsTAfP5XbYGU+f3Te1/9ZLOxvXEudOEjfDbM9S/jg/iPAeGAg592k +EY/GcJpcISDdPeb/LGt57iBhr/Qm1kEp+RDP6GHW9BmvjQkJCUKhxNy0NIaVgdvTHhR YvPdgInTqIRZtQxWkng+eQnRevj5r2ZScsdz6bz11itMAY15sTR0eTFBkGnBZSi96Yzg n1IQ== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20251104; t=1787932795; x=1788537595; 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=PfBBnipMtCzeMdVMDwm1MJXt0ux8RleKhlKvKgsAJxc=; b=e4CYr9BRmFDKhzvIg2TwYmOw8lUKedoWwKdmtihfKycY2zWhVZj3yzyYYTt7KL+U5z p+tjXgXNXblL7tdh+HZ2IHTM5ujqPABmRzldB8PrMxq0Qcwe4q2kjpELAqSEXa1rOP10 +aMMS+KQU3yejML6CAOM5/y9F1Iw5x4a329/nL0KKbME0KlFMh6EChUwBuzZCB9cH+5Q qCjUe5hrsq1Hv657kgtr4O2JPfZH6NreRy6wsxeYJ2SCs0fN3h/OWZgNOLTmsYq/2Vv8 pRgpJ5O28yZG0lK1fi8rg8g08581RmdkYDzfB1mzNtiBV5KlKMvbTnV3LGqQ8ZPbiX3V Yn4A== X-Gm-Message-State: AFuF++nkOLzPK3G9H4dk4BQ5ydAjJhXUjmB2GcsmZZqSEadNuqtRqguN X2xaY/woXbh4jpphO4sy7sYIP94UZEZ4vcJ5+mesTp5CjShIwNpxzpGECIOaPw== X-Gm-Gg: AR+sD13vSuxW6sRh1dW+TZVIFubU0TxKjuifsv8Elxwosp/08ifjj+z0TenwS3gkU5O VqOe+11vmI4AoMqlomQN9lgykMIfB094ftIlUvubyjtE4cHPywBPrYA5WZAEfHj95TNm+bzJYNL SblmeD0M2KpkerMkAqqczEVDYNKXVZQ9Awa6CpX3nw4lC/dI6GbRcUmr9Z496LDuQbQg3T8D2n5 HQOvEyhOURoRvzkoNeGAc8M5nA6MFkh/Y6waltdAZ5rl6KkrPgj4YTSjMuy+MwwLI5NEKfJxZfB REre6ZOWU9a+E59LIOC/cQhjwOfmqHpzVKg6dxhaJs3kheQ44+viGkIfzQIxIudfRbjsk4M/hNz LfnU5tcKpWE8/SJSJTvnSzB85NsvWFhbY1E7tS821jVGBrvxwEZvw2WLTXZy7tX9qcLjur6tU7M hCoPZ0LpAsGCWnXyCTrcG63j7q62mYHJkGIkYZ4+GuomayGlyMsqYNxqhe3Il6b4cuJz7OYA== X-Received: by 2002:a05:6830:618a:b0:7e7:76f:3ec0 with SMTP id 46e09a7af769-7f4f24fb8b2mr10285375a34.15.1787932795438; Fri, 28 Aug 2026 08:59:55 -0700 (PDT) Received: from localhost.localdomain ([2601:283:4b01:ba50::50d4]) by smtp.gmail.com with ESMTPSA id 46e09a7af769-7f4fa96ce0dsm1506165a34.13.2026.08.28.08.59.54 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Fri, 28 Aug 2026 08:59:54 -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 v3 11/11] hashserv: server: Fix upstream get-unihash miss truncating stream Date: Fri, 28 Aug 2026 09:58:18 -0600 Message-ID: <20260828155942.1219468-12-JPEWhacker@gmail.com> X-Mailer: git-send-email 2.55.0 In-Reply-To: <20260828155942.1219468-1-JPEWhacker@gmail.com> References: <20260730183254.793698-1-JPEWhacker@gmail.com> <20260828155942.1219468-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, 28 Aug 2026 15:59:57 -0000 X-Groupsio-URL: https://lists.openembedded.org/g/bitbake-devel/message/20117 From: Michal Sieron When an upstream hash equivalence server is configured, handle_get_stream() resolves each query through get_upstream_result(). On an upstream miss this returned None, but None is the end-of-stream sentinel consumed by _stream_queue_handler(). Returning it there terminated the response stream early, so the downstream client received fewer replies than it sent, timed out and retried. Return an empty string ("") on a miss instead, which matches the non-upstream get-stream handler. AI-Generated: Uses Cursor Signed-off-by: Michal Sieron Signed-off-by: Joshua Watt --- lib/hashserv/server.py | 20 ++++++++++++++++++-- lib/hashserv/tests.py | 37 +++++++++++++++++++++++++++++++++++++ 2 files changed, 55 insertions(+), 2 deletions(-) diff --git a/lib/hashserv/server.py b/lib/hashserv/server.py index 0153730fa..d9c30bdb8 100644 --- a/lib/hashserv/server.py +++ b/lib/hashserv/server.py @@ -232,7 +232,15 @@ def permissions(*permissions, allow_anon=True, allow_self_service=False): class UpstreamQueue(object): UPSTREAM_NONCE = object() - def __init__(self, queue, get_local_result, send_upstream, get_upstream_result): + def __init__( + self, + logger, + queue, + get_local_result, + send_upstream, + get_upstream_result, + ): + self.logger = logger self.queue = queue self.pending = [] self.cond = asyncio.Condition() @@ -256,6 +264,12 @@ class UpstreamQueue(object): if value is self.UPSTREAM_NONCE: value = await self.get_upstream_result(m) + if value is None: + self.logger.error( + "None is not allowed as a stream value. Terminating stream" + ) + return + await self.queue.put(value) finally: await self.queue.put(None) @@ -504,10 +518,11 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): if unihash: method, taskhash = m.split() await self.server.backfill_queue.put((method, taskhash)) - return unihash + return unihash or "" queue = asyncio.Queue() upstream = UpstreamQueue( + self.logger, queue, get_local_result, send_upstream, @@ -543,6 +558,7 @@ class ServerClient(bb.asyncrpc.AsyncServerConnection): queue = asyncio.Queue() upstream = UpstreamQueue( + self.logger, queue, get_local_result, stream.send_query, diff --git a/lib/hashserv/tests.py b/lib/hashserv/tests.py index 6201ce3bd..15ee7ecde 100644 --- a/lib/hashserv/tests.py +++ b/lib/hashserv/tests.py @@ -606,6 +606,43 @@ class HashEquivalenceCommonTests(object): asyncio.run(check_unihashes()) + def test_upstream_get_stream_miss(self): + down_server = self.start_server(upstream=self.server.address) + down_client = self.start_client(down_server.address) + + # Two hashes present upstream (hits) + taskhash1 = '8aa96fcffb5831b3c2c0cb75f0431e3f8b20554a' + outhash1 = 'afe240a439959ce86f5e322f8c208e1fedefea9e813f2140c81af866cc9edf7e' + unihash1 = '5b521d8a12683086cc08bc2c6d94a7a2dcff17eba53b9911e145d51164689380' + self.client.report_unihash(taskhash1, self.METHOD, outhash1, unihash1) + + taskhash2 = 'e3da00593d6a7fb435c7e2114976c59c5fd6d561' + outhash2 = '1cf8713e645f491eb9c959d20b5cae1c47133a292626dda9b10709857cbe688a' + unihash2 = '7aebef07d66a8c0f92d0c4f65ec8b1fbb850a3693c53827b8774b64fa9a8a9fe' + self.client.report_unihash(taskhash2, self.METHOD, outhash2, unihash2) + + # Two taskhashes present nowhere (upstream misses) + miss1 = '0000000000000000000000000000000000000001' + miss2 = '0000000000000000000000000000000000000002' + + # Miss interleaved with hits: a miss must not truncate the stream + results = down_client.get_unihash_batch([ + (self.METHOD, miss1), + (self.METHOD, taskhash1), + (self.METHOD, miss2), + (self.METHOD, taskhash2), + ]) + self.assertEqual(results, [None, unihash1, None, unihash2]) + + # All-miss batch + self.assertEqual( + down_client.get_unihash_batch([(self.METHOD, miss1), (self.METHOD, miss2)]), + [None, None], + ) + + # Singular get-unihash miss + self.assertClientGetHash(down_client, miss1, None) + def test_unihash_exsits(self): taskhash, outhash, unihash = self.create_test_hash(self.client) self.assertTrue(self.client.unihash_exists(unihash))