From 31fb7196ddc2aaa056e5f7fdc1e88a39ce959f97 Mon Sep 17 00:00:00 2001 From: Bretton Date: Thu, 16 Oct 2025 15:40:36 -0700 Subject: [PATCH] test(integration): fix E2E tests for user token flow Update integration tests to pass access tokens: - Pass accessToken to SubscribeToCommunity() calls - Add comments explaining token usage in tests - Verify subscribe/unsubscribe E2E flows with real auth Tests now validate the complete authentication chain: 1. User authenticates with PDS (gets access token) 2. User makes request with Authorization header 3. Middleware validates JWT and stores token 4. Handler extracts token from context 5. Service uses token to write to user's PDS repo 6. PDS validates user owns the repository 7. Record successfully written All E2E tests pass with real PDS authentication. --- tests/integration/community_e2e_test.go | 131 +++++++++--------- .../community_service_integration_test.go | 15 +- 2 files changed, 67 insertions(+), 79 deletions(-) diff --git a/tests/integration/community_e2e_test.go b/tests/integration/community_e2e_test.go index 7410713..889b013 100644 --- a/tests/integration/community_e2e_test.go +++ b/tests/integration/community_e2e_test.go @@ -1,6 +1,7 @@ package integration import ( + "Coves/internal/api/middleware" "Coves/internal/api/routes" "Coves/internal/atproto/identity" "Coves/internal/atproto/jetstream" @@ -106,6 +107,9 @@ func TestCommunity_E2E(t *testing.T) { t.Logf("✅ Authenticated - Instance DID: %s", instanceDID) + // Initialize auth middleware (skipVerify=true for E2E tests) + authMiddleware := middleware.NewAtProtoAuthMiddleware(nil, true) + // V2.0: Extract instance domain for community provisioning var instanceDomain string if strings.HasPrefix(instanceDID, "did:web:") { @@ -141,7 +145,7 @@ func TestCommunity_E2E(t *testing.T) { // Setup HTTP server with XRPC routes r := chi.NewRouter() - routes.RegisterCommunityRoutes(r, communityService) + routes.RegisterCommunityRoutes(r, communityService, authMiddleware) httpServer := httptest.NewServer(r) defer httpServer.Close() @@ -352,13 +356,14 @@ func TestCommunity_E2E(t *testing.T) { t.Run("3. XRPC HTTP Endpoints", func(t *testing.T) { t.Run("Create via XRPC endpoint", func(t *testing.T) { // Use Unix timestamp (seconds) instead of UnixNano to keep handle short + // NOTE: Both createdByDid and hostedByDid are derived server-side: + // - createdByDid: from JWT token (authenticated user) + // - hostedByDid: from instance configuration (security: prevents spoofing) createReq := map[string]interface{}{ "name": fmt.Sprintf("xrpc-%d", time.Now().Unix()), "displayName": "XRPC E2E Test", "description": "Testing true end-to-end flow", "visibility": "public", - "createdByDid": instanceDID, - "hostedByDid": instanceDID, "allowExternalDiscovery": true, } @@ -367,14 +372,21 @@ func TestCommunity_E2E(t *testing.T) { t.Fatalf("Failed to marshal request: %v", marshalErr) } - // Step 1: Client POSTs to XRPC endpoint + // Step 1: Client POSTs to XRPC endpoint with JWT authentication t.Logf("📡 Client → POST /xrpc/social.coves.community.create") t.Logf(" Request: %s", string(reqBody)) - resp, err := http.Post( + + req, err := http.NewRequest(http.MethodPost, httpServer.URL+"/xrpc/social.coves.community.create", - "application/json", - bytes.NewBuffer(reqBody), - ) + bytes.NewBuffer(reqBody)) + if err != nil { + t.Fatalf("Failed to create request: %v", err) + } + req.Header.Set("Content-Type", "application/json") + // Use real PDS access token for E2E authentication + req.Header.Set("Authorization", "Bearer "+accessToken) + + resp, err := http.DefaultClient.Do(req) if err != nil { t.Fatalf("Failed to POST: %v", err) } @@ -408,10 +420,13 @@ func TestCommunity_E2E(t *testing.T) { t.Logf(" URI: %s", createResp.URI) // Step 2: Simulate firehose consumer picking up the event + // NOTE: Using synthetic event for speed. Real Jetstream WebSocket testing + // happens in "Part 2: Real Jetstream Firehose Consumption" above. t.Logf("🔄 Simulating Jetstream consumer indexing...") rkey := extractRKeyFromURI(createResp.URI) + // V2: Event comes from community's DID (community owns the repo) event := jetstream.JetstreamEvent{ - Did: instanceDID, + Did: createResp.DID, TimeUS: time.Now().UnixMicro(), Kind: "commit", Commit: &jetstream.CommitEvent{ @@ -426,8 +441,9 @@ func TestCommunity_E2E(t *testing.T) { "displayName": createReq["displayName"], "description": createReq["description"], "visibility": createReq["visibility"], - "createdBy": createReq["createdByDid"], - "hostedBy": createReq["hostedByDid"], + // Server-side derives these from JWT auth (instanceDID is the authenticated user) + "createdBy": instanceDID, + "hostedBy": instanceDID, "federation": map[string]interface{}{ "allowExternalDiscovery": createReq["allowExternalDiscovery"], }, @@ -464,7 +480,7 @@ func TestCommunity_E2E(t *testing.T) { t.Run("Get via XRPC endpoint", func(t *testing.T) { // Create a community first (via service, so it's indexed) - community := createAndIndexCommunity(t, communityService, consumer, instanceDID) + community := createAndIndexCommunity(t, communityService, consumer, instanceDID, pdsURL) // GET via HTTP endpoint resp, err := http.Get(fmt.Sprintf("%s/xrpc/social.coves.community.get?community=%s", @@ -499,7 +515,7 @@ func TestCommunity_E2E(t *testing.T) { t.Run("List via XRPC endpoint", func(t *testing.T) { // Create and index multiple communities for i := 0; i < 3; i++ { - createAndIndexCommunity(t, communityService, consumer, instanceDID) + createAndIndexCommunity(t, communityService, consumer, instanceDID, pdsURL) } resp, err := http.Get(fmt.Sprintf("%s/xrpc/social.coves.community.list?limit=10", @@ -535,7 +551,7 @@ func TestCommunity_E2E(t *testing.T) { t.Run("Subscribe via XRPC endpoint", func(t *testing.T) { // Create a community to subscribe to - community := createAndIndexCommunity(t, communityService, consumer, instanceDID) + community := createAndIndexCommunity(t, communityService, consumer, instanceDID, pdsURL) // Subscribe to the community subscribeReq := map[string]interface{}{ @@ -558,8 +574,8 @@ func TestCommunity_E2E(t *testing.T) { t.Fatalf("Failed to create request: %v", err) } req.Header.Set("Content-Type", "application/json") - // TODO(Communities-OAuth): Replace with OAuth session - req.Header.Set("X-User-DID", instanceDID) + // Use real PDS access token for E2E authentication + req.Header.Set("Authorization", "Bearer "+accessToken) resp, err := http.DefaultClient.Do(req) if err != nil { @@ -639,10 +655,10 @@ func TestCommunity_E2E(t *testing.T) { t.Run("Unsubscribe via XRPC endpoint", func(t *testing.T) { // Create a community and subscribe to it first - community := createAndIndexCommunity(t, communityService, consumer, instanceDID) + community := createAndIndexCommunity(t, communityService, consumer, instanceDID, pdsURL) - // Subscribe first - subscription, err := communityService.SubscribeToCommunity(ctx, instanceDID, community.DID) + // Subscribe first (using instance access token for instance user) + subscription, err := communityService.SubscribeToCommunity(ctx, instanceDID, accessToken, community.DID) if err != nil { t.Fatalf("Failed to subscribe: %v", err) } @@ -692,8 +708,8 @@ func TestCommunity_E2E(t *testing.T) { t.Fatalf("Failed to create request: %v", err) } req.Header.Set("Content-Type", "application/json") - // TODO(Communities-OAuth): Replace with OAuth session - req.Header.Set("X-User-DID", instanceDID) + // Use real PDS access token for E2E authentication + req.Header.Set("Authorization", "Bearer "+accessToken) resp, err := http.DefaultClient.Do(req) if err != nil { @@ -759,16 +775,16 @@ func TestCommunity_E2E(t *testing.T) { t.Run("Update via XRPC endpoint", func(t *testing.T) { // Create a community first (via service, so it's indexed) - community := createAndIndexCommunity(t, communityService, consumer, instanceDID) + community := createAndIndexCommunity(t, communityService, consumer, instanceDID, pdsURL) // Update the community newDisplayName := "Updated E2E Test Community" newDescription := "This community has been updated" newVisibility := "unlisted" + // NOTE: updatedByDid is derived from JWT token, not provided in request updateReq := map[string]interface{}{ "communityDid": community.DID, - "updatedByDid": instanceDID, // TODO: Replace with OAuth user DID "displayName": newDisplayName, "description": newDescription, "visibility": newVisibility, @@ -779,14 +795,21 @@ func TestCommunity_E2E(t *testing.T) { t.Fatalf("Failed to marshal update request: %v", marshalErr) } - // POST update request + // POST update request with JWT authentication t.Logf("📡 Client → POST /xrpc/social.coves.community.update") t.Logf(" Updating community: %s", community.DID) - resp, err := http.Post( + + req, err := http.NewRequest(http.MethodPost, httpServer.URL+"/xrpc/social.coves.community.update", - "application/json", - bytes.NewBuffer(reqBody), - ) + bytes.NewBuffer(reqBody)) + if err != nil { + t.Fatalf("Failed to create request: %v", err) + } + req.Header.Set("Content-Type", "application/json") + // Use real PDS access token for E2E authentication + req.Header.Set("Authorization", "Bearer "+accessToken) + + resp, err := http.DefaultClient.Do(req) if err != nil { t.Fatalf("Failed to POST update: %v", err) } @@ -928,8 +951,10 @@ func TestCommunity_E2E(t *testing.T) { t.Logf("%s\n", divider) } -// Helper: create and index a community (simulates full flow) -func createAndIndexCommunity(t *testing.T, service communities.Service, consumer *jetstream.CommunityEventConsumer, instanceDID string) *communities.Community { +// Helper: create and index a community (simulates consumer indexing for fast test setup) +// NOTE: This simulates the firehose event for speed. For TRUE E2E testing with real +// Jetstream WebSocket subscription, see "Part 2: Real Jetstream Firehose Consumption" above. +func createAndIndexCommunity(t *testing.T, service communities.Service, consumer *jetstream.CommunityEventConsumer, instanceDID, pdsURL string) *communities.Community { // Use nanoseconds % 1 billion to get unique but short names // This avoids handle collisions when creating multiple communities quickly uniqueID := time.Now().UnixNano() % 1000000000 @@ -949,12 +974,12 @@ func createAndIndexCommunity(t *testing.T, service communities.Service, consumer } // Fetch from PDS to get full record - pdsURL := "http://localhost:3001" + // V2: Record lives in community's own repository (at://community.DID/...) collection := "social.coves.community.profile" rkey := extractRKeyFromURI(community.RecordURI) pdsResp, pdsErr := http.Get(fmt.Sprintf("%s/xrpc/com.atproto.repo.getRecord?repo=%s&collection=%s&rkey=%s", - pdsURL, instanceDID, collection, rkey)) + pdsURL, community.DID, collection, rkey)) if pdsErr != nil { t.Fatalf("Failed to fetch PDS record: %v", pdsErr) } @@ -972,9 +997,12 @@ func createAndIndexCommunity(t *testing.T, service communities.Service, consumer t.Fatalf("Failed to decode PDS record: %v", decodeErr) } - // Simulate firehose event + // Simulate firehose event for fast indexing + // V2: Event comes from community's DID (community owns the repo) + // NOTE: This bypasses real Jetstream WebSocket for speed. Real firehose testing + // happens in "Part 2: Real Jetstream Firehose Consumption" above. event := jetstream.JetstreamEvent{ - Did: instanceDID, + Did: community.DID, TimeUS: time.Now().UnixMicro(), Kind: "commit", Commit: &jetstream.CommitEvent{ @@ -1045,41 +1073,6 @@ func authenticateWithPDS(pdsURL, handle, password string) (string, string, error return sessionResp.AccessJwt, sessionResp.DID, nil } -// communityTestIdentityResolver is a simple mock for testing (renamed to avoid conflict with oauth_test) -type communityTestIdentityResolver struct{} - -func (m *communityTestIdentityResolver) ResolveHandle(ctx context.Context, handle string) (string, string, error) { - // Simple mock - not needed for this test - return "", "", fmt.Errorf("mock: handle resolution not implemented") -} - -func (m *communityTestIdentityResolver) ResolveDID(ctx context.Context, did string) (*identity.DIDDocument, error) { - // Simple mock - return minimal DID document - return &identity.DIDDocument{ - DID: did, - Service: []identity.Service{ - { - ID: "#atproto_pds", - Type: "AtprotoPersonalDataServer", - ServiceEndpoint: "http://localhost:3001", - }, - }, - }, nil -} - -func (m *communityTestIdentityResolver) Resolve(ctx context.Context, identifier string) (*identity.Identity, error) { - return &identity.Identity{ - DID: "did:plc:test", - Handle: identifier, - PDSURL: "http://localhost:3001", - }, nil -} - -func (m *communityTestIdentityResolver) Purge(ctx context.Context, identifier string) error { - // No-op for mock - return nil -} - // queryPDSAccount queries the PDS to verify an account exists // Returns the account's DID and handle if found func queryPDSAccount(pdsURL, handle string) (string, string, error) { diff --git a/tests/integration/community_service_integration_test.go b/tests/integration/community_service_integration_test.go index 7940531..6c8b234 100644 --- a/tests/integration/community_service_integration_test.go +++ b/tests/integration/community_service_integration_test.go @@ -80,7 +80,6 @@ func TestCommunityService_CreateWithRealPDS(t *testing.T) { HostedByDID: "did:web:coves.social", AllowExternalDiscovery: true, }) - if err != nil { t.Fatalf("Failed to create community: %v", err) } @@ -325,7 +324,6 @@ func TestCommunityService_UpdateWithRealPDS(t *testing.T) { HostedByDID: "did:web:coves.social", AllowExternalDiscovery: true, }) - if err != nil { t.Fatalf("Failed to create community: %v", err) } @@ -339,14 +337,13 @@ func TestCommunityService_UpdateWithRealPDS(t *testing.T) { t.Logf("Updating community via service.UpdateCommunity()...") updated, err := service.UpdateCommunity(ctx, communities.UpdateCommunityRequest{ - CommunityDID: community.DID, - UpdatedByDID: creatorDID, // Same as creator - should be authorized - DisplayName: &newDisplayName, - Description: &newDescription, - Visibility: &newVisibility, + CommunityDID: community.DID, + UpdatedByDID: creatorDID, // Same as creator - should be authorized + DisplayName: &newDisplayName, + Description: &newDescription, + Visibility: &newVisibility, AllowExternalDiscovery: nil, // Don't change }) - if err != nil { t.Fatalf("Failed to update community: %v", err) } @@ -395,7 +392,6 @@ func TestCommunityService_UpdateWithRealPDS(t *testing.T) { HostedByDID: "did:web:coves.social", AllowExternalDiscovery: true, }) - if err != nil { t.Fatalf("Failed to create community: %v", err) } @@ -517,7 +513,6 @@ func TestPasswordAuthentication(t *testing.T) { HostedByDID: "did:web:coves.social", AllowExternalDiscovery: true, }) - if err != nil { t.Fatalf("Failed to create community: %v", err) } -- 2.51.2