281 lines
7.3 KiB
Go
281 lines
7.3 KiB
Go
package creator
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"database/sql"
|
|
_ "embed"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5/pgconn"
|
|
_ "github.com/jackc/pgx/v5/stdlib"
|
|
)
|
|
|
|
//go:embed migrations/017_creator.sql
|
|
var migration017 string
|
|
|
|
//go:embed migrations/018_collection_checkpoints.sql
|
|
var migration018 string
|
|
|
|
//go:embed migrations/019_metric_plans.sql
|
|
var migration019 string
|
|
|
|
//go:embed migrations/020_competitor_sync_leases.sql
|
|
var migration020 string
|
|
|
|
//go:embed migrations/021_work_sources.sql
|
|
var migration021 string
|
|
|
|
//go:embed migrations/022_collection_lease_tokens.sql
|
|
var migration022 string
|
|
|
|
//go:embed migrations/023_password_references.sql
|
|
var migration023 string
|
|
|
|
//go:embed migrations/024_event_processing_times.sql
|
|
var migration024 string
|
|
|
|
//go:embed migrations/025_competitor_sync_tokens.sql
|
|
var migration025 string
|
|
|
|
//go:embed migrations/026_event_message_text.sql
|
|
var migration026 string
|
|
|
|
//go:embed migrations/027_creator_listener_state.sql
|
|
var migration027 string
|
|
|
|
//go:embed migrations/028_creator_event_message_type.sql
|
|
var migration028 string
|
|
|
|
//go:embed migrations/029_creator_source_sync_lease.sql
|
|
var migration029 string
|
|
|
|
//go:embed migrations/030_creator_event_gateway_time.sql
|
|
var migration030 string
|
|
|
|
//go:embed migrations/031_xhs_raw_payloads.sql
|
|
var migration031 string
|
|
|
|
//go:embed migrations/032_douyin_release_remediation.sql
|
|
var migration032 string
|
|
|
|
//go:embed migrations/033_competitor_tags.sql
|
|
var migration033 string
|
|
|
|
// Creator migrations use their own history table. Phase-a and hub retain a
|
|
// shared schema_migration table for their schemas, but their numeric versions
|
|
// overlap with the creator migration files and must not suppress one another.
|
|
//
|
|
//go:embed migrations/034_competitor_tags_repair.sql
|
|
var migration034 string
|
|
|
|
//go:embed migrations/035_account_deletion.sql
|
|
var migration035 string
|
|
|
|
//go:embed migrations/036_competitor_unique_id.sql
|
|
var migration036 string
|
|
|
|
//go:embed migrations/037_competitor_share_jobs.sql
|
|
var migration037 string
|
|
|
|
type SecretReference struct {
|
|
ID string
|
|
Provider string
|
|
}
|
|
|
|
type SecretBridge interface {
|
|
Store(context.Context, SecretReference, string, string) error
|
|
Delete(context.Context, SecretReference, string) error
|
|
}
|
|
|
|
type Store struct {
|
|
db *sql.DB
|
|
secrets SecretBridge
|
|
// Automatic writes for one execution account stay serialized while event
|
|
// ingestion and unrelated accounts remain concurrent.
|
|
automaticLocks sync.Map // map[string]*sync.Mutex
|
|
}
|
|
|
|
func Open(ctx context.Context, databaseURL string) (*Store, error) {
|
|
db, err := sql.Open("pgx", databaseURL)
|
|
if err != nil {
|
|
return nil, errors.New("open creator database")
|
|
}
|
|
db.SetMaxOpenConns(10)
|
|
db.SetMaxIdleConns(2)
|
|
db.SetConnMaxIdleTime(5 * time.Minute)
|
|
if err := db.PingContext(ctx); err != nil {
|
|
_ = db.Close()
|
|
return nil, errors.New("connect to creator database")
|
|
}
|
|
store := &Store{db: db}
|
|
if err := store.migrate(ctx); err != nil {
|
|
_ = db.Close()
|
|
return nil, err
|
|
}
|
|
return store, nil
|
|
}
|
|
|
|
func (s *Store) Close() error { return s.db.Close() }
|
|
|
|
func (s *Store) Ping(ctx context.Context) error { return s.db.PingContext(ctx) }
|
|
|
|
func (s *Store) SetSecretBridge(bridge SecretBridge) { s.secrets = bridge }
|
|
|
|
func (s *Store) automaticExecutionLock(accountID string) *sync.Mutex {
|
|
lock, _ := s.automaticLocks.LoadOrStore(accountID, &sync.Mutex{})
|
|
return lock.(*sync.Mutex)
|
|
}
|
|
|
|
func (s *Store) acquireAutomaticExecutionLock(ctx context.Context, accountID string) (func(), error) {
|
|
conn, err := s.db.Conn(ctx)
|
|
if err != nil {
|
|
return nil, databaseError(err)
|
|
}
|
|
if _, err := conn.ExecContext(ctx, `SELECT pg_advisory_lock(hashtextextended($1, 0))`, accountID); err != nil {
|
|
_ = conn.Close()
|
|
return nil, databaseError(err)
|
|
}
|
|
return func() { _ = conn.Close() }, nil
|
|
}
|
|
|
|
func (s *Store) migrate(ctx context.Context) error {
|
|
tx, err := s.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return errors.New("begin creator schema migration")
|
|
}
|
|
defer tx.Rollback()
|
|
if _, err := tx.ExecContext(ctx, `SELECT pg_advisory_xact_lock(1542738017)`); err != nil {
|
|
return errors.New("lock creator schema migration")
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `CREATE TABLE IF NOT EXISTS creator_schema_migration (version integer PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`); err != nil {
|
|
return errors.New("create creator schema migration table")
|
|
}
|
|
migrations := []struct {
|
|
version int
|
|
sql string
|
|
}{
|
|
{version: 17, sql: migration017},
|
|
{version: 18, sql: migration018},
|
|
{version: 19, sql: migration019},
|
|
{version: 20, sql: migration020},
|
|
{version: 21, sql: migration021},
|
|
{version: 22, sql: migration022},
|
|
{version: 23, sql: migration023},
|
|
{version: 24, sql: migration024},
|
|
{version: 25, sql: migration025},
|
|
{version: 26, sql: migration026},
|
|
{version: 27, sql: migration027},
|
|
{version: 28, sql: migration028},
|
|
{version: 29, sql: migration029},
|
|
{version: 30, sql: migration030},
|
|
{version: 31, sql: migration031},
|
|
{version: 32, sql: migration032},
|
|
{version: 33, sql: migration033},
|
|
{version: 34, sql: migration034},
|
|
{version: 35, sql: migration035},
|
|
{version: 36, sql: migration036},
|
|
{version: 37, sql: migration037},
|
|
}
|
|
for _, migration := range migrations {
|
|
var applied bool
|
|
if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM creator_schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil {
|
|
return errors.New("read creator schema migration state")
|
|
}
|
|
if applied {
|
|
continue
|
|
}
|
|
if _, err := tx.ExecContext(ctx, migration.sql); err != nil {
|
|
return fmt.Errorf("apply creator schema migration %d: %w", migration.version, err)
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `INSERT INTO creator_schema_migration (version) VALUES ($1)`, migration.version); err != nil {
|
|
return fmt.Errorf("record creator schema migration %d: %w", migration.version, err)
|
|
}
|
|
}
|
|
if err := tx.Commit(); err != nil {
|
|
return errors.New("commit creator schema migration")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func newID(prefix string) string {
|
|
var bytes [12]byte
|
|
if _, err := rand.Read(bytes[:]); err != nil {
|
|
panic(fmt.Sprintf("generate creator id: %v", err))
|
|
}
|
|
return prefix + "-" + hex.EncodeToString(bytes[:])
|
|
}
|
|
|
|
func databaseError(err error) error {
|
|
if err == nil {
|
|
return nil
|
|
}
|
|
var pgErr *pgconn.PgError
|
|
if errors.As(err, &pgErr) {
|
|
switch pgErr.Code {
|
|
case "23505":
|
|
return ErrConflict
|
|
case "23503", "23514", "22P02":
|
|
return ErrInvalid
|
|
}
|
|
}
|
|
return err
|
|
}
|
|
|
|
func rowError(err error) error {
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return ErrNotFound
|
|
}
|
|
return databaseError(err)
|
|
}
|
|
|
|
func jsonText(value any) (string, error) {
|
|
encoded, err := json.Marshal(value)
|
|
if err != nil {
|
|
return "", fmt.Errorf("encode creator json: %w", err)
|
|
}
|
|
return string(encoded), nil
|
|
}
|
|
|
|
func decodeStringList(encoded []byte) ([]string, error) {
|
|
if len(encoded) == 0 {
|
|
return []string{}, nil
|
|
}
|
|
var result []string
|
|
if err := json.Unmarshal(encoded, &result); err != nil {
|
|
return nil, fmt.Errorf("decode creator string list: %w", err)
|
|
}
|
|
if result == nil {
|
|
result = []string{}
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func nullableString(value string) any {
|
|
if value == "" {
|
|
return nil
|
|
}
|
|
return value
|
|
}
|
|
|
|
func nullableTime(value sql.NullTime) *time.Time {
|
|
if !value.Valid {
|
|
return nil
|
|
}
|
|
result := value.Time.UTC()
|
|
return &result
|
|
}
|
|
|
|
func nullableInt64(value sql.NullInt64) *int64 {
|
|
if !value.Valid {
|
|
return nil
|
|
}
|
|
result := value.Int64
|
|
return &result
|
|
}
|