Files
creator-hub/internal/creator/competitor_share_jobs.go
T

292 lines
11 KiB
Go

package creator
import (
"context"
"database/sql"
"errors"
"net/url"
"strconv"
"strings"
"time"
"unicode/utf8"
"github.com/jackc/pgx/v5/pgtype"
)
const (
CompetitorShareJobQueued = "queued"
CompetitorShareJobProcessing = "processing"
CompetitorShareJobSucceeded = "succeeded"
CompetitorShareJobFailed = "failed"
MaxCompetitorShareJobAttempts = 3
)
func validateCompetitorShareJobInput(input CompetitorShareJobInput) error {
if !ValidatePlatform(input.Platform) || input.ShareURL == "" || utf8.RuneCountInString(input.ShareURL) > 2000 {
return ErrInvalid
}
parsed, err := url.Parse(input.ShareURL)
if err != nil || parsed.Scheme != "https" || parsed.Hostname() == "" || parsed.User != nil || parsed.Port() != "" || parsed.Fragment != "" {
return ErrInvalid
}
return validateCreatorTags(input.Tags)
}
func normalizeCompetitorShareJobInput(input CompetitorShareJobInput) (CompetitorShareJobInput, error) {
input.Platform = strings.TrimSpace(input.Platform)
input.ShareURL = strings.TrimSpace(input.ShareURL)
if input.Tags == nil {
input.Tags = []string{}
}
if err := validateCompetitorShareJobInput(input); err != nil {
return CompetitorShareJobInput{}, err
}
return input, nil
}
func (s *Store) CreateCompetitorShareJob(ctx context.Context, input CompetitorShareJobInput) (CompetitorShareJob, error) {
input, err := normalizeCompetitorShareJobInput(input)
if err != nil {
return CompetitorShareJob{}, err
}
id := newID("competitor-share-job")
if _, err := s.db.ExecContext(ctx, `
INSERT INTO creator_competitor_share_job (job_id, platform, share_url, tags)
VALUES ($1, $2, $3, $4)`, id, input.Platform, input.ShareURL, input.Tags); err != nil {
return CompetitorShareJob{}, databaseError(err)
}
return s.GetCompetitorShareJob(ctx, id)
}
func scanCompetitorShareJob(scanner interface{ Scan(...any) error }) (CompetitorShareJob, error) {
var result CompetitorShareJob
var tags pgtype.FlatArray[string]
var competitorID sql.NullString
var leaseUntil, lastAttemptAt, completedAt sql.NullTime
if err := scanner.Scan(&result.ID, &result.Platform, &result.ShareURL, pgtype.NewMap().SQLScanner(&tags),
&result.Status, &result.Attempts, &competitorID, &result.FailureReason, &leaseUntil, &lastAttemptAt,
&completedAt, &result.CreatedAt, &result.UpdatedAt); err != nil {
return CompetitorShareJob{}, err
}
result.Tags = []string(tags)
if competitorID.Valid {
result.CompetitorID = competitorID.String
}
result.LastAttemptAt = nullableTime(lastAttemptAt)
result.CompletedAt = nullableTime(completedAt)
return result, nil
}
const competitorShareJobSelect = `SELECT share_job.job_id, share_job.platform, share_job.share_url, share_job.tags, share_job.status, share_job.attempts, competitor.competitor_id,
share_job.failure_reason, share_job.lease_until, share_job.last_attempt_at, share_job.completed_at, share_job.created_at, share_job.updated_at
FROM creator_competitor_share_job share_job
LEFT JOIN creator_competitor competitor ON competitor.id = share_job.competitor_id`
func (s *Store) GetCompetitorShareJob(ctx context.Context, id string) (CompetitorShareJob, error) {
result, err := scanCompetitorShareJob(s.db.QueryRowContext(ctx, competitorShareJobSelect+` WHERE share_job.job_id = $1`, id))
return result, rowError(err)
}
// ListCompetitorShareJobsWithAuthor 附带作者名(join competitor 表),供导入页列表展示。
func (s *Store) ListCompetitorShareJobsWithAuthor(ctx context.Context, platform, status string) ([]CompetitorShareJobView, error) {
if platform != "" && !ValidatePlatform(platform) || status != "" && !validateCompetitorShareJobStatus(status) {
return nil, ErrInvalid
}
query, args := competitorShareJobFilter(`SELECT share_job.job_id, share_job.platform, share_job.share_url, share_job.tags, share_job.status, share_job.attempts, competitor.competitor_id,
share_job.failure_reason, share_job.lease_until, share_job.last_attempt_at, share_job.completed_at, share_job.created_at, share_job.updated_at,
competitor.nickname, competitor.avatar_url
FROM creator_competitor_share_job AS share_job
LEFT JOIN creator_competitor AS competitor ON competitor.id = share_job.competitor_id`, platform, status)
rows, err := s.db.QueryContext(ctx, query+` ORDER BY share_job.created_at DESC, share_job.job_id`, args...)
if err != nil {
return nil, databaseError(err)
}
defer rows.Close()
var items []CompetitorShareJobView
for rows.Next() {
var item CompetitorShareJobView
var tags pgtype.FlatArray[string]
var competitorID, nickname, avatarURL sql.NullString
var leaseUntil, lastAttemptAt, completedAt sql.NullTime
if err := rows.Scan(&item.ID, &item.Platform, &item.ShareURL, pgtype.NewMap().SQLScanner(&tags),
&item.Status, &item.Attempts, &competitorID, &item.FailureReason, &leaseUntil, &lastAttemptAt,
&completedAt, &item.CreatedAt, &item.UpdatedAt, &nickname, &avatarURL); err != nil {
return nil, databaseError(err)
}
item.Tags = []string(tags)
if competitorID.Valid {
item.CompetitorID = competitorID.String
}
item.LastAttemptAt = nullableTime(lastAttemptAt)
item.CompletedAt = nullableTime(completedAt)
item.AuthorName = nickname.String
item.AuthorAvatarURL = avatarURL.String
items = append(items, item)
}
return items, rows.Err()
}
func validateCompetitorShareJobStatus(status string) bool {
switch status {
case CompetitorShareJobQueued, CompetitorShareJobProcessing, CompetitorShareJobSucceeded, CompetitorShareJobFailed:
return true
default:
return false
}
}
func competitorShareJobFilter(query, platform, status string) (string, []any) {
conditions := make([]string, 0, 2)
args := make([]any, 0, 2)
if platform != "" {
args = append(args, platform)
conditions = append(conditions, "share_job.platform = $"+strconv.Itoa(len(args)))
}
if status != "" {
args = append(args, status)
conditions = append(conditions, "share_job.status = $"+strconv.Itoa(len(args)))
}
if len(conditions) > 0 {
query += ` WHERE ` + strings.Join(conditions, ` AND `)
}
return query, args
}
func (s *Store) ListCompetitorShareJobs(ctx context.Context, platform, status string) ([]CompetitorShareJob, error) {
if platform != "" && !ValidatePlatform(platform) || status != "" && !validateCompetitorShareJobStatus(status) {
return nil, ErrInvalid
}
query := competitorShareJobSelect
conditions := make([]string, 0, 2)
args := make([]any, 0, 2)
if platform != "" {
args = append(args, platform)
conditions = append(conditions, "share_job.platform = $"+strconv.Itoa(len(args)))
}
if status != "" {
args = append(args, status)
conditions = append(conditions, "share_job.status = $"+strconv.Itoa(len(args)))
}
if len(conditions) > 0 {
query += ` WHERE ` + strings.Join(conditions, ` AND `)
}
query += ` ORDER BY share_job.created_at DESC, share_job.id`
rows, err := s.db.QueryContext(ctx, query, args...)
if err != nil {
return nil, databaseError(err)
}
defer rows.Close()
result := make([]CompetitorShareJob, 0)
for rows.Next() {
item, err := scanCompetitorShareJob(rows)
if err != nil {
return nil, err
}
result = append(result, item)
}
return result, rows.Err()
}
func (s *Store) ListDueCompetitorShareJobs(ctx context.Context, now time.Time) ([]CompetitorShareJob, error) {
if now.IsZero() {
return nil, ErrInvalid
}
rows, err := s.db.QueryContext(ctx, competitorShareJobSelect+` WHERE
(status = 'queued' AND attempts < $1) OR
(status = 'processing' AND lease_until IS NOT NULL AND lease_until <= $2)
ORDER BY created_at, job_id`, MaxCompetitorShareJobAttempts, now.UTC())
if err != nil {
return nil, databaseError(err)
}
defer rows.Close()
result := make([]CompetitorShareJob, 0)
for rows.Next() {
item, err := scanCompetitorShareJob(rows)
if err != nil {
return nil, err
}
result = append(result, item)
}
return result, rows.Err()
}
func (s *Store) ClaimCompetitorShareJob(ctx context.Context, id string, now time.Time) (CompetitorShareJob, string, bool, error) {
if strings.TrimSpace(id) == "" || now.IsZero() {
return CompetitorShareJob{}, "", false, ErrInvalid
}
token := newID("competitor-share-lease")
claimed, err := scanCompetitorShareJob(s.db.QueryRowContext(ctx, `
UPDATE creator_competitor_share_job
SET status = 'processing',
attempts = CASE WHEN status = 'queued' THEN attempts + 1 ELSE attempts END,
lease_token = $3,
lease_until = $2 + interval '10 minutes',
last_attempt_at = $2,
updated_at = $2
WHERE job_id = $1 AND (
(status = 'queued' AND attempts < $4) OR
(status = 'processing' AND lease_until IS NOT NULL AND lease_until <= $2)
)
RETURNING job_id, platform, share_url, tags, status, attempts,
(SELECT competitor.competitor_id FROM creator_competitor competitor WHERE competitor.id = creator_competitor_share_job.competitor_id),
failure_reason, lease_until, last_attempt_at, completed_at, created_at, updated_at`, id, now.UTC(), token, MaxCompetitorShareJobAttempts))
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return CompetitorShareJob{}, "", false, nil
}
return CompetitorShareJob{}, "", false, databaseError(err)
}
return claimed, token, true, nil
}
func (s *Store) MarkCompetitorShareJob(ctx context.Context, id, leaseToken, status, competitorID, failureReason string) error {
if strings.TrimSpace(id) == "" || strings.TrimSpace(leaseToken) == "" ||
(status != CompetitorShareJobQueued && status != CompetitorShareJobSucceeded && status != CompetitorShareJobFailed) {
return ErrInvalid
}
if status == CompetitorShareJobSucceeded && strings.TrimSpace(competitorID) == "" || status != CompetitorShareJobSucceeded && competitorID != "" {
return ErrInvalid
}
result, err := s.db.ExecContext(ctx, `
UPDATE creator_competitor_share_job
SET status = $3, competitor_id = (SELECT competitor.id FROM creator_competitor competitor WHERE competitor.competitor_id = $4), failure_reason = $5,
lease_token = '', lease_until = NULL,
completed_at = CASE WHEN $3 = 'queued' THEN NULL ELSE now() END,
updated_at = now()
WHERE job_id = $1 AND lease_token = $2 AND status = 'processing'`, id, leaseToken, status, nullableString(competitorID), failureReason)
if err != nil {
return databaseError(err)
}
if affected, err := result.RowsAffected(); err != nil {
return databaseError(err)
} else if affected != 1 {
return ErrConflict
}
return nil
}
func (s *Store) RetryCompetitorShareJob(ctx context.Context, id string) (CompetitorShareJob, error) {
id = strings.TrimSpace(id)
if id == "" {
return CompetitorShareJob{}, ErrInvalid
}
result, err := s.db.ExecContext(ctx, `
UPDATE creator_competitor_share_job
SET status = 'queued', attempts = 0, competitor_id = NULL,
lease_token = '', lease_until = NULL, completed_at = NULL, updated_at = now()
WHERE job_id = $1 AND status = 'failed'`, id)
if err != nil {
return CompetitorShareJob{}, databaseError(err)
}
if affected, err := result.RowsAffected(); err != nil {
return CompetitorShareJob{}, databaseError(err)
} else if affected != 1 {
if _, getErr := s.GetCompetitorShareJob(ctx, id); getErr != nil {
return CompetitorShareJob{}, getErr
}
return CompetitorShareJob{}, ErrConflict
}
return s.GetCompetitorShareJob(ctx, id)
}