From 55f8c9a527bca03a5be26e5e1b66ca484b46d068 Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 18 Jul 2026 12:09:27 +0800 Subject: [PATCH] =?UTF-8?q?fix(obs):=20=E5=AF=B9=E9=BD=90=E6=97=A0?= =?UTF-8?q?=E5=85=BC=E5=AE=B9=E5=B1=82=E4=B8=8E=E5=81=A5=E5=BA=B7/UV=20?= =?UTF-8?q?=E6=9D=83=E5=A8=81=E8=AF=AD=E4=B9=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Agent 本地旧观测缓冲直接删除并运行重建;设计文档去掉兼容期表述。 健康当前态以 PG status/message 为准,CH 仅存 status 与连接时序; Zone 曲线标明分桶 UV,顶部为整窗独立访客。 --- docs/changelog/index.md | 3 + docs/design/observability-data-model.md | 110 +++++++++--------- docs/design/observability-design.md | 83 +++++++------ docs/design/observability-transport-model.md | 22 ++-- .../[zoneId]/components/zone-overview.tsx | 20 +++- .../apps/agent/state/observability_buffer.go | 48 +++++++- .../agent/state/observability_buffer_test.go | 92 +++++++++++++++ internal/apps/openflare/agent/helpers.go | 14 +++ 8 files changed, 285 insertions(+), 107 deletions(-) diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 8636e499..7864a72b 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -41,6 +41,9 @@ sidebar: false - 观测存储 M5:新增 `of_node_edge_health` 与 `of_access_log_hourly`;移除已废弃的请求预聚合表与 OpenResty 吞吐观测表,业务趋势优先读访问日志小时汇总。 - 观测清理目标对齐为 `node_edge_health`;边缘健康完整写入 status;提供历史小时汇总回填 SQL。 - 看板与节点 UV 文案改为「24h / 查询窗口独立访客」(整窗 `uniqExact`);小时趋势图仅展示请求/错误,不绘分时 UV。 +- Agent 二进制替换升级时,若本地观测补传缓冲仍是旧格式或损坏,则直接删除该文件并在运行中重建,避免半迁移数据或阻塞心跳。 +- Zone 流量图:顶部独立访客为整窗去重;曲线标明为分桶 UV(桶内去重,不可跨桶相加)。 +- 明确 OpenResty 健康权威:当前 status/message 以节点表为准,ClickHouse 仅存 status 与连接时序(不含 message)。 - Pages 部署包上传体积限制改为系统动态配置,默认仍为 100 MiB,可按环境调整。 - 上传新部署后会按保留策略自动清理超出数量的历史部署包,减少磁盘占用。 - Agent 在同步 Pages 部署时信任控制面已完成的包校验,不再重复限制文件数与展开体积,仅校验下载完整性并安全解压到本地。 diff --git a/docs/design/observability-data-model.md b/docs/design/observability-data-model.md index d97dc0b9..70eb2986 100644 --- a/docs/design/observability-data-model.md +++ b/docs/design/observability-data-model.md @@ -1,6 +1,7 @@ # Agent 上报协议与观测落库数据模型 -你会学到:重构后 Agent 心跳/WS 上报的 **数据结构**、Server **如何解析与写入**、ClickHouse / 关系库 **目标表结构**,以及与旧字段/旧表的兼容关系。 +你会学到:重构后 Agent 心跳/WS 上报的 **数据结构**、Server **如何解析与写入**、ClickHouse / 关系库 **目标表结构**。 +**无协议兼容层**:Agent 以销毁重建或二进制替换升级;旧字段不解析、旧缓冲整文件丢弃。 本设计是 [边缘可观测与业务流量统计重构](./observability-design.md) 的 **协议与存储专章**,实现时以本文字段与 DDL 为准。 @@ -16,7 +17,7 @@ | 一张业务明细表 | 访问日志是 L1 唯一写入路径 | | 聚合在库内/控制面 | 小时汇总由 ClickHouse MV 或查询生成,Agent 不写汇总表 | | 字段不重叠 | `bytes_sent` = 已提供数据;网卡 `network_*` = 宿主机;不再有业务 `openresty_tx` | -| 可演进 | 新字段可选;旧 Agent 缺字段时 Server 填默认值 | +| 可演进 | 新字段可选;缺省数值填 0,不解析已删除的旧协议字段 | --- @@ -79,30 +80,31 @@ | 字段 | 类型 | 必填 | 说明 | | --- | --- | --- | --- | -| `schema_version` | int | 建议 | `2` = 本设计;缺省或 `0/1` 按旧协议兼容解析 | +| `schema_version` | int | 建议 | 固定为 `2`(本设计) | | `node_id` | string | ✅ | 节点 ID | | `name` | string | ✅ | 显示名 | | `ip` | string | ✅ | 上报 IP | | `version` / `ext_version` | string | ✅ | Agent 版本 | | `current_version` | string | | 本地激活配置版本摘要 | | `last_error` | string | | 最近同步/运行错误,可空 | +| `openresty_status` | string | ✅(有 OpenResty 时) | **最新健康态权威字段** → 写 PG 节点表 | +| `openresty_message` | string | | **最新健康说明权威字段** → 写 PG 节点表(**不进 CH**) | | `profile` | object | | 主机概况,变化时上报(可节流) | | `host_metrics` | object | 建议每拍 | L3 资源快照 | -| `edge_health` | object | 建议每拍 | L2 OpenResty 健康 | +| `edge_health` | object | 建议每拍 | L2 连接时序 + 与顶层一致的 status | | `access_logs` | array | | 本拍增量访问明细 | | `buffered` | array | | 离线补传的事实批次(见 §3.6) | | `health_events` | array | | 边缘健康事件 | | `waf_ip_group_checksums` | map | | 差分同步用,非观测湖 | -**协议 v2 删除(不再作为权威,兼容期可忽略):** +**已删除、Server 不再解析的字段(无兼容层):** | 旧字段 | 处置 | | --- | --- | -| `traffic_report` | 忽略,不落业务表 | -| `openresty_observation.openresty_rx_bytes` / `tx` | 忽略 | -| `openresty_status` / `openresty_message`(顶层) | 迁入 `edge_health`;兼容期从旧字段回填 | -| `snapshot` | 重命名为 `host_metrics`;兼容期别名读取 | -| `buffered_observability` | 重命名为 `buffered`;结构见 §3.6 | +| `traffic_report` | 不存在于协议;不落库 | +| `openresty_observation` | 不存在;连接与状态走 `edge_health` | +| `snapshot` | 不存在;仅用 `host_metrics` | +| `buffered_observability` | 不存在;仅用 `buffered` | ### 3.2 `profile` — 主机概况(低频) @@ -174,12 +176,19 @@ | 字段 | 类型 | 语义 | | --- | --- | --- | -| `status` | string | `healthy` / `unhealthy` / `unknown` | -| `message` | string | 状态说明 | +| `status` | string | `healthy` / `unhealthy` / `unknown`(须与顶层 `openresty_status` 一致) | +| `message` | string | 状态说明(上报可带;**仅用于回填 PG 最新态,不进 CH**) | | `connections` | int64 | stub_status Active connections | -`status`/`message` 同步更新关系库节点最新状态;`connections` 写入 CH `of_node_edge_health` 供节点详情曲线(可选)。 +#### 健康状态权威源(收敛) +| 数据 | 权威存储 | 说明 | +| --- | --- | --- | +| **当前** OpenResty 是否健康 + 说明文案 | **PG 节点表** `openresty_status` / `openresty_message` | UI 徽章、列表、告警以这里为准 | +| **时序** 健康 status + 连接数 | **CH** `of_node_edge_health`(`status`, `connections`) | 连接曲线 / 健康状态历史;**无 message 列** | +| Agent 上报 | 顶层 status/message + `edge_health` | Server 归一化后二者 status 对齐;message **只写 PG** | + +因此:查「现在是否 unhealthy」→ 读 PG;查「过去 24h 连接数」→ 读 CH。 ### 3.5 `access_logs[]` — 访问明细(L1,业务唯一事实) Agent:tail access.log → 解析 JSON 行 → 原样字段上报(可截断 path)。 @@ -205,7 +214,7 @@ Agent:tail access.log → 解析 JSON 行 → 原样字段上报(可截断 p | `path` | string | ✅ | `$request_uri`,Agent 可截断 | 路径 | | `status_code` | int | ✅ | `$status` | 状态码 | | `bytes_sent` | int64 | ✅ | **`$body_bytes_sent`** | **已提供数据**(响应体) | -| `request_length` | int64 | 建议 | `$request_length` | **接收数据**;旧 Agent 缺省 0 | +| `request_length` | int64 | 建议 | `$request_length` | **接收数据** | | `request_time_ms` | int64 | 可选 | `$request_time * 1000` | 耗时;缺省 0 | **明确不由 Agent 上报(由 Server 写入):** @@ -256,29 +265,23 @@ Agent:tail access.log → 解析 JSON 行 → 原样字段上报(可截断 p // pkg/protocol/agent.go(目标形态,实现时替换旧类型) type NodePayload struct { - SchemaVersion int `json:"schema_version,omitempty"` - NodeID string `json:"node_id"` - Name string `json:"name"` - IP string `json:"ip"` - Version string `json:"version"` - ExtVersion string `json:"ext_version"` - CurrentVersion string `json:"current_version"` - LastError string `json:"last_error"` - Profile *NodeSystemProfile `json:"profile,omitempty"` - HostMetrics *NodeHostMetrics `json:"host_metrics,omitempty"` - EdgeHealth *NodeEdgeHealth `json:"edge_health,omitempty"` - AccessLogs []NodeAccessLog `json:"access_logs,omitempty"` - Buffered []BufferedFacts `json:"buffered,omitempty"` - HealthEvents []NodeHealthEvent `json:"health_events"` - WAFIPGroupChecksums map[string]string `json:"waf_ip_group_checksums,omitempty"` - - // Deprecated: schema_version < 2 兼容 - Snapshot *NodeHostMetrics `json:"snapshot,omitempty"` - OpenrestyStatus string `json:"openresty_status,omitempty"` - OpenrestyMessage string `json:"openresty_message,omitempty"` - OpenrestyObservation json.RawMessage `json:"openresty_observation,omitempty"` // 仅解析 connections - TrafficReport json.RawMessage `json:"traffic_report,omitempty"` // 忽略 - BufferedObservability []BufferedFacts `json:"buffered_observability,omitempty"` + SchemaVersion int `json:"schema_version,omitempty"` + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + Version string `json:"version"` + ExtVersion string `json:"ext_version"` + CurrentVersion string `json:"current_version"` + LastError string `json:"last_error"` + OpenrestyStatus string `json:"openresty_status"` // PG 最新态权威 + OpenrestyMessage string `json:"openresty_message"` // PG 最新态权威;不进 CH + Profile *NodeSystemProfile `json:"profile,omitempty"` + HostMetrics *NodeHostMetrics `json:"host_metrics,omitempty"` + EdgeHealth *NodeEdgeHealth `json:"edge_health,omitempty"` + AccessLogs []NodeAccessLog `json:"access_logs,omitempty"` + Buffered []BufferedFacts `json:"buffered,omitempty"` + HealthEvents []NodeHealthEvent `json:"health_events"` + WAFIPGroupChecksums map[string]string `json:"waf_ip_group_checksums,omitempty"` } type NodeHostMetrics struct { @@ -562,11 +565,11 @@ SETTINGS index_granularity = 8192; | 列 | 说明 | | --- | --- | -| `status` | 瞬时健康 | +| `status` | 瞬时健康(与 PG 当前态同源;用于时序,非唯一 UI 权威) | | `connections` | 当前连接数 | +**无** `message` 列(说明文案仅 PG 最新态)。 **无** `openresty_rx_bytes` / `openresty_tx_bytes`。 - ### 5.6 关系库(节点最新态,非分析湖) 与观测湖分离,保持「最新一份」: @@ -641,15 +644,15 @@ Agent 解析: --- -## 8. 兼容策略(协议 v1 → v2) +## 8. 升级策略(无兼容层) -| 客户端 | Server 行为 | +| 项 | 策略 | | --- | --- | -| 新 Agent `schema_version=2` | 按本文写入 L1/L2/L3 | -| 旧 Agent 带 `snapshot` + `access_logs` | 映射为 host_metrics;access_logs 无 request_length 则 0 | -| 旧 Agent 带 `traffic_report` | **丢弃** | -| 旧 Agent 带 `openresty_observation` | 只取 `connections` + 顶层 status;rx/tx 丢弃 | -| 读路径 | 业务 API **只读** access_logs(及 hourly);不再读 request_reports / openresty 吞吐 | +| Agent 升级 | **销毁重建**优先;允许**二进制替换** | +| 协议 | 仅 schema v2 字段;旧 JSON 字段不解析 | +| 本地观测缓冲 | 若仍是旧格式(含 `snapshot` / `openresty_observation` / `traffic_report`)或损坏 → **整文件删除**,运行中重建 | +| 读路径 | 业务 API **只读** access_logs(及 hourly);健康当前态读 PG;连接时序读 CH edge_health | +| 旧 Agent | 必须升级;控制面不提供 v1 双读路径 | --- @@ -695,10 +698,11 @@ Agent 解析: **写入:** -1. `of_node_metric_snapshots` 1 行(network_tx=2000 累计) -2. `of_node_edge_health` 1 行(connections=5) -3. `of_node_access_logs` 1 行(bytes_sent=500, request_length=80, region=Server 填充) -4. MV 异步计入 `of_access_log_hourly` +1. PG 节点最新态:`openresty_status` / `openresty_message`(若上报) +2. `of_node_metric_snapshots` 1 行(network_tx=2000 累计) +3. `of_node_edge_health` 1 行(status + connections=5;**无 message**) +4. `of_node_access_logs` 1 行(bytes_sent=500, request_length=80, region=Server 填充) +5. MV 异步计入 `of_access_log_hourly` **查询 24h 已提供数据:** `sum(bytes_sent)` → 至少 500(加历史) **查询宿主机出站:** 对 snapshots 差分,与 500 **无强制相等关系**。 @@ -707,12 +711,12 @@ Agent 解析: ## 10. 实现检查清单 -- [x] `pkg/protocol`:v2 类型与 deprecated 兼容字段 +- [x] `pkg/protocol`:仅 v2 字段,无兼容别名 - [x] Agent:只组 `host_metrics` / `edge_health` / `access_logs` / `buffered` -- [x] Server normalize + 停写 request_reports / openresty rx/tx +- [x] Server:无 request_reports / openresty 吞吐;健康当前态 PG、时序 CH - [x] CH migration:`request_length`、`request_time_ms`、`of_node_edge_health`、`of_access_log_hourly`、hourly 回填 - [x] 看板/Zone API 统一读 access log 聚合 -- [x] 文档与前端文案:已提供数据 ≠ 宿主机网卡出站;UV 策略(整窗 uniqExact / 小时路径 UV=0) +- [x] UV:整窗 uniqExact;Zone 曲线标明分桶 UV;小时趋势不绘 UV --- diff --git a/docs/design/observability-design.md b/docs/design/observability-design.md index db2f60a4..233ff4cf 100644 --- a/docs/design/observability-design.md +++ b/docs/design/observability-design.md @@ -45,7 +45,7 @@ * Agent 保持轻量:解析日志行、读 `/proc`、健康检查;不做业务分析。 * 控制面 API 错误仍走统一信封与 `response.Abort*`。 -* 访问日志字段变更须同时更新 OpenResty `log_format` 与 Agent 解析器,并保证向后兼容至少一个小版本。 +* 访问日志字段变更须同时更新 OpenResty `log_format` 与 Agent 解析器;Agent 与控制面同版本发布,不保留旧协议解析。 --- @@ -228,20 +228,18 @@ flowchart TB | 概念 | 字段 | 展示名 | | --- | --- | --- | -| CPU / 内存 / 磁盘占用 | 现有 snapshot | 保持 | +| CPU / 内存 / 磁盘占用 | `host_metrics` | 保持 | | 网卡累计字节 | `network_rx_bytes` / `network_tx_bytes` | **宿主机网卡入/出站** | | 磁盘 IO 累计 | `disk_read_bytes` / `disk_write_bytes` | 磁盘读/写 | -### 6.2 废弃或降级字段 +### 6.2 已删除字段(无兼容层) -| 现字段 | 处置 | 原因 | +| 原字段 | 处置 | 原因 | | --- | --- | --- | -| `openresty_tx_bytes` | **废弃业务用途**;迁移期可读但 UI 不再展示为业务出站 | 与 `bytes_sent` 重复 | -| `openresty_rx_bytes` | **废弃业务用途**;由 `request_length` 聚合替代 | 与日志重复 | -| `TrafficReport` 全量 | **废弃权威地位**;迁移期可停写或仅兼容旧 Agent | 边缘预聚合 | -| `TrafficReport.top_domains` / `status_codes` / `unique_visitor_count` | 改由 Server 查日志 | 同上 | +| `openresty_tx_bytes` / `openresty_rx_bytes` | **删除** | 业务字节以 access log 为准 | +| `TrafficReport` 及 TopN/窗内 UV | **删除** | 边缘预聚合 | | Agent state 内业务 lifetime 累计 | 删除 | 违背 P1 | -| Lua shared dict 业务吞吐/窗口请求计数 | 删除或仅保留本地诊断 | 非投递主路径 | +| Lua shared dict 业务吞吐/窗口请求计数 | 删除 | 非投递主路径 | ### 6.3 命名对照(前端文案强制) @@ -262,21 +260,22 @@ flowchart TB ```text NodePayload - identity / version / openresty_status + identity / version / openresty_status / openresty_message # 最新态 → PG profile # 主机概况(低频) - snapshot # L3 资源读数(含网卡累计原值) - openresty_connections # L2 瞬时(可挂在精简 observation 或 snapshot 扩展) + host_metrics # L3 资源读数(含网卡累计原值) + edge_health # L2:status + connections(CH 时序;message 不进 CH) access_logs[] # L1 明细(主路径) health_events[] - buffered_observability[] # 缓冲的是上述事实,不是报表 + buffered[] # 缓冲的是上述事实,不是报表 waf_ip_group_checksums ``` -移除或标记 deprecated(兼容窗口内 Server 忽略写入分析权威路径): +协议中已删除(无兼容层): ```text -traffic_report # deprecated -openresty_observation.rx/tx # deprecated(connections 迁出后可删结构) +traffic_report +openresty_observation +snapshot / buffered_observability 别名 ``` ### 7.2 Access log 上报要求 @@ -291,7 +290,7 @@ openresty_observation.rx/tx # deprecated(connections 迁出后可删结构 | `path` | ✅ | 可截断 | | `status_code` | ✅ | | | `bytes_sent` | ✅ | body 字节,已提供数据 | -| `request_length` | ✅(协议补齐) | 接收数据;旧 Agent 可缺省为 0 | +| `request_length` | ✅ | 接收数据 | Agent 职责: @@ -329,10 +328,10 @@ Agent 职责: | 输入 | 表 | 说明 | | --- | --- | --- | -| `access_logs[]` | `of_node_access_logs` | 权威业务明细;补齐 `request_length` 列(若尚无) | -| `snapshot` | `of_node_metric_snapshots` | L3;网卡/磁盘累计 | -| 连接数 / 健康 | 现有节点状态或精简 obs 表 | L2 | -| `traffic_report` / openresty rx/tx | **停止作为权威写入** 或兼容期双写但不读 | 迁移后删除写入 | +| `access_logs[]` | `of_node_access_logs` | 权威业务明细 | +| `host_metrics` | `of_node_metric_snapshots` | L3;网卡/磁盘累计 | +| `openresty_status` / `openresty_message` | **PG 节点表** | L2 **最新态权威**(message 仅此) | +| `edge_health` | `of_node_edge_health` | L2 时序:status + connections(**无 message**) | GeoIP:继续在 Server 入库路径解析 `remote_addr` → `region`,不在 Agent 做。 @@ -407,7 +406,7 @@ of_access_log_hourly } ``` -兼容:旧字段 `bytes_sent` 可在一个版本内作为 `bytes_provided` 的别名返回,文档标注 deprecated。 +API 业务字节字段使用 `bytes_provided` / `bytes_received`(访问日志聚合);不再返回 openresty 吞吐别名。 ### 9.2 看板 @@ -457,28 +456,36 @@ bytes_sent (= $body_bytes_sent), request_length --- -## 11. 兼容与迁移 +## 11. 升级与迁移(无兼容层) -### 11.1 阶段划分 +### 11.1 阶段回顾(已落地) -| 阶段 | 内容 | 结果 | -| --- | --- | --- | -| **M1 读路径切换** | 看板/节点业务趋势改为 access log 聚合;UI 文案改为已提供/接收数据 | 对账立刻成立;旧字段可仍写入 | -| **M2 协议补齐** | AccessLog 上报 `request_length`;Server 入库 | 接收数据可用 | -| **M3 停写预聚合** | Server 忽略/停写 TrafficReport 与 openresty rx/tx 权威路径 | 减负 | -| **M4 Agent 瘦身** | 移除边缘 TrafficReport 构建、Lua 业务计数、lifetime 累计 state | 符合 P1 | -| **M5 清理** | 删除废弃 CH 表/列、API 字段、前端类型 | 无冗余 | +| 阶段 | 内容 | +| --- | --- | +| **M1–M5** | 读路径切 access log;协议 v2;停预聚合;edge_health + access_log_hourly;删旧表与 API 兼容字段 | -### 11.2 兼容策略 +### 11.2 升级策略 -* 旧 Agent 仍发 `TrafficReport`:Server **不用于** 看板业务趋势。 -* 旧 Agent 无 `request_length`:`bytes_received` 为 0 或不展示。 -* 明细缺失时段:业务图为空或仅部分;**不得**回退到 openresty_tx 冒充已提供数据(避免再次双真相)。 +* **Agent:销毁重建优先**;允许二进制替换。 +* 二进制替换时:本地旧观测缓冲(含 `snapshot` / `openresty_observation` / `traffic_report`)**整文件删除**,运行后重建。 +* Server **不**解析 v1 字段,**不**双读 request_reports / openresty 吞吐。 +* 明细缺失时段:业务图为空或仅部分;**不得**用网卡或已删除的 openresty 吞吐冒充已提供数据。 ### 11.3 数据回填 -* 历史「已提供数据」以 access log 为准,无需从 openresty 观测回填。 -* 历史看板 openresty 曲线可保留只读至 TTL,或直接隐藏。 +* 历史「已提供数据」以 access log 为准。 +* `of_access_log_hourly` 创建前历史用 goose 回填 SQL(ANTI JOIN 防重)。 + +### 11.4 健康状态权威 + +* **当前态**:PG `openresty_status` / `openresty_message`。 +* **时序**:CH `of_node_edge_health`(status + connections;无 message)。 + +### 11.5 UV + +* **整窗独立访客**:`uniqExact(remote_addr)`(看板合计、Zone 合计)。 +* **分桶 UV**(Zone 曲线):桶内 uniq,**不可跨桶相加**;UI 须标明。 +* **小时趋势路径**:不绘 / 不填分时 UV(hourly 表不含 UV)。 --- @@ -527,7 +534,7 @@ sum(各 Zone 已提供) + sum(未归属 Host) = 全局已提供 | 明细量大导致 CH 与心跳变重 | 批量、压缩、采样策略评估;Server rollup;限制单次条数 | | 短暂丢失日志导致业务量偏低 | 本地 buffer 与轮转处理;监控 access log 采集滞后 | | 用户仍对比「网卡出站」与「已提供」 | UI 分区与文案强制「宿主机」前缀 | -| 旧 Agent 长期在线 | 兼容期忽略预聚合;文档要求升级 Agent 以获得接收数据 | +| 旧 Agent 长期在线 | **无兼容层**;必须升级/重建 Agent | **为何不保留 Agent 预聚合作为优化?** diff --git a/docs/design/observability-transport-model.md b/docs/design/observability-transport-model.md index a0615c52..03ea1bb1 100644 --- a/docs/design/observability-transport-model.md +++ b/docs/design/observability-transport-model.md @@ -124,11 +124,12 @@ | `health_events` | 事件 | 如 openresty_unhealthy | | `waf_ip_group_checksums` | 同步 | 非观测湖 | -**目标态不再作为权威业务数据(旧字段,兼容期可忽略):** +**协议已删除(无兼容层,旧 Agent 必须升级):** -- `traffic_report`(窗内请求/UV/Top 域名等预聚合) -- `openresty_observation.openresty_rx_bytes` / `openresty_tx_bytes` -- 顶层与「已提供数据」平行的业务「出站」字段 +- `traffic_report` +- `openresty_observation`(含 rx/tx) +- `snapshot` / `buffered_observability` +- 业务含义的 openresty 吞吐字段 --- @@ -293,10 +294,15 @@ Agent GET /openflare/observability | 字段 | 来源 | | --- | --- | -| `status` / `message` | Agent 对 OpenResty 健康探测结果(配置校验/进程等,可与观测口 `ok` 配合) | +| `status` / `message` | Agent 健康探测(配置校验/进程等,可与观测口 `ok` 配合);须与顶层 `openresty_status` / `openresty_message` 对齐 | | `connections` | 观测口 `connections.active` | -落库:关系库节点最新状态 + 可选 CH `of_node_edge_health` 做连接曲线。 +**落库拆分(权威源):** + +| 内容 | 写入 | +| --- | --- | +| 最新 `status` + `message` | **PG 节点表**(UI / 列表 / 告警) | +| 时序 `status` + `connections` | **CH `of_node_edge_health`**(**无 message**) | --- @@ -460,8 +466,9 @@ t=6s 下一轮… | Lua dict 60s 窗 request_count + Agent 10s 拉 + Server sum | **删除**;请求数 = 日志 count | | openresty_tx 当「出站」 | **删除**;已提供数据 = `sum(bytes_sent)` | | 两个口 observability + stub_status | **合并为一个** observability,只返回连接/探活 | -| TrafficReport 预聚合 | **不作为权威**;兼容期可忽略 | +| TrafficReport 预聚合 | **删除**;协议与 API 均无此路径 | | 业务与网卡混称「流量」 | **分文案、分 API、分表** | +| 健康 status/message | **PG 最新态权威**;CH 仅 status+连接时序 | --- @@ -486,3 +493,4 @@ t=6s 下一轮… | 2026-07-18 | 初稿:作为「最新传输模型」单页说明——三层、频率、示例 JSON、采集来源、与旧模型对照 | | 2026-07-18 | 默认上报间隔 3s;离线阈值 60s;补传窗口 60 分钟 | | 2026-07-18 | M5:edge_health 表、access_log_hourly、废弃 request_reports/obs_openresty 吞吐表 | +| 2026-07-18 | 无兼容层:删除「兼容期可忽略」表述;健康 message 仅 PG、CH 无 message | diff --git a/frontend/app/(main)/websites/[zoneId]/components/zone-overview.tsx b/frontend/app/(main)/websites/[zoneId]/components/zone-overview.tsx index 1897aabe..b23d94eb 100644 --- a/frontend/app/(main)/websites/[zoneId]/components/zone-overview.tsx +++ b/frontend/app/(main)/websites/[zoneId]/components/zone-overview.tsx @@ -31,7 +31,7 @@ const rangeOptions: Array<{ value: ZoneStatsRange; label: string }> = [ ]; const visitorsChartConfig = { - value: { label: '唯一访问者', color: 'hsl(217 91% 60%)' }, + value: { label: '分桶 UV', color: 'hsl(217 91% 60%)' }, } satisfies ChartConfig; const requestsChartConfig = { @@ -165,7 +165,8 @@ export function ZoneOverviewPanel({ >; dataKey: string; @@ -242,9 +245,16 @@ function MetricTrendCard({
-
- - {label} +
+
+ + {label} +
+ {description ? ( +

+ {description} +

+ ) : null}

{value}

diff --git a/internal/apps/agent/state/observability_buffer.go b/internal/apps/agent/state/observability_buffer.go index e8e9b4ec..8e2fa5e5 100644 --- a/internal/apps/agent/state/observability_buffer.go +++ b/internal/apps/agent/state/observability_buffer.go @@ -3,10 +3,12 @@ package state import ( "encoding/json" + "log/slog" "os" "path/filepath" "sort" "strconv" + "strings" "sync" "github.com/Rain-kl/Wavelet/internal/apps/agent/protocol" @@ -15,6 +17,8 @@ import ( const observabilityBufferWindowSeconds = 60 // ObservabilityBufferRecord stores observability facts for a single time window. +// Disk JSON is schema-v2 only: host_metrics / edge_health / access_logs. +// Pre-v2 buffers are discarded on load (binary upgrade without data-dir wipe). type ObservabilityBufferRecord struct { WindowStartedAtUnix int64 `json:"window_started_at_unix"` HostMetrics *protocol.NodeMetricSnapshot `json:"host_metrics,omitempty"` @@ -191,10 +195,14 @@ func (s *ObservabilityBufferStore) loadUnlocked() ([]ObservabilityBufferRecord, s.cacheLoaded = true return []ObservabilityBufferRecord{}, nil } - var records []ObservabilityBufferRecord - if err = json.Unmarshal(data, &records); err != nil { - return nil, err + + // Binary upgrade: drop pre-v2 or corrupt buffer entirely; agent rebuilds on subsequent heartbeats. + records, reason, ok := parseObservabilityBufferDisk(data) + if !ok { + s.discardBufferFile(reason) + return []ObservabilityBufferRecord{}, nil } + s.cache = records s.cacheLoaded = true copied := make([]ObservabilityBufferRecord, len(s.cache)) @@ -202,6 +210,38 @@ func (s *ObservabilityBufferStore) loadUnlocked() ([]ObservabilityBufferRecord, return copied, nil } +// parseObservabilityBufferDisk returns v2 records, or ok=false when the on-disk file should be wiped. +func parseObservabilityBufferDisk(data []byte) (records []ObservabilityBufferRecord, reason string, ok bool) { + raw := strings.TrimSpace(string(data)) + if raw == "" { + return []ObservabilityBufferRecord{}, "", true + } + // Valid buffer is a JSON array of window records. + if !strings.HasPrefix(raw, "[") { + return nil, "legacy or unreadable observability buffer", false + } + // Pre-v2 keys: discard whole file (no field migration). + if strings.Contains(raw, `"snapshot"`) || + strings.Contains(raw, `"openresty_observation"`) || + strings.Contains(raw, `"traffic_report"`) { + return nil, "legacy observability buffer format", false + } + if err := json.Unmarshal(data, &records); err != nil { + return nil, "observability buffer JSON decode failed", false + } + return records, "", true +} + +func (s *ObservabilityBufferStore) discardBufferFile(reason string) { + if err := os.Remove(s.path); err != nil && !os.IsNotExist(err) { + slog.Warn("remove observability buffer failed", "path", s.path, "reason", reason, "error", err) + } else { + slog.Info("discarded observability buffer; will rebuild on run", "path", s.path, "reason", reason) + } + s.cache = []ObservabilityBufferRecord{} + s.cacheLoaded = true +} + func (s *ObservabilityBufferStore) saveUnlocked(records []ObservabilityBufferRecord) error { if err := os.MkdirAll(filepath.Dir(s.path), stateDirPerm); err != nil { return err @@ -210,7 +250,7 @@ func (s *ObservabilityBufferStore) saveUnlocked(records []ObservabilityBufferRec if err != nil { return err } - if err := os.WriteFile(s.path, data, stateFilePerm); err != nil { + if err := os.WriteFile(s.path, data, stateFilePerm); err != nil { //nolint:gosec // path is agent-local buffer path from config return err } s.cache = records diff --git a/internal/apps/agent/state/observability_buffer_test.go b/internal/apps/agent/state/observability_buffer_test.go index 2fc9923d..3b0aacf4 100644 --- a/internal/apps/agent/state/observability_buffer_test.go +++ b/internal/apps/agent/state/observability_buffer_test.go @@ -1,7 +1,9 @@ package state import ( + "os" "path/filepath" + "strings" "testing" "github.com/Rain-kl/Wavelet/internal/apps/agent/protocol" @@ -95,3 +97,93 @@ func TestObservabilityWindowStartedAt(t *testing.T) { t.Fatalf("unexpected host-metrics window start: %d", value) } } + +func TestObservabilityBufferStoreDiscardsLegacyDiskJSON(t *testing.T) { + path := filepath.Join(t.TempDir(), "observability-buffer.json") + legacy := `[{ + "window_started_at_unix": 1710403200, + "snapshot": {"captured_at_unix": 1710403205, "cpu_usage_percent": 11.5}, + "openresty_observation": {"captured_at_unix": 1710403206, "openresty_connections": 7}, + "traffic_report": {"request_count": 42}, + "access_logs": [{"logged_at_unix": 1710403201, "path": "/", "status_code": 200}] + }]` + if err := os.WriteFile(path, []byte(legacy), 0o644); err != nil { + t.Fatalf("WriteFile: %v", err) + } + + store := NewObservabilityBufferStore(path) + records, err := store.Replayable(0, 0) + if err != nil { + t.Fatalf("Replayable: %v", err) + } + if len(records) != 0 { + t.Fatalf("expected legacy buffer discarded, got %+v", records) + } + // Replayable may rewrite an empty v2 array; legacy keys must be gone. + if body, err := os.ReadFile(path); err == nil { + raw := string(body) + for _, key := range []string{`"snapshot"`, `"openresty_observation"`, `"traffic_report"`} { + if strings.Contains(raw, key) { + t.Fatalf("legacy key %s still present: %s", key, raw) + } + } + } + + // Fresh upsert after discard should create a clean v2 file. + if err := store.Upsert(ObservabilityBufferRecord{ + WindowStartedAtUnix: 1710403200, + HostMetrics: &protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403205}, + }, 0); err != nil { + t.Fatalf("Upsert after discard: %v", err) + } + records, err = store.Replayable(0, 0) + if err != nil { + t.Fatalf("Replayable after rebuild: %v", err) + } + if len(records) != 1 || records[0].HostMetrics == nil { + t.Fatalf("expected rebuilt buffer, got %+v", records) + } +} + +func TestObservabilityBufferStoreDiscardsCorruptJSON(t *testing.T) { + path := filepath.Join(t.TempDir(), "observability-buffer.json") + if err := os.WriteFile(path, []byte(`{not-json`), 0o644); err != nil { + t.Fatalf("WriteFile: %v", err) + } + store := NewObservabilityBufferStore(path) + records, err := store.Replayable(0, 0) + if err != nil { + t.Fatalf("Replayable should not fail: %v", err) + } + if len(records) != 0 { + t.Fatalf("expected empty after discard, got %+v", records) + } + // Corrupt payload must not remain; empty rewrite is fine. + if body, err := os.ReadFile(path); err == nil && strings.Contains(string(body), "not-json") { + t.Fatalf("corrupt content still on disk: %s", body) + } +} + +func TestObservabilityBufferStoreKeepsModernJSON(t *testing.T) { + path := filepath.Join(t.TempDir(), "observability-buffer.json") + modern := `[{ + "window_started_at_unix": 1710403200, + "host_metrics": {"captured_at_unix": 1710403205, "cpu_usage_percent": 3}, + "edge_health": {"captured_at_unix": 1710403205, "status": "healthy", "connections": 2}, + "access_logs": [] + }]` + if err := os.WriteFile(path, []byte(modern), 0o644); err != nil { + t.Fatalf("WriteFile: %v", err) + } + store := NewObservabilityBufferStore(path) + records, err := store.Replayable(0, 0) + if err != nil { + t.Fatalf("Replayable: %v", err) + } + if len(records) != 1 || records[0].HostMetrics == nil || records[0].HostMetrics.CPUUsagePercent != 3 { + t.Fatalf("modern buffer should be kept: %+v", records) + } + if _, err := os.Stat(path); err != nil { + t.Fatalf("modern buffer file should remain: %v", err) + } +} diff --git a/internal/apps/openflare/agent/helpers.go b/internal/apps/openflare/agent/helpers.go index 5bbca57d..187f382c 100644 --- a/internal/apps/openflare/agent/helpers.go +++ b/internal/apps/openflare/agent/helpers.go @@ -64,6 +64,20 @@ func normalizeNodePayload(payload NodePayload) NodePayload { payload.LastError = truncateForDatabase(payload.LastError, maxDatabaseTextLength) payload.OpenrestyStatus = normalizeOpenrestyStatus(payload.OpenrestyStatus) payload.OpenrestyMessage = truncateForDatabase(payload.OpenrestyMessage, maxDatabaseTextLength) + // Align L2 edge_health with top-level status/message (PG is latest-state authority). + if payload.EdgeHealth != nil { + if s := strings.TrimSpace(payload.EdgeHealth.Status); s != "" { + if payload.OpenrestyStatus == "" || payload.OpenrestyStatus == openrestyStatusUnknown { + payload.OpenrestyStatus = normalizeOpenrestyStatus(s) + } + } + if m := strings.TrimSpace(payload.EdgeHealth.Message); m != "" && payload.OpenrestyMessage == "" { + payload.OpenrestyMessage = truncateForDatabase(m, maxDatabaseTextLength) + } + // CH series status must match the same authority as PG after normalize. + payload.EdgeHealth.Status = payload.OpenrestyStatus + payload.EdgeHealth.Message = payload.OpenrestyMessage + } return payload }