diff --git a/README.md b/README.md index fe568d276..b9b834c53 100644 --- a/README.md +++ b/README.md @@ -63,15 +63,15 @@ Choose the command for your server architecture: **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## Product Tour diff --git a/README_DE.md b/README_DE.md index 5ad2fea14..9f4880951 100644 --- a/README_DE.md +++ b/README_DE.md @@ -63,15 +63,15 @@ Wählen Sie den Befehl für Ihre Serverarchitektur: **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## Produkttour diff --git a/README_ES.md b/README_ES.md index 4a8842d54..80d765ed9 100644 --- a/README_ES.md +++ b/README_ES.md @@ -63,15 +63,15 @@ Elige el comando para la arquitectura de tu servidor: **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## Recorrido del producto diff --git a/README_FR.md b/README_FR.md index 305e8881d..43a6b8c59 100644 --- a/README_FR.md +++ b/README_FR.md @@ -63,15 +63,15 @@ Choisissez la commande adaptée à l’architecture de votre serveur : **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## Tour du produit diff --git a/README_JA.md b/README_JA.md index f95ffd068..4bcca4459 100644 --- a/README_JA.md +++ b/README_JA.md @@ -63,15 +63,15 @@ **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## 製品ツアー diff --git a/README_KO.md b/README_KO.md index b1643bdde..06e40a1e2 100644 --- a/README_KO.md +++ b/README_KO.md @@ -63,15 +63,15 @@ **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## 제품 둘러보기 diff --git a/README_PT.md b/README_PT.md index c577f91cf..cdf134175 100644 --- a/README_PT.md +++ b/README_PT.md @@ -63,15 +63,15 @@ Escolha o comando para a arquitetura do seu servidor: **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## Tour do produto diff --git a/README_RU.md b/README_RU.md index 4ba1896b2..069d3612f 100644 --- a/README_RU.md +++ b/README_RU.md @@ -63,15 +63,15 @@ **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## Обзор продукта diff --git a/README_ZH.md b/README_ZH.md index 781ded731..d627da9d8 100644 --- a/README_ZH.md +++ b/README_ZH.md @@ -63,15 +63,15 @@ **AMD64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-amd64.tar.xz -tar -xf ongrid-v0.17.5-linux-amd64.tar.xz && cd ongrid-v0.17.5-linux-amd64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-amd64.tar.xz +tar -xf ongrid-v0.17.6-linux-amd64.tar.xz && cd ongrid-v0.17.6-linux-amd64 sudo ./install.sh ``` **ARM64** ```bash -wget https://github.com/ongridio/ongrid/releases/download/v0.17.5/ongrid-v0.17.5-linux-arm64.tar.xz -tar -xf ongrid-v0.17.5-linux-arm64.tar.xz && cd ongrid-v0.17.5-linux-arm64 +wget https://github.com/ongridio/ongrid/releases/download/v0.17.6/ongrid-v0.17.6-linux-arm64.tar.xz +tar -xf ongrid-v0.17.6-linux-arm64.tar.xz && cd ongrid-v0.17.6-linux-arm64 sudo ./install.sh ``` @@ -79,10 +79,10 @@ sudo ./install.sh ```bash # AMD64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-amd64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-amd64.tar.xz # ARM64 -wget https://ongrid.cloud/dl/ongrid-v0.17.5-linux-arm64.tar.xz +wget https://ongrid.cloud/dl/ongrid-v0.17.6-linux-arm64.tar.xz ``` ## 产品导览 diff --git a/VERSION b/VERSION index 7a81a8f5b..ef09e0004 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -v0.17.5 +v0.17.6 diff --git a/api/tunnel/v1/tunnel.proto b/api/tunnel/v1/tunnel.proto index e8a8f8bbc..8ee6f5508 100644 --- a/api/tunnel/v1/tunnel.proto +++ b/api/tunnel/v1/tunnel.proto @@ -1,5 +1,14 @@ // These messages are transported over a geminio tunnel (not gRPC). -// Encoded as JSON on the wire in MVP (see ADR-001); may switch to proto binary in Phase 2. +// Encoded as JSON on the wire (see ADR-001). After register_edge advertises +// metrics_compression="snappy", push_host_metrics and push_prom_samples may use +// a six-byte prefix (00 4f 47 4d 53 01) followed by Snappy block-compressed JSON. +// RPC names and response formats stay unchanged. Decoded compressed requests +// are limited to 16 MiB. Batches over 6 MiB of JSON are split before compression, +// leaving room for Geminio's base64 framing within its 10 MiB packet limit. +// Item order and metadata are preserved. An indivisible item or metadata may +// exceed the batch budget only if its encoded payload still fits the tunnel; +// otherwise the caller receives an error, without truncation or acknowledgement. +// Smaller than 1 KiB and non-beneficial requests also keep legacy JSON. // // No gRPC service is declared here — cloud and edge register handlers by method name // on geminio End objects; these messages are just the request / response body shapes. @@ -87,6 +96,8 @@ message RegisterEdgeResponse { uint64 edge_id = 1 [json_name = "edge_id"]; uint64 org_id = 2 [json_name = "org_id"]; google.protobuf.Timestamp server_time = 3 [json_name = "server_time"]; + // Empty in older Managers. Only "snappy" enables metric request compression. + string metrics_compression = 4 [json_name = "metrics_compression"]; } // ============================================================ diff --git a/cmd/ongrid-edge/main.go b/cmd/ongrid-edge/main.go index bd845f45f..f5f11cd3c 100644 --- a/cmd/ongrid-edge/main.go +++ b/cmd/ongrid-edge/main.go @@ -900,6 +900,7 @@ func (a collectorAdapter) CollectAll(ctx context.Context) ([]edgebiz.CollectorOu Source: o.Source, HostPoint: o.HostPoint, HostPointValid: o.HostPointValid, + SnapshotID: o.SnapshotID, Samples: o.Samples, }) } diff --git a/docs/rfc/RFC-004-tunnel-metrics-compression.md b/docs/rfc/RFC-004-tunnel-metrics-compression.md new file mode 100644 index 000000000..fece1b2a0 --- /dev/null +++ b/docs/rfc/RFC-004-tunnel-metrics-compression.md @@ -0,0 +1,71 @@ +# RFC-004:Tunnel 指标压缩、拆批与重复上报优化 + +## 元信息 + +- 状态:已实现,本地验证通过,待发布 +- 作者:Codex(按维护者确认的方案实施) +- 日期:2026-09-30 +- 需求来源:维护者确认实现 JSON + Snappy、按字节拆批、新旧版本兼容,并在保留有效数据的前提下减少重复上报。 +- 关联 issue / ADR:无;本轮会话直接授权实施。 + +## 背景和范围 + +Edge 的 `push_prom_samples` 和 `push_host_metrics` 原来通过 Frontier 发送未压缩 JSON。指标名、标签和时间戳重复,增加网络流量。公共 Tunnel 客户端覆盖普通 metrics、Auto APM、自定义指标、数据库指标、Kubernetes 指标及旧 Collector;在生产者现有分批之后增加字节检查,采集频率和样本范围保持原样。 + +旧 `auto/scrape` Collector 会按上报定时器重复读取同一次抓取的缓存,并重新生成时间戳。普通 metrics 插件的配置也可能包含重复 URL。本轮分别修复这两处可证明的重复,避免按数值相同就省略新观测。当前默认 Collector 模式为 `off`,旧快照优化仅影响显式启用 `auto/scrape` 的设备。 + +## 方案 + +1. Manager 在 `register_edge` 响应中新增可选 `metrics_compression: "snappy"`。旧 Manager 不返回该字段,旧 Edge 忽略该字段。Edge 仅在成功注册并识别此值后启用压缩;连接替换、重新注册时先清除能力,过期连接的响应不能更改新连接能力。 +2. 请求仍使用原 RPC 名。Frontier 按 RPC 名获取服务列表再选 Manager,增加压缩专用 RPC 会改变混合版本集群的路由选择。 +3. JSON 达到 1 KiB、最多 16 MiB 且压缩后确实更小时,发送 `00 4f 47 4d 53 01` 六字节标识 + Snappy block 数据。使用现有 Snappy 依赖,不增加配置。响应仍为原 JSON。 +4. Manager 只在两个指标 Handler 中识别压缩格式。解码前检查压缩体大小和 Snappy 声明的原始大小,均不超过 16 MiB,避免恶意长度触发巨大分配;随后复用原身份检查和写入链路。 +5. 16 MiB 是压缩格式的解码上限,**不是采集总量上限或传输批大小**。Geminio 的实际帧上限为 10 MiB,且会将请求体 Base64 编码。因此客户端按 **6 MiB JSON 字节预算**顺序拆批,再压缩;大小包含公共字段、分隔符和括号。使用 `json.RawMessage` 保留数值精度、未知公共字段和样本编码。单条样本或公共字段不可拆时保留完整内容,再检查编码后的请求体能否装入 Tunnel(预留 1 KiB 协议字段开销,计入 Base64 膨胀);无法装入则明确报错并保留已确认前缀,不截断或虚假确认。回退旧 JSON 时也必须通过此检查。 +6. Manager 回滚且 Edge 连接未断时,旧 JSON 解码器会因首字节 NUL 明确拒绝请求。仅匹配原 Handler 的精确 `invalid character '\x00' looking for beginning of value` 解码错误时,关闭当前连接的压缩并用原 JSON 重发一次。超时、断网、写入错误均原样返回,不进行这种重发。后续成功重新注册可以再次开启压缩。 +7. 拆批串行发送,每批接受数必须等于发送数才能继续。取消、错误或不完整确认会停止后续发送;客户端响应中汇总已完整确认的前缀,即使同时返回错误,也可让旧 Collector 只补发未确认部分。不根据含义不明确的部分接受数跳过样本。 +8. 旧 Collector 每次成功抓取生成不可变快照 ID,展平样本和主机 counter/rate 映射只执行一次,时间戳来自抓取时刻。Agent 分别记录主机摘要和详细指标的确认进度;已确认快照不重复上报,失败部分仍可在下一次定时触发时补发。新快照即使数值相同也会发送。缓存仍返回完整结果,使 CompositeCollector 正确识别已有主机数据源。 +9. 普通 metrics 插件对完全相同的 URL 做稳定去重,保留顺序和不同查询参数。不同目标、标签、时间点之间不做基于值的去重。 +10. Manager 在注册未完成、设备映射缺失或 Kubernetes 集群映射未就绪时返回 `accepted: 0`,避免将未进入写入链路的数据误记为成功快照。Prometheus 显式禁用时保留原有接受并丢弃的行为。 + +## 备选方案 + +- Protobuf + Snappy:会改变全部指标消息编码及映射,本轮保留现有 JSON 语义。 +- gzip / 新增 zstd:前者也可行,但本项目已使用 Snappy;后者没有必要新增依赖。 +- 压缩专用 RPC:会改变 Frontier 的服务选择,混合版本部署可能命中另一 Manager,故保留原 RPC 名。 + +## 影响和回滚 + +修改 Edge Tunnel、旧 Collector 及 Agent 上报确认、普通 metrics URL 解析、Manager 两个接收 Handler 和注册响应。无需修改 Frontier、端口、Helm 或用户配置。检查确认已有 remote_write 使用 Snappy,Trace / Log 使用 gzip。可分别回滚 Edge 或 Manager;无数据库迁移。压缩需要双方支持,连接旧 Manager 时自动使用 JSON。 + +本轮保证压缩和拆批不主动截断、抽样或修改有效样本;不等同于新增持久化可靠投递。旧 Collector 仍只保留最新快照,网络故障期间新快照可覆盖旧快照;进程重启会清除内存确认状态,超时也无法判断服务端是否已接收。确认数沿用现有接收链路语义,并非最终存储的持久化凭证。旧 Manager 在映射缺失时的原有虚假接受响应无法由新 Edge 单独修正。需要故障期间零丢失时,应另行设计持久化队列及服务端幂等确认。 + +压缩增加 CPU 和瞬时缓冲分配,不能等同于总资源开销降低。以含重复标签的 1,000 条合成指标做编码/解码基准,记录体积与耗时;结果不代表生产集群的实际压缩率。 + +## 实施与验收 + +实施窗口:2026-09-30,协议及实现 → 兼容性回归 → 基准和自查。 + +- [x] 验证样本名称、标签、数值、时间戳、顺序、未知公共字段和成功接受总数完全不变。 +- [x] 验证超过 6 MiB 的请求先拆批再压缩,并用 Geminio 实际帧编码检查 Base64 后仍小于 10 MiB;不可拆数据能够装入时保留完整内容,无法装入时返回错误。 +- [x] 验证未协商、小包、无压缩收益使用原 JSON。 +- [x] 以旧协议请求/解码器验证新旧兼容、Manager 回滚和连接代次切换;未部署混合版本生产集群。 +- [x] 验证畸形压缩体、超限声明和身份不匹配不能进入写入。 +- [x] 验证超时/写入错误不会触发压缩回退重发,失败拆批保留已确认进度,取消后不继续发送。 +- [x] 验证重复读取旧快照不修改时间戳和 counter/rate;新抓取、原始 exporter 时间戳及 Composite 主机源选择正常。 +- [x] 验证 host / samples 独立确认、失败只补发未确认部分、Manager 映射恢复后可接受同一数据。 +- [x] 验证 URL 去重保留不同查询参数。 +- [x] 相关包通过 `go test -race`、`go vet`,记录基准;未做生产规模压测。 + +### 本地结果 + +- 本轮 Linux arm64(隔离 Go 1.25 容器):Tunnel、frontierbound、Agent biz、collector、metrics、custommetrics、databasemetrics、autoapm、Kubernetes 九个包 `go test -race` 通过;修改路径 `go vet`、Manager / Edge 编译通过。 +- 本轮 macOS arm64:6 MiB 拆批预算、16 MiB 解压上限、不可拆样本/公共字段、顺序和数值精度、取消、旧 Manager、回滚及部分失败测试通过;`make proto` 通过。 +- Review 复现:约 9 MiB 的原始 JSON 在旧 Manager 模式下未拆批,Geminio 实际编码帧达 12,582,363 字节,超过 10 MiB。修正后同一请求拆成 2 批,旧 JSON 请求体合计 9,436,686 字节,Snappy 请求体合计 442,915 字节,各帧均符合传输上限,所有样本字段完全一致。此用例刻意使用高度重复标签验证大包,不代表生产压缩率。 +- 旧快照连续读取三次仅发送一次;新快照即使数值不变仍会发送;主机摘要成功、详细指标部分失败时,只补发未确认数据。 +- 前一轮压缩基础实现验证:Edge 的 Linux arm64 / amd64 无 CGO 编译、解码器 5 秒模糊测试(约 223,295 次)通过。Manager 需要 CGO,因此不使用关闭 CGO 的交叉编译作为验证结果。 +- 前一轮压缩基准(Apple M3 Pro、Go 1.25.11、1,000 条合成指标):285,187 → 24,986 字节,减少 91.2%。编码 JSON 约 0.689 ms/批,JSON + Snappy 约 0.797 ms/批;解码分别约 2.027 / 2.078 ms/批。压缩路径增加缓冲分配,实际数据和负载下需另测。 + +## 变更记录 + +- 2026-09-30:实现兼容协商与 Snappy 压缩;随后补齐实际字节拆批、成功快照去重、URL 去重及未写入数据的接收确认修正。 +- 2026-09-30 Review:修正混淆 16 MiB 解压上限和 10 MiB Geminio 帧上限的问题;拆批预算改为 6 MiB,覆盖 Base64 膨胀及旧协议回退的大小校验。 diff --git a/internal/edgeagent/biz/agent.go b/internal/edgeagent/biz/agent.go index 36183c85b..f84895538 100644 --- a/internal/edgeagent/biz/agent.go +++ b/internal/edgeagent/biz/agent.go @@ -51,6 +51,7 @@ type NetworkDiscoveryCollector interface { // Samples is the open-set rich path consumed by the new // push_prom_samples wire method. type CollectorOutput struct { + SnapshotID uint64 // nonzero identifies one immutable cached scrape Source string HostPoint tunnel.HostMetricPoint HostPointValid bool @@ -610,13 +611,13 @@ var errTunnelStuck = errors.New("tunnel stuck") // push_prom_samples path. One push per source — multi-target scrape // produces one push_host_metrics + one push_prom_samples per target. // -// On either push failure, the corresponding output is dropped and the -// next tick retries with fresh data; we deliberately do not buffer -// open-set samples on the edge because Prometheus remote_write expects -// timely delivery and stale samples are useless. +// Cached scrapes retain their acquisition ID and timestamp. Only confirmed +// delivery is deduplicated; failed paths remain eligible on the next tick. +// The collector still retains its latest snapshot, not a persistent queue. func (a *Agent) metricsLoop(ctx context.Context) error { t := time.NewTicker(a.cfg.MetricsInterval) defer t.Stop() + delivered := make(map[string]snapshotDelivery) for { select { @@ -629,16 +630,38 @@ func (a *Agent) metricsLoop(ctx context.Context) error { // CollectAll may still return a partial slice on error. } for _, out := range outs { - a.pushOne(ctx, out) + a.pushSnapshot(ctx, out, delivered) } } } } +type snapshotDelivery struct { + id uint64 + host bool + samples int +} + +func (a *Agent) pushSnapshot(ctx context.Context, out CollectorOutput, delivered map[string]snapshotDelivery) { + if out.SnapshotID == 0 { + a.pushOne(ctx, out) + return + } + state := delivered[out.Source] + if state.id != out.SnapshotID { + state = snapshotDelivery{id: out.SnapshotID} + } + out.HostPointValid = out.HostPointValid && !state.host + out.Samples = out.Samples[state.samples:] + host, samples := a.pushOne(ctx, out) + state.host = state.host || host + state.samples += samples + delivered[out.Source] = state +} + // pushOne emits one CollectorOutput's two halves (HostPoint and Samples) -// to cloud. Errors are logged but never propagate — the next tick is -// the only retry strategy here. -func (a *Agent) pushOne(ctx context.Context, out CollectorOutput) { +// to cloud, returning confirmed delivery independently for the two paths. +func (a *Agent) pushOne(ctx context.Context, out CollectorOutput) (hostSent bool, samplesSent int) { // 1) legacy fast path: push_host_metrics with one point, but only // for the selected host source. Component scrape targets should not // populate dashboard/alert fast-path rows. @@ -656,7 +679,10 @@ func (a *Agent) pushOne(ctx context.Context, out CollectorOutput) { slog.String("source", out.Source), slog.Any("err", err), ) + } else if resp1.Accepted != 1 { + a.log.Warn("agent: host metric not accepted", slog.String("source", out.Source), slog.Uint64("accepted", uint64(resp1.Accepted))) } else { + hostSent = true a.log.Debug("agent: pushed host metrics", slog.String("source", out.Source), slog.Int("accepted", int(resp1.Accepted)), @@ -666,7 +692,7 @@ func (a *Agent) pushOne(ctx context.Context, out CollectorOutput) { // 2) open-set rich path: push_prom_samples if len(out.Samples) == 0 { - return + return hostSent, 0 } rctx2, cancel2 := context.WithTimeout(ctx, 15*time.Second) var resp2 tunnel.PushPromSamplesResponse @@ -677,19 +703,28 @@ func (a *Agent) pushOne(ctx context.Context, out CollectorOutput) { Samples: out.Samples, }, &resp2) cancel2() + // Split requests expose only a fully acknowledged prefix even on error. + if resp2.Accepted >= 0 && resp2.Accepted <= len(out.Samples) { + samplesSent = resp2.Accepted + } if err != nil { a.log.Warn("agent: push_prom_samples failed", slog.String("source", out.Source), slog.Int("samples", len(out.Samples)), slog.Any("err", err), ) - return + return hostSent, samplesSent + } + if samplesSent != len(out.Samples) { + a.log.Warn("agent: prom samples not fully accepted", slog.String("source", out.Source), slog.Int("accepted", resp2.Accepted), slog.Int("samples", len(out.Samples))) + return hostSent, 0 // A count alone cannot identify an arbitrary partial subset. } a.log.Debug("agent: pushed prom samples", slog.String("source", out.Source), slog.Int("samples", len(out.Samples)), slog.Int("accepted", resp2.Accepted), ) + return hostSent, samplesSent } // noopCollector is used when the Phase 1 New() constructor is still in diff --git a/internal/edgeagent/biz/agent_metrics_test.go b/internal/edgeagent/biz/agent_metrics_test.go new file mode 100644 index 000000000..850f3488c --- /dev/null +++ b/internal/edgeagent/biz/agent_metrics_test.go @@ -0,0 +1,104 @@ +package biz + +import ( + "context" + "errors" + "io" + "log/slog" + "reflect" + "testing" + + "github.com/ongridio/ongrid/internal/pkg/tunnel" +) + +type snapshotPushClient struct { + tunnel.Client + hostCalls, sampleCalls int + failHost bool + partialFailure bool + rejectSamples bool + sent [][]tunnel.PromSample +} + +func (c *snapshotPushClient) Call(_ context.Context, method string, req, resp any) error { + if method == tunnel.MethodPushHostMetrics { + c.hostCalls++ + if c.failHost { + return errors.New("host unavailable") + } + resp.(*tunnel.PushHostMetricsResponse).Accepted = 1 + return nil + } + c.sampleCalls++ + samples := req.(tunnel.PushPromSamplesRequest).Samples + if c.rejectSamples { + c.rejectSamples = false + resp.(*tunnel.PushPromSamplesResponse).Accepted = 0 + return nil + } + if c.partialFailure { + c.partialFailure = false + c.sent = append(c.sent, samples[:1]) + resp.(*tunnel.PushPromSamplesResponse).Accepted = 1 + return errors.New("second chunk unavailable") + } + c.sent = append(c.sent, samples) + resp.(*tunnel.PushPromSamplesResponse).Accepted = len(samples) + return nil +} + +func TestSnapshotDeliveryDeduplicatesOnlyConfirmedAcquisitions(t *testing.T) { + c := &snapshotPushClient{} + a := &Agent{client: c, edgeID: 42, log: slog.New(slog.NewTextHandler(io.Discard, nil))} + delivered := map[string]snapshotDelivery{} + out := CollectorOutput{Source: "scrape:host", SnapshotID: 1, HostPointValid: true, + Samples: []tunnel.PromSample{{Name: "constant", Value: 7, TsMs: 100}, {Name: "counter", Value: 10, TsMs: 100}}} + for range 3 { + a.pushSnapshot(context.Background(), out, delivered) + } + if c.hostCalls != 1 || c.sampleCalls != 1 { + t.Fatalf("same snapshot repeated: host=%d samples=%d", c.hostCalls, c.sampleCalls) + } + out.SnapshotID = 2 // Same values are still valid observations from a new scrape. + a.pushSnapshot(context.Background(), out, delivered) + if c.hostCalls != 2 || c.sampleCalls != 2 { + t.Fatal("new acquisition with unchanged values was suppressed") + } + out.SnapshotID = 0 // Embedded collectors are sampled afresh on each invocation. + for range 2 { + a.pushSnapshot(context.Background(), out, delivered) + } + if c.hostCalls != 4 || c.sampleCalls != 4 { + t.Fatal("uncached collection was suppressed") + } +} + +func TestSnapshotDeliveryRetriesOnlyUnconfirmedParts(t *testing.T) { + for _, mode := range []string{"partial failure", "host and partial failure", "unresolved identity"} { + failHost := mode == "host and partial failure" + rejected := mode == "unresolved identity" + c := &snapshotPushClient{failHost: failHost, partialFailure: !rejected, rejectSamples: rejected} + a := &Agent{client: c, edgeID: 42, log: slog.New(slog.NewTextHandler(io.Discard, nil))} + delivered := map[string]snapshotDelivery{} + out := CollectorOutput{Source: "scrape:host", SnapshotID: 1, HostPointValid: true, + Samples: []tunnel.PromSample{{Name: "first", Value: 1, TsMs: 123}, {Name: "second", Value: 2, TsMs: 123}}} + a.pushSnapshot(context.Background(), out, delivered) + c.failHost = false + a.pushSnapshot(context.Background(), out, delivered) + a.pushSnapshot(context.Background(), out, delivered) + wantHost := 1 + if failHost { + wantHost = 2 + } + if c.hostCalls != wantHost || c.sampleCalls != 2 { + t.Fatalf("host=%d samples=%d", c.hostCalls, c.sampleCalls) + } + want := [][]tunnel.PromSample{out.Samples[:1], out.Samples[1:]} + if rejected { + want = [][]tunnel.PromSample{out.Samples} + } + if !reflect.DeepEqual(c.sent, want) { + t.Fatalf("confirmed samples repeated or missing: %+v", c.sent) + } + } +} diff --git a/internal/edgeagent/collector/scrape.go b/internal/edgeagent/collector/scrape.go index edbfe772a..59f0db7dc 100644 --- a/internal/edgeagent/collector/scrape.go +++ b/internal/edgeagent/collector/scrape.go @@ -27,15 +27,14 @@ import ( ) // Scraper drives one HTTP scrape goroutine per target and stores the -// most recent successful MetricFamily snapshot per target in memory. +// most recent successful flattened snapshot per target in memory. // // Scrape mode is multi-target: one CollectorOutput is produced per // target on each tick, so the agent loop iterates and pushes them // individually with distinct Source values. // -// HostInfo / GetHostLoad / GetProcessList still use gopsutil — the -// scraper itself doesn't try to derive host load from arbitrary -// upstream metric names. +// HostInfo / GetProcessList use gopsutil; GetHostLoad reads the same immutable +// host-role snapshot used for pushes, without changing counter-rate state. type Scraper struct { cfg *ScrapeConfig log *slog.Logger @@ -43,16 +42,19 @@ type Scraper struct { // per-target HTTP client (TLS / bearer auth applied at construction) clients map[string]*http.Client - mu sync.RWMutex - snapshot map[string]targetSnapshot - mappers map[string]*Mapper + mu sync.RWMutex + snapshot map[string]targetSnapshot + mappers map[string]*Mapper + nextSnapshotID uint64 } type targetSnapshot struct { - families []*dto.MetricFamily - at time.Time - source string - role string + id uint64 + samples []tunnel.PromSample + hostPoint tunnel.HostMetricPoint + at time.Time + source string + role string } // NewScraper builds a Scraper from the parsed config. Run must be called @@ -171,14 +173,21 @@ func (s *Scraper) scrapeOnce(ctx context.Context, t ScrapeTarget) { return } mfs := familiesToSlice(families) + at := time.Now() + snap := targetSnapshot{ + at: at, source: SourceScrapePrefix + t.Name, role: t.Role, + samples: FlattenSamples(at, SourceScrapePrefix+t.Name, mfs, t.StaticLabels), + } + if t.Role == ScrapeRoleHost { + // Map counters exactly once per acquisition. Re-reading a cached + // snapshot must neither re-stamp it nor mutate the rate baseline. + snap.hostPoint = s.mappers[t.Name].MapToHostPoint(at, mfs) + } s.mu.Lock() - s.snapshot[t.Name] = targetSnapshot{ - families: mfs, - at: time.Now(), - source: SourceScrapePrefix + t.Name, - role: t.Role, - } + s.nextSnapshotID++ + snap.id = s.nextSnapshotID + s.snapshot[t.Name] = snap s.mu.Unlock() } @@ -186,38 +195,17 @@ func (s *Scraper) scrapeOnce(ctx context.Context, t ScrapeTarget) { // snapshot. Targets that have not yet produced a successful scrape are // skipped silently — the next tick will retry. func (s *Scraper) CollectAll(ctx context.Context) ([]CollectorOutput, error) { - now := time.Now() s.mu.RLock() - snaps := make([]targetSnapshot, 0, len(s.snapshot)) - names := make([]string, 0, len(s.snapshot)) - for name, snap := range s.snapshot { - names = append(names, name) - snaps = append(snaps, snap) - } - s.mu.RUnlock() - - if len(snaps) == 0 { - return nil, nil - } - out := make([]CollectorOutput, 0, len(snaps)) - for i, snap := range snaps { - name := names[i] - target := s.targetByName(name) - var extras map[string]string - if target != nil { - extras = target.StaticLabels - } - mp := s.mappers[name] - if mp == nil { - mp = NewMapper() - s.mappers[name] = mp - } + defer s.mu.RUnlock() + out := make([]CollectorOutput, 0, len(s.snapshot)) + for _, snap := range s.snapshot { co := CollectorOutput{ - Source: snap.source, - Samples: FlattenSamples(now, snap.source, snap.families, extras), + SnapshotID: snap.id, + Source: snap.source, + Samples: snap.samples, } if snap.role == ScrapeRoleHost { - co.HostPoint = mp.MapToHostPoint(now, snap.families) + co.HostPoint = snap.hostPoint co.HostPointValid = true } out = append(out, co) @@ -280,14 +268,10 @@ func (s *Scraper) GetHostLoad(ctx context.Context) (tunnel.GetHostLoadResponse, if snap.role != ScrapeRoleHost { continue } - if len(snap.families) == 0 { + if len(snap.samples) == 0 { continue } - mp := s.mappers[k] - if mp == nil { - continue - } - hp := mp.MapToHostPoint(now, snap.families) + hp := snap.hostPoint if hp.CPUPct == 0 && hp.MemPct == 0 && hp.Load1 == 0 { continue } @@ -297,20 +281,12 @@ func (s *Scraper) GetHostLoad(ctx context.Context) (tunnel.GetHostLoadResponse, resp.Load1 = hp.Load1 resp.Load5 = hp.Load5 resp.Load15 = hp.Load15 + resp.SampledAt = snap.at.Unix() return resp, nil } return resp, nil } -func (s *Scraper) targetByName(name string) *ScrapeTarget { - for i := range s.cfg.Targets { - if s.cfg.Targets[i].Name == name { - return &s.cfg.Targets[i] - } - } - return nil -} - // GetProcessList delegates to gopsutil — scraped targets don't carry // process tables in any standard form. func (s *Scraper) GetProcessList(ctx context.Context, topN int, sortBy string) (tunnel.GetProcessListResponse, error) { diff --git a/internal/edgeagent/collector/scrape_test.go b/internal/edgeagent/collector/scrape_test.go new file mode 100644 index 000000000..9f504bb6e --- /dev/null +++ b/internal/edgeagent/collector/scrape_test.go @@ -0,0 +1,91 @@ +package collector + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "reflect" + "sync/atomic" + "testing" + "time" +) + +func TestScrapeSnapshotsKeepAcquisitionIdentityAndRates(t *testing.T) { + var calls atomic.Int64 + var fail atomic.Bool + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if fail.Load() { + w.WriteHeader(http.StatusServiceUnavailable) + return + } + n := calls.Add(1) + _, err := fmt.Fprintf(w, "# TYPE node_cpu_seconds_total counter\nnode_cpu_seconds_total{cpu=\"0\",mode=\"idle\"} %d\nnode_cpu_seconds_total{cpu=\"0\",mode=\"user\"} %d\n# TYPE node_load1 gauge\nnode_load1 2.5\n# TYPE fixed gauge\nfixed 7 123456789\n", 100+10*n, 200+10*n) + if err != nil { + t.Error(err) + } + })) + defer srv.Close() + target := ScrapeTarget{Name: "host", URL: srv.URL, Role: ScrapeRoleHost, Interval: 30 * time.Second, Timeout: time.Second} + s := NewScraper(&ScrapeConfig{Targets: []ScrapeTarget{target}}, nil) + read := func() CollectorOutput { + t.Helper() + out, err := s.CollectAll(context.Background()) + if err != nil || len(out) != 1 { + t.Fatalf("outputs=%d err=%v", len(out), err) + } + return out[0] + } + s.scrapeOnce(context.Background(), target) + first := read() + if first.SnapshotID == 0 || !first.HostPointValid { + t.Fatal("missing scrape identity") + } + s.scrapeOnce(context.Background(), target) + second := read() + if second.SnapshotID == first.SnapshotID || second.HostPoint.CPUPct != 50 { + t.Fatalf("fresh snapshot/rate incorrect: %+v", second.HostPoint) + } + for range 3 { + load, err := s.GetHostLoad(context.Background()) + if err != nil || load.CPUPct != 50 || load.SampledAt != second.HostPoint.Ts { + t.Fatalf("load=%+v err=%v", load, err) + } + if got := read(); !reflect.DeepEqual(got, second) { + t.Fatal("cached read changed timestamp, identity or counter rate") + } + } + for _, sample := range second.Samples { + if sample.Name == "fixed" && sample.TsMs != 123456789 { + t.Fatalf("explicit timestamp changed: %d", sample.TsMs) + } + if sample.Name == "node_load1" && sample.TsMs != s.snapshot[target.Name].at.UnixMilli() { + t.Fatal("sample was re-stamped after acquisition") + } + } + fail.Store(true) + s.scrapeOnce(context.Background(), target) + if got := read(); !reflect.DeepEqual(got, second) { + t.Fatal("failed scrape made stale data look fresh") + } + embedded := &countingEmbedded{} + composite := NewComposite(embedded, s, nil) + for range 3 { + if _, err := composite.CollectAll(context.Background()); err != nil { + t.Fatal(err) + } + } + if embedded.calls != 0 { + t.Fatal("cached host source incorrectly activated embedded fallback") + } +} + +type countingEmbedded struct { + Collector + calls int +} + +func (c *countingEmbedded) CollectAll(context.Context) ([]CollectorOutput, error) { + c.calls++ + return nil, nil +} diff --git a/internal/edgeagent/collector/types.go b/internal/edgeagent/collector/types.go index 2c5a67ee1..5a874a1d9 100644 --- a/internal/edgeagent/collector/types.go +++ b/internal/edgeagent/collector/types.go @@ -29,6 +29,9 @@ const ( // - Samples is the flat open-set rich path consumed by the new // push_prom_samples wire method. type CollectorOutput struct { + // SnapshotID is nonzero only for cached scrapes. A new acquisition gets + // a new ID even when every sample value is unchanged. + SnapshotID uint64 Source CollectorSource HostPoint tunnel.HostMetricPoint HostPointValid bool diff --git a/internal/edgeagent/plugins/metrics/scrape.go b/internal/edgeagent/plugins/metrics/scrape.go index c2b2b9b17..b37763429 100644 --- a/internal/edgeagent/plugins/metrics/scrape.go +++ b/internal/edgeagent/plugins/metrics/scrape.go @@ -87,6 +87,17 @@ func parseSpec(spec map[string]interface{}) (specView, error) { if v := stringFrom(spec, "application_metrics_url"); v != "" && !slices.Contains(out.URLs, v) { out.URLs = append(out.URLs, v) } + // One plugin applies the same labels/auth to every URL. Exact duplicates + // would scrape and transmit the same target twice on every tick. + seen := make(map[string]bool, len(out.URLs)) + urls := out.URLs[:0] + for _, u := range out.URLs { + if !seen[u] { + seen[u] = true + urls = append(urls, u) + } + } + out.URLs = urls if v := stringFrom(spec, "scrape_interval"); v != "" { d, err := time.ParseDuration(v) if err != nil { diff --git a/internal/edgeagent/plugins/metrics/scrape_test.go b/internal/edgeagent/plugins/metrics/scrape_test.go index 5c86aaaf3..4f593d087 100644 --- a/internal/edgeagent/plugins/metrics/scrape_test.go +++ b/internal/edgeagent/plugins/metrics/scrape_test.go @@ -470,3 +470,12 @@ func TestApplicationMetricsScrapeTarget(t *testing.T) { t.Fatalf("duplicated custom target %+v %v", custom, err) } } + +func TestDuplicateMetricsURLsAreScrapedOnce(t *testing.T) { + first := "http://127.0.0.1:9464/metrics?scope=first" + second := "http://127.0.0.1:9464/metrics?scope=second" + spec, err := parseSpec(map[string]any{"target_urls": []string{first, second, first, second}, "application_metrics_url": first}) + if err != nil || len(spec.URLs) != 2 || spec.URLs[0] != first || spec.URLs[1] != second { + t.Fatalf("targets=%v err=%v", spec.URLs, err) + } +} diff --git a/internal/manager/service/frontierbound/handlers.go b/internal/manager/service/frontierbound/handlers.go index 3e66ebdb4..565369753 100644 --- a/internal/manager/service/frontierbound/handlers.go +++ b/internal/manager/service/frontierbound/handlers.go @@ -301,8 +301,9 @@ func Install(ctx context.Context, c *Client, w Wiring) error { } c.bindEdgeTransport(edgeID, canonicalEdgeID) out := tunnel.RegisterEdgeResponse{ - EdgeID: canonicalEdgeID, - ServerTime: time.Now().UTC().Unix(), + EdgeID: canonicalEdgeID, + ServerTime: time.Now().UTC().Unix(), + MetricsCompression: tunnel.MetricsCompressionSnappy, } return json.Marshal(out) }); err != nil { @@ -431,7 +432,7 @@ func Install(ctx context.Context, c *Client, w Wiring) error { // push_host_metrics: forward batches to the ingester. if err := c.Register(ctx, tunnel.MethodPushHostMetrics, func(rpcCtx context.Context, edgeID uint64, body []byte) ([]byte, error) { var in tunnel.PushHostMetricsRequest - if err := json.Unmarshal(body, &in); err != nil { + if err := tunnel.DecodeMetricsRequest(body, &in); err != nil { return nil, fmt.Errorf("push_host_metrics: decode: %w", err) } canonicalEdgeID := c.canonicalizeEdgeID(edgeID) @@ -477,7 +478,7 @@ func Install(ctx context.Context, c *Client, w Wiring) error { // silently — the edge has no business knowing the cloud's Prom state. if err := c.Register(ctx, tunnel.MethodPushPromSamples, func(rpcCtx context.Context, edgeID uint64, body []byte) ([]byte, error) { var in tunnel.PushPromSamplesRequest - if err := json.Unmarshal(body, &in); err != nil { + if err := tunnel.DecodeMetricsRequest(body, &in); err != nil { return nil, fmt.Errorf("push_prom_samples: decode: %w", err) } canonicalEdgeID := c.canonicalizeEdgeID(edgeID) @@ -486,10 +487,9 @@ func Install(ctx context.Context, c *Client, w Wiring) error { } n := len(in.Samples) if canonicalEdgeID == 0 { - // Edge hasn't completed register_edge yet. Silent drop to - // avoid leaking the raw transport ID as edge_id label - // (v0.7.39 fix). - return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: n}) + // Registration is incomplete. Do not write a transport ID as an + // edge label or acknowledge data that has not reached ingestion. + return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: 0}) } if w.PromIngester == nil { // Prom disabled / not wired. Quiet drop, return Accepted=n so the @@ -511,7 +511,7 @@ func Install(ctx context.Context, c *Client, w Wiring) error { slog.String("source", in.Source), slog.Int("n", n), ) - return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: n}) + return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: 0}) } if err := w.PromIngester.PushKubernetes(rpcCtx, clusterID, in.Source, in.Samples); err != nil { log.Warn("frontierbound: k8s prom ingest push", @@ -543,16 +543,15 @@ func Install(ctx context.Context, c *Client, w Wiring) error { } return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: n}) } - // Host junction missing — drop rather than pollute the TSDB - // with edge_id-as-device_id (issue #96). Accepted=n so the - // edge does not spin-retry; the link lands on register_edge. + // Host junction missing: leave the snapshot unconfirmed so a + // later collection tick can retry after registration repairs it. log.Warn("frontierbound: push_prom_samples dropped — device_id unresolved (edge_devices host junction missing; edge needs to (re)register)", slog.Uint64("edge_id", canonicalEdgeID), slog.Uint64("transport_edge_id", edgeID), slog.String("source", in.Source), slog.Int("n", n), ) - return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: n}) + return json.Marshal(tunnel.PushPromSamplesResponse{Accepted: 0}) } if err := w.PromIngester.Push(rpcCtx, deviceID, in.Source, in.Samples); err != nil { log.Warn("frontierbound: prom ingest push", diff --git a/internal/manager/service/frontierbound/handlers_test.go b/internal/manager/service/frontierbound/handlers_test.go index d6d7d144d..ed074bb23 100644 --- a/internal/manager/service/frontierbound/handlers_test.go +++ b/internal/manager/service/frontierbound/handlers_test.go @@ -5,10 +5,12 @@ import ( "encoding/json" "errors" "log/slog" + "reflect" "strings" "sync" "testing" + "github.com/golang/snappy" "github.com/singchia/geminio" edgebiz "github.com/ongridio/ongrid/internal/manager/biz/edge" @@ -28,6 +30,7 @@ type fakePromIngester struct { wantK8sErr error pushCnt int pushK8sPushCnt int + gotSamples []tunnel.PromSample } func (f *fakePromIngester) Push(_ context.Context, edgeID uint64, source string, samples []tunnel.PromSample) error { @@ -37,6 +40,7 @@ func (f *fakePromIngester) Push(_ context.Context, edgeID uint64, source string, f.gotEdge = edgeID f.gotSrc = source f.gotN = len(samples) + f.gotSamples = samples return f.wantErr } @@ -52,9 +56,13 @@ func (f *fakePromIngester) PushKubernetes(_ context.Context, clusterID uint64, s // fakeMetricIngester is a minimal stub for the existing MetricIngester // requirement of Install. push_host_metrics tests aren't run here. -type fakeMetricIngester struct{} +type fakeMetricIngester struct { + gotDevice uint64 + gotPoints []tunnel.HostMetricPoint +} -func (f *fakeMetricIngester) Push(_ context.Context, _ uint64, _ []tunnel.HostMetricPoint) error { +func (f *fakeMetricIngester) Push(_ context.Context, deviceID uint64, points []tunnel.HostMetricPoint) error { + f.gotDevice, f.gotPoints = deviceID, points return nil } @@ -459,6 +467,58 @@ func TestInstall_PushPromSamples_HappyPath(t *testing.T) { } } +func TestInstall_CompressedMetricsPreserveIdentityAndData(t *testing.T) { + pi, mi := &fakePromIngester{}, &fakeMetricIngester{} + fs, c, _ := installAndDispatch(t, Wiring{ + PromIngester: pi, MetricIngester: mi, DeviceResolver: &fakeDeviceResolver{id: 81}, + }) + c.bindEdgeTransport(7, 42) + samples := []tunnel.PromSample{{Name: "request_total", Labels: map[string]string{"service_name": "api"}, Value: 17.5, TsMs: 1234567}} + points := []tunnel.HostMetricPoint{{Ts: 1234, CPUPct: 12.5, MemPct: 34.5}} + for _, tc := range []struct { + method string + request any + }{ + {tunnel.MethodPushPromSamples, tunnel.PushPromSamplesRequest{EdgeID: 42, Source: "obi", Samples: samples}}, + {tunnel.MethodPushHostMetrics, tunnel.PushHostMetricsRequest{EdgeID: 42, Points: points}}, + } { + t.Run(tc.method, func(t *testing.T) { + body, err := json.Marshal(tc.request) + if err != nil { + t.Fatal(err) + } + // Encode the documented wire format independently of the Edge encoder. + wire := append([]byte{0, 'O', 'G', 'M', 'S', 1}, snappy.Encode(nil, body)...) + rsp := &fakeResp{} + fs.rpcs[tc.method](context.Background(), &fakeReq{data: wire, clientID: 7}, rsp) + if rsp.err != nil { + t.Fatal(rsp.err) + } + if string(rsp.data) != `{"accepted":1}` { + t.Fatalf("response=%s", rsp.data) + } + }) + } + if pi.gotEdge != 81 || pi.gotSrc != "obi" || !reflect.DeepEqual(pi.gotSamples, samples) { + t.Fatal("Prom metrics identity/data changed") + } + if mi.gotDevice != 81 || !reflect.DeepEqual(mi.gotPoints, points) { + t.Fatal("host metrics identity/data changed") + } + // Compression must not bypass the existing authenticated Edge-ID check. + body, err := json.Marshal(tunnel.PushPromSamplesRequest{EdgeID: 99, Samples: samples}) + if err != nil { + t.Fatal(err) + } + rsp := &fakeResp{} + fs.rpcs[tunnel.MethodPushPromSamples](context.Background(), &fakeReq{ + data: append([]byte{0, 'O', 'G', 'M', 'S', 1}, snappy.Encode(nil, body)...), clientID: 7, + }, rsp) + if rsp.err == nil || pi.pushCnt != 1 { + t.Fatal("mismatched identity reached the ingester") + } +} + func TestInstall_PushPromSamples_DoesNotTrustUnboundBodyEdgeID(t *testing.T) { pi := &fakePromIngester{} _, _, rpc := installAndDispatch(t, Wiring{PromIngester: pi, Log: slog.Default()}) @@ -549,29 +609,45 @@ func TestInstall_PushPromSamples_IngesterError(t *testing.T) { } } -// Issue #96: when the host junction can't be resolved, resolveDeviceID -// returns 0 and the handler MUST drop the batch (never write edge_id as -// the device_id label). The ingester must not be called. -func TestInstall_PushPromSamples_DropsWhenDeviceUnresolved(t *testing.T) { - pi := &fakePromIngester{} - _, c, rpc := installAndDispatch(t, Wiring{ - PromIngester: pi, - DeviceResolver: &fakeDeviceResolver{err: errors.New("no host junction")}, - Log: slog.Default(), - }) - c.bindEdgeTransport(1, 5) - body, _ := json.Marshal(tunnel.PushPromSamplesRequest{ - EdgeID: 5, - Source: "embedded", - Samples: []tunnel.PromSample{{Name: "x", Value: 1, TsMs: 1}}, - }) - rsp := &fakeResp{} - rpc(context.Background(), &fakeReq{data: body, clientID: 1}, rsp) - if rsp.err != nil { - t.Fatalf("drop must not error: %v", rsp.err) - } - if pi.pushCnt != 0 { - t.Fatalf("ingester called %d times, want 0 — must drop, never write edge_id as device_id", pi.pushCnt) +// Missing identity must never write transport IDs as device labels (#96), +// nor falsely confirm a snapshot that the Edge would then deduplicate. +func TestInstall_PushPromSamples_OnlyAcknowledgesResolvedIdentity(t *testing.T) { + for _, missing := range []string{"registration", "device", "cluster"} { + t.Run(missing, func(t *testing.T) { + pi := &fakePromIngester{} + resolver := &fakeDeviceResolver{id: 81} + registry := &fakeK8sRegistry{} + _, c, rpc := installAndDispatch(t, Wiring{PromIngester: pi, DeviceResolver: resolver, K8sRegistry: registry}) + if missing != "registration" { + c.bindEdgeTransport(1, 5) + } + if missing == "device" { + resolver.err = errors.New("no host junction") + } + source := "scrape:host" + if missing == "cluster" { + source = "k8s:kube-state-metrics" + } + body, err := json.Marshal(tunnel.PushPromSamplesRequest{EdgeID: 5, Source: source, + Samples: []tunnel.PromSample{{Name: "x", Value: 1, TsMs: 1}}}) + if err != nil { + t.Fatal(err) + } + for _, accepted := range []int{0, 1} { + rsp := &fakeResp{} + rpc(context.Background(), &fakeReq{data: body, clientID: 1}, rsp) + var out tunnel.PushPromSamplesResponse + if err := json.Unmarshal(rsp.data, &out); err != nil || rsp.err != nil || out.Accepted != accepted { + t.Fatalf("accepted=%d want=%d decode=%v rpc=%v", out.Accepted, accepted, err, rsp.err) + } + if pi.pushCnt+pi.pushK8sPushCnt != accepted { + t.Fatal("acknowledgement does not match ingestion") + } + c.bindEdgeTransport(1, 5) + resolver.err = nil + registry.clusterID = 7 + } + }) } } @@ -650,7 +726,8 @@ func TestInstall_PushPromSamples_K8sSourceBypassesHostDeviceID(t *testing.T) { func TestInstall_PushPromSamples_NilIngesterSilentlyAccepts(t *testing.T) { // Wiring.PromIngester == nil => Prom disabled. - _, _, rpc := installAndDispatch(t, Wiring{PromIngester: nil, Log: slog.Default()}) + _, c, rpc := installAndDispatch(t, Wiring{PromIngester: nil, Log: slog.Default()}) + c.bindEdgeTransport(1, 1) body, _ := json.Marshal(tunnel.PushPromSamplesRequest{ Source: "embedded", diff --git a/internal/pkg/tunnel/client.go b/internal/pkg/tunnel/client.go index 484de4e4c..f3f37f550 100644 --- a/internal/pkg/tunnel/client.go +++ b/internal/pkg/tunnel/client.go @@ -55,12 +55,13 @@ type geminioClient struct { endPtr atomic.Pointer[geminio.End] - connMu sync.Mutex - activeConn net.Conn - activeGeneration uint64 - pendingConn net.Conn - pendingGeneration uint64 - nextGeneration uint64 + connMu sync.Mutex + activeConn net.Conn + activeGeneration uint64 + pendingConn net.Conn + pendingGeneration uint64 + nextGeneration uint64 + metricsCompression bool // guarded by connMu; reset on each connection/registration closeOnce sync.Once closed atomic.Bool @@ -268,31 +269,87 @@ func (c *geminioClient) registerOn(end geminio.End, method string, h Handler) er // errors as-is. RetryEnd owns transport recovery; application errors must // not force a second connection while the current transport is still live. func (c *geminioClient) Call(ctx context.Context, method string, req, resp any) error { - end := c.loadEnd() - if end == nil { - return errors.New("tunnel: not dialed") - } body, err := json.Marshal(req) if err != nil { return fmt.Errorf("marshal %q req: %w", method, err) } - connGeneration := c.connectionGeneration() - rsp, callErr := end.Call(ctx, method, end.NewRequest(body)) + var data []byte + if len(body) > maxMetricsBatchBytes && metricsItemsField(method) != "" { + data, err = c.callMetricsBatches(ctx, method, body) + } else { + data, err = c.callJSON(ctx, method, body) + } + if resp != nil && (err == nil || data != nil) { + if decodeErr := json.Unmarshal(data, resp); decodeErr != nil { + return fmt.Errorf("unmarshal %q resp: %w", method, decodeErr) + } + } + return err +} + +func (c *geminioClient) callJSON(ctx context.Context, method string, body []byte) ([]byte, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + end := c.loadEnd() + if end == nil { + return nil, errors.New("tunnel: not dialed") + } + c.connMu.Lock() + connGeneration := c.activeGeneration + if method == MethodRegisterEdge { + c.metricsCompression = false + } + compress := c.metricsCompression + c.connMu.Unlock() + wireBody := body + if compress { + wireBody = encodeMetricsRequest(method, body) + } + if err := checkMetricsWireSize(method, wireBody); err != nil { + return nil, err + } + rsp, callErr := end.Call(ctx, method, end.NewRequest(wireBody)) + remoteError := callErr == nil + if callErr == nil { + callErr = rsp.Error() + } + if wireBody[0] == 0 && legacyMetricsEncodingError(method, callErr) { + // Manager can roll back without the Edge-to-Frontier connection changing. + // Keep the original RPC name so Frontier retains the same routing key. + c.setMetricsCompression(connGeneration, false) + if err := checkMetricsWireSize(method, body); err != nil { + return nil, err + } + rsp, callErr = end.Call(ctx, method, end.NewRequest(body)) + remoteError = callErr == nil + if callErr == nil { + callErr = rsp.Error() + } + } if callErr != nil { c.recycleBrokenRoute(method, callErr, connGeneration) - return fmt.Errorf("tunnel call %q: %w", method, callErr) - } - if rerr := rsp.Error(); rerr != nil { - c.recycleBrokenRoute(method, rerr, connGeneration) - return fmt.Errorf("tunnel call %q: remote: %w", method, rerr) + if remoteError { + return nil, fmt.Errorf("tunnel call %q: remote: %w", method, callErr) + } + return nil, fmt.Errorf("tunnel call %q: %w", method, callErr) } - if resp == nil { - return nil + if method == MethodRegisterEdge { + var registration RegisterEdgeResponse + if err := json.Unmarshal(rsp.Data(), ®istration); err != nil { + return nil, fmt.Errorf("unmarshal %q resp: %w", method, err) + } + c.setMetricsCompression(connGeneration, registration.MetricsCompression == MetricsCompressionSnappy) } - if err := json.Unmarshal(rsp.Data(), resp); err != nil { - return fmt.Errorf("unmarshal %q resp: %w", method, err) + return rsp.Data(), nil +} + +func (c *geminioClient) setMetricsCompression(generation uint64, enabled bool) { + c.connMu.Lock() + defer c.connMu.Unlock() + if generation == c.activeGeneration { + c.metricsCompression = enabled } - return nil } func (c *geminioClient) trackConnection(conn net.Conn) { @@ -334,6 +391,7 @@ func (c *geminioClient) promotePendingConnection() { } c.activeConn = conn c.activeGeneration = c.pendingGeneration + c.metricsCompression = false c.pendingConn = nil c.pendingGeneration = 0 c.connMu.Unlock() diff --git a/internal/pkg/tunnel/messages.go b/internal/pkg/tunnel/messages.go index 68cd1cd46..e4f4b1823 100644 --- a/internal/pkg/tunnel/messages.go +++ b/internal/pkg/tunnel/messages.go @@ -330,8 +330,9 @@ type RegisterEdgeRequest struct { // RegisterEdgeResponse is what the cloud answers on successful register. type RegisterEdgeResponse struct { - EdgeID uint64 `json:"edge_id"` - ServerTime int64 `json:"server_time"` // unix seconds UTC + EdgeID uint64 `json:"edge_id"` + ServerTime int64 `json:"server_time"` // unix seconds UTC + MetricsCompression string `json:"metrics_compression,omitempty"` } // --------------------------------------------------------------------- diff --git a/internal/pkg/tunnel/metrics_batch.go b/internal/pkg/tunnel/metrics_batch.go new file mode 100644 index 000000000..5f777875b --- /dev/null +++ b/internal/pkg/tunnel/metrics_batch.go @@ -0,0 +1,110 @@ +package tunnel + +import ( + "context" + "encoding/json" + "fmt" +) + +func metricsItemsField(method string) string { + switch method { + case MethodPushPromSamples: + return "samples" + case MethodPushHostMetrics: + return "points" + default: + return "" + } +} + +// forEachMetricsBatch uses the same exact JSON-size accounting as the K8s +// metrics batcher. RawMessage preserves numbers and unknown metadata fields. +// Only oversized requests take this path; one batch buffer is reused after +// each synchronous send. No sample is skipped to make a batch fit. +func forEachMetricsBatch(ctx context.Context, body []byte, field string, limit int, send func([]byte, int) error) error { + if err := ctx.Err(); err != nil { + return err + } + var metadata map[string]json.RawMessage + if err := json.Unmarshal(body, &metadata); err != nil { + return fmt.Errorf("decode metrics request: %w", err) + } + var items []json.RawMessage + if err := json.Unmarshal(metadata[field], &items); err != nil { + return fmt.Errorf("decode metrics %s: %w", field, err) + } + delete(metadata, field) + header, err := json.Marshal(metadata) + if err != nil { + return fmt.Errorf("encode metrics metadata: %w", err) + } + prefix := header[:len(header)-1] // Reopen the object and append its sample array. + if len(metadata) > 0 { + prefix = append(prefix, ',') + } + prefix = append(prefix, '"') + prefix = append(prefix, field...) + prefix = append(prefix, '"', ':', '[') + baseBytes := len(prefix) + 2 // closing array and object + if baseBytes > limit { + // Metadata cannot be split without changing the protocol. Preserve the + // full request; callJSON still checks whether its wire encoding fits. + return send(body, len(items)) + } + batch := make([]byte, 0, min(limit, len(body))) + batch = append(batch, prefix...) + count := 0 + for _, item := range items { + if err := ctx.Err(); err != nil { + return err + } + separator := 0 + if count > 0 { + separator = 1 + } + if count > 0 && len(batch)+separator+len(item)+2 > limit { + if err := send(append(batch, ']', '}'), count); err != nil { + return err + } + batch = append(batch[:0], prefix...) + count = 0 + } + if count > 0 { + batch = append(batch, ',') + } + // An indivisible item travels alone. callJSON compresses it if possible + // and rejects it if it cannot fit the transport; nothing is cut. + batch = append(batch, item...) + count++ + } + if err := ctx.Err(); err != nil { + return err + } + return send(append(batch, ']', '}'), count) +} + +func (c *geminioClient) callMetricsBatches(ctx context.Context, method string, body []byte) ([]byte, error) { + accepted := 0 + err := forEachMetricsBatch(ctx, body, metricsItemsField(method), maxMetricsBatchBytes, func(batch []byte, count int) error { + data, err := c.callJSON(ctx, method, batch) + if err != nil { + return fmt.Errorf("metrics batch after %d confirmed items: %w", accepted, err) + } + var result PushPromSamplesResponse // Both metric RPCs return an integer accepted count. + if err := json.Unmarshal(data, &result); err != nil { + return fmt.Errorf("decode metrics batch response: %w", err) + } + if result.Accepted != count { + return fmt.Errorf("metrics batch accepted %d of %d items after %d confirmed items", result.Accepted, count, accepted) + } + accepted += count + return nil + }) + // Even on failure, expose only the fully acknowledged prefix. Callers can + // resume from it without replaying successful batches; no automatic retry. + data, encodeErr := json.Marshal(PushPromSamplesResponse{Accepted: accepted}) + if encodeErr != nil { + return nil, fmt.Errorf("encode metrics batch response: %w", encodeErr) + } + return data, err +} diff --git a/internal/pkg/tunnel/metrics_batch_test.go b/internal/pkg/tunnel/metrics_batch_test.go new file mode 100644 index 000000000..7f29ef339 --- /dev/null +++ b/internal/pkg/tunnel/metrics_batch_test.go @@ -0,0 +1,301 @@ +package tunnel + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "reflect" + "strings" + "testing" + + "github.com/golang/snappy" + "github.com/singchia/geminio" + "github.com/singchia/geminio/packet" + "github.com/singchia/geminio/pkg/id" +) + +func assertMetricsPacketFits(t *testing.T, method string, body []byte) { + t.Helper() + factory := packet.NewPacketFactory(id.NewIDCounter(id.Inc)) + pkt := factory.NewRequestPacket([]byte(method), body) + pkt.Data.Custom = make([]byte, 8) // Frontier appends the authenticated Edge ID. + encoded, err := pkt.Encode() + if err != nil { + t.Fatal(err) + } + // Geminio base64-encodes Data.Value within its JSON frame. Test the real + // encoded packet, since a raw JSON body below 10 MiB can still exceed it. + if len(encoded) > packet.MaxDecodablePacketLen { + t.Fatalf("encoded Geminio packet exceeds 10 MiB: %d bytes", len(encoded)) + } +} + +func TestMetricsBatchExactBytesAndMetadata(t *testing.T) { + for _, field := range []string{"samples", "points"} { + t.Run(field, func(t *testing.T) { + item := json.RawMessage(`{"name":"测试","value":1.2345678901234567,"ts_ms":9007199254740993}`) + metadata := map[string]json.RawMessage{ + "edge_id": []byte("18446744073709551615"), "source": []byte(`"测试\\\"source"`), + "future": []byte(`{"preserve":true}`), field: []byte("[" + string(item) + "]"), + } + one, err := json.Marshal(metadata) + if err != nil { + t.Fatal(err) + } + encoded, err := json.Marshal([]json.RawMessage{item, item, item}) + if err != nil { + t.Fatal(err) + } + metadata[field] = encoded + body, err := json.Marshal(metadata) + if err != nil { + t.Fatal(err) + } + count := 0 + err = forEachMetricsBatch(context.Background(), body, field, len(one), func(batch []byte, n int) error { + if len(batch) != len(one) || n != 1 || !json.Valid(batch) { + t.Fatalf("invalid batch: %d bytes, %d items", len(batch), n) + } + var got map[string]json.RawMessage + if err := json.Unmarshal(batch, &got); err != nil { + return err + } + for key, value := range metadata { + if key != field && !bytes.Equal(got[key], value) { + t.Fatalf("metadata %s changed: %s", key, got[key]) + } + } + var items []json.RawMessage + if err := json.Unmarshal(got[field], &items); err != nil { + return err + } + // Compare JSON bytes after normal encoding, preserving large integers. + var original []json.RawMessage + if err := json.Unmarshal(encoded, &original); err != nil { + return err + } + if !bytes.Equal(items[0], original[0]) { + t.Fatal("sample encoding changed") + } + count += n + return nil + }) + if err != nil || count != 3 { + t.Fatalf("count=%d err=%v", count, err) + } + }) + } +} + +func TestMetricsBatchOversizedItemAndCancellation(t *testing.T) { + body := []byte(`{"samples":[{"value":1},{"name":"` + strings.Repeat("x", 100) + `"}]}`) + calls := 0 + var got []json.RawMessage + send := func(batch []byte, count int) error { + calls++ + var request struct { + Samples []json.RawMessage `json:"samples"` + } + if err := json.Unmarshal(batch, &request); err != nil { + return err + } + if count != len(request.Samples) { + t.Fatal("incorrect batch count") + } + got = append(got, request.Samples...) + return nil + } + if err := forEachMetricsBatch(context.Background(), body, "samples", 64, send); err != nil || calls != 2 { + t.Fatalf("oversize item lost: err=%v calls=%d", err, calls) + } + if len(got) != 2 || string(got[1]) != `{"name":"`+strings.Repeat("x", 100)+`"}` { + t.Fatal("oversized item truncated") + } + metadataBody := []byte(`{"source":"` + strings.Repeat("x", 100) + `","samples":[1]}`) + if err := forEachMetricsBatch(context.Background(), metadataBody, "samples", 64, func(batch []byte, count int) error { + if !bytes.Equal(batch, metadataBody) || count != 1 { + t.Fatal("oversized metadata changed") + } + return nil + }); err != nil { + t.Fatal(err) + } + calls = 0 + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := forEachMetricsBatch(ctx, body, "samples", 256, send); !errors.Is(err, context.Canceled) || calls != 0 { + t.Fatalf("cancel err=%v calls=%d", err, calls) + } + body = []byte(`{"samples":[1,2,3]}`) + ctx, cancel = context.WithCancel(context.Background()) + defer cancel() + err := forEachMetricsBatch(ctx, body, "samples", len(`{"samples":[1]}`), func([]byte, int) error { calls++; cancel(); return nil }) + if !errors.Is(err, context.Canceled) || calls != 1 { + t.Fatalf("mid-batch cancel err=%v calls=%d", err, calls) + } +} + +func TestSingleLargeMetricPreservesDataWhenItsFrameFits(t *testing.T) { + request := PushPromSamplesRequest{EdgeID: 42, Source: "large", Samples: []PromSample{ + {Name: "before", Labels: map[string]string{"label": strings.Repeat("a", 1024)}, Value: 1, TsMs: 123}, + {Name: "oversized", Labels: map[string]string{"label": strings.Repeat("b", maxMetricsBatchBytes+1024)}, Value: 2, TsMs: 124}, + {Name: "after", Labels: map[string]string{"label": strings.Repeat("c", 1024)}, Value: 3, TsMs: 125}, + }} + for _, compressed := range []bool{false, true} { + client := NewClient(ClientConfig{}).(*geminioClient) + client.setMetricsCompression(0, compressed) + calls := 0 + var end geminio.End = &metricsTestEnd{call: func(_ context.Context, method string, req geminio.Request) (geminio.Response, error) { + if calls >= len(request.Samples) { + t.Fatal("unexpected retry") + } + body := req.Data() + assertMetricsPacketFits(t, method, body) + if (body[0] == 0) != compressed { + t.Fatal("unexpected encoding") + } + var got PushPromSamplesRequest + if err := DecodeMetricsRequest(body, &got); err != nil { + t.Fatal(err) + } + if got.EdgeID != request.EdgeID || got.Source != request.Source || !reflect.DeepEqual(got.Samples, request.Samples[calls:calls+1]) { + t.Fatal("sample or metadata lost") + } + calls++ + return &metricsTestResponse{body: []byte(`{"accepted":1}`)}, nil + }} + client.endPtr.Store(&end) + var result PushPromSamplesResponse + if err := client.Call(context.Background(), MethodPushPromSamples, &request, &result); err != nil || calls != 3 || result.Accepted != 3 { + t.Fatalf("accepted=%d calls=%d err=%v", result.Accepted, calls, err) + } + } +} + +func TestUnsendableMetricPreservesOnlyConfirmedPrefix(t *testing.T) { + for _, mode := range []string{"legacy", "rollback", "decoder limit"} { + t.Run(mode, func(t *testing.T) { + labelBytes := maxMetricsWireBytes + 1024 + if mode == "decoder limit" { + labelBytes = maxMetricsDecodedBytes + } + request := PushPromSamplesRequest{Samples: []PromSample{ + {Name: "first", Value: 1, TsMs: 123}, + {Name: "oversized", Labels: map[string]string{"large": strings.Repeat("x", labelBytes)}, Value: 2, TsMs: 124}, + {Name: "after", Value: 3, TsMs: 125}, + }} + client := NewClient(ClientConfig{}).(*geminioClient) + client.setMetricsCompression(0, mode != "legacy") + calls, received := 0, 0 + var end geminio.End = &metricsTestEnd{call: func(_ context.Context, method string, req geminio.Request) (geminio.Response, error) { + calls++ + assertMetricsPacketFits(t, method, req.Data()) + var got PushPromSamplesRequest + if err := json.Unmarshal(req.Data(), &got); err != nil { + return &metricsTestResponse{err: fmt.Errorf("%s: decode: %w", method, err)}, nil + } + if len(got.Samples) != 1 || got.Samples[0].Name != "first" { + t.Fatal("oversized or later sample was incorrectly sent") + } + received++ + return &metricsTestResponse{body: []byte(`{"accepted":1}`)}, nil + }} + client.endPtr.Store(&end) + var result PushPromSamplesResponse + err := client.Call(context.Background(), MethodPushPromSamples, request, &result) + wantCalls := 1 + if mode == "rollback" { + wantCalls = 2 // The compressed large item was rejected before decoding. + } + if !errors.Is(err, packet.ErrPacketTooLarge) || result.Accepted != 1 || received != 1 || calls != wantCalls { + t.Fatalf("accepted=%d received=%d calls=%d err=%v", result.Accepted, received, calls, err) + } + }) + } +} + +func TestLargeMetricRequestSplitsBeforeCompression(t *testing.T) { + // The raw request is below the 16 MiB decompression limit, but its legacy + // Geminio frame exceeds 10 MiB. Two samples fit the 6 MiB JSON batch budget. + const labelBytes = (3 << 20) - 256 + request := PushPromSamplesRequest{EdgeID: 42, Source: "large", Samples: []PromSample{ + {Name: "first", Labels: map[string]string{"large": strings.Repeat("a", labelBytes)}, Value: 1, TsMs: 123}, + {Name: "second", Labels: map[string]string{"large": strings.Repeat("b", labelBytes)}, Value: 2, TsMs: 124}, + {Name: "third", Labels: map[string]string{"large": strings.Repeat("c", labelBytes)}, Value: 3, TsMs: 125}, + }} + for _, mode := range []string{"snappy", "old manager", "failure", "partial acknowledgement", "rollback"} { + t.Run(mode, func(t *testing.T) { + client := NewClient(ClientConfig{}).(*geminioClient) + client.setMetricsCompression(0, mode != "old manager") + var calls, received, compressed, wireBytes int + failure := errors.New("receiver unavailable") + var end geminio.End = &metricsTestEnd{call: func(_ context.Context, method string, req geminio.Request) (geminio.Response, error) { + calls++ + if method != MethodPushPromSamples { + t.Fatal("routing key changed") + } + body := req.Data() + assertMetricsPacketFits(t, method, body) + wireBytes += len(body) + if body[0] == 0 { + compressed++ + if mode == "rollback" { + var old PushPromSamplesRequest + err := json.Unmarshal(body, &old) + return &metricsTestResponse{err: fmt.Errorf("%s: decode: %w", method, err)}, nil + } + var err error + body, err = snappy.Decode(nil, body[len(metricsSnappyPrefix):]) + if err != nil { + t.Fatal(err) + } + } + if len(body) > maxMetricsBatchBytes { + t.Fatalf("batch exceeds cap: %d", len(body)) + } + var got PushPromSamplesRequest + if err := json.Unmarshal(body, &got); err != nil { + t.Fatal(err) + } + if got.EdgeID != request.EdgeID || got.Source != request.Source || !reflect.DeepEqual(got.Samples, request.Samples[received:received+len(got.Samples)]) { + t.Fatal("metadata, order or samples changed") + } + if received == 2 && mode == "failure" { + return nil, failure + } + if received == 2 && mode == "partial acknowledgement" { + return &metricsTestResponse{body: []byte(`{"accepted":0}`)}, nil + } + received += len(got.Samples) + return &metricsTestResponse{body: []byte(fmt.Sprintf(`{"accepted":%d}`, len(got.Samples)))}, nil + }} + client.endPtr.Store(&end) + var result PushPromSamplesResponse + err := client.Call(context.Background(), MethodPushPromSamples, &request, &result) + if mode == "failure" || mode == "partial acknowledgement" { + if err == nil || received != 2 || result.Accepted != 2 || calls != 2 { + t.Fatalf("accepted=%d received=%d calls=%d err=%v", result.Accepted, received, calls, err) + } + if mode == "failure" && !errors.Is(err, failure) { + t.Fatal(err) + } + return + } + if err != nil || received != 3 || result.Accepted != 3 { + t.Fatalf("accepted=%d received=%d err=%v", result.Accepted, received, err) + } + if mode == "snappy" && compressed != 2 { + t.Fatalf("compressed batches=%d", compressed) + } + if mode == "old manager" && compressed != 0 { + t.Fatal("old manager received compression") + } + if mode == "rollback" && (compressed != 1 || calls != 3) { + t.Fatalf("rollback compressed=%d calls=%d", compressed, calls) + } + t.Logf("%s: requests=%d wire bytes=%d", mode, calls, wireBytes) + }) + } +} diff --git a/internal/pkg/tunnel/metrics_compression.go b/internal/pkg/tunnel/metrics_compression.go new file mode 100644 index 000000000..54cf6529a --- /dev/null +++ b/internal/pkg/tunnel/metrics_compression.go @@ -0,0 +1,73 @@ +package tunnel + +import ( + "bytes" + "encoding/json" + "fmt" + + "github.com/golang/snappy" + "github.com/singchia/geminio/packet" +) + +const ( + // MetricsCompressionSnappy is the registration capability for compressed metrics. + MetricsCompressionSnappy = "snappy" + // The initial NUL makes this unambiguously invalid legacy JSON. Old + // handlers reject it before ingestion, allowing a safe JSON fallback. + metricsSnappyPrefix = "\x00OGMS\x01" + minMetricsCompressionBytes = 1024 + maxMetricsDecodedBytes = 16 << 20 + // Geminio base64-encodes payloads inside a JSON frame capped at 10 MiB. + maxMetricsBatchBytes = 6 << 20 + // Reserve 1 KiB for RPC metadata, deadlines, and Frontier's Edge-ID tail. + maxMetricsWireBytes = (packet.MaxDecodablePacketLen - 1024) / 4 * 3 +) + +func checkMetricsWireSize(method string, body []byte) error { + if metricsItemsField(method) != "" && len(body) > maxMetricsWireBytes { + return fmt.Errorf("%s: encoded metrics payload %d exceeds tunnel budget %d: %w", method, len(body), maxMetricsWireBytes, packet.ErrPacketTooLarge) + } + return nil +} + +func encodeMetricsRequest(method string, body []byte) []byte { + if (method != MethodPushPromSamples && method != MethodPushHostMetrics) || + len(body) < minMetricsCompressionBytes || len(body) > maxMetricsDecodedBytes { + return body + } + compressed := snappy.Encode(nil, body) + if len(metricsSnappyPrefix)+len(compressed) >= len(body) { + return body + } + return append([]byte(metricsSnappyPrefix), compressed...) +} + +// DecodeMetricsRequest accepts legacy JSON or negotiated Snappy-compressed JSON. +// dst is the concrete request type owned by the caller, as with json.Unmarshal. +func DecodeMetricsRequest(body []byte, dst any) error { + if bytes.HasPrefix(body, []byte(metricsSnappyPrefix)) { + compressed := body[len(metricsSnappyPrefix):] + // Inspect the advertised decoded size before Snappy allocates memory. + if len(compressed) > maxMetricsDecodedBytes { + return fmt.Errorf("compressed metrics request exceeds %d bytes", maxMetricsDecodedBytes) + } + n, err := snappy.DecodedLen(compressed) + if err != nil { + return fmt.Errorf("metrics snappy length: %w", err) + } + if n > maxMetricsDecodedBytes { + return fmt.Errorf("decoded metrics request exceeds %d bytes", maxMetricsDecodedBytes) + } + body, err = snappy.Decode(nil, compressed) + if err != nil { + return fmt.Errorf("metrics snappy decode: %w", err) + } + } + return json.Unmarshal(body, dst) +} + +// Only this exact legacy decoder error proves that the batch was not ingested. +// Timeouts, transport failures and ingestion errors must never trigger a resend. +func legacyMetricsEncodingError(method string, err error) bool { + return err != nil && err.Error() == method+": decode: invalid character '\\x00' looking for beginning of value" +} diff --git a/internal/pkg/tunnel/metrics_compression_test.go b/internal/pkg/tunnel/metrics_compression_test.go new file mode 100644 index 000000000..6ab56833f --- /dev/null +++ b/internal/pkg/tunnel/metrics_compression_test.go @@ -0,0 +1,348 @@ +package tunnel + +import ( + "bytes" + "context" + "encoding/binary" + "encoding/json" + "errors" + "fmt" + "io" + "math/rand" + "reflect" + "testing" + + "github.com/golang/snappy" + "github.com/singchia/geminio" + "github.com/singchia/geminio/options" +) + +func metricsTestBatch() PushPromSamplesRequest { + batch := PushPromSamplesRequest{EdgeID: 42, Source: "obi"} + for i := 0; i < 1000; i++ { + batch.Samples = append(batch.Samples, PromSample{ + Name: "http_server_request_duration_seconds_bucket", + Labels: map[string]string{ + "service_name": fmt.Sprintf("micro-api-%d", i%20), + "k8s_namespace_name": "production", "http_request_method": "GET", + "k8s_pod_name": fmt.Sprintf("micro-api-%d-7d96c48d9b-%05d", i%20, i/20), + "http_route": fmt.Sprintf("/api/v1/resource/%d", i%7), "le": "0.5", + }, + Value: float64(i*17) / 10, TsMs: 1790726400000, + }) + } + return batch +} + +func TestMetricsCompressionPreservesSamples(t *testing.T) { + want := metricsTestBatch() + raw, err := json.Marshal(want) + if err != nil { + t.Fatal(err) + } + wire := encodeMetricsRequest(MethodPushPromSamples, raw) + if wire[0] != 0 || len(wire) >= len(raw) { + t.Fatal("large metric batch was not compressed") + } + for _, body := range [][]byte{raw, wire} { + var got PushPromSamplesRequest + if err := DecodeMetricsRequest(body, &got); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(got, want) { + t.Fatal("sample identity, values or timestamps changed") + } + } + t.Logf("synthetic 1000-sample batch: JSON=%d bytes, Snappy wire=%d bytes, reduction=%.1f%%", len(raw), len(wire), 100*(1-float64(len(wire))/float64(len(raw)))) +} + +func TestMetricsCompressionSkipsUnsuitablePayloads(t *testing.T) { + random := make([]byte, 4096) + if _, err := rand.New(rand.NewSource(1)).Read(random); err != nil { + t.Fatal(err) + } + for _, tc := range []struct { + name, method string + body []byte + }{ + {"small", MethodPushPromSamples, []byte(`{"samples":[]}`)}, + {"incompressible", MethodPushPromSamples, random}, + {"indivisible oversized JSON", MethodPushPromSamples, bytes.Repeat([]byte("x"), maxMetricsDecodedBytes+1)}, + {"heartbeat", MethodHeartbeat, bytes.Repeat([]byte("x"), 4096)}, + } { + t.Run(tc.name, func(t *testing.T) { + if !bytes.Equal(encodeMetricsRequest(tc.method, tc.body), tc.body) { + t.Fatal("payload should retain its original encoding") + } + }) + } +} + +func TestMetricsDecoderRejectsInvalidAndOversizedData(t *testing.T) { + oversized := binary.AppendUvarint(nil, maxMetricsDecodedBytes+1) + for _, tc := range []struct { + name string + body []byte + }{ + {"empty", nil}, + {"invalid JSON", []byte(`{"samples":`)}, + {"unknown version", []byte("\x00OGMS\x02")}, + {"truncated header", []byte(metricsSnappyPrefix)}, + {"truncated snappy", append([]byte(metricsSnappyPrefix), 10)}, + {"decoded too large", append([]byte(metricsSnappyPrefix), oversized...)}, + {"compressed too large", append([]byte(metricsSnappyPrefix), make([]byte, maxMetricsDecodedBytes+1)...)}, + {"compressed invalid JSON", append([]byte(metricsSnappyPrefix), snappy.Encode(nil, []byte("bad JSON"))...)}, + } { + t.Run(tc.name, func(t *testing.T) { + var req PushPromSamplesRequest + if err := DecodeMetricsRequest(tc.body, &req); err == nil { + t.Fatal("invalid request accepted") + } + }) + } +} + +// Only the transport operations used by Call are stubbed; the production +// client's negotiation, encoding, decoding and fallback all execute normally. +type metricsTestEnd struct { + geminio.End + call func(context.Context, string, geminio.Request) (geminio.Response, error) +} + +func (e *metricsTestEnd) NewRequest(body []byte, _ ...*options.NewRequestOptions) geminio.Request { + return &metricsTestRequest{body: body} +} +func (e *metricsTestEnd) Call(ctx context.Context, method string, req geminio.Request, _ ...*options.CallOptions) (geminio.Response, error) { + return e.call(ctx, method, req) +} + +type metricsTestRequest struct { + geminio.Request + body []byte +} + +func (r *metricsTestRequest) Data() []byte { return r.body } + +type metricsTestResponse struct { + geminio.Response + body []byte + err error +} + +func (r *metricsTestResponse) Data() []byte { return r.body } +func (r *metricsTestResponse) Error() error { return r.err } + +func TestCallRejectsEmptyResponse(t *testing.T) { + client := NewClient(ClientConfig{}).(*geminioClient) + var end geminio.End = &metricsTestEnd{call: func(context.Context, string, geminio.Request) (geminio.Response, error) { + return &metricsTestResponse{}, nil + }} + client.endPtr.Store(&end) + var response map[string]any + var syntaxError *json.SyntaxError + if err := client.Call(context.Background(), MethodHeartbeat, struct{}{}, &response); !errors.As(err, &syntaxError) { + t.Fatalf("empty response must retain its decode error: %v", err) + } + if err := client.Call(context.Background(), MethodHeartbeat, struct{}{}, nil); err != nil { + t.Fatal("caller did not request a response:", err) + } +} + +func TestMetricCallsNegotiateAndHandleManagerRollback(t *testing.T) { + for _, advertised := range []string{"", MetricsCompressionSnappy, "unknown"} { + t.Run("advertised="+advertised, func(t *testing.T) { + client := NewClient(ClientConfig{}).(*geminioClient) + want := metricsTestBatch() + var legacy bool + var compressed, accepted int + end := &metricsTestEnd{call: func(_ context.Context, method string, req geminio.Request) (geminio.Response, error) { + if method == MethodRegisterEdge { + body, err := json.Marshal(RegisterEdgeResponse{EdgeID: 42, MetricsCompression: advertised}) + return &metricsTestResponse{body: body}, err + } + if method != MethodPushPromSamples { + t.Fatalf("routing method changed: %s", method) + } + if req.Data()[0] == 0 { + compressed++ + } + var got PushPromSamplesRequest + var err error + if legacy { + err = json.Unmarshal(req.Data(), &got) + } else { + err = DecodeMetricsRequest(req.Data(), &got) + } + if err != nil { + return &metricsTestResponse{err: fmt.Errorf("%s: decode: %w", method, err)}, nil + } + if !reflect.DeepEqual(got, want) { + t.Fatal("batch changed across transport") + } + accepted++ + return &metricsTestResponse{body: []byte(`{"accepted":1000}`)}, nil + }} + var transport geminio.End = end + client.endPtr.Store(&transport) + push := func() { + t.Helper() + var resp PushPromSamplesResponse + if err := client.Call(context.Background(), MethodPushPromSamples, want, &resp); err != nil { + t.Fatal(err) + } + if resp.Accepted != len(want.Samples) { + t.Fatal("acceptance count changed") + } + } + push() // No compression before registration. + if compressed != 0 { + t.Fatal("compressed before negotiation") + } + if err := client.Call(context.Background(), MethodRegisterEdge, RegisterEdgeRequest{}, nil); err != nil { + t.Fatal(err) + } + push() + legacy = true // Rollback without a transport reconnect. + push() + push() // Fallback is remembered, so old Manager pays no repeated retry cost. + wantCompressed := 0 + if advertised == MetricsCompressionSnappy { + wantCompressed = 2 + } + if compressed != wantCompressed || accepted != 4 { + t.Fatalf("compressed=%d accepted=%d", compressed, accepted) + } + // A new successful registration can enable compression again, then + // an old Manager's registration must revoke it even without reconnect. + if err := client.Call(context.Background(), MethodRegisterEdge, RegisterEdgeRequest{}, nil); err != nil { + t.Fatal(err) + } + advertised = "" + if err := client.Call(context.Background(), MethodRegisterEdge, RegisterEdgeRequest{}, nil); err != nil { + t.Fatal(err) + } + push() + if compressed != wantCompressed { + t.Fatal("old registration did not disable compression") + } + }) + } +} + +func TestMetricCallsDoNotRetryAmbiguousFailures(t *testing.T) { + for _, failure := range []error{context.DeadlineExceeded, io.EOF, errors.New("push_prom_samples: backend unavailable"), errors.New("push_prom_samples: decode: invalid JSON")} { + for _, remote := range []bool{false, true} { + t.Run(fmt.Sprintf("%v/remote=%v", failure, remote), func(t *testing.T) { + client := NewClient(ClientConfig{}).(*geminioClient) + client.setMetricsCompression(0, true) + calls := 0 + var end geminio.End = &metricsTestEnd{call: func(_ context.Context, _ string, req geminio.Request) (geminio.Response, error) { + calls++ + if req.Data()[0] != 0 { + t.Fatal("expected compressed request") + } + if remote { + return &metricsTestResponse{err: failure}, nil + } + return nil, failure + }} + client.endPtr.Store(&end) + if err := client.Call(context.Background(), MethodPushPromSamples, metricsTestBatch(), nil); !errors.Is(err, failure) { + t.Fatalf("error=%v", err) + } + if calls != 1 { + t.Fatalf("ambiguous failure retried %d times", calls) + } + }) + } + } +} + +func TestMetricsCompressionResetsOnReconnect(t *testing.T) { + client := NewClient(ClientConfig{}).(*geminioClient) + client.trackConnection(&closeSpyConn{}) + client.promotePendingConnection() + oldGeneration := client.connectionGeneration() + client.setMetricsCompression(oldGeneration, true) + client.trackConnection(&closeSpyConn{}) + client.promotePendingConnection() + client.setMetricsCompression(oldGeneration, true) // A late old registration response. + if client.metricsCompression { + t.Fatal("stale capability survived reconnect") + } + client.setMetricsCompression(client.connectionGeneration(), true) + client.setMetricsCompression(oldGeneration, false) // A late old fallback response. + if !client.metricsCompression { + t.Fatal("old fallback disabled the new connection") + } +} + +func BenchmarkMetricsWire(b *testing.B) { + batch := metricsTestBatch() + raw, err := json.Marshal(batch) + if err != nil { + b.Fatal(err) + } + compressed := encodeMetricsRequest(MethodPushPromSamples, raw) + for _, tc := range []struct { + name string + compressed bool + body []byte + }{ + {"JSON", false, raw}, {"Snappy", true, compressed}, + } { + b.Run("encode/"+tc.name, func(b *testing.B) { + b.ReportAllocs() + for b.Loop() { + body, err := json.Marshal(batch) + if err != nil { + b.Fatal(err) + } + if tc.compressed { + body = encodeMetricsRequest(MethodPushPromSamples, body) + } + if len(body) == 0 { + b.Fatal("empty payload") + } + } + b.ReportMetric(float64(len(tc.body)), "wire-bytes/op") + }) + b.Run("decode/"+tc.name, func(b *testing.B) { + b.ReportAllocs() + for b.Loop() { + var got PushPromSamplesRequest + if err := DecodeMetricsRequest(tc.body, &got); err != nil { + b.Fatal(err) + } + } + }) + } +} + +func FuzzMetricsRequestDecoding(f *testing.F) { + f.Add([]byte(`{"edge_id":42,"samples":[{"name":"up","value":1,"ts_ms":123}]}`)) + f.Add([]byte(`{"samples":[{"name":"up","labels":{},"value":1,"ts_ms":123}]}`)) + f.Add(append([]byte(metricsSnappyPrefix), snappy.Encode(nil, []byte(`{"samples":[]}`))...)) + f.Add(append([]byte(metricsSnappyPrefix), binary.AppendUvarint(nil, maxMetricsDecodedBytes+1)...)) + f.Fuzz(func(t *testing.T, body []byte) { + var want PushPromSamplesRequest + if err := DecodeMetricsRequest(body, &want); err != nil { + return + } + raw, err := json.Marshal(want) + if err != nil { + t.Fatal(err) + } + var got PushPromSamplesRequest + if err := DecodeMetricsRequest(encodeMetricsRequest(MethodPushPromSamples, raw), &got); err != nil { + t.Fatal(err) + } + roundTrip, err := json.Marshal(got) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(raw, roundTrip) { + t.Fatal("round trip changed decoded metrics") + } + }) +} diff --git a/internal/pkg/tunnel/types.go b/internal/pkg/tunnel/types.go index 363949372..84942403c 100644 --- a/internal/pkg/tunnel/types.go +++ b/internal/pkg/tunnel/types.go @@ -86,6 +86,8 @@ type Client interface { // automatic on reconnect. RegisterHandler(method string, h Handler) // Call invokes an RPC on cloud (heartbeat, push_host_metrics, ...). + // Oversized metric requests are split in order. On a later batch failure, + // resp.Accepted contains only the fully acknowledged prefix, even on error. Call(ctx context.Context, method string, req, resp any) error // AcceptStream blocks until the cloud opens a new bidirectional // stream against this edge (frontier OpenStream call). Used by