refactor(backend): rename OpenFlare directory to lowercase openflare

This commit is contained in:
ryan
2026-08-30 17:43:23 +08:00
parent 06d5fedbfc
commit c93ff6674f
543 changed files with 819 additions and 819 deletions
@@ -0,0 +1,65 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package chwriter
import (
"sync"
"time"
)
const dedupTTL = 2 * time.Minute
type dedupSet struct {
mu sync.Mutex
keys map[string]time.Time
lastCleanup time.Time
}
func newDedupSet() *dedupSet {
return &dedupSet{
keys: make(map[string]time.Time),
lastCleanup: time.Now(),
}
}
// markIfNew records key when it has not been seen within dedupTTL.
func (s *dedupSet) markIfNew(key string) bool {
if s == nil || key == "" {
return false
}
now := time.Now()
s.mu.Lock()
defer s.mu.Unlock()
s.cleanupExpiredLocked(now)
if expiresAt, exists := s.keys[key]; exists && now.Before(expiresAt) {
return false
}
s.keys[key] = now.Add(dedupTTL)
return true
}
// unmark removes a key so a later enqueue or flush retry may accept it again.
func (s *dedupSet) unmark(key string) {
if s == nil || key == "" {
return
}
s.mu.Lock()
defer s.mu.Unlock()
delete(s.keys, key)
}
func (s *dedupSet) cleanupExpiredLocked(now time.Time) {
if now.Sub(s.lastCleanup) < 30*time.Second {
return
}
for existing, expiresAt := range s.keys {
if now.After(expiresAt) {
delete(s.keys, existing)
}
}
s.lastCleanup = now
}
@@ -0,0 +1,186 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package chwriter
import (
"context"
"errors"
"sync"
"testing"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
"Wavelet/pkg/batchwriter"
)
func TestDedupSetMarkIfNew(t *testing.T) {
t.Parallel()
set := newDedupSet()
if !set.markIfNew("node-a|1") {
t.Fatal("markIfNew() = false, want true on first key")
}
if set.markIfNew("node-a|1") {
t.Fatal("markIfNew() = true, want false on duplicate key")
}
if !set.markIfNew("node-b|1") {
t.Fatal("markIfNew() = false, want true on different key")
}
if set.markIfNew("") {
t.Fatal("markIfNew() = true, want false on empty key")
}
}
func TestDedupSetUnmarkAllowsRetry(t *testing.T) {
t.Parallel()
set := newDedupSet()
if !set.markIfNew("k") {
t.Fatal("markIfNew() = false, want true")
}
set.unmark("k")
if !set.markIfNew("k") {
t.Fatal("markIfNew() after unmark = false, want true")
}
}
func TestQueueWithDedupDoesNotMarkWhenEnqueueFails(t *testing.T) {
t.Parallel()
cfg := batchwriter.DefaultConfig()
cfg.QueueSize = 1
cfg.MaxBatchSize = 10
cfg.FlushInterval = time.Hour
// Block the worker so the queue stays full after one enqueue.
block := make(chan struct{})
writer, err := batchwriter.New[int](cfg, func(context.Context, []int) error {
<-block
return nil
})
if err != nil {
t.Fatalf("New() error = %v", err)
}
writer.Start(context.Background())
t.Cleanup(func() {
close(block)
stopCtx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
_ = writer.Stop(stopCtx)
})
// Fill the channel buffer (and the worker's current receive slot may empty one).
// Keep enqueueing until full so subsequent queueWithDedup fails.
for i := 0; i < cfg.QueueSize+2; i++ {
_ = writer.TryEnqueue(i)
if writer.IsFull() {
break
}
}
if !writer.IsFull() {
t.Fatal("writer not full after filling; cannot test enqueue failure path")
}
dedup := newDedupSet()
queueWithDedup(writer, dedup, "dedup-key", 99)
// Key must not remain marked after failed enqueue.
if !dedup.markIfNew("dedup-key") {
t.Fatal("dedup key still marked after failed enqueue; want unmark")
}
}
func TestQueueWithDedupMarksOnlyOnSuccess(t *testing.T) {
t.Parallel()
cfg := batchwriter.DefaultConfig()
cfg.MaxBatchSize = 100
cfg.FlushInterval = time.Hour
writer, err := batchwriter.New[int](cfg, func(context.Context, []int) error { return nil })
if err != nil {
t.Fatalf("New() error = %v", err)
}
writer.Start(context.Background())
t.Cleanup(func() {
stopCtx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
_ = writer.Stop(stopCtx)
})
dedup := newDedupSet()
queueWithDedup(writer, dedup, "ok-key", 1)
if dedup.markIfNew("ok-key") {
t.Fatal("markIfNew() = true after successful enqueue, want false (key marked)")
}
}
func TestFlushErrorHandlerUnmarksKeys(t *testing.T) {
t.Parallel()
dedup := newDedupSet()
flushErr := errors.New("ch down")
var (
mu sync.Mutex
errCount int
)
cfg := batchwriter.Config{
Name: "test_obs",
QueueSize: 10,
MaxBatchSize: 1,
FlushInterval: time.Hour,
}
keyFn := func(s analyticsmodel.NodeMetricSnapshot) string {
return metricSnapshotKey(s)
}
writer, err := batchwriter.New(
cfg,
func(context.Context, []analyticsmodel.NodeMetricSnapshot) error { return flushErr },
batchwriter.WithFlushErrorHandler[analyticsmodel.NodeMetricSnapshot](func(_ context.Context, items []analyticsmodel.NodeMetricSnapshot, err error) {
mu.Lock()
errCount++
mu.Unlock()
for _, item := range items {
dedup.unmark(keyFn(item))
}
}),
)
if err != nil {
t.Fatalf("New() error = %v", err)
}
writer.Start(context.Background())
t.Cleanup(func() {
stopCtx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
_ = writer.Stop(stopCtx)
})
item := analyticsmodel.NodeMetricSnapshot{
NodeID: "n1",
CapturedAt: time.Unix(1, 0).UTC(),
}
key := keyFn(item)
if !dedup.markIfNew(key) {
t.Fatal("markIfNew failed")
}
if !writer.TryEnqueue(item) {
t.Fatal("TryEnqueue failed")
}
deadline := time.Now().Add(time.Second)
for {
mu.Lock()
ready := errCount >= 1
mu.Unlock()
if ready || time.Now().After(deadline) {
break
}
time.Sleep(5 * time.Millisecond)
}
if !dedup.markIfNew(key) {
t.Fatal("key still marked after flush error unmark; want available for retry")
}
}
@@ -0,0 +1,88 @@
//go:build live_ch
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package chwriter_test
import (
"context"
"testing"
"time"
"Wavelet/openflare/plugins/server/kernel/repository"
"Wavelet/openflare/plugins/server/domain/observability/chwriter"
"Wavelet/openflare/plugins/server/kernel/model"
db "Wavelet/plugins/infra/database"
)
// Run with Docker ClickHouse + config.yaml:
//
// go test -tags live_ch ./internal/apps/openflare/chwriter -run TestLiveAppWritePath -count=1 -timeout 2m
func TestLiveAppWritePath(t *testing.T) {
if db.ChConn == nil {
t.Skip("ClickHouse connection not ready")
}
ctx := context.Background()
chwriter.Init(ctx)
now := time.Now().UTC()
nodeID := "e2e-app-write-" + now.Format("150405")
if err := repository.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{
NodeID: nodeID,
CapturedAt: now,
CPUUsagePercent: 33.3,
MemoryUsedBytes: 111,
MemoryTotalBytes: 1000,
StorageUsedBytes: 222,
StorageTotalBytes: 2000,
DiskReadBytes: 10,
DiskWriteBytes: 20,
NetworkRxBytes: 30,
NetworkTxBytes: 40,
}); err != nil {
t.Fatalf("InsertOpenFlareMetricSnapshot: %v", err)
}
deadline := time.Now().Add(45 * time.Second)
var found bool
for time.Now().Before(deadline) {
rows, err := repository.ListOpenFlareMetricSnapshotsSince(ctx, nodeID, now.Add(-time.Minute), 10)
if err != nil {
t.Fatalf("ListOpenFlareMetricSnapshotsSince: %v", err)
}
if len(rows) > 0 {
found = true
t.Logf("found snapshot id=%d cpu=%.1f after flush", rows[0].ID, rows[0].CPUUsagePercent)
break
}
time.Sleep(2 * time.Second)
}
if !found {
t.Fatal("metric snapshot not visible in ClickHouse after flush wait")
}
latest, err := repository.ListOpenFlareLatestMetricSnapshotsSince(ctx, "", now.Add(-time.Hour))
if err != nil {
t.Fatalf("ListOpenFlareLatestMetricSnapshotsSince: %v", err)
}
var latestOK bool
for _, row := range latest {
if row != nil && row.NodeID == nodeID {
latestOK = true
break
}
}
if !latestOK {
t.Fatalf("latest-per-node query missing node %s (rows=%d)", nodeID, len(latest))
}
stats := chwriter.WriterStats()
if len(stats) == 0 {
t.Fatal("WriterStats empty after Init")
}
for _, s := range stats {
t.Logf("writer %s running=%v depth=%d drops=%d flush_err=%d", s.Name, s.Running, s.Depth, s.Drops, s.FlushErrors)
}
}
@@ -0,0 +1,400 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
// Package chwriter queues OpenFlare ClickHouse writes and flushes them through
// internal/infra/persistence/batchwriter with per-table writer instances.
package chwriter
import (
"context"
"fmt"
"sync"
"time"
analyticsmodel "Wavelet/openflare/plugins/server/kernel/model/analytics"
"Wavelet/openflare/plugins/server/kernel/repository/logstore"
"Wavelet/pkg/batchwriter"
"Wavelet/pkg/logger"
)
const (
// Observability traffic is sparse (heartbeat ~10s/node). Prefer larger batches to
// cut ClickHouse parts/merges; MaxFlushWait bounds visibility lag for single-node labs.
observabilityQueueSize = 5_000
observabilityMaxBatchSize = 500
observabilityMinBatchSize = 20
observabilityFlushEvery = 10 * time.Second
observabilityMaxFlushWait = 30 * time.Second
nodeAccessLogQueueSize = 10_000
nodeAccessLogMaxBatchSize = 1_000
nodeAccessLogMinBatchSize = 50
nodeAccessLogFlushEvery = 2 * time.Second
nodeAccessLogMaxFlushWait = 5 * time.Second
// flushAttempts is total tries (1 initial + short retries) before giving up a batch.
flushAttempts = 2
flushRetryBackoff = 50 * time.Millisecond
)
var (
initOnce sync.Once
metricSnapshotWriter *batchwriter.Writer[analyticsmodel.NodeMetricSnapshot]
edgeHealthWriter *batchwriter.Writer[analyticsmodel.NodeEdgeHealth]
frpsWriter *batchwriter.Writer[analyticsmodel.NodeObsFrps]
frpcWriter *batchwriter.Writer[analyticsmodel.NodeObsFrpc]
nodeAccessLogWriter *batchwriter.Writer[analyticsmodel.NodeAccessLog]
metricSnapshotDedup *dedupSet
edgeHealthDedup *dedupSet
frpsDedup *dedupSet
frpcDedup *dedupSet
)
// Init starts OpenFlare log batch writers. Safe to call multiple times.
// Writers always initialize regardless of ClickHouse.enabled; the active log
// store is resolved via logstore at flush time (PG/SQLite when CH is not active).
func Init(ctx context.Context) {
initOnce.Do(func() {
metricSnapshotDedup = newDedupSet()
edgeHealthDedup = newDedupSet()
frpsDedup = newDedupSet()
frpcDedup = newDedupSet()
metricSnapshotWriter = mustNewObservabilityWriter(
"metric_snapshots",
withFlushRetries(flushNodeMetricSnapshots),
metricSnapshotDedup,
metricSnapshotKey,
)
edgeHealthWriter = mustNewObservabilityWriter(
"edge_health",
withFlushRetries(flushNodeEdgeHealth),
edgeHealthDedup,
edgeHealthKey,
)
frpsWriter = mustNewObservabilityWriter(
"frps_obs",
withFlushRetries(flushNodeObsFrps),
frpsDedup,
frpsKey,
)
frpcWriter = mustNewObservabilityWriter(
"frpc_obs",
withFlushRetries(flushNodeObsFrpc),
frpcDedup,
frpcKey,
)
nodeAccessLogWriter = mustNewNodeAccessLogWriter()
metricSnapshotWriter.Start(ctx)
edgeHealthWriter.Start(ctx)
frpsWriter.Start(ctx)
frpcWriter.Start(ctx)
nodeAccessLogWriter.Start(ctx)
wireModelInsertHooks()
})
}
// Stop drains all OpenFlare ClickHouse writers.
func Stop(ctx context.Context) error {
if !running() {
return nil
}
var firstErr error
for _, writer := range []batchStopper{
metricSnapshotWriter,
edgeHealthWriter,
frpsWriter,
frpcWriter,
nodeAccessLogWriter,
} {
if writer == nil {
continue
}
if err := writer.Stop(ctx); err != nil && firstErr == nil {
firstErr = err
}
}
return firstErr
}
// Drain 等待所有 OpenFlare 日志 writer 的在途批次落库:轮询队列 Depth 归零后
// 再保持一个最大 flush 周期(observabilityFlushEvery)持续为空才返回;
// 不停止 writer(迁移冻结后由 ensureWritable 拒绝新写入)。未初始化时直接返回 nil。
func Drain(ctx context.Context) error {
return drainWriters(ctx, WriterStats, observabilityFlushEvery)
}
// drainWriters 轮询 stats 直至所有队列 Depth=0 并持续 quietPeriod 无新积压。
func drainWriters(ctx context.Context, stats func() []batchwriter.Stats, quietPeriod time.Duration) error {
if !running() {
return nil
}
ticker := time.NewTicker(drainPollInterval)
defer ticker.Stop()
var quietSince time.Time
for {
if allDepthZero(stats()) {
if quietSince.IsZero() {
quietSince = time.Now()
} else if time.Since(quietSince) >= quietPeriod {
return nil
}
} else {
quietSince = time.Time{}
}
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
}
}
}
// drainPollInterval 队列轮询间隔。
const drainPollInterval = 50 * time.Millisecond
func allDepthZero(stats []batchwriter.Stats) bool {
for _, s := range stats {
if s.Depth > 0 {
return false
}
}
return true
}
// WriterStats returns queue depth and failure counters for all OpenFlare writers.
func WriterStats() []batchwriter.Stats {
writers := []statsProvider{
metricSnapshotWriter,
edgeHealthWriter,
frpsWriter,
frpcWriter,
nodeAccessLogWriter,
}
out := make([]batchwriter.Stats, 0, len(writers))
for _, w := range writers {
if w == nil {
continue
}
out = append(out, w.Stats())
}
return out
}
// QueueMetricSnapshot enqueues a metric snapshot for asynchronous flush.
func QueueMetricSnapshot(snapshot analyticsmodel.NodeMetricSnapshot) {
queueWithDedup(metricSnapshotWriter, metricSnapshotDedup, metricSnapshotKey(snapshot), snapshot)
}
// QueueEdgeHealth enqueues an L2 edge health snapshot for asynchronous flush.
func QueueEdgeHealth(row analyticsmodel.NodeEdgeHealth) {
queueWithDedup(edgeHealthWriter, edgeHealthDedup, edgeHealthKey(row), row)
}
// QueueFrpsObservation enqueues an FRPS observation for asynchronous flush.
func QueueFrpsObservation(observation analyticsmodel.NodeObsFrps) {
queueWithDedup(frpsWriter, frpsDedup, frpsKey(observation), observation)
}
// QueueFrpcObservation enqueues an FRPC observation for asynchronous flush.
func QueueFrpcObservation(observation analyticsmodel.NodeObsFrpc) {
queueWithDedup(frpcWriter, frpcDedup, frpcKey(observation), observation)
}
// QueueNodeAccessLogs enqueues node access logs for asynchronous flush.
func QueueNodeAccessLogs(logs []analyticsmodel.NodeAccessLog) {
if nodeAccessLogWriter == nil || len(logs) == 0 {
return
}
for _, logItem := range logs {
nodeAccessLogWriter.TryEnqueue(logItem)
}
}
func queueWithDedup[T any](writer *batchwriter.Writer[T], dedup *dedupSet, key string, item T) {
if writer == nil {
return
}
// Mark first so concurrent duplicates still collapse; release on enqueue failure
// so a full queue does not permanently suppress the item.
if !dedup.markIfNew(key) {
return
}
if !writer.TryEnqueue(item) {
dedup.unmark(key)
}
}
func mustNewObservabilityWriter[T any](
name string,
flush batchwriter.FlushFunc[T],
dedup *dedupSet,
keyFn func(T) string,
) *batchwriter.Writer[T] {
cfg := batchwriter.Config{
Name: name,
QueueSize: observabilityQueueSize,
MaxBatchSize: observabilityMaxBatchSize,
MinBatchSize: observabilityMinBatchSize,
FlushInterval: observabilityFlushEvery,
MaxFlushWait: observabilityMaxFlushWait,
}
writer, err := batchwriter.New(
cfg,
flush,
withObservabilityDropHandler[T](name),
batchwriter.WithFlushErrorHandler[T](func(ctx context.Context, items []T, err error) {
logger.ErrorF(ctx, "[OpenFlare] flush %s failed (batch=%d): %v", name, len(items), err)
if dedup == nil || keyFn == nil {
return
}
for _, item := range items {
dedup.unmark(keyFn(item))
}
}),
)
if err != nil {
panic(fmt.Sprintf("openflare chwriter %s: %v", name, err))
}
return writer
}
func mustNewNodeAccessLogWriter() *batchwriter.Writer[analyticsmodel.NodeAccessLog] {
cfg := batchwriter.Config{
Name: "node_access_logs",
QueueSize: nodeAccessLogQueueSize,
MaxBatchSize: nodeAccessLogMaxBatchSize,
MinBatchSize: nodeAccessLogMinBatchSize,
FlushInterval: nodeAccessLogFlushEvery,
MaxFlushWait: nodeAccessLogMaxFlushWait,
}
writer, err := batchwriter.New[analyticsmodel.NodeAccessLog](
cfg,
withFlushRetries(flushNodeAccessLogs),
batchwriter.WithDropHandler[analyticsmodel.NodeAccessLog](func(item analyticsmodel.NodeAccessLog) {
logger.WarnF(context.Background(), "[OpenFlare] node access log queue full, dropping log for node %s path %s", item.NodeID, item.Path)
}),
batchwriter.WithFlushErrorHandler[analyticsmodel.NodeAccessLog](func(ctx context.Context, items []analyticsmodel.NodeAccessLog, err error) {
logger.ErrorF(ctx, "[OpenFlare] flush node access logs failed (batch=%d): %v", len(items), err)
}),
)
if err != nil {
panic(fmt.Sprintf("openflare chwriter node_access_logs: %v", err))
}
return writer
}
func withObservabilityDropHandler[T any](name string) batchwriter.Option[T] {
return batchwriter.WithDropHandler(func(_ T) {
logger.WarnF(context.Background(), "[OpenFlare] %s queue full, dropping observability item", name)
})
}
// withFlushRetries wraps a flush function with a short retry to ride out brief CH blips.
func withFlushRetries[T any](flush batchwriter.FlushFunc[T]) batchwriter.FlushFunc[T] {
return func(ctx context.Context, items []T) error {
var err error
for attempt := 1; attempt <= flushAttempts; attempt++ {
err = flush(ctx, items)
if err == nil {
return nil
}
if attempt == flushAttempts {
break
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(flushRetryBackoff * time.Duration(attempt)):
}
}
return err
}
}
func wireModelInsertHooks() {
logstore.SetObservabilityHooks(logstore.ObservabilityHooks{
QueueMetricSnapshot: QueueMetricSnapshot,
QueueEdgeHealth: QueueEdgeHealth,
QueueNodeObsFrps: QueueFrpsObservation,
QueueNodeObsFrpc: QueueFrpcObservation,
})
logstore.SetAccessLogHooks(logstore.AccessLogHooks{
QueueNodeAccessLogs: QueueNodeAccessLogs,
})
}
// 以下 flush 函数作为 batchwriter 的落库目标:激活库由 logstore 在 flush 时决定。
func flushNodeMetricSnapshots(ctx context.Context, rows []analyticsmodel.NodeMetricSnapshot) error {
s, err := logstore.Active(ctx)
if err != nil {
return err
}
return s.Observability.BatchInsertNodeMetricSnapshots(ctx, rows)
}
func flushNodeEdgeHealth(ctx context.Context, rows []analyticsmodel.NodeEdgeHealth) error {
s, err := logstore.Active(ctx)
if err != nil {
return err
}
return s.Observability.BatchInsertNodeEdgeHealth(ctx, rows)
}
func flushNodeObsFrps(ctx context.Context, rows []analyticsmodel.NodeObsFrps) error {
s, err := logstore.Active(ctx)
if err != nil {
return err
}
return s.Observability.BatchInsertNodeObsFrps(ctx, rows)
}
func flushNodeObsFrpc(ctx context.Context, rows []analyticsmodel.NodeObsFrpc) error {
s, err := logstore.Active(ctx)
if err != nil {
return err
}
return s.Observability.BatchInsertNodeObsFrpc(ctx, rows)
}
func flushNodeAccessLogs(ctx context.Context, rows []analyticsmodel.NodeAccessLog) error {
s, err := logstore.Active(ctx)
if err != nil {
return err
}
return s.AccessLogs.BatchInsertNodeAccessLogs(ctx, rows)
}
func metricSnapshotKey(snapshot analyticsmodel.NodeMetricSnapshot) string {
return fmt.Sprintf("%s|%d", snapshot.NodeID, snapshot.CapturedAt.UTC().UnixNano())
}
func edgeHealthKey(row analyticsmodel.NodeEdgeHealth) string {
return fmt.Sprintf("%s|%d", row.NodeID, row.CapturedAt.UTC().UnixNano())
}
func frpsKey(observation analyticsmodel.NodeObsFrps) string {
return fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano())
}
func frpcKey(observation analyticsmodel.NodeObsFrpc) string {
return fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano())
}
type batchStopper interface {
Stop(ctx context.Context) error
}
type statsProvider interface {
Stats() batchwriter.Stats
}
func running() bool {
return metricSnapshotWriter != nil && metricSnapshotWriter.Running()
}