feat: continue CreatorHub plan01 implementation
This commit is contained in:
@@ -0,0 +1,892 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/creator-hub/internal/creator"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/douyin"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/hub"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/phasea"
|
||||
"github.com/gofiber/fiber/v3"
|
||||
"github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func registerCreator(app *fiber.App, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge) {
|
||||
registerCreatorWithServices(app, store, phaseAStore, hubStore, credentials, nil, nil, nil)
|
||||
}
|
||||
|
||||
func registerCreatorWithServices(app *fiber.App, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge, executor creator.ActionExecutor, generator creator.TextGenerator, analyzer creator.ThemeAnalyzer) {
|
||||
app.Get("/api/creator/settings", func(c fiber.Ctx) error {
|
||||
settings, err := store.GetSettings(c.Context())
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(settings)
|
||||
})
|
||||
app.Put("/api/creator/settings", func(c fiber.Ctx) error {
|
||||
var input creator.SettingsUpdate
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
settings, err := store.UpdateSettings(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(settings)
|
||||
})
|
||||
|
||||
app.Get("/api/creator/accounts", func(c fiber.Ctx) error {
|
||||
profiles, err := store.ListAccountProfiles(c.Context())
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(profiles)
|
||||
})
|
||||
app.Get("/api/creator/accounts/:id/profile", func(c fiber.Ctx) error {
|
||||
profile, err := store.GetAccountProfile(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(profile)
|
||||
})
|
||||
app.Put("/api/creator/accounts/:id/profile", func(c fiber.Ctx) error {
|
||||
var input creator.AccountProfileUpdate
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
profile, err := store.UpdateAccountProfile(c.Context(), c.Params("id"), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(profile)
|
||||
})
|
||||
app.Post("/api/creator/accounts/:id/login-result", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
Status string `json:"status"`
|
||||
Reason string `json:"reason"`
|
||||
ActualKey string `json:"actual_platform_account_key"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
result, err := store.RecordLoginResult(c.Context(), c.Params("id"), input.Status, input.Reason, input.ActualKey)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(result)
|
||||
})
|
||||
app.Post("/api/creator/accounts/:id/big-account", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
Enabled bool `json:"enabled"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
profile, err := store.SetBigAccount(c.Context(), c.Params("id"), input.Enabled)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(profile)
|
||||
})
|
||||
app.Get("/api/creator/relations", func(c fiber.Ctx) error {
|
||||
relations, err := store.ListRelations(c.Context(), c.Query("big_account_id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(relations)
|
||||
})
|
||||
app.Post("/api/creator/relations", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
creator.Relation
|
||||
Enabled bool `json:"enabled"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
if err := store.SetRelation(c.Context(), input.BigAccountID, input.SmallAccountID, input.Enabled); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.SendStatus(fiber.StatusNoContent)
|
||||
})
|
||||
app.Get("/api/creator/accounts/:id/strategies", func(c fiber.Ctx) error {
|
||||
strategies, err := store.ListStrategies(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(strategies)
|
||||
})
|
||||
app.Post("/api/creator/accounts/:id/strategies", func(c fiber.Ctx) error {
|
||||
var input creator.StrategyInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
strategy, err := store.CreateStrategy(c.Context(), c.Params("id"), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.Status(fiber.StatusCreated).JSON(strategy)
|
||||
})
|
||||
app.Put("/api/creator/strategies/:id", func(c fiber.Ctx) error {
|
||||
var input creator.StrategyInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
strategy, err := store.UpdateStrategy(c.Context(), c.Params("id"), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(strategy)
|
||||
})
|
||||
app.Post("/api/creator/strategies/:id/enable", func(c fiber.Ctx) error { return setStrategyEnabled(c, store, true) })
|
||||
app.Post("/api/creator/strategies/:id/disable", func(c fiber.Ctx) error { return setStrategyEnabled(c, store, false) })
|
||||
app.Delete("/api/creator/strategies/:id", func(c fiber.Ctx) error {
|
||||
if err := store.DeleteStrategy(c.Context(), c.Params("id")); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.SendStatus(fiber.StatusNoContent)
|
||||
})
|
||||
|
||||
app.Get("/api/creator/competitors", func(c fiber.Ctx) error {
|
||||
items, err := store.ListCompetitors(c.Context(), c.Query("platform"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/competitors", func(c fiber.Ctx) error {
|
||||
var input creator.CompetitorInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, err := store.CreateCompetitor(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.Status(fiber.StatusCreated).JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/competitors/:id", func(c fiber.Ctx) error {
|
||||
item, err := store.GetCompetitor(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/competitors/:id/pause", func(c fiber.Ctx) error {
|
||||
item, err := store.SetCompetitorEnabled(c.Context(), c.Params("id"), false)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/competitors/:id/resume", func(c fiber.Ctx) error {
|
||||
item, err := store.SetCompetitorEnabled(c.Context(), c.Params("id"), true)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/competitors/:id/sync", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
AccountID string `json:"account_id"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
report, err := syncCreatorCompetitor(c.Context(), store, phaseAStore, hubStore, credentials, c.Params("id"), input.AccountID)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.Status(fiber.StatusAccepted).JSON(report)
|
||||
})
|
||||
|
||||
app.Get("/api/creator/works", func(c fiber.Ctx) error {
|
||||
filter, err := workFilter(c)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
items, err := store.ListWorks(c.Context(), filter)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/works", func(c fiber.Ctx) error {
|
||||
var input creator.WorkInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, inserted, err := store.UpsertWork(c.Context(), input, time.Now().UTC())
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if inserted {
|
||||
status = fiber.StatusCreated
|
||||
}
|
||||
return c.Status(status).JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/works/:id", func(c fiber.Ctx) error {
|
||||
item, err := store.GetWork(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/works/:id/metrics", func(c fiber.Ctx) error {
|
||||
items, err := store.ListMetrics(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/works/:id/metrics", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
CollectedAt time.Time `json:"collected_at"`
|
||||
Likes *int64 `json:"likes"`
|
||||
CommentsCount *int64 `json:"comments_count"`
|
||||
Shares *int64 `json:"shares"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
if input.CollectedAt.IsZero() {
|
||||
input.CollectedAt = time.Now().UTC()
|
||||
}
|
||||
settings, err := store.GetSettings(c.Context())
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
point, err := store.RecordMetric(c.Context(), creator.MetricInput{WorkID: c.Params("id"), CollectedAt: input.CollectedAt, Likes: input.Likes, CommentsCount: input.CommentsCount, Shares: input.Shares}, settings, time.Now().UTC())
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(point)
|
||||
})
|
||||
app.Get("/api/creator/works/:id/material", func(c fiber.Ctx) error {
|
||||
item, err := store.GetMaterial(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/works/:id/material/select", func(c fiber.Ctx) error {
|
||||
item, inserted, err := store.SelectMaterial(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if inserted {
|
||||
status = fiber.StatusCreated
|
||||
}
|
||||
return c.Status(status).JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/works/:id/material/step", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
Step string `json:"step"`
|
||||
Status string `json:"status"`
|
||||
Reference string `json:"reference"`
|
||||
Reason string `json:"reason"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, err := store.SetMaterialStep(c.Context(), c.Params("id"), input.Step, input.Status, input.Reference, input.Reason)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/works/:id/material/rewrite/confirm", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
Requirement string `json:"requirement"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, err := store.ConfirmRewrite(c.Context(), c.Params("id"), input.Requirement)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Put("/api/creator/works/:id/material/rewrite", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
Title string `json:"title"`
|
||||
Script string `json:"script"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, err := store.SaveRewrite(c.Context(), c.Params("id"), input.Title, input.Script)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
|
||||
app.Get("/api/creator/comments", func(c fiber.Ctx) error {
|
||||
items, err := store.ListComments(c.Context(), c.Query("platform"), c.Query("work_id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/comments", func(c fiber.Ctx) error {
|
||||
var input creator.CommentInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, inserted, err := store.SaveComment(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if inserted {
|
||||
status = fiber.StatusCreated
|
||||
}
|
||||
return c.Status(status).JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/comments/:id", func(c fiber.Ctx) error {
|
||||
item, err := store.GetComment(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/rules", func(c fiber.Ctx) error {
|
||||
items, err := store.ListRules(c.Context(), c.Query("enabled_only") == "true")
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/rules", func(c fiber.Ctx) error {
|
||||
var input creator.LeadRuleInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, err := store.CreateRule(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.Status(fiber.StatusCreated).JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/rules/:id", func(c fiber.Ctx) error {
|
||||
item, err := store.GetRule(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Put("/api/creator/rules/:id", func(c fiber.Ctx) error {
|
||||
var input creator.LeadRuleInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, err := store.UpdateRule(c.Context(), c.Params("id"), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/rules/:id/enable", func(c fiber.Ctx) error {
|
||||
item, err := store.SetRuleEnabled(c.Context(), c.Params("id"), true)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/rules/:id/disable", func(c fiber.Ctx) error {
|
||||
item, err := store.SetRuleEnabled(c.Context(), c.Params("id"), false)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/rule-results", func(c fiber.Ctx) error {
|
||||
items, err := store.ListRuleResults(c.Context(), c.Query("comment_id"), c.Query("rule_id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Get("/api/creator/leads", func(c fiber.Ctx) error {
|
||||
items, err := store.ListLeads(c.Context(), c.Query("platform"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/comments/:id/analyze", func(c fiber.Ctx) error {
|
||||
var input struct {
|
||||
RuleID string `json:"rule_id"`
|
||||
}
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
result, err := store.AnalyzeComment(c.Context(), c.Params("id"), input.RuleID, analyzer)
|
||||
if err != nil && !errors.Is(err, creator.ErrUnavailable) {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if errors.Is(err, creator.ErrUnavailable) {
|
||||
status = fiber.StatusServiceUnavailable
|
||||
}
|
||||
return c.Status(status).JSON(result)
|
||||
})
|
||||
|
||||
app.Get("/api/creator/events", func(c fiber.Ctx) error {
|
||||
items, err := store.ListEvents(c.Context(), c.Query("account_id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/events", func(c fiber.Ctx) error {
|
||||
var input creator.InteractionEvent
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
result, err := store.RecordEvent(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if !result.Duplicate {
|
||||
status = fiber.StatusCreated
|
||||
}
|
||||
return c.Status(status).JSON(result)
|
||||
})
|
||||
app.Post("/api/creator/events/process", func(c fiber.Ctx) error {
|
||||
var input creator.InteractionEvent
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
result, err := store.ProcessAutomaticEvent(c.Context(), input, executor, generator)
|
||||
if err != nil && !errors.Is(err, creator.ErrUnavailable) {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if errors.Is(err, creator.ErrUnavailable) {
|
||||
status = fiber.StatusServiceUnavailable
|
||||
}
|
||||
return c.Status(status).JSON(result)
|
||||
})
|
||||
app.Post("/api/creator/events/:id/display", func(c fiber.Ctx) error {
|
||||
event, err := store.SetEventDisplayed(c.Context(), c.Params("id"), time.Now().UTC())
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(event)
|
||||
})
|
||||
app.Get("/api/creator/operations", func(c fiber.Ctx) error {
|
||||
items, err := store.ListOperations(c.Context(), c.Query("account_id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/operations", func(c fiber.Ctx) error {
|
||||
var input creator.OperationInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, inserted, err := store.CreateOperation(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if inserted {
|
||||
status = fiber.StatusCreated
|
||||
}
|
||||
return c.Status(status).JSON(item)
|
||||
})
|
||||
app.Get("/api/creator/operations/:id", func(c fiber.Ctx) error {
|
||||
item, err := store.GetOperation(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
app.Post("/api/creator/operations/:id/execute", func(c fiber.Ctx) error {
|
||||
item, err := store.ExecuteManualOperation(c.Context(), c.Params("id"), executor)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
})
|
||||
|
||||
app.Get("/api/creator/conversations", func(c fiber.Ctx) error {
|
||||
items, err := store.ListConversations(c.Context(), c.Query("account_id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Get("/api/creator/conversations/:id/messages", func(c fiber.Ctx) error {
|
||||
items, err := store.ListMessages(c.Context(), c.Params("id"))
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(items)
|
||||
})
|
||||
app.Post("/api/creator/messages", func(c fiber.Ctx) error {
|
||||
var input creator.MessageInput
|
||||
if err := decodeCreator(c, &input); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
item, inserted, err := store.SaveMessage(c.Context(), input)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
status := fiber.StatusOK
|
||||
if inserted {
|
||||
status = fiber.StatusCreated
|
||||
}
|
||||
return c.Status(status).JSON(item)
|
||||
})
|
||||
}
|
||||
|
||||
func setStrategyEnabled(c fiber.Ctx, store *creator.Store, enabled bool) error {
|
||||
item, err := store.SetStrategyEnabled(c.Context(), c.Params("id"), enabled)
|
||||
if err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.JSON(item)
|
||||
}
|
||||
|
||||
func workFilter(c fiber.Ctx) (creator.WorkFilter, error) {
|
||||
filter := creator.WorkFilter{Platform: c.Query("platform"), SourceID: c.Query("source_id"), SourceType: c.Query("source_type")}
|
||||
for _, field := range []struct {
|
||||
name string
|
||||
target **int64
|
||||
}{{"min_likes", &filter.MinLikes}, {"min_comments", &filter.MinComments}, {"min_shares", &filter.MinShares}} {
|
||||
value := c.Query(field.name)
|
||||
if value == "" {
|
||||
continue
|
||||
}
|
||||
parsed, err := strconv.ParseInt(value, 10, 64)
|
||||
if err != nil {
|
||||
return creator.WorkFilter{}, creator.ErrInvalid
|
||||
}
|
||||
*field.target = &parsed
|
||||
}
|
||||
for _, field := range []struct {
|
||||
name string
|
||||
target **time.Time
|
||||
}{{"published_after", &filter.PublishedAfter}, {"published_before", &filter.PublishedBefore}} {
|
||||
value := c.Query(field.name)
|
||||
if value == "" {
|
||||
continue
|
||||
}
|
||||
parsed, err := time.Parse(time.RFC3339, value)
|
||||
if err != nil {
|
||||
return creator.WorkFilter{}, creator.ErrInvalid
|
||||
}
|
||||
*field.target = &parsed
|
||||
}
|
||||
return filter, nil
|
||||
}
|
||||
|
||||
func decodeCreator(c fiber.Ctx, destination any) error {
|
||||
decoder := json.NewDecoder(bytes.NewReader(c.Body()))
|
||||
decoder.DisallowUnknownFields()
|
||||
if err := decoder.Decode(destination); err != nil {
|
||||
return creator.ErrInvalid
|
||||
}
|
||||
if err := decoder.Decode(&struct{}{}); !errors.Is(err, io.EOF) {
|
||||
return creator.ErrInvalid
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func creatorError(c fiber.Ctx, err error) error {
|
||||
logrus.WithError(err).WithField("service", "control-plane").Error("creator operation failed")
|
||||
status, message := fiber.StatusInternalServerError, "creator operation failed"
|
||||
switch {
|
||||
case errors.Is(err, creator.ErrInvalid):
|
||||
status, message = fiber.StatusBadRequest, creator.ErrInvalid.Error()
|
||||
case errors.Is(err, creator.ErrConflict):
|
||||
status, message = fiber.StatusConflict, creator.ErrConflict.Error()
|
||||
case errors.Is(err, creator.ErrNotFound):
|
||||
status, message = fiber.StatusNotFound, creator.ErrNotFound.Error()
|
||||
case errors.Is(err, creator.ErrUnavailable):
|
||||
status, message = fiber.StatusServiceUnavailable, creator.ErrUnavailable.Error()
|
||||
case errors.Is(err, creator.ErrUncertain):
|
||||
status, message = fiber.StatusConflict, creator.ErrUncertain.Error()
|
||||
}
|
||||
return c.Status(status).JSON(map[string]string{"error": message})
|
||||
}
|
||||
|
||||
type creatorGatewayBrowser struct {
|
||||
gateway hub.Gateway
|
||||
environment hub.EnvironmentContext
|
||||
}
|
||||
|
||||
func (browser creatorGatewayBrowser) SetCookies(ctx context.Context, cookies []douyin.Cookie) error {
|
||||
payload := gatewayGenerationPayload(browser.environment)
|
||||
payload["cookies"] = cookies
|
||||
status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost, "/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/douyin/cookies", payload, 30*time.Second)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if status != http.StatusNoContent && status != http.StatusNotModified {
|
||||
return fmt.Errorf("douyin cookie injection rejected with HTTP %d: %s", status, string(body))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (browser creatorGatewayBrowser) Get(ctx context.Context, target string) (douyin.Response, error) {
|
||||
payload := gatewayGenerationPayload(browser.environment)
|
||||
payload["url"] = target
|
||||
status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost, "/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/douyin/get", payload, 30*time.Second)
|
||||
if err != nil {
|
||||
return douyin.Response{}, err
|
||||
}
|
||||
if status != http.StatusOK {
|
||||
return douyin.Response{}, fmt.Errorf("douyin browser request rejected with HTTP %d: %s", status, string(body))
|
||||
}
|
||||
var response struct {
|
||||
Status int `json:"status"`
|
||||
Body string `json:"body"`
|
||||
Challenge douyin.Challenge `json:"challenge"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &response); err != nil {
|
||||
return douyin.Response{}, fmt.Errorf("decode douyin browser response: %w", err)
|
||||
}
|
||||
return douyin.Response{Status: response.Status, Body: []byte(response.Body), Challenge: response.Challenge}, nil
|
||||
}
|
||||
|
||||
func syncCreatorCompetitor(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge, competitorID, accountID string) (creator.CollectionReport, error) {
|
||||
return syncCreatorCompetitorWithClaim(ctx, store, phaseAStore, hubStore, credentials, competitorID, accountID, true)
|
||||
}
|
||||
|
||||
func syncCreatorCompetitorDue(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge, competitorID, accountID string) (creator.CollectionReport, error) {
|
||||
return syncCreatorCompetitorWithClaim(ctx, store, phaseAStore, hubStore, credentials, competitorID, accountID, false)
|
||||
}
|
||||
|
||||
func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge, competitorID, accountID string, force bool) (creator.CollectionReport, error) {
|
||||
if store == nil || phaseAStore == nil || hubStore == nil || credentials == nil || accountID == "" {
|
||||
return creator.CollectionReport{}, creator.ErrUnavailable
|
||||
}
|
||||
competitor, err := store.GetCompetitor(ctx, competitorID)
|
||||
if err != nil {
|
||||
return creator.CollectionReport{}, err
|
||||
}
|
||||
settings, err := store.GetSettings(ctx)
|
||||
if err != nil {
|
||||
return creator.CollectionReport{}, err
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
claimed, err := store.ClaimCompetitorSync(ctx, competitorID, force, now)
|
||||
if err != nil {
|
||||
return creator.CollectionReport{}, err
|
||||
}
|
||||
if !claimed {
|
||||
return creator.CollectionReport{}, creator.ErrConflict
|
||||
}
|
||||
blocked := func(blockErr error) (creator.CollectionReport, error) {
|
||||
markErr := store.MarkCompetitorSync(ctx, competitorID, "blocked", "", blockErr.Error(), nil)
|
||||
return creator.CollectionReport{}, errors.Join(blockErr, markErr)
|
||||
}
|
||||
if competitor.Platform != creator.PlatformDouyin {
|
||||
return blocked(fmt.Errorf("%w: 小红书采集器尚未完成平台能力验证", creator.ErrUnavailable))
|
||||
}
|
||||
account, err := phaseAStore.GetAccount(ctx, accountID)
|
||||
if err != nil {
|
||||
return blocked(err)
|
||||
}
|
||||
if account.Platform != creator.PlatformDouyin || account.AuthorizationStatus != "authorized" {
|
||||
return blocked(creator.ErrConflict)
|
||||
}
|
||||
if _, err := store.GetAccountProfile(ctx, accountID); err != nil {
|
||||
return blocked(err)
|
||||
}
|
||||
resolver, ok := credentials.(phasea.CredentialResolver)
|
||||
if !ok {
|
||||
return blocked(creator.ErrUnavailable)
|
||||
}
|
||||
rawCredential, err := phaseAStore.ResolveAccountCredential(ctx, accountID, resolver)
|
||||
if err != nil {
|
||||
return blocked(fmt.Errorf("%w: resolve account credential: %v", creator.ErrUnavailable, err))
|
||||
}
|
||||
cookies, err := douyin.ParseCookieHeader(rawCredential)
|
||||
if err != nil {
|
||||
cookies, err = douyin.ParseCookieBundle(rawCredential)
|
||||
}
|
||||
if err != nil {
|
||||
return blocked(fmt.Errorf("%w: invalid account cookie bundle", creator.ErrConflict))
|
||||
}
|
||||
environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID)
|
||||
if err != nil {
|
||||
return blocked(fmt.Errorf("%w: account environment unavailable: %v", creator.ErrUnavailable, err))
|
||||
}
|
||||
if environment.RuntimeID == "" || environment.RuntimeNetworkID == "" || environment.BindingVersion <= 0 {
|
||||
return blocked(fmt.Errorf("%w: account runtime is not running", creator.ErrUnavailable))
|
||||
}
|
||||
gateway, err := hubStore.GetGateway(ctx, environment.Gateway)
|
||||
if err != nil {
|
||||
return blocked(fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err))
|
||||
}
|
||||
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
|
||||
if err := browser.SetCookies(ctx, cookies); err != nil {
|
||||
return blocked(fmt.Errorf("%w: set account cookies: %v", creator.ErrUnavailable, err))
|
||||
}
|
||||
collector := douyin.CreatorCollector{Browser: browser, AccountKey: competitor.PlatformAccountKey, SourceType: creator.SourceCompetitor, SourceID: competitor.ID}
|
||||
if err := collector.VerifyIdentity(ctx, account.PlatformAccountKey); err != nil {
|
||||
return blocked(fmt.Errorf("%w: account identity verification failed: %v", creator.ErrConflict, err))
|
||||
}
|
||||
report, collectErr := store.CollectSource(ctx, competitor.Platform, creator.SourceCompetitor, competitor.ID, collector, now)
|
||||
if collectErr != nil {
|
||||
status, next := "failed", now.Add(time.Duration(settings.NewWorkIntervalSeconds)*time.Second)
|
||||
if errors.Is(collectErr, creator.ErrUnavailable) || errors.Is(collectErr, creator.ErrConflict) {
|
||||
status, next = "blocked", time.Time{}
|
||||
}
|
||||
var nextAt *time.Time
|
||||
if !next.IsZero() {
|
||||
nextAt = &next
|
||||
}
|
||||
markErr := store.MarkCompetitorSync(ctx, competitorID, status, "", collectErr.Error(), nextAt)
|
||||
return report, errors.Join(collectErr, markErr)
|
||||
}
|
||||
next := time.Now().UTC().Add(time.Duration(settings.NewWorkIntervalSeconds) * time.Second)
|
||||
if err := store.MarkCompetitorSync(ctx, competitorID, "idle", "", "", &next); err != nil {
|
||||
return report, err
|
||||
}
|
||||
return report, nil
|
||||
}
|
||||
|
||||
func runCreatorScheduleOnce(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge) error {
|
||||
if store == nil {
|
||||
return creator.ErrUnavailable
|
||||
}
|
||||
settings, err := store.GetSettings(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
competitors, err := store.ListDueCompetitors(ctx, now)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, competitor := range competitors {
|
||||
accountID, err := creatorCollectionAccount(ctx, store, phaseAStore, hubStore, competitor.Platform)
|
||||
if err != nil {
|
||||
_ = store.MarkCompetitorSync(ctx, competitor.ID, "blocked", "", err.Error(), nil)
|
||||
logrus.WithError(err).WithField("competitor_id", competitor.ID).Warn("creator competitor sync blocked")
|
||||
continue
|
||||
}
|
||||
if _, err := syncCreatorCompetitorDue(ctx, store, phaseAStore, hubStore, credentials, competitor.ID, accountID); err != nil {
|
||||
logrus.WithError(err).WithField("competitor_id", competitor.ID).Warn("creator competitor scheduled sync failed")
|
||||
}
|
||||
}
|
||||
ownedAccounts, err := store.ListDueOwnedAccounts(ctx, now, settings.NewWorkIntervalSeconds)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, accountID := range ownedAccounts {
|
||||
if err := syncCreatorOwned(ctx, store, phaseAStore, hubStore, credentials, accountID, now); err != nil {
|
||||
logrus.WithError(err).WithField("account_id", accountID).Warn("creator owned scheduled sync failed")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func syncCreatorOwned(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge, accountID string, now time.Time) error {
|
||||
if store == nil || phaseAStore == nil || hubStore == nil || credentials == nil || accountID == "" {
|
||||
return creator.ErrUnavailable
|
||||
}
|
||||
account, err := phaseAStore.GetAccount(ctx, accountID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if account.Platform != creator.PlatformDouyin || account.AuthorizationStatus != "authorized" {
|
||||
return creator.ErrConflict
|
||||
}
|
||||
if _, err := store.GetAccountProfile(ctx, accountID); err != nil {
|
||||
return err
|
||||
}
|
||||
resolver, ok := credentials.(phasea.CredentialResolver)
|
||||
if !ok {
|
||||
return creator.ErrUnavailable
|
||||
}
|
||||
rawCredential, err := phaseAStore.ResolveAccountCredential(ctx, accountID, resolver)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%w: resolve account credential: %v", creator.ErrUnavailable, err)
|
||||
}
|
||||
cookies, err := douyin.ParseCookieHeader(rawCredential)
|
||||
if err != nil {
|
||||
cookies, err = douyin.ParseCookieBundle(rawCredential)
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("%w: invalid account cookie bundle", creator.ErrConflict)
|
||||
}
|
||||
environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%w: account environment unavailable: %v", creator.ErrUnavailable, err)
|
||||
}
|
||||
if environment.RuntimeID == "" || environment.RuntimeNetworkID == "" || environment.BindingVersion <= 0 {
|
||||
return fmt.Errorf("%w: account runtime is not running", creator.ErrUnavailable)
|
||||
}
|
||||
gateway, err := hubStore.GetGateway(ctx, environment.Gateway)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err)
|
||||
}
|
||||
browser := creatorGatewayBrowser{gateway: gateway, environment: environment}
|
||||
if err := browser.SetCookies(ctx, cookies); err != nil {
|
||||
return fmt.Errorf("%w: set account cookies: %v", creator.ErrUnavailable, err)
|
||||
}
|
||||
collector := douyin.CreatorCollector{Browser: browser, AccountKey: account.PlatformAccountKey, SourceType: creator.SourceOwned, SourceID: account.ID}
|
||||
if err := collector.VerifyIdentity(ctx, account.PlatformAccountKey); err != nil {
|
||||
return fmt.Errorf("%w: account identity verification failed: %v", creator.ErrConflict, err)
|
||||
}
|
||||
_, err = store.CollectSource(ctx, account.Platform, creator.SourceOwned, account.ID, collector, now)
|
||||
return err
|
||||
}
|
||||
|
||||
func creatorCollectionAccount(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, platform string) (string, error) {
|
||||
if store == nil || phaseAStore == nil || hubStore == nil || platform == "" {
|
||||
return "", creator.ErrUnavailable
|
||||
}
|
||||
accounts, err := phaseAStore.ListAccounts(ctx)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, account := range accounts {
|
||||
if account.Platform != platform || account.AuthorizationStatus != "authorized" {
|
||||
continue
|
||||
}
|
||||
if _, err := store.GetAccountProfile(ctx, account.ID); err != nil {
|
||||
continue
|
||||
}
|
||||
if _, err := hubStore.GetEnvironmentContextForAccount(ctx, account.ID); err != nil {
|
||||
continue
|
||||
}
|
||||
return account.ID, nil
|
||||
}
|
||||
return "", fmt.Errorf("%w: no authorized creator collection account", creator.ErrUnavailable)
|
||||
}
|
||||
|
||||
func runCreatorScheduler(ctx context.Context, store *creator.Store, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge) {
|
||||
ticker := time.NewTicker(30 * time.Second)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
if err := runCreatorScheduleOnce(ctx, store, phaseAStore, hubStore, credentials); err != nil {
|
||||
logrus.WithError(err).Error("creator scheduler failed")
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -106,6 +106,32 @@ func (bridge *persistentCredentialBridge) Delete(ctx context.Context, reference
|
||||
return nil
|
||||
}
|
||||
|
||||
func (bridge *persistentCredentialBridge) Resolve(ctx context.Context, reference phasea.CredentialReference, key string) ([]byte, error) {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !validCredentialReference(reference.Provider, key) {
|
||||
return nil, errors.New("invalid credential reference")
|
||||
}
|
||||
payload, err := os.ReadFile(bridge.path(reference.Provider, key))
|
||||
if err != nil {
|
||||
return nil, errors.New("resolve credential")
|
||||
}
|
||||
aead, err := bridge.aead()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(payload) < 1+aead.NonceSize() || payload[0] != credentialFileVersion {
|
||||
return nil, errors.New("invalid credential")
|
||||
}
|
||||
nonceEnd := 1 + aead.NonceSize()
|
||||
value, err := aead.Open(nil, payload[1:nonceEnd], payload[nonceEnd:], credentialAAD(reference.Provider, key))
|
||||
if err != nil {
|
||||
return nil, errors.New("resolve credential")
|
||||
}
|
||||
return value, nil
|
||||
}
|
||||
|
||||
func (bridge *persistentCredentialBridge) aead() (cipher.AEAD, error) {
|
||||
block, err := aes.NewCipher(bridge.key[:])
|
||||
if err != nil {
|
||||
|
||||
@@ -12,7 +12,7 @@ import (
|
||||
"git.ipao.vip/rogee/creator-hub/internal/hub"
|
||||
)
|
||||
|
||||
const testDouyinIdentityURL = "https://www.douyin.com/aweme/v1/web/user/profile/self/"
|
||||
const testDouyinIdentityURL = "https://www.douyin.com/aweme/v1/web/user/profile/self/?aid=6383&device_platform=webapp"
|
||||
|
||||
func TestDouyinGatewayBrowserFencesAccountGeneration(t *testing.T) {
|
||||
requests := 0
|
||||
|
||||
@@ -106,7 +106,7 @@ func gatewayCreatePayload(environment hub.EnvironmentContext, imageRef string, n
|
||||
fingerprint := environment.Fingerprint
|
||||
fingerprint.ProxyServer = ""
|
||||
fingerprint.DisableNonProxiedUDP = false
|
||||
cmd := append(fingerprint.Args(), "about:blank")
|
||||
cmd := append(fingerprint.Args(), "--remote-allow-origins=*", "about:blank")
|
||||
return map[string]any{
|
||||
"alias": environment.Alias,
|
||||
"name": environment.Name,
|
||||
|
||||
@@ -1026,8 +1026,8 @@ func TestCreateBrowserOrchestratesGateway(t *testing.T) {
|
||||
t.Fatalf("platform must force the bound exit: %#v", payload)
|
||||
}
|
||||
cmd := payload["cmd"].([]any)
|
||||
if len(cmd) != 4 || cmd[0] != "--fingerprint=2024" || cmd[1] != "--fingerprint-platform=windows" ||
|
||||
cmd[2] != "--timezone=Asia/Shanghai" || cmd[3] != "about:blank" {
|
||||
if len(cmd) != 5 || cmd[0] != "--fingerprint=2024" || cmd[1] != "--fingerprint-platform=windows" ||
|
||||
cmd[2] != "--timezone=Asia/Shanghai" || cmd[3] != "--remote-allow-origins=*" || cmd[4] != "about:blank" {
|
||||
t.Fatalf("cmd must carry fingerprint args plus start url: %#v", cmd)
|
||||
}
|
||||
if stored := store.envs["account-a"].Fingerprint; stored.ProxyServer != "" || stored.DisableNonProxiedUDP {
|
||||
|
||||
@@ -17,6 +17,7 @@ import (
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/creator-hub/internal/creator"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/hub"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/phasea"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/taskstate"
|
||||
@@ -74,6 +75,12 @@ func newCommand() *cobra.Command {
|
||||
return err
|
||||
}
|
||||
defer hubStore.Close()
|
||||
creatorStore, err := creator.Open(command.Context(), cfg.databaseURL)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer creatorStore.Close()
|
||||
creatorStore.SetSecretBridge(creatorSecretBridge{bridge: credentials})
|
||||
notify := newAttentionNotifier(os.Stderr)
|
||||
phaseAStore.SetTaskNotifier(notify)
|
||||
hubStore.SetTaskNotifier(notify)
|
||||
@@ -84,11 +91,19 @@ func newCommand() *cobra.Command {
|
||||
defer close(heartbeatDone)
|
||||
runtimeLeaseHeartbeat(heartbeatContext, hubStore)
|
||||
}()
|
||||
listenErr := newHandlerWithCredentialBridge(cfg.webDir, cfg.username, cfg.password, phaseAStore, hubStore, credentials).Listen(cfg.listenAddr, fiber.ListenConfig{
|
||||
creatorScheduleContext, stopCreatorScheduler := context.WithCancel(command.Context())
|
||||
creatorScheduleDone := make(chan struct{})
|
||||
go func() {
|
||||
defer close(creatorScheduleDone)
|
||||
runCreatorScheduler(creatorScheduleContext, creatorStore, phaseAStore, hubStore, credentials)
|
||||
}()
|
||||
listenErr := newHandlerWithCreator(cfg.webDir, cfg.username, cfg.password, phaseAStore, hubStore, credentials, creatorStore).Listen(cfg.listenAddr, fiber.ListenConfig{
|
||||
GracefulContext: command.Context(),
|
||||
DisableStartupMessage: true,
|
||||
})
|
||||
stopCreatorScheduler()
|
||||
stopHeartbeat()
|
||||
<-creatorScheduleDone
|
||||
<-heartbeatDone
|
||||
return listenErr
|
||||
},
|
||||
@@ -212,7 +227,29 @@ func newHandlerWithStores(webDirectory, username, password string, phaseAStore *
|
||||
return newHandlerWithCredentialBridge(webDirectory, username, password, phaseAStore, hubStore, nil)
|
||||
}
|
||||
|
||||
type creatorSecretBridge struct {
|
||||
bridge phasea.CredentialBridge
|
||||
}
|
||||
|
||||
func (b creatorSecretBridge) Store(ctx context.Context, reference creator.SecretReference, key, value string) error {
|
||||
if b.bridge == nil {
|
||||
return errors.New("credential bridge is unavailable")
|
||||
}
|
||||
return b.bridge.Store(ctx, phasea.CredentialReference{ID: reference.ID, Provider: reference.Provider}, key, value)
|
||||
}
|
||||
|
||||
func (b creatorSecretBridge) Delete(ctx context.Context, reference creator.SecretReference, key string) error {
|
||||
if b.bridge == nil {
|
||||
return errors.New("credential bridge is unavailable")
|
||||
}
|
||||
return b.bridge.Delete(ctx, phasea.CredentialReference{ID: reference.ID, Provider: reference.Provider}, key)
|
||||
}
|
||||
|
||||
func newHandlerWithCredentialBridge(webDirectory, username, password string, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge) *fiber.App {
|
||||
return newHandlerWithCreator(webDirectory, username, password, phaseAStore, hubStore, credentials, nil)
|
||||
}
|
||||
|
||||
func newHandlerWithCreator(webDirectory, username, password string, phaseAStore *phasea.Store, hubStore *hub.Store, credentials phasea.CredentialBridge, creatorStore *creator.Store) *fiber.App {
|
||||
app := fiber.New(fiber.Config{
|
||||
AppName: "CreatorHub control plane",
|
||||
BodyLimit: 1 << 20,
|
||||
@@ -231,6 +268,9 @@ func newHandlerWithCredentialBridge(webDirectory, username, password string, pha
|
||||
if phaseAStore != nil {
|
||||
registerPhaseA(app, phaseAStore, hubStore, credentials)
|
||||
}
|
||||
if creatorStore != nil {
|
||||
registerCreator(app, creatorStore, phaseAStore, hubStore, credentials)
|
||||
}
|
||||
app.Get("/*", spaHandler(webDirectory))
|
||||
return app
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"regexp"
|
||||
@@ -15,6 +16,7 @@ import (
|
||||
|
||||
"git.ipao.vip/rogee/creator-hub/internal/douyin"
|
||||
"github.com/gofiber/fiber/v3"
|
||||
"github.com/sirupsen/logrus"
|
||||
"golang.org/x/net/websocket"
|
||||
)
|
||||
|
||||
@@ -23,6 +25,7 @@ const (
|
||||
douyinOriginURL = "https://www.douyin.com/"
|
||||
douyinIdentityPath = "/aweme/v1/web/user/profile/self/"
|
||||
douyinWorksPath = "/aweme/v1/web/aweme/post/"
|
||||
douyinCommentsPath = "/aweme/v1/web/comment/list/"
|
||||
douyinResponseLimit = 1 << 20
|
||||
browserControlTimeout = 15 * time.Second
|
||||
)
|
||||
@@ -71,7 +74,11 @@ func (api gateway) setDouyinCookies(c fiber.Ctx) error {
|
||||
if err := api.requireDouyinGeneration(c.Params("id"), input.douyinGenerationRequest); err != nil {
|
||||
return writeError(c, statusFor(err), err)
|
||||
}
|
||||
if api.browser == nil || api.browser.SetCookies(c.Context(), c.Params("id"), input.Cookies) != nil {
|
||||
if api.browser == nil {
|
||||
return writeError(c, http.StatusBadGateway, errors.New("restricted browser operation failed"))
|
||||
}
|
||||
if err := api.browser.SetCookies(c.Context(), c.Params("id"), input.Cookies); err != nil {
|
||||
logrus.WithError(err).WithField("browser_id", c.Params("id")).Error("restricted browser cookie operation failed")
|
||||
return writeError(c, http.StatusBadGateway, errors.New("restricted browser operation failed"))
|
||||
}
|
||||
if err := api.requireDouyinGeneration(c.Params("id"), input.douyinGenerationRequest); err != nil {
|
||||
@@ -99,6 +106,7 @@ func (api gateway) getDouyin(c fiber.Ctx) error {
|
||||
}
|
||||
response, err := api.browser.Get(c.Context(), c.Params("id"), input.URL)
|
||||
if err != nil {
|
||||
logrus.WithError(err).WithField("browser_id", c.Params("id")).Error("restricted browser fetch operation failed")
|
||||
return writeError(c, http.StatusBadGateway, errors.New("restricted browser operation failed"))
|
||||
}
|
||||
if err := api.requireDouyinGeneration(c.Params("id"), input.douyinGenerationRequest); err != nil {
|
||||
@@ -174,10 +182,22 @@ func validDouyinURL(raw string) bool {
|
||||
query := parsed.Query()
|
||||
switch parsed.Path {
|
||||
case douyinIdentityPath:
|
||||
return parsed.RawQuery == ""
|
||||
return len(query) == 2 && len(query["aid"]) == 1 && query.Get("aid") == "6383" &&
|
||||
len(query["device_platform"]) == 1 && query.Get("device_platform") == "webapp"
|
||||
case douyinWorksPath:
|
||||
return len(query) == 3 && len(query["sec_user_id"]) == 1 && douyinAccountKeyPattern.MatchString(query.Get("sec_user_id")) &&
|
||||
len(query["count"]) == 1 && query.Get("count") == "20" && len(query["max_cursor"]) == 1 && query.Get("max_cursor") == "0"
|
||||
if len(query) != 3 || len(query["sec_user_id"]) != 1 || !douyinAccountKeyPattern.MatchString(query.Get("sec_user_id")) ||
|
||||
len(query["count"]) != 1 || query.Get("count") != "20" || len(query["max_cursor"]) != 1 {
|
||||
return false
|
||||
}
|
||||
cursor, err := strconv.ParseInt(query.Get("max_cursor"), 10, 64)
|
||||
return err == nil && cursor >= 0
|
||||
case douyinCommentsPath:
|
||||
if len(query) != 3 || len(query["aweme_id"]) != 1 || !douyinAccountKeyPattern.MatchString(query.Get("aweme_id")) ||
|
||||
len(query["count"]) != 1 || query.Get("count") != "20" || len(query["cursor"]) != 1 {
|
||||
return false
|
||||
}
|
||||
cursor, err := strconv.ParseInt(query.Get("cursor"), 10, 64)
|
||||
return err == nil && cursor >= 0
|
||||
default:
|
||||
return false
|
||||
}
|
||||
@@ -227,11 +247,15 @@ func (browser cdpBrowser) SetCookies(ctx context.Context, alias string, cookies
|
||||
LoaderID string `json:"loaderId"`
|
||||
ErrorText string `json:"errorText"`
|
||||
}
|
||||
if err := cdpCommand(connection, &commandID, "Page.navigate", map[string]string{"url": douyinOriginURL}, &navigation, &events); err != nil ||
|
||||
navigation.ErrorText != "" || navigation.FrameID == "" || navigation.LoaderID == "" {
|
||||
navigationErr := cdpCommand(connection, &commandID, "Page.navigate", map[string]string{"url": douyinOriginURL}, &navigation, &events)
|
||||
logrus.WithFields(logrus.Fields{"frame_id": navigation.FrameID, "loader_id": navigation.LoaderID,
|
||||
"error_text": navigation.ErrorText, "event_count": len(events), "command_error": navigationErr != nil}).Debug("restricted browser navigation response")
|
||||
if navigationErr != nil || navigation.ErrorText != "" || navigation.FrameID == "" || navigation.LoaderID == "" {
|
||||
logrus.WithFields(logrus.Fields{"frame_id": navigation.FrameID, "loader_id": navigation.LoaderID,
|
||||
"error_text": navigation.ErrorText, "event_count": len(events)}).Error("restricted browser navigation response invalid")
|
||||
return errors.New("restricted browser navigation failed")
|
||||
}
|
||||
if err := waitForDouyinPage(ctx, connection, &commandID, navigation.FrameID, navigation.LoaderID, events); err != nil {
|
||||
if err := waitForDouyinPage(ctx, connection, &commandID, navigation.FrameID, events); err != nil {
|
||||
return err
|
||||
}
|
||||
return cdpCommand(connection, &commandID, "Network.setCookies", map[string]any{"cookies": cdpCookies}, nil, nil)
|
||||
@@ -286,6 +310,18 @@ func (browser cdpBrowser) connect(ctx context.Context, alias string) (*websocket
|
||||
if browser.endpoint != nil {
|
||||
base = browser.endpoint(alias)
|
||||
}
|
||||
baseURL, err := url.Parse(base)
|
||||
if err != nil || baseURL.Scheme != "http" || baseURL.Hostname() == "" || baseURL.Port() == "" {
|
||||
return nil, errors.New("restricted browser unavailable")
|
||||
}
|
||||
if net.ParseIP(baseURL.Hostname()) == nil {
|
||||
addresses, lookupErr := net.LookupIP(baseURL.Hostname())
|
||||
if lookupErr != nil || len(addresses) == 0 {
|
||||
return nil, errors.New("restricted browser unavailable")
|
||||
}
|
||||
baseURL.Host = net.JoinHostPort(addresses[0].String(), baseURL.Port())
|
||||
}
|
||||
base = strings.TrimRight(baseURL.String(), "/")
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodGet, base+"/json/list", nil)
|
||||
if err != nil {
|
||||
return nil, errors.New("restricted browser unavailable")
|
||||
@@ -300,6 +336,7 @@ func (browser cdpBrowser) connect(ctx context.Context, alias string) (*websocket
|
||||
client.CheckRedirect = func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }
|
||||
response, err := client.Do(request)
|
||||
if err != nil {
|
||||
logrus.WithError(err).WithField("browser_id", alias).Error("restricted browser discovery failed")
|
||||
return nil, errors.New("restricted browser unavailable")
|
||||
}
|
||||
defer response.Body.Close()
|
||||
@@ -322,7 +359,6 @@ func (browser cdpBrowser) connect(ctx context.Context, alias string) (*websocket
|
||||
if err := decoder.Decode(&trailing); !errors.Is(err, io.EOF) {
|
||||
return nil, errors.New("restricted browser unavailable")
|
||||
}
|
||||
baseURL, _ := url.Parse(base)
|
||||
pageTarget := ""
|
||||
for _, target := range targets {
|
||||
if target.Type != "page" {
|
||||
@@ -352,6 +388,7 @@ func (browser cdpBrowser) connect(ctx context.Context, alias string) (*websocket
|
||||
}
|
||||
connection, err := config.DialContext(ctx)
|
||||
if err != nil {
|
||||
logrus.WithError(err).WithField("browser_id", alias).Error("restricted browser websocket failed")
|
||||
return nil, errors.New("restricted browser unavailable")
|
||||
}
|
||||
deadline := time.Now().Add(browserControlTimeout)
|
||||
@@ -359,6 +396,7 @@ func (browser cdpBrowser) connect(ctx context.Context, alias string) (*websocket
|
||||
deadline = contextDeadline
|
||||
}
|
||||
_ = connection.SetDeadline(deadline)
|
||||
logrus.WithFields(logrus.Fields{"browser_id": alias, "cdp_endpoint": pageTarget}).Debug("restricted browser websocket connected")
|
||||
return connection, nil
|
||||
}
|
||||
|
||||
@@ -373,11 +411,13 @@ type cdpMessage struct {
|
||||
func cdpCommand(connection *websocket.Conn, commandID *int, method string, parameters any, output any, events *[]cdpMessage) error {
|
||||
*commandID = *commandID + 1
|
||||
if err := websocket.JSON.Send(connection, map[string]any{"id": *commandID, "method": method, "params": parameters}); err != nil {
|
||||
logrus.WithError(err).WithField("cdp_method", method).Error("restricted browser command send failed")
|
||||
return errors.New("restricted browser command failed")
|
||||
}
|
||||
for range 128 {
|
||||
var reply cdpMessage
|
||||
if err := websocket.JSON.Receive(connection, &reply); err != nil {
|
||||
logrus.WithError(err).WithField("cdp_method", method).Error("restricted browser command receive failed")
|
||||
return errors.New("restricted browser command failed")
|
||||
}
|
||||
if reply.ID != *commandID {
|
||||
@@ -387,6 +427,7 @@ func cdpCommand(connection *websocket.Conn, commandID *int, method string, param
|
||||
continue
|
||||
}
|
||||
if len(reply.Error) != 0 || len(reply.Result) == 0 || bytes.Equal(bytes.TrimSpace(reply.Result), []byte("null")) {
|
||||
logrus.WithFields(logrus.Fields{"cdp_method": method, "cdp_error": string(reply.Error), "has_result": len(reply.Result) != 0}).Error("restricted browser command returned failure")
|
||||
return errors.New("restricted browser command failed")
|
||||
}
|
||||
if output != nil && json.Unmarshal(reply.Result, output) != nil {
|
||||
@@ -397,17 +438,32 @@ func cdpCommand(connection *websocket.Conn, commandID *int, method string, param
|
||||
return errors.New("restricted browser command failed")
|
||||
}
|
||||
|
||||
func waitForDouyinPage(ctx context.Context, connection *websocket.Conn, commandID *int, frameID, loaderID string, events []cdpMessage) error {
|
||||
func waitForDouyinPage(ctx context.Context, connection *websocket.Conn, commandID *int, frameID string, events []cdpMessage) error {
|
||||
deadline := time.Now().Add(10 * time.Second)
|
||||
if contextDeadline, ok := ctx.Deadline(); ok && contextDeadline.Before(deadline) {
|
||||
deadline = contextDeadline
|
||||
}
|
||||
_ = connection.SetReadDeadline(deadline)
|
||||
for attempts := 0; attempts < 256; attempts++ {
|
||||
bufferedLifecycle := map[string]bool{}
|
||||
for _, event := range events {
|
||||
if event.Method != "Page.lifecycleEvent" {
|
||||
continue
|
||||
}
|
||||
var lifecycle struct {
|
||||
FrameID string `json:"frameId"`
|
||||
LoaderID string `json:"loaderId"`
|
||||
Name string `json:"name"`
|
||||
}
|
||||
if json.Unmarshal(event.Params, &lifecycle) == nil && lifecycle.FrameID == frameID && (lifecycle.Name == "DOMContentLoaded" || lifecycle.Name == "load") {
|
||||
bufferedLifecycle[lifecycle.LoaderID+":"+lifecycle.Name] = true
|
||||
}
|
||||
}
|
||||
for {
|
||||
var event cdpMessage
|
||||
if len(events) != 0 {
|
||||
event, events = events[0], events[1:]
|
||||
} else if websocket.JSON.Receive(connection, &event) != nil {
|
||||
} else if err := websocket.JSON.Receive(connection, &event); err != nil {
|
||||
logrus.WithError(err).Error("restricted browser lifecycle receive failed")
|
||||
return errors.New("restricted browser navigation failed")
|
||||
}
|
||||
if event.Method != "Page.lifecycleEvent" {
|
||||
@@ -418,7 +474,11 @@ func waitForDouyinPage(ctx context.Context, connection *websocket.Conn, commandI
|
||||
LoaderID string `json:"loaderId"`
|
||||
Name string `json:"name"`
|
||||
}
|
||||
if json.Unmarshal(event.Params, &lifecycle) != nil || lifecycle.FrameID != frameID || lifecycle.LoaderID != loaderID || lifecycle.Name != "load" {
|
||||
if err := json.Unmarshal(event.Params, &lifecycle); err != nil {
|
||||
logrus.WithError(err).Debug("restricted browser lifecycle decode failed")
|
||||
continue
|
||||
}
|
||||
if lifecycle.FrameID != frameID || lifecycle.Name != "DOMContentLoaded" && lifecycle.Name != "load" || bufferedLifecycle[lifecycle.LoaderID+":"+lifecycle.Name] {
|
||||
continue
|
||||
}
|
||||
var evaluated struct {
|
||||
@@ -429,11 +489,12 @@ func waitForDouyinPage(ctx context.Context, connection *websocket.Conn, commandI
|
||||
if err := cdpCommand(connection, commandID, "Runtime.evaluate", map[string]any{
|
||||
"expression": "location.origin", "returnByValue": true,
|
||||
}, &evaluated, nil); err != nil || evaluated.Result.Value != douyinOrigin {
|
||||
logrus.WithFields(logrus.Fields{"frame_id": lifecycle.FrameID, "loader_id": lifecycle.LoaderID,
|
||||
"origin": evaluated.Result.Value}).Error("restricted browser navigation origin failed")
|
||||
return errors.New("restricted browser navigation failed")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
return errors.New("restricted browser navigation failed")
|
||||
}
|
||||
|
||||
func detectDouyinChallenge(status int, body string) douyin.Challenge {
|
||||
|
||||
@@ -16,6 +16,8 @@ import (
|
||||
"golang.org/x/net/websocket"
|
||||
)
|
||||
|
||||
const douyinIdentityURL = "https://www.douyin.com/aweme/v1/web/user/profile/self/?aid=6383&device_platform=webapp"
|
||||
|
||||
type fakeRestrictedBrowser struct {
|
||||
cookies []douyin.Cookie
|
||||
urls []string
|
||||
@@ -90,7 +92,7 @@ func TestGatewayRestrictedDouyinContract(t *testing.T) {
|
||||
}
|
||||
response.Body.Close()
|
||||
|
||||
identityURL := "https://www.douyin.com" + douyinIdentityPath
|
||||
identityURL := douyinIdentityURL
|
||||
getBody := `{` + generation + `,"url":"` + identityURL + `"}`
|
||||
response, err = app.Test(authed(http.MethodPost, "/v1/browsers/account-a/douyin/get", strings.NewReader(getBody)))
|
||||
body, _ := io.ReadAll(response.Body)
|
||||
@@ -222,7 +224,7 @@ func TestCDPBrowserUsesOnlyNarrowCommands(t *testing.T) {
|
||||
}})
|
||||
if !onlyOldLoader {
|
||||
afterReply = append(afterReply, map[string]any{"method": "Page.lifecycleEvent", "params": map[string]any{
|
||||
"frameId": "frame-new", "loaderId": "loader-new", "name": "load",
|
||||
"frameId": "frame-new", "loaderId": "loader-final", "name": "load",
|
||||
}})
|
||||
}
|
||||
case "Runtime.evaluate":
|
||||
@@ -269,7 +271,7 @@ func TestCDPBrowserUsesOnlyNarrowCommands(t *testing.T) {
|
||||
}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result, err := browser.Get(context.Background(), "account-a", "https://www.douyin.com"+douyinIdentityPath)
|
||||
result, err := browser.Get(context.Background(), "account-a", douyinIdentityURL)
|
||||
if err != nil || result.Status != 412 || result.Challenge != douyin.ChallengeCaptcha {
|
||||
t.Fatalf("unexpected CDP response: %#v err=%v", result, err)
|
||||
}
|
||||
@@ -297,13 +299,13 @@ func TestCDPBrowserUsesOnlyNarrowCommands(t *testing.T) {
|
||||
}
|
||||
pageOrigin, fetchMode = douyinOrigin, "redirect"
|
||||
mu.Unlock()
|
||||
if _, err := browser.Get(context.Background(), "account-a", "https://www.douyin.com"+douyinIdentityPath); err == nil {
|
||||
if _, err := browser.Get(context.Background(), "account-a", douyinIdentityURL); err == nil {
|
||||
t.Fatal("accepted a redirected fetch")
|
||||
}
|
||||
mu.Lock()
|
||||
fetchMode = "too_large"
|
||||
mu.Unlock()
|
||||
if _, err := browser.Get(context.Background(), "account-a", "https://www.douyin.com"+douyinIdentityPath); err == nil {
|
||||
if _, err := browser.Get(context.Background(), "account-a", douyinIdentityURL); err == nil {
|
||||
t.Fatal("accepted a response at the 1 MiB limit")
|
||||
}
|
||||
mu.Lock()
|
||||
@@ -395,9 +397,9 @@ func TestCDPDiscoveryRequiresOneJSONValue(t *testing.T) {
|
||||
|
||||
func TestDouyinURLContract(t *testing.T) {
|
||||
for target, want := range map[string]bool{
|
||||
"https://www.douyin.com" + douyinIdentityPath: true,
|
||||
douyinIdentityURL: true,
|
||||
"https://www.douyin.com" + douyinWorksPath + "?sec_user_id=sec-a&count=20&max_cursor=0": true,
|
||||
"https://www.douyin.com" + douyinWorksPath + "?sec_user_id=sec-a&count=20&max_cursor=1": false,
|
||||
"https://www.douyin.com" + douyinWorksPath + "?sec_user_id=sec-a&count=20&max_cursor=1": true,
|
||||
"https://www.douyin.com" + douyinWorksPath + "?sec_user_id=sec-a&count=20&max_cursor=0&method=publish": false,
|
||||
"https://www.douyin.com/aweme/v1/web/commit/item/": false,
|
||||
} {
|
||||
|
||||
@@ -919,6 +919,11 @@ func (api gateway) requireGeneration(id string, input generationRequest) (string
|
||||
}
|
||||
version, _ := strconv.ParseInt(labels[bindingVersionLabel], 10, 64)
|
||||
if input.RuntimeID == "" || runtimeID != input.RuntimeID || version != input.BindingVersion || labels[networkIDLabel] != input.NetworkID {
|
||||
logrus.WithFields(logrus.Fields{
|
||||
"browser_id": id, "input_runtime_id": input.RuntimeID, "actual_runtime_id": runtimeID,
|
||||
"input_binding_version": input.BindingVersion, "actual_binding_version": version,
|
||||
"input_network_id": input.NetworkID, "actual_network_id": labels[networkIDLabel],
|
||||
}).Warn("browser generation mismatch")
|
||||
return "", false, errGenerationConflict
|
||||
}
|
||||
return runtimeID, true, nil
|
||||
|
||||
Reference in New Issue
Block a user