mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-04 07:06:36 +08:00
refactor(structure): group platform, infra, and shared packages
Move process wiring, technical adapters, and cross-cutting contracts out of flat internal/ packages so new code has a clear home without changing business layout.
This commit is contained in:
@@ -0,0 +1,52 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package batchwriter
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
defaultQueueSize = 10_000
|
||||
defaultMaxBatchSize = 1_000
|
||||
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
|
||||
|
||||
// FlushInterval triggers a time-based flush even when the batch is smaller.
|
||||
FlushInterval time.Duration
|
||||
}
|
||||
|
||||
// DefaultConfig returns production-friendly defaults aligned with audit log batching.
|
||||
func DefaultConfig() Config {
|
||||
return Config{
|
||||
QueueSize: defaultQueueSize,
|
||||
MaxBatchSize: defaultMaxBatchSize,
|
||||
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.FlushInterval <= 0 {
|
||||
return fmt.Errorf("batchwriter: flush interval must be positive")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -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")
|
||||
@@ -0,0 +1,215 @@
|
||||
// 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"
|
||||
"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. The batch is discarded
|
||||
// after the handler returns; the worker continues processing.
|
||||
type FlushErrorHandler func(ctx context.Context, batchSize int, err error)
|
||||
|
||||
// Writer buffers items and flushes them by size or interval.
|
||||
type Writer[T any] struct {
|
||||
cfg Config
|
||||
flush FlushFunc[T]
|
||||
|
||||
onFlushError FlushErrorHandler
|
||||
onDrop func(T)
|
||||
|
||||
startOnce sync.Once
|
||||
stopOnce sync.Once
|
||||
|
||||
mu sync.RWMutex
|
||||
ch chan T
|
||||
workerCtx context.Context
|
||||
done chan struct{}
|
||||
}
|
||||
|
||||
// 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) 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
|
||||
}
|
||||
|
||||
func (w *Writer[T]) run() {
|
||||
ticker := time.NewTicker(w.cfg.FlushInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
batch := make([]T, 0, w.cfg.MaxBatchSize)
|
||||
flush := func() {
|
||||
if len(batch) == 0 {
|
||||
return
|
||||
}
|
||||
items := append([]T(nil), batch...)
|
||||
if err := w.flush(w.workerCtx, items); err != nil {
|
||||
if w.onFlushError != nil {
|
||||
w.onFlushError(w.workerCtx, len(items), err)
|
||||
}
|
||||
}
|
||||
batch = batch[:0]
|
||||
}
|
||||
|
||||
defer func() {
|
||||
flush()
|
||||
close(w.done)
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case item, ok := <-w.ch:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
batch = append(batch, item)
|
||||
if len(batch) >= w.cfg.MaxBatchSize {
|
||||
flush()
|
||||
}
|
||||
case <-ticker.C:
|
||||
flush()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (w *Writer[T]) notifyDrop(item T) {
|
||||
if w.onDrop == nil {
|
||||
return
|
||||
}
|
||||
w.onDrop(item)
|
||||
}
|
||||
@@ -0,0 +1,284 @@
|
||||
// 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.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
|
||||
batchSize int
|
||||
)
|
||||
|
||||
writer, err := New[int](cfg, func(context.Context, []int) error {
|
||||
return flushErr
|
||||
}, WithFlushErrorHandler[int](func(_ context.Context, size int, err error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
errCount++
|
||||
batchSize = size
|
||||
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
|
||||
gotSize := batchSize
|
||||
mu.Unlock()
|
||||
|
||||
if gotCount != 1 {
|
||||
t.Fatalf("flush error handler count = %d, want 1", gotCount)
|
||||
}
|
||||
if gotSize != 1 {
|
||||
t.Fatalf("flush error handler batch size = %d, want 1", gotSize)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package db 提供数据库连接与基础设施
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2"
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
clickhouseDriver "gorm.io/driver/clickhouse"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/plugin/opentelemetry/tracing"
|
||||
)
|
||||
|
||||
const (
|
||||
clickhouseMaxExecTime = 60 // ClickHouse 最大执行时间(秒)
|
||||
clickhouseReadTimeoutFactor = 2 // ReadTimeout 为 DialTimeout 的倍数
|
||||
)
|
||||
|
||||
var (
|
||||
// ChConn ClickHouse 原生连接实例,用于批量写入
|
||||
ChConn driver.Conn
|
||||
|
||||
chDB *gorm.DB
|
||||
)
|
||||
|
||||
func init() {
|
||||
if !config.Config.ClickHouse.Enabled {
|
||||
return
|
||||
}
|
||||
|
||||
cfg := config.Config.ClickHouse
|
||||
if cfg.Database == "" {
|
||||
log.Fatalf("[ClickHouse] database name is required (expected: wavelet)\n")
|
||||
}
|
||||
|
||||
opts := buildClickHouseOptions()
|
||||
|
||||
var err error
|
||||
ChConn, err = clickhouse.Open(opts)
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] init connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
if err = ChConn.Ping(context.Background()); err != nil {
|
||||
log.Fatalf("[ClickHouse] ping failed: %v\n", err)
|
||||
}
|
||||
|
||||
chDB, err = gorm.Open(clickhouseDriver.New(clickhouseDriver.Config{
|
||||
DSN: buildClickHouseDSN(),
|
||||
}), &gorm.Config{
|
||||
SkipDefaultTransaction: true,
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] init gorm connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
if err = chDB.Use(
|
||||
tracing.NewPlugin(
|
||||
tracing.WithoutMetrics(),
|
||||
tracing.WithAttributes(
|
||||
attribute.String("db.instance", cfg.Database),
|
||||
attribute.String("db.system", "ClickHouse"),
|
||||
),
|
||||
),
|
||||
); err != nil {
|
||||
log.Fatalf("[ClickHouse] init trace failed: %v\n", err)
|
||||
}
|
||||
|
||||
sqlDB, err := chDB.DB()
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] load sql db failed: %v\n", err)
|
||||
}
|
||||
|
||||
sqlDB.SetMaxIdleConns(cfg.MaxIdleConn)
|
||||
sqlDB.SetMaxOpenConns(cfg.MaxOpenConn)
|
||||
sqlDB.SetConnMaxLifetime(time.Duration(cfg.ConnMaxLifetime) * time.Second)
|
||||
|
||||
log.Println("[ClickHouse] connection established successfully")
|
||||
}
|
||||
|
||||
func buildClickHouseOptions() *clickhouse.Options {
|
||||
cfg := config.Config.ClickHouse
|
||||
|
||||
return &clickhouse.Options{
|
||||
Addr: cfg.Hosts,
|
||||
Auth: clickhouse.Auth{
|
||||
Database: cfg.Database,
|
||||
Username: cfg.Username,
|
||||
Password: cfg.Password,
|
||||
},
|
||||
Settings: clickhouse.Settings{
|
||||
"max_execution_time": clickhouseMaxExecTime,
|
||||
},
|
||||
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,
|
||||
ReadTimeout: time.Duration(cfg.DialTimeout*clickhouseReadTimeoutFactor) * time.Second,
|
||||
BlockBufferSize: cfg.BlockBufferSize,
|
||||
}
|
||||
}
|
||||
|
||||
func buildClickHouseDSN() string {
|
||||
cfg := config.Config.ClickHouse
|
||||
|
||||
chURL := &url.URL{
|
||||
Scheme: "clickhouse",
|
||||
Host: strings.Join(cfg.Hosts, ","),
|
||||
Path: "/" + cfg.Database,
|
||||
}
|
||||
if cfg.Username != "" || cfg.Password != "" {
|
||||
chURL.User = url.UserPassword(cfg.Username, cfg.Password)
|
||||
}
|
||||
|
||||
query := chURL.Query()
|
||||
query.Set("dial_timeout", fmt.Sprintf("%ds", cfg.DialTimeout))
|
||||
query.Set("read_timeout", fmt.Sprintf("%ds", cfg.DialTimeout*clickhouseReadTimeoutFactor))
|
||||
query.Set("max_execution_time", strconv.Itoa(clickhouseMaxExecTime))
|
||||
chURL.RawQuery = query.Encode()
|
||||
|
||||
return chURL.String()
|
||||
}
|
||||
|
||||
// ChDB returns a context-aware GORM ClickHouse instance.
|
||||
func ChDB(ctx context.Context) *gorm.DB {
|
||||
if chDB == nil {
|
||||
return nil
|
||||
}
|
||||
return chDB.WithContext(ctx)
|
||||
}
|
||||
|
||||
// SetChDBForTest sets the package-level ClickHouse GORM instance for testing.
|
||||
func SetChDBForTest(d *gorm.DB) {
|
||||
chDB = d
|
||||
}
|
||||
|
||||
// SetChConnForTest sets the package-level native ClickHouse connection for testing.
|
||||
func SetChConnForTest(c driver.Conn) {
|
||||
ChConn = c
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package db
|
||||
|
||||
const (
|
||||
errRedisHashSetFailed = "failed to set redis hash: %w"
|
||||
errRedisHashDeleteFailed = "failed to delete redis hash field: %w"
|
||||
errUnmarshalDataFailed = "failed to unmarshal data: %w"
|
||||
errMarshalDataFailed = "failed to marshal data: %w"
|
||||
errRedisKeySetFailed = "failed to set redis key: %w"
|
||||
)
|
||||
@@ -0,0 +1,46 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package idgen 提供分布式 ID 生成器
|
||||
package idgen
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/bwmarrin/snowflake"
|
||||
)
|
||||
|
||||
// 2025-12-01 00:00:00 UTC 的毫秒时间戳
|
||||
const epoch int64 = 1764547200000
|
||||
|
||||
const maxNegativeIDRetries = 3
|
||||
|
||||
var node *snowflake.Node
|
||||
|
||||
func init() {
|
||||
snowflake.Epoch = epoch
|
||||
|
||||
nodeID := config.Config.App.NodeID
|
||||
var err error
|
||||
node, err = snowflake.NewNode(nodeID)
|
||||
if err != nil {
|
||||
log.Fatalf("[Snowflake] init failed: %v\n", err)
|
||||
}
|
||||
log.Printf("[Snowflake] initialized with node ID: %d, epoch: 2025-12-01\n", nodeID)
|
||||
}
|
||||
|
||||
// NextUint64ID 生成下一个分布式唯一 ID。
|
||||
// 理论上不应出现负值;若出现则最多重试 maxNegativeIDRetries 次,仍失败则 panic。
|
||||
func NextUint64ID() uint64 {
|
||||
for attempt := 1; attempt <= maxNegativeIDRetries; attempt++ {
|
||||
id := node.Generate().Int64()
|
||||
if id >= 0 {
|
||||
return uint64(id)
|
||||
}
|
||||
log.Printf("[Snowflake] generated negative ID: %d (attempt %d/%d)", id, attempt, maxNegativeIDRetries)
|
||||
}
|
||||
panic(fmt.Sprintf("[Snowflake] generated negative ID after %d attempts", maxNegativeIDRetries))
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package idgen
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestNextUint64ID(t *testing.T) {
|
||||
id := NextUint64ID()
|
||||
assert.NotZero(t, id)
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package migrator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"embed"
|
||||
"io/fs"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/pressly/goose/v3"
|
||||
)
|
||||
|
||||
const (
|
||||
clickhouseMigrationDir = "goose/clickhouse"
|
||||
clickhouseGooseVersionTable = "goose_clickhouse_version"
|
||||
clickhouseMaxExecTime = 60
|
||||
clickhouseReadTimeoutFactor = 2
|
||||
)
|
||||
|
||||
// clickhouseMigrationFS contains SQL migrations under goose/clickhouse.
|
||||
//
|
||||
//go:embed goose/clickhouse/*.sql
|
||||
var clickhouseMigrationFS embed.FS
|
||||
|
||||
// MigrateClickHouse runs goose migrations against ClickHouse when enabled.
|
||||
func MigrateClickHouse() Report {
|
||||
if !config.Config.ClickHouse.Enabled {
|
||||
return Report{Backend: "ClickHouse"}
|
||||
}
|
||||
|
||||
cfg := config.Config.ClickHouse
|
||||
sqlDB := clickhouse.OpenDB(&clickhouse.Options{
|
||||
Addr: cfg.Hosts,
|
||||
Auth: clickhouse.Auth{
|
||||
Database: cfg.Database,
|
||||
Username: cfg.Username,
|
||||
Password: cfg.Password,
|
||||
},
|
||||
Settings: clickhouse.Settings{
|
||||
"max_execution_time": clickhouseMaxExecTime,
|
||||
},
|
||||
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,
|
||||
ReadTimeout: time.Duration(cfg.DialTimeout*clickhouseReadTimeoutFactor) * time.Second,
|
||||
BlockBufferSize: cfg.BlockBufferSize,
|
||||
})
|
||||
|
||||
subFS, err := fs.Sub(clickhouseMigrationFS, "goose/clickhouse")
|
||||
if err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] get sub fs failed: %v\n", err)
|
||||
}
|
||||
|
||||
provider, err := goose.NewProvider(
|
||||
"clickhouse",
|
||||
sqlDB,
|
||||
subFS,
|
||||
goose.WithTableName(clickhouseGooseVersionTable),
|
||||
goose.WithDisableGlobalRegistry(true),
|
||||
)
|
||||
if err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] create goose provider failed: %v\n", err)
|
||||
}
|
||||
previousVersion, err := provider.GetDBVersion(context.Background())
|
||||
if err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] get goose version failed: %v\n", err)
|
||||
}
|
||||
|
||||
if _, err := provider.Up(context.Background()); err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] goose migrate failed: %v\n", err)
|
||||
}
|
||||
currentVersion, err := provider.GetDBVersion(context.Background())
|
||||
if err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] get migrated goose version failed: %v\n", err)
|
||||
}
|
||||
closeClickHouseDB(sqlDB)
|
||||
|
||||
log.Println("[ClickHouse] goose migrate success")
|
||||
return Report{
|
||||
Backend: "ClickHouse",
|
||||
Enabled: true,
|
||||
Version: currentVersion,
|
||||
Applied: currentVersion != previousVersion,
|
||||
}
|
||||
}
|
||||
|
||||
func closeClickHouseDB(sqlDB *sql.DB) {
|
||||
if err := sqlDB.Close(); err != nil {
|
||||
log.Printf("[ClickHouse] close sql db failed: %v\n", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package migrator
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/pressly/goose/v3"
|
||||
)
|
||||
|
||||
func TestClickHouseMigrationFilesEmbedded(t *testing.T) {
|
||||
entries, err := clickhouseMigrationFS.ReadDir(clickhouseMigrationDir)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadDir(%q) error = %v", clickhouseMigrationDir, err)
|
||||
}
|
||||
if len(entries) == 0 {
|
||||
t.Fatal("expected embedded ClickHouse migrations, got none")
|
||||
}
|
||||
|
||||
found := false
|
||||
for _, entry := range entries {
|
||||
if entry.IsDir() {
|
||||
continue
|
||||
}
|
||||
if entry.Name() == "202606190001_create_user_access_logs.sql" {
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Fatal("expected 202606190001_create_user_access_logs.sql in embedded migrations")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClickHouseGooseDialect(t *testing.T) {
|
||||
if err := goose.SetDialect("clickhouse"); err != nil {
|
||||
t.Fatalf("SetDialect(clickhouse) error = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateClickHouseSkipsWhenDisabled(t *testing.T) {
|
||||
previousEnabled := config.Config.ClickHouse.Enabled
|
||||
config.Config.ClickHouse.Enabled = false
|
||||
t.Cleanup(func() {
|
||||
config.Config.ClickHouse.Enabled = previousEnabled
|
||||
})
|
||||
|
||||
MigrateClickHouse()
|
||||
}
|
||||
+21
@@ -0,0 +1,21 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS w_user_access_logs
|
||||
(
|
||||
id UInt64,
|
||||
user_id UInt64,
|
||||
path String,
|
||||
method String,
|
||||
ip String,
|
||||
user_agent String,
|
||||
headers String,
|
||||
status Int32,
|
||||
latency Int64,
|
||||
created_at DateTime
|
||||
)
|
||||
ENGINE = MergeTree()
|
||||
PARTITION BY toYYYYMM(created_at)
|
||||
ORDER BY (created_at, ip, user_id)
|
||||
SETTINGS index_granularity = 8192;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_user_access_logs;
|
||||
@@ -0,0 +1,184 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS users (
|
||||
id BIGINT PRIMARY KEY,
|
||||
username VARCHAR(64) UNIQUE,
|
||||
password VARCHAR(255),
|
||||
nickname VARCHAR(255),
|
||||
email VARCHAR(255),
|
||||
avatar_url VARCHAR(255),
|
||||
is_active BOOLEAN DEFAULT TRUE,
|
||||
is_admin BOOLEAN DEFAULT FALSE,
|
||||
bio VARCHAR(500),
|
||||
phone VARCHAR(32),
|
||||
gender VARCHAR(16),
|
||||
website VARCHAR(255),
|
||||
location VARCHAR(255),
|
||||
last_login_at TIMESTAMPTZ,
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_email ON users (email);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_is_active ON users (is_active);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_last_login_at ON users (last_login_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_created_at ON users (created_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS auth_sources (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
name VARCHAR(80) NOT NULL UNIQUE,
|
||||
type VARCHAR(20) NOT NULL,
|
||||
display_name VARCHAR(100),
|
||||
is_active BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
client_id VARCHAR(255),
|
||||
client_secret VARCHAR(1024),
|
||||
openid_discovery_url VARCHAR(1024),
|
||||
scopes VARCHAR(255),
|
||||
icon_url VARCHAR(1024),
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_auth_sources_is_active ON auth_sources (is_active);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS external_accounts (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
auth_source_id BIGINT,
|
||||
user_id BIGINT NOT NULL,
|
||||
external_id VARCHAR(255) NOT NULL,
|
||||
external_username VARCHAR(255),
|
||||
email VARCHAR(255),
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_external_accounts_auth_source_id ON external_accounts (auth_source_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_external_accounts_user_id ON external_accounts (user_id);
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_external_accounts_source_external ON external_accounts (auth_source_id, external_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS system_configs (
|
||||
key VARCHAR(64) PRIMARY KEY,
|
||||
value TEXT NOT NULL,
|
||||
type VARCHAR(32) NOT NULL DEFAULT 'system',
|
||||
visibility INTEGER NOT NULL DEFAULT 0,
|
||||
description VARCHAR(255),
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS uploads (
|
||||
id BIGINT PRIMARY KEY,
|
||||
user_id BIGINT NOT NULL,
|
||||
file_name VARCHAR(255) NOT NULL,
|
||||
file_path VARCHAR(500) NOT NULL,
|
||||
file_size BIGINT NOT NULL,
|
||||
mime_type VARCHAR(100) NOT NULL,
|
||||
extension VARCHAR(50) NOT NULL,
|
||||
hash VARCHAR(64),
|
||||
storage_driver VARCHAR(50) NOT NULL,
|
||||
type VARCHAR(50) NOT NULL,
|
||||
status VARCHAR(20) NOT NULL,
|
||||
metadata JSONB,
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_user_id ON uploads (user_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_file_path ON uploads (file_path);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_hash ON uploads (hash);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_type ON uploads (type);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS access_tokens (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
user_id BIGINT NOT NULL,
|
||||
name VARCHAR(128) NOT NULL,
|
||||
token_hash VARCHAR(64) NOT NULL UNIQUE,
|
||||
masked_token VARCHAR(64) NOT NULL,
|
||||
last_used_at TIMESTAMPTZ,
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_access_tokens_user_id ON access_tokens (user_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS task_executions (
|
||||
id BIGINT PRIMARY KEY,
|
||||
task_id VARCHAR(128) NOT NULL UNIQUE,
|
||||
task_type VARCHAR(64) NOT NULL,
|
||||
task_name VARCHAR(128),
|
||||
status VARCHAR(32) NOT NULL,
|
||||
retryable BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
max_retry INTEGER NOT NULL DEFAULT 0,
|
||||
retry_count INTEGER NOT NULL DEFAULT 0,
|
||||
log TEXT,
|
||||
error_message TEXT,
|
||||
result TEXT,
|
||||
started_at TIMESTAMPTZ,
|
||||
finished_at TIMESTAMPTZ,
|
||||
duration BIGINT,
|
||||
payload TEXT,
|
||||
triggered_by VARCHAR(32) NOT NULL DEFAULT 'system',
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_task_type ON task_executions (task_type);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_status ON task_executions (status);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_started_at ON task_executions (started_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_created_at ON task_executions (created_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS templates (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
key VARCHAR(80) NOT NULL UNIQUE,
|
||||
name VARCHAR(100) NOT NULL,
|
||||
type VARCHAR(20) NOT NULL DEFAULT 'email',
|
||||
subject VARCHAR(255),
|
||||
content TEXT NOT NULL,
|
||||
description VARCHAR(255),
|
||||
is_system BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_is_system ON templates (is_system);
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_created_at ON templates (created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_updated_at ON templates (updated_at);
|
||||
|
||||
INSERT INTO system_configs (key, value, type, visibility, description, created_at, updated_at) VALUES
|
||||
('cap_login_enabled', 'false', 'system', 1, '是否启用登录人机验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_auto_solve', 'true', 'system', 1, '打开页面后是否自动开始计算,关闭则需用户手动点击触发', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_count', '1', 'system', 0, '客户端需求解的 PoW 难题总数,默认 1,推荐 1~5', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_size', '32', 'system', 0, '人机验证盐值长度', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_difficulty', '4', 'system', 0, '人机验证 PoW 难度(目标前缀长度)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_ttl_seconds', '600', 'system', 0, '人机验证难题有效时间(秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_token_ttl_seconds', '1200', 'system', 0, '人机验证兑换凭证有效时间(秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('server_address', '', 'system', 0, '服务器地址(用于跨域源控制,不设定则允许任意源)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_host', '', 'system', 0, 'SMTP 服务器地址(例如 smtp.example.com)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_port', '587', 'system', 0, 'SMTP 端口(例如 587 或 465)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_username', '', 'system', 0, 'SMTP 账户(如 sender@example.com)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_password', '', 'system', 0, 'SMTP 访问凭证(授权码/密码)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('upload_allowed_extensions', 'jpg,png,webp', 'system', 1, '允许上传的图片扩展名(逗号分隔)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('site_name', 'Wavelet', 'system', 1, '系统平台的展示名称', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('password_login_enabled', 'true', 'system', 1, '是否允许使用账号密码登录', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('registration_enabled', 'true', 'system', 1, '控制普通用户是否可以自主注册(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('password_register_enabled', 'true', 'system', 1, '是否允许通过密码创建本地账号', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('oidc_login_enabled', 'true', 'system', 1, '是否允许使用第三方 OIDC 认证源登录', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('max_api_keys_per_user', '5', 'business', 1, '限制每个普通用户可以创建的 API Key 最大数量', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('email_login_verification_enabled', 'false', 'system', 1, '是否开启邮箱登录验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('email_register_verification_enabled', 'false', 'system', 1, '是否开启邮箱注册验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('menu_display_config', '{}', 'system', 1, '目录显示配置(JSON 字符串,格式为 {url: enabled})', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('search_engine_indexing_enabled', 'false', 'system', 1, '是否允许搜索引擎爬取/检索该站点(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('update_upstream_repository', 'Rain-kl/Wavelet', 'system', 0, 'GitHub Actions Release 上游仓库(owner/repo 或 GitHub 仓库地址)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('storage_config', '{"driver":"local","local":{"root":"."},"s3":{"region":"us-east-1"},"r2":{"region":"auto"},"minio":{"region":"us-east-1","path_style":true},"oss":{},"webdav":{}}', 'system', 0, '文件存储驱动及连接配置(JSON)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
INSERT INTO users (id, username, password, nickname, avatar_url, is_active, is_admin, last_login_at, created_at, updated_at)
|
||||
VALUES (1, 'admin', '12345678', 'Administrator', '', TRUE, TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (username) DO NOTHING;
|
||||
|
||||
INSERT INTO templates (key, name, type, subject, content, description, is_system, created_at, updated_at) VALUES
|
||||
('login_email', '登录验证码邮件', 'email', 'Wavelet 登录验证码', '<h3>Wavelet 登录验证</h3><p>您的登录验证码为:<strong>{{.Code}}</strong>,5分钟内有效,请勿将验证码泄露给他人。</p>', '用户密码登录时发送的验证码邮件模板,支持变量:{{.Code}}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('register_email', '注册验证码邮件', 'email', 'Wavelet 注册验证码', '<h3>Wavelet 注册验证</h3><p>您的注册验证码为:<strong>{{.Code}}</strong>,5分钟内有效,请勿泄露给他人。</p>', '用户注册时发送的验证码邮件模板,支持变量:{{.Code}}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS templates;
|
||||
DROP TABLE IF EXISTS task_executions;
|
||||
DROP TABLE IF EXISTS access_tokens;
|
||||
DROP TABLE IF EXISTS uploads;
|
||||
DROP TABLE IF EXISTS system_configs;
|
||||
DROP TABLE IF EXISTS external_accounts;
|
||||
DROP TABLE IF EXISTS auth_sources;
|
||||
DROP TABLE IF EXISTS users;
|
||||
@@ -0,0 +1,20 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS schedules (
|
||||
id BIGINT PRIMARY KEY,
|
||||
name VARCHAR(128) NOT NULL,
|
||||
task_type VARCHAR(64) NOT NULL,
|
||||
cron VARCHAR(64) NOT NULL,
|
||||
payload TEXT,
|
||||
is_active BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_schedules_is_active ON schedules (is_active);
|
||||
|
||||
-- Seed initial cleanup task
|
||||
INSERT INTO schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at)
|
||||
VALUES (1, '清理未使用上传', 'cleanup_unused_uploads', '0 */2 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (id) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS schedules;
|
||||
+5
@@ -0,0 +1,5 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE access_tokens ADD COLUMN is_admin BOOLEAN NOT NULL DEFAULT FALSE;
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE access_tokens DROP COLUMN IF EXISTS is_admin;
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
-- +goose Up
|
||||
-- +goose StatementBegin
|
||||
ALTER TABLE access_tokens DROP COLUMN IF EXISTS last_used_at;
|
||||
-- +goose StatementEnd
|
||||
|
||||
-- +goose Down
|
||||
-- +goose StatementBegin
|
||||
ALTER TABLE access_tokens ADD COLUMN last_used_at TIMESTAMPTZ;
|
||||
-- +goose StatementEnd
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
-- +goose Up
|
||||
-- +goose StatementBegin
|
||||
ALTER TABLE schedules ALTER COLUMN id ADD GENERATED BY DEFAULT AS IDENTITY (START WITH 100);
|
||||
-- +goose StatementEnd
|
||||
|
||||
-- +goose Down
|
||||
-- +goose StatementBegin
|
||||
ALTER TABLE schedules ALTER COLUMN id DROP IDENTITY IF EXISTS;
|
||||
-- +goose StatementEnd
|
||||
+66
@@ -0,0 +1,66 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE users RENAME TO w_users;
|
||||
ALTER TABLE auth_sources RENAME TO w_auth_sources;
|
||||
ALTER TABLE external_accounts RENAME TO w_external_accounts;
|
||||
ALTER TABLE system_configs RENAME TO w_system_configs;
|
||||
ALTER TABLE uploads RENAME TO w_uploads;
|
||||
ALTER TABLE access_tokens RENAME TO w_access_tokens;
|
||||
ALTER TABLE task_executions RENAME TO w_task_executions;
|
||||
ALTER TABLE templates RENAME TO w_templates;
|
||||
ALTER TABLE schedules RENAME TO w_schedules;
|
||||
|
||||
-- Rename indexes for consistency
|
||||
ALTER INDEX IF EXISTS idx_users_email RENAME TO idx_w_users_email;
|
||||
ALTER INDEX IF EXISTS idx_users_is_active RENAME TO idx_w_users_is_active;
|
||||
ALTER INDEX IF EXISTS idx_users_last_login_at RENAME TO idx_w_users_last_login_at;
|
||||
ALTER INDEX IF EXISTS idx_users_created_at RENAME TO idx_w_users_created_at;
|
||||
ALTER INDEX IF EXISTS idx_auth_sources_is_active RENAME TO idx_w_auth_sources_is_active;
|
||||
ALTER INDEX IF EXISTS idx_external_accounts_auth_source_id RENAME TO idx_w_external_accounts_auth_source_id;
|
||||
ALTER INDEX IF EXISTS idx_external_accounts_user_id RENAME TO idx_w_external_accounts_user_id;
|
||||
ALTER INDEX IF EXISTS idx_external_accounts_source_external RENAME TO idx_w_external_accounts_source_external;
|
||||
ALTER INDEX IF EXISTS idx_uploads_user_id RENAME TO idx_w_uploads_user_id;
|
||||
ALTER INDEX IF EXISTS idx_uploads_file_path RENAME TO idx_w_uploads_file_path;
|
||||
ALTER INDEX IF EXISTS idx_uploads_hash RENAME TO idx_w_uploads_hash;
|
||||
ALTER INDEX IF EXISTS idx_uploads_type RENAME TO idx_w_uploads_type;
|
||||
ALTER INDEX IF EXISTS idx_access_tokens_user_id RENAME TO idx_w_access_tokens_user_id;
|
||||
ALTER INDEX IF EXISTS idx_task_executions_task_type RENAME TO idx_w_task_executions_task_type;
|
||||
ALTER INDEX IF EXISTS idx_task_executions_status RENAME TO idx_w_task_executions_status;
|
||||
ALTER INDEX IF EXISTS idx_task_executions_started_at RENAME TO idx_w_task_executions_started_at;
|
||||
ALTER INDEX IF EXISTS idx_task_executions_created_at RENAME TO idx_w_task_executions_created_at;
|
||||
ALTER INDEX IF EXISTS idx_templates_is_system RENAME TO idx_w_templates_is_system;
|
||||
ALTER INDEX IF EXISTS idx_templates_created_at RENAME TO idx_w_templates_created_at;
|
||||
ALTER INDEX IF EXISTS idx_templates_updated_at RENAME TO idx_w_templates_updated_at;
|
||||
ALTER INDEX IF EXISTS idx_schedules_is_active RENAME TO idx_w_schedules_is_active;
|
||||
|
||||
-- +goose Down
|
||||
ALTER INDEX IF EXISTS idx_w_schedules_is_active RENAME TO idx_schedules_is_active;
|
||||
ALTER INDEX IF EXISTS idx_w_templates_updated_at RENAME TO idx_templates_updated_at;
|
||||
ALTER INDEX IF EXISTS idx_w_templates_created_at RENAME TO idx_templates_created_at;
|
||||
ALTER INDEX IF EXISTS idx_w_templates_is_system RENAME TO idx_templates_is_system;
|
||||
ALTER INDEX IF EXISTS idx_w_task_executions_created_at RENAME TO idx_task_executions_created_at;
|
||||
ALTER INDEX IF EXISTS idx_w_task_executions_started_at RENAME TO idx_task_executions_started_at;
|
||||
ALTER INDEX IF EXISTS idx_w_task_executions_status RENAME TO idx_task_executions_status;
|
||||
ALTER INDEX IF EXISTS idx_w_task_executions_task_type RENAME TO idx_task_executions_task_type;
|
||||
ALTER INDEX IF EXISTS idx_w_access_tokens_user_id RENAME TO idx_access_tokens_user_id;
|
||||
ALTER INDEX IF EXISTS idx_w_uploads_type RENAME TO idx_uploads_type;
|
||||
ALTER INDEX IF EXISTS idx_w_uploads_hash RENAME TO idx_uploads_hash;
|
||||
ALTER INDEX IF EXISTS idx_w_uploads_file_path RENAME TO idx_uploads_file_path;
|
||||
ALTER INDEX IF EXISTS idx_w_uploads_user_id RENAME TO idx_uploads_user_id;
|
||||
ALTER INDEX IF EXISTS idx_w_external_accounts_source_external RENAME TO idx_external_accounts_source_external;
|
||||
ALTER INDEX IF EXISTS idx_w_external_accounts_user_id RENAME TO idx_external_accounts_user_id;
|
||||
ALTER INDEX IF EXISTS idx_w_external_accounts_auth_source_id RENAME TO idx_external_accounts_auth_source_id;
|
||||
ALTER INDEX IF EXISTS idx_w_auth_sources_is_active RENAME TO idx_auth_sources_is_active;
|
||||
ALTER INDEX IF EXISTS idx_w_users_created_at RENAME TO idx_users_created_at;
|
||||
ALTER INDEX IF EXISTS idx_w_users_last_login_at RENAME TO idx_users_last_login_at;
|
||||
ALTER INDEX IF EXISTS idx_w_users_is_active RENAME TO idx_users_is_active;
|
||||
ALTER INDEX IF EXISTS idx_w_users_email RENAME TO idx_users_email;
|
||||
|
||||
ALTER TABLE w_schedules RENAME TO schedules;
|
||||
ALTER TABLE w_templates RENAME TO templates;
|
||||
ALTER TABLE w_task_executions RENAME TO task_executions;
|
||||
ALTER TABLE w_access_tokens RENAME TO access_tokens;
|
||||
ALTER TABLE w_uploads RENAME TO uploads;
|
||||
ALTER TABLE w_system_configs RENAME TO system_configs;
|
||||
ALTER TABLE w_external_accounts RENAME TO external_accounts;
|
||||
ALTER TABLE w_auth_sources RENAME TO auth_sources;
|
||||
ALTER TABLE w_users RENAME TO users;
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES ('file_access_whitelist', '["avatar"]', 'system', 1, '免登录访问的文件业务类型白名单 (JSON 数组格式)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'file_access_whitelist';
|
||||
+10
@@ -0,0 +1,10 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES
|
||||
('disk_cache_max_size_mb', '100', 'system', 0, '磁盘缓存最大空间大小 (MB)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('disk_cache_ttl_minutes', '60', 'system', 0, '磁盘缓存默认有效期 (分钟)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('disk_cache_lru_enabled', 'true', 'system', 0, '是否启用 LRU 淘汰机制', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key IN ('disk_cache_max_size_mb', 'disk_cache_ttl_minutes', 'disk_cache_lru_enabled');
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES
|
||||
('login_session_ttl_hours', '0', 'system', 0, '登录会话过期时间 (小时,0表示浏览器关闭后自动退出,-1表示永不过期)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'login_session_ttl_hours';
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES ('update_upstream_repository', 'Rain-kl/Wavelet', 'system', 0, 'GitHub Actions Release 上游仓库(owner/repo 或 GitHub 仓库地址)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'update_upstream_repository';
|
||||
+6
@@ -0,0 +1,6 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE w_uploads ADD COLUMN access_mode INTEGER NOT NULL DEFAULT 0;
|
||||
UPDATE w_uploads SET access_mode = 1 WHERE type = 'avatar';
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE w_uploads DROP COLUMN access_mode;
|
||||
+5
@@ -0,0 +1,5 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE w_system_configs ALTER COLUMN value TYPE TEXT;
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE w_system_configs ALTER COLUMN value TYPE VARCHAR(255);
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES (
|
||||
'storage_config',
|
||||
'{"driver":"local","local":{"root":"."},"s3":{"region":"us-east-1"},"r2":{"region":"auto"},"minio":{"region":"us-east-1","path_style":true},"oss":{},"webdav":{}}',
|
||||
'system',
|
||||
0,
|
||||
'文件存储驱动及连接配置(JSON)',
|
||||
CURRENT_TIMESTAMP,
|
||||
CURRENT_TIMESTAMP
|
||||
)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'storage_config';
|
||||
+34
@@ -0,0 +1,34 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE w_push_events (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
event_key VARCHAR(80) NOT NULL UNIQUE,
|
||||
name VARCHAR(100) NOT NULL,
|
||||
channels TEXT NOT NULL,
|
||||
targets TEXT NOT NULL,
|
||||
template TEXT NOT NULL,
|
||||
enabled BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
created_at TIMESTAMPTZ NOT NULL,
|
||||
updated_at TIMESTAMPTZ NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX idx_w_push_events_enabled ON w_push_events(enabled);
|
||||
|
||||
CREATE TABLE w_push_histories (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
event_key VARCHAR(80) NOT NULL,
|
||||
channel VARCHAR(50) NOT NULL,
|
||||
target VARCHAR(255) NOT NULL,
|
||||
title VARCHAR(255) NOT NULL,
|
||||
content TEXT NOT NULL,
|
||||
level VARCHAR(20) NOT NULL,
|
||||
status VARCHAR(20) NOT NULL,
|
||||
error_msg TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX idx_w_push_histories_event ON w_push_histories(event_key);
|
||||
CREATE INDEX idx_w_push_histories_created ON w_push_histories(created_at);
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_push_histories;
|
||||
DROP TABLE IF EXISTS w_push_events;
|
||||
@@ -0,0 +1,8 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_users (id, username, password, nickname, avatar_url, is_active, is_admin, last_login_at, created_at, updated_at)
|
||||
VALUES (999, 'system', '*', '系统', '', TRUE, FALSE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (username) DO NOTHING;
|
||||
SELECT setval(pg_get_serial_sequence('w_users', 'id'), COALESCE((SELECT MAX(id) FROM w_users), 1));
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_users WHERE username = 'system';
|
||||
+19
@@ -0,0 +1,19 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE w_push_channels (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
name VARCHAR(80) NOT NULL UNIQUE,
|
||||
description VARCHAR(255),
|
||||
type VARCHAR(50) NOT NULL DEFAULT 'custom',
|
||||
token VARCHAR(100),
|
||||
url TEXT NOT NULL,
|
||||
other TEXT NOT NULL,
|
||||
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
created_at TIMESTAMPTZ NOT NULL,
|
||||
updated_at TIMESTAMPTZ NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX idx_w_push_channels_name ON w_push_channels(name);
|
||||
CREATE INDEX idx_w_push_channels_enabled ON w_push_channels(enabled);
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_push_channels;
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
-- +goose Up
|
||||
UPDATE w_schedules
|
||||
SET name = '系统定期垃圾清理',
|
||||
task_type = 'system_cleanup'
|
||||
WHERE id = 1;
|
||||
|
||||
-- +goose Down
|
||||
UPDATE w_schedules
|
||||
SET name = '清理未使用上传',
|
||||
task_type = 'cleanup_unused_uploads'
|
||||
WHERE id = 1;
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE w_push_events ADD COLUMN task_type VARCHAR(100) NOT NULL DEFAULT '';
|
||||
CREATE INDEX idx_w_push_events_task_type ON w_push_events(task_type);
|
||||
|
||||
-- +goose Down
|
||||
DROP INDEX IF EXISTS idx_w_push_events_task_type;
|
||||
ALTER TABLE w_push_events DROP COLUMN task_type;
|
||||
@@ -0,0 +1,6 @@
|
||||
-- +goose Up
|
||||
DELETE FROM w_system_configs WHERE key = 'push_config';
|
||||
DELETE FROM w_push_events WHERE event_key = 'admin_login';
|
||||
DELETE FROM w_system_configs WHERE key = 'push_global_token';
|
||||
|
||||
-- +goose Down
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
-- +goose Up
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status);
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status);
|
||||
|
||||
-- +goose Down
|
||||
DROP INDEX IF EXISTS idx_w_uploads_hash_file_size_status;
|
||||
DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status;
|
||||
DROP INDEX IF EXISTS idx_w_uploads_status_created_at;
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS w_upload_stats (
|
||||
dimension VARCHAR(32) NOT NULL,
|
||||
stat_key VARCHAR(64) NOT NULL DEFAULT '',
|
||||
file_count BIGINT NOT NULL DEFAULT 0,
|
||||
file_size BIGINT NOT NULL DEFAULT 0,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
PRIMARY KEY (dimension, stat_key)
|
||||
);
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_upload_stats;
|
||||
+67
@@ -0,0 +1,67 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT 'total', '', COUNT(*), COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = EXCLUDED.file_count,
|
||||
file_size = EXCLUDED.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT
|
||||
'type',
|
||||
COALESCE(NULLIF(type, ''), 'generic'),
|
||||
COUNT(*),
|
||||
COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
GROUP BY COALESCE(NULLIF(type, ''), 'generic')
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = EXCLUDED.file_count,
|
||||
file_size = EXCLUDED.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT
|
||||
'category',
|
||||
CASE
|
||||
WHEN LOWER(mime_type) LIKE 'image/%'
|
||||
OR LOWER(extension) IN ('jpg', 'jpeg', 'png', 'webp', 'gif') THEN '图片'
|
||||
WHEN LOWER(mime_type) LIKE 'video/%' THEN '视频'
|
||||
WHEN LOWER(mime_type) LIKE 'audio/%' THEN '音频'
|
||||
WHEN LOWER(extension) IN ('zip', 'rar', '7z', 'tar', 'gz', 'tgz', 'bz2', 'xz')
|
||||
OR LOWER(mime_type) LIKE '%zip%'
|
||||
OR LOWER(mime_type) LIKE '%tar%'
|
||||
OR LOWER(mime_type) LIKE '%gzip%' THEN '压缩包'
|
||||
WHEN LOWER(extension) IN ('pdf', 'doc', 'docx', 'xls', 'xlsx', 'ppt', 'pptx', 'txt', 'md', 'csv', 'json', 'yaml', 'yml', 'xml')
|
||||
OR LOWER(mime_type) LIKE 'text/%'
|
||||
OR LOWER(mime_type) = 'application/pdf' THEN '文档'
|
||||
ELSE '其他'
|
||||
END,
|
||||
COUNT(*),
|
||||
COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
GROUP BY 2
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = EXCLUDED.file_count,
|
||||
file_size = EXCLUDED.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT
|
||||
'trend',
|
||||
TO_CHAR(created_at, 'YYYY-MM-DD'),
|
||||
COUNT(*),
|
||||
COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
GROUP BY TO_CHAR(created_at, 'YYYY-MM-DD')
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = EXCLUDED.file_count,
|
||||
file_size = EXCLUDED.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_upload_stats;
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status;
|
||||
ALTER TABLE w_uploads DROP COLUMN IF EXISTS storage_driver;
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE w_uploads ADD COLUMN IF NOT EXISTS storage_driver VARCHAR(50) NOT NULL DEFAULT 'local';
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status);
|
||||
@@ -0,0 +1,184 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS users (
|
||||
id BIGINT PRIMARY KEY,
|
||||
username VARCHAR(64) UNIQUE,
|
||||
password VARCHAR(255),
|
||||
nickname VARCHAR(255),
|
||||
email VARCHAR(255),
|
||||
avatar_url VARCHAR(255),
|
||||
is_active BOOLEAN DEFAULT TRUE,
|
||||
is_admin BOOLEAN DEFAULT FALSE,
|
||||
bio VARCHAR(500),
|
||||
phone VARCHAR(32),
|
||||
gender VARCHAR(16),
|
||||
website VARCHAR(255),
|
||||
location VARCHAR(255),
|
||||
last_login_at DATETIME,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_email ON users (email);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_is_active ON users (is_active);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_last_login_at ON users (last_login_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_users_created_at ON users (created_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS auth_sources (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name VARCHAR(80) NOT NULL UNIQUE,
|
||||
type VARCHAR(20) NOT NULL,
|
||||
display_name VARCHAR(100),
|
||||
is_active BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
client_id VARCHAR(255),
|
||||
client_secret VARCHAR(1024),
|
||||
openid_discovery_url VARCHAR(1024),
|
||||
scopes VARCHAR(255),
|
||||
icon_url VARCHAR(1024),
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_auth_sources_is_active ON auth_sources (is_active);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS external_accounts (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
auth_source_id BIGINT,
|
||||
user_id BIGINT NOT NULL,
|
||||
external_id VARCHAR(255) NOT NULL,
|
||||
external_username VARCHAR(255),
|
||||
email VARCHAR(255),
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_external_accounts_auth_source_id ON external_accounts (auth_source_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_external_accounts_user_id ON external_accounts (user_id);
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_external_accounts_source_external ON external_accounts (auth_source_id, external_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS system_configs (
|
||||
key VARCHAR(64) PRIMARY KEY,
|
||||
value TEXT NOT NULL,
|
||||
type VARCHAR(32) NOT NULL DEFAULT 'system',
|
||||
visibility INTEGER NOT NULL DEFAULT 0,
|
||||
description VARCHAR(255),
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS uploads (
|
||||
id BIGINT PRIMARY KEY,
|
||||
user_id BIGINT NOT NULL,
|
||||
file_name VARCHAR(255) NOT NULL,
|
||||
file_path VARCHAR(500) NOT NULL,
|
||||
file_size BIGINT NOT NULL,
|
||||
mime_type VARCHAR(100) NOT NULL,
|
||||
extension VARCHAR(50) NOT NULL,
|
||||
hash VARCHAR(64),
|
||||
storage_driver VARCHAR(50) NOT NULL,
|
||||
type VARCHAR(50) NOT NULL,
|
||||
status VARCHAR(20) NOT NULL,
|
||||
metadata JSON,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_user_id ON uploads (user_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_file_path ON uploads (file_path);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_hash ON uploads (hash);
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_type ON uploads (type);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS access_tokens (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
user_id BIGINT NOT NULL,
|
||||
name VARCHAR(128) NOT NULL,
|
||||
token_hash VARCHAR(64) NOT NULL UNIQUE,
|
||||
masked_token VARCHAR(64) NOT NULL,
|
||||
last_used_at DATETIME,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_access_tokens_user_id ON access_tokens (user_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS task_executions (
|
||||
id BIGINT PRIMARY KEY,
|
||||
task_id VARCHAR(128) NOT NULL UNIQUE,
|
||||
task_type VARCHAR(64) NOT NULL,
|
||||
task_name VARCHAR(128),
|
||||
status VARCHAR(32) NOT NULL,
|
||||
retryable BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
max_retry INTEGER NOT NULL DEFAULT 0,
|
||||
retry_count INTEGER NOT NULL DEFAULT 0,
|
||||
log TEXT,
|
||||
error_message TEXT,
|
||||
result TEXT,
|
||||
started_at DATETIME,
|
||||
finished_at DATETIME,
|
||||
duration BIGINT,
|
||||
payload TEXT,
|
||||
triggered_by VARCHAR(32) NOT NULL DEFAULT 'system',
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_task_type ON task_executions (task_type);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_status ON task_executions (status);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_started_at ON task_executions (started_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_created_at ON task_executions (created_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS templates (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
key VARCHAR(80) NOT NULL UNIQUE,
|
||||
name VARCHAR(100) NOT NULL,
|
||||
type VARCHAR(20) NOT NULL DEFAULT 'email',
|
||||
subject VARCHAR(255),
|
||||
content TEXT NOT NULL,
|
||||
description VARCHAR(255),
|
||||
is_system BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_is_system ON templates (is_system);
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_created_at ON templates (created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_updated_at ON templates (updated_at);
|
||||
|
||||
INSERT INTO system_configs (key, value, type, visibility, description, created_at, updated_at) VALUES
|
||||
('cap_login_enabled', 'false', 'system', 1, '是否启用登录人机验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_auto_solve', 'true', 'system', 1, '打开页面后是否自动开始计算,关闭则需用户手动点击触发', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_count', '1', 'system', 0, '客户端需求解的 PoW 难题总数,默认 1,推荐 1~5', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_size', '32', 'system', 0, '人机验证盐值长度', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_difficulty', '4', 'system', 0, '人机验证 PoW 难度(目标前缀长度)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_challenge_ttl_seconds', '600', 'system', 0, '人机验证难题有效时间(秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('cap_token_ttl_seconds', '1200', 'system', 0, '人机验证兑换凭证有效时间(秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('server_address', '', 'system', 0, '服务器地址(用于跨域源控制,不设定则允许任意源)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_host', '', 'system', 0, 'SMTP 服务器地址(例如 smtp.example.com)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_port', '587', 'system', 0, 'SMTP 端口(例如 587 或 465)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_username', '', 'system', 0, 'SMTP 账户(如 sender@example.com)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('smtp_password', '', 'system', 0, 'SMTP 访问凭证(授权码/密码)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('upload_allowed_extensions', 'jpg,png,webp', 'system', 1, '允许上传的图片扩展名(逗号分隔)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('site_name', 'Wavelet', 'system', 1, '系统平台的展示名称', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('password_login_enabled', 'true', 'system', 1, '是否允许使用账号密码登录', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('registration_enabled', 'true', 'system', 1, '控制普通用户是否可以自主注册(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('password_register_enabled', 'true', 'system', 1, '是否允许通过密码创建本地账号', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('oidc_login_enabled', 'true', 'system', 1, '是否允许使用第三方 OIDC 认证源登录', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('max_api_keys_per_user', '5', 'business', 1, '限制每个普通用户可以创建的 API Key 最大数量', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('email_login_verification_enabled', 'false', 'system', 1, '是否开启邮箱登录验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('email_register_verification_enabled', 'false', 'system', 1, '是否开启邮箱注册验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('menu_display_config', '{}', 'system', 1, '目录显示配置(JSON 字符串,格式为 {url: enabled})', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('search_engine_indexing_enabled', 'false', 'system', 1, '是否允许搜索引擎爬取/检索该站点(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('update_upstream_repository', 'Rain-kl/Wavelet', 'system', 0, 'GitHub Actions Release 上游仓库(owner/repo 或 GitHub 仓库地址)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('storage_config', '{"driver":"local","local":{"root":"."},"s3":{"region":"us-east-1"},"r2":{"region":"auto"},"minio":{"region":"us-east-1","path_style":true},"oss":{},"webdav":{}}', 'system', 0, '文件存储驱动及连接配置(JSON)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
INSERT INTO users (id, username, password, nickname, avatar_url, is_active, is_admin, last_login_at, created_at, updated_at)
|
||||
VALUES (1, 'admin', '12345678', 'Administrator', '', TRUE, TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (username) DO NOTHING;
|
||||
|
||||
INSERT INTO templates (key, name, type, subject, content, description, is_system, created_at, updated_at) VALUES
|
||||
('login_email', '登录验证码邮件', 'email', 'Wavelet 登录验证码', '<h3>Wavelet 登录验证</h3><p>您的登录验证码为:<strong>{{.Code}}</strong>,5分钟内有效,请勿将验证码泄露给他人。</p>', '用户密码登录时发送的验证码邮件模板,支持变量:{{.Code}}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('register_email', '注册验证码邮件', 'email', 'Wavelet 注册验证码', '<h3>Wavelet 注册验证</h3><p>您的注册验证码为:<strong>{{.Code}}</strong>,5分钟内有效,请勿泄露给他人。</p>', '用户注册时发送的验证码邮件模板,支持变量:{{.Code}}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS templates;
|
||||
DROP TABLE IF EXISTS task_executions;
|
||||
DROP TABLE IF EXISTS access_tokens;
|
||||
DROP TABLE IF EXISTS uploads;
|
||||
DROP TABLE IF EXISTS system_configs;
|
||||
DROP TABLE IF EXISTS external_accounts;
|
||||
DROP TABLE IF EXISTS auth_sources;
|
||||
DROP TABLE IF EXISTS users;
|
||||
@@ -0,0 +1,20 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS schedules (
|
||||
id BIGINT PRIMARY KEY,
|
||||
name VARCHAR(128) NOT NULL,
|
||||
task_type VARCHAR(64) NOT NULL,
|
||||
cron VARCHAR(64) NOT NULL,
|
||||
payload TEXT,
|
||||
is_active BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_schedules_is_active ON schedules (is_active);
|
||||
|
||||
-- Seed initial cleanup task
|
||||
INSERT INTO schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at)
|
||||
VALUES (1, '清理未使用上传', 'cleanup_unused_uploads', '0 */2 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (id) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS schedules;
|
||||
+5
@@ -0,0 +1,5 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE access_tokens ADD COLUMN is_admin BOOLEAN NOT NULL DEFAULT 0;
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE access_tokens DROP COLUMN is_admin;
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
-- +goose Up
|
||||
-- +goose StatementBegin
|
||||
ALTER TABLE access_tokens DROP COLUMN last_used_at;
|
||||
-- +goose StatementEnd
|
||||
|
||||
-- +goose Down
|
||||
-- +goose StatementBegin
|
||||
ALTER TABLE access_tokens ADD COLUMN last_used_at DATETIME;
|
||||
-- +goose StatementEnd
|
||||
+32
@@ -0,0 +1,32 @@
|
||||
-- +goose Up
|
||||
-- +goose StatementBegin
|
||||
-- 1. Rename existing schedules table
|
||||
ALTER TABLE schedules RENAME TO schedules_old;
|
||||
|
||||
-- 2. Create new schedules table with AUTOINCREMENT
|
||||
CREATE TABLE schedules (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name VARCHAR(128) NOT NULL,
|
||||
task_type VARCHAR(64) NOT NULL,
|
||||
cron VARCHAR(64) NOT NULL,
|
||||
payload TEXT,
|
||||
is_active BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
|
||||
);
|
||||
|
||||
-- 3. Copy existing data
|
||||
INSERT INTO schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at)
|
||||
SELECT id, name, task_type, cron, payload, is_active, created_at, updated_at FROM schedules_old;
|
||||
|
||||
-- 4. Drop the old table
|
||||
DROP TABLE schedules_old;
|
||||
|
||||
-- 5. Recreate index
|
||||
CREATE INDEX IF NOT EXISTS idx_schedules_is_active ON schedules (is_active);
|
||||
-- +goose StatementEnd
|
||||
|
||||
-- +goose Down
|
||||
-- +goose StatementBegin
|
||||
-- Re-creating table with AUTOINCREMENT cannot be undone simply without recreating table again.
|
||||
-- +goose StatementEnd
|
||||
+123
@@ -0,0 +1,123 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE users RENAME TO w_users;
|
||||
ALTER TABLE auth_sources RENAME TO w_auth_sources;
|
||||
ALTER TABLE external_accounts RENAME TO w_external_accounts;
|
||||
ALTER TABLE system_configs RENAME TO w_system_configs;
|
||||
ALTER TABLE uploads RENAME TO w_uploads;
|
||||
ALTER TABLE access_tokens RENAME TO w_access_tokens;
|
||||
ALTER TABLE task_executions RENAME TO w_task_executions;
|
||||
ALTER TABLE templates RENAME TO w_templates;
|
||||
ALTER TABLE schedules RENAME TO w_schedules;
|
||||
|
||||
-- Drop and Recreate SQLite indexes
|
||||
DROP INDEX IF EXISTS idx_users_email;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_users_email ON w_users (email);
|
||||
DROP INDEX IF EXISTS idx_users_is_active;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_users_is_active ON w_users (is_active);
|
||||
DROP INDEX IF EXISTS idx_users_last_login_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_users_last_login_at ON w_users (last_login_at);
|
||||
DROP INDEX IF EXISTS idx_users_created_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_users_created_at ON w_users (created_at);
|
||||
|
||||
DROP INDEX IF EXISTS idx_auth_sources_is_active;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_auth_sources_is_active ON w_auth_sources (is_active);
|
||||
|
||||
DROP INDEX IF EXISTS idx_external_accounts_auth_source_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_external_accounts_auth_source_id ON w_external_accounts (auth_source_id);
|
||||
DROP INDEX IF EXISTS idx_external_accounts_user_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_external_accounts_user_id ON w_external_accounts (user_id);
|
||||
DROP INDEX IF EXISTS idx_external_accounts_source_external;
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_w_external_accounts_source_external ON w_external_accounts (auth_source_id, external_id);
|
||||
|
||||
DROP INDEX IF EXISTS idx_uploads_user_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_user_id ON w_uploads (user_id);
|
||||
DROP INDEX IF EXISTS idx_uploads_file_path;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_file_path ON w_uploads (file_path);
|
||||
DROP INDEX IF EXISTS idx_uploads_hash;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_hash ON w_uploads (hash);
|
||||
DROP INDEX IF EXISTS idx_uploads_type;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_type ON w_uploads (type);
|
||||
|
||||
DROP INDEX IF EXISTS idx_access_tokens_user_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_access_tokens_user_id ON w_access_tokens (user_id);
|
||||
|
||||
DROP INDEX IF EXISTS idx_task_executions_task_type;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_task_executions_task_type ON w_task_executions (task_type);
|
||||
DROP INDEX IF EXISTS idx_task_executions_status;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_task_executions_status ON w_task_executions (status);
|
||||
DROP INDEX IF EXISTS idx_task_executions_started_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_task_executions_started_at ON w_task_executions (started_at);
|
||||
DROP INDEX IF EXISTS idx_task_executions_created_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_task_executions_created_at ON w_task_executions (created_at);
|
||||
|
||||
DROP INDEX IF EXISTS idx_templates_is_system;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_templates_is_system ON w_templates (is_system);
|
||||
DROP INDEX IF EXISTS idx_templates_created_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_templates_created_at ON w_templates (created_at);
|
||||
DROP INDEX IF EXISTS idx_templates_updated_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_templates_updated_at ON w_templates (updated_at);
|
||||
|
||||
DROP INDEX IF EXISTS idx_schedules_is_active;
|
||||
CREATE INDEX IF NOT EXISTS idx_w_schedules_is_active ON w_schedules (is_active);
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE w_schedules RENAME TO schedules;
|
||||
ALTER TABLE w_templates RENAME TO templates;
|
||||
ALTER TABLE w_task_executions RENAME TO task_executions;
|
||||
ALTER TABLE w_access_tokens RENAME TO access_tokens;
|
||||
ALTER TABLE w_uploads RENAME TO uploads;
|
||||
ALTER TABLE w_system_configs RENAME TO system_configs;
|
||||
ALTER TABLE w_external_accounts RENAME TO external_accounts;
|
||||
ALTER TABLE w_auth_sources RENAME TO auth_sources;
|
||||
ALTER TABLE w_users RENAME TO users;
|
||||
|
||||
-- Revert indexes
|
||||
DROP INDEX IF EXISTS idx_w_users_email;
|
||||
CREATE INDEX IF NOT EXISTS idx_users_email ON users (email);
|
||||
DROP INDEX IF EXISTS idx_w_users_is_active;
|
||||
CREATE INDEX IF NOT EXISTS idx_users_is_active ON users (is_active);
|
||||
DROP INDEX IF EXISTS idx_w_users_last_login_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_users_last_login_at ON users (last_login_at);
|
||||
DROP INDEX IF EXISTS idx_w_users_created_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_users_created_at ON users (created_at);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_auth_sources_is_active;
|
||||
CREATE INDEX IF NOT EXISTS idx_auth_sources_is_active ON auth_sources (is_active);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_external_accounts_auth_source_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_external_accounts_auth_source_id ON external_accounts (auth_source_id);
|
||||
DROP INDEX IF EXISTS idx_w_external_accounts_user_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_external_accounts_user_id ON external_accounts (user_id);
|
||||
DROP INDEX IF EXISTS idx_w_external_accounts_source_external;
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_external_accounts_source_external ON external_accounts (auth_source_id, external_id);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_uploads_user_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_user_id ON uploads (user_id);
|
||||
DROP INDEX IF EXISTS idx_w_uploads_file_path;
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_file_path ON uploads (file_path);
|
||||
DROP INDEX IF EXISTS idx_w_uploads_hash;
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_hash ON uploads (hash);
|
||||
DROP INDEX IF EXISTS idx_w_uploads_type;
|
||||
CREATE INDEX IF NOT EXISTS idx_uploads_type ON uploads (type);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_access_tokens_user_id;
|
||||
CREATE INDEX IF NOT EXISTS idx_access_tokens_user_id ON access_tokens (user_id);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_task_executions_task_type;
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_task_type ON task_executions (task_type);
|
||||
DROP INDEX IF EXISTS idx_w_task_executions_status;
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_status ON task_executions (status);
|
||||
DROP INDEX IF EXISTS idx_w_task_executions_started_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_started_at ON task_executions (started_at);
|
||||
DROP INDEX IF EXISTS idx_w_task_executions_created_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_task_executions_created_at ON task_executions (created_at);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_templates_is_system;
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_is_system ON templates (is_system);
|
||||
DROP INDEX IF EXISTS idx_w_templates_created_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_created_at ON templates (created_at);
|
||||
DROP INDEX IF EXISTS idx_w_templates_updated_at;
|
||||
CREATE INDEX IF NOT EXISTS idx_templates_updated_at ON templates (updated_at);
|
||||
|
||||
DROP INDEX IF EXISTS idx_w_schedules_is_active;
|
||||
CREATE INDEX IF NOT EXISTS idx_schedules_is_active ON schedules (is_active);
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES ('file_access_whitelist', '["avatar"]', 'system', 1, '免登录访问的文件业务类型白名单 (JSON 数组格式)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'file_access_whitelist';
|
||||
+10
@@ -0,0 +1,10 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES
|
||||
('disk_cache_max_size_mb', '100', 'system', 0, '磁盘缓存最大空间大小 (MB)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('disk_cache_ttl_minutes', '60', 'system', 0, '磁盘缓存默认有效期 (分钟)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||
('disk_cache_lru_enabled', 'true', 'system', 0, '是否启用 LRU 淘汰机制', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key IN ('disk_cache_max_size_mb', 'disk_cache_ttl_minutes', 'disk_cache_lru_enabled');
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES
|
||||
('login_session_ttl_hours', '0', 'system', 0, '登录会话过期时间 (小时,0表示浏览器关闭后自动退出,-1表示永不过期)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'login_session_ttl_hours';
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES ('update_upstream_repository', 'Rain-kl/Wavelet', 'system', 0, 'GitHub Actions Release 上游仓库(owner/repo 或 GitHub 仓库地址)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'update_upstream_repository';
|
||||
+6
@@ -0,0 +1,6 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE w_uploads ADD COLUMN access_mode INTEGER NOT NULL DEFAULT 0;
|
||||
UPDATE w_uploads SET access_mode = 1 WHERE type = 'avatar';
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE w_uploads DROP COLUMN access_mode;
|
||||
+6
@@ -0,0 +1,6 @@
|
||||
-- +goose Up
|
||||
-- SQLite stores VARCHAR and TEXT with the same TEXT affinity.
|
||||
SELECT 1;
|
||||
|
||||
-- +goose Down
|
||||
SELECT 1;
|
||||
@@ -0,0 +1,15 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES (
|
||||
'storage_config',
|
||||
'{"driver":"local","local":{"root":"."},"s3":{"region":"us-east-1"},"r2":{"region":"auto"},"minio":{"region":"us-east-1","path_style":true},"oss":{},"webdav":{}}',
|
||||
'system',
|
||||
0,
|
||||
'文件存储驱动及连接配置(JSON)',
|
||||
CURRENT_TIMESTAMP,
|
||||
CURRENT_TIMESTAMP
|
||||
)
|
||||
ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'storage_config';
|
||||
@@ -0,0 +1,34 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE w_push_events (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
event_key TEXT NOT NULL UNIQUE,
|
||||
name TEXT NOT NULL,
|
||||
channels TEXT NOT NULL,
|
||||
targets TEXT NOT NULL,
|
||||
template TEXT NOT NULL,
|
||||
enabled BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
created_at DATETIME NOT NULL,
|
||||
updated_at DATETIME NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX idx_w_push_events_enabled ON w_push_events(enabled);
|
||||
|
||||
CREATE TABLE w_push_histories (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
event_key TEXT NOT NULL,
|
||||
channel TEXT NOT NULL,
|
||||
target TEXT NOT NULL,
|
||||
title TEXT NOT NULL,
|
||||
content TEXT NOT NULL,
|
||||
level TEXT NOT NULL,
|
||||
status TEXT NOT NULL,
|
||||
error_msg TEXT,
|
||||
created_at DATETIME NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX idx_w_push_histories_event ON w_push_histories(event_key);
|
||||
CREATE INDEX idx_w_push_histories_created ON w_push_histories(created_at);
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_push_histories;
|
||||
DROP TABLE IF EXISTS w_push_events;
|
||||
@@ -0,0 +1,7 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_users (id, username, password, nickname, avatar_url, is_active, is_admin, last_login_at, created_at, updated_at)
|
||||
VALUES (999, 'system', '*', '系统', '', TRUE, FALSE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT (username) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_users WHERE username = 'system';
|
||||
+19
@@ -0,0 +1,19 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE w_push_channels (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name TEXT NOT NULL UNIQUE,
|
||||
description TEXT,
|
||||
type TEXT NOT NULL DEFAULT 'custom',
|
||||
token TEXT,
|
||||
url TEXT NOT NULL,
|
||||
other TEXT NOT NULL,
|
||||
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
created_at DATETIME NOT NULL,
|
||||
updated_at DATETIME NOT NULL
|
||||
);
|
||||
|
||||
CREATE INDEX idx_w_push_channels_name ON w_push_channels(name);
|
||||
CREATE INDEX idx_w_push_channels_enabled ON w_push_channels(enabled);
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_push_channels;
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
-- +goose Up
|
||||
UPDATE w_schedules
|
||||
SET name = '系统定期垃圾清理',
|
||||
task_type = 'system_cleanup'
|
||||
WHERE id = 1;
|
||||
|
||||
-- +goose Down
|
||||
UPDATE w_schedules
|
||||
SET name = '清理未使用上传',
|
||||
task_type = 'cleanup_unused_uploads'
|
||||
WHERE id = 1;
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
-- +goose Up
|
||||
ALTER TABLE w_push_events ADD COLUMN task_type VARCHAR(100) NOT NULL DEFAULT '';
|
||||
CREATE INDEX idx_w_push_events_task_type ON w_push_events(task_type);
|
||||
|
||||
-- +goose Down
|
||||
DROP INDEX IF EXISTS idx_w_push_events_task_type;
|
||||
-- SQLite does not support DROP COLUMN in older versions easily, but standard ALTER TABLE DROP COLUMN works in SQLite 3.35.0+.
|
||||
-- We can write standard DROP COLUMN.
|
||||
ALTER TABLE w_push_events DROP COLUMN task_type;
|
||||
@@ -0,0 +1,6 @@
|
||||
-- +goose Up
|
||||
DELETE FROM w_system_configs WHERE key = 'push_config';
|
||||
DELETE FROM w_push_events WHERE event_key = 'admin_login';
|
||||
DELETE FROM w_system_configs WHERE key = 'push_global_token';
|
||||
|
||||
-- +goose Down
|
||||
+9
@@ -0,0 +1,9 @@
|
||||
-- +goose Up
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at);
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status);
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status);
|
||||
|
||||
-- +goose Down
|
||||
DROP INDEX IF EXISTS idx_w_uploads_hash_file_size_status;
|
||||
DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status;
|
||||
DROP INDEX IF EXISTS idx_w_uploads_status_created_at;
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS w_upload_stats (
|
||||
dimension VARCHAR(32) NOT NULL,
|
||||
stat_key VARCHAR(64) NOT NULL DEFAULT '',
|
||||
file_count BIGINT NOT NULL DEFAULT 0,
|
||||
file_size BIGINT NOT NULL DEFAULT 0,
|
||||
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
PRIMARY KEY (dimension, stat_key)
|
||||
);
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_upload_stats;
|
||||
+67
@@ -0,0 +1,67 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT 'total', '', COUNT(*), COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = excluded.file_count,
|
||||
file_size = excluded.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT
|
||||
'type',
|
||||
COALESCE(NULLIF(type, ''), 'generic'),
|
||||
COUNT(*),
|
||||
COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
GROUP BY COALESCE(NULLIF(type, ''), 'generic')
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = excluded.file_count,
|
||||
file_size = excluded.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT
|
||||
'category',
|
||||
CASE
|
||||
WHEN LOWER(mime_type) LIKE 'image/%'
|
||||
OR LOWER(extension) IN ('jpg', 'jpeg', 'png', 'webp', 'gif') THEN '图片'
|
||||
WHEN LOWER(mime_type) LIKE 'video/%' THEN '视频'
|
||||
WHEN LOWER(mime_type) LIKE 'audio/%' THEN '音频'
|
||||
WHEN LOWER(extension) IN ('zip', 'rar', '7z', 'tar', 'gz', 'tgz', 'bz2', 'xz')
|
||||
OR LOWER(mime_type) LIKE '%zip%'
|
||||
OR LOWER(mime_type) LIKE '%tar%'
|
||||
OR LOWER(mime_type) LIKE '%gzip%' THEN '压缩包'
|
||||
WHEN LOWER(extension) IN ('pdf', 'doc', 'docx', 'xls', 'xlsx', 'ppt', 'pptx', 'txt', 'md', 'csv', 'json', 'yaml', 'yml', 'xml')
|
||||
OR LOWER(mime_type) LIKE 'text/%'
|
||||
OR LOWER(mime_type) = 'application/pdf' THEN '文档'
|
||||
ELSE '其他'
|
||||
END,
|
||||
COUNT(*),
|
||||
COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
GROUP BY 2
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = excluded.file_count,
|
||||
file_size = excluded.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size)
|
||||
SELECT
|
||||
'trend',
|
||||
STRFTIME('%Y-%m-%d', created_at),
|
||||
COUNT(*),
|
||||
COALESCE(SUM(file_size), 0)
|
||||
FROM w_uploads
|
||||
WHERE status != 'deleted'
|
||||
GROUP BY STRFTIME('%Y-%m-%d', created_at)
|
||||
ON CONFLICT (dimension, stat_key) DO UPDATE SET
|
||||
file_count = excluded.file_count,
|
||||
file_size = excluded.file_size,
|
||||
updated_at = CURRENT_TIMESTAMP;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_upload_stats;
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
-- +goose Up
|
||||
DROP INDEX IF EXISTS idx_w_uploads_storage_driver_status;
|
||||
-- SQLite lacks DROP COLUMN IF EXISTS; goose runs this only after initial schema created the column.
|
||||
ALTER TABLE w_uploads DROP COLUMN storage_driver;
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE w_uploads ADD COLUMN storage_driver VARCHAR(50) NOT NULL DEFAULT 'local';
|
||||
CREATE INDEX IF NOT EXISTS idx_w_uploads_storage_driver_status ON w_uploads (storage_driver, status);
|
||||
@@ -0,0 +1,103 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package migrator 提供数据库迁移功能
|
||||
package migrator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"embed"
|
||||
"log"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/pressly/goose/v3"
|
||||
)
|
||||
|
||||
// migrationFS contains SQL migrations under goose/<dialect>.
|
||||
//
|
||||
//go:embed goose/postgres/*.sql goose/sqlite/*.sql
|
||||
var migrationFS embed.FS
|
||||
|
||||
// dbType 返回当前数据库类型名称(用于日志输出)
|
||||
func dbType() string {
|
||||
if !config.Config.Database.Enabled {
|
||||
return "SQLite"
|
||||
}
|
||||
return "PostgreSQL"
|
||||
}
|
||||
|
||||
const (
|
||||
dialectSqlite = "sqlite3"
|
||||
dialectPostgres = "postgres"
|
||||
cascadeSuffix = " CASCADE"
|
||||
)
|
||||
|
||||
// Report describes the database migration state observed during startup.
|
||||
type Report struct {
|
||||
Backend string
|
||||
Enabled bool
|
||||
Version int64
|
||||
Applied bool
|
||||
}
|
||||
|
||||
func gooseDialect() string {
|
||||
if !config.Config.Database.Enabled {
|
||||
return dialectSqlite
|
||||
}
|
||||
return dialectPostgres
|
||||
}
|
||||
|
||||
func migrationDir() string {
|
||||
if !config.Config.Database.Enabled {
|
||||
return "goose/sqlite"
|
||||
}
|
||||
return "goose/postgres"
|
||||
}
|
||||
|
||||
// Migrate 执行数据库迁移
|
||||
func Migrate() Report {
|
||||
gormDB := db.DB(context.Background())
|
||||
if gormDB == nil {
|
||||
log.Fatalf("[%s] database not initialized\n", dbType())
|
||||
}
|
||||
|
||||
sqlDB, err := gormDB.DB()
|
||||
if err != nil {
|
||||
log.Fatalf("[%s] load sql db failed: %v\n", dbType(), err)
|
||||
}
|
||||
|
||||
goose.SetBaseFS(migrationFS)
|
||||
if err := goose.SetDialect(gooseDialect()); err != nil {
|
||||
log.Fatalf("[%s] set goose dialect failed: %v\n", dbType(), err)
|
||||
}
|
||||
previousVersion, err := goose.GetDBVersion(sqlDB)
|
||||
if err != nil {
|
||||
log.Fatalf("[%s] get goose version failed: %v\n", dbType(), err)
|
||||
}
|
||||
if err := goose.Up(sqlDB, migrationDir()); err != nil {
|
||||
log.Fatalf("[%s] goose migrate failed: %v\n", dbType(), err)
|
||||
}
|
||||
|
||||
clearSystemConfigCache()
|
||||
currentVersion, err := goose.GetDBVersion(sqlDB)
|
||||
if err != nil {
|
||||
log.Fatalf("[%s] get migrated goose version failed: %v\n", dbType(), err)
|
||||
}
|
||||
|
||||
log.Printf("[%s] goose migrate success\n", dbType())
|
||||
return Report{
|
||||
Backend: dbType(),
|
||||
Enabled: true,
|
||||
Version: currentVersion,
|
||||
Applied: currentVersion != previousVersion,
|
||||
}
|
||||
}
|
||||
|
||||
func clearSystemConfigCache() {
|
||||
if err := repository.InvalidateAllSystemConfigCaches(context.Background()); err != nil {
|
||||
log.Printf("[%s] clear system config cache failed: %v\n", dbType(), err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,132 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package migrator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/alicebob/miniredis/v2"
|
||||
"github.com/glebarez/sqlite"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"github.com/redis/go-redis/v9/maintnotifications"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const expectedMigratedSystemConfigCount = 30
|
||||
|
||||
func TestMigrateInitializesSQLiteDatabase(t *testing.T) {
|
||||
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("gorm.Open(sqlite) error = %v", err)
|
||||
}
|
||||
|
||||
mr, err := miniredis.Run()
|
||||
if err != nil {
|
||||
t.Fatalf("miniredis.Run() error = %v", err)
|
||||
}
|
||||
redisClient := redis.NewClient(&redis.Options{
|
||||
Addr: mr.Addr(),
|
||||
MaintNotificationsConfig: &maintnotifications.Config{
|
||||
Mode: maintnotifications.ModeDisabled,
|
||||
},
|
||||
})
|
||||
|
||||
previousDBEnabled := config.Config.Database.Enabled
|
||||
config.Config.Database.Enabled = false
|
||||
db.SetDB(sqliteDB)
|
||||
t.Cleanup(func() {
|
||||
config.Config.Database.Enabled = previousDBEnabled
|
||||
db.SetDB(nil)
|
||||
_ = redisClient.Close()
|
||||
mr.Close()
|
||||
})
|
||||
|
||||
Migrate()
|
||||
|
||||
var systemConfigCount int64
|
||||
if err := sqliteDB.Table("w_system_configs").Count(&systemConfigCount).Error; err != nil {
|
||||
t.Fatalf("Migrate() count w_system_configs error = %v", err)
|
||||
}
|
||||
if systemConfigCount != expectedMigratedSystemConfigCount {
|
||||
t.Errorf("Migrate() w_system_configs count = %d, want %d", systemConfigCount, expectedMigratedSystemConfigCount)
|
||||
}
|
||||
|
||||
var adminCount int64
|
||||
if err := sqliteDB.Table("w_users").Where("username = ?", "admin").Count(&adminCount).Error; err != nil {
|
||||
t.Fatalf("Migrate() count admin user error = %v", err)
|
||||
}
|
||||
if adminCount != 1 {
|
||||
t.Errorf("Migrate() admin user count = %d, want %d", adminCount, 1)
|
||||
}
|
||||
|
||||
var templateCount int64
|
||||
if err := sqliteDB.Table("w_templates").Count(&templateCount).Error; err != nil {
|
||||
t.Fatalf("Migrate() count templates error = %v", err)
|
||||
}
|
||||
if templateCount != 2 {
|
||||
t.Errorf("Migrate() templates count = %d, want %d", templateCount, 2)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateClearsStaleSystemConfigCache(t *testing.T) {
|
||||
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("gorm.Open(sqlite) error = %v", err)
|
||||
}
|
||||
|
||||
mr, err := miniredis.Run()
|
||||
if err != nil {
|
||||
t.Fatalf("miniredis.Run() error = %v", err)
|
||||
}
|
||||
redisClient := redis.NewClient(&redis.Options{Addr: mr.Addr()})
|
||||
|
||||
previousDBEnabled := config.Config.Database.Enabled
|
||||
previousRedis := db.Redis
|
||||
config.Config.Database.Enabled = false
|
||||
db.SetDB(sqliteDB)
|
||||
db.Redis = redisClient
|
||||
t.Cleanup(func() {
|
||||
config.Config.Database.Enabled = previousDBEnabled
|
||||
db.SetDB(nil)
|
||||
db.Redis = previousRedis
|
||||
_ = redisClient.Close()
|
||||
mr.Close()
|
||||
})
|
||||
|
||||
staleConfig := model.SystemConfig{
|
||||
Key: model.ConfigKeyCapLoginEnabled,
|
||||
Value: "true",
|
||||
Type: "system",
|
||||
}
|
||||
if err := db.HSetJSON(context.Background(), repository.SystemConfigRedisHashKey, model.ConfigKeyCapLoginEnabled, &staleConfig); err != nil {
|
||||
t.Fatalf("HSetJSON() error = %v", err)
|
||||
}
|
||||
|
||||
Migrate()
|
||||
|
||||
exists, err := db.Redis.Exists(context.Background(), db.PrefixedKey(repository.SystemConfigRedisHashKey)).Result()
|
||||
if err != nil {
|
||||
t.Fatalf("Redis.Exists() error = %v", err)
|
||||
}
|
||||
if exists != 0 {
|
||||
t.Fatalf("system config cache exists = %d, want 0", exists)
|
||||
}
|
||||
|
||||
enabled, err := repository.GetBoolByKey(context.Background(), model.ConfigKeyCapLoginEnabled)
|
||||
if err != nil {
|
||||
t.Fatalf("GetBoolByKey(%s) error = %v", model.ConfigKeyCapLoginEnabled, err)
|
||||
}
|
||||
if enabled {
|
||||
t.Fatalf("GetBoolByKey(%s) = true, want false", model.ConfigKeyCapLoginEnabled)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,217 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/glebarez/sqlite"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/plugin/dbresolver"
|
||||
"gorm.io/plugin/opentelemetry/tracing"
|
||||
)
|
||||
|
||||
var (
|
||||
db *gorm.DB
|
||||
)
|
||||
|
||||
func init() {
|
||||
if !config.Config.Database.Enabled {
|
||||
// PostgreSQL 禁用,使用 SQLite
|
||||
initSQLite()
|
||||
return
|
||||
}
|
||||
|
||||
initPostgres()
|
||||
}
|
||||
|
||||
// initSQLite 初始化 SQLite 数据库(PostgreSQL 禁用时的后备方案)
|
||||
func initSQLite() {
|
||||
sqlitePath := config.Config.Database.SQLitePath
|
||||
if sqlitePath == "" {
|
||||
sqlitePath = "./data/wavelet.db"
|
||||
}
|
||||
|
||||
var err error
|
||||
db, err = gorm.Open(sqlite.Open(sqlitePath), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
Logger: &gormZapLogger{
|
||||
logLevel: parseLogLevel(config.Config.Database.LogLevel),
|
||||
slowThreshold: config.Config.Database.SlowThreshold,
|
||||
ignoreRecordNotFoundError: config.Config.App.IsProduction(),
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("[SQLite] init connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
// Trace 注入
|
||||
if err = db.Use(
|
||||
tracing.NewPlugin(
|
||||
tracing.WithoutMetrics(),
|
||||
tracing.WithAttributes(
|
||||
attribute.String("db.instance", sqlitePath),
|
||||
attribute.String("db.system", "SQLite"),
|
||||
),
|
||||
),
|
||||
); err != nil {
|
||||
log.Fatalf("[SQLite] init trace failed: %v\n", err)
|
||||
}
|
||||
|
||||
log.Printf("[SQLite] initialized (path: %s)\n", sqlitePath)
|
||||
}
|
||||
|
||||
// initPostgres 初始化 PostgreSQL 数据库
|
||||
func initPostgres() {
|
||||
var err error
|
||||
dbConfig := config.Config.Database
|
||||
|
||||
// 构建主库 DSN 并连接
|
||||
primaryDSN := buildDSN(dbConfig.Host, dbConfig.Port, dbConfig.Username, dbConfig.Password)
|
||||
|
||||
pgConfig := postgres.Config{
|
||||
DSN: primaryDSN,
|
||||
PreferSimpleProtocol: dbConfig.PreferSimpleProtocol,
|
||||
}
|
||||
|
||||
db, err = gorm.Open(postgres.New(pgConfig), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
Logger: &gormZapLogger{
|
||||
logLevel: parseLogLevel(config.Config.Database.LogLevel),
|
||||
slowThreshold: config.Config.Database.SlowThreshold,
|
||||
ignoreRecordNotFoundError: config.Config.App.IsProduction(),
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("[PostgreSQL] init connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
// Trace 注入
|
||||
if err = db.Use(
|
||||
tracing.NewPlugin(
|
||||
tracing.WithoutMetrics(),
|
||||
tracing.WithAttributes(
|
||||
attribute.String("db.instance", dbConfig.Database),
|
||||
attribute.String("db.ip", dbConfig.Host),
|
||||
attribute.String("server.address", net.JoinHostPort(dbConfig.Host, strconv.Itoa(dbConfig.Port))),
|
||||
attribute.String("db.system", "PostgreSQL"),
|
||||
),
|
||||
),
|
||||
); err != nil {
|
||||
log.Fatalf("[PostgreSQL] init trace failed: %v\n", err)
|
||||
}
|
||||
|
||||
if len(dbConfig.Replicas) > 0 {
|
||||
var replicaDialectors []gorm.Dialector
|
||||
for _, replica := range dbConfig.Replicas {
|
||||
username := replica.Username
|
||||
if username == "" {
|
||||
username = dbConfig.Username
|
||||
}
|
||||
password := replica.Password
|
||||
if password == "" {
|
||||
password = dbConfig.Password
|
||||
}
|
||||
replicaDSN := buildDSN(replica.Host, replica.Port, username, password)
|
||||
replicaDialectors = append(replicaDialectors, postgres.New(postgres.Config{
|
||||
DSN: replicaDSN,
|
||||
PreferSimpleProtocol: dbConfig.PreferSimpleProtocol,
|
||||
}))
|
||||
}
|
||||
|
||||
resolver := dbresolver.Register(dbresolver.Config{
|
||||
Replicas: replicaDialectors,
|
||||
Policy: dbresolver.RandomPolicy{},
|
||||
})
|
||||
|
||||
resolver.SetMaxIdleConns(dbConfig.MaxIdleConn).
|
||||
SetMaxOpenConns(dbConfig.MaxOpenConn).
|
||||
SetConnMaxLifetime(time.Duration(dbConfig.ConnMaxLifetime) * time.Second).
|
||||
SetConnMaxIdleTime(time.Duration(dbConfig.ConnMaxIdleTime) * time.Second)
|
||||
|
||||
if err = db.Use(resolver); err != nil {
|
||||
log.Fatalf("[PostgreSQL] init dbresolver failed: %v\n", err)
|
||||
}
|
||||
log.Printf("[PostgreSQL] initialized in Primary-Replica mode (%d replicas)\n", len(dbConfig.Replicas))
|
||||
} else {
|
||||
log.Println("[PostgreSQL] initialized in Standalone mode")
|
||||
}
|
||||
|
||||
// 获取通用数据库对象设置连接池
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
log.Fatalf("[PostgreSQL] load sql db failed: %v\n", err)
|
||||
}
|
||||
|
||||
sqlDB.SetMaxIdleConns(dbConfig.MaxIdleConn)
|
||||
sqlDB.SetMaxOpenConns(dbConfig.MaxOpenConn)
|
||||
sqlDB.SetConnMaxLifetime(time.Duration(dbConfig.ConnMaxLifetime) * time.Second)
|
||||
sqlDB.SetConnMaxIdleTime(time.Duration(dbConfig.ConnMaxIdleTime) * time.Second)
|
||||
|
||||
}
|
||||
|
||||
// buildDSN 构建 PostgreSQL DSN
|
||||
func buildDSN(host string, port int, username, password string) string {
|
||||
cfg := config.Config.Database
|
||||
pqURL := &url.URL{
|
||||
Scheme: "postgres",
|
||||
Host: net.JoinHostPort(host, strconv.Itoa(port)),
|
||||
Path: cfg.Database,
|
||||
}
|
||||
if username != "" {
|
||||
pqURL.User = url.UserPassword(username, password)
|
||||
}
|
||||
|
||||
query := pqURL.Query()
|
||||
sslMode := cfg.SSLMode
|
||||
if sslMode == "" {
|
||||
sslMode = "disable"
|
||||
}
|
||||
query.Set("sslmode", sslMode)
|
||||
if cfg.ApplicationName != "" {
|
||||
query.Set("application_name", cfg.ApplicationName)
|
||||
}
|
||||
if cfg.SearchPath != "" {
|
||||
query.Set("search_path", cfg.SearchPath)
|
||||
}
|
||||
if cfg.DefaultQueryExecMode != "" {
|
||||
query.Set("default_query_exec_mode", cfg.DefaultQueryExecMode)
|
||||
}
|
||||
if cfg.StatementCacheCapacity > 0 {
|
||||
query.Set("statement_cache_capacity", strconv.Itoa(cfg.StatementCacheCapacity))
|
||||
}
|
||||
|
||||
rawQuery := query.Encode()
|
||||
if cfg.TimeZone != "" {
|
||||
if rawQuery != "" {
|
||||
rawQuery += "&"
|
||||
}
|
||||
rawQuery += "TimeZone=" + cfg.TimeZone
|
||||
}
|
||||
pqURL.RawQuery = rawQuery
|
||||
|
||||
return pqURL.String()
|
||||
}
|
||||
|
||||
// DB 返回带上下文追踪的 GORM 数据库实例
|
||||
func DB(ctx context.Context) *gorm.DB {
|
||||
if db == nil {
|
||||
return nil
|
||||
}
|
||||
return db.WithContext(ctx)
|
||||
}
|
||||
|
||||
// SetDB sets the package-level database instance for testing.
|
||||
func SetDB(d *gorm.DB) {
|
||||
db = d
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
gormLogger "gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
// nanoToMilli 纳秒转毫秒的除数
|
||||
const nanoToMilli = 1e6
|
||||
|
||||
type gormZapLogger struct {
|
||||
logLevel gormLogger.LogLevel
|
||||
ignoreRecordNotFoundError bool
|
||||
slowThreshold time.Duration
|
||||
}
|
||||
|
||||
func (l *gormZapLogger) LogMode(level gormLogger.LogLevel) gormLogger.Interface {
|
||||
clone := *l
|
||||
clone.logLevel = level
|
||||
return &clone
|
||||
}
|
||||
|
||||
func (l *gormZapLogger) Info(ctx context.Context, fmt string, args ...interface{}) {
|
||||
if l.logLevel >= gormLogger.Info {
|
||||
logger.InfoF(ctx, fmt, args...)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *gormZapLogger) Warn(ctx context.Context, fmt string, args ...interface{}) {
|
||||
if l.logLevel >= gormLogger.Warn {
|
||||
logger.WarnF(ctx, fmt, args...)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *gormZapLogger) Error(ctx context.Context, fmt string, args ...interface{}) {
|
||||
if l.logLevel >= gormLogger.Error {
|
||||
logger.ErrorF(ctx, fmt, args...)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *gormZapLogger) Trace(ctx context.Context, begin time.Time, fc func() (sql string, rowsAffected int64), err error) {
|
||||
elapsed := time.Since(begin)
|
||||
switch {
|
||||
case err != nil && l.logLevel >= gormLogger.Error && (!errors.Is(err, gorm.ErrRecordNotFound) || !l.ignoreRecordNotFoundError):
|
||||
_, rows := fc()
|
||||
logger.ErrorF(ctx, "database query failed: %s [%.3fms] [rows:%v]", err, float64(elapsed.Nanoseconds())/nanoToMilli, formatRows(rows))
|
||||
case elapsed > l.slowThreshold && l.slowThreshold != 0 && l.logLevel >= gormLogger.Warn:
|
||||
_, rows := fc()
|
||||
slowLog := fmt.Sprintf("SLOW SQL >= %v", l.slowThreshold)
|
||||
logger.WarnF(ctx, "%s [%.3fms] [rows:%v]", slowLog, float64(elapsed.Nanoseconds())/nanoToMilli, formatRows(rows))
|
||||
case l.logLevel == gormLogger.Info:
|
||||
sql, rows := fc()
|
||||
logger.DebugF(ctx, "[%.3fms] [rows:%v] %s", float64(elapsed.Nanoseconds())/nanoToMilli, formatRows(rows), sql)
|
||||
}
|
||||
}
|
||||
|
||||
func formatRows(rows int64) interface{} {
|
||||
if rows == -1 {
|
||||
return "-"
|
||||
}
|
||||
return rows
|
||||
}
|
||||
|
||||
func parseLogLevel(level string) gormLogger.LogLevel {
|
||||
level = strings.ToLower(level)
|
||||
switch level {
|
||||
case "silent":
|
||||
return gormLogger.Silent
|
||||
case "error":
|
||||
return gormLogger.Error
|
||||
case "warn":
|
||||
return gormLogger.Warn
|
||||
case "info":
|
||||
return gormLogger.Info
|
||||
case "debug":
|
||||
return gormLogger.Info
|
||||
default:
|
||||
return gormLogger.Info
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package db
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
gormLogger "gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
func TestParseLogLevel(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
configuredLevel string
|
||||
want gormLogger.LogLevel
|
||||
}{
|
||||
{
|
||||
name: "debug enables SQL trace processing",
|
||||
configuredLevel: "debug",
|
||||
want: gormLogger.Info,
|
||||
},
|
||||
{
|
||||
name: "development preserves configured level",
|
||||
configuredLevel: "warn",
|
||||
want: gormLogger.Warn,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if got := parseLogLevel(tt.configuredLevel); got != tt.want {
|
||||
t.Fatalf("parseLogLevel() = %v, want %v", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,198 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/redis/go-redis/extra/redisotel/v9"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"github.com/redis/go-redis/v9/maintnotifications"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
)
|
||||
|
||||
var (
|
||||
// Redis 全局 Redis 客户端实例
|
||||
Redis redis.UniversalClient
|
||||
)
|
||||
|
||||
func init() {
|
||||
cfg := config.Config.Redis
|
||||
|
||||
if !cfg.Enabled {
|
||||
log.Println("[Redis] is disabled, skipping Redis initialization")
|
||||
return
|
||||
}
|
||||
|
||||
if cfg.ClusterMode {
|
||||
// Cluster 模式
|
||||
Redis = redis.NewClusterClient(&redis.ClusterOptions{
|
||||
Addrs: cfg.Addrs,
|
||||
Username: cfg.Username,
|
||||
Password: cfg.Password,
|
||||
PoolSize: cfg.PoolSize,
|
||||
MinIdleConns: cfg.MinIdleConn,
|
||||
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
||||
ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second,
|
||||
WriteTimeout: time.Duration(cfg.WriteTimeout) * time.Second,
|
||||
MaxRetries: cfg.MaxRetries,
|
||||
PoolTimeout: time.Duration(cfg.PoolTimeout) * time.Second,
|
||||
ConnMaxIdleTime: time.Duration(cfg.ConnMaxIdleTime) * time.Second,
|
||||
MaintNotificationsConfig: redisMaintNotificationsConfig(cfg.MaintNotifications),
|
||||
})
|
||||
log.Println("[Redis] initialized in Cluster mode")
|
||||
} else {
|
||||
// Standalone 或 Sentinel 模式
|
||||
options := &redis.UniversalOptions{
|
||||
Addrs: cfg.Addrs,
|
||||
MasterName: cfg.MasterName, // 非空时启用 Sentinel
|
||||
Username: cfg.Username,
|
||||
Password: cfg.Password,
|
||||
DB: cfg.DB,
|
||||
PoolSize: cfg.PoolSize,
|
||||
MinIdleConns: cfg.MinIdleConn,
|
||||
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
||||
ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second,
|
||||
WriteTimeout: time.Duration(cfg.WriteTimeout) * time.Second,
|
||||
MaxRetries: cfg.MaxRetries,
|
||||
PoolTimeout: time.Duration(cfg.PoolTimeout) * time.Second,
|
||||
ConnMaxIdleTime: time.Duration(cfg.ConnMaxIdleTime) * time.Second,
|
||||
MaintNotificationsConfig: redisMaintNotificationsConfig(cfg.MaintNotifications),
|
||||
}
|
||||
if cfg.MasterName != "" {
|
||||
client := redis.NewFailoverClient(options.Failover())
|
||||
// FailoverOptions 暂不暴露该配置,在首次建连前写入客户端选项。
|
||||
client.Options().MaintNotificationsConfig = redisMaintNotificationsConfig(cfg.MaintNotifications)
|
||||
Redis = client
|
||||
log.Println("[Redis] initialized in Sentinel mode")
|
||||
} else {
|
||||
Redis = redis.NewUniversalClient(options)
|
||||
log.Println("[Redis] initialized in Standalone mode")
|
||||
}
|
||||
}
|
||||
|
||||
// OpenTelemetry 追踪(UniversalClient 兼容)
|
||||
if err := redisotel.InstrumentTracing(
|
||||
Redis,
|
||||
redisotel.WithAttributes(
|
||||
attribute.String("db.instance", fmt.Sprintf("%v", cfg.DB)),
|
||||
attribute.String("db.ip", strings.Join(cfg.Addrs, ",")),
|
||||
attribute.String("db.system", "Redis"),
|
||||
),
|
||||
); err != nil {
|
||||
log.Fatalf("[Redis] failed to init trace: %v\n", err)
|
||||
}
|
||||
|
||||
// 测试连接
|
||||
_, err := Redis.Ping(context.Background()).Result()
|
||||
if err != nil {
|
||||
log.Fatalf("[Redis] failed to connect to redis: %v\n", err)
|
||||
}
|
||||
}
|
||||
|
||||
func redisMaintNotificationsConfig(enabled bool) *maintnotifications.Config {
|
||||
mode := maintnotifications.ModeDisabled
|
||||
if enabled {
|
||||
mode = maintnotifications.ModeAuto
|
||||
}
|
||||
return &maintnotifications.Config{Mode: mode}
|
||||
}
|
||||
|
||||
// PrefixedKey 返回带前缀的 Key
|
||||
func PrefixedKey(key string) string {
|
||||
prefix := config.Config.Redis.KeyPrefix
|
||||
if prefix == "" {
|
||||
return key
|
||||
}
|
||||
return prefix + key
|
||||
}
|
||||
|
||||
// HSetJSON 将泛型数据序列化为 JSON 并设置到 Redis Hash
|
||||
// ctx: 上下文
|
||||
// hashKey: Redis Hash key
|
||||
// fieldKey: Hash field key
|
||||
// data: 要存储的数据(泛型)
|
||||
func HSetJSON[T any](ctx context.Context, hashKey, fieldKey string, data T) error {
|
||||
jsonData, err := json.Marshal(data)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := Redis.HSet(ctx, PrefixedKey(hashKey), fieldKey, jsonData).Err(); err != nil {
|
||||
return fmt.Errorf(errRedisHashSetFailed, err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// HDel removes one or more fields from a Redis Hash.
|
||||
func HDel(ctx context.Context, hashKey string, fieldKeys ...string) error {
|
||||
if Redis == nil || len(fieldKeys) == 0 {
|
||||
return nil
|
||||
}
|
||||
if err := Redis.HDel(ctx, PrefixedKey(hashKey), fieldKeys...).Err(); err != nil {
|
||||
return fmt.Errorf(errRedisHashDeleteFailed, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// HGetJSON 从 Redis Hash 获取数据并反序列化为泛型类型
|
||||
// ctx: 上下文
|
||||
// hashKey: Redis Hash key
|
||||
// fieldKey: Hash field key
|
||||
// data: 用于接收数据的指针(泛型)
|
||||
func HGetJSON[T any](ctx context.Context, hashKey, fieldKey string, data *T) error {
|
||||
val, err := Redis.HGet(ctx, PrefixedKey(hashKey), fieldKey).Result()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := json.Unmarshal([]byte(val), data); err != nil {
|
||||
return fmt.Errorf(errUnmarshalDataFailed, err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetJSON 从Redis获取数据并反序列化为泛型类型
|
||||
// ctx: 上下文
|
||||
// key: Redis key
|
||||
// data: 用于接收数据的指针(泛型)
|
||||
func GetJSON[T any](ctx context.Context, key string, data *T) error {
|
||||
val, err := Redis.Get(ctx, PrefixedKey(key)).Bytes()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := json.Unmarshal(val, data); err != nil {
|
||||
return fmt.Errorf(errUnmarshalDataFailed, err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetJSON 将泛型数据序列化为JSON并设置到Redis
|
||||
// ctx: 上下文
|
||||
// key: Redis key
|
||||
// data: 要存储的数据(泛型)
|
||||
// expiration: 过期时间
|
||||
func SetJSON[T any](ctx context.Context, key string, data T, expiration time.Duration) error {
|
||||
jsonData, err := json.Marshal(data)
|
||||
if err != nil {
|
||||
return fmt.Errorf(errMarshalDataFailed, err)
|
||||
}
|
||||
|
||||
if err := Redis.Set(ctx, PrefixedKey(key), jsonData, expiration).Err(); err != nil {
|
||||
return fmt.Errorf(errRedisKeySetFailed, err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/redis/go-redis/v9/maintnotifications"
|
||||
)
|
||||
|
||||
func TestRedisMaintNotificationsConfig(t *testing.T) {
|
||||
for _, test := range []struct {
|
||||
name string
|
||||
enabled bool
|
||||
want maintnotifications.Mode
|
||||
}{
|
||||
{name: "disabled by default", enabled: false, want: maintnotifications.ModeDisabled},
|
||||
{name: "auto when enabled", enabled: true, want: maintnotifications.ModeAuto},
|
||||
} {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
cfg := redisMaintNotificationsConfig(test.enabled)
|
||||
if cfg.Mode != test.want {
|
||||
t.Fatalf("maintenance notifications mode = %v, want %v", cfg.Mode, test.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user