This repository has no description
1.9 kB
108 lines
1package eventstream
2
3import (
4 "database/sql"
5 "encoding/json"
6 "sync"
7 "time"
8
9 "tangled.org/core/notifier"
10)
11
12type Store interface {
13 Exec(query string, args ...any) (sql.Result, error)
14 Query(query string, args ...any) (*sql.Rows, error)
15}
16
17var (
18 clockMu sync.Mutex
19 lastNanos int64
20)
21
22func HighWater(s Store) (int64, error) {
23 clockMu.Lock()
24 defer clockMu.Unlock()
25
26 rows, err := s.Query(`select coalesce(max(created), 0) from events`)
27 if err != nil {
28 return 0, err
29 }
30 defer rows.Close()
31
32 var created int64
33 if !rows.Next() {
34 if err := rows.Err(); err != nil {
35 return 0, err
36 }
37 return 0, sql.ErrNoRows
38 }
39 if err := rows.Scan(&created); err != nil {
40 return 0, err
41 }
42 if err := rows.Err(); err != nil {
43 return 0, err
44 }
45 if created > lastNanos {
46 lastNanos = created
47 }
48 return lastNanos, nil
49}
50
51func Insert(s Store, ev Event, n *notifier.Notifier) error {
52 clockMu.Lock()
53 defer clockMu.Unlock()
54
55 if ev.Created == 0 {
56 now := time.Now().UnixNano()
57 if now <= lastNanos {
58 now = lastNanos + 1
59 }
60 ev.Created = now
61 }
62 if ev.Created > lastNanos {
63 lastNanos = ev.Created
64 }
65
66 if _, err := s.Exec(
67 `insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`,
68 ev.Rkey,
69 ev.Nsid,
70 []byte(ev.EventJson),
71 ev.Created,
72 ); err != nil {
73 return err
74 }
75 n.NotifyAll()
76 return nil
77}
78
79func List(s Store, cursor int64, limit int) ([]Event, error) {
80 rows, err := s.Query(`
81 select rkey, nsid, event, created
82 from events
83 where created > ?
84 order by created asc
85 limit ?
86 `, cursor, limit)
87 if err != nil {
88 return nil, err
89 }
90 defer rows.Close()
91
92 var out []Event
93 for rows.Next() {
94 var ev Event
95 var eventJsonStr string
96 if err := rows.Scan(&ev.Rkey, &ev.Nsid, &eventJsonStr, &ev.Created); err != nil {
97 return nil, err
98 }
99 ev.EventJson = json.RawMessage(eventJsonStr)
100 out = append(out, ev)
101 }
102
103 if err := rows.Err(); err != nil {
104 return nil, err
105 }
106
107 return out, nil
108}