Files
MAcode/tcs-iptv/internal/sync/sync.go
T

167 lines
4.7 KiB
Go

// Package sync 标识同步服务(方案第三阶段 — 标识同步能力补足)。
//
// 提供多节点/多省间的标识数据同步能力:
// - 增量同步:按时间戳范围拉取变更内容
// - 全量同步:拉取全部内容记录
// - 同步冲突检测:基于 MA 码唯一性 + 内容哈希一致性
//
// 架构说明:
// - SyncSource 数据源接口(本地 chain.Client 实现)
// - SyncSink 数据汇接口(远端节点实现)
// - SyncService 同步编排器:拉取变更 → 冲突检测 → 推送
package sync
import (
"errors"
"fmt"
"time"
"github.com/tcs-iptv/tcs/internal/chain"
"github.com/tcs-iptv/tcs/internal/model"
)
// SyncSource 标识同步数据源接口。
type SyncSource interface {
ListContents(status string) ([]model.Content, error)
QueryContent(maCode string) (model.Content, error)
QueryMappings(maCode string) (chain.MappingsResult, error)
ListEpisodes(maCode string) ([]model.HashBinding, error)
}
// SyncSink 标识同步数据汇接口(远端节点实现)。
type SyncSink interface {
UpsertContent(c model.Content) error
UpsertBinding(maCode string, b model.HashBinding) error
UpsertMapping(maCode string, m model.Mapping) error
}
// ConflictResolver 同步冲突解决策略。
type ConflictResolver int
const (
// ConflictSkip 跳过冲突(保留远端数据)。
ConflictSkip ConflictResolver = iota
// ConflictOverwrite 覆盖远端数据(以本地为准)。
ConflictOverwrite
// ConflictFail 冲突时报错中止。
ConflictFail
)
// SyncRequest 同步请求。
type SyncRequest struct {
Since time.Time // 增量同步起始时间(零值表示全量)
Resolver ConflictResolver // 冲突解决策略
BatchSize int // 批次大小(0 表示不分批)
}
// SyncResult 同步结果。
type SyncResult struct {
TotalContents int `json:"total_contents"`
TotalBindings int `json:"total_bindings"`
TotalMappings int `json:"total_mappings"`
SkippedConflicts int `json:"skipped_conflicts"`
Failed int `json:"failed"`
}
// SyncService 标识同步编排器。
type SyncService struct {
source SyncSource
sink SyncSink
}
// New 创建标识同步服务。
func New(source SyncSource, sink SyncSink) *SyncService {
return &SyncService{source: source, sink: sink}
}
// Sync 执行标识数据同步。
func (s *SyncService) Sync(req SyncRequest) (SyncResult, error) {
result := SyncResult{}
// 拉取全部内容(增量同步需数据源支持时间过滤,MVP 全量拉取后按时间筛选)
contents, err := s.source.ListContents("")
if err != nil {
return result, fmt.Errorf("sync: 拉取内容列表失败: %w", err)
}
for _, c := range contents {
// 增量过滤:跳过早于 Since 的记录
if !req.Since.IsZero() && c.CreatedAt.Before(req.Since) {
continue
}
// 冲突检测:检查远端是否已存在该 MA 码
if req.Resolver == ConflictSkip {
if err := s.sink.UpsertContent(c); err != nil {
if errors.Is(err, ErrConflict) {
result.SkippedConflicts++
continue
}
result.Failed++
continue
}
} else {
if err := s.sink.UpsertContent(c); err != nil {
if errors.Is(err, ErrConflict) && req.Resolver == ConflictFail {
return result, fmt.Errorf("sync: 冲突 MA 码 %s: %w", c.MACode, err)
}
result.Failed++
continue
}
}
result.TotalContents++
// 同步哈希绑定
eps, _ := s.source.ListEpisodes(c.MACode)
for _, b := range eps {
if err := s.sink.UpsertBinding(c.MACode, b); err != nil {
result.Failed++
continue
}
result.TotalBindings++
}
// 同步映射
mr, _ := s.source.QueryMappings(c.MACode)
for _, m := range mr.Mappings {
if err := s.sink.UpsertMapping(c.MACode, m); err != nil {
result.Failed++
continue
}
result.TotalMappings++
}
}
return result, nil
}
// ErrConflict 同步冲突错误。
var ErrConflict = errors.New("sync: content already exists at remote (conflict)")
// ---- chain.Client 适配为 SyncSource ----
// ChainSource 将 chain.Client 适配为 SyncSource。
type ChainSource struct {
Client chain.Client
}
// ListContents 列出全部内容。
func (cs *ChainSource) ListContents(status string) ([]model.Content, error) {
return cs.Client.ListContents(status)
}
// QueryContent 查询内容主记录。
func (cs *ChainSource) QueryContent(maCode string) (model.Content, error) {
return cs.Client.QueryContent(maCode)
}
// QueryMappings 查询映射。
func (cs *ChainSource) QueryMappings(maCode string) (chain.MappingsResult, error) {
return cs.Client.QueryMappings(maCode)
}
// ListEpisodes 列出集级哈希。
func (cs *ChainSource) ListEpisodes(maCode string) ([]model.HashBinding, error) {
return cs.Client.ListEpisodes(maCode)
}