279 lines
10 KiB
Go
279 lines
10 KiB
Go
package api
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"net/http"
|
||
"net/url"
|
||
"strconv"
|
||
"strings"
|
||
"time"
|
||
|
||
"git.ipao.vip/rogee/creator-hub/internal/creator"
|
||
"github.com/gofiber/fiber/v3"
|
||
"github.com/sirupsen/logrus"
|
||
)
|
||
|
||
type PrivateSendResult struct {
|
||
State, ServerID, Error string
|
||
MessageAt *time.Time
|
||
}
|
||
type privateMessageSender func(context.Context, creator.PrivateMessageSendInput, creator.PrivateMessageReservation) PrivateSendResult
|
||
|
||
func registerPrivateMessageRoutes(app *fiber.App, store *creator.Store, send privateMessageSender) {
|
||
app.Get("/api/creator/private-messages/conversations", func(c fiber.Ctx) error {
|
||
page, size, _, err := creatorPageQuery(c)
|
||
if err != nil {
|
||
return creatorError(c, err)
|
||
}
|
||
result, err := store.ListPrivateConversations(c.Context(), c.Query("account_id"), page, size)
|
||
if err != nil {
|
||
return creatorError(c, err)
|
||
}
|
||
return c.JSON(result)
|
||
})
|
||
app.Get("/api/creator/private-messages/messages", func(c fiber.Ctx) error {
|
||
page, size, _, err := creatorPageQuery(c)
|
||
if err != nil {
|
||
return creatorError(c, err)
|
||
}
|
||
result, err := store.ListPrivateMessages(c.Context(), c.Query("account_id"), c.Query("peer_uid"), page, size)
|
||
if err != nil {
|
||
return creatorError(c, err)
|
||
}
|
||
return c.JSON(result)
|
||
})
|
||
app.Get("/api/creator/private-messages/status", func(c fiber.Ctx) error {
|
||
result, err := store.ListPrivateSyncStatus(c.Context())
|
||
if err != nil {
|
||
return creatorError(c, err)
|
||
}
|
||
return c.JSON(result)
|
||
})
|
||
app.Post("/api/creator/private-messages/send", func(c fiber.Ctx) error {
|
||
var input creator.PrivateMessageSendInput
|
||
if err := c.Bind().Body(&input); err != nil {
|
||
return creatorError(c, creator.ErrInvalid)
|
||
}
|
||
reservation, err := store.BeginPrivateMessage(c.Context(), input)
|
||
if err != nil {
|
||
return creatorError(c, err)
|
||
}
|
||
if !reservation.New {
|
||
return c.JSON(reservation.Message)
|
||
}
|
||
fields := logrus.Fields{"account_id": input.AccountID, "peer_uid": input.PeerUID, "request_id": input.RequestID, "message_id": reservation.Message.ID, "generation": reservation.Generation}
|
||
logrus.WithFields(fields).Info("private message send reserved")
|
||
// Completion must be saved even when the client closes its connection.
|
||
ctx, cancel := context.WithTimeout(context.WithoutCancel(c.Context()), 65*time.Second)
|
||
defer cancel()
|
||
result := send(ctx, input, reservation)
|
||
message, err := store.FinishPrivateMessage(ctx, reservation.Message.ID, result.State, result.ServerID, result.Error, result.MessageAt)
|
||
if err != nil {
|
||
logrus.WithError(err).WithFields(fields).Error("private message send outcome persistence failed")
|
||
return creatorError(c, err)
|
||
}
|
||
fields["state"] = message.State
|
||
fields["server_id"] = message.ServerID
|
||
if message.State != "succeeded" {
|
||
logrus.WithFields(fields).WithField("reason", message.Error).Error("private message send not confirmed")
|
||
} else {
|
||
logrus.WithFields(fields).Info("private message platform send confirmed")
|
||
}
|
||
return c.JSON(message)
|
||
})
|
||
}
|
||
func decodePrivateSendResult(raw []byte) PrivateSendResult {
|
||
var result struct {
|
||
Status string `json:"status"`
|
||
Code string `json:"code"`
|
||
StatusCode json.RawMessage `json:"status_code"`
|
||
CheckCode string `json:"check_code"`
|
||
CheckMessage string `json:"check_message"`
|
||
Success *bool `json:"success"`
|
||
Message struct {
|
||
ServerID string `json:"server_id"`
|
||
} `json:"message"`
|
||
}
|
||
if err := json.Unmarshal(raw, &result); err != nil {
|
||
return PrivateSendResult{State: "unknown", Error: "网关响应无法解析,发送结果未确认:" + err.Error()}
|
||
}
|
||
reason := result.Code
|
||
if len(result.StatusCode) > 0 && string(result.StatusCode) != "null" {
|
||
reason += ";status_code=" + string(result.StatusCode)
|
||
}
|
||
if result.CheckCode != "" {
|
||
reason += ";check_code=" + result.CheckCode
|
||
}
|
||
if result.CheckMessage != "" {
|
||
reason += ";" + result.CheckMessage
|
||
}
|
||
if strings.Trim(string(result.StatusCode), "\"") == "1008" {
|
||
return PrivateSendResult{State: "unknown", Error: "聊天客户端网络错误,发送结果未确认:" + reason}
|
||
}
|
||
if result.Status == "failed" || result.Status == "denied" || result.Status == "unsupported" {
|
||
return PrivateSendResult{State: "failed", Error: "发送失败:" + reason}
|
||
}
|
||
if result.Status != "succeeded" || result.Success == nil || !*result.Success || result.Message.ServerID == "" || result.Message.ServerID == "0" {
|
||
return PrivateSendResult{State: "unknown", Error: fmt.Sprintf("发送结果未确认:status=%s %s", result.Status, reason)}
|
||
}
|
||
if _, err := strconv.ParseUint(result.Message.ServerID, 10, 64); err != nil {
|
||
return PrivateSendResult{State: "unknown", Error: "平台返回的消息编号无效"}
|
||
}
|
||
return PrivateSendResult{State: "succeeded", ServerID: result.Message.ServerID}
|
||
}
|
||
func privateListenerState(ctx context.Context, store *creator.Store, accountID string) (creator.ListenerState, error) {
|
||
states, err := store.ListPrivateListeners(ctx)
|
||
if err != nil {
|
||
return creator.ListenerState{}, err
|
||
}
|
||
for _, state := range states {
|
||
if state.AccountID == accountID {
|
||
return state, nil
|
||
}
|
||
}
|
||
return creator.ListenerState{}, creator.ErrNotFound
|
||
}
|
||
|
||
func gatewayPrivateMessageSender(accounts HubStore, store *creator.Store) privateMessageSender {
|
||
return func(ctx context.Context, input creator.PrivateMessageSendInput, r creator.PrivateMessageReservation) PrivateSendResult {
|
||
state, err := privateListenerState(ctx, store, input.AccountID)
|
||
if err != nil {
|
||
return PrivateSendResult{State: "failed", Error: err.Error()}
|
||
}
|
||
if !state.Enabled || state.Generation != r.Generation {
|
||
return PrivateSendResult{State: "failed", Error: "监听已关闭或状态已变化,未发送"}
|
||
}
|
||
session, err := listenerSessionForAccount(ctx, store, accounts, state)
|
||
if err != nil {
|
||
return PrivateSendResult{State: "failed", Error: err.Error()}
|
||
}
|
||
target, environment := session.target, session.environment
|
||
payload := gatewayGenerationPayload(environment)
|
||
payload["expected_uid"] = r.UID
|
||
payload["action"] = "dm"
|
||
payload["target_uid"] = input.PeerUID
|
||
payload["text"] = input.Text
|
||
payload["confirm"] = true
|
||
payload["operation_id"] = r.Message.ID
|
||
status, body, err := gatewayCallWithLimit(ctx, target, http.MethodPost, "/v1/browsers/"+url.PathEscape(environment.Alias)+"/douyin/action", payload, 55*time.Second, largeGatewayResponseLimit)
|
||
if err != nil {
|
||
return PrivateSendResult{State: "unknown", Error: "发送请求已交给网关,但结果未确认:" + err.Error()}
|
||
}
|
||
if status != http.StatusOK {
|
||
state := "unknown"
|
||
if status == 400 || status == 409 {
|
||
state = "failed"
|
||
}
|
||
return PrivateSendResult{State: state, Error: fmt.Sprintf("发送网关 HTTP %d:%s", status, body)}
|
||
}
|
||
return decodePrivateSendResult(body)
|
||
}
|
||
}
|
||
|
||
// Polling uses the IM SDK itself: site notification notices are not reliable
|
||
// evidence of individual chat messages. Only enabled listener generations write.
|
||
func RunPrivateMessageSync(ctx context.Context, accounts HubStore, store *creator.Store) {
|
||
if err := store.RecoverInterruptedPrivateMessages(ctx); err != nil {
|
||
logrus.WithError(err).Error("private message interrupted-send recovery failed")
|
||
return
|
||
}
|
||
states := make(map[string]string)
|
||
checkpoints := make(map[string]map[string]string)
|
||
for {
|
||
listeners, err := store.ListPrivateListeners(ctx)
|
||
if err != nil && !errors.Is(err, context.Canceled) {
|
||
logrus.WithError(err).Error("private message enabled-account query failed")
|
||
}
|
||
for _, listener := range listeners {
|
||
if !listener.Enabled {
|
||
continue
|
||
}
|
||
if ctx.Err() != nil {
|
||
return
|
||
}
|
||
key := listener.AccountID + ":" + listener.Generation
|
||
syncCtx, cancel := context.WithTimeout(ctx, 35*time.Second)
|
||
next, err := syncPrivateInbox(syncCtx, accounts, store, listener, checkpoints[key])
|
||
cancel()
|
||
if err != nil {
|
||
if ctx.Err() != nil {
|
||
return
|
||
}
|
||
if errors.Is(err, creator.ErrConflict) {
|
||
continue
|
||
}
|
||
reason := err.Error()
|
||
if states[key] != reason {
|
||
logrus.WithError(err).WithFields(logrus.Fields{"account_id": listener.AccountID, "generation": listener.Generation}).Error("private message inbox sync failed")
|
||
}
|
||
states[key] = reason
|
||
if statusErr := store.UpdatePrivateSyncStatus(ctx, listener.AccountID, listener.Generation, reason); statusErr != nil && !errors.Is(statusErr, creator.ErrConflict) {
|
||
logrus.WithError(statusErr).WithField("account_id", listener.AccountID).Error("private message sync error persistence failed")
|
||
}
|
||
} else {
|
||
if states[key] != "connected" {
|
||
logrus.WithFields(logrus.Fields{"account_id": listener.AccountID, "generation": listener.Generation}).Info("private message inbox sync connected")
|
||
}
|
||
states[key] = "connected"
|
||
checkpoints[key] = next
|
||
}
|
||
}
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-time.After(5 * time.Second):
|
||
}
|
||
}
|
||
}
|
||
func syncPrivateInbox(ctx context.Context, accounts HubStore, store *creator.Store, listener creator.ListenerState, checkpoints map[string]string) (map[string]string, error) {
|
||
profile, err := store.GetAccountProfile(ctx, listener.AccountID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if profile.LoginStatus != "logged_in" {
|
||
return nil, fmt.Errorf("账号未登录")
|
||
}
|
||
session, err := listenerSessionForAccount(ctx, store, accounts, listener)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
target, environment := session.target, session.environment
|
||
payload := gatewayGenerationPayload(environment)
|
||
payload["expected_uid"] = profile.PlatformAccountKey
|
||
if checkpoints != nil {
|
||
payload["checkpoints"] = checkpoints
|
||
}
|
||
status, body, err := gatewayCallWithLimit(ctx, target, http.MethodPost, "/v1/browsers/"+url.PathEscape(environment.Alias)+"/douyin/inbox", payload, 30*time.Second, largeGatewayResponseLimit)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if status != 200 {
|
||
return nil, fmt.Errorf("收件同步网关 HTTP %d:%s", status, body)
|
||
}
|
||
items, err := creator.DecodePrivateInbox(body, profile.PlatformAccountKey)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
var response struct {
|
||
Checkpoints map[string]string `json:"checkpoints"`
|
||
}
|
||
if err := json.Unmarshal(body, &response); err != nil {
|
||
return nil, err
|
||
}
|
||
if response.Checkpoints == nil {
|
||
return nil, fmt.Errorf("收件同步未返回进度")
|
||
}
|
||
if err := store.SavePrivateInbox(ctx, listener.AccountID, listener.Generation, items); err != nil {
|
||
return nil, err
|
||
}
|
||
for _, peer := range items.Peers {
|
||
if peer.Error != "" {
|
||
logrus.WithFields(logrus.Fields{"account_id": listener.AccountID, "generation": listener.Generation, "peer_uid": peer.PeerUID, "reason": peer.Error}).Warn("private message peer nickname lookup failed")
|
||
}
|
||
}
|
||
return response.Checkpoints, nil
|
||
}
|