store.go 7.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281
  1. // client2server - Persistent storage (SQLite)
  2. //
  3. // Tables:
  4. // users(id, username, password_hash, role, created_at)
  5. // events(id, router_id, event_type, hostname, payload_json, ts, received_at)
  6. // commands(id, router_id, command, args_json, issued_by, status, result, ts)
  7. // alerts(id, router_id, kind, message, created_at, acknowledged_at)
  8. package main
  9. import (
  10. "database/sql"
  11. "encoding/json"
  12. "fmt"
  13. "log"
  14. "sync"
  15. "time"
  16. _ "modernc.org/sqlite"
  17. )
  18. var (
  19. db *sql.DB
  20. dbMu sync.Mutex
  21. )
  22. func initStore(path string) error {
  23. d, err := sql.Open("sqlite", path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)")
  24. if err != nil {
  25. return fmt.Errorf("open sqlite: %w", err)
  26. }
  27. d.SetMaxOpenConns(1) // SQLite single-writer; many readers OK
  28. if err := d.Ping(); err != nil {
  29. return fmt.Errorf("ping sqlite: %w", err)
  30. }
  31. db = d
  32. schema := `
  33. CREATE TABLE IF NOT EXISTS users (
  34. id INTEGER PRIMARY KEY AUTOINCREMENT,
  35. username TEXT UNIQUE NOT NULL,
  36. password_hash TEXT NOT NULL,
  37. role TEXT NOT NULL DEFAULT 'user',
  38. created_at DATETIME DEFAULT CURRENT_TIMESTAMP
  39. );
  40. CREATE TABLE IF NOT EXISTS events (
  41. id TEXT PRIMARY KEY,
  42. router_id TEXT NOT NULL,
  43. event_type TEXT NOT NULL,
  44. hostname TEXT,
  45. payload_json TEXT,
  46. ts DATETIME,
  47. received_at DATETIME DEFAULT CURRENT_TIMESTAMP
  48. );
  49. CREATE INDEX IF NOT EXISTS idx_events_router ON events(router_id);
  50. CREATE INDEX IF NOT EXISTS idx_events_ts ON events(received_at);
  51. CREATE INDEX IF NOT EXISTS idx_events_type ON events(event_type);
  52. CREATE TABLE IF NOT EXISTS commands (
  53. id TEXT PRIMARY KEY,
  54. router_id TEXT NOT NULL,
  55. command TEXT NOT NULL,
  56. args_json TEXT,
  57. issued_by TEXT,
  58. status TEXT NOT NULL,
  59. result_json TEXT,
  60. ts DATETIME DEFAULT CURRENT_TIMESTAMP
  61. );
  62. CREATE INDEX IF NOT EXISTS idx_commands_router ON commands(router_id);
  63. CREATE INDEX IF NOT EXISTS idx_commands_ts ON commands(ts);
  64. CREATE TABLE IF NOT EXISTS alerts (
  65. id INTEGER PRIMARY KEY AUTOINCREMENT,
  66. router_id TEXT NOT NULL,
  67. kind TEXT NOT NULL,
  68. message TEXT NOT NULL,
  69. created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
  70. acknowledged_at DATETIME
  71. );
  72. CREATE INDEX IF NOT EXISTS idx_alerts_unack ON alerts(acknowledged_at);
  73. `
  74. if _, err := db.Exec(schema); err != nil {
  75. return fmt.Errorf("schema: %w", err)
  76. }
  77. // Default admin user (username: admin, password: admin)
  78. // Only if no users exist
  79. var count int
  80. if err := db.QueryRow("SELECT COUNT(*) FROM users").Scan(&count); err != nil {
  81. return err
  82. }
  83. if count == 0 {
  84. hash, err := HashPassword("admin")
  85. if err != nil {
  86. return err
  87. }
  88. _, err = db.Exec(
  89. "INSERT INTO users (username, password_hash, role) VALUES (?, ?, ?)",
  90. "admin", hash, "system_admin",
  91. )
  92. if err != nil {
  93. return err
  94. }
  95. log.Println("created default admin user (username=admin password=admin) — CHANGE IT")
  96. }
  97. return nil
  98. }
  99. func saveEvent(ev RouterEvent) error {
  100. dbMu.Lock()
  101. defer dbMu.Unlock()
  102. payload, _ := json.Marshal(ev.Payload)
  103. var hostname *string
  104. if ev.Hostname != "" {
  105. hostname = &ev.Hostname
  106. }
  107. // Use ReceivedAt as fallback so listEvents never sees a zero time.
  108. ts := ev.ReceivedAt
  109. if !ev.Timestamp.IsZero() {
  110. ts = ev.Timestamp
  111. }
  112. _, err := db.Exec(
  113. `INSERT OR REPLACE INTO events (id, router_id, event_type, hostname, payload_json, ts, received_at)
  114. VALUES (?, ?, ?, ?, ?, ?, ?)`,
  115. ev.ID, ev.RouterID, ev.EventType, hostname, string(payload), ts, ev.ReceivedAt,
  116. )
  117. return err
  118. }
  119. func listEvents(limit int, routerID, eventType string) ([]RouterEvent, error) {
  120. dbMu.Lock()
  121. defer dbMu.Unlock()
  122. q := "SELECT id, router_id, event_type, COALESCE(hostname, ''), COALESCE(payload_json, '{}'), ts, received_at FROM events WHERE 1=1"
  123. args := []interface{}{}
  124. if routerID != "" {
  125. q += " AND router_id = ?"
  126. args = append(args, routerID)
  127. }
  128. if eventType != "" {
  129. q += " AND event_type = ?"
  130. args = append(args, eventType)
  131. }
  132. q += " ORDER BY received_at DESC LIMIT ?"
  133. args = append(args, limit)
  134. rows, err := db.Query(q, args...)
  135. if err != nil {
  136. return nil, err
  137. }
  138. defer rows.Close()
  139. out := make([]RouterEvent, 0, limit)
  140. for rows.Next() {
  141. var ev RouterEvent
  142. var payloadJSON string
  143. var ts time.Time
  144. if err := rows.Scan(&ev.ID, &ev.RouterID, &ev.EventType, &ev.Hostname, &payloadJSON, &ts, &ev.ReceivedAt); err != nil {
  145. return nil, err
  146. }
  147. _ = json.Unmarshal([]byte(payloadJSON), &ev.Payload)
  148. ev.Timestamp = ts
  149. out = append(out, ev)
  150. }
  151. return out, nil
  152. }
  153. func saveCommand(cmd RouterCommand, issuedBy, status string) error {
  154. dbMu.Lock()
  155. defer dbMu.Unlock()
  156. argsJSON, _ := json.Marshal(cmd.Args)
  157. _, err := db.Exec(
  158. `INSERT OR REPLACE INTO commands (id, router_id, command, args_json, issued_by, status) VALUES (?, ?, ?, ?, ?, ?)`,
  159. cmd.ID, cmd.RouterID, cmd.Command, string(argsJSON), issuedBy, status,
  160. )
  161. return err
  162. }
  163. func updateCommandResult(cmdID, status string, result *CommandResult) error {
  164. dbMu.Lock()
  165. defer dbMu.Unlock()
  166. var resultJSON []byte
  167. if result != nil {
  168. resultJSON, _ = json.Marshal(result)
  169. }
  170. _, err := db.Exec(
  171. "UPDATE commands SET status = ?, result_json = ? WHERE id = ?",
  172. status, string(resultJSON), cmdID,
  173. )
  174. return err
  175. }
  176. func listCommands(limit int, routerID string) ([]map[string]interface{}, error) {
  177. dbMu.Lock()
  178. defer dbMu.Unlock()
  179. q := `SELECT id, router_id, command, COALESCE(args_json, '{}'), COALESCE(issued_by, ''), status, COALESCE(result_json, ''), ts
  180. FROM commands WHERE 1=1`
  181. args := []interface{}{}
  182. if routerID != "" {
  183. q += " AND router_id = ?"
  184. args = append(args, routerID)
  185. }
  186. q += " ORDER BY ts DESC LIMIT ?"
  187. args = append(args, limit)
  188. rows, err := db.Query(q, args...)
  189. if err != nil {
  190. return nil, err
  191. }
  192. defer rows.Close()
  193. out := []map[string]interface{}{}
  194. for rows.Next() {
  195. var id, routerID, command, argsJSON, issuedBy, status, resultJSON, ts string
  196. if err := rows.Scan(&id, &routerID, &command, &argsJSON, &issuedBy, &status, &resultJSON, &ts); err != nil {
  197. return nil, err
  198. }
  199. entry := map[string]interface{}{
  200. "id": id, "router_id": routerID, "command": command,
  201. "issued_by": issuedBy, "status": status, "ts": ts,
  202. }
  203. var args map[string]string
  204. _ = json.Unmarshal([]byte(argsJSON), &args)
  205. entry["args"] = args
  206. if resultJSON != "" {
  207. var res map[string]interface{}
  208. _ = json.Unmarshal([]byte(resultJSON), &res)
  209. entry["result"] = res
  210. }
  211. out = append(out, entry)
  212. }
  213. return out, nil
  214. }
  215. func createAlert(routerID, kind, message string) {
  216. dbMu.Lock()
  217. defer dbMu.Unlock()
  218. _, _ = db.Exec(
  219. `INSERT INTO alerts (router_id, kind, message) VALUES (?, ?, ?)`,
  220. routerID, kind, message,
  221. )
  222. }
  223. func listAlerts(unackOnly bool, limit int) ([]map[string]interface{}, error) {
  224. dbMu.Lock()
  225. defer dbMu.Unlock()
  226. q := `SELECT id, router_id, kind, message, created_at, COALESCE(acknowledged_at, '') FROM alerts`
  227. if unackOnly {
  228. q += ` WHERE acknowledged_at IS NULL`
  229. }
  230. q += ` ORDER BY created_at DESC LIMIT ?`
  231. rows, err := db.Query(q, limit)
  232. if err != nil {
  233. return nil, err
  234. }
  235. defer rows.Close()
  236. out := []map[string]interface{}{}
  237. for rows.Next() {
  238. var id int64
  239. var routerID, kind, message, createdAt, ackedAt string
  240. if err := rows.Scan(&id, &routerID, &kind, &message, &createdAt, &ackedAt); err != nil {
  241. return nil, err
  242. }
  243. entry := map[string]interface{}{
  244. "id": id, "router_id": routerID, "kind": kind,
  245. "message": message, "created_at": createdAt,
  246. }
  247. if ackedAt != "" {
  248. entry["acknowledged_at"] = ackedAt
  249. }
  250. out = append(out, entry)
  251. }
  252. return out, nil
  253. }
  254. func acknowledgeAlert(id int64) error {
  255. dbMu.Lock()
  256. defer dbMu.Unlock()
  257. _, err := db.Exec("UPDATE alerts SET acknowledged_at = CURRENT_TIMESTAMP WHERE id = ?", id)
  258. return err
  259. }