fix: align applied listener schema while preserving data

This commit is contained in:
2026-10-06 15:49:34 +08:00
parent 911bde3454
commit 9573e0cbcf
8 changed files with 218 additions and 5 deletions
+1 -1
View File
@@ -7,7 +7,7 @@ stop_signal = "SIGTERM"
kill_delay = 500
ps1 = []
entrypoint = "./tmp/control-plane"
include_ext = ["go"]
include_ext = ["go", "sql"]
exclude_dir = ["tmp", "web", "docs", ".git", ".codegraph"]
exclude_file = []
delay = 200
+1 -1
View File
@@ -29,7 +29,7 @@
## 红线清单(快速阶段也不能省,现在便宜、以后极贵)
1. 数据模型/表结构:认真设计,建表慎重——改表成本远高于写代码。
1. 数据模型/表结构:认真设计,建表慎重——改表成本远高于写代码。已经在数据库执行过的更新脚本禁止修改或复用编号;表名、字段及约束调整必须使用新编号,并验证新建数据库、已更新数据库和重复启动,保留现有数据。
2. 目录结构与模块边界:保持简单清晰,不堆一坨代码。
3. 基础错误日志:出错时至少能看到发生了什么。
4. Git:小步提交,保持历史清晰。
@@ -0,0 +1,120 @@
package environment
import (
"context"
"database/sql"
"os"
"strings"
"testing"
)
func TestListenerSchemaUpgrade(t *testing.T) {
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
if databaseURL == "" {
t.Skip("set CREATORHUB_POSTGRES_TEST_URL")
}
t.Run("fresh and repeated startup", func(t *testing.T) {
ctx := context.Background()
testURL := isolatedDatabaseURL(t, databaseURL)
store := openFullyMigratedHub(t, ctx, testURL)
store.Close()
db, err := sql.Open("pgx", testURL)
if err != nil {
t.Fatal(err)
}
defer db.Close()
assertListenerSchema(t, db)
store = openFullyMigratedHub(t, ctx, testURL)
store.Close()
assertListenerSchema(t, db)
})
t.Run("already applied old names and uniqueness upgrade without data loss", func(t *testing.T) {
ctx := context.Background()
testURL := isolatedDatabaseURL(t, databaseURL)
store := openFullyMigratedHub(t, ctx, testURL)
store.Close()
db, err := sql.Open("pgx", testURL)
if err != nil {
t.Fatal(err)
}
defer db.Close()
_, err = db.Exec(`
DELETE FROM schema_migration WHERE version = 1052;
INSERT INTO social_account(account_id,name,credential_provider,credential_key,platform,platform_account_key)
VALUES('listener-upgrade-account','保留的账号','os_keyring','creatorhub/listener-upgrade','douyin','listener-upgrade-uid');
INSERT INTO creator_account_listener(account_id,enabled,generation,status,last_delivery_id,dm_synced_at,dm_sync_error)
SELECT id,true,'original-generation','gap','original-delivery',TIMESTAMPTZ '2026-10-01T12:00:00Z','original-error' FROM social_account WHERE account_id='listener-upgrade-account';
INSERT INTO creator_account_event(account_id,event_key,event_type,generation,message_text)
SELECT id,'same-event-key','like','original-generation','保留的事件' FROM social_account WHERE account_id='listener-upgrade-account';
INSERT INTO creator_private_message(id,account_id,peer_uid,server_id,direction,message_type,text,state)
SELECT 'preserved-message',id,'12345','67890','inbound','text','保留的私信','succeeded' FROM social_account WHERE account_id='listener-upgrade-account';
ALTER TABLE creator_account_listener RENAME TO creator_listener_state;
ALTER TABLE creator_account_event DROP CONSTRAINT creator_account_event_account_id_event_type_event_key_key;
ALTER TABLE creator_account_event ADD CONSTRAINT creator_account_event_account_id_event_key_key UNIQUE(account_id,event_key);
`)
if err != nil {
t.Fatal(err)
}
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version IN (1050,1051)`, 2)
store, err = Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
store.Close()
assertListenerSchema(t, db)
assertDatabaseCount(t, db, `SELECT count(*) FROM social_account WHERE account_id='listener-upgrade-account' AND name='保留的账号'`, 1)
assertDatabaseCount(t, db, `SELECT count(*) FROM creator_account_listener WHERE enabled AND generation='original-generation' AND status='gap' AND last_delivery_id='original-delivery' AND dm_synced_at=TIMESTAMPTZ '2026-10-01T12:00:00Z' AND dm_sync_error='original-error'`, 1)
assertDatabaseCount(t, db, `SELECT count(*) FROM creator_account_event WHERE message_text='保留的事件'`, 1)
assertDatabaseCount(t, db, `SELECT count(*) FROM creator_private_message WHERE id='preserved-message' AND text='保留的私信'`, 1)
// Distinct event types may share a key; exact event replays still conflict.
if _, err := db.Exec(`INSERT INTO creator_account_event(account_id,event_key,event_type,generation) SELECT id,'same-event-key','comment','original-generation' FROM social_account WHERE account_id='listener-upgrade-account'`); err != nil {
t.Fatal(err)
}
if _, err := db.Exec(`INSERT INTO creator_account_event(account_id,event_key,event_type,generation) SELECT id,'same-event-key','comment','original-generation' FROM social_account WHERE account_id='listener-upgrade-account'`); err == nil {
t.Fatal("duplicate event accepted")
}
store, err = Open(ctx, testURL)
if err != nil {
t.Fatal(err)
}
store.Close()
assertListenerSchema(t, db)
})
for _, fixture := range []struct{ name, sql string }{
{"both table names fail visibly", `CREATE TABLE creator_listener_state(account_id bigint);`},
{"neither table exists fails visibly", `DROP TABLE creator_account_listener;`},
} {
t.Run(fixture.name, func(t *testing.T) {
ctx := context.Background()
testURL := isolatedDatabaseURL(t, databaseURL)
store := openFullyMigratedHub(t, ctx, testURL)
store.Close()
db, err := sql.Open("pgx", testURL)
if err != nil {
t.Fatal(err)
}
defer db.Close()
if _, err := db.Exec(`DELETE FROM schema_migration WHERE version=1052;` + fixture.sql); err != nil {
t.Fatal(err)
}
store, err = Open(ctx, testURL)
if err == nil {
store.Close()
t.Fatal("invalid listener schema accepted")
}
if !strings.Contains(err.Error(), "apply environment schema migration 1052") {
t.Fatal(err)
}
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version=1052`, 0)
})
}
}
func assertListenerSchema(t *testing.T, db *sql.DB) {
t.Helper()
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version=1052`, 1)
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema=current_schema() AND table_name='creator_account_listener'`, 1)
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema=current_schema() AND table_name='creator_listener_state'`, 0)
assertDatabaseCount(t, db, `SELECT count(*) FROM pg_constraint WHERE conrelid='creator_account_event'::regclass AND conname='creator_account_event_account_id_event_key_key'`, 0)
assertDatabaseCount(t, db, `SELECT count(*) FROM pg_constraint WHERE conrelid='creator_account_event'::regclass AND conname='creator_account_event_account_id_event_type_event_key_key'`, 1)
}
+4 -2
View File
@@ -30,7 +30,8 @@ func TestUnifiedAccountMigration(t *testing.T) {
t.Fatal(err)
}
defer db.Close()
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 56)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version <= 1049`, 56)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version = 1050`, 1)
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema()
AND table_name IN ('social_account', 'browser_env', 'network_exit', 'environment_binding')`, 3)
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_image'`, 0)
@@ -42,7 +43,8 @@ func TestUnifiedAccountMigration(t *testing.T) {
store = openFullyMigratedHub(t, ctx, testURL)
store.Close()
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 56)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version <= 1049`, 56)
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version = 1050`, 1)
})
t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) {
@@ -0,0 +1,31 @@
-- Reception only. Do not restore the removed automatic-operation tables.
CREATE TABLE creator_account_listener (
account_id bigint PRIMARY KEY REFERENCES social_account(id) ON DELETE CASCADE,
enabled boolean NOT NULL DEFAULT false,
generation text NOT NULL,
status text NOT NULL DEFAULT 'stopped' CHECK (status IN ('starting', 'ready', 'gap', 'stopping', 'stopped', 'error')),
boundary_at timestamptz,
last_delivery_id text NOT NULL DEFAULT '',
reason text NOT NULL DEFAULT '',
updated_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE creator_account_event (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
account_id bigint NOT NULL REFERENCES social_account(id) ON DELETE CASCADE,
event_key text NOT NULL,
event_type text NOT NULL CHECK (event_type IN ('like', 'comment', 'follow', 'repost', 'dm')),
interactor_uid text NOT NULL DEFAULT '',
comment_id text NOT NULL DEFAULT '',
work_id text NOT NULL DEFAULT '',
message_type text NOT NULL DEFAULT 'text' CHECK (message_type IN ('text', 'image', 'voice', 'video', 'sticker', 'non_text')),
message_text text NOT NULL DEFAULT '',
platform_event_at timestamptz,
gateway_received_at timestamptz,
received_at timestamptz NOT NULL DEFAULT now(),
generation text NOT NULL,
baseline boolean NOT NULL DEFAULT false,
UNIQUE (account_id, event_type, event_key)
);
CREATE INDEX creator_account_event_received_idx ON creator_account_event (received_at DESC, id DESC);
CREATE INDEX creator_account_event_account_received_idx ON creator_account_event (account_id, received_at DESC, id DESC);
@@ -0,0 +1,21 @@
CREATE TABLE creator_private_message (
id TEXT PRIMARY KEY,
account_id BIGINT NOT NULL REFERENCES social_account(id) ON DELETE CASCADE,
peer_uid TEXT NOT NULL,
peer_name TEXT NOT NULL DEFAULT '',
server_id TEXT,
request_id TEXT,
direction TEXT NOT NULL CHECK (direction IN ('inbound', 'outbound')),
message_type TEXT NOT NULL CHECK (message_type IN ('text', 'image', 'audio', 'video', 'sticker', 'unknown')),
text TEXT NOT NULL DEFAULT '',
state TEXT NOT NULL CHECK (state IN ('sending', 'succeeded', 'failed', 'unknown')),
error TEXT NOT NULL DEFAULT '',
message_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (account_id, server_id),
UNIQUE (account_id, request_id)
);
CREATE INDEX creator_private_message_conversation_idx ON creator_private_message(account_id, peer_uid, created_at DESC, id DESC);
ALTER TABLE creator_account_listener ADD COLUMN dm_synced_at TIMESTAMPTZ;
ALTER TABLE creator_account_listener ADD COLUMN dm_sync_error TEXT NOT NULL DEFAULT '';
@@ -0,0 +1,30 @@
-- A deployed revision of 1050 used creator_listener_state and deduplicated
-- events by (account_id, event_key). Its version is already recorded, so editing
-- 1050 cannot update that database. Align it with the current schema forward,
-- preserving the listener row, private-message sync fields and all event data.
DO $$
BEGIN
IF to_regclass('creator_listener_state') IS NOT NULL THEN
IF to_regclass('creator_account_listener') IS NOT NULL THEN
RAISE EXCEPTION 'both listener tables exist; refusing to merge listener states';
END IF;
ALTER TABLE creator_listener_state RENAME TO creator_account_listener;
ELSIF to_regclass('creator_account_listener') IS NULL THEN
RAISE EXCEPTION 'listener table is missing despite applied migration 1050';
END IF;
END $$;
ALTER TABLE creator_account_event
DROP CONSTRAINT IF EXISTS creator_account_event_account_id_event_key_key;
DO $$
BEGIN
IF NOT EXISTS (
SELECT 1 FROM pg_constraint
WHERE conrelid = 'creator_account_event'::regclass
AND conname = 'creator_account_event_account_id_event_type_event_key_key'
) THEN
ALTER TABLE creator_account_event
ADD CONSTRAINT creator_account_event_account_id_event_type_event_key_key
UNIQUE (account_id, event_type, event_key);
END IF;
END $$;
+10 -1
View File
@@ -188,6 +188,15 @@ var migration1048 string
//go:embed migrations/1049_work_covers_static_files.sql
var migration1049 string
//go:embed migrations/1050_account_event_listener.sql
var migration1050 string
//go:embed migrations/1051_private_messages.sql
var migration1051 string
//go:embed migrations/1052_listener_schema_alignment.sql
var migration1052 string
var (
ErrConflict = errors.New("resource conflicts with existing state")
ErrInvalid = errors.New("invalid environment input")
@@ -329,7 +338,7 @@ func (s *Store) migrate(ctx context.Context) error {
{1029, migration1029}, {1030, migration1030}, {1031, migration1031}, {1032, migration1032}, {1033, migration1033}, {1034, migration1034},
{1035, migration1035}, {1036, migration1036}, {1037, migration1037}, {1038, migration1038}, {1039, migration1039}, {1040, migration1040},
{1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {1045, migration1045}, {1046, migration1046},
{43, migration043}, {44, migration044}, {1047, migration1047}, {1048, migration1048}, {1049, migration1049}} {
{43, migration043}, {44, migration044}, {1047, migration1047}, {1048, migration1048}, {1049, migration1049}, {1050, migration1050}, {1051, migration1051}, {1052, migration1052}} {
var applied bool
if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil {
return errors.New("read environment schema migration state")