From e7f6634f0c884a73c604494672c59ed7c6b906fc Mon Sep 17 00:00:00 2001 From: Bruno Moura Date: Fri, 9 Oct 2026 07:53:33 +0100 Subject: [PATCH 1/6] bump chainlink-data-streams --- core/scripts/go.mod | 2 +- core/scripts/go.sum | 4 ++-- deployment/go.mod | 2 +- deployment/go.sum | 4 ++-- go.mod | 2 +- go.sum | 4 ++-- integration-tests/go.mod | 2 +- integration-tests/go.sum | 4 ++-- integration-tests/load/go.mod | 2 +- integration-tests/load/go.sum | 4 ++-- plugins/plugins.public.yaml | 2 +- system-tests/lib/go.mod | 2 +- system-tests/lib/go.sum | 4 ++-- system-tests/tests/go.mod | 2 +- system-tests/tests/go.sum | 4 ++-- 15 files changed, 22 insertions(+), 22 deletions(-) diff --git a/core/scripts/go.mod b/core/scripts/go.mod index c7488e6a9ae..77d6fee0208 100644 --- a/core/scripts/go.mod +++ b/core/scripts/go.mod @@ -509,7 +509,7 @@ require ( github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 // indirect github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa // indirect github.com/smartcontractkit/chainlink-confidential-compute v1.3.0 // indirect - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 // indirect github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 // indirect github.com/smartcontractkit/chainlink-feeds v0.1.2-0.20250227211209-7cd000095135 // indirect github.com/smartcontractkit/chainlink-framework/chains v0.0.0-20260724153515-bb6a2de39bcb // indirect diff --git a/core/scripts/go.sum b/core/scripts/go.sum index e342ecf03c5..d98919a0fd6 100644 --- a/core/scripts/go.sum +++ b/core/scripts/go.sum @@ -1698,8 +1698,8 @@ github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0 github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa/go.mod h1:nfT8vdE44iAlhDWDV6PaBWucxZZAcDsgGkZ+D0kHpAI= github.com/smartcontractkit/chainlink-confidential-compute v1.3.0 h1:y64hOzM9H61E4EYRzQDtRwqunDBlKpsZc1SvFEUlSLE= github.com/smartcontractkit/chainlink-confidential-compute v1.3.0/go.mod h1:fq9n85XoREIUxgdXIX3RyDIRCRyxZKZ5NCBbFDfdySo= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 h1:CeN/QrCpLFCdWBe/7Tzh/HuXjgTtSq/98YYj40/hHlQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3/go.mod h1:7jQ+XcxZjY19FaWPkvmckf/mT8A+p230BvLMVeEjJig= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= diff --git a/deployment/go.mod b/deployment/go.mod index 8ad400331ee..dd59da1204a 100644 --- a/deployment/go.mod +++ b/deployment/go.mod @@ -48,7 +48,7 @@ require ( github.com/smartcontractkit/chainlink-ccip/deployment v0.0.0-20261001213322-a2d686f610ca github.com/smartcontractkit/chainlink-common v0.11.2-0.20261007174852-7fabc093ff84 github.com/smartcontractkit/chainlink-common/keystore v1.3.1-0.20260903141829-ef07b52a737d - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 diff --git a/deployment/go.sum b/deployment/go.sum index 2edf811379f..18ab881b8f4 100644 --- a/deployment/go.sum +++ b/deployment/go.sum @@ -1417,8 +1417,8 @@ github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.202610051 github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa h1:7yIjZpk1mYYu630yXcGR7N+a+pcAHvn4OHlfhJK8GGY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa/go.mod h1:nfT8vdE44iAlhDWDV6PaBWucxZZAcDsgGkZ+D0kHpAI= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 h1:CeN/QrCpLFCdWBe/7Tzh/HuXjgTtSq/98YYj40/hHlQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3/go.mod h1:7jQ+XcxZjY19FaWPkvmckf/mT8A+p230BvLMVeEjJig= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= diff --git a/go.mod b/go.mod index aaadfb7b2b7..24afa53726f 100644 --- a/go.mod +++ b/go.mod @@ -85,7 +85,7 @@ require ( github.com/smartcontractkit/chainlink-common v0.11.2-0.20261007174852-7fabc093ff84 github.com/smartcontractkit/chainlink-common/keystore v1.3.1-0.20260903141829-ef07b52a737d github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 github.com/smartcontractkit/chainlink-evm/gethwrappers v0.0.0-20260915165527-3701875605f4 diff --git a/go.sum b/go.sum index 293e5d3a3c0..6f71010fdfa 100644 --- a/go.sum +++ b/go.sum @@ -1112,8 +1112,8 @@ github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.202610051 github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa h1:7yIjZpk1mYYu630yXcGR7N+a+pcAHvn4OHlfhJK8GGY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa/go.mod h1:nfT8vdE44iAlhDWDV6PaBWucxZZAcDsgGkZ+D0kHpAI= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8/go.mod h1:r8xuVq198XVbQSIfq5N5+W+ubnTmsjht/fn8xwAwYYw= github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 h1:QJiXTG9CmaQAuMRn5JGi+Jhji7fSkehVnKpjc8oNJJY= diff --git a/integration-tests/go.mod b/integration-tests/go.mod index c1ba7167138..df582161cb6 100644 --- a/integration-tests/go.mod +++ b/integration-tests/go.mod @@ -413,7 +413,7 @@ require ( github.com/smartcontractkit/chainlink-ccv v0.14.0 // indirect github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 // indirect github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa // indirect - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 // indirect github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 // indirect github.com/smartcontractkit/chainlink-feeds v0.1.2-0.20250227211209-7cd000095135 // indirect github.com/smartcontractkit/chainlink-framework/chains v0.0.0-20260724153515-bb6a2de39bcb // indirect diff --git a/integration-tests/go.sum b/integration-tests/go.sum index 0626035ea37..840d7a3a27e 100644 --- a/integration-tests/go.sum +++ b/integration-tests/go.sum @@ -1402,8 +1402,8 @@ github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.202610051 github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa h1:7yIjZpk1mYYu630yXcGR7N+a+pcAHvn4OHlfhJK8GGY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa/go.mod h1:nfT8vdE44iAlhDWDV6PaBWucxZZAcDsgGkZ+D0kHpAI= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 h1:CeN/QrCpLFCdWBe/7Tzh/HuXjgTtSq/98YYj40/hHlQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3/go.mod h1:7jQ+XcxZjY19FaWPkvmckf/mT8A+p230BvLMVeEjJig= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= diff --git a/integration-tests/load/go.mod b/integration-tests/load/go.mod index 44892354fb8..3044231563d 100644 --- a/integration-tests/load/go.mod +++ b/integration-tests/load/go.mod @@ -484,7 +484,7 @@ require ( github.com/smartcontractkit/chainlink-common/keystore v1.3.1-0.20260903141829-ef07b52a737d // indirect github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 // indirect github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa // indirect - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 // indirect github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 // indirect github.com/smartcontractkit/chainlink-evm/gethwrappers v0.0.0-20261005112317-b723176adfe8 // indirect github.com/smartcontractkit/chainlink-feeds v0.1.2-0.20250227211209-7cd000095135 // indirect diff --git a/integration-tests/load/go.sum b/integration-tests/load/go.sum index 8bf8a5c90ee..4aef1454001 100644 --- a/integration-tests/load/go.sum +++ b/integration-tests/load/go.sum @@ -1642,8 +1642,8 @@ github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.202610051 github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa h1:7yIjZpk1mYYu630yXcGR7N+a+pcAHvn4OHlfhJK8GGY= github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa/go.mod h1:nfT8vdE44iAlhDWDV6PaBWucxZZAcDsgGkZ+D0kHpAI= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 h1:CeN/QrCpLFCdWBe/7Tzh/HuXjgTtSq/98YYj40/hHlQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3/go.mod h1:7jQ+XcxZjY19FaWPkvmckf/mT8A+p230BvLMVeEjJig= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= diff --git a/plugins/plugins.public.yaml b/plugins/plugins.public.yaml index 5367d0b83f8..2e170497892 100644 --- a/plugins/plugins.public.yaml +++ b/plugins/plugins.public.yaml @@ -43,7 +43,7 @@ plugins: streams: - moduleURI: "github.com/smartcontractkit/chainlink-data-streams" - gitRef: "v1.1.2-0.20261002081259-6c2163b21db9" + gitRef: "v1.1.2-0.20261009064810-68424c714170" installPath: "./mercury/cmd/chainlink-mercury" ton: diff --git a/system-tests/lib/go.mod b/system-tests/lib/go.mod index dc69d1a490e..7aa21a29a74 100644 --- a/system-tests/lib/go.mod +++ b/system-tests/lib/go.mod @@ -493,7 +493,7 @@ require ( github.com/smartcontractkit/chainlink-ccv v0.14.0 // indirect github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 // indirect github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa // indirect - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 // indirect github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 // indirect github.com/smartcontractkit/chainlink-feeds v0.1.2-0.20250227211209-7cd000095135 // indirect github.com/smartcontractkit/chainlink-framework/chains v0.0.0-20260724153515-bb6a2de39bcb // indirect diff --git a/system-tests/lib/go.sum b/system-tests/lib/go.sum index f5e05fac052..b33b7a8ea1c 100644 --- a/system-tests/lib/go.sum +++ b/system-tests/lib/go.sum @@ -1669,8 +1669,8 @@ github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0 github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa/go.mod h1:nfT8vdE44iAlhDWDV6PaBWucxZZAcDsgGkZ+D0kHpAI= github.com/smartcontractkit/chainlink-confidential-compute v1.3.0 h1:y64hOzM9H61E4EYRzQDtRwqunDBlKpsZc1SvFEUlSLE= github.com/smartcontractkit/chainlink-confidential-compute v1.3.0/go.mod h1:fq9n85XoREIUxgdXIX3RyDIRCRyxZKZ5NCBbFDfdySo= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 h1:CeN/QrCpLFCdWBe/7Tzh/HuXjgTtSq/98YYj40/hHlQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3/go.mod h1:7jQ+XcxZjY19FaWPkvmckf/mT8A+p230BvLMVeEjJig= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= diff --git a/system-tests/tests/go.mod b/system-tests/tests/go.mod index 12430de9a6d..af902752f5b 100644 --- a/system-tests/tests/go.mod +++ b/system-tests/tests/go.mod @@ -628,7 +628,7 @@ require ( github.com/smartcontractkit/chainlink-ccip/deployment v0.0.0-20261001213322-a2d686f610ca // indirect github.com/smartcontractkit/chainlink-ccv v0.14.0 // indirect github.com/smartcontractkit/chainlink-common/x/config v0.0.0-20260928140953-a9b0d3a04daa // indirect - github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect + github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 // indirect github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 // indirect github.com/smartcontractkit/chainlink-feeds v0.1.2-0.20250227211209-7cd000095135 // indirect github.com/smartcontractkit/chainlink-framework/chains v0.0.0-20260724153515-bb6a2de39bcb // indirect diff --git a/system-tests/tests/go.sum b/system-tests/tests/go.sum index 55d314fe65c..972d2f3c79e 100644 --- a/system-tests/tests/go.sum +++ b/system-tests/tests/go.sum @@ -1850,8 +1850,8 @@ github.com/smartcontractkit/chainlink-confidential-compute v1.3.0 h1:y64hOzM9H61 github.com/smartcontractkit/chainlink-confidential-compute v1.3.0/go.mod h1:fq9n85XoREIUxgdXIX3RyDIRCRyxZKZ5NCBbFDfdySo= github.com/smartcontractkit/chainlink-confidential-compute/tests/testhelpers v0.0.0-20260812145307-d77342c53d7d h1:CZ0Om7lANyhpFJcmZE8RwQpEP4PyvoNLhOJckjm+zao= github.com/smartcontractkit/chainlink-confidential-compute/tests/testhelpers v0.0.0-20260812145307-d77342c53d7d/go.mod h1:Q5q/ohoAF5N7GXZS6QxgPWSaSLFtS6/53QHJbBUkCmU= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 h1:C7Dv8SRUU97tsUJvDDOmntJI6RZxQflIi88NlvdUnCQ= -github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170 h1:HeQlHO9MrUXTJmxYCNX4xPX14glLAdwIN3igc9Am/g8= +github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261009064810-68424c714170/go.mod h1:ExmEs17bPBrsjV8XUs+Om5vvT/twPHeYfmAMwEnrtQQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3 h1:CeN/QrCpLFCdWBe/7Tzh/HuXjgTtSq/98YYj40/hHlQ= github.com/smartcontractkit/chainlink-deployments-framework v0.123.3/go.mod h1:7jQ+XcxZjY19FaWPkvmckf/mT8A+p230BvLMVeEjJig= github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8 h1:49f3Vk1Igc5UTQc26lYrS1/+du+imLBKoYW2IX/SiCM= From 73758ad0359189b0dd27c3f5473f410a422a6b75 Mon Sep 17 00:00:00 2001 From: Bruno Moura Date: Fri, 9 Oct 2026 07:54:52 +0100 Subject: [PATCH 2/6] synchronization: add the LLOAttributedObservation telemetry type --- core/services/synchronization/common.go | 4 ++++ core/services/synchronization/common_test.go | 6 ++++++ 2 files changed, 10 insertions(+) diff --git a/core/services/synchronization/common.go b/core/services/synchronization/common.go index e0ef75a8bd2..cd0917c7c80 100644 --- a/core/services/synchronization/common.go +++ b/core/services/synchronization/common.go @@ -33,6 +33,8 @@ const ( LLOObservation TelemetryType = "llo-observation" LLOOutcome TelemetryType = "llo-outcome" LLOReport TelemetryType = "llo-report" + + LLOAttributedObservation TelemetryType = "llo-attributed-observation" ) type TelemPayload struct { @@ -91,6 +93,8 @@ func TelemetryTypeToDomainAndEntity(telemType TelemetryType) (domain, entity str return "data-streams.telemetry.llo-observation", "telem.LLOObservationTelemetry", nil case LLOOutcome: return "data-streams.telemetry.llo-outcome", "telem.LLOOutcomeTelemetry", nil + case LLOAttributedObservation: + return "data-streams.telemetry.llo-attributed-observation", "telem.LLOAttributedObservationTelemetry", nil case FunctionsRequests: return "functions.telemetry.functions-requests", "telem.FunctionsRequest", nil case HeadReport: diff --git a/core/services/synchronization/common_test.go b/core/services/synchronization/common_test.go index e6785f3b2d7..95cbe1cee88 100644 --- a/core/services/synchronization/common_test.go +++ b/core/services/synchronization/common_test.go @@ -93,6 +93,12 @@ func TestTelemetryTypeToDomainAndEntity(t *testing.T) { expectedDomain: "data-streams.telemetry.llo-outcome", expectedEntity: "telem.LLOOutcomeTelemetry", }, + { + name: "LLOAttributedObservation", + telemType: LLOAttributedObservation, + expectedDomain: "data-streams.telemetry.llo-attributed-observation", + expectedEntity: "telem.LLOAttributedObservationTelemetry", + }, { name: "FunctionsRequests", telemType: FunctionsRequests, From 58a458f1af8aca1200ee5a5d5d679292fd311171 Mon Sep 17 00:00:00 2001 From: Bruno Moura Date: Fri, 9 Oct 2026 07:58:17 +0100 Subject: [PATCH 3/6] llo/telem: LLOAttributedObservationTelemetry support --- .../llo/observation/data_source_test.go | 5 +- core/services/llo/telem/sampling.go | 26 ++++++ core/services/llo/telem/sampling_test.go | 54 ++++++++++++ core/services/llo/telem/telemetry.go | 22 +++++ core/services/llo/telem/telemetry_test.go | 84 +++++++++++++++++++ 5 files changed, 190 insertions(+), 1 deletion(-) diff --git a/core/services/llo/observation/data_source_test.go b/core/services/llo/observation/data_source_test.go index 1dc2b2db5a5..c78dcf695d4 100644 --- a/core/services/llo/observation/data_source_test.go +++ b/core/services/llo/observation/data_source_test.go @@ -179,7 +179,10 @@ func (m *mockTelemeter) GetOutcomeTelemetryCh() chan<- *lloprotocol.LLOOutcomeTe return nil } func (m *mockTelemeter) GetReportTelemetryCh() chan<- *lloprotocol.LLOReportTelemetry { return nil } -func (m *mockTelemeter) CaptureEATelemetry() bool { return true } +func (m *mockTelemeter) GetAttributedObservationTelemetryCh() chan<- *lloprotocol.LLOAttributedObservationTelemetry { + return nil +} +func (m *mockTelemeter) CaptureEATelemetry() bool { return true } func (m *mockTelemeter) CaptureObservationTelemetry() bool { return true } diff --git a/core/services/llo/telem/sampling.go b/core/services/llo/telem/sampling.go index 03bb17585f6..a5fa0cf9604 100644 --- a/core/services/llo/telem/sampling.go +++ b/core/services/llo/telem/sampling.go @@ -153,6 +153,20 @@ func fingerprint(typ synchronization.TelemetryType, msg proto.Message) (string, hex.EncodeToString(m.ConfigDigest), } return strings.Join(traits, samplerDelimiter), nanosToSec(int64(m.ObservationTimestampNanoseconds)), nil //nolint:gosec // G115 + case synchronization.LLOAttributedObservation: + m, ok := msg.(*lloprotocol.LLOAttributedObservationTelemetry) + if !ok || m == nil { + return "", 0, errors.New("invalid telemetry type, expected LLOAttributedObservationTelemetry") + } + // A large observation is split in parts sharing the observer: key on the + // lowest stream id so each part is sampled on its own. + traits := []string{ + strconv.FormatUint(uint64(m.DonId), 10), + strconv.FormatUint(uint64(m.Observer), 10), + hex.EncodeToString(m.ConfigDigest), + strconv.FormatUint(uint64(lowestStreamID(m.StreamValues)), 10), + } + return strings.Join(traits, samplerDelimiter), nanosToSec(int64(m.AgreedObservationTimestampNanoseconds)), nil //nolint:gosec // G115 case synchronization.PipelineBridge: m, ok := msg.(*LLOBridgeTelemetry) if !ok || m == nil { @@ -171,6 +185,18 @@ func fingerprint(typ synchronization.TelemetryType, msg proto.Message) (string, } } +// lowestStreamID returns the lowest key of values, or 0 when empty. +func lowestStreamID(values map[uint32]*lloprotocol.LLOStreamValue) uint32 { + var lowest uint32 + first := true + for id := range values { + if first || id < lowest { + lowest, first = id, false + } + } + return lowest +} + func nanosToSec(n int64) int32 { return int32(n / int64(time.Second)) //nolint:gosec // G115 } diff --git a/core/services/llo/telem/sampling_test.go b/core/services/llo/telem/sampling_test.go index 0a6e693b430..1dbc2a78976 100644 --- a/core/services/llo/telem/sampling_test.go +++ b/core/services/llo/telem/sampling_test.go @@ -73,6 +73,33 @@ func TestFingerprint(t *testing.T) { ts: int32(ot.Unix()), //nolint:gosec // G115 err: nil, }, + { + name: "successful attributed observation", + msg: &lloprotocol.LLOAttributedObservationTelemetry{ + DonId: donID, + Observer: 3, + ConfigDigest: configDigest, + AgreedObservationTimestampNanoseconds: uint64(ot.UnixNano()), + StreamValues: map[uint32]*lloprotocol.LLOStreamValue{streamID: {}, streamID + 7: {}}, + }, + typ: synchronization.LLOAttributedObservation, + fingerprint: fmt.Sprintf("%d-%d-%x-%d", donID, 3, configDigest, streamID), + ts: int32(ot.Unix()), //nolint:gosec // G115 + err: nil, + }, + { + name: "attributed observation without values", + msg: &lloprotocol.LLOAttributedObservationTelemetry{ + DonId: donID, + Observer: 3, + ConfigDigest: configDigest, + AgreedObservationTimestampNanoseconds: uint64(ot.UnixNano()), + }, + typ: synchronization.LLOAttributedObservation, + fingerprint: fmt.Sprintf("%d-%d-%x-%d", donID, 3, configDigest, 0), + ts: int32(ot.Unix()), //nolint:gosec // G115 + err: nil, + }, { name: "successful bridge", msg: &LLOBridgeTelemetry{ @@ -139,6 +166,33 @@ func TestSample(t *testing.T) { assert.False(t, shouldSend) } +// TestSample_AttributedObservationParts ensures the parts of a split +// observation are sampled independently, and a repeated part is deduplicated. +func TestSample_AttributedObservationParts(t *testing.T) { + t.Parallel() + + samplr := newSampler(logger.TestSugared(t), true) + ts := uint64(time.Unix(1600000000, 0).UnixNano()) + part := func(observer uint32, ids ...uint32) *lloprotocol.LLOAttributedObservationTelemetry { + m := &lloprotocol.LLOAttributedObservationTelemetry{ + DonId: 2, + Observer: observer, + ConfigDigest: []byte("digest"), + AgreedObservationTimestampNanoseconds: ts, + StreamValues: map[uint32]*lloprotocol.LLOStreamValue{}, + } + for _, id := range ids { + m.StreamValues[id] = &lloprotocol.LLOStreamValue{} + } + return m + } + + assert.True(t, samplr.Sample(synchronization.LLOAttributedObservation, part(1, 1, 2))) + assert.True(t, samplr.Sample(synchronization.LLOAttributedObservation, part(1, 3, 4)), "second part of the same observation") + assert.True(t, samplr.Sample(synchronization.LLOAttributedObservation, part(2, 1, 2)), "another observer") + assert.False(t, samplr.Sample(synchronization.LLOAttributedObservation, part(1, 1, 2)), "repeated part") +} + // TestPruningLoop ensures the pruning loop works as expected. func TestPruningLoop(t *testing.T) { if testing.Short() { diff --git a/core/services/llo/telem/telemetry.go b/core/services/llo/telem/telemetry.go index 572427659e4..4ed9ac17f6e 100644 --- a/core/services/llo/telem/telemetry.go +++ b/core/services/llo/telem/telemetry.go @@ -36,6 +36,7 @@ type Telemeter interface { MakeObservationScopedTelemetryCh(opts DSOpts, size int) (ch chan<- any) GetOutcomeTelemetryCh() chan<- *lloprotocol.LLOOutcomeTelemetry GetReportTelemetryCh() chan<- *lloprotocol.LLOReportTelemetry + GetAttributedObservationTelemetryCh() chan<- *lloprotocol.LLOAttributedObservationTelemetry CaptureEATelemetry() bool CaptureObservationTelemetry() bool TrackSeqNr(digest types.ConfigDigest, seqNr uint64) @@ -55,6 +56,8 @@ type TelemeterParams struct { CaptureOutcomeTelemetry bool CaptureReportTelemetry bool SampleTelemetry bool + // CaptureAttributedObservationTelemetry is only honoured by llo/v31. + CaptureAttributedObservationTelemetry bool } func NewTelemeterService(params TelemeterParams) TelemeterService { @@ -91,6 +94,11 @@ func newTelemeter(params TelemeterParams) *telemeter { if params.CaptureReportTelemetry { t.chReportTelemetry = make(chan *lloprotocol.LLOReportTelemetry, (2+2)*lloprotocol.MaxReportCount) // 2 instances+2x size safety buffer } + if params.CaptureAttributedObservationTelemetry { + // One per observer per round from f+1 rotating emitters, more when + // large observations are split. + t.chAttributedObservationTelemetry = make(chan *lloprotocol.LLOAttributedObservationTelemetry, 1000) + } t.Service, t.eng = services.Config{ Name: "LLOTelemeterService", Start: t.start, @@ -107,6 +115,9 @@ func newTelemeter(params TelemeterParams) *telemeter { if t.chReportTelemetry != nil { close(t.chReportTelemetry) } + if t.chAttributedObservationTelemetry != nil { + close(t.chAttributedObservationTelemetry) + } close(t.chTransmissionSeqNr) return nil @@ -135,6 +146,8 @@ type telemeter struct { chOutcomeTelemetry chan *lloprotocol.LLOOutcomeTelemetry chReportTelemetry chan *lloprotocol.LLOReportTelemetry + chAttributedObservationTelemetry chan *lloprotocol.LLOAttributedObservationTelemetry + currentSeqNrMu sync.Mutex currentSeqNr map[string]uint64 chTransmissionSeqNr chan struct { @@ -216,6 +229,10 @@ func (t *telemeter) GetReportTelemetryCh() chan<- *lloprotocol.LLOReportTelemetr return t.chReportTelemetry } +func (t *telemeter) GetAttributedObservationTelemetryCh() chan<- *lloprotocol.LLOAttributedObservationTelemetry { + return t.chAttributedObservationTelemetry +} + func (t *telemeter) CaptureEATelemetry() bool { return t.captureEATelemetry } @@ -251,6 +268,8 @@ func (t *telemeter) start(_ context.Context) error { t.enqueueTelemetry(types.ConfigDigest(rt.ConfigDigest).Hex(), rt.SeqNr, synchronization.LLOOutcome, rt) case rt := <-t.chReportTelemetry: t.enqueueTelemetry(types.ConfigDigest(rt.ConfigDigest).Hex(), rt.SeqNr, synchronization.LLOReport, rt) + case at := <-t.chAttributedObservationTelemetry: + t.enqueueTelemetry(types.ConfigDigest(at.ConfigDigest).Hex(), at.SeqNr, synchronization.LLOAttributedObservation, at) case tx := <-t.chTransmissionSeqNr: // Drain any pending outcome or report telemetry before sending buffered telemetry t.sendBufferedTelemetry(tx.digest, tx.seqNr) @@ -502,6 +521,9 @@ func (t *nullTelemeter) GetOutcomeTelemetryCh() chan<- *lloprotocol.LLOOutcomeTe func (t *nullTelemeter) GetReportTelemetryCh() chan<- *lloprotocol.LLOReportTelemetry { return nil } +func (t *nullTelemeter) GetAttributedObservationTelemetryCh() chan<- *lloprotocol.LLOAttributedObservationTelemetry { + return nil +} func (t *nullTelemeter) CaptureEATelemetry() bool { return false } diff --git a/core/services/llo/telem/telemetry_test.go b/core/services/llo/telem/telemetry_test.go index 01231f6ce90..2bcf2d0bef5 100644 --- a/core/services/llo/telem/telemetry_test.go +++ b/core/services/llo/telem/telemetry_test.go @@ -3,6 +3,7 @@ package telem import ( "encoding/hex" "errors" + "slices" "testing" "time" @@ -1193,3 +1194,86 @@ func Test_Telemeter_reportTelemetry_samplingAtFlushTime(t *testing.T) { "each per-channel report should be admitted (distinct sampler fingerprints)") }) } + +func Test_Telemeter_attributedObservationTelemetry(t *testing.T) { + t.Parallel() + + lggr := logger.TestLogger(t) + donID := uint32(1) + + t.Run("returns nil channel if CaptureAttributedObservationTelemetry is false", func(t *testing.T) { + t.Parallel() + tm := newTelemeter(TelemeterParams{ + Logger: lggr, + DonID: donID, + }) + assert.Nil(t, tm.GetAttributedObservationTelemetryCh()) + assert.Nil(t, NullTelemeter.GetAttributedObservationTelemetryCh()) + }) + + t.Run("buffers every part and flushes on transmission", func(t *testing.T) { + t.Parallel() + m := &mockMonitoringEndpoint{chTypedLogs: make(chan typedLog, 100)} + tm := newTelemeter(TelemeterParams{ + Logger: lggr, + MonitoringEndpoint: m, + DonID: donID, + CaptureAttributedObservationTelemetry: true, + }) + servicetest.Run(t, tm) + ch := tm.GetAttributedObservationTelemetryCh() + require.NotNil(t, ch) + + opts := &mockOpts{} + cd := opts.ConfigDigest() + parts := []*lloprotocol.LLOAttributedObservationTelemetry{ + { + ConfigDigest: cd[:], + SeqNr: opts.SeqNr(), + DonId: donID, + Observer: 2, + Emitter: 3, + OracleObservationTimestampNanoseconds: 4, + AgreedObservationTimestampNanoseconds: 5, + StreamValues: map[uint32]*lloprotocol.LLOStreamValue{1: {Type: 1, Value: []byte{6}}}, + RemoveChannelIds: []uint32{7}, + }, + { + ConfigDigest: cd[:], + SeqNr: opts.SeqNr(), + DonId: donID, + Observer: 2, + Emitter: 3, + OracleObservationTimestampNanoseconds: 4, + AgreedObservationTimestampNanoseconds: 5, + StreamValues: map[uint32]*lloprotocol.LLOStreamValue{8: {Type: 1, Value: []byte{9}}}, + }, + } + for _, p := range parts { + ch <- p + } + + testutils.RequireEventually(t, func() bool { + tm.telemetryBufferMu.Lock() + defer tm.telemetryBufferMu.Unlock() + return len(tm.telemetryBuffer[cd.Hex()][opts.SeqNr()]) == len(parts) + }) + + tm.TrackSeqNr(opts.ConfigDigest(), opts.SeqNr()) + + received := make([]*lloprotocol.LLOAttributedObservationTelemetry, 0, len(parts)) + for range parts { + tLog := <-m.chTypedLogs + assert.Equal(t, synchronization.LLOAttributedObservation, tLog.telemType) + decoded := &lloprotocol.LLOAttributedObservationTelemetry{} + require.NoError(t, proto.Unmarshal(tLog.log, decoded)) + received = append(received, decoded) + } + for i, want := range parts { + idx := slices.IndexFunc(received, func(got *lloprotocol.LLOAttributedObservationTelemetry) bool { + return proto.Equal(want, got) + }) + assert.GreaterOrEqual(t, idx, 0, "part %d not flushed", i) + } + }) +} From 4c04668ffd704329b2efe823ad57d35f5b971719 Mon Sep 17 00:00:00 2001 From: Bruno Moura Date: Fri, 9 Oct 2026 08:00:58 +0100 Subject: [PATCH 4/6] llo/telem: prune buffered telemetry per config digest --- core/services/llo/telem/telemetry.go | 50 ++++++++++++++----- core/services/llo/telem/telemetry_test.go | 58 +++++++++++++++++++++++ 2 files changed, 97 insertions(+), 11 deletions(-) diff --git a/core/services/llo/telem/telemetry.go b/core/services/llo/telem/telemetry.go index 4ed9ac17f6e..8f429834968 100644 --- a/core/services/llo/telem/telemetry.go +++ b/core/services/llo/telem/telemetry.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "sync" + "time" "google.golang.org/protobuf/proto" @@ -26,6 +27,11 @@ import ( const adapterLWBAErrorName = "AdapterLWBAError" +// bufferedTelemetryDigestTTL is how long the buffered telemetry of a config +// digest is kept without a flush. A retired or removed instance stops +// transmitting, so its buffer is evicted after this long. +const bufferedTelemetryDigestTTL = time.Minute + // DSOpts is the shared, version-agnostic LLO data-source options (llo/v30 and // llo/v31 both use llodatasource.DSOpts). Aliased here so the telemetry and // observation paths keep referring to telem.DSOpts. @@ -86,6 +92,8 @@ func newTelemeter(params TelemeterParams) *telemeter { }, 10), currentSeqNr: make(map[string]uint64), telemetryBuffer: make(map[string]map[uint64][]telemetryEntry), + bufferFlushedAt: make(map[string]time.Time), + bufferTTL: bufferedTelemetryDigestTTL, sampler: newSampler(logger.Sugared(params.Logger), params.SampleTelemetry), } if params.CaptureOutcomeTelemetry { @@ -159,6 +167,11 @@ type telemeter struct { // for transmitting rounds sequence numbers telemetryBufferMu sync.Mutex telemetryBuffer map[string]map[uint64][]telemetryEntry + // bufferFlushedAt is when each digest buffer was created or last flushed. + // Each digest is flushed by its own transmissions, so the buffers of + // concurrent instances (blue/green) are independent. + bufferFlushedAt map[string]time.Time + bufferTTL time.Duration sampler *sampler } @@ -312,12 +325,23 @@ func (t *telemeter) sendBufferedTelemetry(digest types.ConfigDigest, seqNr uint6 } } - // drop messages for config digests that are not transmitting - for d := range t.telemetryBuffer { - if d == cd { - continue + // evict the buffers of digests that stopped transmitting + now := time.Now() + t.bufferFlushedAt[cd] = now + var evicted []string + for d, flushedAt := range t.bufferFlushedAt { + if now.Sub(flushedAt) > t.bufferTTL { + delete(t.telemetryBuffer, d) + delete(t.bufferFlushedAt, d) + evicted = append(evicted, d) + } + } + if len(evicted) > 0 { + t.currentSeqNrMu.Lock() + for _, d := range evicted { + delete(t.currentSeqNr, d) } - delete(t.telemetryBuffer, d) + t.currentSeqNrMu.Unlock() } go func() { @@ -359,9 +383,7 @@ func (t *telemeter) enqueueTelemetry(digest string, seqNr uint64, typ synchroniz t.telemetryBufferMu.Lock() defer t.telemetryBufferMu.Unlock() - if _, ok := t.telemetryBuffer[digest]; !ok { - t.telemetryBuffer[digest] = make(map[uint64][]telemetryEntry) - } + t.ensureDigestBuffer(digest) // Sampling is applied at flush time for buffered telemetry t.telemetryBuffer[digest][seqNr] = []telemetryEntry{{ @@ -374,9 +396,7 @@ func (t *telemeter) enqueueTelemetry(digest string, seqNr uint64, typ synchroniz t.telemetryBufferMu.Lock() defer t.telemetryBufferMu.Unlock() - if _, ok := t.telemetryBuffer[digest]; !ok { - t.telemetryBuffer[digest] = make(map[uint64][]telemetryEntry) - } + t.ensureDigestBuffer(digest) // Sampling is applied at flush time for buffered telemetry t.telemetryBuffer[digest][seqNr] = append(t.telemetryBuffer[digest][seqNr], telemetryEntry{ telemType: typ, @@ -385,6 +405,14 @@ func (t *telemeter) enqueueTelemetry(digest string, seqNr uint64, typ synchroniz } } +// ensureDigestBuffer creates the buffer of digest. Callers hold telemetryBufferMu. +func (t *telemeter) ensureDigestBuffer(digest string) { + if _, ok := t.telemetryBuffer[digest]; !ok { + t.telemetryBuffer[digest] = make(map[uint64][]telemetryEntry) + t.bufferFlushedAt[digest] = time.Now() + } +} + func (t *telemeter) prepareObservationTelemetry(p any, opts DSOpts) { var telemType synchronization.TelemetryType var msg proto.Message diff --git a/core/services/llo/telem/telemetry_test.go b/core/services/llo/telem/telemetry_test.go index 2bcf2d0bef5..d8487cea53e 100644 --- a/core/services/llo/telem/telemetry_test.go +++ b/core/services/llo/telem/telemetry_test.go @@ -1277,3 +1277,61 @@ func Test_Telemeter_attributedObservationTelemetry(t *testing.T) { } }) } + +func Test_Telemeter_bufferedTelemetryPerDigest(t *testing.T) { + t.Parallel() + + lggr := logger.TestLogger(t) + outcome := func(digest ocr2types.ConfigDigest, seqNr uint64) *lloprotocol.LLOOutcomeTelemetry { + return &lloprotocol.LLOOutcomeTelemetry{ConfigDigest: digest[:], SeqNr: seqNr} + } + blue, green := ocr2types.ConfigDigest{0xb1}, ocr2types.ConfigDigest{0x91} + + t.Run("a transmission of one digest keeps the other digest buffer", func(t *testing.T) { + t.Parallel() + m := &mockMonitoringEndpoint{chTypedLogs: make(chan typedLog, 10)} + tm := newTelemeter(TelemeterParams{Logger: lggr, MonitoringEndpoint: m}) + + tm.enqueueTelemetry(blue.Hex(), 10, synchronization.LLOOutcome, outcome(blue, 10)) + tm.enqueueTelemetry(green.Hex(), 3, synchronization.LLOOutcome, outcome(green, 3)) + + tm.sendBufferedTelemetry(blue, 10) + <-m.chTypedLogs + + tm.telemetryBufferMu.Lock() + require.Len(t, tm.telemetryBuffer[green.Hex()][3], 1, "green buffer must survive blue's transmission") + tm.telemetryBufferMu.Unlock() + + tm.sendBufferedTelemetry(green, 3) + tLog := <-m.chTypedLogs + decoded := &lloprotocol.LLOOutcomeTelemetry{} + require.NoError(t, proto.Unmarshal(tLog.log, decoded)) + assert.Equal(t, green[:], decoded.ConfigDigest) + assert.Equal(t, uint64(3), decoded.SeqNr) + }) + + t.Run("evicts a digest buffer idle for longer than the TTL", func(t *testing.T) { + t.Parallel() + m := &mockMonitoringEndpoint{chTypedLogs: make(chan typedLog, 10)} + tm := newTelemeter(TelemeterParams{Logger: lggr, MonitoringEndpoint: m}) + tm.bufferTTL = time.Millisecond + + tm.enqueueTelemetry(green.Hex(), 3, synchronization.LLOOutcome, outcome(green, 3)) + tm.sendBufferedTelemetry(green, 2) // tracks green, keeps seqNr 3 buffered + time.Sleep(5 * time.Millisecond) + + tm.enqueueTelemetry(blue.Hex(), 10, synchronization.LLOOutcome, outcome(blue, 10)) + tm.sendBufferedTelemetry(blue, 10) + <-m.chTypedLogs + + tm.telemetryBufferMu.Lock() + assert.NotContains(t, tm.telemetryBuffer, green.Hex()) + assert.NotContains(t, tm.bufferFlushedAt, green.Hex()) + assert.Contains(t, tm.bufferFlushedAt, blue.Hex()) + tm.telemetryBufferMu.Unlock() + + tm.currentSeqNrMu.Lock() + assert.NotContains(t, tm.currentSeqNr, green.Hex()) + tm.currentSeqNrMu.Unlock() + }) +} From bacfa4c0469e5e62759a6c7aeb1ce9a5da3c2faf Mon Sep 17 00:00:00 2001 From: Bruno Moura Date: Fri, 9 Oct 2026 08:11:56 +0100 Subject: [PATCH 5/6] llo: gate plugin telemetry per plugin version --- core/services/llo/delegate.go | 83 +++++++++++++++++++----------- core/services/llo/delegate_test.go | 72 ++++++++++++++++++++++++++ 2 files changed, 125 insertions(+), 30 deletions(-) diff --git a/core/services/llo/delegate.go b/core/services/llo/delegate.go index aaa7f9a899a..8c143f84869 100644 --- a/core/services/llo/delegate.go +++ b/core/services/llo/delegate.go @@ -69,8 +69,10 @@ type DelegateConfig struct { JobName null.String CaptureEATelemetry bool CaptureObservationTelemetry bool - CaptureOutcomeTelemetry bool - CaptureReportTelemetry bool + // CaptureOutcomeTelemetry and CaptureReportTelemetry apply to v30 instances. + // v31 instances read their plugin telemetry flags from V31Config. + CaptureOutcomeTelemetry bool + CaptureReportTelemetry bool // LLO ChannelDefinitionCache llotypes.ChannelDefinitionCache @@ -116,6 +118,22 @@ type DelegateConfig struct { KeyValueDatabaseFactory ocr3_1types.KeyValueDatabaseFactory } +// telemeterParams creates a telemetry channel when either plugin version needs +// it. Each factory is only handed the channels its own version enables. +func (cfg DelegateConfig) telemeterParams(lggr logger.Logger) telem.TelemeterParams { + return telem.TelemeterParams{ + Logger: lggr, + MonitoringEndpoint: cfg.PluginMonitoringEndpoint, + DonID: cfg.DonID, + CaptureEATelemetry: cfg.CaptureEATelemetry, + CaptureObservationTelemetry: cfg.CaptureObservationTelemetry, + CaptureOutcomeTelemetry: cfg.CaptureOutcomeTelemetry || cfg.V31Config.CaptureOutcomeTelemetry, + CaptureReportTelemetry: cfg.CaptureReportTelemetry || cfg.V31Config.CaptureReportTelemetry, + CaptureAttributedObservationTelemetry: cfg.V31Config.CaptureAttributedObservationTelemetry, + SampleTelemetry: cfg.SampleTelemetry, + } +} + // anyV31 reports whether any protocol instance runs the v31 plugin. The // OCR3.1-only dependencies are per job, so one v31 instance requires them. func (cfg DelegateConfig) anyV31() bool { @@ -174,16 +192,7 @@ func NewDelegate(cfg DelegateConfig) (job.ServiceCtx, error) { } reportCodecs := NewReportCodecs(codecLggr, cfg.DonID) - t := telem.NewTelemeterService(telem.TelemeterParams{ - Logger: lggr, - MonitoringEndpoint: cfg.PluginMonitoringEndpoint, - DonID: cfg.DonID, - CaptureEATelemetry: cfg.CaptureEATelemetry, - CaptureObservationTelemetry: cfg.CaptureObservationTelemetry, - CaptureOutcomeTelemetry: cfg.CaptureOutcomeTelemetry, - CaptureReportTelemetry: cfg.CaptureReportTelemetry, - SampleTelemetry: cfg.SampleTelemetry, - }) + t := telem.NewTelemeterService(cfg.telemeterParams(lggr)) ds := observation.NewDataSource(logger.Named(lggr, "DataSource"), cfg.Registry, t) @@ -262,22 +271,7 @@ func (d *delegate) newOracleV30(i int, configTracker ocr2types.ContractConfigTra OffchainKeyring: d.cfg.OffchainKeyring, OnchainKeyring: ocr3shims.OnchainKeyringAsOnchainKeyring2(d.cfg.OnchainKeyring), ReportingPluginFactory: promwrapper.NewReportingPluginFactory( - llov30.NewPluginFactory( - llov30.PluginFactoryParams{ - Config: d.cfg.ReportingPluginConfig, - PredecessorRetirementReportCache: psrrc, - ShouldRetireCache: d.src, - RetirementReportCodec: d.cfg.RetirementReportCodec, - ChannelDefinitionCache: d.cfg.ChannelDefinitionCache, - DataSource: d.ds, - Logger: logger.Named(lggr, "ReportingPlugin"), - OnchainConfigCodec: lloprotocol.EVMOnchainConfigCodec{}, - ReportCodecs: d.reportCodecs, - OutcomeTelemetryCh: d.telem.GetOutcomeTelemetryCh(), - ReportTelemetryCh: d.telem.GetReportTelemetryCh(), - DonID: d.cfg.DonID, - }, - ), + llov30.NewPluginFactory(d.v30FactoryParams(lggr, psrrc)), lggr, "", d.cfg.ChainID, @@ -287,12 +281,40 @@ func (d *delegate) newOracleV30(i int, configTracker ocr2types.ContractConfigTra }) } +// v30FactoryParams assembles the v30 plugin factory params. +func (d *delegate) v30FactoryParams(lggr logger.Logger, psrrc lloprotocol.PredecessorRetirementReportCache) llov30.PluginFactoryParams { + return llov30.PluginFactoryParams{ + Config: d.cfg.ReportingPluginConfig, + PredecessorRetirementReportCache: psrrc, + ShouldRetireCache: d.src, + RetirementReportCodec: d.cfg.RetirementReportCodec, + ChannelDefinitionCache: d.cfg.ChannelDefinitionCache, + DataSource: d.ds, + Logger: logger.Named(lggr, "ReportingPlugin"), + OnchainConfigCodec: lloprotocol.EVMOnchainConfigCodec{}, + ReportCodecs: d.reportCodecs, + OutcomeTelemetryCh: enabledCh(d.cfg.CaptureOutcomeTelemetry, d.telem.GetOutcomeTelemetryCh()), + ReportTelemetryCh: enabledCh(d.cfg.CaptureReportTelemetry, d.telem.GetReportTelemetryCh()), + DonID: d.cfg.DonID, + } +} + +// enabledCh returns ch when enabled and nil otherwise, so a protocol instance only +// emits the plugin telemetry its own version enables. +func enabledCh[T any](enabled bool, ch chan<- T) chan<- T { + if !enabled { + return nil + } + return ch +} + // v31FactoryParams assembles the v31 plugin factory params, mapping the job's // V31Config knobs onto it. Knobs left at zero are forwarded as zero, which the // factory reads as "apply the plugin default". func (d *delegate) v31FactoryParams(lggr logger.Logger, psrrc lloprotocol.PredecessorRetirementReportCache) llov31.PluginFactoryParams { return llov31.PluginFactoryParams{ VerboseLogging: d.cfg.ReportingPluginConfig.VerboseLogging || d.cfg.V31Config.VerboseLogging, + CaptureStagingTelemetry: d.cfg.V31Config.CaptureStagingTelemetry, PredecessorRetirementReportCache: psrrc, ShouldRetireCache: d.src, RetirementReportCodec: d.cfg.RetirementReportCodec, @@ -301,8 +323,9 @@ func (d *delegate) v31FactoryParams(lggr logger.Logger, psrrc lloprotocol.Predec Logger: logger.Named(lggr, "ReportingPlugin"), OnchainConfigCodec: lloprotocol.EVMOnchainConfigCodec{}, ReportCodecs: d.reportCodecs, - OutcomeTelemetryCh: d.telem.GetOutcomeTelemetryCh(), - ReportTelemetryCh: d.telem.GetReportTelemetryCh(), + OutcomeTelemetryCh: enabledCh(d.cfg.V31Config.CaptureOutcomeTelemetry, d.telem.GetOutcomeTelemetryCh()), + ReportTelemetryCh: enabledCh(d.cfg.V31Config.CaptureReportTelemetry, d.telem.GetReportTelemetryCh()), + AttributedObservationTelemetryCh: enabledCh(d.cfg.V31Config.CaptureAttributedObservationTelemetry, d.telem.GetAttributedObservationTelemetryCh()), DonID: d.cfg.DonID, MaxSnapshotRounds: d.cfg.V31Config.MaxSnapshotRounds, BlobLifetimeRounds: d.cfg.V31Config.BlobLifetimeRounds, diff --git a/core/services/llo/delegate_test.go b/core/services/llo/delegate_test.go index 8fa23291277..27369ccfd89 100644 --- a/core/services/llo/delegate_test.go +++ b/core/services/llo/delegate_test.go @@ -15,6 +15,7 @@ import ( lloconfig "github.com/smartcontractkit/chainlink-data-streams/llo/pluginconfig" llov30 "github.com/smartcontractkit/chainlink-data-streams/llo/v30" "github.com/smartcontractkit/chainlink/v2/core/services/llo/telem" + "github.com/smartcontractkit/chainlink/v2/core/services/synchronization" ) // newTestDelegate returns a delegate carrying only what v31FactoryParams reads. @@ -37,9 +38,11 @@ func Test_delegate_v31FactoryParams(t *testing.T) { BlobInFlightWaitFactor: 5, MaxBlobSnapshotAge: lloconfig.Duration(9 * time.Second), MaxRoundPeriod: lloconfig.Duration(time.Minute), + CaptureStagingTelemetry: true, }, }).v31FactoryParams(logger.Test(t), nil) + assert.True(t, params.CaptureStagingTelemetry) assert.Equal(t, uint64(7), params.MaxSnapshotRounds) assert.Equal(t, uint64(11), params.BlobLifetimeRounds) assert.Equal(t, 3*time.Second, params.MaxDurationBlobObservation) @@ -60,6 +63,7 @@ func Test_delegate_v31FactoryParams(t *testing.T) { assert.Zero(t, params.BlobInFlightWaitFactor) assert.Zero(t, params.MaxBlobSnapshotAge) assert.Zero(t, params.MaxRoundPeriod) + assert.Nil(t, params.AttributedObservationTelemetryCh) }) t.Run("forwards a negative MaxBlobSnapshotAge, which disables the age check", func(t *testing.T) { @@ -100,6 +104,74 @@ func Test_delegate_v31FactoryParams(t *testing.T) { }) } +// Test_delegate_pluginTelemetry checks that each plugin version only gets the +// telemetry channels its own config enables: v30 from the node driven flags, v31 +// from V31Config alone. +func Test_delegate_pluginTelemetry(t *testing.T) { + t.Parallel() + + type channels struct{ outcome, report, attributed bool } + for _, tc := range []struct { + name string + cfg DelegateConfig + v30, v31 channels + }{ + { + name: "all off", + }, + { + name: "node flags enable v30 only", + cfg: DelegateConfig{ + CaptureEATelemetry: true, + CaptureOutcomeTelemetry: true, + CaptureReportTelemetry: true, + }, + v30: channels{outcome: true, report: true}, + }, + { + name: "v31 flags enable v31 only, without the node flag", + cfg: DelegateConfig{ + V31Config: lloconfig.V31Config{ + CaptureOutcomeTelemetry: true, + CaptureReportTelemetry: true, + CaptureAttributedObservationTelemetry: true, + }, + }, + v31: channels{outcome: true, report: true, attributed: true}, + }, + { + name: "each flag independently", + cfg: DelegateConfig{ + CaptureEATelemetry: true, + CaptureOutcomeTelemetry: true, + V31Config: lloconfig.V31Config{CaptureReportTelemetry: true}, + }, + v30: channels{outcome: true}, + v31: channels{report: true}, + }, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + tc.cfg.PluginMonitoringEndpoint = nopTypedEndpoint{} + d := &delegate{cfg: tc.cfg, telem: telem.NewTelemeterService(tc.cfg.telemeterParams(logger.Test(t)))} + + v30 := d.v30FactoryParams(logger.Test(t), nil) + assert.Equal(t, tc.v30.outcome, v30.OutcomeTelemetryCh != nil, "v30 outcome") + assert.Equal(t, tc.v30.report, v30.ReportTelemetryCh != nil, "v30 report") + + v31 := d.v31FactoryParams(logger.Test(t), nil) + assert.Equal(t, tc.v31.outcome, v31.OutcomeTelemetryCh != nil, "v31 outcome") + assert.Equal(t, tc.v31.report, v31.ReportTelemetryCh != nil, "v31 report") + assert.Equal(t, tc.v31.attributed, v31.AttributedObservationTelemetryCh != nil, "v31 attributed observation") + }) + } +} + +type nopTypedEndpoint struct{} + +func (nopTypedEndpoint) SendTypedLog(synchronization.TelemetryType, []byte) {} + // stubKeyValueDatabaseFactory and stubBinaryNetworkEndpoint2Factory stand in // for the OCR3.1-only dependencies; validateInstances only checks that they // are present. From a12307dd1d12d7722521f3084757d5b9de7b4b2c Mon Sep 17 00:00:00 2001 From: Bruno Moura Date: Fri, 9 Oct 2026 08:23:34 +0100 Subject: [PATCH 6/6] llo: hardcode unused LLO observation telemetry to false --- core/services/ocr2/delegate.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/services/ocr2/delegate.go b/core/services/ocr2/delegate.go index 0a448310c6c..78d00db844f 100644 --- a/core/services/ocr2/delegate.go +++ b/core/services/ocr2/delegate.go @@ -1810,7 +1810,7 @@ func (d *Delegate) newServicesLLO( CaptureEATelemetry: jb.OCR2OracleSpec.CaptureEATelemetry, // NOTE: These can be turned off/on, or made configurable in future if // necessary - CaptureObservationTelemetry: jb.OCR2OracleSpec.CaptureEATelemetry, + CaptureObservationTelemetry: false, CaptureOutcomeTelemetry: jb.OCR2OracleSpec.CaptureEATelemetry, CaptureReportTelemetry: false,