2
0

count_cache_test.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432
  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 the row
  111. // cache and already-built in-memory indexes coherent without full-table cache
  112. // invalidation.
  113. func TestIncrementalCacheAndIndexMaintenance(t *testing.T) {
  114. kv := newTestKVServer(t)
  115. defer kv.close()
  116. pool := newTestKVPool(kv, 2, 5*time.Second)
  117. defer pool.Close()
  118. schemas := NewSchemaManager(pool, "testdb")
  119. tables := NewTableManager(pool, schemas, "testdb")
  120. err := schemas.CreateTable(&Schema{
  121. Name: "users",
  122. Columns: []Column{
  123. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  124. {Name: "status", Type: "TEXT", Nullable: false},
  125. },
  126. })
  127. if err != nil {
  128. t.Fatalf("create table: %v", err)
  129. }
  130. err = schemas.CreateIndex(&Index{
  131. Name: "idx_users_status",
  132. Table: "users",
  133. Columns: []IndexColumn{
  134. {Name: "status"},
  135. },
  136. })
  137. if err != nil {
  138. t.Fatalf("create index: %v", err)
  139. }
  140. // Seed: 1 active, 1 inactive.
  141. if err := tables.Insert("users", Row{"id": 1, "status": "active"}); err != nil {
  142. t.Fatalf("insert 1: %v", err)
  143. }
  144. if err := tables.Insert("users", Row{"id": 2, "status": "inactive"}); err != nil {
  145. t.Fatalf("insert 2: %v", err)
  146. }
  147. // Load the row cache and the index.
  148. all, err := tables.Select("users", nil)
  149. if err != nil || len(all) != 2 {
  150. t.Fatalf("initial select: len=%d err=%v", len(all), err)
  151. }
  152. active, err := tables.SelectByIndex("users", "idx_users_status", "active")
  153. if err != nil || len(active) != 1 {
  154. t.Fatalf("initial active index: len=%d err=%v", len(active), err)
  155. }
  156. // Bulk insert two more active rows; the index must reflect them.
  157. if n, err := tables.InsertBulk("users", []Row{
  158. {"id": 3, "status": "active"},
  159. {"id": 4, "status": "active"},
  160. }); err != nil || n != 2 {
  161. t.Fatalf("bulk insert: n=%d err=%v", n, err)
  162. }
  163. active, err = tables.SelectByIndex("users", "idx_users_status", "active")
  164. if err != nil || len(active) != 3 {
  165. t.Fatalf("active index after bulk: len=%d err=%v", len(active), err)
  166. }
  167. all, err = tables.Select("users", nil)
  168. if err != nil || len(all) != 4 {
  169. t.Fatalf("select after bulk: len=%d err=%v", len(all), err)
  170. }
  171. // Update one active -> inactive; both index buckets must stay coherent.
  172. if n, err := tables.Update("users", Row{"status": "inactive"}, func(r Row) bool {
  173. return fmt.Sprintf("%v", r["id"]) == "3"
  174. }); err != nil || n != 1 {
  175. t.Fatalf("update: n=%d err=%v", n, err)
  176. }
  177. active, _ = tables.SelectByIndex("users", "idx_users_status", "active")
  178. inactive, _ := tables.SelectByIndex("users", "idx_users_status", "inactive")
  179. if len(active) != 2 || len(inactive) != 2 {
  180. t.Fatalf("indexes after update: active=%d inactive=%d", len(active), len(inactive))
  181. }
  182. // Delete one row; cache and index must shrink.
  183. if n, err := tables.Delete("users", func(r Row) bool {
  184. return fmt.Sprintf("%v", r["id"]) == "4"
  185. }); err != nil || n != 1 {
  186. t.Fatalf("delete: n=%d err=%v", n, err)
  187. }
  188. all, err = tables.Select("users", nil)
  189. if err != nil || len(all) != 3 {
  190. t.Fatalf("select after delete: len=%d err=%v", len(all), err)
  191. }
  192. active, _ = tables.SelectByIndex("users", "idx_users_status", "active")
  193. if len(active) != 1 {
  194. t.Fatalf("active index after delete: len=%d", len(active))
  195. }
  196. }
  197. func TestClearIndexDropsInMemoryEntries(t *testing.T) {
  198. kv := newTestKVServer(t)
  199. defer kv.close()
  200. pool := newTestKVPool(kv, 2, 5*time.Second)
  201. defer pool.Close()
  202. schemas := NewSchemaManager(pool, "testdb")
  203. tables := NewTableManager(pool, schemas, "testdb")
  204. if err := schemas.CreateTable(&Schema{
  205. Name: "users",
  206. Columns: []Column{
  207. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  208. {Name: "status", Type: "TEXT", Nullable: false},
  209. },
  210. }); err != nil {
  211. t.Fatalf("create table: %v", err)
  212. }
  213. if err := schemas.CreateIndex(&Index{Name: "idx_status", Table: "users", Columns: []IndexColumn{{Name: "status"}}}); err != nil {
  214. t.Fatalf("create index: %v", err)
  215. }
  216. if err := tables.Insert("users", Row{"id": 1, "status": "old"}); err != nil {
  217. t.Fatalf("insert: %v", err)
  218. }
  219. if rows, err := tables.SelectByIndex("users", "idx_status", "old"); err != nil || len(rows) != 1 {
  220. t.Fatalf("build index: len=%d err=%v", len(rows), err)
  221. }
  222. if err := tables.ClearIndex("idx_status", "users", []string{"status"}); err != nil {
  223. t.Fatalf("clear index: %v", err)
  224. }
  225. tables.cacheMu.RLock()
  226. _, cached := tables.indexCache["idx_status"]
  227. _, mapped := tables.indexTable["idx_status"]
  228. tables.cacheMu.RUnlock()
  229. if cached || mapped {
  230. t.Fatalf("cleared index remains cached: cache=%v table=%v", cached, mapped)
  231. }
  232. if rows, err := tables.SelectByIndex("users", "idx_status", "old"); err != nil || len(rows) != 0 {
  233. t.Fatalf("cleared index rebuilt during drop: len=%d err=%v", len(rows), err)
  234. }
  235. if err := tables.BuildIndex("idx_status", "users", []string{"status"}); err != nil {
  236. t.Fatalf("rebuild index: %v", err)
  237. }
  238. if rows, err := tables.SelectByIndex("users", "idx_status", "old"); err != nil || len(rows) != 1 {
  239. t.Fatalf("rebuilt index: len=%d err=%v", len(rows), err)
  240. }
  241. }
  242. func TestConcurrentDuplicateInsertKeepsCountExact(t *testing.T) {
  243. kv := newTestKVServer(t)
  244. defer kv.close()
  245. pool := newTestKVPool(kv, 2, 5*time.Second)
  246. defer pool.Close()
  247. schemas := NewSchemaManager(pool, "testdb")
  248. tables := NewTableManager(pool, schemas, "testdb")
  249. if err := schemas.CreateTable(&Schema{
  250. Name: "users",
  251. Columns: []Column{{Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true}},
  252. }); err != nil {
  253. t.Fatalf("create table: %v", err)
  254. }
  255. if got, err := tables.CountFast("users"); err != nil || got != 0 {
  256. t.Fatalf("initial count=%d err=%v", got, err)
  257. }
  258. start := make(chan struct{})
  259. errs := make(chan error, 2)
  260. for i := 0; i < 2; i++ {
  261. go func() {
  262. <-start
  263. errs <- tables.Insert("users", Row{"id": 1})
  264. }()
  265. }
  266. close(start)
  267. successes := 0
  268. for i := 0; i < 2; i++ {
  269. if err := <-errs; err == nil {
  270. successes++
  271. }
  272. }
  273. if successes != 1 {
  274. t.Fatalf("successful inserts=%d, want 1", successes)
  275. }
  276. if got, err := tables.CountFast("users"); err != nil || got != 1 {
  277. t.Fatalf("final count=%d err=%v, want 1", got, err)
  278. }
  279. }
  280. // TestConcurrentFirstLoadAndInsert runs the first cache load (Select on an
  281. // unloaded table) concurrently with inserts, then verifies the cache ends up
  282. // coherent with durable rows. Run under -race to detect map/slice data races.
  283. func TestConcurrentFirstLoadAndInsert(t *testing.T) {
  284. kv := newTestKVServer(t)
  285. defer kv.close()
  286. pool := newTestKVPool(kv, 16, 10*time.Second)
  287. defer pool.Close()
  288. schemas := NewSchemaManager(pool, "testdb")
  289. tables := NewTableManager(pool, schemas, "testdb")
  290. err := schemas.CreateTable(&Schema{
  291. Name: "users",
  292. Columns: []Column{
  293. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  294. {Name: "name", Type: "TEXT", Nullable: true},
  295. },
  296. })
  297. if err != nil {
  298. t.Fatalf("create table: %v", err)
  299. }
  300. const n = 200
  301. start := make(chan struct{})
  302. var wg sync.WaitGroup
  303. for i := 0; i < n; i++ {
  304. i := i
  305. wg.Add(1)
  306. go func() {
  307. defer wg.Done()
  308. <-start
  309. if i%2 == 0 {
  310. if err := tables.Insert("users", Row{"id": int64(i + 1), "name": fmt.Sprintf("u%d", i)}); err != nil {
  311. t.Errorf("insert %d: %v", i, err)
  312. return
  313. }
  314. } else {
  315. if _, err := tables.Select("users", nil); err != nil {
  316. t.Errorf("select: %v", err)
  317. }
  318. }
  319. }()
  320. }
  321. close(start)
  322. wg.Wait()
  323. rows, err := tables.Select("users", nil)
  324. if err != nil {
  325. t.Fatalf("final select: %v", err)
  326. }
  327. if want := n / 2; len(rows) != want {
  328. t.Fatalf("final select = %d rows, want %d", len(rows), want)
  329. }
  330. if got, _ := tables.CountFast("users"); got != n/2 {
  331. t.Fatalf("final count = %d, want %d", got, n/2)
  332. }
  333. }
  334. // TestConcurrentFirstCountFastAndInsert runs first-time count derivation
  335. // concurrently with inserts, then verifies the count ends up exact.
  336. func TestConcurrentFirstCountFastAndInsert(t *testing.T) {
  337. kv := newTestKVServer(t)
  338. defer kv.close()
  339. pool := newTestKVPool(kv, 16, 10*time.Second)
  340. defer pool.Close()
  341. schemas := NewSchemaManager(pool, "testdb")
  342. tables := NewTableManager(pool, schemas, "testdb")
  343. err := schemas.CreateTable(&Schema{
  344. Name: "items",
  345. Columns: []Column{
  346. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  347. },
  348. })
  349. if err != nil {
  350. t.Fatalf("create table: %v", err)
  351. }
  352. const n = 200
  353. start := make(chan struct{})
  354. var wg sync.WaitGroup
  355. for i := 0; i < n; i++ {
  356. i := i
  357. wg.Add(1)
  358. go func() {
  359. defer wg.Done()
  360. <-start
  361. if i%2 == 0 {
  362. if err := tables.Insert("items", Row{"id": int64(i + 1)}); err != nil {
  363. t.Errorf("insert %d: %v", i, err)
  364. }
  365. } else {
  366. if _, err := tables.CountFast("items"); err != nil {
  367. t.Errorf("count: %v", err)
  368. }
  369. }
  370. }()
  371. }
  372. close(start)
  373. wg.Wait()
  374. if got, _ := tables.CountFast("items"); got != n/2 {
  375. t.Fatalf("final count = %d, want %d", got, n/2)
  376. }
  377. rows, err := tables.Select("items", nil)
  378. if err != nil {
  379. t.Fatalf("final select: %v", err)
  380. }
  381. if len(rows) != n/2 {
  382. t.Fatalf("final select = %d rows, want %d", len(rows), n/2)
  383. }
  384. }