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}