package sync import ( "errors" "sync" "testing" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/tcs-iptv/tcs/internal/chain" "github.com/tcs-iptv/tcs/internal/model" ) // memorySink 内存数据汇(测试用)。 type memorySink struct { mu sync.Mutex contents map[string]model.Content bindings map[string][]model.HashBinding mappings map[string][]model.Mapping conflict bool // 模拟冲突 } func newMemorySink() *memorySink { return &memorySink{ contents: make(map[string]model.Content), bindings: make(map[string][]model.HashBinding), mappings: make(map[string][]model.Mapping), } } func (m *memorySink) UpsertContent(c model.Content) error { m.mu.Lock() defer m.mu.Unlock() if m.conflict { if _, exists := m.contents[c.CCCode]; exists { return ErrConflict } } m.contents[c.CCCode] = c return nil } func (m *memorySink) UpsertBinding(ccCode string, b model.HashBinding) error { m.mu.Lock() defer m.mu.Unlock() m.bindings[ccCode] = append(m.bindings[ccCode], b) return nil } func (m *memorySink) UpsertMapping(ccCode string, mp model.Mapping) error { m.mu.Lock() defer m.mu.Unlock() m.mappings[ccCode] = append(m.mappings[ccCode], mp) return nil } func TestSyncService_FullSync(t *testing.T) { src := &ChainSource{Client: chain.NewMemoryChain()} sink := newMemorySink() // 在源端发码 _, err := src.Client.IssueMA(chain.RoleRegulator, chain.IssueRequest{ CCCode: "MA.156.8531.6101/WD/20260000001", ContentTwinID: "ctid-sync-001", FileHash: "fh-sync-001", MerkleRoot: "mr-sync-001", Episodes: []model.EpisodeHash{ {Episode: 1, FileSHA256: "fh-sync-001-E1"}, {Episode: 2, FileSHA256: "fh-sync-001-E2"}, }, Content: model.Content{Title: "同步测试剧", EpisodeCount: 2, MAType: "WD", Issuer: "测试局"}, }) require.NoError(t, err) // 注册映射 _, err = src.Client.RegisterMapping(chain.RoleCP, model.Mapping{ ContentTwinID: "ctid-sync-001", Party: model.PartyCP, PartyID: "PROV-SYNC-001", PartyName: "测试CP", }) require.NoError(t, err) // 执行全量同步 svc := New(src, sink) result, err := svc.Sync(SyncRequest{Resolver: ConflictOverwrite}) require.NoError(t, err) assert.Equal(t, 1, result.TotalContents) assert.Equal(t, 2, result.TotalBindings) assert.Equal(t, 1, result.TotalMappings) assert.Equal(t, 0, result.Failed) // 验证远端数据 assert.Equal(t, "同步测试剧", sink.contents["MA.156.8531.6101/WD/20260000001"].Title) assert.Len(t, sink.bindings["MA.156.8531.6101/WD/20260000001"], 2) } func TestSyncService_ConflictSkip(t *testing.T) { src := &ChainSource{Client: chain.NewMemoryChain()} sink := newMemorySink() sink.conflict = true // 源端发码 _, err := src.Client.IssueMA(chain.RoleRegulator, chain.IssueRequest{ CCCode: "MA.156.8531.6101/WD/20260000002", ContentTwinID: "ctid-sync-002", FileHash: "fh-sync-002", MerkleRoot: "mr-sync-002", Content: model.Content{Title: "冲突测试剧", MAType: "WD", Issuer: "测试局"}, }) require.NoError(t, err) // 远端预先存在同 MA 码 sink.contents["MA.156.8531.6101/WD/20260000002"] = model.Content{ CCCode: "MA.156.8531.6101/WD/20260000002", Title: "远端已有", } // 冲突跳过策略 svc := New(src, sink) result, err := svc.Sync(SyncRequest{Resolver: ConflictSkip}) require.NoError(t, err) assert.Equal(t, 1, result.SkippedConflicts) assert.Equal(t, 0, result.TotalContents) } func TestSyncService_ConflictFail(t *testing.T) { src := &ChainSource{Client: chain.NewMemoryChain()} sink := newMemorySink() sink.conflict = true _, err := src.Client.IssueMA(chain.RoleRegulator, chain.IssueRequest{ CCCode: "MA.156.8531.6101/WD/20260000003", ContentTwinID: "ctid-sync-003", FileHash: "fh-sync-003", MerkleRoot: "mr-sync-003", Content: model.Content{Title: "冲突Fail测试", MAType: "WD", Issuer: "测试局"}, }) require.NoError(t, err) sink.contents["MA.156.8531.6101/WD/20260000003"] = model.Content{ CCCode: "MA.156.8531.6101/WD/20260000003", Title: "远端已有", } svc := New(src, sink) _, err = svc.Sync(SyncRequest{Resolver: ConflictFail}) assert.Error(t, err) assert.True(t, errors.Is(err, ErrConflict)) }