service.go
1package document
2
3import (
4 "database/sql"
5 "embed"
6 "encoding/json"
7 "fmt"
8 "log/slog"
9 "net/http"
10 "os"
11 "path"
12 "sync"
13 "time"
14
15 "git.kilimanjaro.io/rtw/pkg/id"
16 "git.kilimanjaro.io/rtw/pkg/log"
17 "github.com/pressly/goose/v3"
18)
19
20var (
21 ErrSystemIssue error = fmt.Errorf("unexpected system error")
22)
23
24//go:embed schema/*.sql
25var migrations embed.FS
26
27type Document struct {
28 ResourceID int64
29 ID string
30 Content []byte
31 Mtime time.Time
32 Ctime time.Time
33}
34
35type Service struct {
36 db *sql.DB
37 basePath string
38 mu *sync.RWMutex
39 ySweetURL string
40}
41
42type ServiceOption func(*Service) error
43
44// WithBasePath sets the base path for the document service. Database will be located at
45// `{basepath}/data/documents.db`.
46func WithBasePath(basePath string) ServiceOption {
47 return func(s *Service) error {
48 s.basePath = basePath
49 return nil
50 }
51}
52
53// WithYSweetURL sets the URL to the upstream y-sweet API. Defaults to http://127.0.0.1:9000.
54func WithYSweetURL(url string) ServiceOption {
55 return func(s *Service) error {
56 s.ySweetURL = url
57 return nil
58 }
59}
60
61func NewService(opts ...ServiceOption) (*Service, error) {
62 s := &Service{
63 mu: &sync.RWMutex{},
64 ySweetURL: "http://127.0.0.1:9000",
65 }
66 for _, opt := range opts {
67 if err := opt(s); err != nil {
68 return nil, err
69 }
70 }
71 dataDir := path.Join(s.basePath, "data")
72 if err := os.MkdirAll(dataDir, 0755); err != nil {
73 return nil, fmt.Errorf("failed to create data directory: %w", err)
74 }
75 db, err := sql.Open("sqlite3", "file:"+path.Join(dataDir, "documents.db")+"?_journal=WAL&_timeout=5000&_fk=1&cache=shared")
76 if err != nil {
77 return nil, fmt.Errorf("error opening database: %w", err)
78 }
79 s.db = db
80 // apply migrations
81 goose.SetBaseFS(migrations)
82 goose.SetLogger(log.GooseLogger(slog.String("migrations", "document/schema")))
83 if err := goose.SetDialect("sqlite3"); err != nil {
84 return nil, fmt.Errorf("error setting document database type: %w", err)
85 }
86 if err := goose.Up(db, "schema", goose.WithNoVersioning()); err != nil {
87 return nil, fmt.Errorf("error migrating document tables: %w", err)
88 }
89 return s, nil
90}
91
92func (s *Service) DocumentExists(docID string) (id.Key, bool, error) {
93 s.mu.RLock()
94 defer s.mu.RUnlock()
95
96 var result int
97 var resourceID sql.NullInt64
98 if err := s.db.QueryRow("SELECT COUNT(*), resource_id FROM documents WHERE id = ?", docID).Scan(&result, &resourceID); err != nil {
99 return id.NotExist, false, fmt.Errorf("error executing document exists query: %w: %w", ErrSystemIssue, err)
100 }
101 if result == 0 || !resourceID.Valid {
102 return id.NotExist, false, nil
103 }
104 return id.Key(resourceID.Int64), true, nil
105}
106
107func (s *Service) NewDocument(docID string) (id.Key, error) {
108 s.mu.Lock()
109 defer s.mu.Unlock()
110
111 n := id.New()
112 result, err := s.db.Exec("INSERT OR IGNORE INTO documents (resource_id, id) VALUES (?, ?)", n, docID)
113 if err != nil {
114 return id.NotExist, fmt.Errorf("error inserting document: %w: %w", ErrSystemIssue, err)
115 }
116 numRows, err := result.RowsAffected()
117 if err != nil {
118 return id.NotExist, fmt.Errorf("error inserting into document: %w: %w", ErrSystemIssue, err)
119 }
120 if numRows == 0 {
121 var existingID int64
122 if err := s.db.QueryRow("SELECT resource_id FROM documents WHERE id = ?", docID).Scan(&existingID); err != nil {
123 return id.NotExist, fmt.Errorf("error getting existing document: %w: %w", ErrSystemIssue, err)
124 }
125 return id.Key(existingID), nil
126 }
127 return n, nil
128}
129
130func (s *Service) NewDocumentWithID(resourceID id.Key, docID string) error {
131 s.mu.Lock()
132 defer s.mu.Unlock()
133
134 _, err := s.db.Exec("INSERT OR IGNORE INTO documents (resource_id, id) VALUES (?, ?)", resourceID, docID)
135 if err != nil {
136 return fmt.Errorf("error inserting document: %w: %w", ErrSystemIssue, err)
137 }
138 return nil
139}
140
141func (s *Service) ForkDocument(fromDocID string, toDocID string) (id.Key, error) {
142 existingKey, exists, err := s.DocumentExists(toDocID)
143 if exists {
144 return existingKey, nil
145 }
146 if err != nil {
147 return id.NotExist, fmt.Errorf("forking document: doc exists: %w", err)
148 }
149 key, err := s.NewDocument(toDocID)
150 if err != nil {
151 return id.NotExist, fmt.Errorf("forking document: %w", err)
152 }
153 gResp, err := http.Get(fmt.Sprintf("%s/d/%s/as-update", s.ySweetURL, fromDocID))
154 if err != nil || gResp.StatusCode >= 300 {
155 return id.NotExist, fmt.Errorf("forking document: failed to get forked doc: %w", err)
156 }
157 defer gResp.Body.Close()
158 pResp, err := http.Post(fmt.Sprintf("%s/d/%s/update", s.ySweetURL, toDocID), gResp.Header.Get("Content-Type"), gResp.Body)
159 if err != nil || pResp.StatusCode >= 300 {
160 return id.NotExist, fmt.Errorf("forking document: failed to post update: %w", err)
161 }
162 defer pResp.Body.Close()
163 return key, nil
164}
165
166func (s *Service) SetContent(docID string, content []byte) (id.Key, error) {
167 s.mu.Lock()
168 defer s.mu.Unlock()
169
170 var n int64
171 if err := s.db.QueryRow("UPDATE documents SET content = ? WHERE id = ? RETURNING resource_id", content, docID).Scan(&n); err != nil {
172 return id.NotExist, fmt.Errorf("error inserting content for id=%s: %w: %w", docID, ErrSystemIssue, err)
173 }
174 return id.Key(n), nil
175}
176
177func (s *Service) GetContent(docID string) (DocumentContent, error) {
178 s.mu.RLock()
179 defer s.mu.RUnlock()
180
181 var content []byte
182 if err := s.db.QueryRow("SELECT content FROM documents WHERE id = ?", docID).Scan(&content); err != nil {
183 return DocumentContent{}, fmt.Errorf("error getting document content: %w: %w", ErrSystemIssue, err)
184 }
185 if len(content) == 0 {
186 // return empty document
187 return DocumentContent{[]Node{}}, nil
188 }
189 var doc DocumentContent
190 if err := json.Unmarshal(content, &doc); err != nil {
191 return DocumentContent{}, fmt.Errorf("error unmarshaling content: %w: %w", ErrSystemIssue, err)
192 }
193 return doc, nil
194}
195
196func (s *Service) GetContentByResourceID(resourceID id.Key) (DocumentContent, error) {
197 s.mu.RLock()
198 defer s.mu.RUnlock()
199 var content []byte
200 if err := s.db.QueryRow("SELECT content FROM documents WHERE resource_id = ?", resourceID).Scan(&content); err != nil {
201 return DocumentContent{}, fmt.Errorf("error getting document content: %w: %w", ErrSystemIssue, err)
202 }
203 if len(content) == 0 {
204 // return empty document
205 return DocumentContent{[]Node{}}, nil
206 }
207 var doc DocumentContent
208 if err := json.Unmarshal(content, &doc); err != nil {
209 return DocumentContent{}, fmt.Errorf("error unmarshaling content: %w: %w", ErrSystemIssue, err)
210 }
211 return doc, nil
212}
213
214func (s *Service) Shutdown() error {
215 s.mu.Lock()
216 defer s.mu.Unlock()
217
218 if s.db != nil {
219 return s.db.Close()
220 }
221 return nil
222}