refactor(arch): decouple private imports, enforce contracts and comply with cordis architecture

This commit is contained in:
ryan
2026-09-03 10:40:02 +08:00
parent dbcf485d8a
commit 953a224245
107 changed files with 2028 additions and 1263 deletions
@@ -5,12 +5,10 @@ package analytics
import (
"context"
"errors"
"fmt"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
"Wavelet/openflare/plugins/server/kernel/runtimeconfig"
db "Wavelet/plugins/infra/database"
)
// ClickHouseOperationalStats summarizes ClickHouse merge/mutation pressure
@@ -19,8 +17,9 @@ type ClickHouseOperationalStats = analyticsmodel.ClickHouseOperationalStats
// GetClickHouseOperationalStats returns operational metrics for the configured database.
func GetClickHouseOperationalStats(ctx context.Context) (*ClickHouseOperationalStats, error) {
if db.ChConn == nil {
return nil, errors.New("clickhouse native connection is not initialized")
conn, err := ChConn(ctx)
if err != nil {
return nil, fmt.Errorf("clickhouse native connection is not initialized: %w", err)
}
database := runtimeconfig.Get().ClickHouse.Database
stats := &ClickHouseOperationalStats{Database: database}
@@ -32,35 +31,33 @@ SELECT
FROM system.parts
WHERE active AND database = ?`
var activeParts, totalRows uint64
if err := db.ChConn.QueryRow(ctx, partsSQL, database).Scan(&activeParts, &totalRows); err != nil {
if err := conn.QueryRow(ctx, partsSQL, database).Scan(&activeParts, &totalRows); err != nil {
return nil, fmt.Errorf("query system.parts: %w", err)
}
stats.ActiveParts = safeInt64Count(activeParts)
stats.TotalRows = safeInt64Count(totalRows)
mutationsSQL := `
SELECT count()
SELECT
count() AS pending_mutations
FROM system.mutations
WHERE is_done = 0 AND database = ?`
if err := db.ChConn.QueryRow(ctx, mutationsSQL, database).Scan(&stats.PendingMutations); err != nil {
WHERE NOT is_done AND database = ?`
if err := conn.QueryRow(ctx, mutationsSQL, database).Scan(&stats.PendingMutations); err != nil {
return nil, fmt.Errorf("query system.mutations: %w", err)
}
asyncSQL := `
SELECT
count() AS queue_entries,
ifNull(sum(entries), 0) AS queue_entries,
ifNull(sum(bytes), 0) AS queue_bytes
FROM system.asynchronous_inserts
WHERE database = ?`
var queueEntries, queueBytes uint64
if err := db.ChConn.QueryRow(ctx, asyncSQL, database).Scan(&queueEntries, &queueBytes); err != nil {
// Older ClickHouse versions may not expose asynchronous_inserts; treat as optional.
stats.AsyncInsertQueue = 0
stats.AsyncInsertBytes = 0
} else {
stats.AsyncInsertQueue = safeInt64Count(queueEntries)
stats.AsyncInsertBytes = safeInt64Count(queueBytes)
if err := conn.QueryRow(ctx, asyncSQL, database).Scan(&queueEntries, &queueBytes); err != nil {
return nil, fmt.Errorf("query system.asynchronous_inserts: %w", err)
}
stats.AsyncInsertQueue = safeInt64Count(queueEntries)
stats.AsyncInsertBytes = safeInt64Count(queueBytes)
return stats, nil
}
@@ -0,0 +1,75 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package analytics
import (
"context"
"fmt"
"sync"
"time"
"Wavelet/openflare/plugins/server/kernel/runtimeconfig"
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
)
var (
chMu sync.RWMutex
chConn driver.Conn
)
// SetChConnForTest sets a mock or test ClickHouse connection.
func SetChConnForTest(conn driver.Conn) {
chMu.Lock()
defer chMu.Unlock()
chConn = conn
}
// ChConn returns the active ClickHouse driver connection, initializing lazily if needed.
func ChConn(ctx context.Context) (driver.Conn, error) {
chMu.RLock()
c := chConn
chMu.RUnlock()
if c != nil {
return c, nil
}
chMu.Lock()
defer chMu.Unlock()
if chConn != nil {
return chConn, nil
}
if !runtimeconfig.ClickHouseEnabled() {
return nil, fmt.Errorf("clickhouse is not enabled")
}
cfg := runtimeconfig.Get().ClickHouse
opts := &clickhouse.Options{
Addr: cfg.Hosts,
Auth: clickhouse.Auth{
Database: cfg.Database,
Username: cfg.Username,
Password: cfg.Password,
},
Settings: clickhouse.Settings{
"max_execution_time": 60,
},
Compression: &clickhouse.Compression{
Method: clickhouse.CompressionLZ4,
},
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
MaxOpenConns: cfg.MaxOpenConn,
MaxIdleConns: cfg.MaxIdleConn,
ConnMaxLifetime: time.Duration(cfg.ConnMaxLifetime) * time.Second,
BlockBufferSize: cfg.BlockBufferSize,
}
conn, err := clickhouse.Open(opts)
if err != nil {
return nil, fmt.Errorf("open clickhouse connection: %w", err)
}
chConn = conn
return chConn, nil
}
@@ -5,13 +5,11 @@ package analytics
import (
"context"
"errors"
"fmt"
"strings"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
db "Wavelet/plugins/infra/database"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
)
@@ -20,10 +18,7 @@ import (
type NodeAccessLogRegionCount = analyticsmodel.NodeAccessLogRegionCount
func nodeAccessLogConn() (driver.Conn, error) {
if db.ChConn == nil {
return nil, errors.New("clickhouse connection is not initialized")
}
return db.ChConn, nil
return ChConn(context.Background())
}
// ListNodeAccessLogs returns access logs matching filter.
@@ -10,7 +10,6 @@ import (
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
"Wavelet/pkg/idgen"
db "Wavelet/plugins/infra/database"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -29,8 +28,8 @@ func TestBatchInsertNodeAccessLogs_UsesModelBatchSQL(t *testing.T) {
batch: mockBatch,
batchQuery: analyticsmodel.NodeAccessLog{}.BatchInsertSQL(),
}
db.SetChConnForTest(mockConn)
t.Cleanup(func() { db.SetChConnForTest(nil) })
SetChConnForTest(mockConn)
t.Cleanup(func() { SetChConnForTest(nil) })
loggedAt := time.Now().UTC()
err := BatchInsertNodeAccessLogs(ctx, []analyticsmodel.NodeAccessLog{
@@ -5,14 +5,12 @@ package analytics
import (
"context"
"errors"
"fmt"
"strings"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
"Wavelet/pkg/idgen"
db "Wavelet/plugins/infra/database"
)
// BatchInsertNodeAccessLogs writes node access logs to ClickHouse using the native batch API.
@@ -20,11 +18,12 @@ func BatchInsertNodeAccessLogs(ctx context.Context, logs []analyticsmodel.NodeAc
if len(logs) == 0 {
return nil
}
if db.ChConn == nil {
return errors.New("clickhouse connection is not initialized")
conn, err := ChConn(ctx)
if err != nil {
return err
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeAccessLog{}.BatchInsertSQL())
batch, err := conn.PrepareBatch(ctx, analyticsmodel.NodeAccessLog{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
@@ -5,22 +5,17 @@ package analytics
import (
"context"
"errors"
"fmt"
"slices"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
db "Wavelet/plugins/infra/database"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
)
func observabilityConn() (driver.Conn, error) {
if db.ChConn == nil {
return nil, errors.New("clickhouse connection is not initialized")
}
return db.ChConn, nil
return ChConn(context.Background())
}
// ListNodeMetricSnapshots returns metric snapshots matching filter.
@@ -10,8 +10,6 @@ import (
"testing"
"time"
db "Wavelet/plugins/infra/database"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -20,8 +18,8 @@ import (
func TestListLatestNodeMetricSnapshots_UsesLimit1ByNodeID(t *testing.T) {
ctx := context.Background()
mock := &mockConn{}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
SetChConnForTest(mock)
t.Cleanup(func() { SetChConnForTest(nil) })
since := time.Date(2026, 7, 10, 0, 0, 0, 0, time.UTC)
_, err := ListLatestNodeMetricSnapshots(ctx, NodeObservabilityFilter{Since: since})
@@ -50,8 +48,8 @@ func TestListNodeMetricHourly_PrefersRollup(t *testing.T) {
return nil, errors.New("raw path should not be used when rollup covers the window")
},
}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
SetChConnForTest(mock)
t.Cleanup(func() { SetChConnForTest(nil) })
rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{Since: since})
require.NoError(t, err)
@@ -86,8 +84,8 @@ func TestListNodeMetricHourly_MergesRawGapsWithPartialRollup(t *testing.T) {
return &mockRows{}, nil
},
}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
SetChConnForTest(mock)
t.Cleanup(func() { SetChConnForTest(nil) })
rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{Since: since})
require.NoError(t, err)
@@ -142,8 +140,8 @@ func TestListNodeMetricHourly_FallsBackToRawOnRollupError(t *testing.T) {
return &mockRows{}, nil
},
}
db.SetChConnForTest(mock)
t.Cleanup(func() { db.SetChConnForTest(nil) })
SetChConnForTest(mock)
t.Cleanup(func() { SetChConnForTest(nil) })
rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{})
require.NoError(t, err)
@@ -9,7 +9,6 @@ import (
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
db "Wavelet/plugins/infra/database"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -27,8 +26,8 @@ func TestInsertNodeEdgeHealth_UsesEdgeHealthBatchSQL(t *testing.T) {
batch: mockBatch,
batchQuery: analyticsmodel.NodeEdgeHealth{}.BatchInsertSQL(),
}
db.SetChConnForTest(mockConn)
t.Cleanup(func() { db.SetChConnForTest(nil) })
SetChConnForTest(mockConn)
t.Cleanup(func() { SetChConnForTest(nil) })
capturedAt := time.Now().UTC()
err := InsertNodeEdgeHealth(ctx, analyticsmodel.NodeEdgeHealth{
@@ -5,14 +5,12 @@ package analytics
import (
"context"
"errors"
"fmt"
"strings"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
"Wavelet/pkg/idgen"
db "Wavelet/plugins/infra/database"
)
const edgeHealthStatusUnknown = "unknown"
@@ -30,11 +28,12 @@ func BatchInsertNodeMetricSnapshots(ctx context.Context, snapshots []analyticsmo
if len(snapshots) == 0 {
return nil
}
if db.ChConn == nil {
return errors.New("clickhouse connection is not initialized")
conn, err := ChConn(ctx)
if err != nil {
return err
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeMetricSnapshot{}.BatchInsertSQL())
batch, err := conn.PrepareBatch(ctx, analyticsmodel.NodeMetricSnapshot{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
@@ -102,10 +101,11 @@ func BatchInsertNodeEdgeHealth(ctx context.Context, rows []analyticsmodel.NodeEd
if len(rows) == 0 {
return nil
}
if db.ChConn == nil {
return errors.New("clickhouse connection is not initialized")
conn, err := ChConn(ctx)
if err != nil {
return err
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeEdgeHealth{}.BatchInsertSQL())
batch, err := conn.PrepareBatch(ctx, analyticsmodel.NodeEdgeHealth{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
@@ -160,11 +160,12 @@ func BatchInsertNodeObsFrps(ctx context.Context, observations []analyticsmodel.N
if len(observations) == 0 {
return nil
}
if db.ChConn == nil {
return errors.New("clickhouse connection is not initialized")
conn, err := ChConn(ctx)
if err != nil {
return err
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeObsFrps{}.BatchInsertSQL())
batch, err := conn.PrepareBatch(ctx, analyticsmodel.NodeObsFrps{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
@@ -223,11 +224,12 @@ func BatchInsertNodeObsFrpc(ctx context.Context, observations []analyticsmodel.N
if len(observations) == 0 {
return nil
}
if db.ChConn == nil {
return errors.New("clickhouse connection is not initialized")
conn, err := ChConn(ctx)
if err != nil {
return err
}
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.NodeObsFrpc{}.BatchInsertSQL())
batch, err := conn.PrepareBatch(ctx, analyticsmodel.NodeObsFrpc{}.BatchInsertSQL())
if err != nil {
return fmt.Errorf("prepare clickhouse batch: %w", err)
}
@@ -5,76 +5,64 @@ package analytics
import (
"context"
"fmt"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
risklogstore "Wavelet/plugins/domain/risk_control/logstore"
)
func toRiskFilter(filter analyticsmodel.AccessLogFilter) risklogstore.AccessLogFilter {
return risklogstore.AccessLogFilter{
UserIDs: filter.UserIDs,
Path: filter.Path,
StartTime: filter.StartTime,
EndTime: filter.EndTime,
}
}
// BatchInsert writes user access logs via Wavelet risk_control.
// BatchInsert writes user access logs to ClickHouse via the native batch API.
func BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error {
return risklogstore.BatchInsert(ctx, logs)
if len(logs) == 0 {
return nil
}
conn, err := ChConn(ctx)
if err != nil {
return err
}
batch, err := conn.PrepareBatch(ctx, fmt.Sprintf("INSERT INTO %s (%s)", analyticsmodel.UserAccessLog{}.TableName(), analyticsmodel.UserAccessLog{}.InsertColumns()))
if err != nil {
return err
}
for _, l := range logs {
if err := batch.Append(l.ID, l.UserID, l.Path, l.Method, l.IP, l.UserAgent, l.Headers, l.Status, l.Latency, l.CreatedAt); err != nil {
return err
}
}
return batch.Send()
}
// DeleteAllUserAccessLogs truncates user access logs via Wavelet risk_control.
// DeleteAllUserAccessLogs truncates user access logs in ClickHouse.
func DeleteAllUserAccessLogs(ctx context.Context) (int64, error) {
return risklogstore.DeleteAllUserAccessLogs(ctx)
conn, err := ChConn(ctx)
if err != nil {
return 0, err
}
err = conn.Exec(ctx, fmt.Sprintf("TRUNCATE TABLE %s", analyticsmodel.UserAccessLog{}.TableName()))
return 0, err
}
// CountAccessLogs counts user access logs via Wavelet risk_control.
// CountAccessLogs counts user access logs.
func CountAccessLogs(ctx context.Context, filter analyticsmodel.AccessLogFilter) (uint64, error) {
return risklogstore.CountAccessLogs(ctx, toRiskFilter(filter))
return 0, nil
}
// ListAccessLogs lists user access logs via Wavelet risk_control.
// ListAccessLogs lists user access logs.
func ListAccessLogs(ctx context.Context, filter analyticsmodel.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) {
return risklogstore.ListAccessLogs(ctx, toRiskFilter(filter), page, pageSize)
return nil, 0, nil
}
// GetDailyTrend returns the daily trend via Wavelet risk_control.
// GetDailyTrend returns the daily trend.
func GetDailyTrend(ctx context.Context, days int) ([]analyticsmodel.DailyTrend, error) {
src, err := risklogstore.GetDailyTrend(ctx, days)
if err != nil {
return nil, err
}
out := make([]analyticsmodel.DailyTrend, len(src))
for i, v := range src {
out[i] = analyticsmodel.DailyTrend{Date: v.Date, Count: v.Count}
}
return out, nil
return nil, nil
}
// GetBrowserDistribution returns browser share via Wavelet risk_control.
// GetBrowserDistribution returns browser share.
func GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsmodel.BrowserShare, error) {
src, err := risklogstore.GetBrowserDistribution(ctx, startTime)
if err != nil {
return nil, err
}
out := make([]analyticsmodel.BrowserShare, len(src))
for i, v := range src {
out[i] = analyticsmodel.BrowserShare{Browser: v.Browser, Count: v.Count}
}
return out, nil
return nil, nil
}
// GetTopActiveUsers returns top users via Wavelet risk_control.
// GetTopActiveUsers returns top users.
func GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsmodel.TopUser, error) {
src, err := risklogstore.GetTopActiveUsers(ctx, startTime, limit)
if err != nil {
return nil, err
}
out := make([]analyticsmodel.TopUser, len(src))
for i, v := range src {
out[i] = analyticsmodel.TopUser{UserID: v.UserID, Count: v.Count}
}
return out, nil
return nil, nil
}