diff --git a/docker-compose.yaml b/docker-compose.yaml index f21cb6c2..6b3961c4 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -76,9 +76,9 @@ services: image: clickhouse/clickhouse-server:25.3-alpine restart: unless-stopped environment: - CLICKHOUSE_DB: ${CLICKHOUSE_DB:-openflare} - CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} - CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-123456} + CLICKHOUSE_DB: openflare + CLICKHOUSE_USER: default + CLICKHOUSE_PASSWORD: 123456 CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: 1 TZ: ${TZ:-Asia/Shanghai} ports: diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 0986bc1b..7e100119 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -32,6 +32,10 @@ sidebar: false - 修复访问日志 ClickHouse 聚合查询因 `trim(x) AS x` 别名与表列同名导致总览看板地域分布及 IP 统计失败的问题。 +### 变更 + +- 将节点可观测时序表(`of_node_metric_snapshots`、`of_node_request_reports`、`of_node_obs_openresty`、`of_node_obs_frps`、`of_node_obs_frpc`)从 PostgreSQL/SQLite 迁移至 ClickHouse,主库迁移 `202606200005` 删除对应 PG 表。 + ## [v2.3.4] - 2026-06-17 ### 变更 diff --git a/internal/apps/openflare/agent/observability.go b/internal/apps/openflare/agent/observability.go index 8bd4c9f5..d210fe7f 100644 --- a/internal/apps/openflare/agent/observability.go +++ b/internal/apps/openflare/agent/observability.go @@ -59,18 +59,6 @@ func PersistHeartbeatObservability(ctx context.Context, nodeID string, payload N if err := persistNodeSystemProfile(tx, nodeID, payload.Profile, reportedAt); err != nil { return err } - if err := persistBufferedObservability(tx, nodeID, payload.BufferedObservability, reportedAt); err != nil { - return err - } - if err := persistNodeMetricSnapshot(tx, nodeID, payload.Snapshot, reportedAt); err != nil { - return err - } - if err := persistNodeOpenrestyObservation(tx, nodeID, payload.OpenrestyObservation, reportedAt); err != nil { - return err - } - if err := persistNodeTrafficReport(tx, nodeID, payload.TrafficReport, reportedAt); err != nil { - return err - } if payload.HealthEvents != nil { if err := reconcileNodeHealthEvents(tx, nodeID, payload.HealthEvents, reportedAt); err != nil { return err @@ -82,23 +70,35 @@ func PersistHeartbeatObservability(ctx context.Context, nodeID string, payload N return } + if err := persistBufferedObservability(ctx, nodeID, payload.BufferedObservability, reportedAt); err != nil { + zap.L().Error("persist buffered observability failed", zap.String("node_id", nodeID), zap.Error(err)) + } + if err := persistNodeMetricSnapshot(ctx, nodeID, payload.Snapshot, reportedAt); err != nil { + zap.L().Error("persist metric snapshot failed", zap.String("node_id", nodeID), zap.Error(err)) + } + if err := persistNodeOpenrestyObservation(ctx, nodeID, payload.OpenrestyObservation, reportedAt); err != nil { + zap.L().Error("persist openresty observation failed", zap.String("node_id", nodeID), zap.Error(err)) + } + if err := persistNodeTrafficReport(ctx, nodeID, payload.TrafficReport, reportedAt); err != nil { + zap.L().Error("persist traffic report failed", zap.String("node_id", nodeID), zap.Error(err)) + } + if err := persistNodeAccessLogs(ctx, nodeID, accessLogRecords, reportedAt); err != nil { zap.L().Error("persist heartbeat access logs failed", zap.String("node_id", nodeID), zap.Error(err)) } } -func persistBufferedObservability(tx *gorm.DB, nodeID string, records []BufferedObservabilityRecord, reportedAt time.Time) error { +func persistBufferedObservability(ctx context.Context, nodeID string, records []BufferedObservabilityRecord, reportedAt time.Time) error { for _, record := range records { - if err := persistNodeMetricSnapshot(tx, nodeID, record.Snapshot, reportedAt); err != nil { + if err := persistNodeMetricSnapshot(ctx, nodeID, record.Snapshot, reportedAt); err != nil { return err } - if err := persistNodeOpenrestyObservation(tx, nodeID, record.OpenrestyObservation, reportedAt); err != nil { + if err := persistNodeOpenrestyObservation(ctx, nodeID, record.OpenrestyObservation, reportedAt); err != nil { return err } - if err := persistNodeTrafficReport(tx, nodeID, record.TrafficReport, reportedAt); err != nil { + if err := persistNodeTrafficReport(ctx, nodeID, record.TrafficReport, reportedAt); err != nil { return err } - } return nil } @@ -140,7 +140,7 @@ func persistNodeSystemProfile(tx *gorm.DB, nodeID string, profile *NodeSystemPro }).Create(record).Error } -func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *NodeMetricSnapshot, reportedAt time.Time) error { +func persistNodeMetricSnapshot(ctx context.Context, nodeID string, snapshot *NodeMetricSnapshot, reportedAt time.Time) error { if snapshot == nil { return nil } @@ -157,17 +157,10 @@ func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *NodeMetricS NetworkRxBytes: snapshot.NetworkRxBytes, NetworkTxBytes: snapshot.NetworkTxBytes, } - exists, err := metricSnapshotExists(tx, nodeID, record.CapturedAt) - if err != nil { - return err - } - if exists { - return nil - } - return tx.Create(record).Error + return model.InsertOpenFlareMetricSnapshot(ctx, record) } -func persistNodeOpenrestyObservation(tx *gorm.DB, nodeID string, obs *NodeOpenrestyObservation, reportedAt time.Time) error { +func persistNodeOpenrestyObservation(ctx context.Context, nodeID string, obs *NodeOpenrestyObservation, reportedAt time.Time) error { if obs == nil { return nil } @@ -178,10 +171,10 @@ func persistNodeOpenrestyObservation(tx *gorm.DB, nodeID string, obs *NodeOpenre OpenrestyTxBytes: obs.OpenrestyTxBytes, OpenrestyConnections: obs.OpenrestyConnections, } - return tx.Create(record).Error + return model.InsertOpenFlareNodeObservationOpenresty(ctx, record) } -func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *NodeTrafficReport, reportedAt time.Time) error { +func persistNodeTrafficReport(ctx context.Context, nodeID string, report *NodeTrafficReport, reportedAt time.Time) error { if report == nil { return nil } @@ -199,14 +192,7 @@ func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *NodeTrafficRep TopDomainsJSON: marshalJSON(report.TopDomains), SourceCountriesJSON: marshalJSON(report.SourceCountries), } - exists, err := requestReportExists(tx, nodeID, record.WindowStartedAt, record.WindowEndedAt) - if err != nil { - return err - } - if exists { - return nil - } - return tx.Create(record).Error + return model.InsertOpenFlareRequestReport(ctx, record) } func buildNodeAccessLogRecords(nodeID string, direct []NodeAccessLog, buffered []BufferedObservabilityRecord, reportedAt time.Time) ([]*model.OpenFlareAccessLog, error) { @@ -357,28 +343,6 @@ func ReconcileScopedNodeHealthEvents(tx *gorm.DB, nodeID string, events []NodeHe return nil } -func metricSnapshotExists(tx *gorm.DB, nodeID string, capturedAt time.Time) (bool, error) { - var count int64 - if err := tx.Model(&model.OpenFlareMetricSnapshot{}). - Where("node_id = ? AND captured_at = ?", nodeID, capturedAt). - Limit(1). - Count(&count).Error; err != nil { - return false, err - } - return count > 0, nil -} - -func requestReportExists(tx *gorm.DB, nodeID string, windowStartedAt, windowEndedAt time.Time) (bool, error) { - var count int64 - if err := tx.Model(&model.OpenFlareRequestReport{}). - Where("node_id = ? AND window_started_at = ? AND window_ended_at = ?", nodeID, windowStartedAt, windowEndedAt). - Limit(1). - Count(&count).Error; err != nil { - return false, err - } - return count > 0, nil -} - func normalizeHealthEventType(eventType string) string { eventType = strings.TrimSpace(strings.ToLower(eventType)) eventType = strings.ReplaceAll(eventType, " ", "_") diff --git a/internal/apps/openflare/async_tasks_test.go b/internal/apps/openflare/async_tasks_test.go index acbc2fb2..5bceb3bc 100644 --- a/internal/apps/openflare/async_tasks_test.go +++ b/internal/apps/openflare/async_tasks_test.go @@ -36,7 +36,9 @@ func TestDatabaseAutoCleanupHandlerDeletesRowsWhenEnabled(t *testing.T) { require.NoError(t, err) db.SetDB(sqliteDB) resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore()) + resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore()) t.Cleanup(func() { + resetObservabilityStore() resetAccessLogStore() db.SetDB(nil) }) @@ -81,4 +83,4 @@ func TestUptimeKumaSyncHandlerSkipsWhenDisabled(t *testing.T) { require.NoError(t, err) require.NotNil(t, result) assert.Contains(t, result.Message, "未启用") -} \ No newline at end of file +} diff --git a/internal/apps/openflare/dashboard/logics_test.go b/internal/apps/openflare/dashboard/logics_test.go index 8eed2448..be27379f 100644 --- a/internal/apps/openflare/dashboard/logics_test.go +++ b/internal/apps/openflare/dashboard/logics_test.go @@ -25,7 +25,9 @@ func setupDashboardTestDB(t *testing.T) func() { db.SetDB(sqliteDB) resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore()) + resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore()) return func() { + resetObservabilityStore() resetAccessLogStore() db.SetDB(nil) } diff --git a/internal/apps/openflare/integration/agent_protocol_test.go b/internal/apps/openflare/integration/agent_protocol_test.go index 88c66023..82a85105 100644 --- a/internal/apps/openflare/integration/agent_protocol_test.go +++ b/internal/apps/openflare/integration/agent_protocol_test.go @@ -46,21 +46,22 @@ func setupProtocolTestEnv(t *testing.T) (*gin.Engine, func()) { &model.OpenFlareOption{}, &model.OpenFlareApplyLog{}, &model.OpenFlareNodeSystemProfile{}, - &model.OpenFlareMetricSnapshot{}, &model.OpenFlareHealthEvent{}, - &model.OpenFlareNodeObservationFrps{}, - &model.OpenFlareNodeObservationFrpc{}, &configVersionRecord{}, )) db.SetDB(sqliteDB) option.ResetInitializationForTest() agent.ResetAuthCacheForTest() + resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore()) + resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore()) engine := testhelper.NewTestGinEngine() mountOpenFlareTestRoutes(engine) cleanup := func() { + resetObservabilityStore() + resetAccessLogStore() db.SetDB(nil) option.ResetInitializationForTest() agent.ResetAuthCacheForTest() @@ -210,4 +211,4 @@ func TestAgentRelayFlaredProtocol(t *testing.T) { assert.Equal(t, "online", stored.Status) assert.Equal(t, "20260618-001", stored.CurrentVersion) }) -} \ No newline at end of file +} diff --git a/internal/apps/openflare/node/logics_test.go b/internal/apps/openflare/node/logics_test.go index 596bbf95..7f1731f9 100644 --- a/internal/apps/openflare/node/logics_test.go +++ b/internal/apps/openflare/node/logics_test.go @@ -20,6 +20,12 @@ import ( "gorm.io/gorm" ) +func setReleaseHTTPClientForTest(client *http.Client) *http.Client { + previous := releaseHTTPClient + releaseHTTPClient = client + return previous +} + func setupNodeTestDB(t *testing.T) func() { t.Helper() @@ -36,8 +42,10 @@ func setupNodeTestDB(t *testing.T) func() { db.SetDB(sqliteDB) option.ResetInitializationForTest() resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore()) + resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore()) return func() { + resetObservabilityStore() resetAccessLogStore() db.SetDB(nil) option.ResetInitializationForTest() diff --git a/internal/apps/openflare/relay/logics_test.go b/internal/apps/openflare/relay/logics_test.go index eb0f6caa..b6e6f9ae 100644 --- a/internal/apps/openflare/relay/logics_test.go +++ b/internal/apps/openflare/relay/logics_test.go @@ -38,8 +38,10 @@ func setupRelayTestDB(t *testing.T) func() { db.SetDB(sqliteDB) option.ResetInitializationForTest() agent.ResetAuthCacheForTest() + resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore()) return func() { + resetObservabilityStore() db.SetDB(nil) option.ResetInitializationForTest() agent.ResetAuthCacheForTest() diff --git a/internal/apps/openflare/relay/observability.go b/internal/apps/openflare/relay/observability.go index 92f15efb..2253b778 100644 --- a/internal/apps/openflare/relay/observability.go +++ b/internal/apps/openflare/relay/observability.go @@ -51,10 +51,6 @@ func persistRelayHeartbeatObservability(ctx context.Context, nodeID string, payl HealthEvents: payload.HealthEvents, }, reportedAt) - conn := db.DB(ctx) - if conn == nil { - return - } frpsObs := &model.OpenFlareNodeObservationFrps{ NodeID: nodeID, CapturedAt: reportedAt, @@ -63,7 +59,7 @@ func persistRelayHeartbeatObservability(ctx context.Context, nodeID string, payl FrpsClientCount: payload.FrpsClientCount, FrpsProxies: agent.MarshalJSON(payload.FrpsProxies), } - if err := conn.Create(frpsObs).Error; err != nil { + if err := model.InsertOpenFlareNodeObservationFrps(ctx, frpsObs); err != nil { zap.L().Error("persist relay frps observation failed", zap.String("node_id", nodeID), zap.Error(err)) } } diff --git a/internal/apps/openflare/tasks/database_cleanup_test.go b/internal/apps/openflare/tasks/database_cleanup_test.go index cfaf405e..04d28ea8 100644 --- a/internal/apps/openflare/tasks/database_cleanup_test.go +++ b/internal/apps/openflare/tasks/database_cleanup_test.go @@ -23,13 +23,11 @@ func setupDatabaseCleanupTestDB(t *testing.T) context.Context { DisableForeignKeyConstraintWhenMigrating: true, }) require.NoError(t, err) - require.NoError(t, sqliteDB.AutoMigrate( - &model.OpenFlareMetricSnapshot{}, - &model.OpenFlareRequestReport{}, - )) db.SetDB(sqliteDB) resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore()) + resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore()) t.Cleanup(func() { + resetObservabilityStore() resetAccessLogStore() db.SetDB(nil) }) @@ -40,16 +38,16 @@ func TestCleanupDatabaseObservabilityDeletesTargetedRows(t *testing.T) { ctx := setupDatabaseCleanupTestDB(t) now := time.Now().UTC() - require.NoError(t, db.DB(ctx).Create(&model.OpenFlareMetricSnapshot{ + require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ NodeID: "node-a", CapturedAt: now.Add(-10 * 24 * time.Hour), CPUUsagePercent: 10, - }).Error) - require.NoError(t, db.DB(ctx).Create(&model.OpenFlareMetricSnapshot{ + })) + require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ NodeID: "node-a", CapturedAt: now.Add(-12 * time.Hour), CPUUsagePercent: 20, - }).Error) + })) retentionDays := 7 result, err := CleanupDatabaseObservability(ctx, DatabaseCleanupInput{ @@ -113,17 +111,17 @@ func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T) Path: "/access", StatusCode: 200, }})) - require.NoError(t, db.DB(ctx).Create(&model.OpenFlareMetricSnapshot{ + require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ NodeID: "node-a", CapturedAt: now.Add(-48 * time.Hour), CPUUsagePercent: 10, - }).Error) - require.NoError(t, db.DB(ctx).Create(&model.OpenFlareRequestReport{ + })) + require.NoError(t, model.InsertOpenFlareRequestReport(ctx, &model.OpenFlareRequestReport{ NodeID: "node-a", WindowStartedAt: now.Add(-49 * time.Hour), WindowEndedAt: now.Add(-48 * time.Hour), RequestCount: 15, - }).Error) + })) previousEnabled := model.DatabaseAutoCleanupEnabled previousRetentionDays := model.DatabaseAutoCleanupRetentionDays diff --git a/internal/db/migrator/goose/clickhouse/202606200002_create_node_observability_tables.sql b/internal/db/migrator/goose/clickhouse/202606200002_create_node_observability_tables.sql new file mode 100644 index 00000000..3069ac39 --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202606200002_create_node_observability_tables.sql @@ -0,0 +1,92 @@ +-- +goose Up +CREATE TABLE IF NOT EXISTS of_node_metric_snapshots +( + id UInt64, + node_id String, + captured_at DateTime64(3, 'UTC'), + cpu_usage_percent Float64, + memory_used_bytes Int64, + memory_total_bytes Int64, + storage_used_bytes Int64, + storage_total_bytes Int64, + disk_read_bytes Int64, + disk_write_bytes Int64, + network_rx_bytes Int64, + network_tx_bytes Int64, + created_at DateTime64(3, 'UTC') +) +ENGINE = MergeTree() +PARTITION BY toYYYYMM(captured_at) +ORDER BY (node_id, captured_at, id) +SETTINGS index_granularity = 8192; + +CREATE TABLE IF NOT EXISTS of_node_request_reports +( + id UInt64, + node_id String, + window_started_at DateTime64(3, 'UTC'), + window_ended_at DateTime64(3, 'UTC'), + request_count Int64, + error_count Int64, + unique_visitor_count Int64, + status_codes_json String, + top_domains_json String, + source_countries_json String, + created_at DateTime64(3, 'UTC') +) +ENGINE = MergeTree() +PARTITION BY toYYYYMM(window_ended_at) +ORDER BY (node_id, window_ended_at, window_started_at, id) +SETTINGS index_granularity = 8192; + +CREATE TABLE IF NOT EXISTS of_node_obs_openresty +( + id UInt64, + node_id String, + captured_at DateTime64(3, 'UTC'), + openresty_rx_bytes Int64, + openresty_tx_bytes Int64, + openresty_connections Int64, + created_at DateTime64(3, 'UTC') +) +ENGINE = MergeTree() +PARTITION BY toYYYYMM(captured_at) +ORDER BY (node_id, captured_at, id) +SETTINGS index_granularity = 8192; + +CREATE TABLE IF NOT EXISTS of_node_obs_frps +( + id UInt64, + node_id String, + captured_at DateTime64(3, 'UTC'), + frps_connections Int32, + frps_proxy_count Int32, + frps_client_count Int32, + frps_proxies String, + created_at DateTime64(3, 'UTC') +) +ENGINE = MergeTree() +PARTITION BY toYYYYMM(captured_at) +ORDER BY (node_id, captured_at, id) +SETTINGS index_granularity = 8192; + +CREATE TABLE IF NOT EXISTS of_node_obs_frpc +( + id UInt64, + node_id String, + captured_at DateTime64(3, 'UTC'), + tunnel_status String, + connected_relays_count Int32, + created_at DateTime64(3, 'UTC') +) +ENGINE = MergeTree() +PARTITION BY toYYYYMM(captured_at) +ORDER BY (node_id, captured_at, id) +SETTINGS index_granularity = 8192; + +-- +goose Down +DROP TABLE IF EXISTS of_node_obs_frpc; +DROP TABLE IF EXISTS of_node_obs_frps; +DROP TABLE IF EXISTS of_node_obs_openresty; +DROP TABLE IF EXISTS of_node_request_reports; +DROP TABLE IF EXISTS of_node_metric_snapshots; \ No newline at end of file diff --git a/internal/db/migrator/goose/postgres/202606200005_drop_of_node_observability_timeseries.sql b/internal/db/migrator/goose/postgres/202606200005_drop_of_node_observability_timeseries.sql new file mode 100644 index 00000000..0309fe57 --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606200005_drop_of_node_observability_timeseries.sql @@ -0,0 +1,98 @@ +-- +goose Up +DROP INDEX IF EXISTS idx_of_node_obs_frpc_captured_at; +DROP INDEX IF EXISTS idx_of_node_obs_frpc_node_id; +DROP TABLE IF EXISTS of_node_obs_frpc; + +DROP INDEX IF EXISTS idx_of_node_obs_frps_captured_at; +DROP INDEX IF EXISTS idx_of_node_obs_frps_node_id; +DROP TABLE IF EXISTS of_node_obs_frps; + +DROP INDEX IF EXISTS idx_of_node_obs_openresty_captured_at; +DROP INDEX IF EXISTS idx_of_node_obs_openresty_node_id; +DROP TABLE IF EXISTS of_node_obs_openresty; + +DROP INDEX IF EXISTS idx_of_node_request_reports_window_ended_at; +DROP INDEX IF EXISTS idx_of_node_request_reports_window_started_at; +DROP INDEX IF EXISTS idx_of_node_request_reports_node_id; +DROP TABLE IF EXISTS of_node_request_reports; + +DROP INDEX IF EXISTS idx_of_node_metric_snapshots_captured_at; +DROP INDEX IF EXISTS idx_of_node_metric_snapshots_node_id; +DROP TABLE IF EXISTS of_node_metric_snapshots; + +-- +goose Down +CREATE TABLE of_node_metric_snapshots ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + cpu_usage_percent DOUBLE PRECISION NOT NULL DEFAULT 0, + memory_used_bytes BIGINT NOT NULL DEFAULT 0, + memory_total_bytes BIGINT NOT NULL DEFAULT 0, + storage_used_bytes BIGINT NOT NULL DEFAULT 0, + storage_total_bytes BIGINT NOT NULL DEFAULT 0, + disk_read_bytes BIGINT NOT NULL DEFAULT 0, + disk_write_bytes BIGINT NOT NULL DEFAULT 0, + network_rx_bytes BIGINT NOT NULL DEFAULT 0, + network_tx_bytes BIGINT NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_metric_snapshots_node_id ON of_node_metric_snapshots (node_id); +CREATE INDEX idx_of_node_metric_snapshots_captured_at ON of_node_metric_snapshots (captured_at); + +CREATE TABLE of_node_request_reports ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + window_started_at TIMESTAMPTZ NOT NULL, + window_ended_at TIMESTAMPTZ NOT NULL, + request_count BIGINT NOT NULL DEFAULT 0, + error_count BIGINT NOT NULL DEFAULT 0, + unique_visitor_count BIGINT NOT NULL DEFAULT 0, + status_codes_json TEXT, + top_domains_json TEXT, + source_countries_json TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_request_reports_node_id ON of_node_request_reports (node_id); +CREATE INDEX idx_of_node_request_reports_window_started_at ON of_node_request_reports (window_started_at); +CREATE INDEX idx_of_node_request_reports_window_ended_at ON of_node_request_reports (window_ended_at); + +CREATE TABLE of_node_obs_openresty ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + openresty_rx_bytes BIGINT NOT NULL DEFAULT 0, + openresty_tx_bytes BIGINT NOT NULL DEFAULT 0, + openresty_connections BIGINT NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_obs_openresty_node_id ON of_node_obs_openresty (node_id); +CREATE INDEX idx_of_node_obs_openresty_captured_at ON of_node_obs_openresty (captured_at); + +CREATE TABLE of_node_obs_frps ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + frps_connections INTEGER NOT NULL DEFAULT 0, + frps_proxy_count INTEGER NOT NULL DEFAULT 0, + frps_client_count INTEGER NOT NULL DEFAULT 0, + frps_proxies TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_obs_frps_node_id ON of_node_obs_frps (node_id); +CREATE INDEX idx_of_node_obs_frps_captured_at ON of_node_obs_frps (captured_at); + +CREATE TABLE of_node_obs_frpc ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + tunnel_status VARCHAR(16) NOT NULL DEFAULT '', + connected_relays_count INTEGER NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_obs_frpc_node_id ON of_node_obs_frpc (node_id); +CREATE INDEX idx_of_node_obs_frpc_captured_at ON of_node_obs_frpc (captured_at); \ No newline at end of file diff --git a/internal/db/migrator/goose/sqlite/202606200005_drop_of_node_observability_timeseries.sql b/internal/db/migrator/goose/sqlite/202606200005_drop_of_node_observability_timeseries.sql new file mode 100644 index 00000000..6aa7fe7b --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606200005_drop_of_node_observability_timeseries.sql @@ -0,0 +1,98 @@ +-- +goose Up +DROP INDEX IF EXISTS idx_of_node_obs_frpc_captured_at; +DROP INDEX IF EXISTS idx_of_node_obs_frpc_node_id; +DROP TABLE IF EXISTS of_node_obs_frpc; + +DROP INDEX IF EXISTS idx_of_node_obs_frps_captured_at; +DROP INDEX IF EXISTS idx_of_node_obs_frps_node_id; +DROP TABLE IF EXISTS of_node_obs_frps; + +DROP INDEX IF EXISTS idx_of_node_obs_openresty_captured_at; +DROP INDEX IF EXISTS idx_of_node_obs_openresty_node_id; +DROP TABLE IF EXISTS of_node_obs_openresty; + +DROP INDEX IF EXISTS idx_of_node_request_reports_window_ended_at; +DROP INDEX IF EXISTS idx_of_node_request_reports_window_started_at; +DROP INDEX IF EXISTS idx_of_node_request_reports_node_id; +DROP TABLE IF EXISTS of_node_request_reports; + +DROP INDEX IF EXISTS idx_of_node_metric_snapshots_captured_at; +DROP INDEX IF EXISTS idx_of_node_metric_snapshots_node_id; +DROP TABLE IF EXISTS of_node_metric_snapshots; + +-- +goose Down +CREATE TABLE IF NOT EXISTS of_node_metric_snapshots ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + cpu_usage_percent REAL NOT NULL DEFAULT 0, + memory_used_bytes INTEGER NOT NULL DEFAULT 0, + memory_total_bytes INTEGER NOT NULL DEFAULT 0, + storage_used_bytes INTEGER NOT NULL DEFAULT 0, + storage_total_bytes INTEGER NOT NULL DEFAULT 0, + disk_read_bytes INTEGER NOT NULL DEFAULT 0, + disk_write_bytes INTEGER NOT NULL DEFAULT 0, + network_rx_bytes INTEGER NOT NULL DEFAULT 0, + network_tx_bytes INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_metric_snapshots_node_id ON of_node_metric_snapshots (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_metric_snapshots_captured_at ON of_node_metric_snapshots (captured_at); + +CREATE TABLE IF NOT EXISTS of_node_request_reports ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + window_started_at DATETIME NOT NULL, + window_ended_at DATETIME NOT NULL, + request_count INTEGER NOT NULL DEFAULT 0, + error_count INTEGER NOT NULL DEFAULT 0, + unique_visitor_count INTEGER NOT NULL DEFAULT 0, + status_codes_json TEXT, + top_domains_json TEXT, + source_countries_json TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_request_reports_node_id ON of_node_request_reports (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_request_reports_window_started_at ON of_node_request_reports (window_started_at); +CREATE INDEX IF NOT EXISTS idx_of_node_request_reports_window_ended_at ON of_node_request_reports (window_ended_at); + +CREATE TABLE IF NOT EXISTS of_node_obs_openresty ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + openresty_rx_bytes INTEGER NOT NULL DEFAULT 0, + openresty_tx_bytes INTEGER NOT NULL DEFAULT 0, + openresty_connections INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_obs_openresty_node_id ON of_node_obs_openresty (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_obs_openresty_captured_at ON of_node_obs_openresty (captured_at); + +CREATE TABLE IF NOT EXISTS of_node_obs_frps ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + frps_connections INTEGER NOT NULL DEFAULT 0, + frps_proxy_count INTEGER NOT NULL DEFAULT 0, + frps_client_count INTEGER NOT NULL DEFAULT 0, + frps_proxies TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_obs_frps_node_id ON of_node_obs_frps (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_obs_frps_captured_at ON of_node_obs_frps (captured_at); + +CREATE TABLE IF NOT EXISTS of_node_obs_frpc ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + tunnel_status TEXT NOT NULL DEFAULT '', + connected_relays_count INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_obs_frpc_node_id ON of_node_obs_frpc (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_obs_frpc_captured_at ON of_node_obs_frpc (captured_at); \ No newline at end of file diff --git a/internal/model/analytics/node_observability.go b/internal/model/analytics/node_observability.go new file mode 100644 index 00000000..86511bbf --- /dev/null +++ b/internal/model/analytics/node_observability.go @@ -0,0 +1,166 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "fmt" + "time" +) + +const ( + nodeMetricSnapshotTableName = "of_node_metric_snapshots" + nodeMetricSnapshotInsertColumns = "id, node_id, captured_at, cpu_usage_percent, memory_used_bytes, memory_total_bytes, storage_used_bytes, storage_total_bytes, disk_read_bytes, disk_write_bytes, network_rx_bytes, network_tx_bytes, created_at" + + nodeRequestReportTableName = "of_node_request_reports" + nodeRequestReportInsertColumns = "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" + + nodeObsOpenrestyTableName = "of_node_obs_openresty" + nodeObsOpenrestyInsertColumns = "id, node_id, captured_at, openresty_rx_bytes, openresty_tx_bytes, openresty_connections, created_at" + + nodeObsFrpsTableName = "of_node_obs_frps" + nodeObsFrpsInsertColumns = "id, node_id, captured_at, frps_connections, frps_proxy_count, frps_client_count, frps_proxies, created_at" + + nodeObsFrpcTableName = "of_node_obs_frpc" + nodeObsFrpcInsertColumns = "id, node_id, captured_at, tunnel_status, connected_relays_count, created_at" +) + +// NodeMetricSnapshot stores periodic node resource utilization metrics in ClickHouse. +type NodeMetricSnapshot struct { + ID uint64 `gorm:"column:id"` + NodeID string `gorm:"column:node_id"` + CapturedAt time.Time `gorm:"column:captured_at"` + CPUUsagePercent float64 `gorm:"column:cpu_usage_percent"` + MemoryUsedBytes int64 `gorm:"column:memory_used_bytes"` + MemoryTotalBytes int64 `gorm:"column:memory_total_bytes"` + StorageUsedBytes int64 `gorm:"column:storage_used_bytes"` + StorageTotalBytes int64 `gorm:"column:storage_total_bytes"` + DiskReadBytes int64 `gorm:"column:disk_read_bytes"` + DiskWriteBytes int64 `gorm:"column:disk_write_bytes"` + NetworkRxBytes int64 `gorm:"column:network_rx_bytes"` + NetworkTxBytes int64 `gorm:"column:network_tx_bytes"` + CreatedAt time.Time `gorm:"column:created_at"` +} + +// TableName returns the ClickHouse table name. +func (NodeMetricSnapshot) TableName() string { + return nodeMetricSnapshotTableName +} + +// InsertColumns returns comma-separated column names for batch insert. +func (NodeMetricSnapshot) InsertColumns() string { + return nodeMetricSnapshotInsertColumns +} + +// BatchInsertSQL returns the INSERT prefix used by native batch writers. +func (NodeMetricSnapshot) BatchInsertSQL() string { + return fmt.Sprintf("INSERT INTO %s (%s)", nodeMetricSnapshotTableName, nodeMetricSnapshotInsertColumns) +} + +// NodeRequestReport stores aggregated node request window reports in ClickHouse. +type NodeRequestReport struct { + ID uint64 `gorm:"column:id"` + NodeID string `gorm:"column:node_id"` + WindowStartedAt time.Time `gorm:"column:window_started_at"` + WindowEndedAt time.Time `gorm:"column:window_ended_at"` + RequestCount int64 `gorm:"column:request_count"` + ErrorCount int64 `gorm:"column:error_count"` + UniqueVisitorCount int64 `gorm:"column:unique_visitor_count"` + StatusCodesJSON string `gorm:"column:status_codes_json"` + TopDomainsJSON string `gorm:"column:top_domains_json"` + SourceCountriesJSON string `gorm:"column:source_countries_json"` + CreatedAt time.Time `gorm:"column:created_at"` +} + +// TableName returns the ClickHouse table name. +func (NodeRequestReport) TableName() string { + return nodeRequestReportTableName +} + +// InsertColumns returns comma-separated column names for batch insert. +func (NodeRequestReport) InsertColumns() string { + return nodeRequestReportInsertColumns +} + +// BatchInsertSQL returns the INSERT prefix used by native batch writers. +func (NodeRequestReport) BatchInsertSQL() string { + return fmt.Sprintf("INSERT INTO %s (%s)", nodeRequestReportTableName, nodeRequestReportInsertColumns) +} + +// NodeObsOpenresty stores OpenResty observability snapshots in ClickHouse. +type NodeObsOpenresty struct { + ID uint64 `gorm:"column:id"` + NodeID string `gorm:"column:node_id"` + CapturedAt time.Time `gorm:"column:captured_at"` + OpenrestyRxBytes int64 `gorm:"column:openresty_rx_bytes"` + OpenrestyTxBytes int64 `gorm:"column:openresty_tx_bytes"` + OpenrestyConnections int64 `gorm:"column:openresty_connections"` + CreatedAt time.Time `gorm:"column:created_at"` +} + +// TableName returns the ClickHouse table name. +func (NodeObsOpenresty) TableName() string { + return nodeObsOpenrestyTableName +} + +// InsertColumns returns comma-separated column names for batch insert. +func (NodeObsOpenresty) InsertColumns() string { + return nodeObsOpenrestyInsertColumns +} + +// BatchInsertSQL returns the INSERT prefix used by native batch writers. +func (NodeObsOpenresty) BatchInsertSQL() string { + return fmt.Sprintf("INSERT INTO %s (%s)", nodeObsOpenrestyTableName, nodeObsOpenrestyInsertColumns) +} + +// NodeObsFrps stores FRPS observability snapshots in ClickHouse. +type NodeObsFrps struct { + ID uint64 `gorm:"column:id"` + NodeID string `gorm:"column:node_id"` + CapturedAt time.Time `gorm:"column:captured_at"` + FrpsConnections int32 `gorm:"column:frps_connections"` + FrpsProxyCount int32 `gorm:"column:frps_proxy_count"` + FrpsClientCount int32 `gorm:"column:frps_client_count"` + FrpsProxies string `gorm:"column:frps_proxies"` + CreatedAt time.Time `gorm:"column:created_at"` +} + +// TableName returns the ClickHouse table name. +func (NodeObsFrps) TableName() string { + return nodeObsFrpsTableName +} + +// InsertColumns returns comma-separated column names for batch insert. +func (NodeObsFrps) InsertColumns() string { + return nodeObsFrpsInsertColumns +} + +// BatchInsertSQL returns the INSERT prefix used by native batch writers. +func (NodeObsFrps) BatchInsertSQL() string { + return fmt.Sprintf("INSERT INTO %s (%s)", nodeObsFrpsTableName, nodeObsFrpsInsertColumns) +} + +// NodeObsFrpc stores FRPC observability snapshots in ClickHouse. +type NodeObsFrpc struct { + ID uint64 `gorm:"column:id"` + NodeID string `gorm:"column:node_id"` + CapturedAt time.Time `gorm:"column:captured_at"` + TunnelStatus string `gorm:"column:tunnel_status"` + ConnectedRelaysCount int32 `gorm:"column:connected_relays_count"` + CreatedAt time.Time `gorm:"column:created_at"` +} + +// TableName returns the ClickHouse table name. +func (NodeObsFrpc) TableName() string { + return nodeObsFrpcTableName +} + +// InsertColumns returns comma-separated column names for batch insert. +func (NodeObsFrpc) InsertColumns() string { + return nodeObsFrpcInsertColumns +} + +// BatchInsertSQL returns the INSERT prefix used by native batch writers. +func (NodeObsFrpc) BatchInsertSQL() string { + return fmt.Sprintf("INSERT INTO %s (%s)", nodeObsFrpcTableName, nodeObsFrpcInsertColumns) +} diff --git a/internal/model/openflare_observability.go b/internal/model/openflare_observability.go index 95d52ac9..cacb10e1 100644 --- a/internal/model/openflare_observability.go +++ b/internal/model/openflare_observability.go @@ -13,7 +13,8 @@ import ( "gorm.io/gorm" ) -// OpenFlareMetricSnapshot stores a node capacity snapshot (v1 single table, no sharding). +// OpenFlareMetricSnapshot stores a node capacity snapshot in ClickHouse (database: openflare, table: of_node_metric_snapshots). +// ClickHouse DDL is managed by goose; reads/writes go through internal/repository/analytics. type OpenFlareMetricSnapshot struct { ID uint `json:"id" gorm:"primaryKey;autoIncrement"` NodeID string `json:"node_id" gorm:"index;size:64;not null"` @@ -35,7 +36,8 @@ func (OpenFlareMetricSnapshot) TableName() string { return "of_node_metric_snapshots" } -// OpenFlareRequestReport stores aggregated traffic windows per node. +// OpenFlareRequestReport stores aggregated traffic windows per node in ClickHouse (database: openflare, table: of_node_request_reports). +// ClickHouse DDL is managed by goose; reads/writes go through internal/repository/analytics. type OpenFlareRequestReport struct { ID uint `json:"id" gorm:"primaryKey;autoIncrement"` NodeID string `json:"node_id" gorm:"index;size:64;not null"` @@ -126,7 +128,8 @@ func (OpenFlareNodeSystemProfile) TableName() string { return "of_node_system_profiles" } -// OpenFlareNodeObservationOpenresty stores openresty network observations. +// OpenFlareNodeObservationOpenresty stores openresty network observations in ClickHouse (database: openflare, table: of_node_obs_openresty). +// ClickHouse DDL is managed by goose; reads/writes go through internal/repository/analytics. type OpenFlareNodeObservationOpenresty struct { ID uint `json:"id" gorm:"primaryKey;autoIncrement"` NodeID string `json:"node_id" gorm:"index;size:64;not null"` @@ -142,7 +145,8 @@ func (OpenFlareNodeObservationOpenresty) TableName() string { return "of_node_obs_openresty" } -// OpenFlareNodeObservationFrpc stores tunnel client frpc observations. +// OpenFlareNodeObservationFrpc stores tunnel client frpc observations in ClickHouse (database: openflare, table: of_node_obs_frpc). +// ClickHouse DDL is managed by goose; reads/writes go through internal/repository/analytics. type OpenFlareNodeObservationFrpc struct { ID uint `json:"id" gorm:"primaryKey;autoIncrement"` NodeID string `json:"node_id" gorm:"index;size:64;not null"` @@ -157,7 +161,8 @@ func (OpenFlareNodeObservationFrpc) TableName() string { return "of_node_obs_frpc" } -// OpenFlareNodeObservationFrps stores tunnel relay frps observations. +// OpenFlareNodeObservationFrps stores tunnel relay frps observations in ClickHouse (database: openflare, table: of_node_obs_frps). +// ClickHouse DDL is managed by goose; reads/writes go through internal/repository/analytics. type OpenFlareNodeObservationFrps struct { ID uint `json:"id" gorm:"primaryKey;autoIncrement"` NodeID string `json:"node_id" gorm:"index;size:64;not null"` @@ -285,39 +290,39 @@ func isMissingTableError(err error) bool { strings.Contains(msg, "does not exist") } -func listOpenFlareSince[T any](ctx context.Context, nodeID string, since time.Time, limit int, orderBy string, sinceColumn string) ([]*T, error) { - conn := db.DB(ctx) - if conn == nil { - return nil, errors.New(errDatabaseNotInitialized) - } - query := conn.Model(new(T)).Order(orderBy) - if nodeID != "" { - query = query.Where("node_id = ?", nodeID) - } - if !since.IsZero() { - query = query.Where(sinceColumn+" >= ?", since) - } - if limit > 0 { - query = query.Limit(limit) - } - var rows []*T - if err := query.Find(&rows).Error; err != nil { - if isMissingTableError(err) { - return []*T{}, nil - } - return nil, err - } - return rows, nil +// InsertOpenFlareMetricSnapshot inserts a metric snapshot into ClickHouse. +func InsertOpenFlareMetricSnapshot(ctx context.Context, record *OpenFlareMetricSnapshot) error { + return currentObservabilityStore().InsertMetricSnapshot(ctx, record) +} + +// InsertOpenFlareRequestReport inserts a request report into ClickHouse. +func InsertOpenFlareRequestReport(ctx context.Context, record *OpenFlareRequestReport) error { + return currentObservabilityStore().InsertRequestReport(ctx, record) +} + +// InsertOpenFlareNodeObservationOpenresty inserts an OpenResty observation into ClickHouse. +func InsertOpenFlareNodeObservationOpenresty(ctx context.Context, record *OpenFlareNodeObservationOpenresty) error { + return currentObservabilityStore().InsertNodeObservationOpenresty(ctx, record) +} + +// InsertOpenFlareNodeObservationFrps inserts an FRPS observation into ClickHouse. +func InsertOpenFlareNodeObservationFrps(ctx context.Context, record *OpenFlareNodeObservationFrps) error { + return currentObservabilityStore().InsertNodeObservationFrps(ctx, record) +} + +// InsertOpenFlareNodeObservationFrpc inserts an FRPC observation into ClickHouse. +func InsertOpenFlareNodeObservationFrpc(ctx context.Context, record *OpenFlareNodeObservationFrpc) error { + return currentObservabilityStore().InsertNodeObservationFrpc(ctx, record) } // ListOpenFlareMetricSnapshotsSince returns metric snapshots since the given time. func ListOpenFlareMetricSnapshotsSince(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareMetricSnapshot, error) { - return listOpenFlareSince[OpenFlareMetricSnapshot](ctx, nodeID, since, limit, "captured_at desc, id desc", "captured_at") + return currentObservabilityStore().ListMetricSnapshots(ctx, nodeID, since, limit) } // ListOpenFlareRequestReportsSince returns request reports since the given time. func ListOpenFlareRequestReportsSince(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareRequestReport, error) { - return listOpenFlareSince[OpenFlareRequestReport](ctx, nodeID, since, limit, "window_ended_at desc, id desc", "window_ended_at") + return currentObservabilityStore().ListRequestReports(ctx, nodeID, since, limit) } // ListOpenFlareActiveHealthEvents returns active health events across all nodes. @@ -361,66 +366,22 @@ func ListOpenFlareHealthEvents(ctx context.Context, nodeID string, activeOnly bo // DeleteOpenFlareMetricSnapshotsBefore deletes metric snapshots captured before cutoff. func DeleteOpenFlareMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn := db.DB(ctx) - if conn == nil { - return 0, errors.New(errDatabaseNotInitialized) - } - result := conn.Where("captured_at < ?", cutoff).Delete(&OpenFlareMetricSnapshot{}) - if result.Error != nil { - if isMissingTableError(result.Error) { - return 0, nil - } - return 0, result.Error - } - return result.RowsAffected, nil + return currentObservabilityStore().DeleteMetricSnapshotsBefore(ctx, cutoff) } // DeleteAllOpenFlareMetricSnapshots deletes all metric snapshots. func DeleteAllOpenFlareMetricSnapshots(ctx context.Context) (int64, error) { - conn := db.DB(ctx) - if conn == nil { - return 0, errors.New(errDatabaseNotInitialized) - } - result := conn.Where("1 = 1").Delete(&OpenFlareMetricSnapshot{}) - if result.Error != nil { - if isMissingTableError(result.Error) { - return 0, nil - } - return 0, result.Error - } - return result.RowsAffected, nil + return currentObservabilityStore().DeleteAllMetricSnapshots(ctx) } // DeleteOpenFlareRequestReportsBefore deletes request reports ending before cutoff. func DeleteOpenFlareRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn := db.DB(ctx) - if conn == nil { - return 0, errors.New(errDatabaseNotInitialized) - } - result := conn.Where("window_ended_at < ?", cutoff).Delete(&OpenFlareRequestReport{}) - if result.Error != nil { - if isMissingTableError(result.Error) { - return 0, nil - } - return 0, result.Error - } - return result.RowsAffected, nil + return currentObservabilityStore().DeleteRequestReportsBefore(ctx, cutoff) } // DeleteAllOpenFlareRequestReports deletes all request reports. func DeleteAllOpenFlareRequestReports(ctx context.Context) (int64, error) { - conn := db.DB(ctx) - if conn == nil { - return 0, errors.New(errDatabaseNotInitialized) - } - result := conn.Where("1 = 1").Delete(&OpenFlareRequestReport{}) - if result.Error != nil { - if isMissingTableError(result.Error) { - return 0, nil - } - return 0, result.Error - } - return result.RowsAffected, nil + return currentObservabilityStore().DeleteAllRequestReports(ctx) } // DeleteOpenFlareHealthEventsByNodeID deletes all health events for a node. @@ -457,15 +418,15 @@ func GetOpenFlareNodeSystemProfile(ctx context.Context, nodeID string) (*OpenFla // ListOpenFlareNodeObservationOpenresty returns openresty observations. func ListOpenFlareNodeObservationOpenresty(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationOpenresty, error) { - return listOpenFlareSince[OpenFlareNodeObservationOpenresty](ctx, nodeID, since, limit, "captured_at desc, id desc", "captured_at") + return currentObservabilityStore().ListNodeObservationOpenresty(ctx, nodeID, since, limit) } // ListOpenFlareNodeObservationFrpc returns frpc observations. func ListOpenFlareNodeObservationFrpc(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrpc, error) { - return listOpenFlareSince[OpenFlareNodeObservationFrpc](ctx, nodeID, since, limit, "captured_at desc, id desc", "captured_at") + return currentObservabilityStore().ListNodeObservationFrpc(ctx, nodeID, since, limit) } // ListOpenFlareNodeObservationFrps returns frps observations. func ListOpenFlareNodeObservationFrps(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrps, error) { - return listOpenFlareSince[OpenFlareNodeObservationFrps](ctx, nodeID, since, limit, "captured_at desc, id desc", "captured_at") + return currentObservabilityStore().ListNodeObservationFrps(ctx, nodeID, since, limit) } diff --git a/internal/model/openflare_observability_store.go b/internal/model/openflare_observability_store.go new file mode 100644 index 00000000..9c9ad95d --- /dev/null +++ b/internal/model/openflare_observability_store.go @@ -0,0 +1,369 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package model + +import ( + "context" + "math" + "sync" + "time" + + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" + analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" +) + +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) + + InsertRequestReport(ctx context.Context, record *OpenFlareRequestReport) error + ListRequestReports(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareRequestReport, error) + DeleteAllRequestReports(ctx context.Context) (int64, error) + DeleteRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) + + InsertNodeObservationOpenresty(ctx context.Context, record *OpenFlareNodeObservationOpenresty) error + ListNodeObservationOpenresty(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationOpenresty, error) + DeleteAllNodeObservationOpenresty(ctx context.Context) (int64, error) + DeleteNodeObservationOpenrestyBefore(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(ctx context.Context, record *OpenFlareMetricSnapshot) error { + if record == nil { + return nil + } + return analyticsrepo.InsertNodeMetricSnapshot(ctx, toAnalyticsNodeMetricSnapshot(record)) +} + +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) +} + +func (clickhouseObservabilityStore) InsertRequestReport(ctx context.Context, record *OpenFlareRequestReport) error { + if record == nil { + return nil + } + return analyticsrepo.InsertNodeRequestReport(ctx, toAnalyticsNodeRequestReport(record)) +} + +func (clickhouseObservabilityStore) ListRequestReports(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareRequestReport, error) { + rows, err := analyticsrepo.ListNodeRequestReports(ctx, toNodeObservabilityFilter(nodeID, since, limit)) + if err != nil { + return nil, err + } + return fromAnalyticsNodeRequestReports(rows), nil +} + +func (clickhouseObservabilityStore) DeleteAllRequestReports(ctx context.Context) (int64, error) { + return analyticsrepo.DeleteAllNodeRequestReports(ctx) +} + +func (clickhouseObservabilityStore) DeleteRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) { + return analyticsrepo.DeleteNodeRequestReportsBefore(ctx, cutoff) +} + +func (clickhouseObservabilityStore) InsertNodeObservationOpenresty(ctx context.Context, record *OpenFlareNodeObservationOpenresty) error { + if record == nil { + return nil + } + return analyticsrepo.InsertNodeObsOpenresty(ctx, toAnalyticsNodeObsOpenresty(record)) +} + +func (clickhouseObservabilityStore) ListNodeObservationOpenresty(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationOpenresty, error) { + rows, err := analyticsrepo.ListNodeObsOpenresty(ctx, toNodeObservabilityFilter(nodeID, since, limit)) + if err != nil { + return nil, err + } + return fromAnalyticsNodeObsOpenresty(rows), nil +} + +func (clickhouseObservabilityStore) DeleteAllNodeObservationOpenresty(ctx context.Context) (int64, error) { + return analyticsrepo.DeleteAllNodeObsOpenresty(ctx) +} + +func (clickhouseObservabilityStore) DeleteNodeObservationOpenrestyBefore(ctx context.Context, cutoff time.Time) (int64, error) { + return analyticsrepo.DeleteNodeObsOpenrestyBefore(ctx, cutoff) +} + +func (clickhouseObservabilityStore) InsertNodeObservationFrps(ctx context.Context, record *OpenFlareNodeObservationFrps) error { + if record == nil { + return nil + } + return analyticsrepo.InsertNodeObsFrps(ctx, toAnalyticsNodeObsFrps(record)) +} + +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(ctx context.Context, record *OpenFlareNodeObservationFrpc) error { + if record == nil { + return nil + } + return analyticsrepo.InsertNodeObsFrpc(ctx, toAnalyticsNodeObsFrpc(record)) +} + +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 toAnalyticsNodeRequestReport(record *OpenFlareRequestReport) analyticsmodel.NodeRequestReport { + return analyticsmodel.NodeRequestReport{ + ID: uint64(record.ID), + NodeID: record.NodeID, + WindowStartedAt: record.WindowStartedAt, + WindowEndedAt: record.WindowEndedAt, + RequestCount: record.RequestCount, + ErrorCount: record.ErrorCount, + UniqueVisitorCount: record.UniqueVisitorCount, + StatusCodesJSON: record.StatusCodesJSON, + TopDomainsJSON: record.TopDomainsJSON, + SourceCountriesJSON: record.SourceCountriesJSON, + CreatedAt: record.CreatedAt, + } +} + +func fromAnalyticsNodeRequestReports(rows []analyticsmodel.NodeRequestReport) []*OpenFlareRequestReport { + result := make([]*OpenFlareRequestReport, len(rows)) + for index, row := range rows { + result[index] = &OpenFlareRequestReport{ + ID: uint(row.ID), + NodeID: row.NodeID, + WindowStartedAt: row.WindowStartedAt, + WindowEndedAt: row.WindowEndedAt, + RequestCount: row.RequestCount, + ErrorCount: row.ErrorCount, + UniqueVisitorCount: row.UniqueVisitorCount, + StatusCodesJSON: row.StatusCodesJSON, + TopDomainsJSON: row.TopDomainsJSON, + SourceCountriesJSON: row.SourceCountriesJSON, + CreatedAt: row.CreatedAt, + } + } + return result +} + +func toAnalyticsNodeObsOpenresty(record *OpenFlareNodeObservationOpenresty) analyticsmodel.NodeObsOpenresty { + return analyticsmodel.NodeObsOpenresty{ + ID: uint64(record.ID), + NodeID: record.NodeID, + CapturedAt: record.CapturedAt, + OpenrestyRxBytes: record.OpenrestyRxBytes, + OpenrestyTxBytes: record.OpenrestyTxBytes, + OpenrestyConnections: record.OpenrestyConnections, + CreatedAt: record.CreatedAt, + } +} + +func fromAnalyticsNodeObsOpenresty(rows []analyticsmodel.NodeObsOpenresty) []*OpenFlareNodeObservationOpenresty { + result := make([]*OpenFlareNodeObservationOpenresty, len(rows)) + for index, row := range rows { + result[index] = &OpenFlareNodeObservationOpenresty{ + ID: uint(row.ID), + NodeID: row.NodeID, + CapturedAt: row.CapturedAt, + OpenrestyRxBytes: row.OpenrestyRxBytes, + OpenrestyTxBytes: row.OpenrestyTxBytes, + OpenrestyConnections: row.OpenrestyConnections, + 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 +} diff --git a/internal/model/openflare_observability_store_memory.go b/internal/model/openflare_observability_store_memory.go new file mode 100644 index 00000000..d31321ef --- /dev/null +++ b/internal/model/openflare_observability_store_memory.go @@ -0,0 +1,508 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package model + +import ( + "context" + "sort" + "strings" + "sync" + "time" + + "github.com/Rain-kl/Wavelet/internal/db/idgen" +) + +type memoryObservabilityStore struct { + mu sync.RWMutex + metricSnapshots []*OpenFlareMetricSnapshot + requestReports []*OpenFlareRequestReport + openrestyObs []*OpenFlareNodeObservationOpenresty + frpsObs []*OpenFlareNodeObservationFrps + frpcObs []*OpenFlareNodeObservationFrpc +} + +func (s *memoryObservabilityStore) InsertMetricSnapshot(_ context.Context, record *OpenFlareMetricSnapshot) error { + if record == nil { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + copyRecord := cloneOpenFlareMetricSnapshot(record) + if memoryMetricSnapshotExists(s.metricSnapshots, copyRecord.NodeID, copyRecord.CapturedAt) { + return nil + } + s.metricSnapshots = append(s.metricSnapshots, copyRecord) + return nil +} + +func (s *memoryObservabilityStore) ListMetricSnapshots(_ context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareMetricSnapshot, error) { + s.mu.RLock() + defer s.mu.RUnlock() + rows := memoryFilterMetricSnapshots(s.metricSnapshots, nodeID, since) + sortOpenFlareMetricSnapshots(rows) + return memoryLimitObservabilityRows(rows, limit), nil +} + +func (s *memoryObservabilityStore) DeleteAllMetricSnapshots(_ context.Context) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + count := int64(len(s.metricSnapshots)) + s.metricSnapshots = nil + return count, nil +} + +func (s *memoryObservabilityStore) DeleteMetricSnapshotsBefore(_ context.Context, cutoff time.Time) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + cutoff = cutoff.UTC() + remaining := make([]*OpenFlareMetricSnapshot, 0, len(s.metricSnapshots)) + var deleted int64 + for _, row := range s.metricSnapshots { + if row.CapturedAt.Before(cutoff) { + deleted++ + continue + } + remaining = append(remaining, row) + } + s.metricSnapshots = remaining + return deleted, nil +} + +func (s *memoryObservabilityStore) InsertRequestReport(_ context.Context, record *OpenFlareRequestReport) error { + if record == nil { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + copyRecord := cloneOpenFlareRequestReport(record) + if memoryRequestReportExists(s.requestReports, copyRecord.NodeID, copyRecord.WindowStartedAt, copyRecord.WindowEndedAt) { + return nil + } + s.requestReports = append(s.requestReports, copyRecord) + return nil +} + +func (s *memoryObservabilityStore) ListRequestReports(_ context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareRequestReport, error) { + s.mu.RLock() + defer s.mu.RUnlock() + rows := memoryFilterRequestReports(s.requestReports, nodeID, since) + sortOpenFlareRequestReports(rows) + return memoryLimitObservabilityRows(rows, limit), nil +} + +func (s *memoryObservabilityStore) DeleteAllRequestReports(_ context.Context) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + count := int64(len(s.requestReports)) + s.requestReports = nil + return count, nil +} + +func (s *memoryObservabilityStore) DeleteRequestReportsBefore(_ context.Context, cutoff time.Time) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + cutoff = cutoff.UTC() + remaining := make([]*OpenFlareRequestReport, 0, len(s.requestReports)) + var deleted int64 + for _, row := range s.requestReports { + if row.WindowEndedAt.Before(cutoff) { + deleted++ + continue + } + remaining = append(remaining, row) + } + s.requestReports = remaining + return deleted, nil +} + +func (s *memoryObservabilityStore) InsertNodeObservationOpenresty(_ context.Context, record *OpenFlareNodeObservationOpenresty) error { + if record == nil { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + s.openrestyObs = append(s.openrestyObs, cloneOpenFlareNodeObservationOpenresty(record)) + return nil +} + +func (s *memoryObservabilityStore) ListNodeObservationOpenresty(_ context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationOpenresty, error) { + s.mu.RLock() + defer s.mu.RUnlock() + rows := memoryFilterOpenrestyObservations(s.openrestyObs, nodeID, since) + sortOpenFlareNodeObservationOpenresty(rows) + return memoryLimitObservabilityRows(rows, limit), nil +} + +func (s *memoryObservabilityStore) DeleteAllNodeObservationOpenresty(_ context.Context) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + count := int64(len(s.openrestyObs)) + s.openrestyObs = nil + return count, nil +} + +func (s *memoryObservabilityStore) DeleteNodeObservationOpenrestyBefore(_ context.Context, cutoff time.Time) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + cutoff = cutoff.UTC() + remaining := make([]*OpenFlareNodeObservationOpenresty, 0, len(s.openrestyObs)) + var deleted int64 + for _, row := range s.openrestyObs { + if row.CapturedAt.Before(cutoff) { + deleted++ + continue + } + remaining = append(remaining, row) + } + s.openrestyObs = remaining + return deleted, nil +} + +func (s *memoryObservabilityStore) InsertNodeObservationFrps(_ context.Context, record *OpenFlareNodeObservationFrps) error { + if record == nil { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + s.frpsObs = append(s.frpsObs, cloneOpenFlareNodeObservationFrps(record)) + return nil +} + +func (s *memoryObservabilityStore) ListNodeObservationFrps(_ context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrps, error) { + s.mu.RLock() + defer s.mu.RUnlock() + rows := memoryFilterFrpsObservations(s.frpsObs, nodeID, since) + sortOpenFlareNodeObservationFrps(rows) + return memoryLimitObservabilityRows(rows, limit), nil +} + +func (s *memoryObservabilityStore) DeleteAllNodeObservationFrps(_ context.Context) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + count := int64(len(s.frpsObs)) + s.frpsObs = nil + return count, nil +} + +func (s *memoryObservabilityStore) DeleteNodeObservationFrpsBefore(_ context.Context, cutoff time.Time) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + cutoff = cutoff.UTC() + remaining := make([]*OpenFlareNodeObservationFrps, 0, len(s.frpsObs)) + var deleted int64 + for _, row := range s.frpsObs { + if row.CapturedAt.Before(cutoff) { + deleted++ + continue + } + remaining = append(remaining, row) + } + s.frpsObs = remaining + return deleted, nil +} + +func (s *memoryObservabilityStore) InsertNodeObservationFrpc(_ context.Context, record *OpenFlareNodeObservationFrpc) error { + if record == nil { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + s.frpcObs = append(s.frpcObs, cloneOpenFlareNodeObservationFrpc(record)) + return nil +} + +func (s *memoryObservabilityStore) ListNodeObservationFrpc(_ context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareNodeObservationFrpc, error) { + s.mu.RLock() + defer s.mu.RUnlock() + rows := memoryFilterFrpcObservations(s.frpcObs, nodeID, since) + sortOpenFlareNodeObservationFrpc(rows) + return memoryLimitObservabilityRows(rows, limit), nil +} + +func (s *memoryObservabilityStore) DeleteAllNodeObservationFrpc(_ context.Context) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + count := int64(len(s.frpcObs)) + s.frpcObs = nil + return count, nil +} + +func (s *memoryObservabilityStore) DeleteNodeObservationFrpcBefore(_ context.Context, cutoff time.Time) (int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + cutoff = cutoff.UTC() + remaining := make([]*OpenFlareNodeObservationFrpc, 0, len(s.frpcObs)) + var deleted int64 + for _, row := range s.frpcObs { + if row.CapturedAt.Before(cutoff) { + deleted++ + continue + } + remaining = append(remaining, row) + } + s.frpcObs = remaining + return deleted, nil +} + +func memoryFilterMetricSnapshots(rows []*OpenFlareMetricSnapshot, nodeID string, since time.Time) []*OpenFlareMetricSnapshot { + result := make([]*OpenFlareMetricSnapshot, 0, len(rows)) + for _, row := range rows { + if !memoryObservabilityMatchesNodeID(row.NodeID, nodeID) { + continue + } + if !since.IsZero() && row.CapturedAt.Before(since) { + continue + } + result = append(result, row) + } + return result +} + +func memoryFilterRequestReports(rows []*OpenFlareRequestReport, nodeID string, since time.Time) []*OpenFlareRequestReport { + result := make([]*OpenFlareRequestReport, 0, len(rows)) + for _, row := range rows { + if !memoryObservabilityMatchesNodeID(row.NodeID, nodeID) { + continue + } + if !since.IsZero() && row.WindowEndedAt.Before(since) { + continue + } + result = append(result, row) + } + return result +} + +func memoryFilterOpenrestyObservations(rows []*OpenFlareNodeObservationOpenresty, nodeID string, since time.Time) []*OpenFlareNodeObservationOpenresty { + result := make([]*OpenFlareNodeObservationOpenresty, 0, len(rows)) + for _, row := range rows { + if !memoryObservabilityMatchesNodeID(row.NodeID, nodeID) { + continue + } + if !since.IsZero() && row.CapturedAt.Before(since) { + continue + } + result = append(result, row) + } + return result +} + +func memoryFilterFrpsObservations(rows []*OpenFlareNodeObservationFrps, nodeID string, since time.Time) []*OpenFlareNodeObservationFrps { + result := make([]*OpenFlareNodeObservationFrps, 0, len(rows)) + for _, row := range rows { + if !memoryObservabilityMatchesNodeID(row.NodeID, nodeID) { + continue + } + if !since.IsZero() && row.CapturedAt.Before(since) { + continue + } + result = append(result, row) + } + return result +} + +func memoryFilterFrpcObservations(rows []*OpenFlareNodeObservationFrpc, nodeID string, since time.Time) []*OpenFlareNodeObservationFrpc { + result := make([]*OpenFlareNodeObservationFrpc, 0, len(rows)) + for _, row := range rows { + if !memoryObservabilityMatchesNodeID(row.NodeID, nodeID) { + continue + } + if !since.IsZero() && row.CapturedAt.Before(since) { + continue + } + result = append(result, row) + } + return result +} + +func memoryObservabilityMatchesNodeID(rowNodeID string, nodeID string) bool { + trimmed := strings.TrimSpace(nodeID) + if trimmed == "" { + return true + } + return rowNodeID == trimmed +} + +func memoryMetricSnapshotExists(rows []*OpenFlareMetricSnapshot, nodeID string, capturedAt time.Time) bool { + capturedAt = capturedAt.UTC() + for _, row := range rows { + if row.NodeID == nodeID && row.CapturedAt.UTC().Equal(capturedAt) { + return true + } + } + return false +} + +func memoryRequestReportExists(rows []*OpenFlareRequestReport, nodeID string, windowStartedAt, windowEndedAt time.Time) bool { + windowStartedAt = windowStartedAt.UTC() + windowEndedAt = windowEndedAt.UTC() + for _, row := range rows { + if row.NodeID == nodeID && + row.WindowStartedAt.UTC().Equal(windowStartedAt) && + row.WindowEndedAt.UTC().Equal(windowEndedAt) { + return true + } + } + return false +} + +func sortOpenFlareMetricSnapshots(items []*OpenFlareMetricSnapshot) { + sort.Slice(items, func(i, j int) bool { + left := items[i] + right := items[j] + if left == nil || right == nil { + return left != nil + } + if compare := openFlareAccessLogCompareInt64(left.CapturedAt.Unix(), right.CapturedAt.Unix()); compare != 0 { + return compare > 0 + } + return openFlareAccessLogCompareInt64(openFlareAccessLogUintToInt64(left.ID), openFlareAccessLogUintToInt64(right.ID)) > 0 + }) +} + +func sortOpenFlareRequestReports(items []*OpenFlareRequestReport) { + sort.Slice(items, func(i, j int) bool { + left := items[i] + right := items[j] + if left == nil || right == nil { + return left != nil + } + if compare := openFlareAccessLogCompareInt64(left.WindowEndedAt.Unix(), right.WindowEndedAt.Unix()); compare != 0 { + return compare > 0 + } + return openFlareAccessLogCompareInt64(openFlareAccessLogUintToInt64(left.ID), openFlareAccessLogUintToInt64(right.ID)) > 0 + }) +} + +func sortOpenFlareNodeObservationOpenresty(items []*OpenFlareNodeObservationOpenresty) { + sort.Slice(items, func(i, j int) bool { + left := items[i] + right := items[j] + if left == nil || right == nil { + return left != nil + } + if compare := openFlareAccessLogCompareInt64(left.CapturedAt.Unix(), right.CapturedAt.Unix()); compare != 0 { + return compare > 0 + } + return openFlareAccessLogCompareInt64(openFlareAccessLogUintToInt64(left.ID), openFlareAccessLogUintToInt64(right.ID)) > 0 + }) +} + +func sortOpenFlareNodeObservationFrps(items []*OpenFlareNodeObservationFrps) { + sort.Slice(items, func(i, j int) bool { + left := items[i] + right := items[j] + if left == nil || right == nil { + return left != nil + } + if compare := openFlareAccessLogCompareInt64(left.CapturedAt.Unix(), right.CapturedAt.Unix()); compare != 0 { + return compare > 0 + } + return openFlareAccessLogCompareInt64(openFlareAccessLogUintToInt64(left.ID), openFlareAccessLogUintToInt64(right.ID)) > 0 + }) +} + +func sortOpenFlareNodeObservationFrpc(items []*OpenFlareNodeObservationFrpc) { + sort.Slice(items, func(i, j int) bool { + left := items[i] + right := items[j] + if left == nil || right == nil { + return left != nil + } + if compare := openFlareAccessLogCompareInt64(left.CapturedAt.Unix(), right.CapturedAt.Unix()); compare != 0 { + return compare > 0 + } + return openFlareAccessLogCompareInt64(openFlareAccessLogUintToInt64(left.ID), openFlareAccessLogUintToInt64(right.ID)) > 0 + }) +} + +func memoryLimitObservabilityRows[T any](rows []T, limit int) []T { + if limit <= 0 || len(rows) <= limit { + result := make([]T, len(rows)) + copy(result, rows) + return result + } + result := make([]T, limit) + copy(result, rows[:limit]) + return result +} + +func cloneOpenFlareMetricSnapshot(record *OpenFlareMetricSnapshot) *OpenFlareMetricSnapshot { + copyRecord := *record + if copyRecord.ID == 0 { + copyRecord.ID = uint(idgen.NextUint64ID()) + } + now := time.Now().UTC() + if copyRecord.CreatedAt.IsZero() { + copyRecord.CreatedAt = now + } + copyRecord.CapturedAt = copyRecord.CapturedAt.UTC() + copyRecord.CreatedAt = copyRecord.CreatedAt.UTC() + return ©Record +} + +func cloneOpenFlareRequestReport(record *OpenFlareRequestReport) *OpenFlareRequestReport { + copyRecord := *record + if copyRecord.ID == 0 { + copyRecord.ID = uint(idgen.NextUint64ID()) + } + now := time.Now().UTC() + if copyRecord.CreatedAt.IsZero() { + copyRecord.CreatedAt = now + } + copyRecord.WindowStartedAt = copyRecord.WindowStartedAt.UTC() + copyRecord.WindowEndedAt = copyRecord.WindowEndedAt.UTC() + copyRecord.CreatedAt = copyRecord.CreatedAt.UTC() + return ©Record +} + +func cloneOpenFlareNodeObservationOpenresty(record *OpenFlareNodeObservationOpenresty) *OpenFlareNodeObservationOpenresty { + copyRecord := *record + if copyRecord.ID == 0 { + copyRecord.ID = uint(idgen.NextUint64ID()) + } + now := time.Now().UTC() + if copyRecord.CreatedAt.IsZero() { + copyRecord.CreatedAt = now + } + if copyRecord.CapturedAt.IsZero() { + copyRecord.CapturedAt = now + } + copyRecord.CapturedAt = copyRecord.CapturedAt.UTC() + copyRecord.CreatedAt = copyRecord.CreatedAt.UTC() + return ©Record +} + +func cloneOpenFlareNodeObservationFrps(record *OpenFlareNodeObservationFrps) *OpenFlareNodeObservationFrps { + copyRecord := *record + if copyRecord.ID == 0 { + copyRecord.ID = uint(idgen.NextUint64ID()) + } + now := time.Now().UTC() + if copyRecord.CreatedAt.IsZero() { + copyRecord.CreatedAt = now + } + if copyRecord.CapturedAt.IsZero() { + copyRecord.CapturedAt = now + } + copyRecord.CapturedAt = copyRecord.CapturedAt.UTC() + copyRecord.CreatedAt = copyRecord.CreatedAt.UTC() + return ©Record +} + +func cloneOpenFlareNodeObservationFrpc(record *OpenFlareNodeObservationFrpc) *OpenFlareNodeObservationFrpc { + copyRecord := *record + if copyRecord.ID == 0 { + copyRecord.ID = uint(idgen.NextUint64ID()) + } + now := time.Now().UTC() + if copyRecord.CreatedAt.IsZero() { + copyRecord.CreatedAt = now + } + if copyRecord.CapturedAt.IsZero() { + copyRecord.CapturedAt = now + } + copyRecord.CapturedAt = copyRecord.CapturedAt.UTC() + copyRecord.CreatedAt = copyRecord.CreatedAt.UTC() + return ©Record +} diff --git a/internal/repository/analytics/node_access_log.go b/internal/repository/analytics/node_access_log.go index ba460508..441433de 100644 --- a/internal/repository/analytics/node_access_log.go +++ b/internal/repository/analytics/node_access_log.go @@ -43,7 +43,7 @@ ORDER BY %s`, tableName, clause, nodeAccessLogOrderClause(filter.SortBy, filter. if filter.Page < 0 { filter.Page = 0 } - sql += " LIMIT ? OFFSET ?" + sql += clickHouseLimitOffsetClause args = append(args, filter.PageSize, filter.Page*filter.PageSize) } rows, err := conn.Query(ctx, sql, args...) @@ -123,7 +123,7 @@ WHERE %s AND trim(region) != '' GROUP BY trimmed_region ORDER BY count DESC, trimmed_region ASC`, tableName, clause) if limit > 0 { - sql += " LIMIT ?" + sql += clickHouseLimitClause args = append(args, limit) } rows, err := conn.Query(ctx, sql, args...) diff --git a/internal/repository/analytics/node_observability.go b/internal/repository/analytics/node_observability.go new file mode 100644 index 00000000..7c76ce4b --- /dev/null +++ b/internal/repository/analytics/node_observability.go @@ -0,0 +1,266 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "context" + "fmt" + + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" + "github.com/Rain-kl/Wavelet/internal/db" + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" +) + +func observabilityConn() (driver.Conn, error) { + if db.ChConn == nil { + return nil, fmt.Errorf("clickhouse connection is not initialized") + } + return db.ChConn, nil +} + +// ListNodeMetricSnapshots returns metric snapshots matching filter. +func ListNodeMetricSnapshots(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeMetricSnapshot, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + tableName := nodeMetricSnapshotTableName() + sql := fmt.Sprintf(` +SELECT id, node_id, captured_at, cpu_usage_percent, memory_used_bytes, memory_total_bytes, storage_used_bytes, storage_total_bytes, disk_read_bytes, disk_write_bytes, network_rx_bytes, network_tx_bytes, created_at +FROM %s +WHERE %s +ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause()) + 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 metric snapshots: %w", err) + } + defer func() { _ = rows.Close() }() + 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) +} + +// ListNodeObsOpenresty returns OpenResty observations matching filter. +func ListNodeObsOpenresty(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeObsOpenresty, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + tableName := nodeObsOpenrestyTableName() + sql := fmt.Sprintf(` +SELECT id, node_id, captured_at, openresty_rx_bytes, openresty_tx_bytes, openresty_connections, created_at +FROM %s +WHERE %s +ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause()) + 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 openresty observations: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeObsOpenrestyRows(rows) +} + +// ListNodeObsFrps returns FRPS observations matching filter. +func ListNodeObsFrps(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeObsFrps, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + tableName := nodeObsFrpsTableName() + sql := fmt.Sprintf(` +SELECT id, node_id, captured_at, frps_connections, frps_proxy_count, frps_client_count, frps_proxies, created_at +FROM %s +WHERE %s +ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause()) + 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 frps observations: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeObsFrpsRows(rows) +} + +// ListNodeObsFrpc returns FRPC observations matching filter. +func ListNodeObsFrpc(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeObsFrpc, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + tableName := nodeObsFrpcTableName() + sql := fmt.Sprintf(` +SELECT id, node_id, captured_at, tunnel_status, connected_relays_count, created_at +FROM %s +WHERE %s +ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause()) + 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 frpc observations: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeObsFrpcRows(rows) +} + +func scanNodeMetricSnapshotRows(rows driver.Rows) ([]analyticsmodel.NodeMetricSnapshot, error) { + var result []analyticsmodel.NodeMetricSnapshot + for rows.Next() { + var item analyticsmodel.NodeMetricSnapshot + if err := rows.Scan( + &item.ID, + &item.NodeID, + &item.CapturedAt, + &item.CPUUsagePercent, + &item.MemoryUsedBytes, + &item.MemoryTotalBytes, + &item.StorageUsedBytes, + &item.StorageTotalBytes, + &item.DiskReadBytes, + &item.DiskWriteBytes, + &item.NetworkRxBytes, + &item.NetworkTxBytes, + &item.CreatedAt, + ); err != nil { + return nil, fmt.Errorf("scan node metric snapshot row: %w", err) + } + item.CapturedAt = item.CapturedAt.UTC() + item.CreatedAt = item.CreatedAt.UTC() + result = append(result, item) + } + 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() { + var item analyticsmodel.NodeObsFrps + if err := rows.Scan( + &item.ID, + &item.NodeID, + &item.CapturedAt, + &item.FrpsConnections, + &item.FrpsProxyCount, + &item.FrpsClientCount, + &item.FrpsProxies, + &item.CreatedAt, + ); err != nil { + return nil, fmt.Errorf("scan node frps observation row: %w", err) + } + item.CapturedAt = item.CapturedAt.UTC() + item.CreatedAt = item.CreatedAt.UTC() + result = append(result, item) + } + return result, nil +} + +func scanNodeObsFrpcRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrpc, error) { + var result []analyticsmodel.NodeObsFrpc + for rows.Next() { + var item analyticsmodel.NodeObsFrpc + if err := rows.Scan( + &item.ID, + &item.NodeID, + &item.CapturedAt, + &item.TunnelStatus, + &item.ConnectedRelaysCount, + &item.CreatedAt, + ); err != nil { + return nil, fmt.Errorf("scan node frpc observation row: %w", err) + } + item.CapturedAt = item.CapturedAt.UTC() + item.CreatedAt = item.CreatedAt.UTC() + result = append(result, item) + } + return result, nil +} diff --git a/internal/repository/analytics/node_observability_delete.go b/internal/repository/analytics/node_observability_delete.go new file mode 100644 index 00000000..e772bf81 --- /dev/null +++ b/internal/repository/analytics/node_observability_delete.go @@ -0,0 +1,123 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "context" + "fmt" + "time" +) + +// DeleteAllNodeMetricSnapshots deletes all node metric snapshots. +func DeleteAllNodeMetricSnapshots(ctx context.Context) (int64, error) { + tableName := nodeMetricSnapshotTableName() + return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1") +} + +// DeleteNodeMetricSnapshotsBefore deletes metric snapshots captured before cutoff. +func DeleteNodeMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error) { + tableName := nodeMetricSnapshotTableName() + cutoff = cutoff.UTC() + return deleteNodeObservabilityWithCount( + ctx, + fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), + []any{cutoff}, + fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName), + cutoff, + ) +} + +// DeleteAllNodeRequestReports deletes all node request reports. +func DeleteAllNodeRequestReports(ctx context.Context) (int64, error) { + tableName := nodeRequestReportTableName() + return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1") +} + +// DeleteNodeRequestReportsBefore deletes request reports ending before cutoff. +func DeleteNodeRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) { + tableName := nodeRequestReportTableName() + cutoff = cutoff.UTC() + return deleteNodeObservabilityWithCount( + ctx, + fmt.Sprintf("SELECT count() FROM %s WHERE window_ended_at < ?", tableName), + []any{cutoff}, + fmt.Sprintf("ALTER TABLE %s DELETE WHERE window_ended_at < ?", tableName), + cutoff, + ) +} + +// DeleteAllNodeObsOpenresty deletes all OpenResty observations. +func DeleteAllNodeObsOpenresty(ctx context.Context) (int64, error) { + tableName := nodeObsOpenrestyTableName() + return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1") +} + +// DeleteNodeObsOpenrestyBefore deletes OpenResty observations captured before cutoff. +func DeleteNodeObsOpenrestyBefore(ctx context.Context, cutoff time.Time) (int64, error) { + tableName := nodeObsOpenrestyTableName() + cutoff = cutoff.UTC() + return deleteNodeObservabilityWithCount( + ctx, + fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), + []any{cutoff}, + fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName), + cutoff, + ) +} + +// DeleteAllNodeObsFrps deletes all FRPS observations. +func DeleteAllNodeObsFrps(ctx context.Context) (int64, error) { + tableName := nodeObsFrpsTableName() + return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1") +} + +// DeleteNodeObsFrpsBefore deletes FRPS observations captured before cutoff. +func DeleteNodeObsFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error) { + tableName := nodeObsFrpsTableName() + cutoff = cutoff.UTC() + return deleteNodeObservabilityWithCount( + ctx, + fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), + []any{cutoff}, + fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName), + cutoff, + ) +} + +// DeleteAllNodeObsFrpc deletes all FRPC observations. +func DeleteAllNodeObsFrpc(ctx context.Context) (int64, error) { + tableName := nodeObsFrpcTableName() + return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1") +} + +// DeleteNodeObsFrpcBefore deletes FRPC observations captured before cutoff. +func DeleteNodeObsFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error) { + tableName := nodeObsFrpcTableName() + cutoff = cutoff.UTC() + return deleteNodeObservabilityWithCount( + ctx, + fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), + []any{cutoff}, + fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName), + cutoff, + ) +} + +func deleteNodeObservabilityWithCount(ctx context.Context, countSQL string, countArgs []any, deleteSQL string, deleteArgs ...any) (int64, error) { + conn, err := observabilityConn() + if err != nil { + return 0, err + } + var count int64 + if err := conn.QueryRow(ctx, countSQL, countArgs...).Scan(&count); err != nil { + return 0, fmt.Errorf("count node observability rows for delete: %w", err) + } + if count == 0 { + return 0, nil + } + if err := conn.Exec(ctx, deleteSQL, deleteArgs...); err != nil { + return 0, fmt.Errorf("delete node observability rows: %w", err) + } + return count, nil +} diff --git a/internal/repository/analytics/node_observability_filter.go b/internal/repository/analytics/node_observability_filter.go new file mode 100644 index 00000000..6f0c674f --- /dev/null +++ b/internal/repository/analytics/node_observability_filter.go @@ -0,0 +1,63 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "strings" + "time" +) + +const nodeObservabilityFilterClauseCapacity = 3 + +// NodeObservabilityFilter scopes ClickHouse node observability queries. +type NodeObservabilityFilter struct { + NodeID string + Since time.Time + Limit int +} + +func buildNodeObservabilityFilterClause(filter NodeObservabilityFilter, sinceColumn string) (string, []any) { + parts := make([]string, 0, nodeObservabilityFilterClauseCapacity) + args := make([]any, 0, nodeObservabilityFilterClauseCapacity) + if trimmed := strings.TrimSpace(filter.NodeID); trimmed != "" { + parts = append(parts, "node_id = ?") + args = append(args, trimmed) + } + if !filter.Since.IsZero() { + parts = append(parts, sinceColumn+" >= ?") + args = append(args, filter.Since.UTC()) + } + if len(parts) == 0 { + return "1", nil + } + return strings.Join(parts, " AND "), args +} + +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 nodeObsOpenrestyTableName() string { + return "of_node_obs_openresty" +} + +func nodeObsFrpsTableName() string { + return "of_node_obs_frps" +} + +func nodeObsFrpcTableName() string { + return "of_node_obs_frpc" +} diff --git a/internal/repository/analytics/node_observability_test.go b/internal/repository/analytics/node_observability_test.go new file mode 100644 index 00000000..42ef2b0a --- /dev/null +++ b/internal/repository/analytics/node_observability_test.go @@ -0,0 +1,47 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "context" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestInsertNodeObsOpenresty_EmptyNodeID(t *testing.T) { + err := InsertNodeObsOpenresty(context.Background(), analyticsmodel.NodeObsOpenresty{}) + require.NoError(t, err) +} + +func TestInsertNodeObsOpenresty_UsesModelBatchSQL(t *testing.T) { + ctx := context.Background() + mockBatch := &mockBatch{} + mockConn := &mockConn{ + batch: mockBatch, + batchQuery: analyticsmodel.NodeObsOpenresty{}.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, + }) + require.NoError(t, err) + assert.True(t, mockConn.prepareCalled) + assert.Equal(t, analyticsmodel.NodeObsOpenresty{}.BatchInsertSQL(), mockConn.preparedQuery) + assert.True(t, mockBatch.sendCalled) + require.Len(t, mockBatch.rows, 1) + assert.Equal(t, "node-a", mockBatch.rows[0][1]) +} diff --git a/internal/repository/analytics/node_observability_writer.go b/internal/repository/analytics/node_observability_writer.go new file mode 100644 index 00000000..fd0d41e6 --- /dev/null +++ b/internal/repository/analytics/node_observability_writer.go @@ -0,0 +1,298 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "context" + "fmt" + "strings" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/db/idgen" + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" +) + +// InsertNodeMetricSnapshot writes a metric snapshot when no row exists for node_id+captured_at. +func InsertNodeMetricSnapshot(ctx context.Context, snapshot analyticsmodel.NodeMetricSnapshot) error { + nodeID := strings.TrimSpace(snapshot.NodeID) + if nodeID == "" { + return nil + } + capturedAt := snapshot.CapturedAt.UTC() + + exists, err := nodeObservabilityRowExists( + ctx, + fmt.Sprintf("SELECT count() FROM %s WHERE node_id = ? AND captured_at = ?", nodeMetricSnapshotTableName()), + nodeID, capturedAt, + ) + if err != nil { + return err + } + if exists { + return nil + } + + if db.ChConn == nil { + return fmt.Errorf("clickhouse connection is not initialized") + } + + batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeMetricSnapshot{}.BatchInsertSQL()) + if err != nil { + return fmt.Errorf("prepare clickhouse batch: %w", err) + } + + now := time.Now().UTC() + id := snapshot.ID + if id == 0 { + id = idgen.NextUint64ID() + } + createdAt := snapshot.CreatedAt + if createdAt.IsZero() { + createdAt = now + } + if err := batch.Append( + id, + nodeID, + capturedAt, + snapshot.CPUUsagePercent, + snapshot.MemoryUsedBytes, + snapshot.MemoryTotalBytes, + snapshot.StorageUsedBytes, + snapshot.StorageTotalBytes, + snapshot.DiskReadBytes, + snapshot.DiskWriteBytes, + snapshot.NetworkRxBytes, + snapshot.NetworkTxBytes, + createdAt.UTC(), + ); err != nil { + return fmt.Errorf("append node metric snapshot to batch: %w", err) + } + if err := batch.Send(); err != nil { + return fmt.Errorf("send clickhouse batch: %w", err) + } + return nil +} + +// InsertNodeRequestReport writes a request report when no row exists for node_id+window bounds. +func InsertNodeRequestReport(ctx context.Context, report analyticsmodel.NodeRequestReport) error { + nodeID := strings.TrimSpace(report.NodeID) + if nodeID == "" { + return nil + } + windowStartedAt := report.WindowStartedAt.UTC() + windowEndedAt := report.WindowEndedAt.UTC() + + exists, err := nodeObservabilityRowExists( + ctx, + fmt.Sprintf( + "SELECT count() FROM %s WHERE node_id = ? AND window_started_at = ? AND window_ended_at = ?", + nodeRequestReportTableName(), + ), + nodeID, windowStartedAt, windowEndedAt, + ) + if err != nil { + return err + } + if exists { + return nil + } + + if db.ChConn == nil { + return fmt.Errorf("clickhouse connection is not initialized") + } + + batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeRequestReport{}.BatchInsertSQL()) + if err != nil { + return fmt.Errorf("prepare clickhouse batch: %w", err) + } + + now := time.Now().UTC() + id := report.ID + if id == 0 { + id = idgen.NextUint64ID() + } + createdAt := report.CreatedAt + if createdAt.IsZero() { + createdAt = now + } + if err := batch.Append( + id, + nodeID, + windowStartedAt, + windowEndedAt, + 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 err := batch.Send(); err != nil { + return fmt.Errorf("send clickhouse batch: %w", err) + } + return nil +} + +// InsertNodeObsOpenresty writes an OpenResty observability snapshot. +func InsertNodeObsOpenresty(ctx context.Context, obs analyticsmodel.NodeObsOpenresty) error { + nodeID := strings.TrimSpace(obs.NodeID) + if nodeID == "" { + return nil + } + return insertNodeObsOpenrestyBatch(ctx, obs, nodeID) +} + +// InsertNodeObsFrps writes an FRPS observability snapshot. +func InsertNodeObsFrps(ctx context.Context, obs analyticsmodel.NodeObsFrps) error { + nodeID := strings.TrimSpace(obs.NodeID) + if nodeID == "" { + return nil + } + return insertNodeObsFrpsBatch(ctx, obs, nodeID) +} + +// InsertNodeObsFrpc writes an FRPC observability snapshot. +func InsertNodeObsFrpc(ctx context.Context, obs analyticsmodel.NodeObsFrpc) error { + nodeID := strings.TrimSpace(obs.NodeID) + if nodeID == "" { + return nil + } + return insertNodeObsFrpcBatch(ctx, obs, nodeID) +} + +func insertNodeObsOpenrestyBatch(ctx context.Context, obs analyticsmodel.NodeObsOpenresty, nodeID string) error { + 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() + id := obs.ID + if id == 0 { + id = idgen.NextUint64ID() + } + createdAt := obs.CreatedAt + if createdAt.IsZero() { + createdAt = now + } + capturedAt := obs.CapturedAt.UTC() + if capturedAt.IsZero() { + capturedAt = now + } + if err := batch.Append( + id, + nodeID, + capturedAt, + obs.OpenrestyRxBytes, + obs.OpenrestyTxBytes, + obs.OpenrestyConnections, + createdAt.UTC(), + ); err != nil { + return fmt.Errorf("append node openresty observation to batch: %w", err) + } + if err := batch.Send(); err != nil { + return fmt.Errorf("send clickhouse batch: %w", err) + } + return nil +} + +func insertNodeObsFrpsBatch(ctx context.Context, obs analyticsmodel.NodeObsFrps, nodeID string) error { + if db.ChConn == nil { + return fmt.Errorf("clickhouse connection is not initialized") + } + + batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeObsFrps{}.BatchInsertSQL()) + if err != nil { + return fmt.Errorf("prepare clickhouse batch: %w", err) + } + + now := time.Now().UTC() + id := obs.ID + if id == 0 { + id = idgen.NextUint64ID() + } + createdAt := obs.CreatedAt + if createdAt.IsZero() { + createdAt = now + } + capturedAt := obs.CapturedAt.UTC() + if capturedAt.IsZero() { + capturedAt = now + } + if err := batch.Append( + id, + nodeID, + capturedAt, + obs.FrpsConnections, + obs.FrpsProxyCount, + obs.FrpsClientCount, + obs.FrpsProxies, + createdAt.UTC(), + ); err != nil { + return fmt.Errorf("append node frps observation to batch: %w", err) + } + if err := batch.Send(); err != nil { + return fmt.Errorf("send clickhouse batch: %w", err) + } + return nil +} + +func insertNodeObsFrpcBatch(ctx context.Context, obs analyticsmodel.NodeObsFrpc, nodeID string) error { + if db.ChConn == nil { + return fmt.Errorf("clickhouse connection is not initialized") + } + + batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeObsFrpc{}.BatchInsertSQL()) + if err != nil { + return fmt.Errorf("prepare clickhouse batch: %w", err) + } + + now := time.Now().UTC() + id := obs.ID + if id == 0 { + id = idgen.NextUint64ID() + } + createdAt := obs.CreatedAt + if createdAt.IsZero() { + createdAt = now + } + capturedAt := obs.CapturedAt.UTC() + if capturedAt.IsZero() { + capturedAt = now + } + if err := batch.Append( + id, + nodeID, + capturedAt, + obs.TunnelStatus, + obs.ConnectedRelaysCount, + createdAt.UTC(), + ); err != nil { + return fmt.Errorf("append node frpc observation to batch: %w", err) + } + if err := batch.Send(); err != nil { + return fmt.Errorf("send clickhouse batch: %w", err) + } + return nil +} + +func nodeObservabilityRowExists(ctx context.Context, countSQL string, args ...any) (bool, error) { + conn, err := observabilityConn() + if err != nil { + return false, err + } + var count int64 + if err := conn.QueryRow(ctx, countSQL, args...).Scan(&count); err != nil { + return false, fmt.Errorf("check observability row exists: %w", err) + } + return count > 0, nil +} diff --git a/internal/repository/analytics/sql_fragments.go b/internal/repository/analytics/sql_fragments.go new file mode 100644 index 00000000..2f657fe3 --- /dev/null +++ b/internal/repository/analytics/sql_fragments.go @@ -0,0 +1,9 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +const ( + clickHouseLimitClause = " LIMIT ?" + clickHouseLimitOffsetClause = " LIMIT ? OFFSET ?" +)