diff --git a/Makefile b/Makefile index c70c9bbc1..6b48c3904 100644 --- a/Makefile +++ b/Makefile @@ -310,7 +310,7 @@ dev-test: && PKG_CONFIG_PATH=$(PKG_CONFIG_PATH) \ LD_LIBRARY_PATH=$(shell realpath $(BUILDDIR))/lib \ CGO_LDFLAGS="-lm" \ - bash -euo pipefail -c "go test -p 1 -timeout 300s ./pkg/... -v | tee /dev/stderr | go-junit-report -out test.xml" + bash -euo pipefail -c "go test -p 1 -timeout 30m ./pkg/... -v | tee /dev/stderr | go-junit-report -out test.xml" .PHONY: iroh-test iroh-test: diff --git a/pkg/atproto/moderation_test.go b/pkg/atproto/moderation_test.go new file mode 100644 index 000000000..b039ccac9 --- /dev/null +++ b/pkg/atproto/moderation_test.go @@ -0,0 +1,239 @@ +package atproto + +import ( + "context" + "fmt" + "strings" + "testing" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/util" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/constants" + "stream.place/streamplace/pkg/devenv" + "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/statedb" + "stream.place/streamplace/pkg/streamplace" +) + +func TestDelegatedModeration(t *testing.T) { + dev := devenv.WithDevEnv(t) + t.Logf("dev: %+v", dev) + cli := config.CLI{ + BroadcasterHost: "example.com", + DBURL: ":memory:", + RelayHost: strings.ReplaceAll(dev.PDSURL, "http://", "ws://"), + PLCURL: dev.PLCURL, + } + t.Logf("cli: %+v", cli) + b := bus.NewBus() + cli.DataDir = t.TempDir() + mod, err := model.MakeDB(":memory:") + require.NoError(t, err) + state, err := statedb.MakeDB(context.Background(), &cli, nil, mod) + require.NoError(t, err) + atsync := &ATProtoSynchronizer{ + CLI: &cli, + StatefulDB: state, + Model: mod, + Bus: b, + } + + ctx, cancel := context.WithCancel(context.Background()) + + done := make(chan struct{}) + + go func() { + err := atsync.StartFirehose(ctx) + require.NoError(t, err) + close(done) + }() + + streamer := dev.CreateAccount(t) + moderator := dev.CreateAccount(t) + user := dev.CreateAccount(t) + + // Test 1: Create delegation record with expiration time + t.Log("Test 1: Creating delegation record with expiration time") + expirationTime := time.Now().Add(24 * time.Hour).Format(time.RFC3339) + delegationRecord := &streamplace.ModerationPermission{ + LexiconTypeID: "place.stream.moderation.permission", + Moderator: moderator.DID, + Permissions: []string{"ban", "hide", "livestream.manage"}, + CreatedAt: time.Now().Format(util.ISO8601), + ExpirationTime: &expirationTime, + } + + _, err = comatproto.RepoCreateRecord(ctx, streamer.XRPC, &comatproto.RepoCreateRecord_Input{ + Collection: constants.PLACE_STREAM_MODERATION_PERMISSION, + Repo: streamer.DID, + Record: &lexutil.LexiconTypeDecoder{Val: delegationRecord}, + }) + require.NoError(t, err) + + // Wait for firehose ingestion + err = untilNoErrors(t, func() error { + delegation, err := mod.GetModerationDelegation(ctx, streamer.DID, moderator.DID) + if err != nil { + return err + } + if delegation == nil { + return fmt.Errorf("delegation not found") + } + return nil + }) + require.NoError(t, err) + t.Log("✓ Delegation record ingested successfully") + + // Verify delegation details including expiration time + view, err := mod.GetModerationDelegation(ctx, streamer.DID, moderator.DID) + require.NoError(t, err) + require.NotNil(t, view) + require.Equal(t, streamer.DID, view.Author.Did) + + delegation := view.Record.Val.(*streamplace.ModerationPermission) + require.NotNil(t, delegation) + require.Equal(t, moderator.DID, delegation.Moderator) + require.NotNil(t, delegation.ExpirationTime, "expiration time should be set") + + exp, err := time.Parse(time.RFC3339, *delegation.ExpirationTime) + require.NoError(t, err) + require.True(t, exp.After(time.Now()), "expiration time should be in the future") + t.Log("✓ Delegation details verified (including expiration time)") + + // Test 2: Create block (ban user) + t.Log("Test 2: Creating block record") + block := &bsky.GraphBlock{ + Subject: user.DID, + CreatedAt: time.Now().UTC().Format(time.RFC3339), + } + + blockRec, err := comatproto.RepoCreateRecord(ctx, streamer.XRPC, &comatproto.RepoCreateRecord_Input{ + Collection: constants.APP_BSKY_GRAPH_BLOCK, + Repo: streamer.DID, + Record: &lexutil.LexiconTypeDecoder{Val: block}, + }) + require.NoError(t, err) + t.Logf("✓ Block record created: %s", blockRec.Uri) + + // Wait for firehose to process block + blockRkey := strings.TrimPrefix(blockRec.Uri, fmt.Sprintf("at://%s/%s/", streamer.DID, constants.APP_BSKY_GRAPH_BLOCK)) + err = untilNoErrors(t, func() error { + block, err := mod.GetBlock(ctx, blockRkey) + if err != nil { + return err + } + if block == nil { + return fmt.Errorf("block not found") + } + return nil + }) + require.NoError(t, err) + t.Log("✓ Block record ingested successfully") + + // Test 3: Create chat message and gate + t.Log("Test 3: Creating chat message and gate") + msg := &streamplace.ChatMessage{ + LexiconTypeID: "place.stream.chat.message", + Text: "Test message to be hidden", + CreatedAt: time.Now().Format(util.ISO8601), + Streamer: streamer.DID, + } + + msgRec, err := comatproto.RepoCreateRecord(ctx, user.XRPC, &comatproto.RepoCreateRecord_Input{ + Collection: constants.PLACE_STREAM_CHAT_MESSAGE, + Repo: user.DID, + Record: &lexutil.LexiconTypeDecoder{Val: msg}, + }) + require.NoError(t, err) + t.Logf("✓ Chat message created: %s", msgRec.Uri) + + // Create gate to hide the message + gate := &streamplace.ChatGate{ + LexiconTypeID: "place.stream.chat.gate", + HiddenMessage: msgRec.Uri, + } + + gateRec, err := comatproto.RepoCreateRecord(ctx, streamer.XRPC, &comatproto.RepoCreateRecord_Input{ + Collection: constants.PLACE_STREAM_CHAT_GATE, + Repo: streamer.DID, + Record: &lexutil.LexiconTypeDecoder{Val: gate}, + }) + require.NoError(t, err) + t.Logf("✓ Gate record created: %s", gateRec.Uri) + + // Wait for firehose to process gate + gateRkey := strings.TrimPrefix(gateRec.Uri, fmt.Sprintf("at://%s/%s/", streamer.DID, constants.PLACE_STREAM_CHAT_GATE)) + err = untilNoErrors(t, func() error { + gate, err := mod.GetGate(ctx, gateRkey) + if err != nil { + return err + } + if gate == nil { + return fmt.Errorf("gate not found") + } + return nil + }) + require.NoError(t, err) + t.Log("✓ Gate record ingested successfully") + + // Test 4: Delete gate (unhide message) + t.Log("Test 4: Deleting gate record") + + _, err = comatproto.RepoDeleteRecord(ctx, streamer.XRPC, &comatproto.RepoDeleteRecord_Input{ + Collection: constants.PLACE_STREAM_CHAT_GATE, + Repo: streamer.DID, + Rkey: gateRkey, + }) + require.NoError(t, err) + + // Wait for firehose to process deletion + err = untilNoErrors(t, func() error { + gate, err := mod.GetGate(ctx, gateRkey) + if err != nil { + return err + } + if gate != nil { + return fmt.Errorf("gate still exists") + } + return nil + }) + require.NoError(t, err) + t.Log("✓ Gate record deleted successfully") + + // Test 5: Delete delegation (revoke permissions) + t.Log("Test 5: Deleting delegation record") + uri, err := syntax.ParseATURI(view.Uri) + require.NoError(t, err) + + _, err = comatproto.RepoDeleteRecord(ctx, streamer.XRPC, &comatproto.RepoDeleteRecord_Input{ + Collection: constants.PLACE_STREAM_MODERATION_PERMISSION, + Repo: streamer.DID, + Rkey: uri.RecordKey().String(), + }) + require.NoError(t, err) + + // Wait for firehose to process deletion + err = untilNoErrors(t, func() error { + delegation, err := mod.GetModerationDelegation(ctx, streamer.DID, moderator.DID) + if err != nil { + return err + } + if delegation != nil { + return fmt.Errorf("delegation still exists") + } + return nil + }) + require.NoError(t, err) + t.Log("✓ Delegation record deleted successfully") + + t.Log("All moderation tests passed!") + cancel() + <-done +} diff --git a/pkg/moderation/permissions.go b/pkg/moderation/permissions.go index 64e2fe56f..81b0a9d10 100644 --- a/pkg/moderation/permissions.go +++ b/pkg/moderation/permissions.go @@ -5,10 +5,13 @@ import ( "fmt" "time" - "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/streamplace" ) +type delegationGetter interface { + GetModerationDelegations(ctx context.Context, streamerDID, moderatorDID string) ([]*streamplace.ModerationDefs_PermissionView, error) +} + // Permission scope constants const ( PermissionBan = "ban" @@ -27,11 +30,11 @@ var ActionPermissions = map[string]string{ // PermissionChecker validates moderation permissions type PermissionChecker struct { - model model.Model + model delegationGetter } // NewPermissionChecker creates a new permission checker -func NewPermissionChecker(m model.Model) *PermissionChecker { +func NewPermissionChecker(m delegationGetter) *PermissionChecker { return &PermissionChecker{model: m} } diff --git a/pkg/moderation/permissions_test.go b/pkg/moderation/permissions_test.go new file mode 100644 index 000000000..5c68d51a9 --- /dev/null +++ b/pkg/moderation/permissions_test.go @@ -0,0 +1,243 @@ +package moderation + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/bluesky-social/indigo/api/bsky" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/streamplace" +) + +func TestPermissionChecker_CheckPermission_StreamerSelfModeration(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + streamerDID := "did:plc:streamer123" + + err := pc.CheckPermission(context.Background(), streamerDID, streamerDID, "createBlock") + require.NoError(t, err, "streamer should have permission for self-moderation") + + err = pc.CheckPermission(context.Background(), streamerDID, streamerDID, "createGate") + require.NoError(t, err, "streamer should have permission for self-moderation") + + err = pc.CheckPermission(context.Background(), streamerDID, streamerDID, "updateLivestream") + require.NoError(t, err, "streamer should have permission for self-moderation") +} + +func TestPermissionChecker_CheckPermission_WithCorrectPermission(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + mod.addPermissionView(streamerDID, moderatorDID, []string{"ban", "hide"}, nil) + + err := pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.NoError(t, err, "moderator with 'ban' permission should be able to createBlock") + + err = pc.CheckPermission(ctx, moderatorDID, streamerDID, "createGate") + require.NoError(t, err, "moderator with 'hide' permission should be able to createGate") +} + +func TestPermissionChecker_CheckPermission_WithWrongPermission(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + mod.addPermissionView(streamerDID, moderatorDID, []string{"hide"}, nil) + + err := pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.Error(t, err, "moderator with only 'hide' permission should not be able to createBlock") + require.Contains(t, err.Error(), "does not have permission 'ban'") +} + +func TestPermissionChecker_CheckPermission_WithoutAnyPermission(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + err := pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.Error(t, err, "moderator without any delegation should be denied") + require.Contains(t, err.Error(), "does not have permission") +} + +func TestPermissionChecker_CheckPermission_UnknownAction(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + streamerDID := "did:plc:streamer123" + + err := pc.CheckPermission(context.Background(), streamerDID, streamerDID, "unknownAction") + require.Error(t, err) + require.Contains(t, err.Error(), "unknown action") +} + +func TestPermissionChecker_HasPermission(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + mod.addPermissionView(streamerDID, moderatorDID, []string{"ban", "hide"}, nil) + + has, err := pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionBan) + require.NoError(t, err) + require.True(t, has, "should have 'ban' permission") + + has, err = pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionHide) + require.NoError(t, err) + require.True(t, has, "should have 'hide' permission") + + has, err = pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionLivestreamManage) + require.NoError(t, err) + require.False(t, has, "should not have 'livestream.manage' permission") +} + +func TestActionPermissions_Mapping(t *testing.T) { + require.Equal(t, PermissionBan, ActionPermissions["createBlock"]) + require.Equal(t, PermissionBan, ActionPermissions["deleteBlock"]) + require.Equal(t, PermissionHide, ActionPermissions["createGate"]) + require.Equal(t, PermissionHide, ActionPermissions["deleteGate"]) + require.Equal(t, PermissionLivestreamManage, ActionPermissions["updateLivestream"]) +} + +func TestPermissionChecker_HasPermission_MultipleSeparateRecords(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + mod.addPermissionView(streamerDID, moderatorDID, []string{"ban"}, nil) + mod.addPermissionView(streamerDID, moderatorDID, []string{"hide"}, nil) + + hasBan, err := pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionBan) + require.NoError(t, err) + require.True(t, hasBan, "should have 'ban' permission from first record") + + hasHide, err := pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionHide) + require.NoError(t, err) + require.True(t, hasHide, "should have 'hide' permission from second record") + + err = pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.NoError(t, err, "should allow createBlock with 'ban' permission from separate record") + + err = pc.CheckPermission(ctx, moderatorDID, streamerDID, "createGate") + require.NoError(t, err, "should allow createGate with 'hide' permission from separate record") +} + +func TestPermissionChecker_HasPermission_ExpiredDelegation(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + expiredTime := time.Now().Add(-1 * time.Hour) + mod.addPermissionView(streamerDID, moderatorDID, []string{"ban"}, &expiredTime) + + hasPermission, err := pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionBan) + require.NoError(t, err) + require.False(t, hasPermission, "should deny permission for expired delegation") + + err = pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.Error(t, err, "should deny action for expired delegation") + require.Contains(t, err.Error(), "does not have permission") +} + +func TestPermissionChecker_HasPermission_NotYetExpired(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + futureTime := time.Now().Add(1 * time.Hour) + mod.addPermissionView(streamerDID, moderatorDID, []string{"ban", "hide"}, &futureTime) + + hasPermission, err := pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionBan) + require.NoError(t, err) + require.True(t, hasPermission, "should allow permission for not-yet-expired delegation") + + err = pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.NoError(t, err, "should allow action for not-yet-expired delegation") +} + +func TestPermissionChecker_HasPermission_NoExpiration(t *testing.T) { + mod := newMockModel() + pc := NewPermissionChecker(mod) + + ctx := context.Background() + streamerDID := "did:plc:streamer123" + moderatorDID := "did:plc:moderator456" + + mod.addPermissionView(streamerDID, moderatorDID, []string{"ban", "hide"}, nil) + + hasPermission, err := pc.HasPermission(ctx, moderatorDID, streamerDID, PermissionBan) + require.NoError(t, err) + require.True(t, hasPermission, "should allow permission for delegation with no expiration") + + err = pc.CheckPermission(ctx, moderatorDID, streamerDID, "createBlock") + require.NoError(t, err, "should allow action for delegation with no expiration") +} + +type mockModel struct { + delegations map[string][]*streamplace.ModerationDefs_PermissionView +} + +func newMockModel() *mockModel { + return &mockModel{ + delegations: make(map[string][]*streamplace.ModerationDefs_PermissionView), + } +} + +func (m *mockModel) addPermissionView(streamerDID, moderatorDID string, permissions []string, expirationTime *time.Time) { + key := streamerDID + "_" + moderatorDID + + var expTimeStr *string + if expirationTime != nil { + str := expirationTime.Format(time.RFC3339) + expTimeStr = &str + } + + permRecord := &streamplace.ModerationPermission{ + Moderator: moderatorDID, + Permissions: permissions, + ExpirationTime: expTimeStr, + } + + view := &streamplace.ModerationDefs_PermissionView{ + Uri: fmt.Sprintf("at://%s/place.stream.moderation.permission/test", streamerDID), + Cid: "bafytest", + Author: &bsky.ActorDefs_ProfileViewBasic{Did: streamerDID}, + Record: &lexutil.LexiconTypeDecoder{Val: permRecord}, + } + + m.delegations[key] = append(m.delegations[key], view) +} + +func (m *mockModel) GetModerationDelegations(ctx context.Context, streamerDID, moderatorDID string) ([]*streamplace.ModerationDefs_PermissionView, error) { + key := streamerDID + "_" + moderatorDID + delegations, exists := m.delegations[key] + if !exists { + return nil, nil + } + return delegations, nil +}