pizzakv_integration_test.go 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  1. package storage
  2. import (
  3. "bytes"
  4. "fmt"
  5. "os"
  6. "os/exec"
  7. "path/filepath"
  8. "sync"
  9. "sync/atomic"
  10. "testing"
  11. "time"
  12. "github.com/goccy/go-json"
  13. )
  14. func startPizzaKVTest(t testing.TB, binary, socket, database string) func() {
  15. t.Helper()
  16. cmd := exec.Command(binary, "-unix="+socket, "-path="+database)
  17. cmd.Stdout = os.Stderr
  18. cmd.Stderr = os.Stderr
  19. if err := cmd.Start(); err != nil {
  20. t.Fatalf("start PizzaKV: %v", err)
  21. }
  22. var once sync.Once
  23. stop := func() {
  24. once.Do(func() {
  25. _ = cmd.Process.Kill()
  26. _ = cmd.Wait()
  27. })
  28. }
  29. t.Cleanup(stop)
  30. return stop
  31. }
  32. func waitPizzaKVPool(t testing.TB, socket string) *KVPool {
  33. t.Helper()
  34. addr := "unix:" + socket
  35. deadline := time.Now().Add(10 * time.Second)
  36. for time.Now().Before(deadline) {
  37. pool, err := NewKVPool(addr, 2, 5*time.Second)
  38. if err == nil {
  39. return pool
  40. }
  41. time.Sleep(20 * time.Millisecond)
  42. }
  43. t.Fatal("PizzaKV did not become ready")
  44. return nil
  45. }
  46. func shortPizzaKVSocket(t testing.TB) string {
  47. t.Helper()
  48. dir, err := os.MkdirTemp("/tmp", "pkv-")
  49. if err != nil {
  50. t.Fatal(err)
  51. }
  52. t.Cleanup(func() { _ = os.RemoveAll(dir) })
  53. return filepath.Join(dir, "s")
  54. }
  55. func TestPizzaKVIntegration(t *testing.T) {
  56. binary := os.Getenv("PIZZAKV_BIN")
  57. if binary == "" {
  58. t.Skip("PIZZAKV_BIN is not set")
  59. }
  60. dir := t.TempDir()
  61. socket := shortPizzaKVSocket(t)
  62. database := filepath.Join(dir, "integration.pkvdb")
  63. stop := startPizzaKVTest(t, binary, socket, database)
  64. pool := waitPizzaKVPool(t, socket)
  65. schemas := NewSchemaManager(pool, "integration")
  66. tables := NewTableManager(pool, schemas, "integration")
  67. if err := schemas.CreateTable(&Schema{
  68. Name: "events",
  69. Columns: []Column{
  70. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  71. {Name: "payload", Type: "BLOB", Nullable: true},
  72. },
  73. }); err != nil {
  74. t.Fatalf("create table: %v", err)
  75. }
  76. for i := int64(1); i <= 20; i++ {
  77. payload := []byte{byte(i), 0, '|', '\r', '\n'}
  78. if err := tables.Insert("events", Row{"id": i, "payload": payload}); err != nil {
  79. t.Fatalf("insert %d: %v", i, err)
  80. }
  81. }
  82. largePayload := bytes.Repeat([]byte{0xab}, 2*1024*1024)
  83. if err := tables.Insert("events", Row{"id": int64(21), "payload": largePayload}); err != nil {
  84. t.Fatalf("insert large row: %v", err)
  85. }
  86. largeRows, err := tables.Select("events", func(row Row) bool { return row["id"] == int64(21) })
  87. if err != nil || len(largeRows) != 1 {
  88. t.Fatalf("scan large row: len=%d err=%v", len(largeRows), err)
  89. }
  90. if got, ok := largeRows[0]["payload"].([]byte); !ok || !bytes.Equal(got, largePayload) {
  91. t.Fatalf("large binary payload mismatch: len=%d type=%T", len(got), largeRows[0]["payload"])
  92. }
  93. rows, err := tables.SelectWithLimit("events", nil, 3, 2)
  94. if err != nil || len(rows) != 3 {
  95. t.Fatalf("limited select: len=%d err=%v", len(rows), err)
  96. }
  97. row, err := tables.GetByPK("events", "1")
  98. if err != nil {
  99. t.Fatalf("point read: %v", err)
  100. }
  101. if got, ok := row["payload"].([]byte); !ok || !bytes.Equal(got, []byte{1, 0, '|', '\r', '\n'}) {
  102. t.Fatalf("binary payload = %v (%T)", row["payload"], row["payload"])
  103. }
  104. if count, err := tables.CountFast("events"); err != nil || count != 21 {
  105. t.Fatalf("count = %d, err=%v", count, err)
  106. }
  107. if err := pool.Close(); err != nil {
  108. t.Fatalf("close pool: %v", err)
  109. }
  110. stop()
  111. startPizzaKVTest(t, binary, socket, database)
  112. pool = waitPizzaKVPool(t, socket)
  113. defer pool.Close()
  114. schemas = NewSchemaManager(pool, "integration")
  115. tables = NewTableManager(pool, schemas, "integration")
  116. row, err = tables.GetByPK("events", "1")
  117. if err != nil {
  118. t.Fatalf("point read after restart: %v", err)
  119. }
  120. if got, ok := row["payload"].([]byte); !ok || !bytes.Equal(got, []byte{1, 0, '|', '\r', '\n'}) {
  121. t.Fatalf("binary payload after restart = %v (%T)", row["payload"], row["payload"])
  122. }
  123. if count, err := tables.CountFast("events"); err != nil || count != 21 {
  124. t.Fatalf("count after restart = %d, err=%v", count, err)
  125. }
  126. if err := schemas.DropTable("events"); err != nil {
  127. t.Fatalf("drop table: %v", err)
  128. }
  129. if schemas.TableExists("events") {
  130. t.Fatal("table still exists")
  131. }
  132. }
  133. func TestPizzaKVCompareBatchIntegration(t *testing.T) {
  134. binary := os.Getenv("PIZZAKV_BIN")
  135. if binary == "" {
  136. t.Skip("PIZZAKV_BIN is not set")
  137. }
  138. dir := t.TempDir()
  139. socket := shortPizzaKVSocket(t)
  140. database := filepath.Join(dir, "compare.pkvdb")
  141. startPizzaKVTest(t, binary, socket, database)
  142. pool := waitPizzaKVPool(t, socket)
  143. defer pool.Close()
  144. var seededLSN uint64
  145. if err := pool.WithClient(func(c *KVClient) error {
  146. lsn, committed, err := c.CompareBatchWrite(
  147. []CompareCheck{{Key: []byte("k"), LSN: 0}},
  148. []BatchOp{{Op: batchPut, Key: []byte("k"), Value: []byte("v")}},
  149. nil,
  150. )
  151. if err != nil {
  152. return err
  153. }
  154. if !committed || lsn == 0 {
  155. return fmt.Errorf("absent check expected commit, committed=%v lsn=%d", committed, lsn)
  156. }
  157. seededLSN = lsn
  158. return nil
  159. }); err != nil {
  160. t.Fatalf("seed: %v", err)
  161. }
  162. if err := pool.WithClient(func(c *KVClient) error {
  163. lsn, committed, err := c.CompareBatchWrite(
  164. []CompareCheck{{Key: []byte("k"), LSN: 0}},
  165. []BatchOp{{Op: batchPut, Key: []byte("k"), Value: []byte("x")}},
  166. nil,
  167. )
  168. if err != nil {
  169. return err
  170. }
  171. if committed || lsn != 0 {
  172. return fmt.Errorf("expected conflict, committed=%v lsn=%d", committed, lsn)
  173. }
  174. return nil
  175. }); err != nil {
  176. t.Fatalf("stale conflict: %v", err)
  177. }
  178. const workers = 8
  179. var wins atomic.Int32
  180. var wg sync.WaitGroup
  181. errCh := make(chan error, workers)
  182. for i := 0; i < workers; i++ {
  183. wg.Add(1)
  184. go func() {
  185. defer wg.Done()
  186. if err := pool.WithClient(func(c *KVClient) error {
  187. lsn, committed, err := c.CompareBatchWrite(
  188. []CompareCheck{{Key: []byte("k"), LSN: seededLSN}},
  189. []BatchOp{{Op: batchPut, Key: []byte("k"), Value: []byte("winner")}},
  190. nil,
  191. )
  192. if err != nil {
  193. return err
  194. }
  195. if committed {
  196. if lsn <= seededLSN {
  197. return fmt.Errorf("winner lsn %d did not advance past %d", lsn, seededLSN)
  198. }
  199. wins.Add(1)
  200. }
  201. return nil
  202. }); err != nil {
  203. errCh <- err
  204. }
  205. }()
  206. }
  207. wg.Wait()
  208. close(errCh)
  209. for err := range errCh {
  210. if err != nil {
  211. t.Fatalf("concurrent: %v", err)
  212. }
  213. }
  214. if got := wins.Load(); got != 1 {
  215. t.Fatalf("exactly one transaction should win, got %d", got)
  216. }
  217. }
  218. func TestPizzaKVLegacyMigrationIntegration(t *testing.T) {
  219. binary := os.Getenv("PIZZAKV_BIN")
  220. if binary == "" {
  221. t.Skip("PIZZAKV_BIN is not set")
  222. }
  223. dir := t.TempDir()
  224. source := filepath.Join(dir, "legacy.db")
  225. destination := filepath.Join(dir, "legacy.pkvdb")
  226. socket := shortPizzaKVSocket(t)
  227. schema, err := json.Marshal(&Schema{
  228. Name: "events",
  229. Columns: []Column{{Name: "id", Type: "INTEGER", PrimaryKey: true}, {Name: "name", Type: "TEXT", Nullable: true}},
  230. PrimaryKey: "id",
  231. NextRowID: 2,
  232. })
  233. if err != nil {
  234. t.Fatal(err)
  235. }
  236. legacy := []byte(fmt.Sprintf(
  237. "W|legacy:_schema:events|%s\rW|legacy:_sys:tables|[\"events\"]\rW|legacy:_data:events:1|{\"id\":1,\"name\":\"old\",\"_rowid_\":1}\r",
  238. schema,
  239. ))
  240. if err := os.WriteFile(source, legacy, 0o600); err != nil {
  241. t.Fatal(err)
  242. }
  243. if output, err := exec.Command(binary, "-migrate="+source, "-path="+destination).CombinedOutput(); err != nil {
  244. t.Fatalf("migrate legacy database: %v\n%s", err, output)
  245. }
  246. unchanged, err := os.ReadFile(source)
  247. if err != nil {
  248. t.Fatal(err)
  249. }
  250. if !bytes.Equal(unchanged, legacy) {
  251. t.Fatal("migration modified the legacy source")
  252. }
  253. startPizzaKVTest(t, binary, socket, destination)
  254. pool := waitPizzaKVPool(t, socket)
  255. defer pool.Close()
  256. schemas := NewSchemaManager(pool, "legacy")
  257. tables := NewTableManager(pool, schemas, "legacy")
  258. row, err := tables.GetByPK("events", "1")
  259. if err != nil || row["name"] != "old" {
  260. t.Fatalf("read migrated row: row=%v err=%v", row, err)
  261. }
  262. if err := tables.Insert("events", Row{"id": int64(2), "name": "new"}); err != nil {
  263. t.Fatalf("insert binary row after migration: %v", err)
  264. }
  265. rows, err := tables.Select("events", nil)
  266. if err != nil || len(rows) != 2 {
  267. t.Fatalf("mixed legacy/binary scan: len=%d err=%v", len(rows), err)
  268. }
  269. }
  270. func BenchmarkPizzaKVStorage(b *testing.B) {
  271. binary := os.Getenv("PIZZAKV_BIN")
  272. if binary == "" {
  273. b.Skip("PIZZAKV_BIN is not set")
  274. }
  275. dir := b.TempDir()
  276. socket := shortPizzaKVSocket(b)
  277. startPizzaKVTest(b, binary, socket, filepath.Join(dir, "benchmark.pkvdb"))
  278. pool := waitPizzaKVPool(b, socket)
  279. defer pool.Close()
  280. schemas := NewSchemaManager(pool, "benchmark")
  281. tables := NewTableManager(pool, schemas, "benchmark")
  282. if err := schemas.CreateTable(&Schema{
  283. Name: "events",
  284. Columns: []Column{
  285. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  286. {Name: "symbol", Type: "TEXT", Nullable: false},
  287. {Name: "price", Type: "REAL", Nullable: false},
  288. },
  289. }); err != nil {
  290. b.Fatal(err)
  291. }
  292. rows := make([]Row, 10_000)
  293. for i := range rows {
  294. rows[i] = Row{"id": int64(i + 1), "symbol": fmt.Sprintf("PIZZA-%03d", i%100), "price": float64(i) / 100}
  295. }
  296. if n, err := tables.InsertBulk("events", rows); err != nil || n != len(rows) {
  297. b.Fatalf("seed: n=%d err=%v", n, err)
  298. }
  299. b.Run("point_read", func(b *testing.B) {
  300. b.ReportAllocs()
  301. for b.Loop() {
  302. if _, err := tables.GetByPK("events", "5000"); err != nil {
  303. b.Fatal(err)
  304. }
  305. }
  306. })
  307. b.Run("limit_10", func(b *testing.B) {
  308. b.ReportAllocs()
  309. for b.Loop() {
  310. if _, err := tables.SelectWithLimit("events", nil, 10, 0); err != nil {
  311. b.Fatal(err)
  312. }
  313. }
  314. })
  315. b.Run("scan_10000", func(b *testing.B) {
  316. b.ReportAllocs()
  317. for b.Loop() {
  318. if rows, err := tables.Select("events", nil); err != nil || len(rows) != 10_000 {
  319. b.Fatalf("scan: len=%d err=%v", len(rows), err)
  320. }
  321. }
  322. })
  323. }