schema_test.go 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335
  1. package storage
  2. import (
  3. "bufio"
  4. "fmt"
  5. "net"
  6. "strings"
  7. "sync"
  8. "testing"
  9. "time"
  10. )
  11. type testKVServer struct {
  12. mu sync.Mutex
  13. data map[string]string
  14. writes map[string]int
  15. closers []net.Conn
  16. }
  17. func newTestKVServer(t *testing.T) *testKVServer {
  18. t.Helper()
  19. return &testKVServer{
  20. data: make(map[string]string),
  21. writes: make(map[string]int),
  22. }
  23. }
  24. func newTestKVPool(kv *testKVServer, size int, timeout time.Duration) *KVPool {
  25. pool := &KVPool{
  26. pool: make(chan *KVClient, size),
  27. size: size,
  28. timeout: timeout,
  29. }
  30. for i := 0; i < size; i++ {
  31. pool.pool <- kv.client()
  32. }
  33. return pool
  34. }
  35. func (s *testKVServer) close() {
  36. s.mu.Lock()
  37. closers := append([]net.Conn(nil), s.closers...)
  38. s.mu.Unlock()
  39. for _, conn := range closers {
  40. _ = conn.Close()
  41. }
  42. }
  43. func (s *testKVServer) client() *KVClient {
  44. clientConn, serverConn := net.Pipe()
  45. s.mu.Lock()
  46. s.closers = append(s.closers, clientConn, serverConn)
  47. s.mu.Unlock()
  48. go s.handle(serverConn)
  49. return &KVClient{
  50. conn: clientConn,
  51. reader: bufio.NewReader(clientConn),
  52. writer: bufio.NewWriter(clientConn),
  53. }
  54. }
  55. func (s *testKVServer) writeCount(prefix string) int {
  56. s.mu.Lock()
  57. defer s.mu.Unlock()
  58. var count int
  59. for key, writes := range s.writes {
  60. if strings.Contains(key, prefix) {
  61. count += writes
  62. }
  63. }
  64. return count
  65. }
  66. func (s *testKVServer) hasKey(key string) bool {
  67. s.mu.Lock()
  68. defer s.mu.Unlock()
  69. _, ok := s.data[key]
  70. return ok
  71. }
  72. func (s *testKVServer) handle(conn net.Conn) {
  73. defer conn.Close()
  74. r := bufio.NewReader(conn)
  75. for {
  76. cmd, err := r.ReadString('\r')
  77. if err != nil {
  78. return
  79. }
  80. cmd = strings.TrimSuffix(cmd, "\r")
  81. resp := s.execute(cmd)
  82. if _, err := fmt.Fprintf(conn, "%s\r", resp); err != nil {
  83. return
  84. }
  85. }
  86. }
  87. func (s *testKVServer) execute(cmd string) string {
  88. s.mu.Lock()
  89. defer s.mu.Unlock()
  90. switch {
  91. case strings.HasPrefix(cmd, "write "):
  92. parts := strings.SplitN(strings.TrimPrefix(cmd, "write "), "|", 2)
  93. if len(parts) != 2 {
  94. return "error"
  95. }
  96. s.data[parts[0]] = parts[1]
  97. s.writes[parts[0]]++
  98. return "success"
  99. case strings.HasPrefix(cmd, "read "):
  100. key := strings.TrimPrefix(cmd, "read ")
  101. value, ok := s.data[key]
  102. if !ok {
  103. return "error"
  104. }
  105. return value
  106. case strings.HasPrefix(cmd, "delete "):
  107. key := strings.TrimPrefix(cmd, "delete ")
  108. delete(s.data, key)
  109. return "success"
  110. case strings.HasPrefix(cmd, "reads "):
  111. prefix := strings.TrimPrefix(cmd, "reads ")
  112. values := make([]string, 0)
  113. for key, value := range s.data {
  114. if strings.HasPrefix(key, prefix) {
  115. values = append(values, value)
  116. }
  117. }
  118. return strings.Join(values, "\n")
  119. default:
  120. return "error"
  121. }
  122. }
  123. func TestInsertDoesNotRewriteSchemaForRowIDUpdates(t *testing.T) {
  124. kv := newTestKVServer(t)
  125. defer kv.close()
  126. pool := newTestKVPool(kv, 2, 5*time.Second)
  127. defer pool.Close()
  128. schemas := NewSchemaManager(pool, "testdb")
  129. tables := NewTableManager(pool, schemas, "testdb")
  130. err := schemas.CreateTable(&Schema{
  131. Name: "users",
  132. Columns: []Column{
  133. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  134. {Name: "name", Type: "TEXT", Nullable: true},
  135. },
  136. })
  137. if err != nil {
  138. t.Fatalf("create table: %v", err)
  139. }
  140. initialSchemaWrites := kv.writeCount(":_schema:")
  141. if initialSchemaWrites != 1 {
  142. t.Fatalf("expected create table to write schema once, got %d", initialSchemaWrites)
  143. }
  144. for i := int64(1); i <= 3; i++ {
  145. err := tables.Insert("users", Row{"id": i, "name": fmt.Sprintf("user-%d", i)})
  146. if err != nil {
  147. t.Fatalf("insert %d: %v", i, err)
  148. }
  149. }
  150. if got := kv.writeCount(":_schema:"); got != initialSchemaWrites {
  151. t.Fatalf("expected inserts not to rewrite schema, got %d schema writes", got)
  152. }
  153. if got := kv.writeCount(":_sys:rowid:"); got != 0 {
  154. t.Fatalf("expected no rowid counter writes for inserts, got %d", got)
  155. }
  156. }
  157. func TestRowIDIsDerivedFromRowsAfterRestart(t *testing.T) {
  158. kv := newTestKVServer(t)
  159. defer kv.close()
  160. pool := newTestKVPool(kv, 2, 5*time.Second)
  161. defer pool.Close()
  162. schemas := NewSchemaManager(pool, "testdb")
  163. tables := NewTableManager(pool, schemas, "testdb")
  164. err := schemas.CreateTable(&Schema{
  165. Name: "events",
  166. Columns: []Column{
  167. {Name: "name", Type: "TEXT", Nullable: true},
  168. },
  169. })
  170. if err != nil {
  171. t.Fatalf("create table: %v", err)
  172. }
  173. for i := 1; i <= 2; i++ {
  174. err := tables.Insert("events", Row{"name": fmt.Sprintf("event-%d", i)})
  175. if err != nil {
  176. t.Fatalf("insert %d: %v", i, err)
  177. }
  178. }
  179. // Simulate a process restart: new managers have empty in-memory ROWID state
  180. // but the same durable KV rows.
  181. restartedSchemas := NewSchemaManager(pool, "testdb")
  182. restartedTables := NewTableManager(pool, restartedSchemas, "testdb")
  183. if err := restartedTables.Insert("events", Row{"name": "event-3"}); err != nil {
  184. t.Fatalf("insert after restart: %v", err)
  185. }
  186. if !kv.hasKey("testdb:_data:events:3") {
  187. t.Fatalf("expected restart insert to continue at rowid 3")
  188. }
  189. if got := kv.writeCount(":_sys:rowid:"); got != 0 {
  190. t.Fatalf("expected no rowid counter writes, got %d", got)
  191. }
  192. }
  193. func TestInsertDoesNotWriteDurableIndexEntries(t *testing.T) {
  194. kv := newTestKVServer(t)
  195. defer kv.close()
  196. pool := newTestKVPool(kv, 2, 5*time.Second)
  197. defer pool.Close()
  198. schemas := NewSchemaManager(pool, "testdb")
  199. tables := NewTableManager(pool, schemas, "testdb")
  200. err := schemas.CreateTable(&Schema{
  201. Name: "users",
  202. Columns: []Column{
  203. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  204. {Name: "status", Type: "TEXT", Nullable: false},
  205. },
  206. })
  207. if err != nil {
  208. t.Fatalf("create table: %v", err)
  209. }
  210. err = schemas.CreateIndex(&Index{
  211. Name: "idx_users_status",
  212. Table: "users",
  213. Columns: []IndexColumn{
  214. {Name: "status"},
  215. },
  216. })
  217. if err != nil {
  218. t.Fatalf("create index: %v", err)
  219. }
  220. for i := int64(1); i <= 3; i++ {
  221. status := "active"
  222. if i == 2 {
  223. status = "inactive"
  224. }
  225. err := tables.Insert("users", Row{"id": i, "status": status})
  226. if err != nil {
  227. t.Fatalf("insert %d: %v", i, err)
  228. }
  229. }
  230. if got := kv.writeCount(":idx:"); got != 0 {
  231. t.Fatalf("expected no durable index entry writes, got %d", got)
  232. }
  233. rows, err := tables.SelectByIndex("users", "idx_users_status", "active")
  234. if err != nil {
  235. t.Fatalf("select by index: %v", err)
  236. }
  237. if len(rows) != 2 {
  238. t.Fatalf("expected 2 active rows from derived index, got %d", len(rows))
  239. }
  240. }
  241. func TestIndexIsDerivedFromRowsAfterRestart(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. err := schemas.CreateTable(&Schema{
  249. Name: "users",
  250. Columns: []Column{
  251. {Name: "id", Type: "INTEGER", Nullable: false, PrimaryKey: true},
  252. {Name: "status", Type: "TEXT", Nullable: false},
  253. },
  254. })
  255. if err != nil {
  256. t.Fatalf("create table: %v", err)
  257. }
  258. err = schemas.CreateIndex(&Index{
  259. Name: "idx_users_status",
  260. Table: "users",
  261. Columns: []IndexColumn{
  262. {Name: "status"},
  263. },
  264. })
  265. if err != nil {
  266. t.Fatalf("create index: %v", err)
  267. }
  268. for i := int64(1); i <= 3; i++ {
  269. status := "active"
  270. if i == 3 {
  271. status = "inactive"
  272. }
  273. if err := tables.Insert("users", Row{"id": i, "status": status}); err != nil {
  274. t.Fatalf("insert %d: %v", i, err)
  275. }
  276. }
  277. restartedSchemas := NewSchemaManager(pool, "testdb")
  278. restartedTables := NewTableManager(pool, restartedSchemas, "testdb")
  279. rows, err := restartedTables.SelectByIndex("users", "idx_users_status", "active")
  280. if err != nil {
  281. t.Fatalf("select by index after restart: %v", err)
  282. }
  283. if len(rows) != 2 {
  284. t.Fatalf("expected 2 active rows from restart-derived index, got %d", len(rows))
  285. }
  286. if got := kv.writeCount(":idx:"); got != 0 {
  287. t.Fatalf("expected no durable index entry writes, got %d", got)
  288. }
  289. }