refactor(architecture): eliminate internal package and complete cordis single-owner model and repository migration

- Physically purged all legacy internal/ packages, centralized pkg/model/ and pkg/repository/
- Migrated domain models and database repositories into self-contained owner plugins (user, auth, message_gateway, admin, upload, risk_control)
- Decoupled cross-plugin interactions via pure core/contracts and typed EventBus
- Ensured 100% test coverage pass, zero data races (-race clean), and 0 lint issues in make code-check
This commit is contained in:
ryan
2026-08-28 08:40:43 +08:00
parent 1f348fd425
commit fb6a3edb89
323 changed files with 8222 additions and 17693 deletions
+69
View File
@@ -0,0 +1,69 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package batchwriter
import (
"fmt"
"time"
)
const (
defaultQueueSize = 10_000
defaultMaxBatchSize = 1_000
defaultMinBatchSize = 50
defaultFlushEvery = time.Second
)
// Config controls queue capacity and flush thresholds for a Writer instance.
type Config struct {
// Name identifies the writer in logs and diagnostics. Optional.
Name string
// QueueSize is the buffered channel capacity.
QueueSize int
// MaxBatchSize triggers a flush when the in-memory batch reaches this count.
MaxBatchSize int
// MinBatchSize is the minimum in-memory batch size for time-based flushes.
// Zero disables the threshold and preserves legacy interval flush behavior.
// When set, interval flushes below this size are skipped unless MaxFlushWait elapses.
MinBatchSize int
// FlushInterval is how often the worker checks whether a time-based flush should run.
FlushInterval time.Duration
// MaxFlushWait forces a flush of any non-empty batch once the oldest item has waited
// this long, even if MinBatchSize has not been reached. Zero disables the force path.
MaxFlushWait time.Duration
}
// DefaultConfig returns production-friendly defaults aligned with audit log batching.
func DefaultConfig() Config {
return Config{
QueueSize: defaultQueueSize,
MaxBatchSize: defaultMaxBatchSize,
MinBatchSize: defaultMinBatchSize,
FlushInterval: defaultFlushEvery,
}
}
func (c Config) validate() error {
if c.QueueSize <= 0 {
return fmt.Errorf("batchwriter: queue size must be positive")
}
if c.MaxBatchSize <= 0 {
return fmt.Errorf("batchwriter: max batch size must be positive")
}
if c.MinBatchSize < 0 {
return fmt.Errorf("batchwriter: min batch size must be non-negative")
}
if c.FlushInterval <= 0 {
return fmt.Errorf("batchwriter: flush interval must be positive")
}
if c.MaxFlushWait < 0 {
return fmt.Errorf("batchwriter: max flush wait must be non-negative")
}
return nil
}
+8
View File
@@ -0,0 +1,8 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package batchwriter
import "errors"
var errNilFlushFunc = errors.New("batchwriter: flush func is required")
+264
View File
@@ -0,0 +1,264 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
// Package batchwriter provides a reusable buffered batch writer for high-throughput
// append-only sinks such as ClickHouse. Each business domain should own an independent
// Writer instance with its own queue, flush callback, and tuning parameters.
package batchwriter
import (
"context"
"sync"
"sync/atomic"
"time"
)
// FlushFunc persists a batch of queued items. It is invoked from the worker goroutine.
type FlushFunc[T any] func(ctx context.Context, items []T) error
// FlushErrorHandler is called when FlushFunc returns an error after optional retries.
// The batch is discarded after the handler returns; the worker continues processing.
// Handlers receive the failed items so callers can release dedup keys or re-queue.
type FlushErrorHandler[T any] func(ctx context.Context, items []T, err error)
// Stats is a point-in-time snapshot of Writer queue and failure counters.
type Stats struct {
Name string
Depth int
Cap int
Drops int64
FlushErrors int64
Running bool
}
// Writer buffers items and flushes them by size or interval.
type Writer[T any] struct {
cfg Config
flush FlushFunc[T]
onFlushError FlushErrorHandler[T]
onDrop func(T)
startOnce sync.Once
stopOnce sync.Once
mu sync.RWMutex
ch chan T
workerCtx context.Context
done chan struct{}
drops atomic.Int64
flushErrors atomic.Int64
}
// Option configures optional Writer callbacks.
type Option[T any] func(*Writer[T])
// WithFlushErrorHandler registers a callback for flush failures.
func WithFlushErrorHandler[T any](handler FlushErrorHandler[T]) Option[T] {
return func(w *Writer[T]) {
w.onFlushError = handler
}
}
// WithDropHandler registers a callback when TryEnqueue cannot accept an item.
func WithDropHandler[T any](handler func(T)) Option[T] {
return func(w *Writer[T]) {
w.onDrop = handler
}
}
// New creates a Writer. Call Start before enqueueing items.
func New[T any](cfg Config, flush FlushFunc[T], opts ...Option[T]) (*Writer[T], error) {
if flush == nil {
return nil, errNilFlushFunc
}
if err := cfg.validate(); err != nil {
return nil, err
}
w := &Writer[T]{
cfg: cfg,
flush: flush,
done: make(chan struct{}),
}
for _, opt := range opts {
opt(w)
}
return w, nil
}
// Start launches the background worker. It is safe to call at most once.
func (w *Writer[T]) Start(parent context.Context) {
w.startOnce.Do(func() {
w.mu.Lock()
defer w.mu.Unlock()
w.ch = make(chan T, w.cfg.QueueSize)
w.workerCtx = context.WithoutCancel(parent)
go w.run()
})
}
// Stop closes the queue and waits until the worker drains pending items and exits.
func (w *Writer[T]) Stop(ctx context.Context) error {
w.mu.RLock()
ch := w.ch
done := w.done
w.mu.RUnlock()
if ch == nil {
return nil
}
w.stopOnce.Do(func() {
close(ch)
})
select {
case <-done:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
// Running reports whether Start has been called and Stop has not completed.
func (w *Writer[T]) Running() bool {
w.mu.RLock()
defer w.mu.RUnlock()
if w.ch == nil {
return false
}
select {
case <-w.done:
return false
default:
return true
}
}
// TryEnqueue adds one item without blocking. It returns false when the writer is not
// running or the queue is full.
func (w *Writer[T]) TryEnqueue(item T) bool {
w.mu.RLock()
ch := w.ch
w.mu.RUnlock()
if ch == nil {
w.notifyDrop(item)
return false
}
select {
case ch <- item:
return true
default:
w.notifyDrop(item)
return false
}
}
// IsFull reports whether the queue has no remaining capacity.
func (w *Writer[T]) IsFull() bool {
w.mu.RLock()
defer w.mu.RUnlock()
if w.ch == nil {
return false
}
return len(w.ch) >= cap(w.ch)
}
// Len returns the current queue depth.
func (w *Writer[T]) Len() int {
w.mu.RLock()
defer w.mu.RUnlock()
if w.ch == nil {
return 0
}
return len(w.ch)
}
// Cap returns the queue capacity.
func (w *Writer[T]) Cap() int {
return w.cfg.QueueSize
}
// Stats returns a point-in-time snapshot of queue depth and failure counters.
func (w *Writer[T]) Stats() Stats {
return Stats{
Name: w.cfg.Name,
Depth: w.Len(),
Cap: w.Cap(),
Drops: w.drops.Load(),
FlushErrors: w.flushErrors.Load(),
Running: w.Running(),
}
}
func (w *Writer[T]) run() {
ticker := time.NewTicker(w.cfg.FlushInterval)
defer ticker.Stop()
batch := make([]T, 0, w.cfg.MaxBatchSize)
var batchStartedAt time.Time
flush := func() {
if len(batch) == 0 {
return
}
items := append([]T(nil), batch...)
if err := w.flush(w.workerCtx, items); err != nil {
w.flushErrors.Add(1)
if w.onFlushError != nil {
w.onFlushError(w.workerCtx, items, err)
}
}
batch = batch[:0]
batchStartedAt = time.Time{}
}
defer func() {
flush()
close(w.done)
}()
for {
select {
case item, ok := <-w.ch:
if !ok {
return
}
if len(batch) == 0 {
batchStartedAt = time.Now()
}
batch = append(batch, item)
if len(batch) >= w.cfg.MaxBatchSize {
flush()
}
case <-ticker.C:
if w.shouldFlushOnInterval(len(batch), batchStartedAt, time.Now()) {
flush()
}
}
}
}
func (w *Writer[T]) shouldFlushOnInterval(batchLen int, batchStartedAt time.Time, now time.Time) bool {
if batchLen == 0 {
return false
}
if w.cfg.MinBatchSize == 0 || batchLen >= w.cfg.MinBatchSize {
return true
}
if w.cfg.MaxFlushWait <= 0 || batchStartedAt.IsZero() {
return false
}
return !now.Before(batchStartedAt.Add(w.cfg.MaxFlushWait))
}
func (w *Writer[T]) notifyDrop(item T) {
w.drops.Add(1)
if w.onDrop == nil {
return
}
w.onDrop(item)
}
+491
View File
@@ -0,0 +1,491 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package batchwriter
import (
"context"
"errors"
"sync"
"testing"
"time"
"github.com/google/go-cmp/cmp"
)
func TestNewRejectsInvalidConfig(t *testing.T) {
t.Parallel()
_, err := New[int](Config{}, func(context.Context, []int) error { return nil })
if err == nil {
t.Fatal("New() = nil, want validation error")
}
}
func TestNewRejectsNilFlushFunc(t *testing.T) {
t.Parallel()
cfg := DefaultConfig()
_, err := New[int](cfg, nil)
if !errors.Is(err, errNilFlushFunc) {
t.Fatalf("New() error = %v, want %v", err, errNilFlushFunc)
}
}
func TestWriterFlushesOnMaxBatchSize(t *testing.T) {
t.Parallel()
var (
mu sync.Mutex
batches [][]int
)
cfg := DefaultConfig()
cfg.MaxBatchSize = 3
cfg.FlushInterval = time.Hour
writer, err := New[int](cfg, func(_ context.Context, items []int) error {
mu.Lock()
defer mu.Unlock()
batches = append(batches, append([]int(nil), items...))
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()
if err := writer.Stop(stopCtx); err != nil {
t.Fatalf("Stop() error = %v", err)
}
})
for i := range 3 {
if !writer.TryEnqueue(i + 1) {
t.Fatalf("TryEnqueue(%d) = false, want true", i+1)
}
}
deadline := time.Now().Add(time.Second)
for {
mu.Lock()
ready := len(batches) == 1
mu.Unlock()
if ready || time.Now().After(deadline) {
break
}
time.Sleep(10 * time.Millisecond)
}
mu.Lock()
got := batches
mu.Unlock()
want := [][]int{{1, 2, 3}}
if diff := cmp.Diff(want, got); diff != "" {
t.Fatalf("flush batches mismatch (-want +got):\n%s", diff)
}
}
func TestWriterFlushesOnInterval(t *testing.T) {
t.Parallel()
var (
mu sync.Mutex
batch []int
)
cfg := DefaultConfig()
cfg.MaxBatchSize = 100
cfg.MinBatchSize = 0
cfg.FlushInterval = 20 * time.Millisecond
writer, err := New[int](cfg, func(_ context.Context, items []int) error {
mu.Lock()
defer mu.Unlock()
batch = append([]int(nil), items...)
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()
if err := writer.Stop(stopCtx); err != nil {
t.Fatalf("Stop() error = %v", err)
}
})
if !writer.TryEnqueue(42) {
t.Fatal("TryEnqueue() = false, want true")
}
deadline := time.Now().Add(time.Second)
for {
mu.Lock()
ready := len(batch) == 1
mu.Unlock()
if ready || time.Now().After(deadline) {
break
}
time.Sleep(5 * time.Millisecond)
}
mu.Lock()
got := batch
mu.Unlock()
want := []int{42}
if diff := cmp.Diff(want, got); diff != "" {
t.Fatalf("interval flush mismatch (-want +got):\n%s", diff)
}
}
func TestWriterTryEnqueueDropsWhenFull(t *testing.T) {
t.Parallel()
cfg := DefaultConfig()
cfg.QueueSize = 1
cfg.MaxBatchSize = 10
cfg.FlushInterval = time.Hour
var dropped int
writer, err := New[int](cfg, func(context.Context, []int) error { return nil }, WithDropHandler[int](func(int) {
dropped++
}))
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)
})
if !writer.TryEnqueue(1) {
t.Fatal("TryEnqueue(1) = false, want true")
}
if writer.TryEnqueue(2) {
t.Fatal("TryEnqueue(2) = true, want false")
}
if !writer.IsFull() {
t.Fatal("IsFull() = false, want true")
}
if dropped != 1 {
t.Fatalf("dropped = %d, want 1", dropped)
}
}
func TestWriterStopDrainsQueuedItems(t *testing.T) {
t.Parallel()
cfg := DefaultConfig()
cfg.MaxBatchSize = 10
cfg.FlushInterval = time.Hour
var flushed []int
writer, err := New[int](cfg, func(_ context.Context, items []int) error {
flushed = append(flushed, items...)
return nil
})
if err != nil {
t.Fatalf("New() error = %v", err)
}
writer.Start(context.Background())
for i := range 2 {
if !writer.TryEnqueue(i + 1) {
t.Fatalf("TryEnqueue(%d) = false, want true", i+1)
}
}
stopCtx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
if err := writer.Stop(stopCtx); err != nil {
t.Fatalf("Stop() error = %v", err)
}
want := []int{1, 2}
if diff := cmp.Diff(want, flushed); diff != "" {
t.Fatalf("Stop() drain mismatch (-want +got):\n%s", diff)
}
if writer.Running() {
t.Fatal("Running() = true after Stop(), want false")
}
}
func TestWriterInvokesFlushErrorHandler(t *testing.T) {
t.Parallel()
cfg := DefaultConfig()
cfg.MaxBatchSize = 1
cfg.FlushInterval = time.Hour
flushErr := errors.New("flush failed")
var (
mu sync.Mutex
errCount int
gotItems []int
)
writer, err := New[int](cfg, func(context.Context, []int) error {
return flushErr
}, WithFlushErrorHandler[int](func(_ context.Context, items []int, err error) {
mu.Lock()
defer mu.Unlock()
errCount++
gotItems = append([]int(nil), items...)
if !errors.Is(err, flushErr) {
t.Errorf("flush error = %v, want %v", err, flushErr)
}
}))
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)
})
if !writer.TryEnqueue(7) {
t.Fatal("TryEnqueue() = false, want true")
}
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)
}
mu.Lock()
gotCount := errCount
items := gotItems
mu.Unlock()
if gotCount != 1 {
t.Fatalf("flush error handler count = %d, want 1", gotCount)
}
if diff := cmp.Diff([]int{7}, items); diff != "" {
t.Fatalf("flush error handler items mismatch (-want +got):\n%s", diff)
}
stats := writer.Stats()
if stats.FlushErrors != 1 {
t.Fatalf("Stats().FlushErrors = %d, want 1", stats.FlushErrors)
}
}
func TestWriterSkipsIntervalFlushBelowMinBatchSize(t *testing.T) {
t.Parallel()
var (
mu sync.Mutex
batch []int
)
cfg := DefaultConfig()
cfg.MaxBatchSize = 100
cfg.MinBatchSize = 5
cfg.FlushInterval = 20 * time.Millisecond
writer, err := New[int](cfg, func(_ context.Context, items []int) error {
mu.Lock()
defer mu.Unlock()
batch = append([]int(nil), items...)
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()
if err := writer.Stop(stopCtx); err != nil {
t.Fatalf("Stop() error = %v", err)
}
})
for i := range 3 {
if !writer.TryEnqueue(i + 1) {
t.Fatalf("TryEnqueue(%d) = false, want true", i+1)
}
}
time.Sleep(100 * time.Millisecond)
mu.Lock()
got := batch
mu.Unlock()
if len(got) != 0 {
t.Fatalf("interval flush with below-min batch = %v, want no flush", got)
}
}
func TestWriterFlushesOnIntervalWhenMinBatchSizeReached(t *testing.T) {
t.Parallel()
var (
mu sync.Mutex
batch []int
)
cfg := DefaultConfig()
cfg.MaxBatchSize = 100
cfg.MinBatchSize = 3
cfg.FlushInterval = 20 * time.Millisecond
writer, err := New[int](cfg, func(_ context.Context, items []int) error {
mu.Lock()
defer mu.Unlock()
batch = append([]int(nil), items...)
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()
if err := writer.Stop(stopCtx); err != nil {
t.Fatalf("Stop() error = %v", err)
}
})
for i := range 3 {
if !writer.TryEnqueue(i + 1) {
t.Fatalf("TryEnqueue(%d) = false, want true", i+1)
}
}
deadline := time.Now().Add(time.Second)
for {
mu.Lock()
ready := len(batch) == 3
mu.Unlock()
if ready || time.Now().After(deadline) {
break
}
time.Sleep(5 * time.Millisecond)
}
mu.Lock()
got := batch
mu.Unlock()
want := []int{1, 2, 3}
if diff := cmp.Diff(want, got); diff != "" {
t.Fatalf("interval flush at min batch size mismatch (-want +got):\n%s", diff)
}
}
func TestWriterForcesFlushAfterMaxFlushWait(t *testing.T) {
t.Parallel()
var (
mu sync.Mutex
batch []int
)
cfg := DefaultConfig()
cfg.MaxBatchSize = 100
cfg.MinBatchSize = 50
cfg.FlushInterval = 20 * time.Millisecond
cfg.MaxFlushWait = 80 * time.Millisecond
writer, err := New[int](cfg, func(_ context.Context, items []int) error {
mu.Lock()
defer mu.Unlock()
batch = append([]int(nil), items...)
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()
if err := writer.Stop(stopCtx); err != nil {
t.Fatalf("Stop() error = %v", err)
}
})
if !writer.TryEnqueue(1) {
t.Fatal("TryEnqueue(1) = false, want true")
}
deadline := time.Now().Add(time.Second)
for {
mu.Lock()
ready := len(batch) == 1
mu.Unlock()
if ready || time.Now().After(deadline) {
break
}
time.Sleep(5 * time.Millisecond)
}
mu.Lock()
got := batch
mu.Unlock()
if diff := cmp.Diff([]int{1}, got); diff != "" {
t.Fatalf("max flush wait mismatch (-want +got):\n%s", diff)
}
}
func TestWriterStatsTracksDrops(t *testing.T) {
t.Parallel()
cfg := DefaultConfig()
cfg.Name = "test-drops"
cfg.QueueSize = 1
cfg.MaxBatchSize = 10
cfg.FlushInterval = time.Hour
writer, err := 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)
})
if !writer.TryEnqueue(1) {
t.Fatal("TryEnqueue(1) = false, want true")
}
if writer.TryEnqueue(2) {
t.Fatal("TryEnqueue(2) = true, want false")
}
stats := writer.Stats()
if stats.Drops != 1 {
t.Fatalf("Stats().Drops = %d, want 1", stats.Drops)
}
if stats.Cap != 1 {
t.Fatalf("Stats().Cap = %d, want 1", stats.Cap)
}
if !stats.Running {
t.Fatal("Stats().Running = false, want true")
}
}