table.go 58 KB

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