mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
fix(log): PG 日志库批量写入为零 ID 行生成雪花 ID
PostgreSQL 日志表 id 为 NOT NULL 且无默认值,而 GORM 将零值 uint64 主键视为自增并省略 id 列,导致 node access log / 可观测指标等批量 落库持续报 "null value in column id violates not-null constraint"。 在 BatchInsert* 落库前为零 ID 行生成雪花 ID(与 ClickHouse 写入路径 一致),并新增单元回归与 PG 集成回归测试覆盖六张日志表。
This commit is contained in:
@@ -20,6 +20,7 @@ sidebar: false
|
|||||||
|
|
||||||
### 🛠 修复
|
### 🛠 修复
|
||||||
- 修复源站错误页「仅针对 GET 请求」未生效:`error_page` 内部重定向会把请求方法改写成 GET,导致内部 Lua 无法识别 POST/PUT 等原始方法、仍返回自定义错误页;现改为命名 location(`@__openflare_origin_error`)承载错误页,保留原始请求方法与错误状态码,非 GET 请求不再返回自定义错误页。
|
- 修复源站错误页「仅针对 GET 请求」未生效:`error_page` 内部重定向会把请求方法改写成 GET,导致内部 Lua 无法识别 POST/PUT 等原始方法、仍返回自定义错误页;现改为命名 location(`@__openflare_origin_error`)承载错误页,保留原始请求方法与错误状态码,非 GET 请求不再返回自定义错误页。
|
||||||
|
- 修复 PostgreSQL 作为日志库时节点访问日志/可观测指标/用户访问日志批量写入失败:GORM 对零值 `uint64` 主键会省略 `id` 列,而 PG 日志表 `id` 无默认值,导致持续报「null value in column id violates not-null constraint」;现于落库前为零 ID 行生成雪花 ID(与 ClickHouse 写入路径一致),并新增回归测试覆盖六张日志表。
|
||||||
|
|
||||||
## [v3.5.1] - 2026-08-09
|
## [v3.5.1] - 2026-08-09
|
||||||
|
|
||||||
|
|||||||
@@ -459,6 +459,159 @@ func TestDropExpiredPartitionsTimezoneSafety(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestBatchInsertGeneratesIDsPostgres 回归:PG 日志表 id BIGINT NOT NULL 且无默认值;
|
||||||
|
// GORM 把零值 uint64 主键视为自增并省略 id 列,直接插入会报 23502 not-null 违例。
|
||||||
|
// 验证 6 张日志表 BatchInsert* 为零 ID 行生成雪花 ID 后正常落库(修复前本测试失败)。
|
||||||
|
func TestBatchInsertGeneratesIDsPostgres(t *testing.T) {
|
||||||
|
dsn := strings.TrimSpace(os.Getenv("TEST_POSTGRES_DSN"))
|
||||||
|
if dsn == "" {
|
||||||
|
t.Skip("TEST_POSTGRES_DSN is not set")
|
||||||
|
}
|
||||||
|
|
||||||
|
gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{
|
||||||
|
DisableForeignKeyConstraintWhenMigrating: true,
|
||||||
|
Logger: logger.Default.LogMode(logger.Silent),
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("open postgres: %v", err)
|
||||||
|
}
|
||||||
|
sqlDB, err := gdb.DB()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("sql db: %v", err)
|
||||||
|
}
|
||||||
|
sqlDB.SetMaxOpenConns(1)
|
||||||
|
|
||||||
|
schema := fmt.Sprintf("logstore_ids_%d", time.Now().UnixNano())
|
||||||
|
if !regexp.MustCompile(`^[a-z0-9_]+$`).MatchString(schema) {
|
||||||
|
t.Fatalf("invalid schema: %s", schema)
|
||||||
|
}
|
||||||
|
if err := gdb.Exec(`CREATE SCHEMA "` + schema + `"`).Error; err != nil {
|
||||||
|
t.Fatalf("create schema: %v", err)
|
||||||
|
}
|
||||||
|
if err := gdb.Exec(`SET search_path TO "` + schema + `"`).Error; err != nil {
|
||||||
|
t.Fatalf("set search_path: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() {
|
||||||
|
_ = gdb.Exec("SET search_path TO public").Error
|
||||||
|
_ = gdb.Exec(`DROP SCHEMA IF EXISTS "` + schema + `" CASCADE`).Error
|
||||||
|
_ = sqlDB.Close()
|
||||||
|
})
|
||||||
|
|
||||||
|
for _, ddl := range []string{
|
||||||
|
postgresNodeAccessLogsDDL,
|
||||||
|
postgresUserAccessLogsDDL,
|
||||||
|
postgresMetricSnapshotsDDL,
|
||||||
|
postgresEdgeHealthDDL,
|
||||||
|
postgresObsFrpsDDL,
|
||||||
|
postgresObsFrpcDDL,
|
||||||
|
} {
|
||||||
|
if err := gdb.Exec(ddl).Error; err != nil {
|
||||||
|
t.Fatalf("create table: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
ResetForTest()
|
||||||
|
SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil })
|
||||||
|
defer ResetForTest()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
store := newGormStore(gdb)
|
||||||
|
ua := newUserAccessLogGormStore(gdb)
|
||||||
|
|
||||||
|
now := time.Now().UTC()
|
||||||
|
if err := store.EnsurePartitions(ctx, now, now.AddDate(0, 1, 0)); err != nil {
|
||||||
|
t.Fatalf("EnsurePartitions: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
nodeRows := []analyticsmodel.NodeAccessLog{
|
||||||
|
{NodeID: "n1", LoggedAt: now, RemoteAddr: "1.1.1.1", StatusCode: 200},
|
||||||
|
{NodeID: "n1", LoggedAt: now.Add(time.Second), RemoteAddr: "2.2.2.2", StatusCode: 500},
|
||||||
|
}
|
||||||
|
if err := store.BatchInsertNodeAccessLogs(ctx, nodeRows); err != nil {
|
||||||
|
t.Fatalf("insert node access logs with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if nodeRows[0].ID == 0 || nodeRows[1].ID == 0 || nodeRows[0].ID == nodeRows[1].ID {
|
||||||
|
t.Fatalf("node access log ids not generated: %+v", nodeRows)
|
||||||
|
}
|
||||||
|
|
||||||
|
metricRows := []analyticsmodel.NodeMetricSnapshot{
|
||||||
|
{NodeID: "n1", CapturedAt: now},
|
||||||
|
{NodeID: "n2", CapturedAt: now},
|
||||||
|
}
|
||||||
|
if err := store.BatchInsertNodeMetricSnapshots(ctx, metricRows); err != nil {
|
||||||
|
t.Fatalf("insert metric snapshots with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if metricRows[0].ID == 0 || metricRows[1].ID == 0 || metricRows[0].ID == metricRows[1].ID {
|
||||||
|
t.Fatalf("metric snapshot ids not generated: %+v", metricRows)
|
||||||
|
}
|
||||||
|
|
||||||
|
edgeRows := []analyticsmodel.NodeEdgeHealth{
|
||||||
|
{NodeID: "n1", CapturedAt: now, Status: "ok"},
|
||||||
|
{NodeID: "n2", CapturedAt: now, Status: "ok"},
|
||||||
|
}
|
||||||
|
if err := store.BatchInsertNodeEdgeHealth(ctx, edgeRows); err != nil {
|
||||||
|
t.Fatalf("insert edge health with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if edgeRows[0].ID == 0 || edgeRows[1].ID == 0 || edgeRows[0].ID == edgeRows[1].ID {
|
||||||
|
t.Fatalf("edge health ids not generated: %+v", edgeRows)
|
||||||
|
}
|
||||||
|
|
||||||
|
frpsRows := []analyticsmodel.NodeObsFrps{
|
||||||
|
{NodeID: "n1", CapturedAt: now, FrpsConnections: 1},
|
||||||
|
{NodeID: "n2", CapturedAt: now, FrpsConnections: 2},
|
||||||
|
}
|
||||||
|
if err := store.BatchInsertNodeObsFrps(ctx, frpsRows); err != nil {
|
||||||
|
t.Fatalf("insert obs frps with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if frpsRows[0].ID == 0 || frpsRows[1].ID == 0 || frpsRows[0].ID == frpsRows[1].ID {
|
||||||
|
t.Fatalf("obs frps ids not generated: %+v", frpsRows)
|
||||||
|
}
|
||||||
|
|
||||||
|
frpcRows := []analyticsmodel.NodeObsFrpc{
|
||||||
|
{NodeID: "n1", CapturedAt: now, TunnelStatus: "online"},
|
||||||
|
{NodeID: "n2", CapturedAt: now, TunnelStatus: "online"},
|
||||||
|
}
|
||||||
|
if err := store.BatchInsertNodeObsFrpc(ctx, frpcRows); err != nil {
|
||||||
|
t.Fatalf("insert obs frpc with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if frpcRows[0].ID == 0 || frpcRows[1].ID == 0 || frpcRows[0].ID == frpcRows[1].ID {
|
||||||
|
t.Fatalf("obs frpc ids not generated: %+v", frpcRows)
|
||||||
|
}
|
||||||
|
|
||||||
|
userRows := []analyticsmodel.UserAccessLog{
|
||||||
|
{UserID: 101, Path: "/a", CreatedAt: now},
|
||||||
|
{UserID: 102, Path: "/b", CreatedAt: now},
|
||||||
|
}
|
||||||
|
if err := ua.BatchInsert(ctx, userRows); err != nil {
|
||||||
|
t.Fatalf("insert user access logs with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if userRows[0].ID == 0 || userRows[1].ID == 0 || userRows[0].ID == userRows[1].ID {
|
||||||
|
t.Fatalf("user access log ids not generated: %+v", userRows)
|
||||||
|
}
|
||||||
|
|
||||||
|
expect := []struct {
|
||||||
|
name string
|
||||||
|
model any
|
||||||
|
want int64
|
||||||
|
}{
|
||||||
|
{"of_node_access_logs", &analyticsmodel.NodeAccessLog{}, 2},
|
||||||
|
{"of_node_metric_snapshots", &analyticsmodel.NodeMetricSnapshot{}, 2},
|
||||||
|
{"of_node_edge_health", &analyticsmodel.NodeEdgeHealth{}, 2},
|
||||||
|
{"of_node_obs_frps", &analyticsmodel.NodeObsFrps{}, 2},
|
||||||
|
{"of_node_obs_frpc", &analyticsmodel.NodeObsFrpc{}, 2},
|
||||||
|
{"w_user_access_logs", &analyticsmodel.UserAccessLog{}, 2},
|
||||||
|
}
|
||||||
|
for _, e := range expect {
|
||||||
|
var got int64
|
||||||
|
if err := gdb.Model(e.model).Count(&got).Error; err != nil {
|
||||||
|
t.Fatalf("count %s: %v", e.name, err)
|
||||||
|
}
|
||||||
|
if got != e.want {
|
||||||
|
t.Fatalf("%s count = %d, want %d", e.name, got, e.want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// postgresNodeAccessLogsDDL 与 goose/postgres/202608080001_create_log_tables.sql 对齐。
|
// postgresNodeAccessLogsDDL 与 goose/postgres/202608080001_create_log_tables.sql 对齐。
|
||||||
const postgresNodeAccessLogsDDL = `
|
const postgresNodeAccessLogsDDL = `
|
||||||
CREATE TABLE IF NOT EXISTS of_node_access_logs (
|
CREATE TABLE IF NOT EXISTS of_node_access_logs (
|
||||||
@@ -494,3 +647,54 @@ CREATE TABLE IF NOT EXISTS w_user_access_logs (
|
|||||||
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||||
PRIMARY KEY (id, created_at)
|
PRIMARY KEY (id, created_at)
|
||||||
) PARTITION BY RANGE (created_at)`
|
) PARTITION BY RANGE (created_at)`
|
||||||
|
|
||||||
|
// postgresMetricSnapshotsDDL / postgresEdgeHealthDDL / postgresObsFrpsDDL / postgresObsFrpcDDL
|
||||||
|
// 与 goose/postgres/202608080001_create_log_tables.sql 对齐(普通表,无分区)。
|
||||||
|
const postgresMetricSnapshotsDDL = `
|
||||||
|
CREATE TABLE IF NOT EXISTS of_node_metric_snapshots (
|
||||||
|
id BIGINT NOT NULL PRIMARY KEY,
|
||||||
|
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||||
|
captured_at TIMESTAMPTZ NOT NULL,
|
||||||
|
cpu_usage_percent DOUBLE PRECISION NOT NULL DEFAULT 0,
|
||||||
|
memory_used_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
memory_total_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
storage_used_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
storage_total_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
disk_read_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
disk_write_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
network_rx_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
network_tx_bytes BIGINT NOT NULL DEFAULT 0,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||||
|
)`
|
||||||
|
|
||||||
|
const postgresEdgeHealthDDL = `
|
||||||
|
CREATE TABLE IF NOT EXISTS of_node_edge_health (
|
||||||
|
id BIGINT NOT NULL PRIMARY KEY,
|
||||||
|
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||||
|
captured_at TIMESTAMPTZ NOT NULL,
|
||||||
|
status VARCHAR(64) NOT NULL DEFAULT '',
|
||||||
|
connections BIGINT NOT NULL DEFAULT 0,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||||
|
)`
|
||||||
|
|
||||||
|
const postgresObsFrpsDDL = `
|
||||||
|
CREATE TABLE IF NOT EXISTS of_node_obs_frps (
|
||||||
|
id BIGINT NOT NULL PRIMARY KEY,
|
||||||
|
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||||
|
captured_at TIMESTAMPTZ NOT NULL,
|
||||||
|
frps_connections INTEGER NOT NULL DEFAULT 0,
|
||||||
|
frps_proxy_count INTEGER NOT NULL DEFAULT 0,
|
||||||
|
frps_client_count INTEGER NOT NULL DEFAULT 0,
|
||||||
|
frps_proxies TEXT NOT NULL DEFAULT '',
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||||
|
)`
|
||||||
|
|
||||||
|
const postgresObsFrpcDDL = `
|
||||||
|
CREATE TABLE IF NOT EXISTS of_node_obs_frpc (
|
||||||
|
id BIGINT NOT NULL PRIMARY KEY,
|
||||||
|
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||||
|
captured_at TIMESTAMPTZ NOT NULL,
|
||||||
|
tunnel_status VARCHAR(16) NOT NULL DEFAULT '',
|
||||||
|
connected_relays_count INTEGER NOT NULL DEFAULT 0,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||||
|
)`
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/infra/persistence/idgen"
|
||||||
"github.com/Rain-kl/Wavelet/internal/model"
|
"github.com/Rain-kl/Wavelet/internal/model"
|
||||||
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
@@ -113,6 +114,13 @@ func (s *gormLogStore) BatchInsertNodeAccessLogs(ctx context.Context, rows []ana
|
|||||||
if err := s.ensureWritable(ctx); err != nil {
|
if err := s.ensureWritable(ctx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
// GORM 对零值 uint64 主键会省略 id 列;PG 日志表 id 为 NOT NULL 且无默认值,
|
||||||
|
// 须在落库前显式生成雪花 ID(与 CH 写入路径 analytics repo BatchInsert* 行为一致)。
|
||||||
|
for i := range rows {
|
||||||
|
if rows[i].ID == 0 {
|
||||||
|
rows[i].ID = idgen.NextUint64ID()
|
||||||
|
}
|
||||||
|
}
|
||||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1000,6 +1008,11 @@ func (s *gormLogStore) BatchInsertNodeMetricSnapshots(ctx context.Context, rows
|
|||||||
if err := s.ensureWritable(ctx); err != nil {
|
if err := s.ensureWritable(ctx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
for i := range rows {
|
||||||
|
if rows[i].ID == 0 {
|
||||||
|
rows[i].ID = idgen.NextUint64ID()
|
||||||
|
}
|
||||||
|
}
|
||||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1200,6 +1213,11 @@ func (s *gormLogStore) BatchInsertNodeEdgeHealth(ctx context.Context, rows []ana
|
|||||||
if err := s.ensureWritable(ctx); err != nil {
|
if err := s.ensureWritable(ctx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
for i := range rows {
|
||||||
|
if rows[i].ID == 0 {
|
||||||
|
rows[i].ID = idgen.NextUint64ID()
|
||||||
|
}
|
||||||
|
}
|
||||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1250,6 +1268,11 @@ func (s *gormLogStore) BatchInsertNodeObsFrps(ctx context.Context, rows []analyt
|
|||||||
if err := s.ensureWritable(ctx); err != nil {
|
if err := s.ensureWritable(ctx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
for i := range rows {
|
||||||
|
if rows[i].ID == 0 {
|
||||||
|
rows[i].ID = idgen.NextUint64ID()
|
||||||
|
}
|
||||||
|
}
|
||||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1300,6 +1323,11 @@ func (s *gormLogStore) BatchInsertNodeObsFrpc(ctx context.Context, rows []analyt
|
|||||||
if err := s.ensureWritable(ctx); err != nil {
|
if err := s.ensureWritable(ctx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
for i := range rows {
|
||||||
|
if rows[i].ID == 0 {
|
||||||
|
rows[i].ID = idgen.NextUint64ID()
|
||||||
|
}
|
||||||
|
}
|
||||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1365,6 +1393,11 @@ func (s *userAccessLogGormStore) BatchInsert(ctx context.Context, logs []analyti
|
|||||||
if err := s.ensureWritable(ctx); err != nil {
|
if err := s.ensureWritable(ctx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
for i := range logs {
|
||||||
|
if logs[i].ID == 0 {
|
||||||
|
logs[i].ID = idgen.NextUint64ID()
|
||||||
|
}
|
||||||
|
}
|
||||||
return s.db.WithContext(ctx).CreateInBatches(logs, insertBatchSize).Error
|
return s.db.WithContext(ctx).CreateInBatches(logs, insertBatchSize).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -4,9 +4,11 @@
|
|||||||
package logstore
|
package logstore
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -45,6 +47,70 @@ func TestGormBatchInsertAndCount(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// testLogCaptureWriter 捕获 GORM logger 输出(logger.Writer 需实现 Printf)。
|
||||||
|
type testLogCaptureWriter struct {
|
||||||
|
buf *bytes.Buffer
|
||||||
|
}
|
||||||
|
|
||||||
|
func (w testLogCaptureWriter) Write(p []byte) (int, error) { return w.buf.Write(p) }
|
||||||
|
func (w testLogCaptureWriter) Printf(format string, args ...any) {
|
||||||
|
fmt.Fprintf(w.buf, format, args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestGormBatchInsertFillsZeroIDs 回归测试:PG 日志表 id 为 NOT NULL 且无默认值,GORM 对零值
|
||||||
|
// uint64 主键(视为自增)会省略 id 列,导致 PG 插入报 not-null 违例(SQLSTATE 23502)。
|
||||||
|
// 验证 BatchInsert* 落库前为零 ID 行生成雪花 ID:INSERT 语句必须包含 id 列且回填非零、唯一 ID。
|
||||||
|
func TestGormBatchInsertFillsZeroIDs(t *testing.T) {
|
||||||
|
ResetForTest()
|
||||||
|
SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil })
|
||||||
|
defer ResetForTest()
|
||||||
|
|
||||||
|
var buf bytes.Buffer
|
||||||
|
dsn := fmt.Sprintf("file:logstore-idtest-%d?mode=memory&cache=shared", atomic.AddInt64(&testGormStoreSeq, 1))
|
||||||
|
db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{
|
||||||
|
Logger: logger.New(testLogCaptureWriter{&buf}, logger.Config{LogLevel: logger.Info, Colorful: false}),
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("open sqlite: %v", err)
|
||||||
|
}
|
||||||
|
if err := db.AutoMigrate(&analyticsmodel.NodeAccessLog{}); err != nil {
|
||||||
|
t.Fatalf("automigrate: %v", err)
|
||||||
|
}
|
||||||
|
s := newGormStore(db)
|
||||||
|
|
||||||
|
now := time.Now()
|
||||||
|
rows := []analyticsmodel.NodeAccessLog{
|
||||||
|
{NodeID: "n1", LoggedAt: now, RemoteAddr: "1.1.1.1", StatusCode: 200},
|
||||||
|
{NodeID: "n1", LoggedAt: now.Add(time.Second), RemoteAddr: "2.2.2.2", StatusCode: 500},
|
||||||
|
}
|
||||||
|
if err := s.BatchInsertNodeAccessLogs(context.Background(), rows); err != nil {
|
||||||
|
t.Fatalf("insert with zero ids: %v", err)
|
||||||
|
}
|
||||||
|
if rows[0].ID == 0 || rows[1].ID == 0 {
|
||||||
|
t.Fatalf("zero ids not filled: %+v %+v", rows[0], rows[1])
|
||||||
|
}
|
||||||
|
if rows[0].ID == rows[1].ID {
|
||||||
|
t.Fatalf("ids not unique: %d == %d", rows[0].ID, rows[1].ID)
|
||||||
|
}
|
||||||
|
// 捕获日志含 CREATE TABLE 等其它语句,仅校验 INSERT 语句的列清单(而非 RETURNING 子句,
|
||||||
|
// 后者无论是否省略列都含 id)。
|
||||||
|
var insertStmt string
|
||||||
|
for _, line := range strings.Split(buf.String(), "\n") {
|
||||||
|
if strings.Contains(line, "INSERT INTO") {
|
||||||
|
insertStmt = line
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
start := strings.Index(insertStmt, "(")
|
||||||
|
end := strings.Index(insertStmt, ") VALUES")
|
||||||
|
if start < 0 || end <= start {
|
||||||
|
t.Fatalf("cannot parse insert statement: %s", insertStmt)
|
||||||
|
}
|
||||||
|
if columns := insertStmt[start+1 : end]; !strings.Contains(columns, "`id`") {
|
||||||
|
t.Fatalf("insert SQL omits id column: %s", insertStmt)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// TestGormNodeAccessLogPagination 验证节点访问日志分页与 CH ListNodeAccessLogs 一致(0-based):
|
// TestGormNodeAccessLogPagination 验证节点访问日志分页与 CH ListNodeAccessLogs 一致(0-based):
|
||||||
// Page=1 size=2 → OFFSET 2;Page=0 视为第 0 页;PageSize<=0 时与 CH 一致不分页(返回全部匹配行)。
|
// Page=1 size=2 → OFFSET 2;Page=0 视为第 0 页;PageSize<=0 时与 CH 一致不分页(返回全部匹配行)。
|
||||||
func TestGormNodeAccessLogPagination(t *testing.T) {
|
func TestGormNodeAccessLogPagination(t *testing.T) {
|
||||||
|
|||||||
Reference in New Issue
Block a user