table.go 55 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029
  1. package storage
  2. import (
  3. "errors"
  4. "fmt"
  5. "hash/fnv"
  6. "math"
  7. "strconv"
  8. "strings"
  9. "sync"
  10. "time"
  11. )
  12. // Row represents a database row.
  13. type Row map[string]interface{}
  14. // TableManager manages table data operations.
  15. type TableManager struct {
  16. pool *KVPool
  17. schema *SchemaManager
  18. database string
  19. cacheMu sync.RWMutex
  20. indexCache map[string]map[string][]int64 // index name → indexed value → rowids
  21. indexTable map[string]string // index name → table name
  22. rowKeyCache map[string]map[int64]string // table name → rowid → primary key
  23. // disabledIndexes prevents a concurrent lookup from rebuilding an index
  24. // after DROP has cleared it but before the schema entry is removed.
  25. disabledIndexes map[string]bool
  26. // counts holds exact per-table row counts for the COUNT(*) fast path.
  27. // It is derived lazily from durable rows on first use and maintained
  28. // incrementally by Insert/InsertBulk/Delete thereafter.
  29. counts map[string]int
  30. countsInit map[string]bool
  31. countGeneration map[string]time.Time
  32. // stripes are deterministic per-key locks used by point operations
  33. // (GetByPK/Insert/UpdateByPK/DeleteByPK). They replace the table-wide lock
  34. // so point operations on different keys of the same table proceed
  35. // concurrently. Index is derived from the full data key (database+table+pk).
  36. stripes [64]sync.Mutex
  37. // generations tracks full-table scans. predicateGenerations narrows indexed
  38. // equality reads to one index value so unrelated writes do not conflict.
  39. genMu sync.Mutex
  40. generations map[string]uint64
  41. predicateGenerations map[string]uint64
  42. }
  43. // NewTableManager creates a new table manager.
  44. func NewTableManager(pool *KVPool, schema *SchemaManager, database string) *TableManager {
  45. return &TableManager{
  46. pool: pool,
  47. schema: schema,
  48. database: database,
  49. indexCache: make(map[string]map[string][]int64),
  50. indexTable: make(map[string]string),
  51. rowKeyCache: make(map[string]map[int64]string),
  52. disabledIndexes: make(map[string]bool),
  53. counts: make(map[string]int),
  54. countsInit: make(map[string]bool),
  55. countGeneration: make(map[string]time.Time),
  56. generations: make(map[string]uint64),
  57. predicateGenerations: make(map[string]uint64),
  58. }
  59. }
  60. func (m *TableManager) tableLock(key string) *sync.RWMutex {
  61. return m.schema.tableLock(key)
  62. }
  63. // stripeKey returns the deterministic striped lock for a point operation on the
  64. // given full data key. Point operations on different keys therefore serialize
  65. // independently, while operations on the same key are mutually exclusive.
  66. func (m *TableManager) stripeKey(key string) *sync.Mutex {
  67. h := fnv.New32a()
  68. h.Write([]byte(key))
  69. return &m.stripes[h.Sum32()%uint32(len(m.stripes))]
  70. }
  71. // generation returns the current in-process generation for a table. It is
  72. // bumped on every committed write and captured by transaction scans.
  73. func (m *TableManager) generation(table string) uint64 {
  74. key := strings.ToLower(table)
  75. m.genMu.Lock()
  76. g := m.generations[key]
  77. m.genMu.Unlock()
  78. return g
  79. }
  80. // bumpGeneration advances a table's generation. Callers hold the table gate in
  81. // shared mode for point writes or exclusive mode for scan-based writes.
  82. func (m *TableManager) bumpGeneration(table string) {
  83. key := strings.ToLower(table)
  84. m.genMu.Lock()
  85. m.generations[key]++
  86. m.genMu.Unlock()
  87. }
  88. func indexPredicateKey(table, indexName, value string) string {
  89. return strings.ToLower(table) + "\x00" + strings.ToLower(indexName) + "\x00" + value
  90. }
  91. func indexPredicateWildcardKey(table string) string {
  92. return strings.ToLower(table) + "\x00*"
  93. }
  94. func (m *TableManager) predicateSnapshot(table, indexName, value string) (string, uint64, string, uint64) {
  95. valueKey := indexPredicateKey(table, indexName, value)
  96. wildcardKey := indexPredicateWildcardKey(table)
  97. m.genMu.Lock()
  98. valueGen := m.predicateGenerations[valueKey]
  99. wildcardGen := m.predicateGenerations[wildcardKey]
  100. m.genMu.Unlock()
  101. return valueKey, valueGen, wildcardKey, wildcardGen
  102. }
  103. func (m *TableManager) predicateGeneration(key string) uint64 {
  104. m.genMu.Lock()
  105. gen := m.predicateGenerations[key]
  106. m.genMu.Unlock()
  107. return gen
  108. }
  109. func (m *TableManager) bumpIndexPredicates(table string, rows ...Row) {
  110. indexes, err := m.schema.ListTableIndexes(table)
  111. if err != nil || len(indexes) == 0 {
  112. return
  113. }
  114. m.genMu.Lock()
  115. defer m.genMu.Unlock()
  116. for _, row := range rows {
  117. if row == nil {
  118. continue
  119. }
  120. for _, index := range indexes {
  121. columns := make([]string, len(index.Columns))
  122. for i, column := range index.Columns {
  123. columns[i] = column.Name
  124. }
  125. value := formatIndexValue(m.buildIndexValue(row, columns))
  126. m.predicateGenerations[indexPredicateKey(table, index.Name, value)]++
  127. }
  128. }
  129. }
  130. func (m *TableManager) bumpIndexPredicateWildcard(table string) {
  131. m.genMu.Lock()
  132. m.predicateGenerations[indexPredicateWildcardKey(table)]++
  133. m.genMu.Unlock()
  134. }
  135. // compareWritePoint issues a single CompareBatchWrite against the pooled KV.
  136. // It returns whether the compare checks held and the ops committed.
  137. func (m *TableManager) compareWritePoint(checks []CompareCheck, ops []BatchOp) (bool, error) {
  138. var committed bool
  139. err := m.pool.WithClient(func(c *KVClient) error {
  140. _, ok, err := c.CompareBatchWrite(checks, ops, nil)
  141. committed = ok
  142. return err
  143. })
  144. return committed, err
  145. }
  146. // invalidateCache removes a table's derived in-memory indexes.
  147. func (m *TableManager) invalidateCache(table string) {
  148. m.cacheMu.Lock()
  149. key := strings.ToLower(table)
  150. delete(m.rowKeyCache, key)
  151. delete(m.counts, key)
  152. delete(m.countsInit, key)
  153. delete(m.countGeneration, key)
  154. for indexName, tableName := range m.indexTable {
  155. if tableName == key {
  156. delete(m.indexCache, indexName)
  157. delete(m.indexTable, indexName)
  158. }
  159. }
  160. m.cacheMu.Unlock()
  161. }
  162. // InvalidateCache is the exported version for use by the executor.
  163. func (m *TableManager) InvalidateCache(table string) {
  164. m.invalidateCache(table)
  165. }
  166. // rowVisitFunc is invoked for each decoded row in a streaming scan. Return
  167. // stop=true to end the scan early; a non-nil error aborts the scan.
  168. type rowVisitFunc func(Row) (stop bool, err error)
  169. // scanRows streams the rows of table by scanning durable KV rows one page at a
  170. // time under a single pooled client. Each page is decoded as it arrives and
  171. // passed to fn, which may stop the scan early. The cursor is always closed and
  172. // the client always returned to the pool, even on error.
  173. func (m *TableManager) scanRows(table string, fn rowVisitFunc) error {
  174. return m.scanRowsWithPageSize(table, scanPageSize, fn)
  175. }
  176. func (m *TableManager) scanRowsWithPageSize(table string, pageSize uint32, fn rowVisitFunc) error {
  177. return m.pool.WithClient(func(client *KVClient) (retErr error) {
  178. cursor, err := client.ScanWithLimit([]byte(m.dataPrefix(table)), pageSize)
  179. if err != nil {
  180. return err
  181. }
  182. defer func() {
  183. if err := cursor.Close(); retErr == nil {
  184. retErr = err
  185. }
  186. }()
  187. for {
  188. entries, done, err := cursor.Next()
  189. if err != nil {
  190. return err
  191. }
  192. for _, e := range entries {
  193. row, err := decodeRow(e.Value)
  194. if err != nil {
  195. return err
  196. }
  197. stop, err := fn(row)
  198. if err != nil {
  199. return err
  200. }
  201. if stop {
  202. return nil
  203. }
  204. }
  205. if done {
  206. return nil
  207. }
  208. }
  209. })
  210. }
  211. // rowWithLSNVisitFunc is like rowVisitFunc but also passes the durable row LSN.
  212. type rowWithLSNVisitFunc func(row Row, lsn uint64) (stop bool, err error)
  213. // scanRowsWithLSN streams a table's rows together with their durable LSNs. It
  214. // is used by buffered transactions to capture a per-row read set for
  215. // compare-and-swap validation at commit.
  216. func (m *TableManager) scanRowsWithLSN(table string, fn rowWithLSNVisitFunc) error {
  217. return m.pool.WithClient(func(client *KVClient) (retErr error) {
  218. cursor, err := client.ScanWithLimit([]byte(m.dataPrefix(table)), scanPageSize)
  219. if err != nil {
  220. return err
  221. }
  222. defer func() {
  223. if err := cursor.Close(); retErr == nil {
  224. retErr = err
  225. }
  226. }()
  227. for {
  228. entries, done, err := cursor.Next()
  229. if err != nil {
  230. return err
  231. }
  232. for _, e := range entries {
  233. row, err := decodeRow(e.Value)
  234. if err != nil {
  235. return err
  236. }
  237. stop, err := fn(row, e.LSN)
  238. if err != nil {
  239. return err
  240. }
  241. if stop {
  242. return nil
  243. }
  244. }
  245. if done {
  246. return nil
  247. }
  248. }
  249. })
  250. }
  251. // scanCountKeys counts the durable rows of table using a key-only scan so row
  252. // values are never pulled across the wire. It is used for first-time COUNT(*)
  253. // derivation.
  254. func (m *TableManager) scanCountKeys(table string) (int, error) {
  255. count := 0
  256. err := m.pool.WithClient(func(client *KVClient) (retErr error) {
  257. cursor, err := client.ScanKeys([]byte(m.dataPrefix(table)))
  258. if err != nil {
  259. return err
  260. }
  261. defer func() {
  262. if err := cursor.Close(); retErr == nil {
  263. retErr = err
  264. }
  265. }()
  266. for {
  267. entries, done, err := cursor.Next()
  268. if err != nil {
  269. return err
  270. }
  271. count += len(entries)
  272. if done {
  273. return nil
  274. }
  275. }
  276. })
  277. return count, err
  278. }
  279. // CountFast returns the exact number of rows in a table. The count is derived
  280. // from durable rows on first use (recovering across restarts) and then
  281. // maintained incrementally by the write paths, so repeated COUNT(*) queries
  282. // avoid a full table scan. It intentionally does not persist a counter to KV:
  283. // the KV layer has no atomic increment primitive, and a durable counter that
  284. // could diverge from the rows on crash would be worse than a lazily-derived,
  285. // always-exact value. The cost is one key-only scan the first time COUNT(*) is
  286. // issued after startup.
  287. func (m *TableManager) CountFast(table string) (int, error) {
  288. key := strings.ToLower(table)
  289. tableSchema, err := m.schema.GetSchema(table)
  290. if err != nil {
  291. return 0, err
  292. }
  293. // Cached read path: if the count is initialized for the current table
  294. // generation, return it without taking any table lock.
  295. m.cacheMu.Lock()
  296. if m.countGeneration[key].Equal(tableSchema.CreatedAt) {
  297. if m.countsInit[key] {
  298. n := m.counts[key]
  299. m.cacheMu.Unlock()
  300. return n, nil
  301. }
  302. } else {
  303. delete(m.counts, key)
  304. delete(m.countsInit, key)
  305. m.countGeneration[key] = tableSchema.CreatedAt
  306. }
  307. m.cacheMu.Unlock()
  308. // First derivation: hold the table write lock so the key-only scan is
  309. // exact against concurrent writes.
  310. tl := m.tableLock(key)
  311. tl.Lock()
  312. defer tl.Unlock()
  313. m.cacheMu.Lock()
  314. if !m.countGeneration[key].Equal(tableSchema.CreatedAt) {
  315. delete(m.counts, key)
  316. delete(m.countsInit, key)
  317. m.countGeneration[key] = tableSchema.CreatedAt
  318. }
  319. if m.countsInit[key] {
  320. n := m.counts[key]
  321. m.cacheMu.Unlock()
  322. return n, nil
  323. }
  324. m.cacheMu.Unlock()
  325. count, err := m.scanCountKeys(table)
  326. if err != nil {
  327. return 0, err
  328. }
  329. m.cacheMu.Lock()
  330. m.counts[key] = count
  331. m.countsInit[key] = true
  332. m.cacheMu.Unlock()
  333. return count, nil
  334. }
  335. // countInitialized reports whether the derived count cache is initialized for
  336. // the table at this instant. Write paths capture it before their KV write so a
  337. // concurrent first derivation does not double-count the just-written row.
  338. func (m *TableManager) countInitialized(table string) bool {
  339. key := strings.ToLower(table)
  340. m.cacheMu.Lock()
  341. init := m.countsInit[key]
  342. m.cacheMu.Unlock()
  343. return init
  344. }
  345. // incrCount adjusts the derived per-table row count. wasInit reports whether
  346. // the count was already initialized before the corresponding write, so a count
  347. // that was not yet initialized is left to be re-derived from durable rows
  348. // (which already reflect the write) on next use.
  349. func (m *TableManager) incrCount(table string, generation time.Time, delta int, wasInit bool) {
  350. key := strings.ToLower(table)
  351. m.cacheMu.Lock()
  352. if !m.countGeneration[key].Equal(generation) {
  353. delete(m.counts, key)
  354. delete(m.countsInit, key)
  355. m.countGeneration[key] = generation
  356. }
  357. if wasInit && m.countsInit[key] {
  358. m.counts[key] += delta
  359. }
  360. if !wasInit {
  361. // The count was not initialized before the write, so a concurrent first
  362. // derivation may have missed the just-written row. Invalidate to force
  363. // an exact re-derivation on the next COUNT(*).
  364. delete(m.countsInit, key)
  365. }
  366. m.cacheMu.Unlock()
  367. }
  368. // dataKey returns the key for a row.
  369. func (m *TableManager) dataKey(table, pk string) string {
  370. return fmt.Sprintf("%s:_data:%s:%s", m.database, strings.ToLower(table), pk)
  371. }
  372. // dataPrefix returns the prefix for all rows in a table.
  373. func (m *TableManager) dataPrefix(table string) string {
  374. return fmt.Sprintf("%s:_data:%s:", m.database, strings.ToLower(table))
  375. }
  376. // prepareInsert validates and normalizes an insert row, generating the ROWID
  377. // when required. It returns the normalized row and the full data key without
  378. // writing anything, so both the autocommit path (compare-and-swap) and the
  379. // buffered transaction path (staging) can share it. The input row has its
  380. // primary key populated as a side effect.
  381. func (m *TableManager) prepareInsert(table string, row Row) (Row, string, error) {
  382. schema, err := m.schema.GetSchema(table)
  383. if err != nil {
  384. return nil, "", err
  385. }
  386. // Get primary key value
  387. pkValue, ok := row[schema.PrimaryKey]
  388. if !ok {
  389. for k, v := range row {
  390. if strings.EqualFold(k, schema.PrimaryKey) {
  391. pkValue = v
  392. ok = true
  393. break
  394. }
  395. }
  396. }
  397. pkCol, _ := schema.GetColumn(schema.PrimaryKey)
  398. isIntegerPK := pkCol != nil && isIntegerType(pkCol.Type)
  399. var rowid int64
  400. if !ok || pkValue == nil {
  401. if isIntegerPK || !ok {
  402. rowid, err = m.schema.GetNextRowID(table)
  403. if err != nil {
  404. return nil, "", err
  405. }
  406. pkValue = rowid
  407. row[schema.PrimaryKey] = rowid
  408. ok = true
  409. } else {
  410. return nil, "", fmt.Errorf("missing primary key: %s", schema.PrimaryKey)
  411. }
  412. } else if isIntegerPK {
  413. switch v := pkValue.(type) {
  414. case int64:
  415. rowid = v
  416. case float64:
  417. if math.Trunc(v) != v {
  418. return nil, "", fmt.Errorf("invalid integer primary key: %v", v)
  419. }
  420. rowid = int64(v)
  421. case int:
  422. rowid = int64(v)
  423. default:
  424. rowid = 0
  425. }
  426. if rowid > 0 {
  427. if err := m.schema.UpdateMaxRowID(table, rowid); err != nil {
  428. return nil, "", err
  429. }
  430. }
  431. }
  432. pk := fmt.Sprintf("%v", pkValue)
  433. for _, col := range schema.Columns {
  434. if !col.Nullable && col.Default == nil {
  435. val, hasVal := row[col.Name]
  436. if !hasVal {
  437. for k, v := range row {
  438. if strings.EqualFold(k, col.Name) {
  439. val = v
  440. hasVal = true
  441. break
  442. }
  443. }
  444. }
  445. if !hasVal || val == nil {
  446. return nil, "", fmt.Errorf("missing required column: %s", col.Name)
  447. }
  448. }
  449. }
  450. normalizedRow := make(Row)
  451. for _, col := range schema.Columns {
  452. for k, v := range row {
  453. if strings.EqualFold(k, col.Name) {
  454. normalizedRow[col.Name] = v
  455. break
  456. }
  457. }
  458. }
  459. for _, col := range schema.Columns {
  460. if _, ok := normalizedRow[col.Name]; !ok && col.Default != nil {
  461. normalizedRow[col.Name] = col.Default
  462. }
  463. }
  464. if rowid > 0 {
  465. normalizedRow["_rowid_"] = rowid
  466. } else {
  467. newRowID, err := m.schema.GetNextRowID(table)
  468. if err != nil {
  469. return nil, "", err
  470. }
  471. normalizedRow["_rowid_"] = newRowID
  472. }
  473. return normalizedRow, m.dataKey(table, pk), nil
  474. }
  475. // Insert inserts a new row. The duplicate check and write are one atomic
  476. // compare-and-swap so concurrent inserts of the same primary key cannot both
  477. // persist. The shared table gate is acquired before the key's striped lock so
  478. // a queued scan writer cannot invert the lock order with point operations.
  479. func (m *TableManager) Insert(table string, row Row) error {
  480. _, err := m.InsertWithRowID(table, row)
  481. return err
  482. }
  483. // InsertWithRowID returns the actual row identifier, never a concurrent counter.
  484. func (m *TableManager) InsertWithRowID(table string, row Row) (int64, error) {
  485. nr, key, err := m.prepareInsert(table, row)
  486. if err != nil {
  487. return 0, err
  488. }
  489. rowid, _ := rowIDFromRow(nr)
  490. data, err := encodeRow(nr)
  491. if err != nil {
  492. return 0, fmt.Errorf("failed to serialize row: %w", err)
  493. }
  494. wasInit := m.countInitialized(table)
  495. // Point writers share this gate with each other. Transaction commits and
  496. // scan-based writes take it exclusively, so generation validation and cache
  497. // publication are ordered without serializing writes to different keys.
  498. // A UNIQUE index forces the exclusive gate so the validating scan below
  499. // cannot race another writer (or a concurrent CREATE UNIQUE INDEX).
  500. unlock, uniq, err := m.lockForWrite(table)
  501. if err != nil {
  502. return 0, err
  503. }
  504. defer unlock()
  505. if uniq {
  506. if err := m.validateUniqueRows(table, []Row{nr}, nil); err != nil {
  507. return 0, err
  508. }
  509. }
  510. st := m.stripeKey(key)
  511. st.Lock()
  512. defer st.Unlock()
  513. committed, err := m.compareWritePoint(
  514. []CompareCheck{{Key: []byte(key), LSN: 0}},
  515. []BatchOp{{Op: batchPut, Key: []byte(key), Value: data}},
  516. )
  517. if err != nil {
  518. return 0, err
  519. }
  520. if !committed {
  521. return 0, fmt.Errorf("duplicate primary key: %v", row[schemaPrimaryKey(m.schema, table)])
  522. }
  523. m.updateIndexesForRow(table, nr, true)
  524. // Publish derived index state before its generations. A reader that races
  525. // with publication either sees the old generation and aborts or sees the
  526. // complete new state.
  527. m.bumpIndexPredicates(table, nr)
  528. m.bumpGeneration(table)
  529. if schema, serr := m.schema.GetSchema(table); serr == nil {
  530. m.incrCount(table, schema.CreatedAt, 1, wasInit)
  531. }
  532. return rowid, nil
  533. }
  534. func schemaPrimaryKey(s *SchemaManager, table string) string {
  535. schema, err := s.GetSchema(table)
  536. if err != nil {
  537. return "_rowid_"
  538. }
  539. return schema.PrimaryKey
  540. }
  541. // bulkBatchByteBudget bounds a single atomic BATCH_WRITE payload below the
  542. // PKBFI frame limit so a bulk insert never emits a frame the server rejects.
  543. // Each op contributes 12 header bytes plus its key and value.
  544. const bulkBatchByteBudget = 60 * 1024 * 1024
  545. // chunkBatchOps splits ops into atomic BATCH_WRITE chunks bounded by both the
  546. // PKBFI operation-count limit and the frame-size limit. Each chunk is a slice
  547. // of the backing array, valid until the next append to ops.
  548. func chunkBatchOps(ops []BatchOp) [][]BatchOp {
  549. var chunks [][]BatchOp
  550. for i := 0; i < len(ops); {
  551. end := i + maxOperations
  552. if end > len(ops) {
  553. end = len(ops)
  554. }
  555. bytes := 0
  556. j := i
  557. for j < end {
  558. sz := 12 + len(ops[j].Key) + len(ops[j].Value)
  559. if j > i && bytes+sz > bulkBatchByteBudget {
  560. break
  561. }
  562. bytes += sz
  563. j++
  564. }
  565. if j == i {
  566. j = i + 1
  567. }
  568. chunks = append(chunks, ops[i:j])
  569. i = j
  570. }
  571. return chunks
  572. }
  573. // InsertBulk inserts multiple rows efficiently using atomic BATCH_WRITE chunks
  574. // bounded by the PKBFI operation-count and frame-size limits. Skips per-row
  575. // duplicate checks (caller must ensure uniqueness). Used by INSERT ... SELECT.
  576. func (m *TableManager) InsertBulk(table string, rows []Row) (int, error) {
  577. count, _, err := m.InsertBulkWithLastRowID(table, rows)
  578. return count, err
  579. }
  580. // InsertBulkWithLastRowID is InsertBulk, additionally returning the ROWID of the
  581. // last persisted row (0 when nothing persisted).
  582. func (m *TableManager) InsertBulkWithLastRowID(table string, rows []Row) (int, int64, error) {
  583. if len(rows) == 0 {
  584. return 0, 0, nil
  585. }
  586. tl := m.tableLock(table)
  587. tl.Lock()
  588. defer tl.Unlock()
  589. schema, err := m.schema.GetSchema(table)
  590. if err != nil {
  591. return 0, 0, err
  592. }
  593. wasInit := m.countInitialized(table)
  594. pkCol, _ := schema.GetColumn(schema.PrimaryKey)
  595. isIntegerPK := pkCol != nil && isIntegerType(pkCol.Type)
  596. // Normalize rows and assign _rowid_.
  597. normalized := make([]Row, 0, len(rows))
  598. var maxRowID int64
  599. for _, row := range rows {
  600. nr := make(Row)
  601. for _, col := range schema.Columns {
  602. for k, v := range row {
  603. if strings.EqualFold(k, col.Name) {
  604. nr[col.Name] = v
  605. break
  606. }
  607. }
  608. }
  609. for _, col := range schema.Columns {
  610. if _, ok := nr[col.Name]; !ok && col.Default != nil {
  611. nr[col.Name] = col.Default
  612. }
  613. }
  614. var rowid int64
  615. var hasRowid bool
  616. if isIntegerPK {
  617. switch v := nr[schema.PrimaryKey].(type) {
  618. case float64:
  619. if math.Trunc(v) != v {
  620. return 0, 0, fmt.Errorf("invalid integer primary key: %v", v)
  621. }
  622. rowid = int64(v)
  623. hasRowid = true
  624. case int64:
  625. rowid = v
  626. hasRowid = true
  627. case int:
  628. rowid = int64(v)
  629. hasRowid = true
  630. }
  631. }
  632. if !hasRowid {
  633. // Auto-generate for an INTEGER primary key or the synthetic _rowid_,
  634. // mirroring prepareInsert so INSERT ... SELECT works on AUTOINCREMENT
  635. // tables.
  636. if isIntegerPK || schema.PrimaryKey == "_rowid_" {
  637. rowid, err = m.schema.GetNextRowID(table)
  638. if err != nil {
  639. return 0, 0, err
  640. }
  641. if isIntegerPK {
  642. nr[schema.PrimaryKey] = rowid
  643. } else {
  644. nr["_rowid_"] = rowid
  645. }
  646. } else {
  647. pk, ok := nr[schema.PrimaryKey]
  648. if !ok || pk == nil {
  649. return 0, 0, fmt.Errorf("missing primary key: %s", schema.PrimaryKey)
  650. }
  651. rowid, err = m.schema.GetNextRowID(table)
  652. if err != nil {
  653. return 0, 0, err
  654. }
  655. }
  656. }
  657. nr["_rowid_"] = rowid
  658. if rowid > maxRowID {
  659. maxRowID = rowid
  660. }
  661. normalized = append(normalized, nr)
  662. }
  663. if maxRowID > 0 {
  664. m.schema.UpdateMaxRowID(table, maxRowID)
  665. }
  666. // Serialize all rows into batch operations. A parallel slice keeps the
  667. // normalized Row for each op for index maintenance after the write.
  668. ops := make([]BatchOp, 0, len(normalized))
  669. encoded := make([]Row, 0, len(normalized))
  670. for _, nr := range normalized {
  671. pk := fmt.Sprintf("%v", nr[schema.PrimaryKey])
  672. data, err := encodeRow(nr)
  673. if err != nil {
  674. return 0, 0, err
  675. }
  676. ops = append(ops, BatchOp{Op: batchPut, Key: []byte(m.dataKey(table, pk)), Value: data})
  677. encoded = append(encoded, nr)
  678. }
  679. keys := make([][]byte, len(ops))
  680. seen := make(map[string]struct{}, len(ops))
  681. for i, op := range ops {
  682. key := string(op.Key)
  683. if _, duplicate := seen[key]; duplicate {
  684. return 0, 0, fmt.Errorf("duplicate primary key: %s", key)
  685. }
  686. seen[key] = struct{}{}
  687. keys[i] = op.Key
  688. }
  689. existing := make([]bool, len(keys))
  690. if err := m.pool.WithClient(func(client *KVClient) error {
  691. for start := 0; start < len(keys); start += maxOperations {
  692. end := start + maxOperations
  693. if end > len(keys) {
  694. end = len(keys)
  695. }
  696. found, err := client.ExistsMany(keys[start:end])
  697. if err != nil {
  698. return err
  699. }
  700. copy(existing[start:end], found)
  701. }
  702. return nil
  703. }); err != nil {
  704. return 0, 0, err
  705. }
  706. for i, found := range existing {
  707. if found {
  708. return 0, 0, fmt.Errorf("duplicate primary key: %s", ops[i].Key)
  709. }
  710. }
  711. if err := m.validateUniqueRows(table, encoded, nil); err != nil {
  712. return 0, 0, err
  713. }
  714. // Write rows in atomic BATCH_WRITE chunks. Maintain in-memory indexes only
  715. // for rows that actually persisted, so a partial failure cannot leave an
  716. // already-built index stale.
  717. var firstErr error
  718. numOK := 0
  719. for _, chunk := range chunkBatchOps(ops) {
  720. err := m.pool.WithClient(func(c *KVClient) error {
  721. _, err := c.BatchWrite(chunk, nil)
  722. return err
  723. })
  724. if err != nil {
  725. if firstErr == nil {
  726. firstErr = err
  727. }
  728. break
  729. }
  730. numOK += len(chunk)
  731. }
  732. // Maintain indexes for the rows that persisted (the first numOK ops).
  733. for i := 0; i < numOK; i++ {
  734. m.updateIndexesForRow(table, encoded[i], true)
  735. }
  736. if numOK > 0 {
  737. m.bumpGeneration(table)
  738. m.bumpIndexPredicates(table, encoded[:numOK]...)
  739. }
  740. m.incrCount(table, schema.CreatedAt, numOK, wasInit)
  741. var lastRowID int64
  742. if numOK > 0 {
  743. lastRowID, _ = rowIDFromRow(encoded[numOK-1])
  744. }
  745. return numOK, lastRowID, firstErr
  746. }
  747. // updateIndexesForRow adds or removes entries from already-built in-memory
  748. // indexes. Index entries are rebuildable from durable row data, so this method
  749. // intentionally does not write idx:* keys to KV.
  750. func (m *TableManager) updateIndexesForRow(table string, row Row, add bool) {
  751. indexes, err := m.schema.ListTableIndexes(table)
  752. if err != nil || len(indexes) == 0 {
  753. return
  754. }
  755. rowid, ok := rowIDFromRow(row)
  756. if !ok {
  757. return
  758. }
  759. tableKey := strings.ToLower(table)
  760. tableSchema, schemaErr := m.schema.GetSchema(table)
  761. if schemaErr == nil {
  762. m.cacheMu.Lock()
  763. if keys, initialized := m.rowKeyCache[tableKey]; initialized {
  764. if add {
  765. keys[rowid] = fmt.Sprintf("%v", row[tableSchema.PrimaryKey])
  766. } else {
  767. delete(keys, rowid)
  768. }
  769. }
  770. m.cacheMu.Unlock()
  771. }
  772. for _, idx := range indexes {
  773. indexName := strings.ToLower(idx.Name)
  774. m.cacheMu.RLock()
  775. _, initialized := m.indexCache[indexName]
  776. m.cacheMu.RUnlock()
  777. if !initialized {
  778. continue
  779. }
  780. columns := make([]string, len(idx.Columns))
  781. for i, col := range idx.Columns {
  782. columns[i] = col.Name
  783. }
  784. colValue := m.buildIndexValue(row, columns)
  785. if add {
  786. m.AddIndexEntry(idx.Name, colValue, rowid)
  787. } else {
  788. m.RemoveIndexEntry(idx.Name, colValue, rowid)
  789. }
  790. }
  791. }
  792. // Select retrieves rows from a table by scanning durable rows and collecting
  793. // only matching rows.
  794. func (m *TableManager) Select(table string, filter func(Row) bool) ([]Row, error) {
  795. tl := m.tableLock(table)
  796. tl.RLock()
  797. defer tl.RUnlock()
  798. return m.selectRows(table, filter)
  799. }
  800. // selectRows scans a table while the caller holds its shared or exclusive gate.
  801. func (m *TableManager) selectRows(table string, filter func(Row) bool) ([]Row, error) {
  802. if !m.schema.TableExists(table) {
  803. return nil, fmt.Errorf("table not found: %s", table)
  804. }
  805. rows := make([]Row, 0)
  806. err := m.scanRows(table, func(row Row) (bool, error) {
  807. if filter == nil || filter(row) {
  808. rows = append(rows, row)
  809. }
  810. return false, nil
  811. })
  812. if err != nil {
  813. return nil, err
  814. }
  815. return rows, nil
  816. }
  817. func cloneRow(row Row) Row {
  818. if row == nil {
  819. return nil
  820. }
  821. cloned := make(Row, len(row))
  822. for k, v := range row {
  823. cloned[k] = v
  824. }
  825. return cloned
  826. }
  827. // SelectWithLimit retrieves rows with limit and offset, applying filter/offset
  828. // while scanning and closing the scan early once the limit is reached.
  829. func (m *TableManager) SelectWithLimit(table string, filter func(Row) bool, limit, offset int) ([]Row, error) {
  830. tl := m.tableLock(table)
  831. tl.RLock()
  832. defer tl.RUnlock()
  833. if !m.schema.TableExists(table) {
  834. return nil, fmt.Errorf("table not found: %s", table)
  835. }
  836. rows := make([]Row, 0)
  837. skipped := 0
  838. pageSize := scanPageSize
  839. if limit > 0 {
  840. desired := limit
  841. if offset > 0 {
  842. if offset >= scanPageSize-desired {
  843. desired = scanPageSize
  844. } else {
  845. desired += offset
  846. }
  847. }
  848. if desired < 1 {
  849. desired = 1
  850. }
  851. if desired < pageSize {
  852. pageSize = desired
  853. }
  854. }
  855. err := m.scanRowsWithPageSize(table, uint32(pageSize), func(row Row) (bool, error) {
  856. if filter != nil && !filter(row) {
  857. return false, nil
  858. }
  859. if skipped < offset {
  860. skipped++
  861. return false, nil
  862. }
  863. rows = append(rows, row)
  864. return limit > 0 && len(rows) >= limit, nil
  865. })
  866. if err != nil {
  867. return nil, err
  868. }
  869. return rows, nil
  870. }
  871. // Update updates rows matching the filter.
  872. func (m *TableManager) Update(table string, updates Row, filter func(Row) bool) (int, error) {
  873. tl := m.tableLock(strings.ToLower(table))
  874. tl.Lock()
  875. defer tl.Unlock()
  876. rows, err := m.selectRows(table, filter)
  877. if err != nil {
  878. return 0, err
  879. }
  880. schema, err := m.schema.GetSchema(table)
  881. if err != nil {
  882. return 0, err
  883. }
  884. count := 0
  885. for _, row := range rows {
  886. // Snapshot the pre-update row so removed index entries can be restored
  887. // if persistence fails.
  888. oldRow := cloneRow(row)
  889. m.updateIndexesForRow(table, row, false)
  890. // Apply updates
  891. for k, v := range updates {
  892. // Normalize column name
  893. for _, col := range schema.Columns {
  894. if strings.EqualFold(k, col.Name) {
  895. row[col.Name] = v
  896. break
  897. }
  898. }
  899. }
  900. // Get primary key
  901. pkValue := row[schema.PrimaryKey]
  902. pk := fmt.Sprintf("%v", pkValue)
  903. // Serialize row
  904. data, err := encodeRow(row)
  905. if err != nil {
  906. m.updateIndexesForRow(table, oldRow, true)
  907. continue
  908. }
  909. if err := m.validateUniqueRows(table, []Row{row}, excludedKey(m.dataKey(table, fmt.Sprintf("%v", oldRow[schema.PrimaryKey])))); err != nil {
  910. m.updateIndexesForRow(table, oldRow, true)
  911. return count, err
  912. }
  913. // Write back
  914. key := m.dataKey(table, pk)
  915. err = m.pool.WithClient(func(c *KVClient) error {
  916. _, err := c.Put([]byte(key), data)
  917. return err
  918. })
  919. if err == nil {
  920. // Add new index entries after update
  921. m.updateIndexesForRow(table, row, true)
  922. m.bumpIndexPredicates(table, oldRow, row)
  923. count++
  924. } else {
  925. m.updateIndexesForRow(table, oldRow, true)
  926. }
  927. }
  928. if count > 0 {
  929. m.bumpGeneration(table)
  930. }
  931. return count, nil
  932. }
  933. // UpdateFunc updates rows matching the filter using a function to compute new values.
  934. // The updateFn receives the current row and returns the updates to apply.
  935. func (m *TableManager) UpdateFunc(table string, updateFn func(Row) (Row, error), filter func(Row) bool) (int, error) {
  936. tl := m.tableLock(strings.ToLower(table))
  937. tl.Lock()
  938. defer tl.Unlock()
  939. rows, err := m.selectRows(table, filter)
  940. if err != nil {
  941. return 0, err
  942. }
  943. schema, err := m.schema.GetSchema(table)
  944. if err != nil {
  945. return 0, err
  946. }
  947. count := 0
  948. for _, row := range rows {
  949. oldRow := cloneRow(row)
  950. m.updateIndexesForRow(table, row, false)
  951. // Compute updates using the provided function
  952. updates, err := updateFn(row)
  953. if err != nil {
  954. m.updateIndexesForRow(table, oldRow, true)
  955. return count, err
  956. }
  957. // Apply updates
  958. for k, v := range updates {
  959. // Normalize column name
  960. for _, col := range schema.Columns {
  961. if strings.EqualFold(k, col.Name) {
  962. row[col.Name] = v
  963. break
  964. }
  965. }
  966. }
  967. // Get primary key
  968. pkValue := row[schema.PrimaryKey]
  969. pk := fmt.Sprintf("%v", pkValue)
  970. // Serialize row
  971. data, err := encodeRow(row)
  972. if err != nil {
  973. m.updateIndexesForRow(table, oldRow, true)
  974. continue
  975. }
  976. if err := m.validateUniqueRows(table, []Row{row}, excludedKey(m.dataKey(table, fmt.Sprintf("%v", oldRow[schema.PrimaryKey])))); err != nil {
  977. m.updateIndexesForRow(table, oldRow, true)
  978. return count, err
  979. }
  980. // Write back
  981. key := m.dataKey(table, pk)
  982. err = m.pool.WithClient(func(c *KVClient) error {
  983. _, err := c.Put([]byte(key), data)
  984. return err
  985. })
  986. if err == nil {
  987. // Add new index entries after update
  988. m.updateIndexesForRow(table, row, true)
  989. m.bumpIndexPredicates(table, oldRow, row)
  990. count++
  991. } else {
  992. m.updateIndexesForRow(table, oldRow, true)
  993. }
  994. }
  995. if count > 0 {
  996. m.bumpGeneration(table)
  997. }
  998. return count, nil
  999. }
  1000. // UpdateByPK updates one row without scanning the table. It uses the key's
  1001. // striped lock and a compare-and-swap write so a concurrent modification of the
  1002. // same row fails with a serialization error instead of being silently lost.
  1003. func (m *TableManager) UpdateByPK(table, pk string, updateFn func(Row) (Row, error)) (Row, bool, error) {
  1004. schema, err := m.schema.GetSchema(table)
  1005. if err != nil {
  1006. return nil, false, err
  1007. }
  1008. key := m.dataKey(table, pk)
  1009. unlock, uniq, err := m.lockForWrite(table)
  1010. if err != nil {
  1011. return nil, false, err
  1012. }
  1013. defer unlock()
  1014. st := m.stripeKey(key)
  1015. st.Lock()
  1016. defer st.Unlock()
  1017. row, lsn, err := m.getByPKWithLSN(table, pk)
  1018. if err == ErrKeyNotFound {
  1019. return nil, false, nil
  1020. }
  1021. if err != nil {
  1022. return nil, false, err
  1023. }
  1024. oldRow := cloneRow(row)
  1025. updates, err := updateFn(row)
  1026. if err != nil {
  1027. return nil, false, err
  1028. }
  1029. for name, value := range updates {
  1030. for _, column := range schema.Columns {
  1031. if strings.EqualFold(name, column.Name) {
  1032. row[column.Name] = value
  1033. break
  1034. }
  1035. }
  1036. }
  1037. data, err := encodeRow(row)
  1038. if err != nil {
  1039. return nil, false, err
  1040. }
  1041. if uniq {
  1042. if err := m.validateUniqueRows(table, []Row{row}, excludedKey(key)); err != nil {
  1043. return nil, false, err
  1044. }
  1045. }
  1046. committed, err := m.compareWritePoint(
  1047. []CompareCheck{{Key: []byte(key), LSN: lsn}},
  1048. []BatchOp{{Op: batchPut, Key: []byte(key), Value: data}},
  1049. )
  1050. if err != nil {
  1051. return nil, false, err
  1052. }
  1053. if !committed {
  1054. return nil, false, ErrSerialization
  1055. }
  1056. m.updateIndexesForRow(table, oldRow, false)
  1057. m.updateIndexesForRow(table, row, true)
  1058. m.bumpIndexPredicates(table, oldRow, row)
  1059. m.bumpGeneration(table)
  1060. return oldRow, true, nil
  1061. }
  1062. // Delete deletes rows matching the filter.
  1063. func (m *TableManager) Delete(table string, filter func(Row) bool) (int, error) {
  1064. tl := m.tableLock(strings.ToLower(table))
  1065. tl.Lock()
  1066. defer tl.Unlock()
  1067. rows, err := m.selectRows(table, filter)
  1068. if err != nil {
  1069. return 0, err
  1070. }
  1071. schema, err := m.schema.GetSchema(table)
  1072. if err != nil {
  1073. return 0, err
  1074. }
  1075. wasInit := m.countInitialized(table)
  1076. count := 0
  1077. for _, row := range rows {
  1078. // Remove index entries before deleting row
  1079. m.updateIndexesForRow(table, row, false)
  1080. pkValue := row[schema.PrimaryKey]
  1081. pk := fmt.Sprintf("%v", pkValue)
  1082. key := m.dataKey(table, pk)
  1083. err = m.pool.WithClient(func(c *KVClient) error {
  1084. _, err := c.Del([]byte(key))
  1085. return err
  1086. })
  1087. if err == nil {
  1088. m.bumpIndexPredicates(table, row)
  1089. count++
  1090. } else {
  1091. // Restore the index entries removed above.
  1092. m.updateIndexesForRow(table, row, true)
  1093. }
  1094. }
  1095. if count > 0 {
  1096. m.bumpGeneration(table)
  1097. }
  1098. m.incrCount(table, schema.CreatedAt, -count, wasInit)
  1099. return count, nil
  1100. }
  1101. // DeleteByPK deletes one row without scanning the table, using the key's
  1102. // striped lock and a compare-and-swap delete.
  1103. func (m *TableManager) DeleteByPK(table, pk string) (Row, bool, error) {
  1104. schema, err := m.schema.GetSchema(table)
  1105. if err != nil {
  1106. return nil, false, err
  1107. }
  1108. key := m.dataKey(table, pk)
  1109. tl := m.tableLock(table)
  1110. tl.RLock()
  1111. defer tl.RUnlock()
  1112. st := m.stripeKey(key)
  1113. st.Lock()
  1114. defer st.Unlock()
  1115. wasInit := m.countInitialized(table)
  1116. row, lsn, err := m.getByPKWithLSN(table, pk)
  1117. if err == ErrKeyNotFound {
  1118. return nil, false, nil
  1119. }
  1120. if err != nil {
  1121. return nil, false, err
  1122. }
  1123. committed, err := m.compareWritePoint(
  1124. []CompareCheck{{Key: []byte(key), LSN: lsn}},
  1125. []BatchOp{{Op: batchDelete, Key: []byte(key)}},
  1126. )
  1127. if err != nil {
  1128. return nil, false, err
  1129. }
  1130. if !committed {
  1131. return nil, false, ErrSerialization
  1132. }
  1133. m.updateIndexesForRow(table, row, false)
  1134. m.bumpIndexPredicates(table, row)
  1135. m.bumpGeneration(table)
  1136. m.incrCount(table, schema.CreatedAt, -1, wasInit)
  1137. return row, true, nil
  1138. }
  1139. // GetByPK retrieves a row by primary key. Point reads take only the key's
  1140. // striped lock so reads of different keys progress concurrently.
  1141. func (m *TableManager) GetByPK(table string, pk string) (Row, error) {
  1142. key := m.dataKey(table, pk)
  1143. st := m.stripeKey(key)
  1144. st.Lock()
  1145. defer st.Unlock()
  1146. if !m.schema.TableExists(table) {
  1147. return nil, fmt.Errorf("table not found: %s", table)
  1148. }
  1149. row, _, err := m.getByPKWithLSN(table, pk)
  1150. return row, err
  1151. }
  1152. // getByPKWithLSN reads a row by primary key and returns its KV LSN (0 when
  1153. // absent). It performs no locking; callers must hold the appropriate striped
  1154. // or table lock.
  1155. func (m *TableManager) getByPKWithLSN(table, pk string) (Row, uint64, error) {
  1156. key := m.dataKey(table, pk)
  1157. var value []byte
  1158. var lsn uint64
  1159. err := m.pool.WithClient(func(c *KVClient) error {
  1160. res, err := c.Get([]byte(key))
  1161. if err != nil {
  1162. return err
  1163. }
  1164. value = res.Value
  1165. lsn = res.LSN
  1166. return nil
  1167. })
  1168. if err != nil {
  1169. if err == ErrKeyNotFound {
  1170. return nil, 0, ErrKeyNotFound
  1171. }
  1172. return nil, 0, err
  1173. }
  1174. row, err := decodeRow(value)
  1175. if err != nil {
  1176. return nil, 0, fmt.Errorf("failed to parse row: %w", err)
  1177. }
  1178. return row, lsn, nil
  1179. }
  1180. // Count returns the number of rows in a table matching the filter.
  1181. func (m *TableManager) Count(table string, filter func(Row) bool) (int, error) {
  1182. tl := m.tableLock(table)
  1183. tl.RLock()
  1184. defer tl.RUnlock()
  1185. if !m.schema.TableExists(table) {
  1186. return 0, fmt.Errorf("table not found: %s", table)
  1187. }
  1188. count := 0
  1189. err := m.scanRows(table, func(row Row) (bool, error) {
  1190. if filter == nil || filter(row) {
  1191. count++
  1192. }
  1193. return false, nil
  1194. })
  1195. if err != nil {
  1196. return 0, err
  1197. }
  1198. return count, nil
  1199. }
  1200. // Truncate removes all rows from a table.
  1201. func (m *TableManager) Truncate(table string) (int, error) {
  1202. return m.Delete(table, nil)
  1203. }
  1204. // isIntegerType checks if a type name is an integer type.
  1205. func isIntegerType(typeName string) bool {
  1206. t := strings.ToUpper(typeName)
  1207. switch t {
  1208. case "INTEGER", "INT", "SMALLINT", "BIGINT", "TINYINT", "MEDIUMINT":
  1209. return true
  1210. }
  1211. return false
  1212. }
  1213. // IsRowIDColumn checks if a column name is a ROWID alias.
  1214. func IsRowIDColumn(name string) bool {
  1215. n := strings.ToLower(name)
  1216. return n == "rowid" || n == "oid" || n == "_rowid_"
  1217. }
  1218. // Index entry methods - leveraging radix trie for prefix-based lookups
  1219. // Format: {database}:idx:{index_name}:{column_value} → JSON array of rowids
  1220. // indexEntryKey returns the key for an index entry.
  1221. func (m *TableManager) indexEntryKey(indexName string, colValue interface{}) string {
  1222. return fmt.Sprintf("%s:idx:%s:%s", m.database, strings.ToLower(indexName), formatIndexValue(colValue))
  1223. }
  1224. // indexPrefix returns the prefix for all entries of an index.
  1225. func (m *TableManager) indexPrefix(indexName string) string {
  1226. return fmt.Sprintf("%s:idx:%s:", m.database, strings.ToLower(indexName))
  1227. }
  1228. func formatIndexValue(value interface{}) string {
  1229. switch v := value.(type) {
  1230. case float64:
  1231. if v == float64(int64(v)) {
  1232. return fmt.Sprintf("%d", int64(v))
  1233. }
  1234. return fmt.Sprintf("%f", v)
  1235. case int64:
  1236. return fmt.Sprintf("%d", v)
  1237. case int:
  1238. return fmt.Sprintf("%d", v)
  1239. default:
  1240. return fmt.Sprintf("%v", v)
  1241. }
  1242. }
  1243. func rowIDFromRow(row Row) (int64, bool) {
  1244. switch v := row["_rowid_"].(type) {
  1245. case int64:
  1246. return v, true
  1247. case int:
  1248. return int64(v), true
  1249. case float64:
  1250. return int64(v), true
  1251. default:
  1252. return 0, false
  1253. }
  1254. }
  1255. func (m *TableManager) ensureIndex(index *Index) error {
  1256. indexKey := strings.ToLower(index.Name)
  1257. m.cacheMu.RLock()
  1258. disabled := m.disabledIndexes[indexKey]
  1259. _, initialized := m.indexCache[indexKey]
  1260. m.cacheMu.RUnlock()
  1261. if disabled {
  1262. return nil
  1263. }
  1264. if initialized {
  1265. return nil
  1266. }
  1267. // Serialize index build against writes to the same table so the derived
  1268. // entries cannot miss a concurrently-inserted row.
  1269. table := index.Table
  1270. key := strings.ToLower(table)
  1271. tl := m.tableLock(key)
  1272. tl.Lock()
  1273. defer tl.Unlock()
  1274. m.cacheMu.RLock()
  1275. disabled = m.disabledIndexes[indexKey]
  1276. _, initialized = m.indexCache[indexKey]
  1277. m.cacheMu.RUnlock()
  1278. if disabled {
  1279. return nil
  1280. }
  1281. if initialized {
  1282. return nil
  1283. }
  1284. columns := make([]string, len(index.Columns))
  1285. for i, col := range index.Columns {
  1286. columns[i] = col.Name
  1287. }
  1288. tableSchema, err := m.schema.GetSchema(table)
  1289. if err != nil {
  1290. return err
  1291. }
  1292. values := make(map[string][]int64)
  1293. rowKeys := make(map[int64]string)
  1294. if err := m.scanRows(table, func(row Row) (bool, error) {
  1295. rowid, ok := rowIDFromRow(row)
  1296. if !ok {
  1297. return false, nil
  1298. }
  1299. colValue := m.buildIndexValue(row, columns)
  1300. valueKey := formatIndexValue(colValue)
  1301. values[valueKey] = append(values[valueKey], rowid)
  1302. rowKeys[rowid] = fmt.Sprintf("%v", row[tableSchema.PrimaryKey])
  1303. return false, nil
  1304. }); err != nil {
  1305. return err
  1306. }
  1307. m.cacheMu.Lock()
  1308. if _, initialized := m.indexCache[indexKey]; !initialized {
  1309. m.indexCache[indexKey] = values
  1310. m.indexTable[indexKey] = key
  1311. m.rowKeyCache[key] = rowKeys
  1312. }
  1313. m.cacheMu.Unlock()
  1314. return nil
  1315. }
  1316. // AddIndexEntry adds a rowid to an in-memory index entry.
  1317. func (m *TableManager) AddIndexEntry(indexName string, colValue interface{}, rowid int64) error {
  1318. indexKey := strings.ToLower(indexName)
  1319. valueKey := formatIndexValue(colValue)
  1320. m.cacheMu.Lock()
  1321. defer m.cacheMu.Unlock()
  1322. values, ok := m.indexCache[indexKey]
  1323. if !ok {
  1324. return nil
  1325. }
  1326. rowids := values[valueKey]
  1327. for _, r := range rowids {
  1328. if r == rowid {
  1329. return nil
  1330. }
  1331. }
  1332. values[valueKey] = append(rowids, rowid)
  1333. return nil
  1334. }
  1335. // RemoveIndexEntry removes a rowid from an in-memory index entry.
  1336. func (m *TableManager) RemoveIndexEntry(indexName string, colValue interface{}, rowid int64) error {
  1337. indexKey := strings.ToLower(indexName)
  1338. valueKey := formatIndexValue(colValue)
  1339. m.cacheMu.Lock()
  1340. defer m.cacheMu.Unlock()
  1341. values, ok := m.indexCache[indexKey]
  1342. if !ok {
  1343. return nil
  1344. }
  1345. rowids := values[valueKey]
  1346. newRowids := make([]int64, 0, len(rowids))
  1347. for _, r := range rowids {
  1348. if r != rowid {
  1349. newRowids = append(newRowids, r)
  1350. }
  1351. }
  1352. if len(newRowids) == 0 {
  1353. delete(values, valueKey)
  1354. return nil
  1355. }
  1356. values[valueKey] = newRowids
  1357. return nil
  1358. }
  1359. // LookupIndex returns rowids matching a column value using the index.
  1360. func (m *TableManager) LookupIndex(indexName string, colValue interface{}) ([]int64, error) {
  1361. index, err := m.schema.GetIndex(indexName)
  1362. if err != nil {
  1363. return nil, err
  1364. }
  1365. if err := m.ensureIndex(index); err != nil {
  1366. return nil, err
  1367. }
  1368. indexKey := strings.ToLower(indexName)
  1369. valueKey := formatIndexValue(colValue)
  1370. m.cacheMu.RLock()
  1371. rowids := append([]int64(nil), m.indexCache[indexKey][valueKey]...)
  1372. m.cacheMu.RUnlock()
  1373. return rowids, nil
  1374. }
  1375. // ClearIndex removes an index's in-memory entries and marks it disabled so a
  1376. // concurrent lookup cannot rebuild it after DROP but before the schema entry
  1377. // is removed. Index entries are derived from durable rows, so there is nothing
  1378. // durable to delete here.
  1379. func (m *TableManager) ClearIndex(indexName, tableName string, columns []string) error {
  1380. indexKey := strings.ToLower(indexName)
  1381. tableKey := strings.ToLower(tableName)
  1382. m.cacheMu.Lock()
  1383. delete(m.indexCache, indexKey)
  1384. delete(m.indexTable, indexKey)
  1385. m.disabledIndexes[indexKey] = true
  1386. rowKeysNeeded := false
  1387. for _, indexedTable := range m.indexTable {
  1388. if indexedTable == tableKey {
  1389. rowKeysNeeded = true
  1390. break
  1391. }
  1392. }
  1393. if !rowKeysNeeded {
  1394. delete(m.rowKeyCache, tableKey)
  1395. }
  1396. m.cacheMu.Unlock()
  1397. return nil
  1398. }
  1399. // BuildIndex builds index entries for all existing rows in a table.
  1400. func (m *TableManager) BuildIndex(indexName, tableName string, columns []string) error {
  1401. indexKey := strings.ToLower(indexName)
  1402. m.cacheMu.Lock()
  1403. delete(m.disabledIndexes, indexKey)
  1404. delete(m.indexCache, indexKey)
  1405. delete(m.indexTable, indexKey)
  1406. m.cacheMu.Unlock()
  1407. index, err := m.schema.GetIndex(indexName)
  1408. if err == nil {
  1409. return m.ensureIndex(index)
  1410. }
  1411. if !errors.Is(err, ErrIndexNotFound) {
  1412. return err
  1413. }
  1414. tableSchema, schemaErr := m.schema.GetSchema(tableName)
  1415. if schemaErr != nil {
  1416. return schemaErr
  1417. }
  1418. values := make(map[string][]int64)
  1419. rowKeys := make(map[int64]string)
  1420. if err := m.scanRows(tableName, func(row Row) (bool, error) {
  1421. rowid, ok := rowIDFromRow(row)
  1422. if !ok {
  1423. return false, nil
  1424. }
  1425. colValue := m.buildIndexValue(row, columns)
  1426. values[formatIndexValue(colValue)] = append(values[formatIndexValue(colValue)], rowid)
  1427. rowKeys[rowid] = fmt.Sprintf("%v", row[tableSchema.PrimaryKey])
  1428. return false, nil
  1429. }); err != nil {
  1430. return err
  1431. }
  1432. m.cacheMu.Lock()
  1433. m.indexCache[indexKey] = values
  1434. m.indexTable[indexKey] = strings.ToLower(tableName)
  1435. m.rowKeyCache[strings.ToLower(tableName)] = rowKeys
  1436. m.cacheMu.Unlock()
  1437. return nil
  1438. }
  1439. // buildIndexValue creates the index key value from row columns.
  1440. func (m *TableManager) buildIndexValue(row Row, columns []string) string {
  1441. formatValue := func(v interface{}) string {
  1442. switch val := v.(type) {
  1443. case float64:
  1444. // Check if it's actually an integer value
  1445. if val == float64(int64(val)) {
  1446. return fmt.Sprintf("%d", int64(val))
  1447. }
  1448. return fmt.Sprintf("%f", val)
  1449. case int64:
  1450. return fmt.Sprintf("%d", val)
  1451. case int:
  1452. return fmt.Sprintf("%d", val)
  1453. default:
  1454. return fmt.Sprintf("%v", val)
  1455. }
  1456. }
  1457. if len(columns) == 1 {
  1458. return formatValue(row[columns[0]])
  1459. }
  1460. // Multi-column index: concatenate values with separator
  1461. var parts []string
  1462. for _, col := range columns {
  1463. parts = append(parts, formatValue(row[col]))
  1464. }
  1465. return strings.Join(parts, "\x00")
  1466. }
  1467. // ── UNIQUE index enforcement ────────────────────────────────────────────────
  1468. //
  1469. // Uniqueness is enforced by scanning durable rows rather than by persisting a
  1470. // separate claim key. There is therefore no storage migration and no new durable
  1471. // state: a table with a UNIQUE index takes its exclusive table gate for writes so
  1472. // a validating scan can never race a concurrent writer. NULL values are exempt
  1473. // (SQLite permits multiple NULLs in a unique index), and composite values are
  1474. // encoded with length prefixes so "a\x00b" in one column can never collide with
  1475. // "a", "b" across two columns.
  1476. // uniqueIndexes returns the UNIQUE indexes defined on a table.
  1477. func (m *TableManager) uniqueIndexes(table string) ([]*Index, error) {
  1478. indexes, err := m.schema.ListTableIndexes(table)
  1479. if err != nil {
  1480. return nil, err
  1481. }
  1482. var uniq []*Index
  1483. for _, idx := range indexes {
  1484. if idx.Unique {
  1485. uniq = append(uniq, idx)
  1486. }
  1487. }
  1488. return uniq, nil
  1489. }
  1490. // hasUniqueIndex reports whether the table has any UNIQUE index. An IO error is
  1491. // surfaced so callers never silently skip uniqueness enforcement.
  1492. func (m *TableManager) hasUniqueIndex(table string) (bool, error) {
  1493. indexes, err := m.uniqueIndexes(table)
  1494. if err != nil {
  1495. return false, err
  1496. }
  1497. return len(indexes) > 0, nil
  1498. }
  1499. // HasUniqueIndex reports whether the table has any UNIQUE index, surfacing IO
  1500. // errors. It is the exported form used by the executor to decide whether to run a
  1501. // DML statement inside an implicit transaction for statement-level atomicity.
  1502. func (m *TableManager) HasUniqueIndex(table string) (bool, error) {
  1503. return m.hasUniqueIndex(table)
  1504. }
  1505. // encodeUniqueValue encodes the indexed value of a row with per-column type tags
  1506. // and length prefixes. It returns isNull=true when any indexed column is NULL, in
  1507. // which case the row is exempt from uniqueness. The length prefix makes composite
  1508. // encodings collision-free.
  1509. func encodeUniqueValue(row Row, columns []string) (string, bool) {
  1510. var sb strings.Builder
  1511. for _, col := range columns {
  1512. v, ok := row[col]
  1513. if !ok {
  1514. for k, val := range row {
  1515. if strings.EqualFold(k, col) {
  1516. v = val
  1517. ok = true
  1518. break
  1519. }
  1520. }
  1521. }
  1522. if !ok || v == nil {
  1523. return "", true
  1524. }
  1525. s := encodeUniqueScalar(v)
  1526. sb.WriteString(strconv.Itoa(len(s)))
  1527. sb.WriteByte(':')
  1528. sb.WriteString(s)
  1529. }
  1530. return sb.String(), false
  1531. }
  1532. // encodeUniqueScalar encodes a scalar for uniqueness comparison. Integral
  1533. // numerics are canonicalized across their integer/float/unsigned Go
  1534. // representations so a computed value such as 1.5-0.5 (float64) collides with
  1535. // the literal 1 (int64) the same way SQLite's numeric affinity does. Non-integral
  1536. // reals keep a full-precision tag so no distinct value is lost. TEXT and BLOB
  1537. // remain distinct from numbers.
  1538. func encodeUniqueScalar(v interface{}) string {
  1539. switch t := v.(type) {
  1540. case bool:
  1541. if t {
  1542. return "b1"
  1543. }
  1544. return "b0"
  1545. case int:
  1546. return "i" + strconv.FormatInt(int64(t), 10)
  1547. case int64:
  1548. return "i" + strconv.FormatInt(t, 10)
  1549. case uint:
  1550. return "i" + strconv.FormatUint(uint64(t), 10)
  1551. case uint8:
  1552. return "i" + strconv.FormatUint(uint64(t), 10)
  1553. case uint16:
  1554. return "i" + strconv.FormatUint(uint64(t), 10)
  1555. case uint32:
  1556. return "i" + strconv.FormatUint(uint64(t), 10)
  1557. case uint64:
  1558. return "i" + strconv.FormatUint(t, 10)
  1559. case uintptr:
  1560. return "i" + strconv.FormatUint(uint64(t), 10)
  1561. case float64:
  1562. if s, ok := encodeIntegralFloat(t); ok {
  1563. return s
  1564. }
  1565. return "r" + strconv.FormatFloat(t, 'g', -1, 64)
  1566. case string:
  1567. return "s" + t
  1568. case []byte:
  1569. return "x" + string(t)
  1570. default:
  1571. return "s" + fmt.Sprintf("%v", t)
  1572. }
  1573. }
  1574. // encodeIntegralFloat canonicalizes a whole float64 to the same "i" encoding used
  1575. // by integer values, without precision loss, so integral reals and integers are
  1576. // not treated as distinct under a unique index. It reports false (and returns
  1577. // nothing) for non-integral values or values outside the int64 range.
  1578. func encodeIntegralFloat(v float64) (string, bool) {
  1579. if v != math.Trunc(v) {
  1580. return "", false
  1581. }
  1582. if v < -9223372036854775808.0 || v >= 9223372036854775808.0 {
  1583. return "", false
  1584. }
  1585. return "i" + strconv.FormatInt(int64(v), 10), true
  1586. }
  1587. // validateUniqueRows checks that pending rows do not conflict with one another or
  1588. // with durable rows on any UNIQUE index. excludedKeys lists the actual durable
  1589. // data keys whose rows are being replaced in this operation (delete, or an update
  1590. // of a non-indexed column), so those rows' own unique values do not self-conflict
  1591. // and swapping two unique values remains possible. Matching is done against the
  1592. // raw scanned KV key rather than a re-stringified primary key, so numeric and
  1593. // composite keys are unambiguous. The caller must hold the table's exclusive gate.
  1594. func (m *TableManager) validateUniqueRows(table string, pending []Row, excludedKeys map[string]bool) error {
  1595. indexes, err := m.uniqueIndexes(table)
  1596. if err != nil {
  1597. return err
  1598. }
  1599. if len(indexes) == 0 {
  1600. return nil
  1601. }
  1602. // seen tracks, per unique index, the encoded values already claimed by the
  1603. // pending rows. Keying by index name (not just the encoded value) is essential:
  1604. // a single row can legitimately carry the same value in two different unique
  1605. // indexes (e.g. name and lower_name both "alice"), which must not collide with
  1606. // itself. Two rows only conflict when they share a value within the SAME index.
  1607. seen := make(map[string]map[string]bool, len(indexes))
  1608. for _, row := range pending {
  1609. for _, idx := range indexes {
  1610. columns := indexColumnNames(idx)
  1611. encoded, isNull := encodeUniqueValue(row, columns)
  1612. if isNull {
  1613. continue
  1614. }
  1615. perIndex := seen[idx.Name]
  1616. if perIndex == nil {
  1617. perIndex = make(map[string]bool)
  1618. seen[idx.Name] = perIndex
  1619. }
  1620. if perIndex[encoded] {
  1621. return fmt.Errorf("UNIQUE constraint failed: %s", idx.Name)
  1622. }
  1623. perIndex[encoded] = true
  1624. }
  1625. }
  1626. return m.pool.WithClient(func(client *KVClient) (retErr error) {
  1627. cursor, err := client.ScanWithLimit([]byte(m.dataPrefix(table)), scanPageSize)
  1628. if err != nil {
  1629. return err
  1630. }
  1631. defer func() {
  1632. if err := cursor.Close(); retErr == nil {
  1633. retErr = err
  1634. }
  1635. }()
  1636. for {
  1637. entries, done, err := cursor.Next()
  1638. if err != nil {
  1639. return err
  1640. }
  1641. for _, e := range entries {
  1642. if excludedKeys != nil && excludedKeys[string(e.Key)] {
  1643. continue
  1644. }
  1645. row, err := decodeRow(e.Value)
  1646. if err != nil {
  1647. return err
  1648. }
  1649. for _, idx := range indexes {
  1650. columns := indexColumnNames(idx)
  1651. encoded, isNull := encodeUniqueValue(row, columns)
  1652. if isNull {
  1653. continue
  1654. }
  1655. if perIndex := seen[idx.Name]; perIndex != nil && perIndex[encoded] {
  1656. return fmt.Errorf("UNIQUE constraint failed: %s", idx.Name)
  1657. }
  1658. }
  1659. }
  1660. if done {
  1661. return nil
  1662. }
  1663. }
  1664. })
  1665. }
  1666. // excludedKey returns the exclusion set that exempts a single row's durable value
  1667. // from the uniqueness scan, keyed by its actual durable data key.
  1668. func excludedKey(dataKey string) map[string]bool {
  1669. return map[string]bool{dataKey: true}
  1670. }
  1671. // lockForWrite acquires the table gate appropriate for a point write: exclusive
  1672. // when the table has a UNIQUE index (so the validating scan cannot race another
  1673. // writer), shared otherwise. The uniqueness decision is re-checked under the lock
  1674. // so a concurrent CREATE UNIQUE INDEX cannot slip in between the unlocked probe
  1675. // and the lock acquisition; the returned unlock closes the gate.
  1676. func (m *TableManager) lockForWrite(table string) (unlock func(), unique bool, err error) {
  1677. tl := m.tableLock(table)
  1678. for {
  1679. uniq, err := m.hasUniqueIndex(table)
  1680. if err != nil {
  1681. return nil, false, err
  1682. }
  1683. if uniq {
  1684. tl.Lock()
  1685. // Re-check under the exclusive gate (authoritative): CreateUniqueIndex
  1686. // also takes the exclusive gate, so it cannot run while we hold it.
  1687. uniq, err = m.hasUniqueIndex(table)
  1688. if err != nil {
  1689. tl.Unlock()
  1690. return nil, false, err
  1691. }
  1692. if uniq {
  1693. return tl.Unlock, true, nil
  1694. }
  1695. tl.Unlock()
  1696. continue
  1697. }
  1698. tl.RLock()
  1699. // A unique index could have appeared before we acquired the shared gate;
  1700. // re-check and escalate if so.
  1701. uniq, err = m.hasUniqueIndex(table)
  1702. if err != nil {
  1703. tl.RUnlock()
  1704. return nil, false, err
  1705. }
  1706. if uniq {
  1707. tl.RUnlock()
  1708. continue
  1709. }
  1710. return tl.RUnlock, false, nil
  1711. }
  1712. }
  1713. func indexColumnNames(idx *Index) []string {
  1714. names := make([]string, len(idx.Columns))
  1715. for i, c := range idx.Columns {
  1716. names[i] = c.Name
  1717. }
  1718. return names
  1719. }
  1720. // ValidateUniqueIndex rejects a UNIQUE index definition if existing rows already
  1721. // contain duplicate non-NULL values for the indexed columns.
  1722. func (m *TableManager) ValidateUniqueIndex(index *Index) error {
  1723. columns := indexColumnNames(index)
  1724. seen := make(map[string]bool)
  1725. return m.scanRows(index.Table, func(row Row) (bool, error) {
  1726. encoded, isNull := encodeUniqueValue(row, columns)
  1727. if isNull {
  1728. return false, nil
  1729. }
  1730. if seen[encoded] {
  1731. return true, fmt.Errorf("UNIQUE constraint failed: %s", index.Name)
  1732. }
  1733. seen[encoded] = true
  1734. return false, nil
  1735. })
  1736. }
  1737. // CreateUniqueIndex validates existing rows and registers the index while
  1738. // holding the table's exclusive gate, so a concurrent writer cannot insert a
  1739. // conflicting value between validation and registration.
  1740. func (m *TableManager) CreateUniqueIndex(index *Index) error {
  1741. tl := m.tableLock(index.Table)
  1742. tl.Lock()
  1743. defer tl.Unlock()
  1744. if err := m.ValidateUniqueIndex(index); err != nil {
  1745. return err
  1746. }
  1747. return m.schema.CreateIndex(index)
  1748. }
  1749. type indexedRowVersion struct {
  1750. row Row
  1751. key string
  1752. lsn uint64
  1753. }
  1754. type indexPredicateSnapshot struct {
  1755. valueKey string
  1756. valueGen uint64
  1757. wildcardKey string
  1758. wildcardGen uint64
  1759. }
  1760. // selectByIndexWithLSN retrieves indexed rows and their durable versions. The
  1761. // transaction layer uses the versions for optimistic commit validation.
  1762. func (m *TableManager) selectByIndexWithLSN(table, indexName string, colValue interface{}) ([]indexedRowVersion, indexPredicateSnapshot, error) {
  1763. index, err := m.schema.GetIndex(indexName)
  1764. if err != nil {
  1765. return nil, indexPredicateSnapshot{}, err
  1766. }
  1767. if err := m.ensureIndex(index); err != nil {
  1768. return nil, indexPredicateSnapshot{}, err
  1769. }
  1770. tableKey := strings.ToLower(table)
  1771. tl := m.tableLock(tableKey)
  1772. tl.RLock()
  1773. defer tl.RUnlock()
  1774. if !m.schema.TableExists(table) {
  1775. return nil, indexPredicateSnapshot{}, fmt.Errorf("table not found: %s", table)
  1776. }
  1777. indexKey := strings.ToLower(indexName)
  1778. valueKey := formatIndexValue(colValue)
  1779. predicateKey, predicateGen, wildcardKey, wildcardGen := m.predicateSnapshot(table, indexName, valueKey)
  1780. snapshot := indexPredicateSnapshot{
  1781. valueKey: predicateKey, valueGen: predicateGen,
  1782. wildcardKey: wildcardKey, wildcardGen: wildcardGen,
  1783. }
  1784. m.cacheMu.RLock()
  1785. rowids := append([]int64(nil), m.indexCache[indexKey][valueKey]...)
  1786. primaryKeys := make([]string, 0, len(rowids))
  1787. missingRowID := int64(0)
  1788. missingRowKey := false
  1789. for _, rowid := range rowids {
  1790. if primaryKey, ok := m.rowKeyCache[tableKey][rowid]; ok {
  1791. primaryKeys = append(primaryKeys, primaryKey)
  1792. } else {
  1793. missingRowID = rowid
  1794. missingRowKey = true
  1795. break
  1796. }
  1797. }
  1798. m.cacheMu.RUnlock()
  1799. if missingRowKey {
  1800. return nil, indexPredicateSnapshot{}, fmt.Errorf("index %s is missing rowid %d", indexName, missingRowID)
  1801. }
  1802. // If no rowids found, return empty result
  1803. if len(rowids) == 0 {
  1804. return []indexedRowVersion{}, snapshot, nil
  1805. }
  1806. rows := make([]indexedRowVersion, 0, len(primaryKeys))
  1807. err = m.pool.WithClient(func(client *KVClient) error {
  1808. keys := make([][]byte, len(primaryKeys))
  1809. for i, primaryKey := range primaryKeys {
  1810. keys[i] = []byte(m.dataKey(table, primaryKey))
  1811. }
  1812. results, err := client.MultiGet(keys)
  1813. if err != nil {
  1814. return err
  1815. }
  1816. for i, result := range results {
  1817. if !result.Found {
  1818. primaryKey := primaryKeys[i]
  1819. return fmt.Errorf("index %s references missing primary key %s", indexName, primaryKey)
  1820. }
  1821. row, err := decodeRow(result.Value)
  1822. if err != nil {
  1823. return err
  1824. }
  1825. rows = append(rows, indexedRowVersion{
  1826. row: row,
  1827. key: m.dataKey(table, primaryKeys[i]),
  1828. lsn: result.LSN,
  1829. })
  1830. }
  1831. return nil
  1832. })
  1833. if err != nil {
  1834. return nil, indexPredicateSnapshot{}, err
  1835. }
  1836. return rows, snapshot, nil
  1837. }
  1838. // SelectByIndex retrieves rows using an in-memory equality index followed by a
  1839. // single MultiGet for the matching primary keys.
  1840. func (m *TableManager) SelectByIndex(table, indexName string, colValue interface{}) ([]Row, error) {
  1841. versions, _, err := m.selectByIndexWithLSN(table, indexName, colValue)
  1842. if err != nil {
  1843. return nil, err
  1844. }
  1845. rows := make([]Row, len(versions))
  1846. for i := range versions {
  1847. rows[i] = versions[i].row
  1848. }
  1849. return rows, nil
  1850. }