feat(obs): 访问日志 SSOT 与 edge_health,去掉协议兼容层

Agent 仅上报 host_metrics/edge_health/access_logs;业务流量与 UV 由
Server 侧访问日志聚合。新增 of_node_edge_health 与 of_access_log_hourly,
删除 request_reports/openresty 吞吐路径;API 不再暴露 traffic_reports
与 openresty_rx|tx。心跳/离线默认阈值与回填迁移一并入库。
This commit is contained in:
ryan
2026-07-18 11:53:11 +08:00
parent 9a0974cce8
commit f0e234df1f
60 changed files with 1799 additions and 2164 deletions
@@ -17,9 +17,7 @@ const (
TableTTLDaysNodeAccessLogs = 90
// TableTTLDaysNodeMetricSnapshots is the of_node_metric_snapshots TTL (30 days).
TableTTLDaysNodeMetricSnapshots = 30
// TableTTLDaysNodeRequestReports is the of_node_request_reports TTL (30 days).
TableTTLDaysNodeRequestReports = 30
// TableTTLDaysNodeObs is the of_node_obs_* TTL (30 days).
// TableTTLDaysNodeObs is the of_node_edge_health / of_node_obs_frps / of_node_obs_frpc TTL (30 days).
TableTTLDaysNodeObs = 30
// TableTTLDaysUserAccessLogs is the w_user_access_logs TTL (180 days).
TableTTLDaysUserAccessLogs = 180
@@ -31,7 +31,6 @@ func TestCleanupModeConstants(t *testing.T) {
func TestTableTTLDaysMatchDDL(t *testing.T) {
assert.Equal(t, 90, TableTTLDaysNodeAccessLogs)
assert.Equal(t, 30, TableTTLDaysNodeMetricSnapshots)
assert.Equal(t, 30, TableTTLDaysNodeRequestReports)
assert.Equal(t, 30, TableTTLDaysNodeObs)
assert.Equal(t, 180, TableTTLDaysUserAccessLogs)
}
@@ -6,6 +6,7 @@ package analytics
import (
"context"
"fmt"
"strings"
"time"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
@@ -35,7 +36,7 @@ func ListNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) ([]anal
clause, args := buildNodeAccessLogFilterClause(filter)
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT id, node_id, logged_at, remote_addr, region, host, path, status_code, bytes_sent, created_at
SELECT id, node_id, logged_at, remote_addr, region, host, path, status_code, bytes_sent, request_length, request_time_ms, created_at
FROM %s
WHERE %s
ORDER BY %s`, tableName, clause, nodeAccessLogOrderClause(filter.SortBy, filter.SortOrder))
@@ -68,6 +69,8 @@ func scanNodeAccessLogRows(rows driver.Rows) ([]analyticsmodel.NodeAccessLog, er
&item.Path,
&item.StatusCode,
&item.BytesSent,
&item.RequestLength,
&item.RequestTimeMs,
&item.CreatedAt,
); err != nil {
return nil, fmt.Errorf("scan node access log row: %w", err)
@@ -143,3 +146,157 @@ ORDER BY count DESC, trimmed_region ASC`, tableName, clause)
}
return result, nil
}
// NodeAccessLogTrafficSummary is a window-level access log traffic summary.
type NodeAccessLogTrafficSummary struct {
RequestCount int64
ErrorCount int64
UniqueIPCount int64
BytesSent int64
RequestLength int64
NodeCount int64
}
// NodeAccessLogValueCount is a grouped value count (status_code, host, ...).
type NodeAccessLogValueCount struct {
Value string
Count int64
}
// NodeAccessLogNodeAggregate is per-node traffic over a window.
type NodeAccessLogNodeAggregate struct {
NodeID string
RequestCount int64
ErrorCount int64
UniqueIPCount int64
}
// TrafficSummaryNodeAccessLogs returns request/error/UV/bytes/node counts for the filter.
func TrafficSummaryNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) (NodeAccessLogTrafficSummary, error) {
conn, err := nodeAccessLogConn()
if err != nil {
return NodeAccessLogTrafficSummary{}, err
}
clause, args := buildNodeAccessLogFilterClause(filter)
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT
count() AS request_count,
countIf(status_code >= 500) AS error_count,
uniqExactIf(remote_addr, remote_addr != '') AS unique_ips,
sum(bytes_sent) AS bytes_sent,
sum(request_length) AS request_length,
uniqExactIf(node_id, node_id != '') AS node_count
FROM %s
WHERE %s`, tableName, clause)
var requestCount, errorCount, uniqueIPs, bytesSent, requestLength, nodeCount uint64
if err := conn.QueryRow(ctx, sql, args...).Scan(
&requestCount, &errorCount, &uniqueIPs, &bytesSent, &requestLength, &nodeCount,
); err != nil {
return NodeAccessLogTrafficSummary{}, fmt.Errorf("traffic summary node access logs: %w", err)
}
return NodeAccessLogTrafficSummary{
RequestCount: safeInt64Count(requestCount),
ErrorCount: safeInt64Count(errorCount),
UniqueIPCount: safeInt64Count(uniqueIPs),
BytesSent: safeInt64Count(bytesSent),
RequestLength: safeInt64Count(requestLength),
NodeCount: safeInt64Count(nodeCount),
}, nil
}
// ValueCountsNodeAccessLogs groups logs by a single dimension column.
// Allowed columns: status_code, host.
func ValueCountsNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter, column string, limit int) ([]NodeAccessLogValueCount, error) {
conn, err := nodeAccessLogConn()
if err != nil {
return nil, err
}
col := strings.TrimSpace(strings.ToLower(column))
switch col {
case nodeAccessLogColumnStatusCode, nodeAccessLogColumnHost:
default:
return nil, fmt.Errorf("unsupported value count column: %s", column)
}
clause, args := buildNodeAccessLogFilterClause(filter)
tableName := nodeAccessLogTableName()
// status_code is numeric; cast to string for a uniform Value field.
var valueExpr string
if col == nodeAccessLogColumnStatusCode {
valueExpr = "toString(" + nodeAccessLogColumnStatusCode + ")"
} else {
valueExpr = "trim(" + nodeAccessLogColumnHost + ")"
}
sql := fmt.Sprintf(`
SELECT %s AS value, count() AS count
FROM %s
WHERE %s AND %s != ''
GROUP BY value
ORDER BY count DESC, value ASC`, valueExpr, tableName, clause, valueExpr)
if limit > 0 {
sql += clickHouseLimitClause
args = append(args, limit)
}
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("value counts node access logs: %w", err)
}
defer func() { _ = rows.Close() }()
var result []NodeAccessLogValueCount
for rows.Next() {
var (
value string
count uint64
)
if err := rows.Scan(&value, &count); err != nil {
return nil, fmt.Errorf("scan value count row: %w", err)
}
result = append(result, NodeAccessLogValueCount{
Value: value,
Count: safeInt64Count(count),
})
}
return result, nil
}
// NodeAggregatesNodeAccessLogs returns per-node request/error/UV aggregates.
func NodeAggregatesNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) ([]NodeAccessLogNodeAggregate, error) {
conn, err := nodeAccessLogConn()
if err != nil {
return nil, err
}
clause, args := buildNodeAccessLogFilterClause(filter)
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT
node_id,
count() AS request_count,
countIf(status_code >= 500) AS error_count,
uniqExactIf(remote_addr, remote_addr != '') AS unique_ips
FROM %s
WHERE %s AND node_id != ''
GROUP BY node_id
ORDER BY request_count DESC, node_id ASC`, tableName, clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("node aggregates node access logs: %w", err)
}
defer func() { _ = rows.Close() }()
var result []NodeAccessLogNodeAggregate
for rows.Next() {
var (
nodeID string
requestCount, errorCount, uniqueIPs uint64
)
if err := rows.Scan(&nodeID, &requestCount, &errorCount, &uniqueIPs); err != nil {
return nil, fmt.Errorf("scan node aggregate row: %w", err)
}
result = append(result, NodeAccessLogNodeAggregate{
NodeID: nodeID,
RequestCount: safeInt64Count(requestCount),
ErrorCount: safeInt64Count(errorCount),
UniqueIPCount: safeInt64Count(uniqueIPs),
})
}
return result, nil
}
@@ -17,6 +17,10 @@ const (
nodeAccessLogSortAscInput = "asc"
nodeAccessLogColumnRemoteAddr = "remote_addr"
nodeAccessLogColumnStatusCode = "status_code"
nodeAccessLogColumnHost = "host"
nodeAccessLogColumnPath = "path"
nodeAccessLogColumnLoggedAt = "logged_at"
)
// NodeAccessLogFilter scopes ClickHouse node access log queries.
@@ -88,21 +92,21 @@ func nodeAccessLogOrderClause(sortBy string, sortOrder string) string {
if normalizeNodeAccessLogSortOrder(sortOrder) == nodeAccessLogSortAscInput {
direction = nodeAccessLogSortAsc
}
column := "logged_at"
column := nodeAccessLogColumnLoggedAt
switch strings.TrimSpace(sortBy) {
case "status_code":
column = "status_code"
case nodeAccessLogColumnStatusCode:
column = nodeAccessLogColumnStatusCode
case nodeAccessLogColumnRemoteAddr:
column = nodeAccessLogColumnRemoteAddr
case "host":
column = "host"
case "path":
column = "path"
case nodeAccessLogColumnHost:
column = nodeAccessLogColumnHost
case nodeAccessLogColumnPath:
column = nodeAccessLogColumnPath
}
if column == "logged_at" {
if column == nodeAccessLogColumnLoggedAt {
return column + " " + direction + ", id " + direction
}
return column + " " + direction + ", logged_at " + direction + ", id " + direction
return column + " " + direction + ", " + nodeAccessLogColumnLoggedAt + " " + direction + ", id " + direction
}
func normalizeNodeAccessLogRemoteAddr(value string) string {
@@ -48,7 +48,8 @@ SELECT
countIf(status_code >= 500) AS server_error_count,
uniqExactIf(remote_addr, remote_addr != '') AS unique_ip_count,
uniqExactIf(host, host != '') AS unique_host_count,
sum(bytes_sent) AS bytes_sent
sum(bytes_sent) AS bytes_sent,
sum(request_length) AS request_length
FROM %s
WHERE %s
GROUP BY bucket_epoch
@@ -69,10 +70,10 @@ ORDER BY %s`, bucketExpr, tableName, clause, nodeAccessLogBucketOrderClause(filt
var result []NodeAccessLogBucketAggregate
for rows.Next() {
var (
bucketEpoch int64
requestCount, successCount, clientErrorCount, serverErrorCount, uniqueIPCount, uniqueHostCount, bytesSent uint64
bucketEpoch int64
requestCount, successCount, clientErrorCount, serverErrorCount, uniqueIPCount, uniqueHostCount, bytesSent, requestLength uint64
)
if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &uniqueIPCount, &uniqueHostCount, &bytesSent); err != nil {
if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &uniqueIPCount, &uniqueHostCount, &bytesSent, &requestLength); err != nil {
return nil, fmt.Errorf("scan bucket aggregate row: %w", err)
}
result = append(result, NodeAccessLogBucketAggregate{
@@ -84,6 +85,7 @@ ORDER BY %s`, bucketExpr, tableName, clause, nodeAccessLogBucketOrderClause(filt
UniqueIPCount: safeInt64Count(uniqueIPCount),
UniqueHostCount: safeInt64Count(uniqueHostCount),
BytesSent: safeInt64Count(bytesSent),
RequestLength: safeInt64Count(requestLength),
})
}
return result, nil
@@ -49,6 +49,8 @@ func TestBatchInsertNodeAccessLogs_UsesModelBatchSQL(t *testing.T) {
assert.True(t, mockBatch.sendCalled)
require.Len(t, mockBatch.rows, 1)
assert.Equal(t, "node-a", mockBatch.rows[0][1])
require.Len(t, mockBatch.rows[0], 10)
assert.Equal(t, uint64(2048), mockBatch.rows[0][8])
require.Len(t, mockBatch.rows[0], 12)
assert.Equal(t, uint64(2048), mockBatch.rows[0][8]) // bytes_sent
assert.Equal(t, uint64(0), mockBatch.rows[0][9]) // request_length
assert.Equal(t, uint32(0), mockBatch.rows[0][10]) // request_time_ms
}
@@ -48,6 +48,8 @@ func BatchInsertNodeAccessLogs(ctx context.Context, logs []analyticsmodel.NodeAc
logItem.Path,
logItem.StatusCode,
logItem.BytesSent,
logItem.RequestLength,
logItem.RequestTimeMs,
createdAt.UTC(),
); err != nil {
return fmt.Errorf("append node access log to batch: %w", err)
@@ -67,62 +67,16 @@ ORDER BY %s%s`, nodeMetricSnapshotTableName(), clause, nodeObservabilityCaptured
return scanNodeMetricSnapshotRows(rows)
}
// ListNodeRequestReports returns request reports matching filter.
func ListNodeRequestReports(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeRequestReport, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
}
clause, args := buildNodeObservabilityFilterClause(filter, "window_ended_at")
tableName := nodeRequestReportTableName()
sql := fmt.Sprintf(`
SELECT id, node_id, window_started_at, window_ended_at, request_count, error_count, unique_visitor_count, status_codes_json, top_domains_json, source_countries_json, created_at
FROM %s
WHERE %s
ORDER BY %s`, tableName, clause, nodeObservabilityWindowEndedAtOrderClause())
if filter.Limit > 0 {
sql += clickHouseLimitClause
args = append(args, filter.Limit)
}
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list node request reports: %w", err)
}
defer func() { _ = rows.Close() }()
return scanNodeRequestReportRows(rows)
}
// ListLatestNodeRequestReports returns the latest request report per node_id.
// Uses ClickHouse LIMIT 1 BY so dashboard traffic health is not skewed by a global raw LIMIT.
func ListLatestNodeRequestReports(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeRequestReport, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
}
clause, args := buildNodeObservabilityFilterClause(filter, "window_ended_at")
sql := fmt.Sprintf(`
SELECT id, node_id, window_started_at, window_ended_at, request_count, error_count, unique_visitor_count, status_codes_json, top_domains_json, source_countries_json, created_at
FROM %s
WHERE %s
ORDER BY %s%s`, nodeRequestReportTableName(), clause, nodeObservabilityWindowEndedAtOrderClause(), clickHouseLimit1ByNodeIDClause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list latest node request reports: %w", err)
}
defer func() { _ = rows.Close() }()
return scanNodeRequestReportRows(rows)
}
// ListNodeObsOpenresty returns OpenResty observations matching filter.
func ListNodeObsOpenresty(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeObsOpenresty, error) {
// ListNodeEdgeHealth returns L2 OpenResty health snapshots.
func ListNodeEdgeHealth(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeEdgeHealth, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
}
clause, args := buildNodeObservabilityFilterClause(filter, "captured_at")
tableName := nodeObsOpenrestyTableName()
tableName := nodeEdgeHealthTableName()
sql := fmt.Sprintf(`
SELECT id, node_id, captured_at, openresty_rx_bytes, openresty_tx_bytes, openresty_connections, created_at
SELECT id, node_id, captured_at, status, connections, created_at
FROM %s
WHERE %s
ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause())
@@ -132,10 +86,27 @@ ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause())
}
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list node openresty observations: %w", err)
return nil, fmt.Errorf("list node edge health: %w", err)
}
defer func() { _ = rows.Close() }()
return scanNodeObsOpenrestyRows(rows)
var result []analyticsmodel.NodeEdgeHealth
for rows.Next() {
var item analyticsmodel.NodeEdgeHealth
if err := rows.Scan(
&item.ID,
&item.NodeID,
&item.CapturedAt,
&item.Status,
&item.Connections,
&item.CreatedAt,
); err != nil {
return nil, fmt.Errorf("scan node edge health row: %w", err)
}
item.CapturedAt = item.CapturedAt.UTC()
item.CreatedAt = item.CreatedAt.UTC()
result = append(result, item)
}
return result, nil
}
// ListNodeObsFrps returns FRPS observations matching filter.
@@ -216,55 +187,6 @@ func scanNodeMetricSnapshotRows(rows driver.Rows) ([]analyticsmodel.NodeMetricSn
return result, nil
}
func scanNodeRequestReportRows(rows driver.Rows) ([]analyticsmodel.NodeRequestReport, error) {
var result []analyticsmodel.NodeRequestReport
for rows.Next() {
var item analyticsmodel.NodeRequestReport
if err := rows.Scan(
&item.ID,
&item.NodeID,
&item.WindowStartedAt,
&item.WindowEndedAt,
&item.RequestCount,
&item.ErrorCount,
&item.UniqueVisitorCount,
&item.StatusCodesJSON,
&item.TopDomainsJSON,
&item.SourceCountriesJSON,
&item.CreatedAt,
); err != nil {
return nil, fmt.Errorf("scan node request report row: %w", err)
}
item.WindowStartedAt = item.WindowStartedAt.UTC()
item.WindowEndedAt = item.WindowEndedAt.UTC()
item.CreatedAt = item.CreatedAt.UTC()
result = append(result, item)
}
return result, nil
}
func scanNodeObsOpenrestyRows(rows driver.Rows) ([]analyticsmodel.NodeObsOpenresty, error) {
var result []analyticsmodel.NodeObsOpenresty
for rows.Next() {
var item analyticsmodel.NodeObsOpenresty
if err := rows.Scan(
&item.ID,
&item.NodeID,
&item.CapturedAt,
&item.OpenrestyRxBytes,
&item.OpenrestyTxBytes,
&item.OpenrestyConnections,
&item.CreatedAt,
); err != nil {
return nil, fmt.Errorf("scan node openresty observation row: %w", err)
}
item.CapturedAt = item.CapturedAt.UTC()
item.CreatedAt = item.CreatedAt.UTC()
result = append(result, item)
}
return result, nil
}
func scanNodeObsFrpsRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrps, error) {
var result []analyticsmodel.NodeObsFrps
for rows.Next() {
@@ -288,13 +210,10 @@ func scanNodeObsFrpsRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrps, error)
return result, nil
}
const nodeTrafficHourlyTableName = "of_node_traffic_hourly"
// NodeTrafficHourly is an hourly traffic rollup row.
//
// UniqueVisitorCount is a peak per-window estimate from short request reports
// (MV uses max()), not true distinct visitors across the hour. SummingMergeTree
// may still inflate residual unmerged parts; do not present as exact UV.
// UniqueVisitorCount is always 0 when sourced from of_access_log_hourly
// (true UV requires raw uniqExact on access logs).
type NodeTrafficHourly struct {
NodeID string
Hour time.Time
@@ -319,16 +238,40 @@ type NodeMetricHourly struct {
ReportedNodes int
}
// NodeOpenrestyHourly is an hourly OpenResty observation aggregation row.
type NodeOpenrestyHourly struct {
Hour time.Time
OpenrestyRxBytes int64
OpenrestyTxBytes int64
ReportedNodes int
// ListNodeTrafficHourly returns hourly traffic from of_access_log_hourly (M5).
// UniqueVisitorCount is always 0 here (UV requires raw uniqExact on access logs).
func ListNodeTrafficHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeTrafficHourly, error) {
rows, err := ListAccessLogHourly(ctx, filter)
if err != nil {
return nil, err
}
// Aggregate across hosts per node/hour.
type key struct {
node string
hour int64
}
merged := make(map[key]*NodeTrafficHourly)
order := make([]key, 0)
for _, row := range rows {
k := key{node: row.NodeID, hour: row.Hour.UTC().Unix()}
item := merged[k]
if item == nil {
item = &NodeTrafficHourly{NodeID: row.NodeID, Hour: row.Hour.UTC()}
merged[k] = item
order = append(order, k)
}
item.RequestCount += row.RequestCount
item.ErrorCount += row.ErrorCount
}
result := make([]NodeTrafficHourly, 0, len(order))
for _, k := range order {
result = append(result, *merged[k])
}
return result, nil
}
// ListNodeTrafficHourly returns hourly traffic rollup rows matching filter.
func ListNodeTrafficHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeTrafficHourly, error) {
// ListAccessLogHourly returns Server-side access log hourly rollups.
func ListAccessLogHourly(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.AccessLogHourly, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
@@ -338,32 +281,42 @@ func ListNodeTrafficHourly(ctx context.Context, filter NodeObservabilityFilter)
SELECT
node_id,
hour,
host,
sum(request_count) AS request_count,
sum(error_count) AS error_count,
max(unique_visitor_count) AS unique_visitor_count
sum(bytes_sent) AS bytes_sent,
sum(request_length) AS request_length
FROM %s
WHERE %s
GROUP BY node_id, hour
ORDER BY hour ASC`, nodeTrafficHourlyTableName, clause)
GROUP BY node_id, hour, host
ORDER BY hour ASC, node_id ASC, host ASC`, accessLogHourlyTableName(), clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list node traffic hourly: %w", err)
return nil, fmt.Errorf("list access log hourly: %w", err)
}
defer func() { _ = rows.Close() }()
result := make([]NodeTrafficHourly, 0)
var result []analyticsmodel.AccessLogHourly
for rows.Next() {
var (
item NodeTrafficHourly
requestCount, errorCount, uniqueVisitorCount uint64
item analyticsmodel.AccessLogHourly
requestCount, errorCount, bytesSent, requestLength uint64
)
if err := rows.Scan(&item.NodeID, &item.Hour, &requestCount, &errorCount, &uniqueVisitorCount); err != nil {
return nil, fmt.Errorf("scan node traffic hourly row: %w", err)
if err := rows.Scan(
&item.NodeID,
&item.Hour,
&item.Host,
&requestCount,
&errorCount,
&bytesSent,
&requestLength,
); err != nil {
return nil, fmt.Errorf("scan access log hourly row: %w", err)
}
item.Hour = item.Hour.UTC()
item.RequestCount = safeInt64Count(requestCount)
item.ErrorCount = safeInt64Count(errorCount)
item.UniqueVisitorCount = safeInt64Count(uniqueVisitorCount)
item.BytesSent = safeInt64Count(bytesSent)
item.RequestLength = safeInt64Count(requestLength)
result = append(result, item)
}
return result, nil
@@ -556,138 +509,6 @@ func scanNodeMetricHourlyRows(rows driver.Rows) ([]NodeMetricHourly, error) {
return result, nil
}
// ListNodeOpenrestyHourly returns hourly OpenResty observation aggregates matching filter.
// Same rollup-first / per-hour merge strategy as ListNodeMetricHourly.
func ListNodeOpenrestyHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) {
rollup, rollupErr := listNodeOpenrestyHourlyFromRollup(ctx, filter)
if rollupErr == nil && len(rollup) > 0 && hourlyRollupCoversWindow(rollup[0].Hour, filter.Since) {
return rollup, nil
}
raw, rawErr := listNodeOpenrestyHourlyFromRaw(ctx, filter)
if rawErr != nil {
if rollupErr == nil && len(rollup) > 0 {
return rollup, nil
}
return nil, rawErr
}
if len(rollup) == 0 {
return raw, nil
}
return mergeNodeOpenrestyHourlyPreferRollup(rollup, raw), nil
}
func mergeNodeOpenrestyHourlyPreferRollup(rollup, raw []NodeOpenrestyHourly) []NodeOpenrestyHourly {
byHour := make(map[int64]NodeOpenrestyHourly, len(raw)+len(rollup))
order := make([]int64, 0, len(raw)+len(rollup))
add := func(row NodeOpenrestyHourly, overwrite bool) {
key := row.Hour.UTC().Truncate(time.Hour).Unix()
if _, exists := byHour[key]; !exists {
order = append(order, key)
byHour[key] = row
return
}
if overwrite {
byHour[key] = row
}
}
for _, row := range raw {
add(row, false)
}
for _, row := range rollup {
add(row, true)
}
sort.Slice(order, func(i, j int) bool { return order[i] < order[j] })
result := make([]NodeOpenrestyHourly, 0, len(order))
for _, key := range order {
result = append(result, byHour[key])
}
return result
}
func listNodeOpenrestyHourlyFromRollup(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
}
clause, args := buildNodeObservabilityFilterClause(filter, "hour")
sql := fmt.Sprintf(`
SELECT
hour,
sum(greatest(openresty_rx_max - openresty_rx_min, 0)) AS openresty_rx_bytes,
sum(greatest(openresty_tx_max - openresty_tx_min, 0)) AS openresty_tx_bytes,
toUInt64(uniqExact(node_id)) AS reported_nodes
FROM %s
WHERE %s
GROUP BY hour
ORDER BY hour ASC`, nodeOpenrestyHourlyTableName(), clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list node openresty hourly from rollup: %w", err)
}
defer func() { _ = rows.Close() }()
return scanNodeOpenrestyHourlyRows(rows)
}
func listNodeOpenrestyHourlyFromRaw(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
}
clause, args := buildNodeObservabilityFilterClause(filter, "captured_at")
tableName := nodeObsOpenrestyTableName()
sql := fmt.Sprintf(`
SELECT
hour,
sum(if(openresty_rx_delta >= 0, openresty_rx_delta, 0)) AS openresty_rx_bytes,
sum(if(openresty_tx_delta >= 0, openresty_tx_delta, 0)) AS openresty_tx_bytes,
toUInt64(uniqExact(node_id)) AS reported_nodes
FROM (
SELECT
node_id,
toStartOfHour(captured_at) AS hour,
openresty_rx_bytes - lagInFrame(openresty_rx_bytes, 1, openresty_rx_bytes) OVER (
PARTITION BY node_id ORDER BY captured_at, id
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
) AS openresty_rx_delta,
openresty_tx_bytes - lagInFrame(openresty_tx_bytes, 1, openresty_tx_bytes) OVER (
PARTITION BY node_id ORDER BY captured_at, id
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
) AS openresty_tx_delta
FROM %s
WHERE %s
)
GROUP BY hour
ORDER BY hour ASC`, tableName, clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list node openresty hourly: %w", err)
}
defer func() { _ = rows.Close() }()
return scanNodeOpenrestyHourlyRows(rows)
}
func scanNodeOpenrestyHourlyRows(rows driver.Rows) ([]NodeOpenrestyHourly, error) {
result := make([]NodeOpenrestyHourly, 0)
for rows.Next() {
var (
item NodeOpenrestyHourly
reportedNodes uint64
rx int64
tx int64
)
if err := rows.Scan(&item.Hour, &rx, &tx, &reportedNodes); err != nil {
return nil, fmt.Errorf("scan node openresty hourly row: %w", err)
}
item.Hour = item.Hour.UTC()
item.OpenrestyRxBytes = rx
item.OpenrestyTxBytes = tx
item.ReportedNodes = int(safeInt64Count(reportedNodes))
result = append(result, item)
}
return result, nil
}
func scanNodeObsFrpcRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrpc, error) {
var result []analyticsmodel.NodeObsFrpc
for rows.Next() {
@@ -51,78 +51,35 @@ func MaterializeNodeMetricSnapshotsTTL(ctx context.Context) (CleanupOutcome, err
)
}
// DeleteAllNodeRequestReports hard-deletes all node request reports via TRUNCATE.
func DeleteAllNodeRequestReports(ctx context.Context) (int64, error) {
// DeleteAllNodeEdgeHealth truncates of_node_edge_health.
func DeleteAllNodeEdgeHealth(ctx context.Context) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
outcome, err := truncateClickHouseTable(ctx, conn, nodeRequestReportTableName())
outcome, err := truncateClickHouseTable(ctx, conn, nodeEdgeHealthTableName())
if err != nil {
return 0, err
}
return outcome.DeletedCount, nil
}
// DeleteNodeRequestReportsBefore force-materializes of_node_request_reports table TTL.
// cutoff is ignored; see MaterializeNodeRequestReportsTTL.
func DeleteNodeRequestReportsBefore(ctx context.Context, _ time.Time) (int64, error) {
outcome, err := MaterializeNodeRequestReportsTTL(ctx)
// DeleteNodeEdgeHealthBefore force-materializes of_node_edge_health TTL.
func DeleteNodeEdgeHealthBefore(ctx context.Context, _ time.Time) (int64, error) {
outcome, err := MaterializeNodeEdgeHealthTTL(ctx)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// MaterializeNodeRequestReportsTTL force-materializes table TTL and reports an honest outcome.
func MaterializeNodeRequestReportsTTL(ctx context.Context) (CleanupOutcome, error) {
// MaterializeNodeEdgeHealthTTL force-materializes of_node_edge_health table TTL.
func MaterializeNodeEdgeHealthTTL(ctx context.Context) (CleanupOutcome, error) {
conn, err := observabilityConn()
if err != nil {
return CleanupOutcome{}, err
}
tableName := nodeRequestReportTableName()
ttlDays := TableTTLDaysNodeRequestReports
cutoff := tableTTLCutoff(ttlDays, time.Now())
return materializeExpiredByTableTTL(
ctx,
conn,
tableName,
ttlDays,
fmt.Sprintf("SELECT count() FROM %s WHERE window_ended_at < ?", tableName),
[]any{cutoff},
)
}
// DeleteAllNodeObsOpenresty hard-deletes all OpenResty observations via TRUNCATE.
func DeleteAllNodeObsOpenresty(ctx context.Context) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
outcome, err := truncateClickHouseTable(ctx, conn, nodeObsOpenrestyTableName())
if err != nil {
return 0, err
}
return outcome.DeletedCount, nil
}
// DeleteNodeObsOpenrestyBefore force-materializes of_node_obs_openresty table TTL.
// cutoff is ignored; see MaterializeNodeObsOpenrestyTTL.
func DeleteNodeObsOpenrestyBefore(ctx context.Context, _ time.Time) (int64, error) {
outcome, err := MaterializeNodeObsOpenrestyTTL(ctx)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// MaterializeNodeObsOpenrestyTTL force-materializes table TTL and reports an honest outcome.
func MaterializeNodeObsOpenrestyTTL(ctx context.Context) (CleanupOutcome, error) {
conn, err := observabilityConn()
if err != nil {
return CleanupOutcome{}, err
}
tableName := nodeObsOpenrestyTableName()
tableName := nodeEdgeHealthTableName()
ttlDays := TableTTLDaysNodeObs
cutoff := tableTTLCutoff(ttlDays, time.Now())
return materializeExpiredByTableTTL(
@@ -38,20 +38,16 @@ func nodeObservabilityCapturedAtOrderClause() string {
return "captured_at DESC, id DESC"
}
func nodeObservabilityWindowEndedAtOrderClause() string {
return "window_ended_at DESC, id DESC"
}
func nodeMetricSnapshotTableName() string {
return "of_node_metric_snapshots"
}
func nodeRequestReportTableName() string {
return "of_node_request_reports"
func nodeEdgeHealthTableName() string {
return "of_node_edge_health"
}
func nodeObsOpenrestyTableName() string {
return "of_node_obs_openresty"
func accessLogHourlyTableName() string {
return "of_access_log_hourly"
}
func nodeObsFrpsTableName() string {
@@ -66,9 +62,5 @@ func nodeMetricCapacityHourlyTableName() string {
return "of_node_metric_capacity_hourly"
}
func nodeOpenrestyHourlyTableName() string {
return "of_node_openresty_hourly"
}
// clickHouseLimit1ByNodeIDClause selects the first row per node_id after ORDER BY.
const clickHouseLimit1ByNodeIDClause = " LIMIT 1 BY node_id"
@@ -35,23 +35,6 @@ func TestListLatestNodeMetricSnapshots_UsesLimit1ByNodeID(t *testing.T) {
assert.Equal(t, since, mock.queryArgs[0][0])
}
func TestListLatestNodeRequestReports_UsesLimit1ByNodeID(t *testing.T) {
ctx := context.Background()
mock := &mockConn{}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
_, err := ListLatestNodeRequestReports(ctx, NodeObservabilityFilter{NodeID: "node-a"})
require.NoError(t, err)
require.Len(t, mock.queries, 1)
assert.Contains(t, mock.queries[0], "LIMIT 1 BY node_id")
assert.Contains(t, mock.queries[0], nodeRequestReportTableName())
assert.Contains(t, mock.queries[0], "window_ended_at DESC")
require.Len(t, mock.queryArgs, 1)
require.Len(t, mock.queryArgs[0], 1)
assert.Equal(t, "node-a", mock.queryArgs[0][0])
}
func TestListNodeMetricHourly_PrefersRollup(t *testing.T) {
ctx := context.Background()
hour := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC)
@@ -170,55 +153,3 @@ func TestListNodeMetricHourly_FallsBackToRawOnRollupError(t *testing.T) {
assert.Contains(t, mock.queries[0], nodeMetricCapacityHourlyTableName())
assert.Contains(t, mock.queries[1], "lagInFrame")
}
func TestListNodeOpenrestyHourly_PrefersRollup(t *testing.T) {
ctx := context.Background()
hour := time.Date(2026, 7, 10, 14, 0, 0, 0, time.UTC)
since := hour.Add(-1 * time.Hour)
mock := &mockConn{
queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) {
if strings.Contains(query, nodeOpenrestyHourlyTableName()) {
return &mockRows{data: [][]any{{
hour, int64(50), int64(70), uint64(3),
}}}, nil
}
return nil, errors.New("raw path should not be used when rollup covers the window")
},
}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
rows, err := ListNodeOpenrestyHourly(ctx, NodeObservabilityFilter{Since: since})
require.NoError(t, err)
require.Len(t, rows, 1)
assert.Equal(t, int64(50), rows[0].OpenrestyRxBytes)
assert.Equal(t, int64(70), rows[0].OpenrestyTxBytes)
assert.Equal(t, 3, rows[0].ReportedNodes)
}
func TestListNodeOpenrestyHourly_FallsBackToRawOnEmptyRollup(t *testing.T) {
ctx := context.Background()
hour := time.Date(2026, 7, 10, 15, 0, 0, 0, time.UTC)
mock := &mockConn{
queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) {
if strings.Contains(query, nodeOpenrestyHourlyTableName()) {
return &mockRows{}, nil
}
if strings.Contains(query, nodeObsOpenrestyTableName()) {
return &mockRows{data: [][]any{{
hour, int64(9), int64(8), uint64(1),
}}}, nil
}
return &mockRows{}, nil
},
}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
rows, err := ListNodeOpenrestyHourly(ctx, NodeObservabilityFilter{})
require.NoError(t, err)
require.Len(t, rows, 1)
assert.Equal(t, int64(9), rows[0].OpenrestyRxBytes)
require.GreaterOrEqual(t, len(mock.queries), 2)
assert.Contains(t, mock.queries[1], "lagInFrame")
}
@@ -14,34 +14,35 @@ import (
"github.com/stretchr/testify/require"
)
func TestInsertNodeObsOpenresty_EmptyNodeID(t *testing.T) {
err := InsertNodeObsOpenresty(context.Background(), analyticsmodel.NodeObsOpenresty{})
func TestInsertNodeEdgeHealth_EmptyNodeID(t *testing.T) {
err := InsertNodeEdgeHealth(context.Background(), analyticsmodel.NodeEdgeHealth{})
require.NoError(t, err)
}
func TestInsertNodeObsOpenresty_UsesModelBatchSQL(t *testing.T) {
func TestInsertNodeEdgeHealth_UsesEdgeHealthBatchSQL(t *testing.T) {
ctx := context.Background()
mockBatch := &mockBatch{}
mockConn := &mockConn{
batch: mockBatch,
batchQuery: analyticsmodel.NodeObsOpenresty{}.BatchInsertSQL(),
batchQuery: analyticsmodel.NodeEdgeHealth{}.BatchInsertSQL(),
}
db.SetChConnForTest(mockConn)
t.Cleanup(func() { db.SetChConnForTest(nil) })
capturedAt := time.Now().UTC()
err := InsertNodeObsOpenresty(ctx, analyticsmodel.NodeObsOpenresty{
NodeID: "node-a",
CapturedAt: capturedAt,
OpenrestyRxBytes: 100,
OpenrestyTxBytes: 200,
OpenrestyConnections: 3,
CreatedAt: capturedAt,
err := InsertNodeEdgeHealth(ctx, analyticsmodel.NodeEdgeHealth{
NodeID: "node-a",
CapturedAt: capturedAt,
Status: "",
Connections: 3,
CreatedAt: capturedAt,
})
require.NoError(t, err)
assert.True(t, mockConn.prepareCalled)
assert.Equal(t, analyticsmodel.NodeObsOpenresty{}.BatchInsertSQL(), mockConn.preparedQuery)
assert.Equal(t, analyticsmodel.NodeEdgeHealth{}.BatchInsertSQL(), mockConn.preparedQuery)
assert.True(t, mockBatch.sendCalled)
require.Len(t, mockBatch.rows, 1)
assert.Equal(t, "node-a", mockBatch.rows[0][1])
assert.Equal(t, "unknown", mockBatch.rows[0][3]) // status default
assert.Equal(t, int64(3), mockBatch.rows[0][4]) // connections
}
@@ -14,6 +14,8 @@ import (
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
)
const edgeHealthStatusUnknown = "unknown"
// InsertNodeMetricSnapshot writes a single metric snapshot via the batch API.
func InsertNodeMetricSnapshot(ctx context.Context, snapshot analyticsmodel.NodeMetricSnapshot) error {
if strings.TrimSpace(snapshot.NodeID) == "" {
@@ -78,105 +80,49 @@ func BatchInsertNodeMetricSnapshots(ctx context.Context, snapshots []analyticsmo
return nil
}
// InsertNodeRequestReport writes a single request report via the batch API.
func InsertNodeRequestReport(ctx context.Context, report analyticsmodel.NodeRequestReport) error {
if strings.TrimSpace(report.NodeID) == "" {
return nil
func normalizeEdgeHealthStatus(status string) string {
status = strings.TrimSpace(status)
if status == "" {
return edgeHealthStatusUnknown
}
return BatchInsertNodeRequestReports(ctx, []analyticsmodel.NodeRequestReport{report})
return status
}
// BatchInsertNodeRequestReports writes request reports to ClickHouse.
func BatchInsertNodeRequestReports(ctx context.Context, reports []analyticsmodel.NodeRequestReport) error {
if len(reports) == 0 {
// InsertNodeEdgeHealth writes a single edge health snapshot.
func InsertNodeEdgeHealth(ctx context.Context, row analyticsmodel.NodeEdgeHealth) error {
if strings.TrimSpace(row.NodeID) == "" {
return nil
}
return BatchInsertNodeEdgeHealth(ctx, []analyticsmodel.NodeEdgeHealth{row})
}
// BatchInsertNodeEdgeHealth writes L2 OpenResty health snapshots to ClickHouse.
func BatchInsertNodeEdgeHealth(ctx context.Context, rows []analyticsmodel.NodeEdgeHealth) error {
if len(rows) == 0 {
return nil
}
if db.ChConn == nil {
return fmt.Errorf("clickhouse connection is not initialized")
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeRequestReport{}.BatchInsertSQL())
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeEdgeHealth{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
now := time.Now().UTC()
for _, report := range reports {
nodeID := strings.TrimSpace(report.NodeID)
for _, row := range rows {
nodeID := strings.TrimSpace(row.NodeID)
if nodeID == "" {
continue
}
id := report.ID
id := row.ID
if id == 0 {
id = idgen.NextUint64ID()
}
createdAt := report.CreatedAt
createdAt := row.CreatedAt
if createdAt.IsZero() {
createdAt = now
}
if err := batch.Append(
id,
nodeID,
report.WindowStartedAt.UTC(),
report.WindowEndedAt.UTC(),
report.RequestCount,
report.ErrorCount,
report.UniqueVisitorCount,
report.StatusCodesJSON,
report.TopDomainsJSON,
report.SourceCountriesJSON,
createdAt.UTC(),
); err != nil {
return fmt.Errorf("append node request report to batch: %w", err)
}
}
if batch.Rows() == 0 {
return nil
}
if err := batch.Send(); err != nil {
return fmt.Errorf("send clickhouse batch: %w", err)
}
return nil
}
// InsertNodeObsOpenresty writes a single OpenResty observation via the batch API.
func InsertNodeObsOpenresty(ctx context.Context, obs analyticsmodel.NodeObsOpenresty) error {
if strings.TrimSpace(obs.NodeID) == "" {
return nil
}
return BatchInsertNodeObsOpenresty(ctx, []analyticsmodel.NodeObsOpenresty{obs})
}
// BatchInsertNodeObsOpenresty writes OpenResty observations to ClickHouse.
func BatchInsertNodeObsOpenresty(ctx context.Context, observations []analyticsmodel.NodeObsOpenresty) error {
if len(observations) == 0 {
return nil
}
if db.ChConn == nil {
return fmt.Errorf("clickhouse connection is not initialized")
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeObsOpenresty{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
now := time.Now().UTC()
for _, obs := range observations {
nodeID := strings.TrimSpace(obs.NodeID)
if nodeID == "" {
continue
}
id := obs.ID
if id == 0 {
id = idgen.NextUint64ID()
}
createdAt := obs.CreatedAt
if createdAt.IsZero() {
createdAt = now
}
capturedAt := obs.CapturedAt.UTC()
capturedAt := row.CapturedAt.UTC()
if capturedAt.IsZero() {
capturedAt = now
}
@@ -184,15 +130,13 @@ func BatchInsertNodeObsOpenresty(ctx context.Context, observations []analyticsmo
id,
nodeID,
capturedAt,
obs.OpenrestyRxBytes,
obs.OpenrestyTxBytes,
obs.OpenrestyConnections,
normalizeEdgeHealthStatus(row.Status),
row.Connections,
createdAt.UTC(),
); err != nil {
return fmt.Errorf("append node openresty observation to batch: %w", err)
return fmt.Errorf("append node edge health to batch: %w", err)
}
}
if batch.Rows() == 0 {
return nil
}