diff --git a/.air.toml b/.air.toml index 258be26..217fdb2 100644 --- a/.air.toml +++ b/.air.toml @@ -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 diff --git a/AGENTS.md b/AGENTS.md index 6cd3772..c8dbcc1 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -29,7 +29,7 @@ ## 红线清单(快速阶段也不能省,现在便宜、以后极贵) -1. 数据模型/表结构:认真设计,建表慎重——改表成本远高于写代码。 +1. 数据模型/表结构:认真设计,建表慎重——改表成本远高于写代码。已经在数据库执行过的更新脚本禁止修改或复用编号;表名、字段及约束调整必须使用新编号,并验证新建数据库、已更新数据库和重复启动,保留现有数据。 2. 目录结构与模块边界:保持简单清晰,不堆一坨代码。 3. 基础错误日志:出错时至少能看到发生了什么。 4. Git:小步提交,保持历史清晰。 diff --git a/internal/environment/listener_schema_upgrade_test.go b/internal/environment/listener_schema_upgrade_test.go new file mode 100644 index 0000000..fda6d42 --- /dev/null +++ b/internal/environment/listener_schema_upgrade_test.go @@ -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) +} diff --git a/internal/environment/migration_test.go b/internal/environment/migration_test.go index aa4273a..e15458f 100644 --- a/internal/environment/migration_test.go +++ b/internal/environment/migration_test.go @@ -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) { diff --git a/internal/environment/migrations/1050_account_event_listener.sql b/internal/environment/migrations/1050_account_event_listener.sql new file mode 100644 index 0000000..4137e20 --- /dev/null +++ b/internal/environment/migrations/1050_account_event_listener.sql @@ -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); diff --git a/internal/environment/migrations/1051_private_messages.sql b/internal/environment/migrations/1051_private_messages.sql new file mode 100644 index 0000000..baeb72e --- /dev/null +++ b/internal/environment/migrations/1051_private_messages.sql @@ -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 ''; diff --git a/internal/environment/migrations/1052_listener_schema_alignment.sql b/internal/environment/migrations/1052_listener_schema_alignment.sql new file mode 100644 index 0000000..99b790a --- /dev/null +++ b/internal/environment/migrations/1052_listener_schema_alignment.sql @@ -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 $$; diff --git a/internal/environment/store.go b/internal/environment/store.go index 195351f..ed55226 100644 --- a/internal/environment/store.go +++ b/internal/environment/store.go @@ -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")