2026-07-10 15:32:34 +08:00
|
|
|
package storage
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"os"
|
2026-07-10 16:05:19 +08:00
|
|
|
"sync"
|
2026-07-10 15:32:34 +08:00
|
|
|
"testing"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func TestStore_InsertAndCount(t *testing.T) {
|
|
|
|
|
dbFile := "test_telegram.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
s, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s.Close()
|
|
|
|
|
|
|
|
|
|
if err := s.Insert("ZCZC TEST NNNN"); err != nil {
|
|
|
|
|
t.Fatalf("Insert() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
count, err := s.CountUnprocessed()
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("CountUnprocessed() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if count != 1 {
|
|
|
|
|
t.Errorf("expected 1 unprocessed, got %d", count)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestStore_LoadUnprocessed(t *testing.T) {
|
|
|
|
|
dbFile := "test_load.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
s, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s.Close()
|
|
|
|
|
|
|
|
|
|
s.Insert("ZCZC MSG1 NNNN")
|
|
|
|
|
s.Insert("ZCZC MSG2 NNNN")
|
|
|
|
|
|
|
|
|
|
telegrams, err := s.LoadUnprocessed()
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("LoadUnprocessed() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if len(telegrams) != 2 {
|
|
|
|
|
t.Errorf("expected 2 telegrams, got %d", len(telegrams))
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestStore_MarkProcessed(t *testing.T) {
|
|
|
|
|
dbFile := "test_mark.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
s, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s.Close()
|
|
|
|
|
|
|
|
|
|
s.Insert("ZCZC TEST NNNN")
|
|
|
|
|
|
|
|
|
|
telegrams, _ := s.LoadUnprocessed()
|
|
|
|
|
if len(telegrams) != 1 {
|
|
|
|
|
t.Fatalf("expected 1 telegram, got %d", len(telegrams))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if err := s.MarkProcessed(telegrams[0].ID); err != nil {
|
|
|
|
|
t.Fatalf("MarkProcessed() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
count, _ := s.CountUnprocessed()
|
|
|
|
|
if count != 0 {
|
|
|
|
|
t.Errorf("expected 0 unprocessed after marking, got %d", count)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestStore_Empty(t *testing.T) {
|
|
|
|
|
dbFile := "test_empty.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
s, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s.Close()
|
|
|
|
|
|
|
|
|
|
telegrams, err := s.LoadUnprocessed()
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("LoadUnprocessed() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if len(telegrams) != 0 {
|
|
|
|
|
t.Errorf("expected 0 telegrams, got %d", len(telegrams))
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-10 16:05:19 +08:00
|
|
|
|
|
|
|
|
func TestStore_ConcurrentInsert(t *testing.T) {
|
|
|
|
|
dbFile := "test_concurrent.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
s, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s.Close()
|
|
|
|
|
|
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
|
const numGoroutines = 10
|
|
|
|
|
const insertsPerGoroutine = 10
|
|
|
|
|
|
|
|
|
|
for i := 0; i < numGoroutines; i++ {
|
|
|
|
|
wg.Add(1)
|
|
|
|
|
go func(id int) {
|
|
|
|
|
defer wg.Done()
|
|
|
|
|
for j := 0; j < insertsPerGoroutine; j++ {
|
|
|
|
|
telegram := "ZCZC CONCURRENT MSG"
|
|
|
|
|
if err := s.Insert(telegram); err != nil {
|
|
|
|
|
t.Errorf("concurrent Insert() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}(i)
|
|
|
|
|
}
|
|
|
|
|
wg.Wait()
|
|
|
|
|
|
|
|
|
|
count, err := s.CountUnprocessed()
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("CountUnprocessed() after concurrent inserts failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
expected := int64(numGoroutines * insertsPerGoroutine)
|
|
|
|
|
if count != expected {
|
|
|
|
|
t.Errorf("expected %d unprocessed, got %d", expected, count)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestStore_CloseAndReopen(t *testing.T) {
|
|
|
|
|
dbFile := "test_reopen.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
// Create and insert
|
|
|
|
|
s, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if err := s.Insert("ZCZC PERSIST NNNN"); err != nil {
|
|
|
|
|
t.Fatalf("Insert() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
s.Close()
|
|
|
|
|
|
|
|
|
|
// Reopen without init flag
|
|
|
|
|
s2, err := New(dbFile, false)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() reopen failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s2.Close()
|
|
|
|
|
|
|
|
|
|
count, err := s2.CountUnprocessed()
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("CountUnprocessed() on reopened db failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if count != 1 {
|
|
|
|
|
t.Errorf("expected 1 unprocessed after reopen, got %d", count)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestStore_InitClearsDatabase(t *testing.T) {
|
|
|
|
|
dbFile := "test_init_clear.db"
|
|
|
|
|
defer os.Remove(dbFile)
|
|
|
|
|
|
|
|
|
|
// Create and insert
|
|
|
|
|
s, err := New(dbFile, false)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
s.Insert("ZCZC WILL BE CLEARED NNNN")
|
|
|
|
|
s.Close()
|
|
|
|
|
|
|
|
|
|
// Reopen with init=true (should recreate)
|
|
|
|
|
s2, err := New(dbFile, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("New() with init=true failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer s2.Close()
|
|
|
|
|
|
|
|
|
|
count, err := s2.CountUnprocessed()
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("CountUnprocessed() failed: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if count != 0 {
|
|
|
|
|
t.Errorf("expected 0 unprocessed after init, got %d", count)
|
|
|
|
|
}
|
|
|
|
|
}
|