140 lines
4.1 KiB
Go
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))
|
|
}
|