Files

231 lines
7.1 KiB
Go

// Package nationalcatalog 全国目录中心服务(方案第三阶段 — 全国统一标识与审核协同)。
//
// 功能:
// - 汇聚多省标识数据,建立全国统一目录
// - 提供跨省标识查询(按 MA 码、Hash、省级编码查询)
// - 省级目录库 ↔ 全国目录中心数据同步
// - 各省内容编码与全国统一 MA 之间的映射同步
// - 跨省审核记录共享
//
// 架构说明:
// - NationalCatalog 全国目录中心,聚合多省数据
// - ProvinceNode 省级节点抽象,通过 SyncService 向全国中心同步
// - 基于 internal/sync 包的同步能力
package nationalcatalog
import (
"fmt"
stdsync "sync"
"github.com/tcs-iptv/tcs/internal/model"
syncpkg "github.com/tcs-iptv/tcs/internal/sync"
)
// NationalCatalog 全国目录中心。
// 聚合多省标识数据,提供统一查询入口。
type NationalCatalog struct {
mu stdsync.RWMutex
contents map[string]model.Content // ccCode -> Content
bindings map[string][]model.HashBinding // ccCode -> bindings
mappings map[string][]model.Mapping // ccCode -> mappings
provinces map[string]*ProvinceInfo // provinceCode -> 省级信息
auditShared map[string][]model.ProvenanceEvent // ccCode -> 跨省共享审核记录
}
// ProvinceInfo 省级节点信息。
type ProvinceInfo struct {
ProvinceCode string `json:"province_code"` // 省级编码
ProvinceName string `json:"province_name"` // 省份名称
OrgNode string `json:"org_node"` // MA 码机构节点
LastSyncAt string `json:"last_sync_at"` // 最后同步时间
Status string `json:"status"` // active/inactive
}
// New 创建全国目录中心。
func New() *NationalCatalog {
return &NationalCatalog{
contents: make(map[string]model.Content),
bindings: make(map[string][]model.HashBinding),
mappings: make(map[string][]model.Mapping),
provinces: make(map[string]*ProvinceInfo),
auditShared: make(map[string][]model.ProvenanceEvent),
}
}
// RegisterProvince 注册省级节点。
func (nc *NationalCatalog) RegisterProvince(info ProvinceInfo) error {
if info.ProvinceCode == "" {
return fmt.Errorf("national: 省级编码不能为空")
}
nc.mu.Lock()
defer nc.mu.Unlock()
info.Status = "active"
nc.provinces[info.ProvinceCode] = &info
return nil
}
// ListProvinces 列出已注册的省级节点。
func (nc *NationalCatalog) ListProvinces() []ProvinceInfo {
nc.mu.RLock()
defer nc.mu.RUnlock()
out := make([]ProvinceInfo, 0, len(nc.provinces))
for _, p := range nc.provinces {
out = append(out, *p)
}
return out
}
// ---- 实现 sync.SyncSink 接口 ----
// UpsertContent 写入/更新内容记录。
func (nc *NationalCatalog) UpsertContent(c model.Content) error {
nc.mu.Lock()
defer nc.mu.Unlock()
if _, exists := nc.contents[c.CCCode]; exists {
return syncpkg.ErrConflict
}
nc.contents[c.CCCode] = c
return nil
}
// UpsertBinding 写入/更新哈希绑定。
func (nc *NationalCatalog) UpsertBinding(ccCode string, b model.HashBinding) error {
nc.mu.Lock()
defer nc.mu.Unlock()
nc.bindings[ccCode] = append(nc.bindings[ccCode], b)
return nil
}
// UpsertMapping 写入/更新映射。
func (nc *NationalCatalog) UpsertMapping(ccCode string, m model.Mapping) error {
nc.mu.Lock()
defer nc.mu.Unlock()
nc.mappings[ccCode] = append(nc.mappings[ccCode], m)
return nil
}
// ---- 全国统一查询 ----
// QueryByMA 按 MA 码查询(全国维度)。
func (nc *NationalCatalog) QueryByMA(ccCode string) (model.ContentQueryResult, error) {
nc.mu.RLock()
defer nc.mu.RUnlock()
c, ok := nc.contents[ccCode]
if !ok {
return model.ContentQueryResult{Found: false}, fmt.Errorf("national: MA 码 %s 未找到", ccCode)
}
return model.ContentQueryResult{
Found: true,
Content: c,
Bindings: nc.bindings[ccCode],
Mappings: nc.mappings[ccCode],
}, nil
}
// QueryByHash 按 Hash 查询(全国维度)。
func (nc *NationalCatalog) QueryByHash(fileHash string) (model.ContentQueryResult, error) {
nc.mu.RLock()
defer nc.mu.RUnlock()
// 先搜索 Content.FileHash(整剧主哈希)
for ccCode, c := range nc.contents {
if c.FileHash == fileHash {
return model.ContentQueryResult{
Found: true,
Content: c,
Bindings: nc.bindings[ccCode],
Mappings: nc.mappings[ccCode],
}, nil
}
}
// 再搜索 bindings 中的 HashValue(集级/转码版哈希)
for ccCode, bindings := range nc.bindings {
for _, b := range bindings {
if b.HashValue == fileHash {
return model.ContentQueryResult{
Found: true,
Content: nc.contents[ccCode],
Bindings: bindings,
Mappings: nc.mappings[ccCode],
}, nil
}
}
}
return model.ContentQueryResult{Found: false}, fmt.Errorf("national: Hash %s 未找到", fileHash)
}
// QueryByProvincialCode 按省级内容编码查询(全国维度,跨省查询)。
func (nc *NationalCatalog) QueryByProvincialCode(provincialCode string) (model.ContentQueryResult, error) {
nc.mu.RLock()
defer nc.mu.RUnlock()
for ccCode, mappings := range nc.mappings {
for _, mp := range mappings {
if mp.Party == model.PartyCP && mp.PartyID == provincialCode {
return model.ContentQueryResult{
Found: true,
Content: nc.contents[ccCode],
Bindings: nc.bindings[ccCode],
Mappings: mappings,
}, nil
}
}
}
return model.ContentQueryResult{Found: false}, fmt.Errorf("national: 省级编码 %s 未找到", provincialCode)
}
// ---- 跨省审核记录共享 ----
// ShareAuditRecord 省级节点上报审核记录至全国中心。
func (nc *NationalCatalog) ShareAuditRecord(ccCode string, event model.ProvenanceEvent) {
nc.mu.Lock()
defer nc.mu.Unlock()
nc.auditShared[ccCode] = append(nc.auditShared[ccCode], event)
}
// QuerySharedAudit 查询跨省共享的审核记录。
func (nc *NationalCatalog) QuerySharedAudit(ccCode string) []model.ProvenanceEvent {
nc.mu.RLock()
defer nc.mu.RUnlock()
return nc.auditShared[ccCode]
}
// ---- 全国统计 ----
// NationalStats 全国统计信息。
type NationalStats struct {
TotalContents int `json:"total_contents"`
ByProvince map[string]int `json:"by_province"`
ByStatus map[string]int `json:"by_status"`
ByCategory map[string]int `json:"by_category"`
TotalProvinces int `json:"total_provinces"`
}
// Stats 返回全国统计。
func (nc *NationalCatalog) Stats() NationalStats {
nc.mu.RLock()
defer nc.mu.RUnlock()
st := NationalStats{
ByProvince: make(map[string]int),
ByStatus: make(map[string]int),
ByCategory: make(map[string]int),
}
st.TotalContents = len(nc.contents)
st.TotalProvinces = len(nc.provinces)
for _, c := range nc.contents {
st.ByStatus[c.Status]++
st.ByCategory[c.MAType]++
// 按机构节点统计省份
for _, mp := range nc.mappings[c.CCCode] {
if mp.Party == model.PartyCP {
st.ByProvince[mp.PartyName]++
}
}
}
return st
}
// SyncFromProvince 从省级节点同步数据至全国中心。
func (nc *NationalCatalog) SyncFromProvince(source syncpkg.SyncSource, resolver syncpkg.ConflictResolver) (syncpkg.SyncResult, error) {
svc := syncpkg.New(source, nc)
return svc.Sync(syncpkg.SyncRequest{Resolver: resolver})
}