mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 14:06:36 +08:00
f0e234df1f
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。心跳/离线默认阈值与回填迁移一并入库。
354 lines
12 KiB
Go
354 lines
12 KiB
Go
// Copyright 2026 Arctel.net
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
|
|
package model
|
|
|
|
import (
|
|
"context"
|
|
"math"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
|
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
|
|
)
|
|
|
|
// ObservabilityInsertHooks queues observability rows for async ClickHouse write.
|
|
// Wired from openflare/chwriter.Init so model never imports the apps layer.
|
|
type ObservabilityInsertHooks struct {
|
|
QueueMetricSnapshot func(analyticsmodel.NodeMetricSnapshot)
|
|
QueueEdgeHealth func(analyticsmodel.NodeEdgeHealth)
|
|
QueueFrpsObservation func(analyticsmodel.NodeObsFrps)
|
|
QueueFrpcObservation func(analyticsmodel.NodeObsFrpc)
|
|
}
|
|
|
|
var (
|
|
observabilityInsertHooksMu sync.RWMutex
|
|
observabilityInsertHooks ObservabilityInsertHooks
|
|
)
|
|
|
|
// SetObservabilityInsertHooks registers async queue callbacks for observability inserts.
|
|
func SetObservabilityInsertHooks(hooks ObservabilityInsertHooks) {
|
|
observabilityInsertHooksMu.Lock()
|
|
observabilityInsertHooks = hooks
|
|
observabilityInsertHooksMu.Unlock()
|
|
}
|
|
|
|
func currentObservabilityInsertHooks() ObservabilityInsertHooks {
|
|
observabilityInsertHooksMu.RLock()
|
|
defer observabilityInsertHooksMu.RUnlock()
|
|
return observabilityInsertHooks
|
|
}
|
|
|
|
type observabilityStore interface {
|
|
InsertMetricSnapshot(ctx context.Context, record *OpenFlareMetricSnapshot) error
|
|
ListMetricSnapshots(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareMetricSnapshot, error)
|
|
DeleteAllMetricSnapshots(ctx context.Context) (int64, error)
|
|
DeleteMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error)
|
|
|
|
InsertEdgeHealth(ctx context.Context, record *OpenFlareEdgeHealth) error
|
|
ListEdgeHealth(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareEdgeHealth, error)
|
|
DeleteAllEdgeHealth(ctx context.Context) (int64, error)
|
|
DeleteEdgeHealthBefore(ctx context.Context, cutoff time.Time) (int64, error)
|
|
|
|
InsertNodeObservationFrps(ctx context.Context, record *OpenFlareNodeObservationFrps) error
|
|
ListNodeObservationFrps(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrps, error)
|
|
DeleteAllNodeObservationFrps(ctx context.Context) (int64, error)
|
|
DeleteNodeObservationFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error)
|
|
|
|
InsertNodeObservationFrpc(ctx context.Context, record *OpenFlareNodeObservationFrpc) error
|
|
ListNodeObservationFrpc(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrpc, error)
|
|
DeleteAllNodeObservationFrpc(ctx context.Context) (int64, error)
|
|
DeleteNodeObservationFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error)
|
|
}
|
|
|
|
var (
|
|
observabilityStoreMu sync.RWMutex
|
|
observabilityStoreHolder observabilityStore
|
|
)
|
|
|
|
func currentObservabilityStore() observabilityStore {
|
|
observabilityStoreMu.RLock()
|
|
defer observabilityStoreMu.RUnlock()
|
|
if observabilityStoreHolder != nil {
|
|
return observabilityStoreHolder
|
|
}
|
|
return clickhouseObservabilityStore{}
|
|
}
|
|
|
|
// SetObservabilityStoreForTest swaps the observability store implementation for unit tests.
|
|
func SetObservabilityStoreForTest(store observabilityStore) func() {
|
|
observabilityStoreMu.Lock()
|
|
previous := observabilityStoreHolder
|
|
observabilityStoreHolder = store
|
|
observabilityStoreMu.Unlock()
|
|
return func() {
|
|
observabilityStoreMu.Lock()
|
|
observabilityStoreHolder = previous
|
|
observabilityStoreMu.Unlock()
|
|
}
|
|
}
|
|
|
|
// NewMemoryObservabilityStore returns an in-memory observability store for unit tests.
|
|
func NewMemoryObservabilityStore() observabilityStore {
|
|
return &memoryObservabilityStore{}
|
|
}
|
|
|
|
type clickhouseObservabilityStore struct{}
|
|
|
|
func (clickhouseObservabilityStore) InsertMetricSnapshot(_ context.Context, record *OpenFlareMetricSnapshot) error {
|
|
if record == nil {
|
|
return nil
|
|
}
|
|
if hook := currentObservabilityInsertHooks().QueueMetricSnapshot; hook != nil {
|
|
hook(toAnalyticsNodeMetricSnapshot(record))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) ListMetricSnapshots(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareMetricSnapshot, error) {
|
|
rows, err := analyticsrepo.ListNodeMetricSnapshots(ctx, toNodeObservabilityFilter(nodeID, since, limit))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return fromAnalyticsNodeMetricSnapshots(rows), nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteAllMetricSnapshots(ctx context.Context) (int64, error) {
|
|
return analyticsrepo.DeleteAllNodeMetricSnapshots(ctx)
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
|
|
return analyticsrepo.DeleteNodeMetricSnapshotsBefore(ctx, cutoff)
|
|
}
|
|
|
|
const edgeHealthStatusUnknown = "unknown"
|
|
|
|
func normalizeEdgeHealthStatus(status string) string {
|
|
status = strings.TrimSpace(status)
|
|
if status == "" {
|
|
return edgeHealthStatusUnknown
|
|
}
|
|
return status
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) InsertEdgeHealth(_ context.Context, record *OpenFlareEdgeHealth) error {
|
|
if record == nil {
|
|
return nil
|
|
}
|
|
if hook := currentObservabilityInsertHooks().QueueEdgeHealth; hook != nil {
|
|
hook(toAnalyticsNodeEdgeHealth(record))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) ListEdgeHealth(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareEdgeHealth, error) {
|
|
rows, err := analyticsrepo.ListNodeEdgeHealth(ctx, toNodeObservabilityFilter(nodeID, since, limit))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return fromAnalyticsNodeEdgeHealth(rows), nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteAllEdgeHealth(ctx context.Context) (int64, error) {
|
|
return analyticsrepo.DeleteAllNodeEdgeHealth(ctx)
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteEdgeHealthBefore(ctx context.Context, cutoff time.Time) (int64, error) {
|
|
return analyticsrepo.DeleteNodeEdgeHealthBefore(ctx, cutoff)
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) InsertNodeObservationFrps(_ context.Context, record *OpenFlareNodeObservationFrps) error {
|
|
if record == nil {
|
|
return nil
|
|
}
|
|
if hook := currentObservabilityInsertHooks().QueueFrpsObservation; hook != nil {
|
|
hook(toAnalyticsNodeObsFrps(record))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) ListNodeObservationFrps(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrps, error) {
|
|
rows, err := analyticsrepo.ListNodeObsFrps(ctx, toNodeObservabilityFilter(nodeID, since, limit))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return fromAnalyticsNodeObsFrps(rows), nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteAllNodeObservationFrps(ctx context.Context) (int64, error) {
|
|
return analyticsrepo.DeleteAllNodeObsFrps(ctx)
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteNodeObservationFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
|
|
return analyticsrepo.DeleteNodeObsFrpsBefore(ctx, cutoff)
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) InsertNodeObservationFrpc(_ context.Context, record *OpenFlareNodeObservationFrpc) error {
|
|
if record == nil {
|
|
return nil
|
|
}
|
|
if hook := currentObservabilityInsertHooks().QueueFrpcObservation; hook != nil {
|
|
hook(toAnalyticsNodeObsFrpc(record))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) ListNodeObservationFrpc(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrpc, error) {
|
|
rows, err := analyticsrepo.ListNodeObsFrpc(ctx, toNodeObservabilityFilter(nodeID, since, limit))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return fromAnalyticsNodeObsFrpc(rows), nil
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteAllNodeObservationFrpc(ctx context.Context) (int64, error) {
|
|
return analyticsrepo.DeleteAllNodeObsFrpc(ctx)
|
|
}
|
|
|
|
func (clickhouseObservabilityStore) DeleteNodeObservationFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error) {
|
|
return analyticsrepo.DeleteNodeObsFrpcBefore(ctx, cutoff)
|
|
}
|
|
|
|
func toNodeObservabilityFilter(nodeID string, since time.Time, limit int) analyticsrepo.NodeObservabilityFilter {
|
|
return analyticsrepo.NodeObservabilityFilter{
|
|
NodeID: nodeID,
|
|
Since: since,
|
|
Limit: limit,
|
|
}
|
|
}
|
|
|
|
func toAnalyticsNodeMetricSnapshot(record *OpenFlareMetricSnapshot) analyticsmodel.NodeMetricSnapshot {
|
|
return analyticsmodel.NodeMetricSnapshot{
|
|
ID: uint64(record.ID),
|
|
NodeID: record.NodeID,
|
|
CapturedAt: record.CapturedAt,
|
|
CPUUsagePercent: record.CPUUsagePercent,
|
|
MemoryUsedBytes: record.MemoryUsedBytes,
|
|
MemoryTotalBytes: record.MemoryTotalBytes,
|
|
StorageUsedBytes: record.StorageUsedBytes,
|
|
StorageTotalBytes: record.StorageTotalBytes,
|
|
DiskReadBytes: record.DiskReadBytes,
|
|
DiskWriteBytes: record.DiskWriteBytes,
|
|
NetworkRxBytes: record.NetworkRxBytes,
|
|
NetworkTxBytes: record.NetworkTxBytes,
|
|
CreatedAt: record.CreatedAt,
|
|
}
|
|
}
|
|
|
|
func fromAnalyticsNodeMetricSnapshots(rows []analyticsmodel.NodeMetricSnapshot) []*OpenFlareMetricSnapshot {
|
|
result := make([]*OpenFlareMetricSnapshot, len(rows))
|
|
for index, row := range rows {
|
|
result[index] = &OpenFlareMetricSnapshot{
|
|
ID: uint(row.ID),
|
|
NodeID: row.NodeID,
|
|
CapturedAt: row.CapturedAt,
|
|
CPUUsagePercent: row.CPUUsagePercent,
|
|
MemoryUsedBytes: row.MemoryUsedBytes,
|
|
MemoryTotalBytes: row.MemoryTotalBytes,
|
|
StorageUsedBytes: row.StorageUsedBytes,
|
|
StorageTotalBytes: row.StorageTotalBytes,
|
|
DiskReadBytes: row.DiskReadBytes,
|
|
DiskWriteBytes: row.DiskWriteBytes,
|
|
NetworkRxBytes: row.NetworkRxBytes,
|
|
NetworkTxBytes: row.NetworkTxBytes,
|
|
CreatedAt: row.CreatedAt,
|
|
}
|
|
}
|
|
return result
|
|
}
|
|
|
|
func toAnalyticsNodeEdgeHealth(record *OpenFlareEdgeHealth) analyticsmodel.NodeEdgeHealth {
|
|
return analyticsmodel.NodeEdgeHealth{
|
|
ID: uint64(record.ID),
|
|
NodeID: record.NodeID,
|
|
CapturedAt: record.CapturedAt,
|
|
Status: normalizeEdgeHealthStatus(record.Status),
|
|
Connections: record.Connections,
|
|
CreatedAt: record.CreatedAt,
|
|
}
|
|
}
|
|
|
|
func fromAnalyticsNodeEdgeHealth(rows []analyticsmodel.NodeEdgeHealth) []*OpenFlareEdgeHealth {
|
|
result := make([]*OpenFlareEdgeHealth, len(rows))
|
|
for index, row := range rows {
|
|
result[index] = &OpenFlareEdgeHealth{
|
|
ID: uint(row.ID),
|
|
NodeID: row.NodeID,
|
|
CapturedAt: row.CapturedAt,
|
|
Status: normalizeEdgeHealthStatus(row.Status),
|
|
Connections: row.Connections,
|
|
CreatedAt: row.CreatedAt,
|
|
}
|
|
}
|
|
return result
|
|
}
|
|
|
|
func toAnalyticsNodeObsFrps(record *OpenFlareNodeObservationFrps) analyticsmodel.NodeObsFrps {
|
|
return analyticsmodel.NodeObsFrps{
|
|
ID: uint64(record.ID),
|
|
NodeID: record.NodeID,
|
|
CapturedAt: record.CapturedAt,
|
|
FrpsConnections: openFlareObservabilityIntToInt32(record.FrpsConnections),
|
|
FrpsProxyCount: openFlareObservabilityIntToInt32(record.FrpsProxyCount),
|
|
FrpsClientCount: openFlareObservabilityIntToInt32(record.FrpsClientCount),
|
|
FrpsProxies: record.FrpsProxies,
|
|
CreatedAt: record.CreatedAt,
|
|
}
|
|
}
|
|
|
|
func fromAnalyticsNodeObsFrps(rows []analyticsmodel.NodeObsFrps) []*OpenFlareNodeObservationFrps {
|
|
result := make([]*OpenFlareNodeObservationFrps, len(rows))
|
|
for index, row := range rows {
|
|
result[index] = &OpenFlareNodeObservationFrps{
|
|
ID: uint(row.ID),
|
|
NodeID: row.NodeID,
|
|
CapturedAt: row.CapturedAt,
|
|
FrpsConnections: int(row.FrpsConnections),
|
|
FrpsProxyCount: int(row.FrpsProxyCount),
|
|
FrpsClientCount: int(row.FrpsClientCount),
|
|
FrpsProxies: row.FrpsProxies,
|
|
CreatedAt: row.CreatedAt,
|
|
}
|
|
}
|
|
return result
|
|
}
|
|
|
|
func toAnalyticsNodeObsFrpc(record *OpenFlareNodeObservationFrpc) analyticsmodel.NodeObsFrpc {
|
|
return analyticsmodel.NodeObsFrpc{
|
|
ID: uint64(record.ID),
|
|
NodeID: record.NodeID,
|
|
CapturedAt: record.CapturedAt,
|
|
TunnelStatus: record.TunnelStatus,
|
|
ConnectedRelaysCount: openFlareObservabilityIntToInt32(record.ConnectedRelaysCount),
|
|
CreatedAt: record.CreatedAt,
|
|
}
|
|
}
|
|
|
|
func openFlareObservabilityIntToInt32(value int) int32 {
|
|
switch {
|
|
case value > math.MaxInt32:
|
|
return math.MaxInt32
|
|
case value < math.MinInt32:
|
|
return math.MinInt32
|
|
default:
|
|
return int32(value)
|
|
}
|
|
}
|
|
|
|
func fromAnalyticsNodeObsFrpc(rows []analyticsmodel.NodeObsFrpc) []*OpenFlareNodeObservationFrpc {
|
|
result := make([]*OpenFlareNodeObservationFrpc, len(rows))
|
|
for index, row := range rows {
|
|
result[index] = &OpenFlareNodeObservationFrpc{
|
|
ID: uint(row.ID),
|
|
NodeID: row.NodeID,
|
|
CapturedAt: row.CapturedAt,
|
|
TunnelStatus: row.TunnelStatus,
|
|
ConnectedRelaysCount: int(row.ConnectedRelaysCount),
|
|
CreatedAt: row.CreatedAt,
|
|
}
|
|
}
|
|
return result
|
|
}
|