diff --git a/internal/api/handlers/community/subscribe.go b/internal/api/handlers/community/subscribe.go index 060e598..68cdbd3 100644 --- a/internal/api/handlers/community/subscribe.go +++ b/internal/api/handlers/community/subscribe.go @@ -34,7 +34,7 @@ func (h *SubscribeHandler) HandleSubscribe(w http.ResponseWriter, r *http.Reques // Parse request body var req struct { - Community string `json:"community"` // DID only (per lexicon) + Community string `json:"community"` // DID only (per lexicon) ContentVisibility int `json:"contentVisibility"` // Optional: 1-5 scale, defaults to 3 } diff --git a/internal/core/communities/community.go b/internal/core/communities/community.go index 93fc8e5..0d359cf 100644 --- a/internal/core/communities/community.go +++ b/internal/core/communities/community.go @@ -55,12 +55,12 @@ type Subscription struct { // CommunityBlock represents a user blocking a community // Block records live in the user's repository (at://user_did/social.coves.community.block/{rkey}) type CommunityBlock struct { - ID int `json:"id" db:"id"` + BlockedAt time.Time `json:"blockedAt" db:"blocked_at"` UserDID string `json:"userDid" db:"user_did"` CommunityDID string `json:"communityDid" db:"community_did"` - BlockedAt time.Time `json:"blockedAt" db:"blocked_at"` RecordURI string `json:"recordUri,omitempty" db:"record_uri"` RecordCID string `json:"recordCid,omitempty" db:"record_cid"` + ID int `json:"id" db:"id"` } // Membership represents active participation with reputation tracking diff --git a/internal/core/communities/service.go b/internal/core/communities/service.go index 3013c72..137523f 100644 --- a/internal/core/communities/service.go +++ b/internal/core/communities/service.go @@ -5,6 +5,7 @@ import ( "bytes" "context" "encoding/json" + "errors" "fmt" "io" "log" @@ -565,10 +566,16 @@ func (s *communityService) BlockCommunity(ctx context.Context, userDID, userAcce // Block exists in our index - return it return existingBlock, nil } - // Race condition: PDS has the block but Jetstream hasn't indexed it yet - // Return typed conflict error so handler can return 409 instead of 500 - // This is normal in eventually-consistent systems - return nil, ErrBlockAlreadyExists + // Only treat as "already exists" if the error is ErrBlockNotFound (race condition) + // Any other error (DB outage, connection failure, etc.) should bubble up + if errors.Is(getErr, ErrBlockNotFound) { + // Race condition: PDS has the block but Jetstream hasn't indexed it yet + // Return typed conflict error so handler can return 409 instead of 500 + // This is normal in eventually-consistent systems + return nil, ErrBlockAlreadyExists + } + // Real datastore error - bubble it up so operators see the failure + return nil, fmt.Errorf("PDS reported duplicate block but failed to fetch from index: %w", getErr) } return nil, fmt.Errorf("failed to create block on PDS: %w", err) } @@ -724,22 +731,6 @@ func (s *communityService) validateCreateRequest(req CreateCommunityRequest) err // PDS write-forward helpers -func (s *communityService) createRecordOnPDS(ctx context.Context, repoDID, collection, rkey string, record map[string]interface{}) (string, string, error) { - endpoint := fmt.Sprintf("%s/xrpc/com.atproto.repo.createRecord", strings.TrimSuffix(s.pdsURL, "/")) - - payload := map[string]interface{}{ - "repo": repoDID, - "collection": collection, - "record": record, - } - - if rkey != "" { - payload["rkey"] = rkey - } - - return s.callPDS(ctx, "POST", endpoint, payload) -} - // createRecordOnPDSAs creates a record with a specific access token (for V2 community auth) func (s *communityService) createRecordOnPDSAs(ctx context.Context, repoDID, collection, rkey string, record map[string]interface{}, accessToken string) (string, string, error) { endpoint := fmt.Sprintf("%s/xrpc/com.atproto.repo.createRecord", strings.TrimSuffix(s.pdsURL, "/")) @@ -771,21 +762,8 @@ func (s *communityService) putRecordOnPDSAs(ctx context.Context, repoDID, collec return s.callPDSWithAuth(ctx, "POST", endpoint, payload, accessToken) } -func (s *communityService) deleteRecordOnPDS(ctx context.Context, repoDID, collection, rkey string) error { - endpoint := fmt.Sprintf("%s/xrpc/com.atproto.repo.deleteRecord", strings.TrimSuffix(s.pdsURL, "/")) - - payload := map[string]interface{}{ - "repo": repoDID, - "collection": collection, - "rkey": rkey, - } - - _, _, err := s.callPDS(ctx, "POST", endpoint, payload) - return err -} - // deleteRecordOnPDSAs deletes a record with a specific access token (for user-scoped deletions) -func (s *communityService) deleteRecordOnPDSAs(ctx context.Context, repoDID, collection, rkey string, accessToken string) error { +func (s *communityService) deleteRecordOnPDSAs(ctx context.Context, repoDID, collection, rkey, accessToken string) error { endpoint := fmt.Sprintf("%s/xrpc/com.atproto.repo.deleteRecord", strings.TrimSuffix(s.pdsURL, "/")) payload := map[string]interface{}{ @@ -798,11 +776,6 @@ func (s *communityService) deleteRecordOnPDSAs(ctx context.Context, repoDID, col return err } -func (s *communityService) callPDS(ctx context.Context, method, endpoint string, payload map[string]interface{}) (string, string, error) { - // Use instance's access token - return s.callPDSWithAuth(ctx, method, endpoint, payload, s.pdsAccessToken) -} - // callPDSWithAuth makes a PDS call with a specific access token (V2: for community authentication) func (s *communityService) callPDSWithAuth(ctx context.Context, method, endpoint string, payload map[string]interface{}, accessToken string) (string, string, error) { jsonData, err := json.Marshal(payload) @@ -870,4 +843,3 @@ func (s *communityService) callPDSWithAuth(ctx context.Context, method, endpoint } // Helper functions - diff --git a/internal/db/postgres/community_repo_blocks.go b/internal/db/postgres/community_repo_blocks.go index fbb800d..f1b172d 100644 --- a/internal/db/postgres/community_repo_blocks.go +++ b/internal/db/postgres/community_repo_blocks.go @@ -171,4 +171,3 @@ func (r *postgresCommunityRepo) IsBlocked(ctx context.Context, userDID, communit return exists, nil } - diff --git a/tests/integration/community_blocking_test.go b/tests/integration/community_blocking_test.go index 98db011..c7ca753 100644 --- a/tests/integration/community_blocking_test.go +++ b/tests/integration/community_blocking_test.go @@ -431,12 +431,12 @@ func createBlockingTestCommunityRepo(t *testing.T, db *sql.DB) communities.Repos func createBlockingTestCommunity(t *testing.T, repo communities.Repository, name, did string) *communities.Community { community := &communities.Community{ - DID: did, - Handle: fmt.Sprintf("!%s@coves.test", name), - Name: name, - DisplayName: fmt.Sprintf("Test Community %s", name), - Description: "Test community for blocking tests", - OwnerDID: did, + DID: did, + Handle: fmt.Sprintf("!%s@coves.test", name), + Name: name, + DisplayName: fmt.Sprintf("Test Community %s", name), + Description: "Test community for blocking tests", + OwnerDID: did, CreatedByDID: "did:plc:test-creator", HostedByDID: "did:plc:test-instance", Visibility: "public", diff --git a/tests/integration/community_e2e_test.go b/tests/integration/community_e2e_test.go index bc5e885..7ec14e1 100644 --- a/tests/integration/community_e2e_test.go +++ b/tests/integration/community_e2e_test.go @@ -687,7 +687,7 @@ func TestCommunity_E2E(t *testing.T) { CID: subscribeResp.CID, Record: map[string]interface{}{ "$type": "social.coves.community.subscription", - "subject": community.DID, + "subject": community.DID, "contentVisibility": float64(5), // JSON numbers are float64 "createdAt": time.Now().Format(time.RFC3339), }, @@ -771,7 +771,7 @@ func TestCommunity_E2E(t *testing.T) { CID: subscription.RecordCID, Record: map[string]interface{}{ "$type": "social.coves.community.subscription", - "subject": community.DID, + "subject": community.DID, "contentVisibility": float64(3), "createdAt": time.Now().Format(time.RFC3339), }, @@ -893,8 +893,8 @@ func TestCommunity_E2E(t *testing.T) { Operation: "delete", Collection: "social.coves.community.subscription", RKey: rkey, - CID: "", // No CID on deletes - Record: nil, // No record data on deletes + CID: "", // No CID on deletes + Record: nil, // No record data on deletes }, } if handleErr := consumer.HandleEvent(context.Background(), &deleteEvent); handleErr != nil { @@ -1505,7 +1505,6 @@ func createAndIndexCommunity(t *testing.T, service communities.Service, consumer return community } - // authenticateWithPDS authenticates with the PDS and returns access token and DID func authenticateWithPDS(pdsURL, handle, password string) (string, string, error) { // Call com.atproto.server.createSession diff --git a/tests/integration/subscription_indexing_test.go b/tests/integration/subscription_indexing_test.go index 123c411..1350fbf 100644 --- a/tests/integration/subscription_indexing_test.go +++ b/tests/integration/subscription_indexing_test.go @@ -46,8 +46,8 @@ func TestSubscriptionIndexing_ContentVisibility(t *testing.T) { RKey: rkey, CID: "bafytest123", Record: map[string]interface{}{ - "$type": "social.coves.community.subscription", - "subject": community.DID, + "$type": "social.coves.community.subscription", + "subject": community.DID, "createdAt": time.Now().Format(time.RFC3339), "contentVisibility": float64(5), // JSON numbers decode as float64 }, @@ -101,8 +101,8 @@ func TestSubscriptionIndexing_ContentVisibility(t *testing.T) { RKey: rkey, CID: "bafydefault", Record: map[string]interface{}{ - "$type": "social.coves.community.subscription", - "subject": community.DID, + "$type": "social.coves.community.subscription", + "subject": community.DID, "createdAt": time.Now().Format(time.RFC3339), // contentVisibility NOT provided }, @@ -130,9 +130,9 @@ func TestSubscriptionIndexing_ContentVisibility(t *testing.T) { t.Run("clamps contentVisibility to valid range (1-5)", func(t *testing.T) { testCases := []struct { + name string input float64 expected int - name string }{ {input: 0, expected: 1, name: "zero clamped to 1"}, {input: -5, expected: 1, name: "negative clamped to 1"}, @@ -201,8 +201,8 @@ func TestSubscriptionIndexing_ContentVisibility(t *testing.T) { RKey: rkey, CID: "bafyidempotent", Record: map[string]interface{}{ - "$type": "social.coves.community.subscription", - "subject": community.DID, + "$type": "social.coves.community.subscription", + "subject": community.DID, "createdAt": time.Now().Format(time.RFC3339), "contentVisibility": float64(4), }, @@ -268,8 +268,8 @@ func TestSubscriptionIndexing_DeleteOperations(t *testing.T) { RKey: rkey, CID: "bafycreate", Record: map[string]interface{}{ - "$type": "social.coves.community.subscription", - "subject": community.DID, + "$type": "social.coves.community.subscription", + "subject": community.DID, "createdAt": time.Now().Format(time.RFC3339), "contentVisibility": float64(3), }, @@ -298,7 +298,7 @@ func TestSubscriptionIndexing_DeleteOperations(t *testing.T) { Operation: "delete", Collection: "social.coves.community.subscription", RKey: rkey, - CID: "", // No CID on deletes + CID: "", // No CID on deletes Record: nil, // No record data on deletes }, } @@ -390,8 +390,8 @@ func TestSubscriptionIndexing_SubscriberCount(t *testing.T) { RKey: rkey, CID: "bafycount", Record: map[string]interface{}{ - "$type": "social.coves.community.subscription", - "subject": community.DID, + "$type": "social.coves.community.subscription", + "subject": community.DID, "createdAt": time.Now().Format(time.RFC3339), "contentVisibility": float64(3), },