diff --git a/tests/integration/community_e2e_test.go b/tests/integration/community_e2e_test.go index e89c78d..cd08baf 100644 --- a/tests/integration/community_e2e_test.go +++ b/tests/integration/community_e2e_test.go @@ -187,6 +187,9 @@ func TestCommunity_E2E(t *testing.T) { } t.Logf("\n📝 Creating community via service: %s", communityName) + // Capture a Jetstream replay cursor BEFORE the write so the Part 2 subscription + // (opened after the write) cannot miss the resulting firehose commit. + communityCreateCursor := jetstreamCursorNow() community, err := communityService.CreateCommunity(ctx, createReq) if err != nil { t.Fatalf("Failed to create community: %v", err) @@ -305,7 +308,7 @@ func TestCommunity_E2E(t *testing.T) { // Start Jetstream consumer in background go func() { - err := subscribeToJetstream(ctx, jetstreamURL, community.DID, consumer, eventChan, errorChan, done) + err := subscribeToJetstream(ctx, withJetstreamCursor(jetstreamURL, communityCreateCursor), community.DID, consumer, eventChan, errorChan, done) if err != nil { errorChan <- err } diff --git a/tests/integration/community_service_integration_test.go b/tests/integration/community_service_integration_test.go index 271d1a4..4454ba6 100644 --- a/tests/integration/community_service_integration_test.go +++ b/tests/integration/community_service_integration_test.go @@ -69,9 +69,8 @@ func TestCommunityService_CreateWithRealPDS(t *testing.T) { // Generate unique community name (keep short for DNS label limit) // Must start with letter, can contain alphanumeric and hyphens - // Use full Unix seconds + nanoseconds remainder for better uniqueness across runs - now := time.Now() - uniqueName := fmt.Sprintf("svc%d%d", now.Unix()%100000, now.UnixNano()%10000) + // Collision-free across runs and within a run (base36 seconds + atomic counter) + uniqueName := fmt.Sprintf("svc%s", uniqueTestID()) // Create community via service (FULL PRODUCTION CODE PATH) t.Logf("Creating community via service.CreateCommunity()...") @@ -321,9 +320,8 @@ func TestCommunityService_UpdateWithRealPDS(t *testing.T) { t.Run("updates community with real PDS", func(t *testing.T) { // First, create a community - // Use full Unix seconds + nanoseconds remainder for better uniqueness across runs - now := time.Now() - uniqueName := fmt.Sprintf("upd%d%d", now.Unix()%100000, now.UnixNano()%10000) + // Collision-free across runs and within a run (base36 seconds + atomic counter) + uniqueName := fmt.Sprintf("upd%s", uniqueTestID()) creatorDID := "did:plc:updatetestuser" t.Logf("Creating community to update...") @@ -393,9 +391,8 @@ func TestCommunityService_UpdateWithRealPDS(t *testing.T) { t.Run("rejects unauthorized updates", func(t *testing.T) { // Create a community - // Use full Unix seconds + nanoseconds remainder for better uniqueness across runs - now := time.Now() - uniqueName := fmt.Sprintf("auth%d%d", now.Unix()%100000, now.UnixNano()%10000) + // Collision-free across runs and within a run (base36 seconds + atomic counter) + uniqueName := fmt.Sprintf("auth%s", uniqueTestID()) creatorDID := "did:plc:creator123" community, err := service.CreateCommunity(ctx, communities.CreateCommunityRequest{ @@ -518,9 +515,8 @@ func TestPasswordAuthentication(t *testing.T) { t.Run("generated password works for session creation", func(t *testing.T) { // Create a community with PDS-generated password - // Use full Unix seconds + nanoseconds remainder for better uniqueness across runs - now := time.Now() - uniqueName := fmt.Sprintf("pwd%d%d", now.Unix()%100000, now.UnixNano()%10000) + // Collision-free across runs and within a run (base36 seconds + atomic counter) + uniqueName := fmt.Sprintf("pwd%s", uniqueTestID()) t.Logf("Creating community with generated password...") community, err := service.CreateCommunity(ctx, communities.CreateCommunityRequest{ diff --git a/tests/integration/helpers.go b/tests/integration/helpers.go index 27cb82d..98ddf56 100644 --- a/tests/integration/helpers.go +++ b/tests/integration/helpers.go @@ -149,6 +149,33 @@ func uniqueTestID() string { return strconv.FormatInt(time.Now().Unix(), 36) + strconv.FormatUint(n, 36) } +// jetstreamCursorNow returns the current time as a Jetstream replay cursor (unix +// microseconds). Capture it IMMEDIATELY BEFORE a PDS write, then pass it through +// withJetstreamCursor when opening the subscription afterwards. +// +// Why: the firehose subscriptions used in these tests are otherwise cursorless, so +// they only stream commits emitted after the socket is dialed — there is no replay. +// A subscription opened after the write therefore races the PDS→Jetstream relay and +// silently drops the event under load (the "subscribe-after-write" flake). Jetstream +// stamps each event's time_us when it ingests the commit, which is always after our +// write, so a cursor captured just before the write is guaranteed to be < the event's +// time_us (we receive it) and > any earlier event's time_us (no stale duplicates). +// Test and Jetstream share the same host clock, so there is no skew to compensate for. +func jetstreamCursorNow() int64 { + return time.Now().UnixMicro() +} + +// withJetstreamCursor appends a replay cursor (unix microseconds, from +// jetstreamCursorNow) to a Jetstream subscribe URL, handling whether the URL already +// carries query parameters (e.g. wantedCollections). +func withJetstreamCursor(baseURL string, cursorMicros int64) string { + sep := "?" + if strings.Contains(baseURL, "?") { + sep = "&" + } + return fmt.Sprintf("%s%scursor=%d", baseURL, sep, cursorMicros) +} + // createPDSAccount creates a new account on PDS and returns access token + DID // This is used for E2E tests that need real PDS accounts func createPDSAccount(pdsURL, handle, email, password string) (accessToken, did string, err error) { diff --git a/tests/integration/post_e2e_test.go b/tests/integration/post_e2e_test.go index fa652a6..8261937 100644 --- a/tests/integration/post_e2e_test.go +++ b/tests/integration/post_e2e_test.go @@ -483,7 +483,10 @@ func TestPostCreation_E2E_LivePDS(t *testing.T) { token := e2eAuth.AddUser(author.DID) req.Header.Set("Authorization", "Bearer "+token) - // Execute request through auth middleware + handler + // Execute request through auth middleware + handler. + // Capture a Jetstream replay cursor BEFORE the write so the Part 2 subscription + // (opened after the write) cannot miss the resulting firehose commit. + postCreateCursor := jetstreamCursorNow() rr := httptest.NewRecorder() handler := e2eAuth.RequireAuth(http.HandlerFunc(createHandler.HandleCreate)) handler.ServeHTTP(rr, req) @@ -543,7 +546,7 @@ func TestPostCreation_E2E_LivePDS(t *testing.T) { // Start Jetstream WebSocket subscriber in background // This creates its own WebSocket connection to Jetstream go func() { - err := subscribeToJetstreamForPost(ctx, jetstreamURL, community.DID, postConsumer, eventChan, errorChan, done) + err := subscribeToJetstreamForPost(ctx, withJetstreamCursor(jetstreamURL, postCreateCursor), community.DID, postConsumer, eventChan, errorChan, done) if err != nil { errorChan <- err } diff --git a/tests/integration/user_journey_e2e_test.go b/tests/integration/user_journey_e2e_test.go index 55282a4..dfc683b 100644 --- a/tests/integration/user_journey_e2e_test.go +++ b/tests/integration/user_journey_e2e_test.go @@ -149,8 +149,10 @@ func TestFullUserJourney_E2E(t *testing.T) { httpServer := httptest.NewServer(r) defer httpServer.Close() - // Cleanup test data from previous runs (clean up ALL journey test data) - timestamp := time.Now().Unix() + // Cleanup test data from previous runs (clean up ALL journey test data). + // A single collision-free testID is shared by every handle/community name and the + // deferred cleanup patterns below (uniqueTestID stays short for PDS handle limits). + testID := uniqueTestID() // Clean up previous test runs - use pattern that matches journey test data // Handles are now shorter: alice{4-digit}.local.coves.dev, bob{4-digit}.local.coves.dev _, _ = db.Exec("DELETE FROM votes WHERE voter_did LIKE '%alice%.local.coves.dev%' OR voter_did LIKE '%bob%.local.coves.dev%'") @@ -162,10 +164,9 @@ func TestFullUserJourney_E2E(t *testing.T) { // Defer cleanup for current test run using specific timestamp pattern defer func() { - shortTS := timestamp % 10000 - alicePattern := fmt.Sprintf("%%alice%d%%", shortTS) - bobPattern := fmt.Sprintf("%%bob%d%%", shortTS) - gjPattern := fmt.Sprintf("%%gj%d%%", shortTS) + alicePattern := fmt.Sprintf("%%alice%s%%", testID) + bobPattern := fmt.Sprintf("%%bob%s%%", testID) + gjPattern := fmt.Sprintf("%%gj%s%%", testID) _, _ = db.Exec("DELETE FROM votes WHERE voter_did LIKE $1 OR voter_did LIKE $2", alicePattern, bobPattern) _, _ = db.Exec("DELETE FROM comments WHERE commenter_did LIKE $1 OR commenter_did LIKE $2", alicePattern, bobPattern) _, _ = db.Exec("DELETE FROM posts WHERE community_did LIKE $1", gjPattern) @@ -199,9 +200,8 @@ func TestFullUserJourney_E2E(t *testing.T) { t.Log("\n👤 Part 1: User A creates account and authenticates...") // Use short handle format to stay under PDS 34-char limit - shortTS := timestamp % 10000 // Use last 4 digits - userAHandle = fmt.Sprintf("alice%d.local.coves.dev", shortTS) - email := fmt.Sprintf("alice%d@test.com", shortTS) + userAHandle = fmt.Sprintf("alice%s.local.coves.dev", testID) + email := fmt.Sprintf("alice%s@test.com", testID) password := "test-password-alice-123" // Create account on PDS @@ -231,8 +231,7 @@ func TestFullUserJourney_E2E(t *testing.T) { // Community handle will be {name}c-{name}.coves.social // Max 34 chars total, so name must be short (34 - 23 = 11 chars max) - shortTS := timestamp % 10000 - communityName := fmt.Sprintf("gj%d", shortTS) // "gj9261" = 6 chars -> handle = 29 chars + communityName := fmt.Sprintf("gj%s", testID) // short prefix + base36 testID keeps the community handle well under limits createReq := map[string]interface{}{ "name": communityName, @@ -249,6 +248,9 @@ func TestFullUserJourney_E2E(t *testing.T) { req.Header.Set("Content-Type", "application/json") req.Header.Set("Authorization", "Bearer "+userAAPIToken) + // Capture a Jetstream replay cursor BEFORE the write so the subscription opened + // afterwards cannot miss the resulting firehose commit (subscribe-after-write race). + communityCreateCursor := jetstreamCursorNow() resp, err := http.DefaultClient.Do(req) require.NoError(t, err) defer func() { _ = resp.Body.Close() }() @@ -279,7 +281,7 @@ func TestFullUserJourney_E2E(t *testing.T) { jetstreamFilterURL := fmt.Sprintf("%s?wantedCollections=social.coves.community.profile", jetstreamURL) go func() { - err := subscribeToJetstreamForCommunity(ctx, jetstreamFilterURL, communityDID, communityConsumer, eventChan, errorChan, done) + err := subscribeToJetstreamForCommunity(ctx, withJetstreamCursor(jetstreamFilterURL, communityCreateCursor), communityDID, communityConsumer, eventChan, errorChan, done) if err != nil { errorChan <- err } @@ -333,6 +335,8 @@ func TestFullUserJourney_E2E(t *testing.T) { req.Header.Set("Content-Type", "application/json") req.Header.Set("Authorization", "Bearer "+userAAPIToken) + // Capture a Jetstream replay cursor BEFORE the write (subscribe-after-write race). + postCreateCursor := jetstreamCursorNow() resp, err := http.DefaultClient.Do(req) require.NoError(t, err) defer func() { _ = resp.Body.Close() }() @@ -357,7 +361,7 @@ func TestFullUserJourney_E2E(t *testing.T) { jetstreamFilterURL := fmt.Sprintf("%s?wantedCollections=social.coves.community.post", jetstreamURL) go func() { - err := subscribeToJetstreamForPost(ctx, jetstreamFilterURL, communityDID, postConsumer, eventChan, errorChan, done) + err := subscribeToJetstreamForPost(ctx, withJetstreamCursor(jetstreamFilterURL, postCreateCursor), communityDID, postConsumer, eventChan, errorChan, done) if err != nil { errorChan <- err } @@ -399,9 +403,8 @@ func TestFullUserJourney_E2E(t *testing.T) { t.Log("\n👤 Part 4: User B creates account and authenticates...") // Use short handle format to stay under PDS 34-char limit - shortTS := timestamp % 10000 // Use last 4 digits - userBHandle = fmt.Sprintf("bob%d.local.coves.dev", shortTS) - email := fmt.Sprintf("bob%d@test.com", shortTS) + userBHandle = fmt.Sprintf("bob%s.local.coves.dev", testID) + email := fmt.Sprintf("bob%s@test.com", testID) password := "test-password-bob-123" // Create account on PDS diff --git a/tests/integration/user_profile_avatar_e2e_test.go b/tests/integration/user_profile_avatar_e2e_test.go index 0f476db..981d5ad 100644 --- a/tests/integration/user_profile_avatar_e2e_test.go +++ b/tests/integration/user_profile_avatar_e2e_test.go @@ -129,14 +129,13 @@ func TestUserProfileAvatarE2E_UpdateWithAvatar(t *testing.T) { defer httpServer.Close() // Cleanup old test data - timestamp := time.Now().Unix() - shortTS := timestamp % 10000 + testID := uniqueTestID() _, _ = db.Exec("DELETE FROM users WHERE handle LIKE 'avatartest%.local.coves.dev'") t.Run("update profile with avatar via real PDS and Jetstream", func(t *testing.T) { // Create test user account on PDS - userHandle := fmt.Sprintf("avatartest%d.local.coves.dev", shortTS) - email := fmt.Sprintf("avatartest%d@test.com", shortTS) + userHandle := fmt.Sprintf("avatartest%s.local.coves.dev", testID) + email := fmt.Sprintf("avatartest%s@test.com", testID) password := "test-password-avatar-123" t.Logf("\n Creating test user account on PDS: %s", userHandle) @@ -423,14 +422,13 @@ func TestUserProfileAvatarE2E_UpdateWithBanner(t *testing.T) { httpServer := httptest.NewServer(r) defer httpServer.Close() - timestamp := time.Now().Unix() - shortTS := timestamp % 10000 + testID := uniqueTestID() _, _ = db.Exec("DELETE FROM users WHERE handle LIKE 'bannertest%.local.coves.dev'") t.Run("update profile with banner via real PDS and Jetstream", func(t *testing.T) { // Create test user account on PDS - userHandle := fmt.Sprintf("bannertest%d.local.coves.dev", shortTS) - email := fmt.Sprintf("bannertest%d@test.com", shortTS) + userHandle := fmt.Sprintf("bannertest%s.local.coves.dev", testID) + email := fmt.Sprintf("bannertest%s@test.com", testID) password := "test-password-banner-123" t.Logf("\n Creating test user account on PDS: %s", userHandle) @@ -668,13 +666,12 @@ func TestUserProfileAvatarE2E_UpdateDisplayNameAndBio(t *testing.T) { httpServer := httptest.NewServer(r) defer httpServer.Close() - timestamp := time.Now().Unix() - shortTS := timestamp % 10000 + testID := uniqueTestID() t.Run("update display name and bio without blobs", func(t *testing.T) { // Create test user account on PDS - userHandle := fmt.Sprintf("texttest%d.local.coves.dev", shortTS) - email := fmt.Sprintf("texttest%d@test.com", shortTS) + userHandle := fmt.Sprintf("texttest%s.local.coves.dev", testID) + email := fmt.Sprintf("texttest%s@test.com", testID) password := "test-password-text-123" userToken, userDID, err := createPDSAccount(pdsURL, userHandle, email, password) @@ -873,22 +870,31 @@ func TestUserProfileAvatarE2E_ReplaceAvatar(t *testing.T) { httpServer := httptest.NewServer(r) defer httpServer.Close() - timestamp := time.Now().Unix() - shortTS := timestamp % 10000 - - // Helper to wait for Jetstream event and extract avatar CID - waitForProfileEvent := func(t *testing.T, userDID string, timeout time.Duration) (string, *jetstream.JetstreamEvent) { + testID := uniqueTestID() + + // subscribeForProfileEvent opens the Jetstream subscription and returns a wait + // function that blocks until a profile commit for userDID arrives (or times out), + // returning the avatar CID extracted from the commit record. + // + // It MUST be called BEFORE the PDS write. The firehose subscription is cursorless + // (see jetstreamURL), so it only streams commits emitted after the socket is + // established — there is no replay. Dialing after the write (the previous helper's + // behavior) races the PDS→firehose relay and silently drops the event under load. + subscribeForProfileEvent := func(t *testing.T, userDID string, timeout time.Duration) func() (string, *jetstream.JetstreamEvent) { eventChan := make(chan *jetstream.JetstreamEvent, 10) done := make(chan bool) + ready := make(chan struct{}) subscribeCtx, cancelSubscribe := context.WithTimeout(ctx, timeout) - defer cancelSubscribe() go func() { conn, _, dialErr := websocket.DefaultDialer.Dial(jetstreamURL, nil) if dialErr != nil { + t.Logf("Failed to connect to Jetstream: %v", dialErr) + close(ready) return } defer func() { _ = conn.Close() }() + close(ready) // socket dialed; safe for the caller to write consecutiveTimeouts := 0 for { @@ -924,30 +930,39 @@ func TestUserProfileAvatarE2E_ReplaceAvatar(t *testing.T) { } }() - select { - case event := <-eventChan: - close(done) - var avatarCID string - if event.Commit.Record != nil { - if avatarMap, ok := event.Commit.Record["avatar"].(map[string]interface{}); ok { - if ref, ok := avatarMap["ref"].(map[string]interface{}); ok { - if link, ok := ref["$link"].(string); ok { - avatarCID = link + // Block until the socket is dialed, then give Jetstream a moment to register + // the subscription, so the caller's subsequent write is guaranteed to land + // after we are listening. + <-ready + time.Sleep(500 * time.Millisecond) + + return func() (string, *jetstream.JetstreamEvent) { + defer cancelSubscribe() + select { + case event := <-eventChan: + close(done) + var avatarCID string + if event.Commit.Record != nil { + if avatarMap, ok := event.Commit.Record["avatar"].(map[string]interface{}); ok { + if ref, ok := avatarMap["ref"].(map[string]interface{}); ok { + if link, ok := ref["$link"].(string); ok { + avatarCID = link + } } } } + return avatarCID, event + case <-time.After(timeout): + close(done) + return "", nil } - return avatarCID, event - case <-time.After(timeout): - close(done) - return "", nil } } t.Run("replace existing avatar with new one", func(t *testing.T) { // Create test user account on PDS - userHandle := fmt.Sprintf("replaceav%d.local.coves.dev", shortTS) - email := fmt.Sprintf("replaceav%d@test.com", shortTS) + userHandle := fmt.Sprintf("replaceav%s.local.coves.dev", testID) + email := fmt.Sprintf("replaceav%s@test.com", testID) password := "test-password-replace-123" userToken, userDID, err := createPDSAccount(pdsURL, userHandle, email, password) @@ -970,10 +985,8 @@ func TestUserProfileAvatarE2E_ReplaceAvatar(t *testing.T) { AvatarMimeType: "image/png", } - // Start listening before update - go func() { - time.Sleep(500 * time.Millisecond) - }() + // Subscribe to the firehose BEFORE the write (cursorless: no replay). + waitInitial := subscribeForProfileEvent(t, userDID, 30*time.Second) reqBody, _ := json.Marshal(updateReq) req, _ := http.NewRequest(http.MethodPost, @@ -988,7 +1001,7 @@ func TestUserProfileAvatarE2E_ReplaceAvatar(t *testing.T) { require.Equal(t, http.StatusOK, resp.StatusCode) // Wait for initial avatar event - initialAvatarCID, initialEvent := waitForProfileEvent(t, userDID, 15*time.Second) + initialAvatarCID, initialEvent := waitInitial() require.NotNil(t, initialEvent, "Should receive initial avatar event") require.NotEmpty(t, initialAvatarCID, "Initial avatar CID should not be empty") @@ -1017,6 +1030,9 @@ func TestUserProfileAvatarE2E_ReplaceAvatar(t *testing.T) { AvatarMimeType: "image/png", } + // Subscribe BEFORE the replacement write (cursorless: no replay). + waitReplacement := subscribeForProfileEvent(t, userDID, 30*time.Second) + reqBody2, _ := json.Marshal(updateReq2) req2, _ := http.NewRequest(http.MethodPost, httpServer.URL+"/xrpc/social.coves.actor.updateProfile", @@ -1030,7 +1046,7 @@ func TestUserProfileAvatarE2E_ReplaceAvatar(t *testing.T) { require.Equal(t, http.StatusOK, resp2.StatusCode) // Wait for replacement avatar event - newAvatarCID, newEvent := waitForProfileEvent(t, userDID, 15*time.Second) + newAvatarCID, newEvent := waitReplacement() require.NotNil(t, newEvent, "Should receive replacement avatar event") require.NotEmpty(t, newAvatarCID, "New avatar CID should not be empty")