count_cache_test.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431
  1. package storage
  2. import (
  3. "fmt"
  4. "sync"
  5. "testing"
  6. "time"
  7. )
  8. // TestCountFastLifecycle verifies the COUNT(*) metadata counter stays exact
  9. // across insert, update, delete, bulk insert, and truncate.
  10. func TestCountFastLifecycle(t *testing.T) {
  11. kv := newTestKVServer(t)
  12. defer kv.close()
  13. pool := newTestKVPool(kv, 2, 5*time.Second)
  14. defer pool.Close()
  15. schemas := NewSchemaManager(pool, "testdb")
  16. tables := NewTableManager(pool, schemas, "testdb")
  17. err := schemas.CreateTable(&Schema{
  18. Name: "users",
  19. Columns: []Column{
  20. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  21. {Name: "name", Type: "TEXT", Nullable: true},
  22. },
  23. })
  24. if err != nil {
  25. t.Fatalf("create table: %v", err)
  26. }
  27. assertCount := func(want int) {
  28. t.Helper()
  29. got, err := tables.CountFast("users")
  30. if err != nil {
  31. t.Fatalf("CountFast: %v", err)
  32. }
  33. if got != want {
  34. t.Fatalf("CountFast = %d, want %d", got, want)
  35. }
  36. }
  37. assertCount(0)
  38. for i := int64(1); i <= 3; i++ {
  39. if err := tables.Insert("users", Row{"id": i, "name": fmt.Sprintf("u%d", i)}); err != nil {
  40. t.Fatalf("insert %d: %v", i, err)
  41. }
  42. }
  43. assertCount(3)
  44. // UPDATE does not change the count.
  45. n, err := tables.Update("users", Row{"name": "renamed"}, func(r Row) bool {
  46. return fmt.Sprintf("%v", r["id"]) == "1"
  47. })
  48. if err != nil || n != 1 {
  49. t.Fatalf("update: n=%d err=%v", n, err)
  50. }
  51. assertCount(3)
  52. // DELETE decrements.
  53. n, err = tables.Delete("users", func(r Row) bool {
  54. return fmt.Sprintf("%v", r["id"]) == "2"
  55. })
  56. if err != nil || n != 1 {
  57. t.Fatalf("delete: n=%d err=%v", n, err)
  58. }
  59. assertCount(2)
  60. // Bulk insert adds len(rows).
  61. n, err = tables.InsertBulk("users", []Row{
  62. {"id": 4, "name": "u4"},
  63. {"id": 5, "name": "u5"},
  64. {"id": 6, "name": "u6"},
  65. })
  66. if err != nil || n != 3 {
  67. t.Fatalf("bulk insert: n=%d err=%v", n, err)
  68. }
  69. assertCount(5)
  70. // Truncate zeroes the count.
  71. if _, err := tables.Truncate("users"); err != nil {
  72. t.Fatalf("truncate: %v", err)
  73. }
  74. assertCount(0)
  75. }
  76. // TestCountFastDerivedAfterRestart verifies the counter is re-derived from
  77. // durable rows when a fresh TableManager (process restart) has no in-memory
  78. // count yet.
  79. func TestCountFastDerivedAfterRestart(t *testing.T) {
  80. kv := newTestKVServer(t)
  81. defer kv.close()
  82. pool := newTestKVPool(kv, 2, 5*time.Second)
  83. defer pool.Close()
  84. schemas := NewSchemaManager(pool, "testdb")
  85. tables := NewTableManager(pool, schemas, "testdb")
  86. err := schemas.CreateTable(&Schema{
  87. Name: "events",
  88. Columns: []Column{
  89. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  90. },
  91. })
  92. if err != nil {
  93. t.Fatalf("create table: %v", err)
  94. }
  95. for i := int64(1); i <= 4; i++ {
  96. if err := tables.Insert("events", Row{"id": i}); err != nil {
  97. t.Fatalf("insert %d: %v", i, err)
  98. }
  99. }
  100. restartedSchemas := NewSchemaManager(pool, "testdb")
  101. restartedTables := NewTableManager(pool, restartedSchemas, "testdb")
  102. got, err := restartedTables.CountFast("events")
  103. if err != nil {
  104. t.Fatalf("CountFast after restart: %v", err)
  105. }
  106. if got != 4 {
  107. t.Fatalf("CountFast after restart = %d, want 4", got)
  108. }
  109. }
  110. // TestIncrementalCacheAndIndexMaintenance verifies that writes keep
  111. // already-built in-memory indexes coherent without full-table invalidation.
  112. func TestIncrementalCacheAndIndexMaintenance(t *testing.T) {
  113. kv := newTestKVServer(t)
  114. defer kv.close()
  115. pool := newTestKVPool(kv, 2, 5*time.Second)
  116. defer pool.Close()
  117. schemas := NewSchemaManager(pool, "testdb")
  118. tables := NewTableManager(pool, schemas, "testdb")
  119. err := schemas.CreateTable(&Schema{
  120. Name: "users",
  121. Columns: []Column{
  122. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  123. {Name: "status", Type: "TEXT", Nullable: false},
  124. },
  125. })
  126. if err != nil {
  127. t.Fatalf("create table: %v", err)
  128. }
  129. err = schemas.CreateIndex(&Index{
  130. Name: "idx_users_status",
  131. Table: "users",
  132. Columns: []IndexColumn{
  133. {Name: "status"},
  134. },
  135. })
  136. if err != nil {
  137. t.Fatalf("create index: %v", err)
  138. }
  139. // Seed: 1 active, 1 inactive.
  140. if err := tables.Insert("users", Row{"id": 1, "status": "active"}); err != nil {
  141. t.Fatalf("insert 1: %v", err)
  142. }
  143. if err := tables.Insert("users", Row{"id": 2, "status": "inactive"}); err != nil {
  144. t.Fatalf("insert 2: %v", err)
  145. }
  146. // Build the in-memory index (Select only streams rows).
  147. all, err := tables.Select("users", nil)
  148. if err != nil || len(all) != 2 {
  149. t.Fatalf("initial select: len=%d err=%v", len(all), err)
  150. }
  151. active, err := tables.SelectByIndex("users", "idx_users_status", "active")
  152. if err != nil || len(active) != 1 {
  153. t.Fatalf("initial active index: len=%d err=%v", len(active), err)
  154. }
  155. // Bulk insert two more active rows; the index must reflect them.
  156. if n, err := tables.InsertBulk("users", []Row{
  157. {"id": 3, "status": "active"},
  158. {"id": 4, "status": "active"},
  159. }); err != nil || n != 2 {
  160. t.Fatalf("bulk insert: n=%d err=%v", n, err)
  161. }
  162. active, err = tables.SelectByIndex("users", "idx_users_status", "active")
  163. if err != nil || len(active) != 3 {
  164. t.Fatalf("active index after bulk: len=%d err=%v", len(active), err)
  165. }
  166. all, err = tables.Select("users", nil)
  167. if err != nil || len(all) != 4 {
  168. t.Fatalf("select after bulk: len=%d err=%v", len(all), err)
  169. }
  170. // Update one active -> inactive; both index buckets must stay coherent.
  171. if n, err := tables.Update("users", Row{"status": "inactive"}, func(r Row) bool {
  172. return fmt.Sprintf("%v", r["id"]) == "3"
  173. }); err != nil || n != 1 {
  174. t.Fatalf("update: n=%d err=%v", n, err)
  175. }
  176. active, _ = tables.SelectByIndex("users", "idx_users_status", "active")
  177. inactive, _ := tables.SelectByIndex("users", "idx_users_status", "inactive")
  178. if len(active) != 2 || len(inactive) != 2 {
  179. t.Fatalf("indexes after update: active=%d inactive=%d", len(active), len(inactive))
  180. }
  181. // Delete one row; cache and index must shrink.
  182. if n, err := tables.Delete("users", func(r Row) bool {
  183. return fmt.Sprintf("%v", r["id"]) == "4"
  184. }); err != nil || n != 1 {
  185. t.Fatalf("delete: n=%d err=%v", n, err)
  186. }
  187. all, err = tables.Select("users", nil)
  188. if err != nil || len(all) != 3 {
  189. t.Fatalf("select after delete: len=%d err=%v", len(all), err)
  190. }
  191. active, _ = tables.SelectByIndex("users", "idx_users_status", "active")
  192. if len(active) != 1 {
  193. t.Fatalf("active index after delete: len=%d", len(active))
  194. }
  195. }
  196. func TestClearIndexDropsInMemoryEntries(t *testing.T) {
  197. kv := newTestKVServer(t)
  198. defer kv.close()
  199. pool := newTestKVPool(kv, 2, 5*time.Second)
  200. defer pool.Close()
  201. schemas := NewSchemaManager(pool, "testdb")
  202. tables := NewTableManager(pool, schemas, "testdb")
  203. if err := schemas.CreateTable(&Schema{
  204. Name: "users",
  205. Columns: []Column{
  206. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  207. {Name: "status", Type: "TEXT", Nullable: false},
  208. },
  209. }); err != nil {
  210. t.Fatalf("create table: %v", err)
  211. }
  212. if err := schemas.CreateIndex(&Index{Name: "idx_status", Table: "users", Columns: []IndexColumn{{Name: "status"}}}); err != nil {
  213. t.Fatalf("create index: %v", err)
  214. }
  215. if err := tables.Insert("users", Row{"id": 1, "status": "old"}); err != nil {
  216. t.Fatalf("insert: %v", err)
  217. }
  218. if rows, err := tables.SelectByIndex("users", "idx_status", "old"); err != nil || len(rows) != 1 {
  219. t.Fatalf("build index: len=%d err=%v", len(rows), err)
  220. }
  221. if err := tables.ClearIndex("idx_status", "users", []string{"status"}); err != nil {
  222. t.Fatalf("clear index: %v", err)
  223. }
  224. tables.cacheMu.RLock()
  225. _, cached := tables.indexCache["idx_status"]
  226. _, mapped := tables.indexTable["idx_status"]
  227. tables.cacheMu.RUnlock()
  228. if cached || mapped {
  229. t.Fatalf("cleared index remains cached: cache=%v table=%v", cached, mapped)
  230. }
  231. if rows, err := tables.SelectByIndex("users", "idx_status", "old"); err != nil || len(rows) != 0 {
  232. t.Fatalf("cleared index rebuilt during drop: len=%d err=%v", len(rows), err)
  233. }
  234. if err := tables.BuildIndex("idx_status", "users", []string{"status"}); err != nil {
  235. t.Fatalf("rebuild index: %v", err)
  236. }
  237. if rows, err := tables.SelectByIndex("users", "idx_status", "old"); err != nil || len(rows) != 1 {
  238. t.Fatalf("rebuilt index: len=%d err=%v", len(rows), err)
  239. }
  240. }
  241. func TestConcurrentDuplicateInsertKeepsCountExact(t *testing.T) {
  242. kv := newTestKVServer(t)
  243. defer kv.close()
  244. pool := newTestKVPool(kv, 2, 5*time.Second)
  245. defer pool.Close()
  246. schemas := NewSchemaManager(pool, "testdb")
  247. tables := NewTableManager(pool, schemas, "testdb")
  248. if err := schemas.CreateTable(&Schema{
  249. Name: "users",
  250. Columns: []Column{{Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true}},
  251. }); err != nil {
  252. t.Fatalf("create table: %v", err)
  253. }
  254. if got, err := tables.CountFast("users"); err != nil || got != 0 {
  255. t.Fatalf("initial count=%d err=%v", got, err)
  256. }
  257. start := make(chan struct{})
  258. errs := make(chan error, 2)
  259. for i := 0; i < 2; i++ {
  260. go func() {
  261. <-start
  262. errs <- tables.Insert("users", Row{"id": 1})
  263. }()
  264. }
  265. close(start)
  266. successes := 0
  267. for i := 0; i < 2; i++ {
  268. if err := <-errs; err == nil {
  269. successes++
  270. }
  271. }
  272. if successes != 1 {
  273. t.Fatalf("successful inserts=%d, want 1", successes)
  274. }
  275. if got, err := tables.CountFast("users"); err != nil || got != 1 {
  276. t.Fatalf("final count=%d err=%v, want 1", got, err)
  277. }
  278. }
  279. // TestConcurrentFirstLoadAndInsert runs a first scan (Select on a table)
  280. // concurrently with inserts, then verifies the table ends up coherent with
  281. // durable rows. Run under -race to detect map/slice data races.
  282. func TestConcurrentFirstLoadAndInsert(t *testing.T) {
  283. kv := newTestKVServer(t)
  284. defer kv.close()
  285. pool := newTestKVPool(kv, 16, 10*time.Second)
  286. defer pool.Close()
  287. schemas := NewSchemaManager(pool, "testdb")
  288. tables := NewTableManager(pool, schemas, "testdb")
  289. err := schemas.CreateTable(&Schema{
  290. Name: "users",
  291. Columns: []Column{
  292. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  293. {Name: "name", Type: "TEXT", Nullable: true},
  294. },
  295. })
  296. if err != nil {
  297. t.Fatalf("create table: %v", err)
  298. }
  299. const n = 200
  300. start := make(chan struct{})
  301. var wg sync.WaitGroup
  302. for i := 0; i < n; i++ {
  303. i := i
  304. wg.Add(1)
  305. go func() {
  306. defer wg.Done()
  307. <-start
  308. if i%2 == 0 {
  309. if err := tables.Insert("users", Row{"id": int64(i + 1), "name": fmt.Sprintf("u%d", i)}); err != nil {
  310. t.Errorf("insert %d: %v", i, err)
  311. return
  312. }
  313. } else {
  314. if _, err := tables.Select("users", nil); err != nil {
  315. t.Errorf("select: %v", err)
  316. }
  317. }
  318. }()
  319. }
  320. close(start)
  321. wg.Wait()
  322. rows, err := tables.Select("users", nil)
  323. if err != nil {
  324. t.Fatalf("final select: %v", err)
  325. }
  326. if want := n / 2; len(rows) != want {
  327. t.Fatalf("final select = %d rows, want %d", len(rows), want)
  328. }
  329. if got, _ := tables.CountFast("users"); got != n/2 {
  330. t.Fatalf("final count = %d, want %d", got, n/2)
  331. }
  332. }
  333. // TestConcurrentFirstCountFastAndInsert runs first-time count derivation
  334. // concurrently with inserts, then verifies the count ends up exact.
  335. func TestConcurrentFirstCountFastAndInsert(t *testing.T) {
  336. kv := newTestKVServer(t)
  337. defer kv.close()
  338. pool := newTestKVPool(kv, 16, 10*time.Second)
  339. defer pool.Close()
  340. schemas := NewSchemaManager(pool, "testdb")
  341. tables := NewTableManager(pool, schemas, "testdb")
  342. err := schemas.CreateTable(&Schema{
  343. Name: "items",
  344. Columns: []Column{
  345. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  346. },
  347. })
  348. if err != nil {
  349. t.Fatalf("create table: %v", err)
  350. }
  351. const n = 200
  352. start := make(chan struct{})
  353. var wg sync.WaitGroup
  354. for i := 0; i < n; i++ {
  355. i := i
  356. wg.Add(1)
  357. go func() {
  358. defer wg.Done()
  359. <-start
  360. if i%2 == 0 {
  361. if err := tables.Insert("items", Row{"id": int64(i + 1)}); err != nil {
  362. t.Errorf("insert %d: %v", i, err)
  363. }
  364. } else {
  365. if _, err := tables.CountFast("items"); err != nil {
  366. t.Errorf("count: %v", err)
  367. }
  368. }
  369. }()
  370. }
  371. close(start)
  372. wg.Wait()
  373. if got, _ := tables.CountFast("items"); got != n/2 {
  374. t.Fatalf("final count = %d, want %d", got, n/2)
  375. }
  376. rows, err := tables.Select("items", nil)
  377. if err != nil {
  378. t.Fatalf("final select: %v", err)
  379. }
  380. if len(rows) != n/2 {
  381. t.Fatalf("final select = %d rows, want %d", len(rows), n/2)
  382. }
  383. }