// 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) }