Files
mixmaker/backend/internal/adapter/postgres/store_integration_test.go
lemintare 6b189bfce4
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
CI / compose (push) Has been cancelled
Add global role synchronization for Discord with configurable interval
This commit introduces a new feature for global role synchronization in Discord, allowing for periodic reconciliation of managed roles. A new environment variable, `DISCORD_GLOBAL_SYNC_INTERVAL`, has been added to configure the synchronization interval, defaulting to 5 minutes. The `RoleWorker` has been updated to schedule global sync jobs, ensuring that missing managed roles are restored and extra assignments are removed without affecting unrelated server roles. Database schema changes support the new synchronization logic, and tests have been added to validate the functionality of the global reconciliation process.
2026-07-19 12:04:33 +03:00

431 lines
15 KiB
Go

package postgres
import (
"context"
"fmt"
"os"
"testing"
"time"
"mixmaker/backend/internal/application"
"mixmaker/backend/internal/domain"
)
type discardPublisher struct{}
func (discardPublisher) Publish(string, any) {}
func TestMigrationsAndReadiness(t *testing.T) {
databaseURL := os.Getenv("TEST_DATABASE_URL")
if databaseURL == "" {
t.Skip("TEST_DATABASE_URL is not configured")
}
ctx := context.Background()
if err := Migrate(ctx, databaseURL, "../../../migrations"); err != nil {
t.Fatal(err)
}
store, err := Open(ctx, databaseURL)
if err != nil {
t.Fatal(err)
}
defer store.Close()
if err := store.Ready(ctx); err != nil {
t.Fatal(err)
}
var workflowColumns int
if err := store.pool.QueryRow(ctx, `SELECT count(*) FROM information_schema.columns
WHERE table_name='events' AND column_name IN ('state','version','ruleset_id','active_series_id','tournament_id')`).Scan(&workflowColumns); err != nil {
t.Fatal(err)
}
if workflowColumns != 5 {
t.Fatalf("workflow migration is incomplete: found %d event columns", workflowColumns)
}
var rosterTable bool
if err := store.pool.QueryRow(ctx, `SELECT to_regclass('public.event_rosters') IS NOT NULL`).Scan(&rosterTable); err != nil {
t.Fatal(err)
}
if !rosterTable {
t.Fatal("event_rosters table is missing")
}
var bracketTable bool
if err := store.pool.QueryRow(ctx, `SELECT to_regclass('public.event_bracket_drafts') IS NOT NULL`).Scan(&bracketTable); err != nil {
t.Fatal(err)
}
if !bracketTable {
t.Fatal("event_bracket_drafts table is missing")
}
for _, table := range []string{"discord_managed_roles", "discord_role_assignments", "discord_role_sync_jobs"} {
var exists bool
if err := store.pool.QueryRow(ctx, `SELECT to_regclass('public.' || $1) IS NOT NULL`, table).Scan(&exists); err != nil {
t.Fatal(err)
}
if !exists {
t.Fatalf("%s table is missing", table)
}
}
}
func TestDiscordRoleSyncPersistence(t *testing.T) {
databaseURL := os.Getenv("TEST_DATABASE_URL")
if databaseURL == "" {
t.Skip("TEST_DATABASE_URL is not configured")
}
ctx := context.Background()
if err := Migrate(ctx, databaseURL, "../../../migrations"); err != nil {
t.Fatal(err)
}
store, err := Open(ctx, databaseURL)
if err != nil {
t.Fatal(err)
}
defer store.Close()
eventID := fmt.Sprintf("discord-sync-%d", time.Now().UnixNano())
role := application.DiscordManagedRole{
Scope: application.DiscordRoleScopeEvent, EventID: eventID, TeamID: "team-1",
Kind: application.DiscordRoleKindTeam, DiscordRoleID: eventID + "-role", RoleName: "Alpha",
}
defer func() {
_, _ = store.pool.Exec(ctx, `DELETE FROM discord_role_sync_jobs WHERE event_id=$1`, eventID)
_, _ = store.pool.Exec(ctx, `DELETE FROM discord_managed_roles WHERE event_id=$1`, eventID)
}()
if err = store.UpsertDiscordManagedRole(ctx, role); err != nil {
t.Fatal(err)
}
if err = store.UpsertDiscordRoleAssignment(ctx, role.DiscordRoleID, "discord-user"); err != nil {
t.Fatal(err)
}
assignments, err := store.ListDiscordRoleAssignments(ctx, role.DiscordRoleID)
if err != nil || len(assignments) != 1 || assignments[0] != "discord-user" {
t.Fatalf("unexpected assignments: %v err=%v", assignments, err)
}
if _, err = store.pool.Exec(ctx, `UPDATE discord_role_sync_jobs SET status='completed' WHERE status='pending'`); err != nil {
t.Fatal(err)
}
if err = store.EnqueueDiscordRoleSync(ctx, eventID, application.DiscordRoleActionReconcile, 1); err != nil {
t.Fatal(err)
}
job, err := store.ClaimDiscordRoleSyncJob(ctx)
if err != nil || job.EventID != eventID || len(job.RoleSnapshot) != 1 || job.Attempts != 1 {
t.Fatalf("unexpected claimed job: %+v err=%v", job, err)
}
if err = store.RetryDiscordRoleSyncJob(ctx, job.ID, "retry", time.Now().Add(-time.Second)); err != nil {
t.Fatal(err)
}
job, err = store.ClaimDiscordRoleSyncJob(ctx)
if err != nil || job.Attempts != 2 {
t.Fatalf("unexpected retried job: %+v err=%v", job, err)
}
if err = store.CompleteDiscordRoleSyncJob(ctx, job.ID, "warning"); err != nil {
t.Fatal(err)
}
var status, warning string
if err = store.pool.QueryRow(ctx, `SELECT status,last_error FROM discord_role_sync_jobs WHERE id=$1`, job.ID).Scan(&status, &warning); err != nil {
t.Fatal(err)
}
if status != "completed" || warning != "warning" {
t.Fatalf("unexpected completed job state: status=%s warning=%s", status, warning)
}
bucket := time.Now().UnixNano()
defer func() {
_, _ = store.pool.Exec(ctx, `DELETE FROM discord_role_sync_jobs
WHERE event_id='__global__' AND action='full_reconcile' AND generation=$1`, bucket)
}()
if err = store.ScheduleGlobalDiscordRoleSync(ctx, bucket); err != nil {
t.Fatal(err)
}
if err = store.ScheduleGlobalDiscordRoleSync(ctx, bucket); err != nil {
t.Fatal(err)
}
var bucketJobs int
if err = store.pool.QueryRow(ctx, `SELECT count(*) FROM discord_role_sync_jobs
WHERE event_id='__global__' AND action='full_reconcile' AND generation=$1`, bucket).Scan(&bucketJobs); err != nil {
t.Fatal(err)
}
if bucketJobs != 1 {
t.Fatalf("expected one global bucket job, got %d", bucketJobs)
}
job, err = store.ClaimDiscordRoleSyncJob(ctx)
if err != nil || job.Action != application.DiscordRoleActionFullReconcile || job.Generation != bucket {
t.Fatalf("unexpected global job: %+v err=%v", job, err)
}
}
func TestFullScrimPipeline(t *testing.T) {
databaseURL := os.Getenv("TEST_DATABASE_URL")
if databaseURL == "" {
t.Skip("TEST_DATABASE_URL is not configured")
}
ctx := context.Background()
if err := Migrate(ctx, databaseURL, "../../../migrations"); err != nil {
t.Fatal(err)
}
store, err := Open(ctx, databaseURL)
if err != nil {
t.Fatal(err)
}
defer store.Close()
suffix := fmt.Sprintf("%d", time.Now().UnixNano())
accountID, eventID := "pipeline-admin-"+suffix, "pipeline-event-"+suffix
playerIDs, seriesIDs := make([]string, 0, 10), make([]string, 0, 1)
_, err = store.pool.Exec(ctx, `INSERT INTO accounts(id,discord_id,username,avatar_url,role,created_at)
VALUES($1,$2,'Pipeline Admin','','admin',now())`, accountID, accountID)
if err != nil {
t.Fatal(err)
}
defer func() {
_, _ = store.pool.Exec(ctx, `DELETE FROM audit_log WHERE actor_account_id=$1`, accountID)
_, _ = store.pool.Exec(ctx, `DELETE FROM events WHERE id=$1`, eventID)
_, _ = store.pool.Exec(ctx, `DELETE FROM discord_role_sync_jobs WHERE event_id=$1`, eventID)
for _, seriesID := range seriesIDs {
_, _ = store.pool.Exec(ctx, `DELETE FROM series WHERE id=$1`, seriesID)
}
for _, playerID := range playerIDs {
_, _ = store.pool.Exec(ctx, `DELETE FROM players WHERE id=$1`, playerID)
}
_, _ = store.pool.Exec(ctx, `DELETE FROM accounts WHERE id=$1`, accountID)
}()
service := application.New(store, discardPublisher{}, nil)
admin := domain.Account{ID: accountID, Role: domain.RoleAdmin}
now := time.Now().UTC()
event, err := service.CreateEvent(ctx, admin, domain.Event{
ID: eventID, Name: "Pipeline", StartsAt: now.Add(24 * time.Hour), EndsAt: now.Add(28 * time.Hour),
RegistrationDeadline: now.Add(23 * time.Hour), RulesetID: "standard-control-hybrid-control",
}, application.CreateEventOptions{})
if err != nil {
t.Fatal(err)
}
// CreateEvent owns the identifier, so use the returned value from here onward.
eventID = event.ID
for i := 0; i < 10; i++ {
player, _, createErr := service.CreateParticipant(ctx, admin, eventID, fmt.Sprintf("Player %d", i+1),
domain.Ratings{Tank: 18 + i%5, Damage: 19 + i%5, Support: 20 + i%5}, domain.Going)
if createErr != nil {
t.Fatal(createErr)
}
playerIDs = append(playerIDs, player.ID)
}
event, err = service.CloseRegistration(ctx, admin, eventID, event.RulesetID, event.Version)
if err != nil {
t.Fatal(err)
}
balance, err := service.GenerateBalance(ctx, admin, eventID, event.Version)
if err != nil || len(balance.Candidates) == 0 {
t.Fatalf("balance failed: candidates=%d err=%v", len(balance.Candidates), err)
}
roster, err := service.SelectWorkflowBalance(ctx, admin, eventID, balance.Candidates[0], balance.Event.Version)
if err != nil {
t.Fatal(err)
}
for _, team := range roster.Teams {
roster, err = service.SetRosterCaptain(ctx, admin, eventID, team.ID, team.Slots[0].PlayerID, roster.Version)
if err != nil {
t.Fatal(err)
}
}
event, err = store.GetEvent(ctx, eventID)
if err != nil {
t.Fatal(err)
}
event, err = service.ConfirmRosters(ctx, admin, eventID, event.Version, roster.Version)
if err != nil {
t.Fatal(err)
}
draft, err := service.InitializeBracket(ctx, admin, eventID, event.Version)
if err != nil {
t.Fatal(err)
}
draft, err = service.ConfirmBracket(ctx, admin, eventID, draft.Version)
if err != nil {
t.Fatal(err)
}
event, err = store.GetEvent(ctx, eventID)
if err != nil {
t.Fatal(err)
}
started, err := service.StartScrim(ctx, admin, eventID, event.Version)
if err != nil || started.Series == nil {
t.Fatalf("start failed: %#v, %v", started, err)
}
seriesIDs = append(seriesIDs, started.Series.ID)
series, err := service.TossSeriesCoin(ctx, admin, domain.Player{}, started.Series.ID, "integration-seed", started.Series.TeamAID, started.Series.Version)
if err != nil {
t.Fatal(err)
}
for mapNumber := 0; mapNumber < 2; mapNumber++ {
if series.Phase == domain.MapPickPhase {
series, err = service.PickSeriesMap(ctx, admin, domain.Player{}, series.ID, series.NextMapPickerID, series.AvailableMaps[0], series.Version)
if err != nil {
t.Fatal(err)
}
}
for series.Phase != domain.HeroBanPhase {
name := firstUnbanned(series.MapDraft.Pool, series.MapDraft.Banned)
series, err = service.BanSeriesMap(ctx, admin, domain.Player{}, series.ID, series.MapDraft.NextTeam(), name, series.Version)
if err != nil {
t.Fatal(err)
}
}
for series.Phase != domain.PlayingPhase {
teamID := series.HeroDraft.NextTeam()
hero := firstAvailableHero(series.HeroDraft, teamID)
series, err = service.BanSeriesHero(ctx, admin, domain.Player{}, series.ID, teamID, hero, series.Version)
if err != nil {
t.Fatal(err)
}
}
series, err = service.RecordSeriesResult(ctx, admin, domain.Player{}, series.ID, series.TeamAID, domain.TeamAWin, series.Version)
if err != nil {
t.Fatal(err)
}
}
if series.Phase != domain.SeriesComplete {
t.Fatalf("series did not complete: %s", series.Phase)
}
event, err = store.GetEvent(ctx, eventID)
if err != nil || event.State != domain.Completed {
t.Fatalf("event did not complete: state=%s err=%v", event.State, err)
}
}
func TestParallelTournamentHydratesLiveSeriesAndFinal(t *testing.T) {
databaseURL := os.Getenv("TEST_DATABASE_URL")
if databaseURL == "" {
t.Skip("TEST_DATABASE_URL is not configured")
}
ctx := context.Background()
if err := Migrate(ctx, databaseURL, "../../../migrations"); err != nil {
t.Fatal(err)
}
store, err := Open(ctx, databaseURL)
if err != nil {
t.Fatal(err)
}
defer store.Close()
suffix := fmt.Sprintf("%d", time.Now().UnixNano())
accountID, eventID, tournamentID := "parallel-admin-"+suffix, "parallel-event-"+suffix, "parallel-tournament-"+suffix
_, err = store.pool.Exec(ctx, `INSERT INTO accounts(id,discord_id,username,avatar_url,role,created_at) VALUES($1,$1,'Parallel Admin','','admin',now())`, accountID)
if err != nil {
t.Fatal(err)
}
defer func() {
_, _ = store.pool.Exec(ctx, `DELETE FROM series WHERE tournament_id=$1`, tournamentID)
_, _ = store.pool.Exec(ctx, `DELETE FROM events WHERE id=$1`, eventID)
_, _ = store.pool.Exec(ctx, `DELETE FROM discord_role_sync_jobs WHERE event_id=$1`, eventID)
_, _ = store.pool.Exec(ctx, `DELETE FROM accounts WHERE id=$1`, accountID)
}()
service := application.New(store, discardPublisher{}, nil)
event, err := service.CreateEvent(ctx, domain.Account{ID: accountID, Role: domain.RoleAdmin}, domain.Event{
ID: eventID, Name: "Parallel tournament", StartsAt: time.Now().Add(time.Hour), EndsAt: time.Now().Add(3 * time.Hour),
RegistrationDeadline: time.Now(), RulesetID: "standard-control-hybrid-control",
}, application.CreateEventOptions{})
if err != nil {
t.Fatal(err)
}
eventID = event.ID
rules, err := store.GetRuleset(ctx, event.RulesetID)
if err != nil {
t.Fatal(err)
}
draft, err := domain.NewBracketDraft(eventID, []string{"team-a", "team-b", "team-c", "team-d"})
if err != nil {
t.Fatal(err)
}
tournament, err := domain.NewGraphTournament(tournamentID, eventID, event.Name, *draft)
if err != nil {
t.Fatal(err)
}
for match, slot := range tournament.ReadyMatches() {
series, createErr := domain.NewSeries(slot.ID, eventID, tournamentID, [2]string{slot.TeamAID, slot.TeamBID}, rules)
if createErr != nil {
t.Fatal(createErr)
}
if match == 0 {
if createErr = series.Toss("parallel-seed", accountID, time.Now().UTC(), rules); createErr != nil {
t.Fatal(createErr)
}
}
if _, createErr = store.SaveSeries(ctx, *series); createErr != nil {
t.Fatal(createErr)
}
for index := range tournament.Matches {
if tournament.Matches[index].ID == slot.ID {
tournament.Matches[index].SeriesID = series.ID
tournament.Matches[index].Series = series
}
}
}
tournament.BuildRounds()
if _, err = store.SaveTournament(ctx, *tournament); err != nil {
t.Fatal(err)
}
hydrated, err := store.GetTournamentByEvent(ctx, eventID)
if err != nil {
t.Fatal(err)
}
if hydrated.Rounds[0][0].Phase != domain.MapBanPhase || hydrated.Rounds[0][1].Phase != domain.CoinTossPending {
t.Fatalf("parallel series were not hydrated independently: %s / %s", hydrated.Rounds[0][0].Phase, hydrated.Rounds[0][1].Phase)
}
for index := range hydrated.Matches {
if hydrated.Matches[index].Round != 0 || hydrated.Matches[index].Series == nil {
continue
}
completed := *hydrated.Matches[index].Series
completed.WinnerTeamID, completed.Phase, completed.Version = completed.TeamAID, domain.SeriesComplete, completed.Version+1
if _, err = store.SaveSeries(ctx, completed); err != nil {
t.Fatal(err)
}
if err = hydrated.ApplySeries(completed); err != nil {
t.Fatal(err)
}
if _, err = store.SaveTournament(ctx, hydrated); err != nil {
t.Fatal(err)
}
}
finalSlot := hydrated.ReadyMatches()[0]
final, err := domain.NewSeries(finalSlot.ID, eventID, tournamentID, [2]string{finalSlot.TeamAID, finalSlot.TeamBID}, rules)
if err != nil {
t.Fatal(err)
}
if _, err = store.SaveSeries(ctx, *final); err != nil {
t.Fatal(err)
}
hydrated, err = store.GetTournamentByEvent(ctx, eventID)
if err != nil || hydrated.Matches[len(hydrated.Matches)-1].Series == nil || hydrated.Matches[len(hydrated.Matches)-1].Series.Phase != domain.CoinTossPending {
t.Fatalf("final was not persisted and hydrated: err=%v", err)
}
}
func firstUnbanned(pool, banned []string) string {
for _, name := range pool {
found := false
for _, item := range banned {
found = found || item == name
}
if !found {
return name
}
}
return ""
}
func firstAvailableHero(draft *domain.HeroDraft, teamID string) string {
for _, hero := range draft.Heroes {
usedRole, usedHero, ownRepeat := false, false, false
for _, role := range draft.CurrentRoles[teamID] {
usedRole = usedRole || role == hero.Role
}
for _, ban := range draft.CurrentBans {
usedHero = usedHero || ban.Value == hero.Name
}
for _, name := range draft.SeriesBans[teamID] {
ownRepeat = ownRepeat || name == hero.Name
}
if !usedRole && !usedHero && !ownRepeat {
return hero.Name
}
}
return ""
}