persistence.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428
  1. package sharedspace
  2. /*
  3. SharedSpace persistence
  4. Persistent spaces write through to the ArozOS system database and
  5. keep their blobs under a durable storage root, so chat history,
  6. shared files and collaborative documents survive server restarts.
  7. Layout (values are JSON, one bolt transaction per write):
  8. table "sharedspace": key "space/<spaceID>" -> spaceRecord
  9. (key "conf/<name>" is reserved for the
  10. admin configuration written by the main
  11. package)
  12. table "sharedspaceitem": key "<spaceID>/<itemID>" -> itemRecord
  13. table "sharedspacedoc": key "<spaceID>/<docID>" -> docRecord
  14. Blobs live at <persistRoot>/<spaceID>/<itemID>; disk paths are never
  15. persisted - they are rebuilt from the root at reload. Reload is
  16. self-healing: records whose space vanished, blob went missing or
  17. JSON no longer parses are deleted, and blob directories without a
  18. live space are removed.
  19. */
  20. import (
  21. "encoding/json"
  22. "os"
  23. "path/filepath"
  24. "sort"
  25. "strconv"
  26. "strings"
  27. "time"
  28. "imuslab.com/arozos/mod/info/logger"
  29. )
  30. const (
  31. dbTableSpaces = "sharedspace"
  32. dbTableItems = "sharedspaceitem"
  33. dbTableDocs = "sharedspacedoc"
  34. dbSpaceKeyPrefix = "space/"
  35. )
  36. // spaceRecord is the persisted form of a Space.
  37. type spaceRecord struct {
  38. ID string `json:"id"`
  39. Name string `json:"name"`
  40. Owner string `json:"owner"`
  41. CreatedAt int64 `json:"createdat"`
  42. Access string `json:"access"`
  43. Metadata map[string]string `json:"metadata"`
  44. Members map[string]string `json:"members"`
  45. MaxItems int `json:"maxitems"`
  46. }
  47. // itemRecord is the persisted form of an Item. Disk paths are rebuilt from
  48. // the persistent root at reload and never stored.
  49. type itemRecord struct {
  50. ID string `json:"id"`
  51. Type string `json:"type"`
  52. Name string `json:"name"`
  53. Text string `json:"text"`
  54. Size int64 `json:"size"`
  55. Uploader string `json:"uploader"`
  56. Origin string `json:"origin"`
  57. Seq int64 `json:"seq"`
  58. CreatedAt int64 `json:"createdat"`
  59. HasBlob bool `json:"hasblob"`
  60. }
  61. // docRecord is the persisted form of a document (current revision only; the
  62. // in-memory history ring is not persisted).
  63. type docRecord struct {
  64. ID string `json:"id"`
  65. Name string `json:"name"`
  66. Creator string `json:"creator"`
  67. Content string `json:"content"`
  68. Revision int64 `json:"revision"`
  69. CreatedAt int64 `json:"createdat"`
  70. UpdatedAt int64 `json:"updatedat"`
  71. UpdatedBy string `json:"updatedby"`
  72. }
  73. // initPersistenceTables creates the sharedspace tables when missing.
  74. func (m *Manager) initPersistenceTables() {
  75. m.db.NewTable(dbTableSpaces)
  76. m.db.NewTable(dbTableItems)
  77. m.db.NewTable(dbTableDocs)
  78. }
  79. // record snapshots the space's persisted state.
  80. func (s *Space) record() *spaceRecord {
  81. s.mu.Lock()
  82. defer s.mu.Unlock()
  83. metadata := make(map[string]string, len(s.metadata))
  84. for key, value := range s.metadata {
  85. metadata[key] = value
  86. }
  87. members := make(map[string]string, len(s.members))
  88. for username, role := range s.members {
  89. members[username] = role
  90. }
  91. return &spaceRecord{
  92. ID: s.ID,
  93. Name: s.Name,
  94. Owner: s.Owner,
  95. CreatedAt: s.CreatedAt.Unix(),
  96. Access: s.access,
  97. Metadata: metadata,
  98. Members: members,
  99. MaxItems: s.maxItems,
  100. }
  101. }
  102. // persistSpace writes the space record through to the database. No-op for
  103. // ephemeral spaces or persistence-less managers.
  104. func (m *Manager) persistSpace(s *Space) {
  105. if s == nil || !s.Persistent || m.db == nil {
  106. return
  107. }
  108. err := m.db.Write(dbTableSpaces, dbSpaceKeyPrefix+s.ID, s.record())
  109. if err != nil {
  110. logger.PrintAndLog("SharedSpace", "Failed to persist space "+s.ID, err)
  111. }
  112. }
  113. // mgrPersistSpace lets Space methods write their record through without
  114. // holding any lock. Nil-safe for spaces built directly in tests.
  115. func (s *Space) mgrPersistSpace() {
  116. if s.mgr != nil {
  117. s.mgr.persistSpace(s)
  118. }
  119. }
  120. // mgrPersistItem writes an item record through to the database.
  121. func (s *Space) mgrPersistItem(item *Item) {
  122. if s.mgr == nil || !s.Persistent || s.mgr.db == nil {
  123. return
  124. }
  125. record := &itemRecord{
  126. ID: item.ID,
  127. Type: item.Type,
  128. Name: item.Name,
  129. Text: item.Text,
  130. Size: item.Size,
  131. Uploader: item.Uploader,
  132. Origin: item.Origin,
  133. Seq: item.Seq,
  134. CreatedAt: item.CreatedAt.Unix(),
  135. HasBlob: item.DiskPath != "",
  136. }
  137. err := s.mgr.db.Write(dbTableItems, s.ID+"/"+item.ID, record)
  138. if err != nil {
  139. logger.PrintAndLog("SharedSpace", "Failed to persist item "+item.ID, err)
  140. }
  141. }
  142. // mgrDeleteItemRecord removes an item record from the database.
  143. func (s *Space) mgrDeleteItemRecord(itemID string) {
  144. if s.mgr == nil || !s.Persistent || s.mgr.db == nil {
  145. return
  146. }
  147. s.mgr.db.Delete(dbTableItems, s.ID+"/"+itemID)
  148. }
  149. // mgrPersistDoc writes a document record through to the database.
  150. func (s *Space) mgrPersistDoc(snapshot *DocSnapshot) {
  151. if s.mgr == nil || !s.Persistent || s.mgr.db == nil {
  152. return
  153. }
  154. record := &docRecord{
  155. ID: snapshot.ID,
  156. Name: snapshot.Name,
  157. Creator: snapshot.Creator,
  158. Content: snapshot.Content,
  159. Revision: snapshot.Revision,
  160. CreatedAt: snapshot.CreatedAt.Unix(),
  161. UpdatedAt: snapshot.UpdatedAt.Unix(),
  162. UpdatedBy: snapshot.UpdatedBy,
  163. }
  164. err := s.mgr.db.Write(dbTableDocs, s.ID+"/"+snapshot.ID, record)
  165. if err != nil {
  166. logger.PrintAndLog("SharedSpace", "Failed to persist document "+snapshot.ID, err)
  167. }
  168. }
  169. // mgrDeleteDocRecord removes a document record from the database.
  170. func (s *Space) mgrDeleteDocRecord(docID string) {
  171. if s.mgr == nil || !s.Persistent || s.mgr.db == nil {
  172. return
  173. }
  174. s.mgr.db.Delete(dbTableDocs, s.ID+"/"+docID)
  175. }
  176. // deleteSpaceRecords removes every database record belonging to a space.
  177. func (m *Manager) deleteSpaceRecords(spaceID string) {
  178. if m.db == nil {
  179. return
  180. }
  181. m.db.Delete(dbTableSpaces, dbSpaceKeyPrefix+spaceID)
  182. for _, table := range []string{dbTableItems, dbTableDocs} {
  183. entries, err := m.db.ListTable(table)
  184. if err != nil {
  185. continue
  186. }
  187. for _, entry := range entries {
  188. key := string(entry[0])
  189. if strings.HasPrefix(key, spaceID+"/") {
  190. m.db.Delete(table, key)
  191. }
  192. }
  193. }
  194. }
  195. // loadPersistedSpaces rebuilds every persistent space from the database at
  196. // construction time, restoring items (with blobs) and documents.
  197. func (m *Manager) loadPersistedSpaces() {
  198. entries, err := m.db.ListTable(dbTableSpaces)
  199. if err != nil {
  200. logger.PrintAndLog("SharedSpace", "Failed to list persisted spaces", err)
  201. return
  202. }
  203. //1. Space records
  204. for _, entry := range entries {
  205. key := string(entry[0])
  206. if !strings.HasPrefix(key, dbSpaceKeyPrefix) {
  207. continue //conf/* and future keys
  208. }
  209. record := spaceRecord{}
  210. if err := json.Unmarshal(entry[1], &record); err != nil || record.ID == "" {
  211. logger.PrintAndLog("SharedSpace", "Dropping corrupted space record "+key, err)
  212. m.db.Delete(dbTableSpaces, key)
  213. continue
  214. }
  215. if !validAccessMode(record.Access) {
  216. record.Access = AccessOpen
  217. }
  218. if record.MaxItems <= 0 {
  219. record.MaxItems = m.defaultMaxItems
  220. }
  221. members := record.Members
  222. if members == nil {
  223. members = map[string]string{}
  224. }
  225. members[record.Owner] = RoleOwner
  226. metadata := record.Metadata
  227. if metadata == nil {
  228. metadata = map[string]string{}
  229. }
  230. space := &Space{
  231. ID: record.ID,
  232. Name: record.Name,
  233. Owner: record.Owner,
  234. access: record.Access,
  235. Persistent: true,
  236. CreatedAt: time.Unix(record.CreatedAt, 0),
  237. storageDir: filepath.Join(m.persistRoot, record.ID),
  238. metadata: metadata,
  239. members: members,
  240. docs: make(map[string]*Doc),
  241. items: []*Item{},
  242. itemIdx: make(map[string]*Item),
  243. listeners: make(map[string]func(*Item)),
  244. evListeners: make(map[string]func(*SpaceEvent)),
  245. nextSeq: 1,
  246. maxItems: record.MaxItems,
  247. trimOldest: true,
  248. mgr: m,
  249. }
  250. m.spaces[space.ID] = space
  251. }
  252. //2. Item records
  253. itemEntries, err := m.db.ListTable(dbTableItems)
  254. if err == nil {
  255. for _, entry := range itemEntries {
  256. key := string(entry[0])
  257. slash := strings.Index(key, "/")
  258. if slash <= 0 {
  259. m.db.Delete(dbTableItems, key)
  260. continue
  261. }
  262. spaceID := key[:slash]
  263. space, exists := m.spaces[spaceID]
  264. if !exists {
  265. m.db.Delete(dbTableItems, key) //space vanished: self-heal
  266. continue
  267. }
  268. record := itemRecord{}
  269. if err := json.Unmarshal(entry[1], &record); err != nil || record.ID == "" {
  270. logger.PrintAndLog("SharedSpace", "Dropping corrupted item record "+key, err)
  271. m.db.Delete(dbTableItems, key)
  272. continue
  273. }
  274. item := &Item{
  275. ID: record.ID,
  276. Type: record.Type,
  277. Name: record.Name,
  278. Text: record.Text,
  279. Size: record.Size,
  280. Uploader: record.Uploader,
  281. Origin: record.Origin,
  282. Seq: record.Seq,
  283. CreatedAt: time.Unix(record.CreatedAt, 0),
  284. }
  285. if record.HasBlob {
  286. item.DiskPath = filepath.Join(m.persistRoot, spaceID, record.ID)
  287. if _, err := os.Stat(item.DiskPath); err != nil {
  288. m.db.Delete(dbTableItems, key) //blob vanished: self-heal
  289. continue
  290. }
  291. }
  292. space.items = append(space.items, item)
  293. space.itemIdx[item.ID] = item
  294. }
  295. }
  296. //3. Document records
  297. docEntries, err := m.db.ListTable(dbTableDocs)
  298. if err == nil {
  299. for _, entry := range docEntries {
  300. key := string(entry[0])
  301. slash := strings.Index(key, "/")
  302. if slash <= 0 {
  303. m.db.Delete(dbTableDocs, key)
  304. continue
  305. }
  306. spaceID := key[:slash]
  307. space, exists := m.spaces[spaceID]
  308. if !exists {
  309. m.db.Delete(dbTableDocs, key)
  310. continue
  311. }
  312. record := docRecord{}
  313. if err := json.Unmarshal(entry[1], &record); err != nil || record.ID == "" {
  314. logger.PrintAndLog("SharedSpace", "Dropping corrupted document record "+key, err)
  315. m.db.Delete(dbTableDocs, key)
  316. continue
  317. }
  318. space.docs[record.ID] = &Doc{
  319. id: record.ID,
  320. name: record.Name,
  321. creator: record.Creator,
  322. createdAt: time.Unix(record.CreatedAt, 0),
  323. content: record.Content,
  324. revision: record.Revision,
  325. updatedAt: time.Unix(record.UpdatedAt, 0),
  326. updatedBy: record.UpdatedBy,
  327. }
  328. }
  329. }
  330. //4. Restore chronological ordering and the per-space sequence counter
  331. loadedSpaces := 0
  332. for _, space := range m.spaces {
  333. sort.Slice(space.items, func(i, j int) bool {
  334. return space.items[i].Seq < space.items[j].Seq
  335. })
  336. if count := len(space.items); count > 0 {
  337. space.nextSeq = space.items[count-1].Seq + 1
  338. }
  339. loadedSpaces++
  340. }
  341. //5. Remove blob directories that no longer belong to a live space
  342. if dirEntries, err := os.ReadDir(m.persistRoot); err == nil {
  343. for _, dirEntry := range dirEntries {
  344. if !dirEntry.IsDir() {
  345. continue
  346. }
  347. if _, exists := m.spaces[dirEntry.Name()]; !exists {
  348. os.RemoveAll(filepath.Join(m.persistRoot, dirEntry.Name()))
  349. }
  350. }
  351. }
  352. if loadedSpaces > 0 {
  353. logger.PrintAndLog("SharedSpace", "Restored "+strconv.Itoa(loadedSpaces)+" persistent shared space(s)", nil)
  354. }
  355. }
  356. // SweepStaleSpaces deletes persistent spaces whose most recent activity
  357. // (newest item, newest document update, or creation) is older than maxAge
  358. // and that have no live channel subscribers. It returns the IDs it deleted.
  359. // Ephemeral spaces are left to their owning subsystems (e.g. MeetRoom).
  360. func (m *Manager) SweepStaleSpaces(maxAge time.Duration) []string {
  361. if maxAge <= 0 {
  362. return nil
  363. }
  364. cutoff := time.Now().Add(-maxAge)
  365. var stale []string
  366. m.mu.RLock()
  367. for id, space := range m.spaces {
  368. if !space.Persistent {
  369. continue
  370. }
  371. space.mu.Lock()
  372. last := space.CreatedAt
  373. for _, item := range space.items {
  374. if item.CreatedAt.After(last) {
  375. last = item.CreatedAt
  376. }
  377. }
  378. for _, doc := range space.docs {
  379. if doc.updatedAt.After(last) {
  380. last = doc.updatedAt
  381. }
  382. }
  383. channel := space.channel
  384. space.mu.Unlock()
  385. if channel != nil && channel.Count() > 0 {
  386. continue
  387. }
  388. if last.Before(cutoff) {
  389. stale = append(stale, id)
  390. }
  391. }
  392. m.mu.RUnlock()
  393. for _, id := range stale {
  394. m.DeleteSpace(id)
  395. logger.PrintAndLog("SharedSpace", "Removed stale shared space "+id, nil)
  396. }
  397. return stale
  398. }