Files

140 lines
4.1 KiB
Go

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