store.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610
  1. package metadata
  2. /*
  3. ACMS local store: in-memory maps backed by cluster.db tables.
  4. Tables (all registered as cluster tables, wiped on leave):
  5. meta_files ID -> FileRecord
  6. meta_paths Path -> ID
  7. meta_volumes ID -> Volume
  8. meta_policy Folder -> Policy
  9. meta_log %020d -> Entry
  10. meta_pending uuid -> Entry (changes not yet accepted by a leader)
  11. meta_state lease / applied / term
  12. */
  13. import (
  14. "encoding/json"
  15. "fmt"
  16. "sort"
  17. "strings"
  18. "sync"
  19. "imuslab.com/arozos/mod/database"
  20. )
  21. const (
  22. tableFiles = "meta_files"
  23. tablePaths = "meta_paths"
  24. tableVolumes = "meta_volumes"
  25. tablePolicy = "meta_policy"
  26. tableJobs = "meta_jobs"
  27. tableSettings = "meta_settings"
  28. tableLog = "meta_log"
  29. tablePending = "meta_pending"
  30. tableState = "meta_state"
  31. stateLease = "lease"
  32. stateApplied = "applied"
  33. stateTerm = "term"
  34. defaultPolicy = 1
  35. )
  36. // Tables lists every cluster.db table used by the metadata store.
  37. var Tables = []string{tableFiles, tablePaths, tableVolumes, tablePolicy, tableJobs, tableSettings, tableLog, tablePending, tableState}
  38. type store struct {
  39. db *database.Database
  40. persist *persister //disk writes happen here, never under mu
  41. mu sync.RWMutex
  42. files map[string]*FileRecord
  43. paths map[string]string
  44. volumes map[string]*Volume
  45. policies map[string]*Policy
  46. jobs map[string]*Job
  47. settings map[string]*Setting
  48. pending map[string]Entry
  49. lastSeq uint64
  50. applied uint64
  51. term uint64
  52. //The lease has its own lock: renewing it must never wait behind the
  53. //namespace, or a busy leader loses its lease.
  54. leaseMu sync.RWMutex
  55. lease Lease
  56. }
  57. func logKey(seq uint64) string { return fmt.Sprintf("%020d", seq) }
  58. func newStore(db *database.Database) *store {
  59. s := &store{
  60. db: db,
  61. files: map[string]*FileRecord{},
  62. paths: map[string]string{},
  63. volumes: map[string]*Volume{},
  64. policies: map[string]*Policy{},
  65. jobs: map[string]*Job{},
  66. settings: map[string]*Setting{},
  67. pending: map[string]Entry{},
  68. }
  69. for _, t := range Tables {
  70. db.NewTable(t)
  71. }
  72. s.load()
  73. s.persist = newPersister(db)
  74. return s
  75. }
  76. // close commits pending disk writes and stops the writer.
  77. func (s *store) close() {
  78. s.persist.close()
  79. }
  80. // discard drops pending disk writes; used when the tables were wiped.
  81. func (s *store) discard() {
  82. s.persist.discard()
  83. }
  84. func (s *store) load() {
  85. if entries, err := s.db.ListTable(tableFiles); err == nil {
  86. for _, kv := range entries {
  87. var rec FileRecord
  88. if json.Unmarshal(kv[1], &rec) == nil && rec.ID != "" {
  89. s.files[rec.ID] = &rec
  90. if !rec.Removed {
  91. s.paths[rec.Path] = rec.ID
  92. }
  93. }
  94. }
  95. }
  96. if entries, err := s.db.ListTable(tableVolumes); err == nil {
  97. for _, kv := range entries {
  98. var v Volume
  99. if json.Unmarshal(kv[1], &v) == nil && v.ID != "" {
  100. s.volumes[v.ID] = &v
  101. }
  102. }
  103. }
  104. if entries, err := s.db.ListTable(tablePolicy); err == nil {
  105. for _, kv := range entries {
  106. var p Policy
  107. if json.Unmarshal(kv[1], &p) == nil && p.Folder != "" {
  108. s.policies[p.Folder] = &p
  109. }
  110. }
  111. }
  112. if entries, err := s.db.ListTable(tableJobs); err == nil {
  113. for _, kv := range entries {
  114. var j Job
  115. if json.Unmarshal(kv[1], &j) == nil && j.ID != "" {
  116. s.jobs[j.ID] = &j
  117. }
  118. }
  119. }
  120. if entries, err := s.db.ListTable(tableSettings); err == nil {
  121. for _, kv := range entries {
  122. var st Setting
  123. if json.Unmarshal(kv[1], &st) == nil && st.Key != "" {
  124. s.settings[st.Key] = &st
  125. }
  126. }
  127. }
  128. if entries, err := s.db.ListTable(tablePending); err == nil {
  129. for _, kv := range entries {
  130. var e Entry
  131. if json.Unmarshal(kv[1], &e) == nil {
  132. s.pending[string(kv[0])] = e
  133. }
  134. }
  135. }
  136. if entries, err := s.db.ListTable(tableLog); err == nil && len(entries) > 0 {
  137. var last Entry
  138. if json.Unmarshal(entries[len(entries)-1][1], &last) == nil {
  139. s.lastSeq = last.Seq
  140. }
  141. }
  142. s.db.Read(tableState, stateApplied, &s.applied)
  143. s.db.Read(tableState, stateTerm, &s.term)
  144. s.db.Read(tableState, stateLease, &s.lease)
  145. }
  146. /*
  147. Files
  148. */
  149. // putFile merges a record last-writer-wins. Returns true when stored.
  150. func (s *store) putFile(in *FileRecord) bool {
  151. if in == nil || in.ID == "" || in.Path == "" {
  152. return false
  153. }
  154. s.mu.Lock()
  155. defer s.mu.Unlock()
  156. rec := in.Clone()
  157. rec.Path = NormalizePath(rec.Path)
  158. existing, ok := s.files[rec.ID]
  159. if ok && rec.Version <= existing.Version {
  160. return false
  161. }
  162. if ok && existing.Path != rec.Path && s.paths[existing.Path] == rec.ID {
  163. delete(s.paths, existing.Path)
  164. s.persist.del(tablePaths, existing.Path)
  165. }
  166. s.files[rec.ID] = rec
  167. s.persist.put(tableFiles, rec.ID, rec)
  168. if rec.Removed {
  169. if s.paths[rec.Path] == rec.ID {
  170. delete(s.paths, rec.Path)
  171. s.persist.del(tablePaths, rec.Path)
  172. }
  173. } else {
  174. //A newer record for the same path replaces an older one's claim
  175. if otherID, taken := s.paths[rec.Path]; taken && otherID != rec.ID {
  176. if other, ok := s.files[otherID]; ok && other.Version < rec.Version {
  177. s.paths[rec.Path] = rec.ID
  178. s.persist.put(tablePaths, rec.Path, rec.ID)
  179. }
  180. } else {
  181. s.paths[rec.Path] = rec.ID
  182. s.persist.put(tablePaths, rec.Path, rec.ID)
  183. }
  184. }
  185. return true
  186. }
  187. func (s *store) getFileByID(id string) (*FileRecord, bool) {
  188. s.mu.RLock()
  189. defer s.mu.RUnlock()
  190. rec, ok := s.files[id]
  191. if !ok {
  192. return nil, false
  193. }
  194. return rec.Clone(), true
  195. }
  196. func (s *store) getFileByPath(p string) (*FileRecord, bool) {
  197. p = NormalizePath(p)
  198. s.mu.RLock()
  199. defer s.mu.RUnlock()
  200. id, ok := s.paths[p]
  201. if !ok {
  202. return nil, false
  203. }
  204. rec, ok := s.files[id]
  205. if !ok || rec.Removed {
  206. return nil, false
  207. }
  208. return rec.Clone(), true
  209. }
  210. // listDir returns the live children of a directory path, sorted by name.
  211. func (s *store) listDir(dir string) []FileRecord {
  212. dir = NormalizePath(dir)
  213. prefix := dir + "/"
  214. if dir == "/" {
  215. prefix = "/"
  216. }
  217. s.mu.RLock()
  218. defer s.mu.RUnlock()
  219. out := []FileRecord{}
  220. for p, id := range s.paths {
  221. if !strings.HasPrefix(p, prefix) || p == dir {
  222. continue
  223. }
  224. rest := strings.TrimPrefix(p, prefix)
  225. if rest == "" || strings.Contains(rest, "/") {
  226. continue
  227. }
  228. if rec, ok := s.files[id]; ok && !rec.Removed {
  229. out = append(out, *rec.Clone())
  230. }
  231. }
  232. sort.Slice(out, func(i, j int) bool { return out[i].Path < out[j].Path })
  233. return out
  234. }
  235. // listSubtree returns every live record under a directory (recursive).
  236. func (s *store) listSubtree(dir string) []FileRecord {
  237. dir = NormalizePath(dir)
  238. prefix := dir + "/"
  239. if dir == "/" {
  240. prefix = "/"
  241. }
  242. s.mu.RLock()
  243. defer s.mu.RUnlock()
  244. out := []FileRecord{}
  245. for p, id := range s.paths {
  246. if !strings.HasPrefix(p, prefix) || p == dir {
  247. continue
  248. }
  249. if rec, ok := s.files[id]; ok && !rec.Removed {
  250. out = append(out, *rec.Clone())
  251. }
  252. }
  253. sort.Slice(out, func(i, j int) bool { return out[i].Path < out[j].Path })
  254. return out
  255. }
  256. func (s *store) allFiles() []FileRecord {
  257. s.mu.RLock()
  258. defer s.mu.RUnlock()
  259. out := make([]FileRecord, 0, len(s.files))
  260. for _, rec := range s.files {
  261. out = append(out, *rec.Clone())
  262. }
  263. sort.Slice(out, func(i, j int) bool { return out[i].Path < out[j].Path })
  264. return out
  265. }
  266. func (s *store) counts() (files int, dirs int) {
  267. s.mu.RLock()
  268. defer s.mu.RUnlock()
  269. for _, id := range s.paths {
  270. if rec, ok := s.files[id]; ok {
  271. if rec.IsDir {
  272. dirs++
  273. } else {
  274. files++
  275. }
  276. }
  277. }
  278. return
  279. }
  280. /*
  281. Volumes and policies
  282. */
  283. func (s *store) putVolume(in *Volume) bool {
  284. if in == nil || in.ID == "" {
  285. return false
  286. }
  287. s.mu.Lock()
  288. defer s.mu.Unlock()
  289. if existing, ok := s.volumes[in.ID]; ok && in.Version <= existing.Version {
  290. return false
  291. }
  292. v := *in
  293. s.volumes[v.ID] = &v
  294. s.persist.put(tableVolumes, v.ID, v)
  295. return true
  296. }
  297. func (s *store) getVolume(id string) (*Volume, bool) {
  298. s.mu.RLock()
  299. defer s.mu.RUnlock()
  300. v, ok := s.volumes[id]
  301. if !ok {
  302. return nil, false
  303. }
  304. c := *v
  305. return &c, true
  306. }
  307. func (s *store) allVolumes() []Volume {
  308. s.mu.RLock()
  309. defer s.mu.RUnlock()
  310. out := make([]Volume, 0, len(s.volumes))
  311. for _, v := range s.volumes {
  312. out = append(out, *v)
  313. }
  314. sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
  315. return out
  316. }
  317. func (s *store) putPolicy(in *Policy) bool {
  318. if in == nil || in.Folder == "" {
  319. return false
  320. }
  321. s.mu.Lock()
  322. defer s.mu.Unlock()
  323. p := *in
  324. p.Folder = TopFolder(p.Folder)
  325. if existing, ok := s.policies[p.Folder]; ok && p.Version <= existing.Version {
  326. return false
  327. }
  328. s.policies[p.Folder] = &p
  329. s.persist.put(tablePolicy, p.Folder, p)
  330. return true
  331. }
  332. func (s *store) allPolicies() []Policy {
  333. s.mu.RLock()
  334. defer s.mu.RUnlock()
  335. out := make([]Policy, 0, len(s.policies))
  336. for _, p := range s.policies {
  337. out = append(out, *p)
  338. }
  339. sort.Slice(out, func(i, j int) bool { return out[i].Folder < out[j].Folder })
  340. return out
  341. }
  342. func (s *store) policyFor(p string) Policy {
  343. folder := TopFolder(p)
  344. s.mu.RLock()
  345. defer s.mu.RUnlock()
  346. if pol, ok := s.policies[folder]; ok && pol.Replicas > 0 {
  347. return *pol
  348. }
  349. return Policy{Folder: folder, Replicas: defaultPolicy}
  350. }
  351. /*
  352. Jobs
  353. */
  354. func (s *store) putJob(in *Job) bool {
  355. if in == nil || in.ID == "" {
  356. return false
  357. }
  358. s.mu.Lock()
  359. defer s.mu.Unlock()
  360. if existing, ok := s.jobs[in.ID]; ok && in.Version <= existing.Version {
  361. return false
  362. }
  363. j := *in
  364. s.jobs[j.ID] = &j
  365. s.persist.put(tableJobs, j.ID, j)
  366. return true
  367. }
  368. func (s *store) getJob(id string) (*Job, bool) {
  369. s.mu.RLock()
  370. defer s.mu.RUnlock()
  371. j, ok := s.jobs[id]
  372. if !ok {
  373. return nil, false
  374. }
  375. c := *j
  376. return &c, true
  377. }
  378. func (s *store) allJobs() []Job {
  379. s.mu.RLock()
  380. defer s.mu.RUnlock()
  381. out := make([]Job, 0, len(s.jobs))
  382. for _, j := range s.jobs {
  383. out = append(out, *j)
  384. }
  385. sort.Slice(out, func(i, j int) bool { return out[i].Created > out[j].Created })
  386. return out
  387. }
  388. // gcJobs drops finished job records older than cutoff.
  389. func (s *store) gcJobs(cutoff int64, keepStatus map[string]bool) int {
  390. s.mu.Lock()
  391. defer s.mu.Unlock()
  392. n := 0
  393. for id, j := range s.jobs {
  394. if j.Created < cutoff && !keepStatus[j.Status] {
  395. delete(s.jobs, id)
  396. s.persist.del(tableJobs, id)
  397. n++
  398. }
  399. }
  400. return n
  401. }
  402. func (s *store) putSetting(in *Setting) bool {
  403. if in == nil || in.Key == "" {
  404. return false
  405. }
  406. s.mu.Lock()
  407. defer s.mu.Unlock()
  408. if existing, ok := s.settings[in.Key]; ok && in.Version <= existing.Version {
  409. return false
  410. }
  411. st := *in
  412. s.settings[st.Key] = &st
  413. s.persist.put(tableSettings, st.Key, st)
  414. return true
  415. }
  416. func (s *store) getSetting(key string) (*Setting, bool) {
  417. s.mu.RLock()
  418. defer s.mu.RUnlock()
  419. st, ok := s.settings[key]
  420. if !ok {
  421. return nil, false
  422. }
  423. c := *st
  424. return &c, true
  425. }
  426. /*
  427. Log and state
  428. */
  429. // appendLog assigns the next sequence number in the given term and persists.
  430. func (s *store) appendLog(e Entry, term uint64) Entry {
  431. s.mu.Lock()
  432. defer s.mu.Unlock()
  433. s.lastSeq++
  434. e.Seq = s.lastSeq
  435. e.Term = term
  436. s.persist.put(tableLog, logKey(e.Seq), e)
  437. return e
  438. }
  439. // logAfter returns up to max entries with Seq > after. ok is false when the
  440. // requested range was compacted away.
  441. func (s *store) logAfter(after uint64, max int) ([]Entry, bool) {
  442. //The log is read back from disk, so commit what is still queued first
  443. s.persist.flush()
  444. lastSeq := s.getLastSeq()
  445. entries, err := s.db.ListTableWithPrefix(tableLog, "")
  446. if err != nil {
  447. return nil, false
  448. }
  449. out := []Entry{}
  450. var first uint64
  451. for _, kv := range entries {
  452. var e Entry
  453. if json.Unmarshal(kv[1], &e) != nil {
  454. continue
  455. }
  456. if first == 0 {
  457. first = e.Seq
  458. }
  459. if e.Seq > after {
  460. out = append(out, e)
  461. if len(out) >= max {
  462. break
  463. }
  464. }
  465. }
  466. if after > 0 && (first > after+1 || (len(entries) == 0 && lastSeq > after)) {
  467. return nil, false
  468. }
  469. return out, true
  470. }
  471. // compactLog keeps only the newest keepLast entries.
  472. func (s *store) compactLog(keepLast int) {
  473. s.persist.flush()
  474. entries, err := s.db.ListTableWithPrefix(tableLog, "")
  475. if err != nil || len(entries) <= keepLast {
  476. return
  477. }
  478. for _, kv := range entries[:len(entries)-keepLast] {
  479. s.persist.del(tableLog, string(kv[0]))
  480. }
  481. }
  482. func (s *store) getLastSeq() uint64 {
  483. s.mu.RLock()
  484. defer s.mu.RUnlock()
  485. return s.lastSeq
  486. }
  487. func (s *store) setLastSeq(seq uint64) {
  488. s.mu.Lock()
  489. defer s.mu.Unlock()
  490. if seq > s.lastSeq {
  491. s.lastSeq = seq
  492. }
  493. }
  494. func (s *store) getApplied() (uint64, uint64) {
  495. s.mu.RLock()
  496. defer s.mu.RUnlock()
  497. return s.applied, s.term
  498. }
  499. func (s *store) setApplied(seq uint64, term uint64) {
  500. s.mu.Lock()
  501. defer s.mu.Unlock()
  502. s.applied = seq
  503. s.term = term
  504. s.persist.put(tableState, stateApplied, seq)
  505. s.persist.put(tableState, stateTerm, term)
  506. }
  507. func (s *store) getLease() Lease {
  508. s.leaseMu.RLock()
  509. defer s.leaseMu.RUnlock()
  510. return s.lease
  511. }
  512. func (s *store) setLease(l Lease) {
  513. s.leaseMu.Lock()
  514. s.lease = l
  515. s.leaseMu.Unlock()
  516. s.persist.put(tableState, stateLease, l)
  517. }
  518. func (s *store) addPending(id string, e Entry) {
  519. s.mu.Lock()
  520. defer s.mu.Unlock()
  521. s.pending[id] = e
  522. s.persist.put(tablePending, id, e)
  523. }
  524. func (s *store) removePending(id string) {
  525. s.mu.Lock()
  526. defer s.mu.Unlock()
  527. delete(s.pending, id)
  528. s.persist.del(tablePending, id)
  529. }
  530. func (s *store) pendingEntries() map[string]Entry {
  531. s.mu.RLock()
  532. defer s.mu.RUnlock()
  533. out := map[string]Entry{}
  534. for k, v := range s.pending {
  535. out[k] = v
  536. }
  537. return out
  538. }
  539. func (s *store) snapshot() Snapshot {
  540. s.mu.RLock()
  541. defer s.mu.RUnlock()
  542. snap := Snapshot{LastSeq: s.lastSeq, Term: s.term, Files: []FileRecord{}, Volumes: []Volume{}, Policies: []Policy{}, Jobs: []Job{}, Settings: []Setting{}}
  543. for _, rec := range s.files {
  544. snap.Files = append(snap.Files, *rec.Clone())
  545. }
  546. for _, v := range s.volumes {
  547. snap.Volumes = append(snap.Volumes, *v)
  548. }
  549. for _, p := range s.policies {
  550. snap.Policies = append(snap.Policies, *p)
  551. }
  552. for _, j := range s.jobs {
  553. snap.Jobs = append(snap.Jobs, *j)
  554. }
  555. for _, st := range s.settings {
  556. snap.Settings = append(snap.Settings, *st)
  557. }
  558. return snap
  559. }