/
codespawn
/
tgstream
Обзор
Документация
Войти
/
codespawn
/
tgstream
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
16
CI/CD
Аналитика
Безопасность
master
internal/adapter/sqlite/sqlite_test.go
183 строки
7 KB
codespawn
fix(sqlite): serialize writers to prevent SQLITE_BUSY under concurrency
29 июл 2026, 20:00
29 июл 2026, 20:00
d3172ee
Код
Авторство
О чём код?
package sqlite_test import ( "context" "fmt" "path/filepath" "strconv" "sync" "testing" "time" "gitverse.ru/codespawn/tgstream/internal/adapter/sqlite" "gitverse.ru/codespawn/tgstream/internal/domain" "gitverse.ru/codespawn/tgstream/internal/port" "gitverse.ru/codespawn/tgstream/internal/port/porttest" ) // newMemoryStore — изолированное in-memory хранилище для контракт-тестов. // :memory: + единое соединение (см. sqlite.New) даёт приватную БД на инстанс. func newMemoryStore(t *testing.T) port.Store { t.Helper() s, err := sqlite.New(context.Background(), sqlite.Config{Path: ":memory:"}) if err != nil { t.Fatalf("sqlite.New: %v", err) } t.Cleanup(func() { _ = s.Close() }) return s } // Контракт порта (L3): sqlite проходит ту же suite, что и фейк. func TestStore_Contract(t *testing.T) { porttest.RunStoreContract(t, func(t *testing.T) port.Store { return newMemoryStore(t) }) } // Конкурентные upsert — под -race. Воспроизводит прод-симптом из логов // (SQLITE_BUSY при max_concurrency=4): без сериализации писателей тяжёлые // конкурентные транзакции теряют весь батч (rollback по BUSY), и в истории // канала образуются дыры. // // Детерминизм RED-состояния (до writeMu) достигается тремя рычагами: // - BusyTimeout=1ms: ожидающий писатель почти не ждёт и тут же падает с BUSY; // - batch=50 постов на транзакцию: каждый UpsertPosts удерживает write-lock // достаточно долго, чтобы конкурент успел в него упереться; // - барьер start: все горутины выпускаются одновременно, максимизируя // конкуренцию за единственный write-lock SQLite. // // С writeMu писатели сериализуются на уровне приложения, SQLite видит ровно // одного, и BUSY не возникает ни при каком таймауте. func TestStore_ConcurrentUpsert(t *testing.T) { dir := t.TempDir() s, err := sqlite.New(context.Background(), sqlite.Config{ Path: filepath.Join(dir, "test.db"), BusyTimeout: 1 * time.Millisecond, }) if err != nil { t.Fatalf("sqlite.New: %v", err) } t.Cleanup(func() { _ = s.Close() }) ctx := context.Background() const ( goroutines = 16 batch = 50 // постов на транзакцию — «тяжёлая» транзакция ) base := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) start := make(chan struct{}) // барьер одновременного старта var wg sync.WaitGroup wg.Add(goroutines) for g := 0; g < goroutines; g++ { ch := fmt.Sprintf("ch%d", g) go func(ch string) { defer wg.Done() <-start posts := make([]domain.Post, 0, batch) for i := 0; i < batch; i++ { posts = append(posts, domain.Post{ ID: fmt.Sprintf("%s-%d", ch, i), Channel: ch, Text: "x", PublishedAt: base.Add(time.Duration(i) * time.Second), URL: fmt.Sprintf("https://t.me/%s/%d", ch, i), }) } if err := s.UpsertPosts(ctx, ch, posts); err != nil { t.Errorf("UpsertPosts(%s): %v", ch, err) } }(ch) } close(start) wg.Wait() // Ни одна транзакция не должна была упасть от SQLITE_BUSY: каждый батч // обязан целиком попасть в БД. Потеря любого батча = дыра в истории — // тот самый прод-симптом. total := 0 for g := 0; g < goroutines; g++ { ch := fmt.Sprintf("ch%d", g) got, err := s.GetPosts(ctx, ch, port.PostQuery{}) if err != nil { t.Fatalf("GetPosts(%s): %v", ch, err) } if len(got) != batch { t.Errorf("%s: got %d posts, want %d (батч потерян из-за SQLITE_BUSY?)", ch, len(got), batch) } total += len(got) } if want := goroutines * batch; total != want { t.Fatalf("total posts: got %d, want %d (потеряно %d при конкурентной записи)", total, want, want-total) } } // Trim оставляет N свежих постов по published_at DESC; прочие удаляются. // Индекс idx_posts_channel_time поддерживает порядок (новее→старее). func TestStore_Trim_KeepsNewest(t *testing.T) { s, err := sqlite.New(context.Background(), sqlite.Config{Path: ":memory:"}) if err != nil { t.Fatalf("sqlite.New: %v", err) } t.Cleanup(func() { _ = s.Close() }) ctx := context.Background() // 5 постов с возрастающими PublishedAt (id = порядок). Ожидаем, что Trim(3) // оставит id 4,5,2 — ошибочное предположение, если порядок перепутан. Ниже // фиксируем явно: остаются ТРИ НОВЕЙШИХ (3,4,5). base := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) for i := 1; i <= 5; i++ { p := domain.Post{ ID: strconv.Itoa(i), Channel: "durov", Text: strconv.Itoa(i), PublishedAt: base.AddDate(0, 0, i), URL: "https://t.me/durov/" + strconv.Itoa(i), } if err := s.UpsertPosts(ctx, "durov", []domain.Post{p}); err != nil { t.Fatalf("UpsertPosts(%d): %v", i, err) } } if err := s.Trim(ctx, "durov", 3); err != nil { t.Fatalf("Trim: %v", err) } posts, err := s.GetPosts(ctx, "durov", port.PostQuery{}) if err != nil { t.Fatalf("GetPosts: %v", err) } if len(posts) != 3 { t.Fatalf("after Trim(3): got %d posts, want 3", len(posts)) } // Новейшие = 5,4,3 (published_at DESC). Зафиксируем сохранённые id. want := map[string]bool{"3": true, "4": true, "5": true} for _, p := range posts { if !want[p.ID] { t.Errorf("Trim оставил лишний/неправильный пост id=%s; want новейшие 3,4,5", p.ID) } } // Trim не должен задевать другие каналы. if err := s.UpsertPosts(ctx, "other", []domain.Post{{ ID: "x", Channel: "other", Text: "x", PublishedAt: base, URL: "https://t.me/x/1", }}); err != nil { t.Fatalf("UpsertPosts(other): %v", err) } if err := s.Trim(ctx, "durov", 3); err != nil { t.Fatalf("Trim(again): %v", err) } other, err := s.GetPosts(ctx, "other", port.PostQuery{}) if err != nil { t.Fatalf("GetPosts(other): %v", err) } if len(other) != 1 { t.Errorf("Trim задел другой канал: other has %d posts, want 1", len(other)) } // max<=0 — no-op. if err := s.Trim(ctx, "durov", 0); err != nil { t.Errorf("Trim(0): %v", err) } }