1594 lines
65 KiB
Go
1594 lines
65 KiB
Go
package v1
|
|
|
|
import (
|
|
"context"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net/http"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/gin-gonic/gin"
|
|
"github.com/gochat/gochat/internal/middleware"
|
|
"github.com/gochat/gochat/internal/model"
|
|
channelmodel "github.com/gochat/gochat/internal/model/channel"
|
|
"github.com/gochat/gochat/internal/service"
|
|
"github.com/gochat/gochat/internal/worker"
|
|
"github.com/google/uuid"
|
|
"gorm.io/datatypes"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
type ShangwutongConnectorHandler struct {
|
|
db *gorm.DB
|
|
messageSvc *service.MessageService
|
|
worker *worker.WorkerPool
|
|
}
|
|
|
|
type shangwutongContactMetadataRequest struct {
|
|
CID string `json:"cid"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateContactMetadata(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
if strings.TrimSpace(c.Param("source_id")) == "" {
|
|
h.connectorError(c, http.StatusBadRequest, "invalid_source_id", "source_id is required", false)
|
|
return
|
|
}
|
|
var request shangwutongContactMetadataRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_contact_metadata", "cid is required", false)
|
|
return
|
|
}
|
|
request.CID = strings.TrimSpace(request.CID)
|
|
if request.CID == "" || len(request.CID) > 255 {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_contact_metadata", "cid is required", false)
|
|
return
|
|
}
|
|
var updated bool
|
|
err := h.db.WithContext(c.Request.Context()).Transaction(func(tx *gorm.DB) error {
|
|
var contactInbox model.ContactInbox
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("inbox_id = ? AND source_id = ?", inbox.ID, c.Param("source_id")).First(&contactInbox).Error; err != nil {
|
|
return err
|
|
}
|
|
metadata := map[string]any{}
|
|
if len(contactInbox.ChannelMetadata) > 0 {
|
|
if err := json.Unmarshal(contactInbox.ChannelMetadata, &metadata); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if metadata == nil {
|
|
metadata = map[string]any{}
|
|
}
|
|
if metadata["cid"] == request.CID {
|
|
return nil
|
|
}
|
|
metadata["cid"] = request.CID
|
|
encoded, err := json.Marshal(metadata)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Model(&contactInbox).Update("channel_metadata", encoded).Error; err != nil {
|
|
return err
|
|
}
|
|
updated = true
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "contact source not found", false)
|
|
} else {
|
|
var syntaxErr *json.SyntaxError
|
|
if errors.As(err, &syntaxErr) {
|
|
h.connectorError(c, http.StatusInternalServerError, "contact_metadata_invalid", "stored contact metadata is invalid", true)
|
|
} else {
|
|
h.connectorError(c, http.StatusInternalServerError, "contact_metadata_update_failed", "failed to update contact metadata", true)
|
|
}
|
|
}
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"updated": updated})
|
|
}
|
|
|
|
type shangwutongContactOperationStatusRequest struct {
|
|
EventID string `json:"event_id"`
|
|
Operation string `json:"operation"`
|
|
ContactID uint `json:"contact_id"`
|
|
Name string `json:"name"`
|
|
Status string `json:"status"`
|
|
ErrorCode string `json:"error_code"`
|
|
ErrorMessage string `json:"error_message"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateContactOperationStatus(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
sourceID := strings.TrimSpace(c.Param("source_id"))
|
|
if sourceID == "" {
|
|
h.connectorError(c, http.StatusBadRequest, "invalid_source_id", "source_id is required", false)
|
|
return
|
|
}
|
|
var request shangwutongContactOperationStatusRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_contact_operation_status", "contact operation status is invalid", false)
|
|
return
|
|
}
|
|
request.EventID = strings.TrimSpace(request.EventID)
|
|
request.Operation = strings.TrimSpace(request.Operation)
|
|
request.Name = strings.TrimSpace(request.Name)
|
|
request.Status = strings.TrimSpace(request.Status)
|
|
request.ErrorCode = strings.TrimSpace(request.ErrorCode)
|
|
request.ErrorMessage = strings.TrimSpace(request.ErrorMessage)
|
|
if request.EventID == "" || request.Operation != "change_contact_name" || request.ContactID == 0 || !oneOf(request.Status, "succeeded", "failed", "uncertain") {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_contact_operation_status", "contact operation status is invalid", false)
|
|
return
|
|
}
|
|
if request.Name == "" || len(request.Name) > 255 || len(request.ErrorCode) > 128 || len(request.ErrorMessage) > 1024 || request.Status == "failed" && request.ErrorCode == "" {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_contact_operation_status", "contact operation error is invalid", false)
|
|
return
|
|
}
|
|
expectedKey := fmt.Sprintf("swt-contact-operation:%d:%s", inbox.ID, request.EventID)
|
|
if c.GetHeader("Idempotency-Key") != expectedKey {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_idempotency_key", "Idempotency-Key does not match event_id", false)
|
|
return
|
|
}
|
|
var contactInbox model.ContactInbox
|
|
if err := h.db.WithContext(c.Request.Context()).Preload("Contact").Where("inbox_id = ? AND source_id = ?", inbox.ID, sourceID).First(&contactInbox).Error; err != nil || contactInbox.Contact.ID == 0 || contactInbox.Contact.ID != request.ContactID || contactInbox.Contact.AccountID != inbox.AccountID {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "contact source not found", false)
|
|
return
|
|
}
|
|
metadata := map[string]any{}
|
|
if len(contactInbox.ChannelMetadata) > 0 {
|
|
if err := json.Unmarshal(contactInbox.ChannelMetadata, &metadata); err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "contact_metadata_invalid", "stored contact metadata is invalid", true)
|
|
return
|
|
}
|
|
}
|
|
if metadata == nil {
|
|
metadata = map[string]any{}
|
|
}
|
|
existing, ok := metadata["swt_contact_name_operation"].(map[string]any)
|
|
if !ok {
|
|
h.connectorError(c, http.StatusConflict, "unknown_operation", "contact operation is not pending", false)
|
|
return
|
|
}
|
|
storedEventID, eventOK := classificationStateString(existing, "event_id")
|
|
storedOperation, operationOK := classificationStateString(existing, "operation")
|
|
storedName, nameOK := classificationStateString(existing, "name")
|
|
storedSourceID, sourceOK := classificationStateString(existing, "source_id")
|
|
storedAccountID, accountOK := classificationStateUint(existing, "account_id")
|
|
storedInboxID, inboxOK := classificationStateUint(existing, "inbox_id")
|
|
storedContactInboxID, contactInboxOK := classificationStateUint(existing, "contact_inbox_id")
|
|
storedContactID, contactOK := classificationStateUint(existing, "contact_id")
|
|
if !eventOK || !operationOK || !nameOK || !sourceOK || !accountOK || !inboxOK || !contactInboxOK || !contactOK || storedOperation != request.Operation {
|
|
h.connectorError(c, http.StatusConflict, "unknown_operation", "contact operation identity is incomplete", false)
|
|
return
|
|
}
|
|
if storedAccountID != inbox.AccountID || storedInboxID != inbox.ID || storedContactInboxID != contactInbox.ID || storedContactID != request.ContactID || storedSourceID != sourceID {
|
|
h.connectorError(c, http.StatusConflict, "stale_operation", "contact operation targets a different contact", false)
|
|
return
|
|
}
|
|
if storedName != request.Name {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "contact operation changes its target name", false)
|
|
return
|
|
}
|
|
status, _ := existing["status"].(string)
|
|
code, _ := existing["error_code"].(string)
|
|
message, _ := existing["error_message"].(string)
|
|
if storedEventID != request.EventID {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "event_id": storedEventID, "status": status})
|
|
return
|
|
}
|
|
if status != "pending" && status != request.Status && !(status == "uncertain" && code == "connector_delivery_uncertain") {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "contact operation result conflicts with its terminal state", false)
|
|
return
|
|
}
|
|
if status == request.Status && code == request.ErrorCode && message == request.ErrorMessage {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "event_id": request.EventID, "status": status})
|
|
return
|
|
}
|
|
updated, staleEventID, staleStatus, conflict, err := h.persistContactOperationStatus(c.Request.Context(), inbox.ID, sourceID, request)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "contact_metadata_update_failed", "failed to update contact operation status", true)
|
|
return
|
|
}
|
|
if conflict {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "contact operation result conflicts with its terminal state", false)
|
|
return
|
|
}
|
|
if staleEventID != "" {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "event_id": staleEventID, "status": staleStatus})
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"updated": updated, "event_id": request.EventID, "status": request.Status})
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) persistContactOperationStatus(ctx context.Context, inboxID uint, sourceID string, request shangwutongContactOperationStatusRequest) (updated bool, staleEventID, staleStatus string, conflict bool, err error) {
|
|
err = h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
var current model.ContactInbox
|
|
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).Where("inbox_id = ? AND source_id = ?", inboxID, sourceID).First(¤t).Error; err != nil {
|
|
return err
|
|
}
|
|
metadata := map[string]any{}
|
|
if len(current.ChannelMetadata) > 0 {
|
|
if err := json.Unmarshal(current.ChannelMetadata, &metadata); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if metadata == nil {
|
|
metadata = map[string]any{}
|
|
}
|
|
if current.ContactID != request.ContactID {
|
|
conflict = true
|
|
return nil
|
|
}
|
|
existing, ok := metadata["swt_contact_name_operation"].(map[string]any)
|
|
if !ok {
|
|
conflict = true
|
|
return nil
|
|
}
|
|
storedEventID, eventOK := classificationStateString(existing, "event_id")
|
|
storedOperation, operationOK := classificationStateString(existing, "operation")
|
|
storedName, nameOK := classificationStateString(existing, "name")
|
|
storedSourceID, sourceOK := classificationStateString(existing, "source_id")
|
|
storedAccountID, accountOK := classificationStateUint(existing, "account_id")
|
|
storedInboxID, inboxOK := classificationStateUint(existing, "inbox_id")
|
|
storedContactInboxID, contactInboxOK := classificationStateUint(existing, "contact_inbox_id")
|
|
storedContactID, contactOK := classificationStateUint(existing, "contact_id")
|
|
if !eventOK || !operationOK || !nameOK || !sourceOK || !accountOK || !inboxOK || !contactInboxOK || !contactOK || storedOperation != request.Operation || storedAccountID == 0 || storedInboxID != inboxID || storedContactInboxID != current.ID || storedContactID != current.ContactID || storedSourceID != sourceID {
|
|
conflict = true
|
|
return nil
|
|
}
|
|
if storedEventID != request.EventID {
|
|
staleEventID, staleStatus = storedEventID, classificationMapStringValue(existing, "status")
|
|
return nil
|
|
}
|
|
if storedName != request.Name {
|
|
conflict = true
|
|
return nil
|
|
}
|
|
status, statusOK := classificationStateString(existing, "status")
|
|
if !statusOK {
|
|
conflict = true
|
|
return nil
|
|
}
|
|
existingCode := classificationMapStringValue(existing, "error_code")
|
|
existingMessage := classificationMapStringValue(existing, "error_message")
|
|
if status != "pending" {
|
|
if status == request.Status && existingCode == request.ErrorCode && existingMessage == request.ErrorMessage {
|
|
return nil
|
|
}
|
|
if !(status == "uncertain" && existingCode == "connector_delivery_uncertain") {
|
|
conflict = true
|
|
return nil
|
|
}
|
|
}
|
|
next := make(map[string]any, len(existing)+3)
|
|
for key, value := range existing {
|
|
next[key] = value
|
|
}
|
|
next["status"] = request.Status
|
|
next["error_code"] = request.ErrorCode
|
|
next["error_message"] = request.ErrorMessage
|
|
next["updated_at"] = time.Now().UTC()
|
|
metadata["swt_contact_name_operation"] = next
|
|
encoded, err := json.Marshal(metadata)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := tx.WithContext(ctx).Model(¤t).Update("channel_metadata", encoded).Error; err != nil {
|
|
return err
|
|
}
|
|
updated = true
|
|
return nil
|
|
})
|
|
return
|
|
}
|
|
|
|
func classificationMapStringValue(values map[string]any, key string) string {
|
|
value, _ := values[key].(string)
|
|
return strings.TrimSpace(value)
|
|
}
|
|
|
|
func NewShangwutongConnectorHandler(db *gorm.DB, messageSvc *service.MessageService, pools ...*worker.WorkerPool) *ShangwutongConnectorHandler {
|
|
var pool *worker.WorkerPool
|
|
if len(pools) > 0 {
|
|
pool = pools[0]
|
|
}
|
|
return &ShangwutongConnectorHandler{db: db, messageSvc: messageSvc, worker: pool}
|
|
}
|
|
|
|
type shangwutongConnectorInbox struct {
|
|
SchemaVersion int64 `json:"schema_version"`
|
|
AccountID uint `json:"account_id"`
|
|
InboxID uint `json:"inbox_id"`
|
|
InboxIdentifier string `json:"inbox_identifier"`
|
|
Enabled bool `json:"enabled"`
|
|
DesiredPresence string `json:"desired_presence"`
|
|
ConfigVersion int64 `json:"config_version"`
|
|
Credentials shangwutongConnectorCredentials `json:"credentials"`
|
|
UpdatedAt time.Time `json:"updated_at"`
|
|
}
|
|
|
|
type shangwutongConnectorCredentials struct {
|
|
SessionID string `json:"session_id"`
|
|
Username string `json:"username"`
|
|
Password string `json:"password"`
|
|
HMACToken string `json:"hmac_token"`
|
|
WebhookSecret string `json:"webhook_secret"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) ListInboxes(c *gin.Context) {
|
|
limit := 100
|
|
if raw := c.Query("limit"); raw != "" {
|
|
parsed, err := strconv.Atoi(raw)
|
|
if err != nil || parsed <= 0 || parsed > 100 {
|
|
h.connectorError(c, http.StatusBadRequest, "invalid_limit", "limit must be between 1 and 100", false)
|
|
return
|
|
}
|
|
limit = parsed
|
|
}
|
|
cursor, err := decodeShangwutongCursor(c.Query("cursor"))
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusBadRequest, "invalid_cursor", "cursor is invalid", false)
|
|
return
|
|
}
|
|
accountIDs, err := h.grantedAccountIDs(c)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "grant_lookup_failed", "failed to load connector grants", true)
|
|
return
|
|
}
|
|
items := make([]shangwutongConnectorInbox, 0)
|
|
nextCursor := ""
|
|
if len(accountIDs) > 0 {
|
|
var inboxes []model.Inbox
|
|
if err := h.db.WithContext(c.Request.Context()).Where(
|
|
"account_id IN ? AND channel_type = ? AND id > ?", accountIDs, "shangwutong", cursor,
|
|
).Order("id ASC").Limit(limit + 1).Find(&inboxes).Error; err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "config_lookup_failed", "failed to load inbox configurations", true)
|
|
return
|
|
}
|
|
if len(inboxes) > limit {
|
|
nextCursor = encodeShangwutongCursor(inboxes[limit-1].ID)
|
|
inboxes = inboxes[:limit]
|
|
}
|
|
for i := range inboxes {
|
|
item, err := h.inboxItem(c, &inboxes[i])
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "config_lookup_failed", "failed to load inbox configuration", true)
|
|
return
|
|
}
|
|
items = append(items, item)
|
|
}
|
|
}
|
|
c.Header("Cache-Control", "no-store")
|
|
c.JSON(http.StatusOK, gin.H{"data": items, "next_cursor": nextCursor})
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) GetInbox(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
item, err := h.inboxItem(c, inbox)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "config_lookup_failed", "failed to load inbox configuration", true)
|
|
return
|
|
}
|
|
c.Header("Cache-Control", "no-store")
|
|
c.JSON(http.StatusOK, item)
|
|
}
|
|
|
|
type shangwutongStatusRequest struct {
|
|
ConfigVersion int64 `json:"config_version"`
|
|
ActualPresence string `json:"actual_presence"`
|
|
ConnectionStatus string `json:"connection_status"`
|
|
CredentialStatus string `json:"credential_status"`
|
|
LastHeartbeatAt *time.Time `json:"last_heartbeat_at"`
|
|
LastErrorCode *string `json:"last_error_code"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateInboxStatus(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
var request shangwutongStatusRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil || !validShangwutongStatus(request) {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_status", "status payload is invalid", false)
|
|
return
|
|
}
|
|
var config model.ChannelShangwutongConfig
|
|
if err := h.db.WithContext(c.Request.Context()).Where("inbox_id = ?", inbox.ID).First(&config).Error; err != nil {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "inbox not found", false)
|
|
return
|
|
}
|
|
if request.ConfigVersion > config.ConfigVersion {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "future_config_version", "config_version is newer than the inbox configuration", false)
|
|
return
|
|
}
|
|
if request.ConfigVersion < config.ConfigVersion {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "config_version": config.ConfigVersion})
|
|
return
|
|
}
|
|
updates := map[string]any{
|
|
"actual_presence": request.ActualPresence, "connection_status": request.ConnectionStatus,
|
|
"credential_status": request.CredentialStatus, "last_error_code": request.LastErrorCode,
|
|
"status_updated_at": time.Now().UTC(),
|
|
}
|
|
if request.LastHeartbeatAt != nil {
|
|
updates["last_heartbeat_at"] = request.LastHeartbeatAt.UTC()
|
|
}
|
|
if err := h.db.WithContext(c.Request.Context()).Model(&model.ChannelShangwutongConfig{}).Where(
|
|
"inbox_id = ? AND config_version = ?", inbox.ID, request.ConfigVersion,
|
|
).Updates(updates).Error; err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "status_update_failed", "failed to update inbox status", true)
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"updated": true, "config_version": config.ConfigVersion})
|
|
}
|
|
|
|
type shangwutongMessageResultRequest struct {
|
|
ResultVersion int64 `json:"result_version"`
|
|
Status string `json:"status"`
|
|
ExternalID *string `json:"external_id"`
|
|
ExternalIDs []string `json:"external_ids"`
|
|
ErrorCode *string `json:"error_code"`
|
|
ErrorMessage *string `json:"error_message"`
|
|
OccurredAt *time.Time `json:"occurred_at"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateMessageStatus(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
messageID, err := strconv.ParseUint(c.Param("message_id"), 10, 64)
|
|
if err != nil || messageID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "message not found", false)
|
|
return
|
|
}
|
|
var request shangwutongMessageResultRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil || !validShangwutongMessageResult(request) {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_message_result", "message result payload is invalid", false)
|
|
return
|
|
}
|
|
expectedKey := fmt.Sprintf("swt-delivery:%d:%d:%d", inbox.ID, messageID, request.ResultVersion)
|
|
if c.GetHeader("Idempotency-Key") != expectedKey {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_idempotency_key", "Idempotency-Key does not match result_version", false)
|
|
return
|
|
}
|
|
if h.messageSvc == nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "message_service_unavailable", "message status service is unavailable", true)
|
|
return
|
|
}
|
|
result := service.ShangwutongMessageResult{
|
|
ResultVersion: request.ResultVersion, Status: request.Status, ExternalID: request.ExternalID,
|
|
ExternalIDs: request.ExternalIDs,
|
|
ErrorCode: request.ErrorCode, ErrorMessage: request.ErrorMessage, OccurredAt: request.OccurredAt.UTC(),
|
|
}
|
|
message, applied, err := h.messageSvc.ApplyShangwutongMessageResult(
|
|
c.Request.Context(), inbox.AccountID, inbox.ID, uint(messageID), result,
|
|
)
|
|
if err != nil {
|
|
switch {
|
|
case errors.Is(err, gorm.ErrRecordNotFound):
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "message not found", false)
|
|
case errors.Is(err, service.ErrShangwutongMessageResultConflict):
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", err.Error(), false)
|
|
case errors.Is(err, service.ErrShangwutongMessageNotEligible):
|
|
h.connectorError(c, http.StatusForbidden, "message_not_eligible", err.Error(), false)
|
|
default:
|
|
h.connectorError(c, http.StatusInternalServerError, "message_result_update_failed", "failed to update message result", true)
|
|
}
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{
|
|
"updated": applied, "message_id": message.ID, "status": message.Status,
|
|
"result_version": request.ResultVersion,
|
|
})
|
|
}
|
|
|
|
type shangwutongClassificationCacheResponse struct {
|
|
InboxID uint `json:"inbox_id"`
|
|
ConversationKinds json.RawMessage `json:"conversation_kinds"`
|
|
CustomerColorKinds json.RawMessage `json:"customer_color_kinds"`
|
|
SyncStatus string `json:"sync_status"`
|
|
SyncedAt *time.Time `json:"synced_at,omitempty"`
|
|
LastErrorCode *string `json:"last_error_code,omitempty"`
|
|
LastErrorMessage *string `json:"last_error_message,omitempty"`
|
|
}
|
|
|
|
type shangwutongClassificationCallbackRequest struct {
|
|
EventID string `json:"event_id"`
|
|
ConversationKinds []shangwutongConversationKind `json:"conversation_kinds"`
|
|
CustomerColorKinds []shangwutongCustomerColorKind `json:"customer_colors"`
|
|
}
|
|
|
|
type shangwutongConversationKind struct {
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
IconIndex int `json:"icon_index"`
|
|
}
|
|
|
|
type shangwutongCustomerColorKind struct {
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) ListClassifications(c *gin.Context) {
|
|
inbox, ok := h.userInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
var cache model.ShangwutongClassificationCache
|
|
err := h.db.WithContext(c.Request.Context()).Where("inbox_id = ?", inbox.ID).First(&cache).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
c.JSON(http.StatusOK, shangwutongClassificationCacheResponse{
|
|
InboxID: inbox.ID, ConversationKinds: json.RawMessage(`[]`), CustomerColorKinds: json.RawMessage(`[]`), SyncStatus: "never",
|
|
})
|
|
return
|
|
}
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "classification_cache_lookup_failed", "failed to load classifications", true)
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, classificationCacheResponse(&cache))
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) SyncClassifications(c *gin.Context) {
|
|
inbox, ok := h.userInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
if h.worker == nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "classification_sync_unavailable", "classification sync is unavailable", true)
|
|
return
|
|
}
|
|
eventID := uuid.NewString()
|
|
job, created, err := h.enqueueClassificationSync(c.Request.Context(), inbox, eventID)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "classification_sync_queue_failed", "failed to persist and queue classification sync", true)
|
|
return
|
|
}
|
|
if created {
|
|
h.worker.Publish(c.Request.Context(), job)
|
|
}
|
|
c.JSON(http.StatusAccepted, gin.H{"sync_id": eventID, "sync_status": "pending"})
|
|
}
|
|
|
|
type shangwutongConversationClassificationRequest struct {
|
|
ChatKindID string `json:"chat_kind_id"`
|
|
CustomerColorID string `json:"customer_color_id"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateConversationClassification(c *gin.Context) {
|
|
accountID, err := strconv.ParseUint(c.Param("account_id"), 10, 64)
|
|
if err != nil || accountID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "conversation not found", false)
|
|
return
|
|
}
|
|
conversationID, err := strconv.ParseUint(c.Param("conversation_id"), 10, 64)
|
|
if err != nil || conversationID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "conversation not found", false)
|
|
return
|
|
}
|
|
var request shangwutongConversationClassificationRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_change", "classification change payload is invalid", false)
|
|
return
|
|
}
|
|
request.ChatKindID, request.CustomerColorID = strings.TrimSpace(request.ChatKindID), strings.TrimSpace(request.CustomerColorID)
|
|
if (request.ChatKindID == "") == (request.CustomerColorID == "") {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_change", "exactly one classification is required", false)
|
|
return
|
|
}
|
|
var conversation model.Conversation
|
|
lookup := h.db.WithContext(c.Request.Context()).Where(
|
|
"id = ? AND account_id = ? AND (display_id IS NULL OR display_id = 0)", uint(conversationID), uint(accountID),
|
|
).First(&conversation)
|
|
if errors.Is(lookup.Error, gorm.ErrRecordNotFound) {
|
|
lookup = h.db.WithContext(c.Request.Context()).Where(
|
|
"display_id = ? AND account_id = ?", uint(conversationID), uint(accountID),
|
|
).First(&conversation)
|
|
}
|
|
if lookup.Error != nil {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "conversation not found", false)
|
|
return
|
|
}
|
|
var inbox model.Inbox
|
|
if err := h.db.WithContext(c.Request.Context()).Where("id = ? AND account_id = ? AND channel_type = ?", conversation.InboxID, uint(accountID), "shangwutong").First(&inbox).Error; err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "not_shangwutong_conversation", "conversation is not a Shangwutong conversation", false)
|
|
return
|
|
}
|
|
var contactInbox model.ContactInbox
|
|
if request.CustomerColorID != "" && (conversation.ContactInboxID == nil || *conversation.ContactInboxID == 0) {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "conversation_target_unavailable", "conversation has no Shangwutong session", false)
|
|
return
|
|
}
|
|
contactInboxQuery := h.db.WithContext(c.Request.Context()).Where("contact_id = ? AND inbox_id = ?", conversation.ContactID, inbox.ID)
|
|
if conversation.ContactInboxID != nil && *conversation.ContactInboxID != 0 {
|
|
contactInboxQuery = h.db.WithContext(c.Request.Context()).Where(
|
|
"id = ? AND contact_id = ? AND inbox_id = ?", *conversation.ContactInboxID, conversation.ContactID, inbox.ID,
|
|
)
|
|
}
|
|
if err := contactInboxQuery.First(&contactInbox).Error; err != nil || strings.TrimSpace(contactInbox.SourceID) == "" {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "conversation_target_unavailable", "conversation has no Shangwutong session", false)
|
|
return
|
|
}
|
|
var cache model.ShangwutongClassificationCache
|
|
if err := h.db.WithContext(c.Request.Context()).Where("inbox_id = ?", inbox.ID).First(&cache).Error; err != nil {
|
|
h.connectorError(c, http.StatusConflict, "classification_not_synced", "sync Shangwutong classifications first", false)
|
|
return
|
|
}
|
|
if cache.SyncStatus != "succeeded" {
|
|
h.connectorError(c, http.StatusConflict, "classification_not_synced", "sync Shangwutong classifications first", false)
|
|
return
|
|
}
|
|
var conversationKinds []shangwutongConversationKind
|
|
var customerColors []shangwutongCustomerColorKind
|
|
if err := json.Unmarshal(cache.ConversationKinds, &conversationKinds); err != nil || json.Unmarshal(cache.CustomerColorKinds, &customerColors) != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "classification_cache_invalid", "stored classifications are invalid", true)
|
|
return
|
|
}
|
|
colorName := ""
|
|
if request.ChatKindID != "" {
|
|
if !containsConversationKind(conversationKinds, request.ChatKindID) {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "classification_not_found", "conversation classification is not in the synced catalog", false)
|
|
return
|
|
}
|
|
} else {
|
|
for _, color := range customerColors {
|
|
if color.ID == request.CustomerColorID {
|
|
colorName = color.Name
|
|
break
|
|
}
|
|
}
|
|
if colorName == "" {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "classification_not_found", "customer classification is not in the synced catalog", false)
|
|
return
|
|
}
|
|
var metadata struct {
|
|
CID string `json:"cid"`
|
|
}
|
|
if err := json.Unmarshal(contactInbox.ChannelMetadata, &metadata); err != nil || strings.TrimSpace(metadata.CID) == "" {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "conversation_target_unavailable", "conversation has no customer cid", false)
|
|
return
|
|
}
|
|
request.CustomerColorID = strings.TrimSpace(request.CustomerColorID)
|
|
if h.worker == nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "classification_queue_unavailable", "classification queue is unavailable", true)
|
|
return
|
|
}
|
|
eventID := uuid.NewString()
|
|
job, created, err := h.enqueueClassificationChange(c.Request.Context(), &inbox, &conversation, contactInbox.ID, contactInbox.SourceID, metadata.CID, "", request.CustomerColorID, colorName, eventID)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "classification_queue_failed", "failed to queue classification change", true)
|
|
return
|
|
}
|
|
if created {
|
|
h.worker.Publish(c.Request.Context(), job)
|
|
}
|
|
c.JSON(http.StatusAccepted, gin.H{"sync_id": eventID, "status": "pending", "customer_color_id": request.CustomerColorID})
|
|
return
|
|
}
|
|
if h.worker == nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "classification_queue_unavailable", "classification queue is unavailable", true)
|
|
return
|
|
}
|
|
eventID := uuid.NewString()
|
|
job, created, err := h.enqueueClassificationChange(c.Request.Context(), &inbox, &conversation, contactInbox.ID, contactInbox.SourceID, "", request.ChatKindID, "", "", eventID)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusServiceUnavailable, "classification_queue_failed", "failed to queue classification change", true)
|
|
return
|
|
}
|
|
if created {
|
|
h.worker.Publish(c.Request.Context(), job)
|
|
}
|
|
c.JSON(http.StatusAccepted, gin.H{"sync_id": eventID, "status": "pending", "chat_kind_id": request.ChatKindID})
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) enqueueClassificationSync(ctx context.Context, inbox *model.Inbox, eventID string) (*model.BackgroundJob, bool, error) {
|
|
var job *model.BackgroundJob
|
|
var created bool
|
|
err := h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
var cache model.ShangwutongClassificationCache
|
|
err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).Where("inbox_id = ?", inbox.ID).First(&cache).Error
|
|
switch {
|
|
case errors.Is(err, gorm.ErrRecordNotFound):
|
|
cache = model.ShangwutongClassificationCache{
|
|
InboxID: inbox.ID, ConversationKinds: datatypes.JSON([]byte(`[]`)), CustomerColorKinds: datatypes.JSON([]byte(`[]`)),
|
|
SyncStatus: "pending", LastSyncEventID: eventID,
|
|
}
|
|
if err := tx.WithContext(ctx).Clauses(clause.OnConflict{
|
|
Columns: []clause.Column{{Name: "inbox_id"}},
|
|
DoUpdates: clause.Assignments(map[string]any{
|
|
"sync_status": "pending", "last_sync_event_id": eventID,
|
|
"last_error_code": nil, "last_error_message": nil, "synced_at": nil,
|
|
}),
|
|
}).Create(&cache).Error; err != nil {
|
|
return err
|
|
}
|
|
case err != nil:
|
|
return err
|
|
default:
|
|
if err := tx.WithContext(ctx).Model(&cache).Updates(map[string]any{
|
|
"sync_status": "pending", "last_sync_event_id": eventID, "last_error_code": nil, "last_error_message": nil, "synced_at": nil,
|
|
}).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
job, created, err = service.EnqueueShangwutongClassificationSyncInTransaction(ctx, h.worker, tx, inbox, eventID)
|
|
return err
|
|
})
|
|
return job, created, err
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) enqueueClassificationChange(ctx context.Context, inbox *model.Inbox, conversation *model.Conversation, contactInboxID uint, sid, cid, chatKindID, customerColorID, customerColorName, eventID string) (*model.BackgroundJob, bool, error) {
|
|
var job *model.BackgroundJob
|
|
var created bool
|
|
err := h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
if err := h.markClassificationStateInTransaction(ctx, tx, conversation, contactInboxID, sid, cid, classificationOperation(chatKindID, customerColorID), eventID, classificationValue(chatKindID, customerColorID), "pending", "", ""); err != nil {
|
|
return err
|
|
}
|
|
var err error
|
|
job, created, err = service.EnqueueShangwutongClassificationChangeInTransaction(ctx, h.worker, tx, inbox, conversation.ID, sid, cid, chatKindID, customerColorID, customerColorName, eventID)
|
|
return err
|
|
})
|
|
return job, created, err
|
|
}
|
|
|
|
func classificationOperation(chatKindID, customerColorID string) string {
|
|
if strings.TrimSpace(chatKindID) != "" {
|
|
return "set_chat_kind"
|
|
}
|
|
return "set_customer_color"
|
|
}
|
|
|
|
func classificationValue(chatKindID, customerColorID string) string {
|
|
if strings.TrimSpace(chatKindID) != "" {
|
|
return strings.TrimSpace(chatKindID)
|
|
}
|
|
return strings.TrimSpace(customerColorID)
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) markClassificationStateInTransaction(ctx context.Context, tx *gorm.DB, conversation *model.Conversation, contactInboxID uint, sid, cid, operation, eventID, value, status, errorCode, errorMessage string) error {
|
|
var current model.Conversation
|
|
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", conversation.ID).First(¤t).Error; err != nil {
|
|
return err
|
|
}
|
|
attributes := map[string]any{}
|
|
if len(current.AdditionalAttributes) > 0 {
|
|
if err := json.Unmarshal(current.AdditionalAttributes, &attributes); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if attributes == nil {
|
|
attributes = map[string]any{}
|
|
}
|
|
operations, _ := attributes["swt_classification_operations"].(map[string]any)
|
|
if operations == nil {
|
|
operations = map[string]any{}
|
|
}
|
|
if existing, ok := operations[operation].(map[string]any); ok {
|
|
if existingEventID, _ := existing["event_id"].(string); existingEventID == eventID && existing["status"] == status && existing["value"] == value {
|
|
return nil
|
|
}
|
|
}
|
|
state := map[string]any{
|
|
"event_id": eventID, "operation": operation, "status": status, "value": value,
|
|
"account_id": current.AccountID, "inbox_id": current.InboxID, "conversation_id": current.ID,
|
|
"contact_id": current.ContactID, "contact_inbox_id": contactInboxID,
|
|
"swt_session_id": strings.TrimSpace(sid), "cid": strings.TrimSpace(cid),
|
|
"error_code": errorCode, "error_message": errorMessage,
|
|
"updated_at": time.Now().UTC(),
|
|
}
|
|
operations[operation] = state
|
|
attributes["swt_classification_operations"] = operations
|
|
attributes["swt_classification_status"] = status
|
|
attributes["swt_classification_event_id"] = eventID
|
|
if errorCode != "" {
|
|
attributes["swt_classification_error_code"] = errorCode
|
|
} else {
|
|
delete(attributes, "swt_classification_error_code")
|
|
}
|
|
if errorMessage != "" {
|
|
attributes["swt_classification_error"] = errorMessage
|
|
} else {
|
|
delete(attributes, "swt_classification_error")
|
|
}
|
|
encoded, err := json.Marshal(attributes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return tx.WithContext(ctx).Model(¤t).Update("additional_attributes", encoded).Error
|
|
}
|
|
|
|
func containsConversationKind(kinds []shangwutongConversationKind, id string) bool {
|
|
for _, kind := range kinds {
|
|
if kind.ID == id {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateClassificationCatalog(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
var request shangwutongClassificationCallbackRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classifications", "classification payload is invalid", false)
|
|
return
|
|
}
|
|
request.EventID = strings.TrimSpace(request.EventID)
|
|
if request.EventID == "" || len(request.EventID) > 128 {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classifications", "classification payload is invalid", false)
|
|
return
|
|
}
|
|
expectedKey := fmt.Sprintf("swt-classification-sync:%d:%s", inbox.ID, request.EventID)
|
|
if c.GetHeader("Idempotency-Key") != expectedKey {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_idempotency_key", "Idempotency-Key does not match event_id", false)
|
|
return
|
|
}
|
|
if !validClassificationCatalog(request) {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classifications", "classification payload is invalid", false)
|
|
return
|
|
}
|
|
conversationKinds, err := json.Marshal(request.ConversationKinds)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classifications", "classification payload is invalid", false)
|
|
return
|
|
}
|
|
customerColors, err := json.Marshal(request.CustomerColorKinds)
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classifications", "classification payload is invalid", false)
|
|
return
|
|
}
|
|
var catalogStale bool
|
|
var catalogNoop bool
|
|
var catalogConflict bool
|
|
var catalogEventID, catalogStatus string
|
|
err = h.db.WithContext(c.Request.Context()).Transaction(func(tx *gorm.DB) error {
|
|
var cache model.ShangwutongClassificationCache
|
|
err := tx.WithContext(c.Request.Context()).Clauses(clause.Locking{Strength: "UPDATE"}).Where("inbox_id = ?", inbox.ID).First(&cache).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
catalogEventID, catalogStatus = "", "never"
|
|
catalogConflict = true
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if cache.LastSyncEventID != request.EventID {
|
|
catalogStale = true
|
|
catalogEventID, catalogStatus = cache.LastSyncEventID, cache.SyncStatus
|
|
return nil
|
|
}
|
|
if cache.SyncStatus == "succeeded" {
|
|
if string(cache.ConversationKinds) == string(conversationKinds) && string(cache.CustomerColorKinds) == string(customerColors) {
|
|
catalogNoop = true
|
|
return nil
|
|
}
|
|
catalogConflict = true
|
|
return nil
|
|
}
|
|
now := time.Now().UTC()
|
|
if err := tx.WithContext(c.Request.Context()).Model(&cache).Updates(map[string]any{
|
|
"conversation_kinds": conversationKinds, "customer_color_kinds": customerColors,
|
|
"sync_status": "succeeded", "synced_at": &now, "last_error_code": nil,
|
|
"last_error_message": nil, "last_sync_event_id": request.EventID,
|
|
}).Error; err != nil {
|
|
return err
|
|
}
|
|
catalogStatus = "succeeded"
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "classification_cache_update_failed", "failed to save classifications", true)
|
|
return
|
|
}
|
|
if catalogConflict {
|
|
h.connectorError(c, http.StatusConflict, "unknown_operation", "classification sync operation is not pending", false)
|
|
return
|
|
}
|
|
if catalogStale {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "sync_status": catalogStatus, "event_id": catalogEventID})
|
|
return
|
|
}
|
|
if catalogNoop {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "sync_status": "succeeded", "event_id": request.EventID})
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"updated": true, "sync_status": "succeeded"})
|
|
}
|
|
|
|
type shangwutongClassificationSyncStatusRequest struct {
|
|
EventID string `json:"event_id"`
|
|
Status string `json:"status"`
|
|
ErrorCode string `json:"error_code"`
|
|
ErrorMessage string `json:"error_message"`
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateClassificationSyncStatus(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
var request shangwutongClassificationSyncStatusRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_sync_status", "classification sync status is invalid", false)
|
|
return
|
|
}
|
|
request.EventID = strings.TrimSpace(request.EventID)
|
|
request.Status = strings.TrimSpace(request.Status)
|
|
request.ErrorCode = strings.TrimSpace(request.ErrorCode)
|
|
request.ErrorMessage = strings.TrimSpace(request.ErrorMessage)
|
|
if request.EventID == "" || request.Status != "failed" || len(request.ErrorCode) > 128 || len(request.ErrorMessage) > 1024 {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_sync_status", "classification sync status is invalid", false)
|
|
return
|
|
}
|
|
expectedKey := fmt.Sprintf("swt-classification-sync:%d:%s", inbox.ID, request.EventID)
|
|
if c.GetHeader("Idempotency-Key") != expectedKey {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_idempotency_key", "Idempotency-Key does not match event_id", false)
|
|
return
|
|
}
|
|
var syncStale bool
|
|
var syncNoop bool
|
|
var syncConflict bool
|
|
var syncEventID, syncStatus string
|
|
err := h.db.WithContext(c.Request.Context()).Transaction(func(tx *gorm.DB) error {
|
|
var cache model.ShangwutongClassificationCache
|
|
err := tx.WithContext(c.Request.Context()).Clauses(clause.Locking{Strength: "UPDATE"}).Where("inbox_id = ?", inbox.ID).First(&cache).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
syncConflict = true
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if cache.LastSyncEventID != request.EventID {
|
|
syncStale = true
|
|
syncEventID, syncStatus = cache.LastSyncEventID, cache.SyncStatus
|
|
return nil
|
|
}
|
|
if cache.SyncStatus == "succeeded" {
|
|
syncConflict = true
|
|
return nil
|
|
}
|
|
if cache.SyncStatus == "failed" && dereferenceString(cache.LastErrorCode) == request.ErrorCode && dereferenceString(cache.LastErrorMessage) == request.ErrorMessage {
|
|
syncNoop = true
|
|
return nil
|
|
}
|
|
if cache.SyncStatus == "failed" {
|
|
syncConflict = true
|
|
return nil
|
|
}
|
|
if err := tx.WithContext(c.Request.Context()).Model(&cache).Updates(map[string]any{
|
|
"last_sync_event_id": request.EventID, "sync_status": request.Status,
|
|
"last_error_code": stringPointer(request.ErrorCode), "last_error_message": stringPointer(request.ErrorMessage),
|
|
}).Error; err != nil {
|
|
return err
|
|
}
|
|
syncStatus = request.Status
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "classification_cache_update_failed", "failed to save classification sync status", true)
|
|
return
|
|
}
|
|
if syncConflict {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "classification sync result conflicts with its operation", false)
|
|
return
|
|
}
|
|
if syncStale {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "sync_status": syncStatus, "event_id": syncEventID})
|
|
return
|
|
}
|
|
if syncNoop {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "sync_status": "failed", "event_id": request.EventID})
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"updated": true, "sync_status": syncStatus})
|
|
}
|
|
|
|
func stringPointer(value string) *string {
|
|
value = strings.TrimSpace(value)
|
|
if value == "" {
|
|
return nil
|
|
}
|
|
return &value
|
|
}
|
|
|
|
func dereferenceString(value *string) string {
|
|
if value == nil {
|
|
return ""
|
|
}
|
|
return strings.TrimSpace(*value)
|
|
}
|
|
|
|
type shangwutongClassificationStatusRequest struct {
|
|
EventID string `json:"event_id"`
|
|
Operation string `json:"operation"`
|
|
Status string `json:"status"`
|
|
ChatKindID string `json:"chat_kind_id"`
|
|
CustomerColorID string `json:"customer_color_id"`
|
|
ErrorCode string `json:"error_code"`
|
|
ErrorMessage string `json:"error_message"`
|
|
}
|
|
|
|
func classificationStateString(state map[string]any, key string) (string, bool) {
|
|
value, ok := state[key].(string)
|
|
value = strings.TrimSpace(value)
|
|
return value, ok && value != ""
|
|
}
|
|
|
|
func classificationStateUint(state map[string]any, key string) (uint, bool) {
|
|
value, ok := state[key]
|
|
if !ok {
|
|
return 0, false
|
|
}
|
|
switch value := value.(type) {
|
|
case float64:
|
|
if value <= 0 || value != float64(uint(value)) {
|
|
return 0, false
|
|
}
|
|
return uint(value), true
|
|
case int:
|
|
return uint(value), value > 0
|
|
case int64:
|
|
return uint(value), value > 0
|
|
case uint:
|
|
return value, value > 0
|
|
case uint64:
|
|
return uint(value), value > 0 && uint64(uint(value)) == value
|
|
default:
|
|
return 0, false
|
|
}
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) UpdateClassificationStatus(c *gin.Context) {
|
|
inbox, ok := h.authorizedInbox(c)
|
|
if !ok {
|
|
return
|
|
}
|
|
conversationID, err := strconv.ParseUint(c.Params.ByName("conversation_id"), 10, 64)
|
|
if err != nil || conversationID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "conversation not found", false)
|
|
return
|
|
}
|
|
var request shangwutongClassificationStatusRequest
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_status", "classification status is invalid", false)
|
|
return
|
|
}
|
|
request.EventID = strings.TrimSpace(request.EventID)
|
|
request.Operation = strings.TrimSpace(request.Operation)
|
|
request.Status = strings.TrimSpace(request.Status)
|
|
request.ChatKindID = strings.TrimSpace(request.ChatKindID)
|
|
request.CustomerColorID = strings.TrimSpace(request.CustomerColorID)
|
|
request.ErrorCode = strings.TrimSpace(request.ErrorCode)
|
|
request.ErrorMessage = strings.TrimSpace(request.ErrorMessage)
|
|
if request.EventID == "" || !oneOf(request.Status, "succeeded", "failed", "uncertain") || len(request.ErrorCode) > 128 || len(request.ErrorMessage) > 1024 {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_status", "classification status is invalid", false)
|
|
return
|
|
}
|
|
expectedKey := fmt.Sprintf("swt-classification-operation:%d:%s", inbox.ID, request.EventID)
|
|
if c.GetHeader("Idempotency-Key") != expectedKey {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_idempotency_key", "Idempotency-Key does not match event_id", false)
|
|
return
|
|
}
|
|
if request.Operation != "set_chat_kind" && request.Operation != "set_customer_color" ||
|
|
(request.Operation == "set_chat_kind" && request.ChatKindID == "") ||
|
|
(request.Operation == "set_customer_color" && request.CustomerColorID == "") ||
|
|
(request.ChatKindID != "" && request.CustomerColorID != "") {
|
|
h.connectorError(c, http.StatusUnprocessableEntity, "invalid_classification_status", "classification operation and target are invalid", false)
|
|
return
|
|
}
|
|
var conversation model.Conversation
|
|
if err := h.db.WithContext(c.Request.Context()).Where("id = ? AND inbox_id = ?", uint(conversationID), inbox.ID).First(&conversation).Error; err != nil {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "conversation not found", false)
|
|
return
|
|
}
|
|
attributes := map[string]any{}
|
|
if len(conversation.AdditionalAttributes) > 0 {
|
|
if err := json.Unmarshal(conversation.AdditionalAttributes, &attributes); err != nil {
|
|
h.connectorError(c, http.StatusInternalServerError, "conversation_attributes_invalid", "conversation attributes are invalid", true)
|
|
return
|
|
}
|
|
}
|
|
if attributes == nil {
|
|
attributes = map[string]any{}
|
|
}
|
|
operationValue := request.ChatKindID
|
|
if request.Operation == "set_customer_color" {
|
|
operationValue = request.CustomerColorID
|
|
}
|
|
operations, _ := attributes["swt_classification_operations"].(map[string]any)
|
|
if operations == nil {
|
|
operations = map[string]any{}
|
|
}
|
|
existing, ok := operations[request.Operation].(map[string]any)
|
|
if !ok {
|
|
h.connectorError(c, http.StatusConflict, "unknown_operation", "classification operation is not pending", false)
|
|
return
|
|
}
|
|
existingEventID, _ := existing["event_id"].(string)
|
|
if existingEventID == "" {
|
|
h.connectorError(c, http.StatusConflict, "unknown_operation", "classification operation is not pending", false)
|
|
return
|
|
}
|
|
if existingContactID, ok := existing["contact_id"].(float64); ok && uint(existingContactID) != conversation.ContactID {
|
|
h.connectorError(c, http.StatusConflict, "stale_operation", "classification operation targets a different contact", false)
|
|
return
|
|
}
|
|
if existingContactInboxID, ok := existing["contact_inbox_id"].(float64); ok && conversation.ContactInboxID != nil && uint(existingContactInboxID) != *conversation.ContactInboxID {
|
|
h.connectorError(c, http.StatusConflict, "stale_operation", "classification operation targets a different contact inbox", false)
|
|
return
|
|
}
|
|
if existingEventID != request.EventID {
|
|
status, _ := existing["status"].(string)
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "event_id": existingEventID, "status": status})
|
|
return
|
|
}
|
|
existingStatus, _ := existing["status"].(string)
|
|
existingValue, _ := existing["value"].(string)
|
|
existingCode, _ := existing["error_code"].(string)
|
|
existingMessage, _ := existing["error_message"].(string)
|
|
if existingValue != operationValue {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "classification result changes its target value", false)
|
|
return
|
|
}
|
|
if existingStatus != "pending" && existingStatus != "" && existingStatus != request.Status &&
|
|
!(existingStatus == "uncertain" && existingCode == "connector_delivery_uncertain") {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "classification result conflicts with its terminal state", false)
|
|
return
|
|
}
|
|
if existingStatus == request.Status && existingValue == operationValue && existingCode == request.ErrorCode && existingMessage == request.ErrorMessage {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "status": existingStatus, "event_id": existingEventID})
|
|
return
|
|
}
|
|
var staleEventID, staleStatus string
|
|
var classificationConflict bool
|
|
var classificationErrorCode, classificationErrorMessage string
|
|
var classificationNoop bool
|
|
if err := h.db.WithContext(c.Request.Context()).Transaction(func(tx *gorm.DB) error {
|
|
var conversations []model.Conversation
|
|
var colorContactInboxIDs map[uint]struct{}
|
|
query := tx.WithContext(c.Request.Context()).Clauses(clause.Locking{Strength: "UPDATE"}).Where(
|
|
"account_id = ? AND inbox_id = ?", conversation.AccountID, conversation.InboxID,
|
|
)
|
|
if request.Operation == "set_customer_color" && request.Status == "succeeded" {
|
|
currentCID := ""
|
|
if conversation.ContactInboxID == nil || *conversation.ContactInboxID == 0 {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation has no active customer binding"
|
|
return nil
|
|
}
|
|
var currentBinding model.ContactInbox
|
|
if err := tx.WithContext(c.Request.Context()).Clauses(clause.Locking{Strength: "UPDATE"}).Where(
|
|
"id = ? AND contact_id = ? AND inbox_id = ?", *conversation.ContactInboxID, conversation.ContactID, inbox.ID,
|
|
).First(¤tBinding).Error; err != nil {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation has no active customer binding"
|
|
return nil
|
|
}
|
|
var currentMetadata struct {
|
|
CID string `json:"cid"`
|
|
}
|
|
if len(currentBinding.ChannelMetadata) > 0 && json.Unmarshal(currentBinding.ChannelMetadata, ¤tMetadata) == nil {
|
|
currentCID = strings.TrimSpace(currentMetadata.CID)
|
|
}
|
|
if currentCID == "" {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation has no active customer binding"
|
|
return nil
|
|
}
|
|
var contactInboxes []model.ContactInbox
|
|
if err := tx.WithContext(c.Request.Context()).Clauses(clause.Locking{Strength: "UPDATE"}).Where("inbox_id = ?", inbox.ID).Find(&contactInboxes).Error; err != nil {
|
|
return err
|
|
}
|
|
targetContactInboxIDs := make([]uint, 0, len(contactInboxes))
|
|
colorContactInboxIDs = make(map[uint]struct{}, len(contactInboxes))
|
|
for _, candidate := range contactInboxes {
|
|
var metadata struct {
|
|
CID string `json:"cid"`
|
|
}
|
|
if len(candidate.ChannelMetadata) == 0 || json.Unmarshal(candidate.ChannelMetadata, &metadata) != nil {
|
|
continue
|
|
}
|
|
if strings.TrimSpace(metadata.CID) == currentCID {
|
|
targetContactInboxIDs = append(targetContactInboxIDs, candidate.ID)
|
|
colorContactInboxIDs[candidate.ID] = struct{}{}
|
|
}
|
|
}
|
|
if len(targetContactInboxIDs) == 0 {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation has no active customer binding"
|
|
return nil
|
|
}
|
|
query = query.Where("contact_inbox_id IN ?", targetContactInboxIDs)
|
|
} else {
|
|
query = query.Where("id = ?", conversation.ID)
|
|
}
|
|
if err := query.Order("id ASC").Find(&conversations).Error; err != nil {
|
|
return err
|
|
}
|
|
var source *model.Conversation
|
|
for i := range conversations {
|
|
if conversations[i].ID == conversation.ID {
|
|
source = &conversations[i]
|
|
break
|
|
}
|
|
}
|
|
if source == nil {
|
|
return gorm.ErrRecordNotFound
|
|
}
|
|
|
|
sourceAttributes := map[string]any{}
|
|
if len(source.AdditionalAttributes) > 0 {
|
|
if err := json.Unmarshal(source.AdditionalAttributes, &sourceAttributes); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if sourceAttributes == nil {
|
|
sourceAttributes = map[string]any{}
|
|
}
|
|
sourceOperations, _ := sourceAttributes["swt_classification_operations"].(map[string]any)
|
|
sourceState, ok := sourceOperations[request.Operation].(map[string]any)
|
|
if !ok {
|
|
classificationErrorCode = "unknown_operation"
|
|
classificationErrorMessage = "classification operation is not pending"
|
|
return nil
|
|
}
|
|
storedEventID, eventOK := classificationStateString(sourceState, "event_id")
|
|
storedOperation, operationOK := classificationStateString(sourceState, "operation")
|
|
storedValue, valueOK := classificationStateString(sourceState, "value")
|
|
storedStatus, statusOK := classificationStateString(sourceState, "status")
|
|
storedAccountID, accountOK := classificationStateUint(sourceState, "account_id")
|
|
storedInboxID, inboxOK := classificationStateUint(sourceState, "inbox_id")
|
|
storedConversationID, conversationOK := classificationStateUint(sourceState, "conversation_id")
|
|
storedContactID, contactOK := classificationStateUint(sourceState, "contact_id")
|
|
storedContactInboxID, contactInboxOK := classificationStateUint(sourceState, "contact_inbox_id")
|
|
storedSessionID, sessionOK := classificationStateString(sourceState, "swt_session_id")
|
|
storedCID, _ := sourceState["cid"].(string)
|
|
storedCID = strings.TrimSpace(storedCID)
|
|
if !eventOK || !operationOK || !valueOK || !statusOK || !accountOK || !inboxOK || !conversationOK || !contactOK || !contactInboxOK || !sessionOK || storedOperation != request.Operation {
|
|
classificationErrorCode = "unknown_operation"
|
|
classificationErrorMessage = "classification operation identity is incomplete"
|
|
return nil
|
|
}
|
|
if storedAccountID != inbox.AccountID || storedInboxID != inbox.ID || storedConversationID != source.ID || storedContactID != source.ContactID {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation targets a different object"
|
|
return nil
|
|
}
|
|
if source.ContactInboxID != nil && *source.ContactInboxID != storedContactInboxID {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation targets a different contact inbox"
|
|
return nil
|
|
}
|
|
var contactInbox model.ContactInbox
|
|
if err := tx.WithContext(c.Request.Context()).Where(
|
|
"id = ? AND contact_id = ? AND inbox_id = ?", storedContactInboxID, source.ContactID, inbox.ID,
|
|
).First(&contactInbox).Error; err != nil || strings.TrimSpace(contactInbox.SourceID) != storedSessionID {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation targets a different remote session"
|
|
return nil
|
|
}
|
|
var contactMetadata struct {
|
|
CID string `json:"cid"`
|
|
}
|
|
if len(contactInbox.ChannelMetadata) > 0 && json.Unmarshal(contactInbox.ChannelMetadata, &contactMetadata) != nil {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation metadata is invalid"
|
|
return nil
|
|
}
|
|
if request.Operation == "set_customer_color" && (storedCID == "" || strings.TrimSpace(contactMetadata.CID) != storedCID) {
|
|
classificationErrorCode = "stale_operation"
|
|
classificationErrorMessage = "classification operation targets a different customer"
|
|
return nil
|
|
}
|
|
if storedEventID != request.EventID {
|
|
staleEventID, staleStatus = storedEventID, storedStatus
|
|
return nil
|
|
}
|
|
if storedValue != operationValue {
|
|
classificationErrorCode = "idempotency_conflict"
|
|
classificationErrorMessage = "classification result changes its target value"
|
|
return nil
|
|
}
|
|
if storedStatus != "pending" {
|
|
existingCode, _ := sourceState["error_code"].(string)
|
|
existingMessage, _ := sourceState["error_message"].(string)
|
|
if storedStatus == request.Status && existingCode == request.ErrorCode && existingMessage == request.ErrorMessage {
|
|
classificationNoop = true
|
|
return nil
|
|
}
|
|
if !(storedStatus == "uncertain" && existingCode == "connector_delivery_uncertain") {
|
|
classificationConflict = true
|
|
return nil
|
|
}
|
|
}
|
|
|
|
nextState := make(map[string]any, len(sourceState)+2)
|
|
for key, value := range sourceState {
|
|
nextState[key] = value
|
|
}
|
|
nextState["status"] = request.Status
|
|
nextState["error_code"] = request.ErrorCode
|
|
nextState["error_message"] = request.ErrorMessage
|
|
nextState["updated_at"] = time.Now().UTC()
|
|
sourceOperations[request.Operation] = nextState
|
|
for i := range conversations {
|
|
isSource := conversations[i].ID == source.ID
|
|
isColorTarget := isSource
|
|
if request.Operation == "set_customer_color" && request.Status == "succeeded" && !isSource && conversations[i].ContactInboxID != nil {
|
|
_, isColorTarget = colorContactInboxIDs[*conversations[i].ContactInboxID]
|
|
}
|
|
if !isSource && !isColorTarget {
|
|
continue
|
|
}
|
|
current := map[string]any{}
|
|
if len(conversations[i].AdditionalAttributes) > 0 {
|
|
if err := json.Unmarshal(conversations[i].AdditionalAttributes, ¤t); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if current == nil {
|
|
current = map[string]any{}
|
|
}
|
|
if request.Status == "succeeded" {
|
|
if request.Operation == "set_chat_kind" {
|
|
current["swt_chat_kind"] = request.ChatKindID
|
|
} else {
|
|
current["swt_label_color"] = request.CustomerColorID
|
|
}
|
|
}
|
|
if isSource {
|
|
currentOperations, _ := current["swt_classification_operations"].(map[string]any)
|
|
if currentOperations == nil {
|
|
currentOperations = map[string]any{}
|
|
}
|
|
currentOperations[request.Operation] = nextState
|
|
current["swt_classification_operations"] = currentOperations
|
|
current["swt_classification_status"] = request.Status
|
|
current["swt_classification_event_id"] = request.EventID
|
|
if request.ErrorCode != "" {
|
|
current["swt_classification_error_code"] = request.ErrorCode
|
|
} else {
|
|
delete(current, "swt_classification_error_code")
|
|
}
|
|
if request.ErrorMessage != "" {
|
|
current["swt_classification_error"] = request.ErrorMessage
|
|
} else {
|
|
delete(current, "swt_classification_error")
|
|
}
|
|
}
|
|
updated, err := json.Marshal(current)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := tx.WithContext(c.Request.Context()).Model(&conversations[i]).Update("additional_attributes", updated).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}); err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "conversation not found", false)
|
|
} else {
|
|
h.connectorError(c, http.StatusInternalServerError, "conversation_update_failed", "failed to update conversation classification", true)
|
|
}
|
|
return
|
|
}
|
|
if classificationErrorCode != "" {
|
|
h.connectorError(c, http.StatusConflict, classificationErrorCode, classificationErrorMessage, false)
|
|
return
|
|
}
|
|
if classificationConflict {
|
|
h.connectorError(c, http.StatusConflict, "idempotency_conflict", "classification result conflicts with its terminal state", false)
|
|
return
|
|
}
|
|
if classificationNoop {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "status": request.Status, "event_id": request.EventID})
|
|
return
|
|
}
|
|
if staleEventID != "" {
|
|
c.JSON(http.StatusOK, gin.H{"updated": false, "stale": true, "event_id": staleEventID, "status": staleStatus})
|
|
return
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"updated": true, "status": request.Status, "event_id": request.EventID})
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) userInbox(c *gin.Context) (*model.Inbox, bool) {
|
|
accountID, err := strconv.ParseUint(c.Param("account_id"), 10, 64)
|
|
if err != nil || accountID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "inbox not found", false)
|
|
return nil, false
|
|
}
|
|
inboxID, err := strconv.ParseUint(c.Param("inbox_id"), 10, 64)
|
|
if err != nil || inboxID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "inbox not found", false)
|
|
return nil, false
|
|
}
|
|
var inbox model.Inbox
|
|
if err := h.db.WithContext(c.Request.Context()).Where("id = ? AND account_id = ? AND channel_type = ?", uint(inboxID), uint(accountID), "shangwutong").First(&inbox).Error; err != nil {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "inbox not found", false)
|
|
return nil, false
|
|
}
|
|
return &inbox, true
|
|
}
|
|
|
|
func classificationCacheResponse(cache *model.ShangwutongClassificationCache) shangwutongClassificationCacheResponse {
|
|
conversationKinds := json.RawMessage(cache.ConversationKinds)
|
|
if len(conversationKinds) == 0 {
|
|
conversationKinds = json.RawMessage(`[]`)
|
|
}
|
|
customerColors := json.RawMessage(cache.CustomerColorKinds)
|
|
if len(customerColors) == 0 {
|
|
customerColors = json.RawMessage(`[]`)
|
|
}
|
|
return shangwutongClassificationCacheResponse{
|
|
InboxID: cache.InboxID, ConversationKinds: conversationKinds, CustomerColorKinds: customerColors,
|
|
SyncStatus: cache.SyncStatus, SyncedAt: cache.SyncedAt, LastErrorCode: cache.LastErrorCode, LastErrorMessage: cache.LastErrorMessage,
|
|
}
|
|
}
|
|
|
|
func validClassificationCatalog(request shangwutongClassificationCallbackRequest) bool {
|
|
if len(request.ConversationKinds) > 256 || len(request.CustomerColorKinds) > 256 {
|
|
return false
|
|
}
|
|
seen := make(map[string]struct{}, len(request.ConversationKinds)+len(request.CustomerColorKinds))
|
|
for _, kind := range request.ConversationKinds {
|
|
if strings.TrimSpace(kind.ID) == "" || strings.TrimSpace(kind.Name) == "" || len(kind.ID) > 128 || len(kind.Name) > 255 {
|
|
return false
|
|
}
|
|
key := "conversation:" + kind.ID
|
|
if _, ok := seen[key]; ok {
|
|
return false
|
|
}
|
|
seen[key] = struct{}{}
|
|
}
|
|
for _, color := range request.CustomerColorKinds {
|
|
if strings.TrimSpace(color.ID) == "" || strings.TrimSpace(color.Name) == "" || len(color.ID) > 128 || len(color.Name) > 255 {
|
|
return false
|
|
}
|
|
key := "color:" + color.ID
|
|
if _, ok := seen[key]; ok {
|
|
return false
|
|
}
|
|
seen[key] = struct{}{}
|
|
}
|
|
return true
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) inboxItem(c *gin.Context, inbox *model.Inbox) (shangwutongConnectorInbox, error) {
|
|
var config model.ChannelShangwutongConfig
|
|
if err := h.db.WithContext(c.Request.Context()).Where("inbox_id = ?", inbox.ID).First(&config).Error; err != nil {
|
|
return shangwutongConnectorInbox{}, err
|
|
}
|
|
var channelAPI channelmodel.ChannelAPI
|
|
if err := h.db.WithContext(c.Request.Context()).Where("inbox_id = ?", inbox.ID).First(&channelAPI).Error; err != nil {
|
|
return shangwutongConnectorInbox{}, err
|
|
}
|
|
return shangwutongConnectorInbox{
|
|
SchemaVersion: 1, AccountID: inbox.AccountID, InboxID: inbox.ID,
|
|
InboxIdentifier: channelAPI.Identifier, Enabled: inbox.Enabled,
|
|
DesiredPresence: config.DesiredPresence, ConfigVersion: config.ConfigVersion,
|
|
Credentials: shangwutongConnectorCredentials{
|
|
SessionID: config.SessionID, Username: config.Username, Password: config.Password,
|
|
HMACToken: channelAPI.HMACToken, WebhookSecret: channelAPI.Secret,
|
|
},
|
|
UpdatedAt: config.UpdatedAt.UTC(),
|
|
}, nil
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) authorizedInbox(c *gin.Context) (*model.Inbox, bool) {
|
|
inboxID, err := strconv.ParseUint(c.Param("inbox_id"), 10, 64)
|
|
if err != nil || inboxID == 0 {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "inbox not found", false)
|
|
return nil, false
|
|
}
|
|
platformAppID := middleware.ConnectorPlatformAppID(c)
|
|
var inbox model.Inbox
|
|
err = h.db.WithContext(c.Request.Context()).Table("inboxes").Select("inboxes.*").Joins(
|
|
"JOIN permissibles ON permissibles.permissible_id = inboxes.account_id AND permissibles.permissible_type = ?", model.PermissibleTypeAccount,
|
|
).Where(
|
|
"permissibles.platform_app_id = ? AND inboxes.id = ? AND inboxes.channel_type = ?", platformAppID, uint(inboxID), "shangwutong",
|
|
).First(&inbox).Error
|
|
if err != nil {
|
|
h.connectorError(c, http.StatusNotFound, "not_found", "inbox not found", false)
|
|
return nil, false
|
|
}
|
|
return &inbox, true
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) grantedAccountIDs(c *gin.Context) ([]uint, error) {
|
|
var accountIDs []uint
|
|
err := h.db.WithContext(c.Request.Context()).Model(&model.Permissible{}).Where(
|
|
"platform_app_id = ? AND permissible_type = ?", middleware.ConnectorPlatformAppID(c), model.PermissibleTypeAccount,
|
|
).Pluck("permissible_id", &accountIDs).Error
|
|
return accountIDs, err
|
|
}
|
|
|
|
func (h *ShangwutongConnectorHandler) connectorError(c *gin.Context, status int, code, message string, retryable bool) {
|
|
c.JSON(status, gin.H{"error": gin.H{
|
|
"code": code, "message": message, "retryable": retryable, "request_id": c.GetString("request_id"),
|
|
}})
|
|
}
|
|
|
|
func validShangwutongStatus(request shangwutongStatusRequest) bool {
|
|
return request.ConfigVersion > 0 && oneOf(request.ActualPresence, "online", "busy", "away", "offline") &&
|
|
oneOf(request.ConnectionStatus, "pending", "logging_in", "connected", "degraded", "relogin_required", "verification_required", "auth_failed", "disabled", "offline") &&
|
|
oneOf(request.CredentialStatus, "pending", "verifying", "applied", "rejected", "verification_required")
|
|
}
|
|
|
|
func validShangwutongMessageResult(request shangwutongMessageResultRequest) bool {
|
|
if request.ResultVersion <= 0 || request.OccurredAt == nil || request.OccurredAt.IsZero() || !oneOf(request.Status, "sent", "failed", "uncertain") {
|
|
return false
|
|
}
|
|
if request.ExternalID != nil {
|
|
if request.Status != "sent" || strings.TrimSpace(*request.ExternalID) == "" {
|
|
return false
|
|
}
|
|
if _, err := strconv.ParseUint(strings.TrimSpace(*request.ExternalID), 10, 64); err != nil {
|
|
return false
|
|
}
|
|
}
|
|
if len(request.ExternalIDs) > 100 || (len(request.ExternalIDs) > 0 && request.Status != "sent") {
|
|
return false
|
|
}
|
|
seen := make(map[string]struct{}, len(request.ExternalIDs))
|
|
for _, externalID := range request.ExternalIDs {
|
|
externalID = strings.TrimSpace(externalID)
|
|
if _, err := strconv.ParseUint(externalID, 10, 64); err != nil {
|
|
return false
|
|
}
|
|
if _, duplicate := seen[externalID]; duplicate {
|
|
return false
|
|
}
|
|
seen[externalID] = struct{}{}
|
|
}
|
|
return request.Status != "failed" || (request.ErrorCode != nil && strings.TrimSpace(*request.ErrorCode) != "")
|
|
}
|
|
|
|
func oneOf(value string, allowed ...string) bool {
|
|
for _, candidate := range allowed {
|
|
if value == candidate {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func encodeShangwutongCursor(id uint) string {
|
|
return base64.RawURLEncoding.EncodeToString([]byte(strconv.FormatUint(uint64(id), 10)))
|
|
}
|
|
|
|
func decodeShangwutongCursor(cursor string) (uint, error) {
|
|
if strings.TrimSpace(cursor) == "" {
|
|
return 0, nil
|
|
}
|
|
decoded, err := base64.RawURLEncoding.DecodeString(cursor)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
value, err := strconv.ParseUint(string(decoded), 10, 64)
|
|
if err != nil || value == 0 {
|
|
return 0, errors.New("invalid cursor")
|
|
}
|
|
return uint(value), nil
|
|
}
|