Eric Bower
·
2026-08-17
1package postgres
2
3import (
4 "database/sql"
5 "errors"
6 "fmt"
7 "log/slog"
8 "math"
9 "slices"
10 "sort"
11 "strings"
12 "sync"
13 "time"
14
15 "github.com/jmoiron/sqlx"
16 _ "github.com/lib/pq"
17 "github.com/picosh/pico/pkg/db"
18 "github.com/picosh/pico/pkg/shared"
19)
20
21// mobileUserAgentExpr is a SQL expression to detect mobile user agents.
22const mobileUserAgentExpr = `user_agent ILIKE '%mobile%' OR user_agent ILIKE '%android%' OR user_agent ILIKE '%iphone%' OR user_agent ILIKE '%ipad%' OR user_agent ILIKE '%ipod%' OR user_agent ILIKE '%blackberry%' OR user_agent ILIKE '%windows phone%'`
23
24var PAGER_SIZE = 15
25
26var SelectPost = `
27 posts.id, user_id, app_users.name, filename, slug, title, text, description,
28 posts.created_at, publish_at, posts.updated_at, hidden, file_size, mime_type, shasum, data, expires_at, views`
29
30type PsqlDB struct {
31 Logger *slog.Logger
32 Db *sqlx.DB
33}
34
35type RowScanner interface {
36 Scan(dest ...any) error
37}
38
39func CreatePostWithTagsByRow(r RowScanner) (*db.Post, error) {
40 post := &db.Post{}
41 tagStr := ""
42 err := r.Scan(
43 &post.ID,
44 &post.UserID,
45 &post.Username,
46 &post.Filename,
47 &post.Slug,
48 &post.Title,
49 &post.Text,
50 &post.Description,
51 &post.CreatedAt,
52 &post.PublishAt,
53 &post.UpdatedAt,
54 &post.Hidden,
55 &post.FileSize,
56 &post.MimeType,
57 &post.Shasum,
58 &post.Data,
59 &post.ExpiresAt,
60 &post.Views,
61 &tagStr,
62 )
63 if err != nil {
64 return nil, err
65 }
66
67 tags := strings.Split(tagStr, ",")
68 for _, tag := range tags {
69 tg := strings.TrimSpace(tag)
70 if tg == "" {
71 continue
72 }
73 post.Tags = append(post.Tags, tg)
74 }
75
76 return post, nil
77}
78
79func NewDB(databaseUrl string, logger *slog.Logger) *PsqlDB {
80 var err error
81 d := &PsqlDB{
82 Logger: logger,
83 }
84 d.Logger.Info("Connecting to postgres", "databaseUrl", databaseUrl)
85
86 db, err := sqlx.Connect("postgres", databaseUrl)
87 if err != nil {
88 d.Logger.Error(err.Error())
89 }
90 d.Db = db
91 return d
92}
93
94func (me *PsqlDB) shouldBlockSingup(ip string) error {
95 blocked := &db.BlockSignups{}
96 err := me.Db.Get(blocked, `SELECT * FROM block_signups WHERE ip = $1`, ip)
97 // an error in this case means the result is empty (no record found)
98 if err != nil {
99 return nil
100 }
101 return fmt.Errorf("your IP address has been blocked: %s", blocked.Reason)
102}
103
104func (me *PsqlDB) RegisterUser(username, pubkey, comment, ip string) (*db.User, error) {
105 lowerName := strings.ToLower(username)
106 valid, err := me.validateName(lowerName)
107 if !valid {
108 return nil, err
109 }
110
111 me.Logger.Info("checking if ip is in block list", "ip", ip, "username", username)
112 err = me.shouldBlockSingup(ip)
113 if err != nil {
114 me.Logger.Warn("user has been blocked from signing up", "ip", ip, "username", username, "err", err)
115 return nil, err
116 }
117
118 tx, err := me.Db.Beginx()
119 if err != nil {
120 return nil, err
121 }
122 defer func() {
123 _ = tx.Rollback()
124 }()
125
126 var id string
127 err = tx.QueryRow(`INSERT INTO app_users (name) VALUES($1) returning id`, lowerName).Scan(&id)
128 if err != nil {
129 return nil, err
130 }
131
132 err = me.insertPublicKeyWithTx(id, pubkey, comment, tx)
133 if err != nil {
134 return nil, err
135 }
136
137 err = tx.Commit()
138 if err != nil {
139 return nil, err
140 }
141
142 return me.FindUserByKey(username, pubkey)
143}
144
145func (me *PsqlDB) insertPublicKeyWithTx(userID, key, name string, tx *sqlx.Tx) error {
146 pk, _ := me.findPublicKeyByKey(key)
147 if pk != nil {
148 return db.ErrPublicKeyTaken
149 }
150 query := `INSERT INTO public_keys (user_id, public_key, name) VALUES ($1, $2, $3)`
151 _, err := tx.Exec(query, userID, key, name)
152 return err
153}
154
155func (me *PsqlDB) InsertPublicKey(userID, key, name string) error {
156 pk, _ := me.findPublicKeyByKey(key)
157 if pk != nil {
158 return db.ErrPublicKeyTaken
159 }
160 query := `INSERT INTO public_keys (user_id, public_key, name) VALUES ($1, $2, $3)`
161 _, err := me.Db.Exec(query, userID, key, name)
162 return err
163}
164
165func (me *PsqlDB) UpdatePublicKey(pubkeyID, name string) (*db.PublicKey, error) {
166 pk, err := me.findPublicKey(pubkeyID)
167 if err != nil {
168 return nil, err
169 }
170
171 query := `UPDATE public_keys SET name=$1 WHERE id=$2;`
172 _, err = me.Db.Exec(query, name, pk.ID)
173 if err != nil {
174 return nil, err
175 }
176
177 pk, err = me.findPublicKey(pubkeyID)
178 if err != nil {
179 return nil, err
180 }
181 return pk, nil
182}
183
184func (me *PsqlDB) findPublicKeyByKey(key string) (*db.PublicKey, error) {
185 var keys []*db.PublicKey
186 rs, err := me.Db.Queryx(`SELECT id, user_id, name, public_key, created_at FROM public_keys WHERE public_key = $1`, key)
187 if err != nil {
188 return nil, err
189 }
190 defer func() { _ = rs.Close() }()
191
192 for rs.Next() {
193 pk := &db.PublicKey{}
194 err := rs.Scan(&pk.ID, &pk.UserID, &pk.Name, &pk.Key, &pk.CreatedAt)
195 if err != nil {
196 return nil, err
197 }
198
199 keys = append(keys, pk)
200 }
201
202 if rs.Err() != nil {
203 return nil, rs.Err()
204 }
205
206 if len(keys) == 0 {
207 return nil, fmt.Errorf("pubkey not found in our database: [%s]", key)
208 }
209
210 // When we run PublicKeyByKey and there are multiple public keys returned from the database
211 // that should mean that we don't have the correct username for this public key.
212 // When that happens we need to reject the authentication and ask the user to provide the correct
213 // username when using ssh. So instead of `ssh <domain>` it should be `ssh user@<domain>`
214 if len(keys) > 1 {
215 return nil, &db.ErrMultiplePublicKeys{}
216 }
217
218 return keys[0], nil
219}
220
221func (me *PsqlDB) findPublicKey(pubkeyID string) (*db.PublicKey, error) {
222 pk := &db.PublicKey{}
223 err := me.Db.Get(pk, `SELECT * FROM public_keys WHERE id = $1`, pubkeyID)
224 if err != nil {
225 return nil, err
226 }
227 return pk, nil
228}
229
230func (me *PsqlDB) FindKeysByUser(user *db.User) ([]*db.PublicKey, error) {
231 var keys []*db.PublicKey
232 err := me.Db.Select(&keys, `SELECT * FROM public_keys WHERE user_id = $1 ORDER BY created_at ASC`, user.ID)
233 if err != nil {
234 return nil, err
235 }
236 return keys, nil
237}
238
239func (me *PsqlDB) RemoveKeys(keyIDs []string) error {
240 param := "{" + strings.Join(keyIDs, ",") + "}"
241 _, err := me.Db.Exec(`DELETE FROM public_keys WHERE id = ANY($1::uuid[])`, param)
242 return err
243}
244
245func (me *PsqlDB) FindUsersWithPost(space string) ([]*db.User, error) {
246 var users []*db.User
247 rs, err := me.Db.Queryx(
248 `SELECT u.id, u.name, u.created_at
249 FROM app_users u
250 INNER JOIN posts ON u.id=posts.user_id
251 WHERE cur_space='feeds'
252 GROUP BY u.id, u.name, u.created_at
253 ORDER BY name ASC`,
254 )
255 if err != nil {
256 return users, err
257 }
258 defer func() { _ = rs.Close() }()
259 for rs.Next() {
260 var name sql.NullString
261 user := &db.User{}
262 err := rs.Scan(
263 &user.ID,
264 &name,
265 &user.CreatedAt,
266 )
267 if err != nil {
268 return users, err
269 }
270 user.Name = name.String
271
272 users = append(users, user)
273 }
274 if rs.Err() != nil {
275 return users, rs.Err()
276 }
277 return users, nil
278}
279
280func (me *PsqlDB) FindUserByKey(username string, key string) (*db.User, error) {
281 me.Logger.Info("attempting to find user with only public key", "key", key)
282 pk, err := me.findPublicKeyByKey(key)
283 if err == nil {
284 me.Logger.Info("found pubkey, looking for user", "key", key, "userId", pk.UserID)
285 user, err := me.FindUser(pk.UserID)
286 if err != nil {
287 return nil, err
288 }
289 user.PublicKey = pk
290 return user, nil
291 }
292
293 if errors.Is(err, &db.ErrMultiplePublicKeys{}) {
294 me.Logger.Info("detected multiple users with same public key", "user", username)
295 user, err := me.findUserForNameAndKey(username, key)
296 if err != nil {
297 me.Logger.Info("could not find user by username and public key", "user", username, "key", key)
298 // this is a little hacky but if we cannot find a user by name and public key
299 // then we return the multiple keys detected error so the user knows to specify their
300 // when logging in
301 return nil, &db.ErrMultiplePublicKeys{}
302 }
303 return user, nil
304 }
305
306 return nil, err
307}
308
309func (me *PsqlDB) FindUserByPubkey(key string) (*db.User, error) {
310 me.Logger.Info("attempting to find user with only public key", "key", key)
311 pk, err := me.findPublicKeyByKey(key)
312 if err != nil {
313 return nil, err
314 }
315
316 me.Logger.Info("found pubkey, looking for user", "key", key, "userId", pk.UserID)
317 user, err := me.FindUser(pk.UserID)
318 if err != nil {
319 return nil, err
320 }
321 user.PublicKey = pk
322 return user, nil
323}
324
325func (me *PsqlDB) FindUser(userID string) (*db.User, error) {
326 user := &db.User{}
327 err := me.Db.Get(user, `SELECT id, COALESCE(name, '') as name, created_at FROM app_users WHERE id = $1`, userID)
328 if err != nil {
329 return nil, err
330 }
331 return user, nil
332}
333
334func (me *PsqlDB) validateName(name string) (bool, error) {
335 lower := strings.ToLower(name)
336 if slices.Contains(db.DenyList, lower) {
337 return false, fmt.Errorf("%s is on deny list: %w", lower, db.ErrNameDenied)
338 }
339 v := db.NameValidator.MatchString(lower)
340 if !v {
341 return false, fmt.Errorf("%s is invalid: %w", lower, db.ErrNameInvalid)
342 }
343 user, _ := me.FindUserByName(lower)
344 if user == nil {
345 return true, nil
346 }
347 return false, fmt.Errorf("%s already taken: %w", lower, db.ErrNameTaken)
348}
349
350func (me *PsqlDB) FindUserByName(name string) (*db.User, error) {
351 user := &db.User{}
352 err := me.Db.Get(user, `SELECT * FROM app_users WHERE name = $1`, strings.ToLower(name))
353 if err != nil {
354 return nil, err
355 }
356 return user, nil
357}
358
359func (me *PsqlDB) findUserForNameAndKey(name string, key string) (*db.User, error) {
360 user := &db.User{}
361 pk := &db.PublicKey{}
362
363 r := me.Db.QueryRow(`SELECT app_users.id, app_users.name, app_users.created_at, public_keys.id as pk_id, public_keys.public_key, public_keys.created_at as pk_created_at FROM app_users LEFT JOIN public_keys ON public_keys.user_id = app_users.id WHERE app_users.name = $1 AND public_keys.public_key = $2`, strings.ToLower(name), key)
364 err := r.Scan(&user.ID, &user.Name, &user.CreatedAt, &pk.ID, &pk.Key, &pk.CreatedAt)
365 if err != nil {
366 return nil, err
367 }
368
369 user.PublicKey = pk
370 return user, nil
371}
372
373func (me *PsqlDB) FindUserByToken(token string) (*db.User, error) {
374 user := &db.User{}
375 err := me.Db.Get(user, `
376 SELECT app_users.id, app_users.name, app_users.created_at
377 FROM app_users
378 LEFT JOIN tokens ON tokens.user_id = app_users.id
379 WHERE tokens.token = $1 AND tokens.expires_at > NOW()`, token)
380 if err != nil {
381 return nil, err
382 }
383 return user, nil
384}
385
386func (me *PsqlDB) FindPostWithFilename(filename string, persona_id string, space string) (*db.Post, error) {
387 query := fmt.Sprintf(`
388 SELECT %s, STRING_AGG(coalesce(post_tags.name, ''), ',') tags
389 FROM posts
390 LEFT JOIN app_users ON app_users.id = posts.user_id
391 LEFT JOIN post_tags ON post_tags.post_id = posts.id
392 WHERE filename = $1 AND user_id = $2 AND cur_space = $3
393 GROUP BY %s`, SelectPost, SelectPost)
394 r := me.Db.QueryRow(query, filename, persona_id, space)
395 post, err := CreatePostWithTagsByRow(r)
396 if err != nil {
397 return nil, err
398 }
399
400 return post, nil
401}
402
403func (me *PsqlDB) FindPostWithSlug(slug string, user_id string, space string) (*db.Post, error) {
404 query := fmt.Sprintf(`
405 SELECT %s, STRING_AGG(coalesce(post_tags.name, ''), ',') tags
406 FROM posts
407 LEFT JOIN app_users ON app_users.id = posts.user_id
408 LEFT JOIN post_tags ON post_tags.post_id = posts.id
409 WHERE slug = $1 AND user_id = $2 AND cur_space = $3
410 GROUP BY %s`, SelectPost, SelectPost)
411 r := me.Db.QueryRow(query, slug, user_id, space)
412 post, err := CreatePostWithTagsByRow(r)
413 if err != nil {
414 // attempt to find post inside post_aliases
415 alias := me.Db.QueryRow(
416 `SELECT post_aliases.post_id FROM post_aliases
417 INNER JOIN posts ON posts.id = post_aliases.post_id
418 WHERE post_aliases.slug = $1 AND posts.user_id = $2`,
419 slug, user_id,
420 )
421 postID := ""
422 err := alias.Scan(&postID)
423 if err != nil {
424 return nil, err
425 }
426
427 return me.FindPost(postID)
428 }
429
430 return post, nil
431}
432
433func (me *PsqlDB) FindPost(postID string) (*db.Post, error) {
434 post := &db.Post{}
435 query := fmt.Sprintf(`
436 SELECT %s
437 FROM posts
438 LEFT JOIN app_users ON app_users.id = posts.user_id
439 WHERE posts.id = $1`, SelectPost)
440 err := me.Db.Get(post, query, postID)
441 if err != nil {
442 return nil, err
443 }
444 return post, nil
445}
446
447func (me *PsqlDB) postPager(rs *sqlx.Rows, pageNum int, space string, tag string) (*db.Paginate[*db.Post], error) {
448 var posts []*db.Post
449 for rs.Next() {
450 post := &db.Post{}
451 err := rs.Scan(
452 &post.ID,
453 &post.UserID,
454 &post.Filename,
455 &post.Slug,
456 &post.Title,
457 &post.Text,
458 &post.Description,
459 &post.PublishAt,
460 &post.Username,
461 &post.UpdatedAt,
462 &post.MimeType,
463 )
464 if err != nil {
465 return nil, err
466 }
467
468 posts = append(posts, post)
469 }
470 if rs.Err() != nil {
471 return nil, rs.Err()
472 }
473
474 var count int
475 var err error
476 if tag == "" {
477 err = me.Db.QueryRow(`SELECT count(id) FROM posts WHERE hidden = FALSE AND cur_space=$1`, space).Scan(&count)
478 } else {
479 err = me.Db.QueryRow(`
480 SELECT count(posts.id)
481 FROM posts
482 LEFT JOIN post_tags ON post_tags.post_id = posts.id
483 WHERE hidden = FALSE AND cur_space=$1 and post_tags.name = $2`, space, tag).Scan(&count)
484 }
485 if err != nil {
486 return nil, err
487 }
488
489 pager := &db.Paginate[*db.Post]{
490 Data: posts,
491 Total: int(math.Ceil(float64(count) / float64(pageNum))),
492 }
493
494 return pager, nil
495}
496
497func (me *PsqlDB) FindPostsByFeed(page *db.Pager, space string) (*db.Paginate[*db.Post], error) {
498 query := `
499 SELECT *
500 FROM (
501 SELECT DISTINCT ON (posts.user_id)
502 posts.id,
503 posts.user_id,
504 posts.filename,
505 posts.slug,
506 posts.title,
507 posts.text,
508 posts.description,
509 posts.publish_at,
510 app_users.name AS username,
511 posts.updated_at,
512 posts.mime_type
513 FROM posts
514 LEFT JOIN app_users ON app_users.id = posts.user_id
515 WHERE
516 hidden = FALSE
517 AND publish_at::date <= CURRENT_DATE
518 AND cur_space = $3
519 ORDER BY posts.user_id, publish_at DESC
520 ) AS latest_posts
521 ORDER BY publish_at DESC
522 LIMIT $1 OFFSET $2`
523 rs, err := me.Db.Queryx(query, page.Num, page.Num*page.Page, space)
524 if err != nil {
525 return nil, err
526 }
527 defer func() { _ = rs.Close() }()
528 return me.postPager(rs, page.Num, space, "")
529}
530
531func (me *PsqlDB) InsertPost(post *db.Post) (*db.Post, error) {
532 var id string
533 query := `
534 INSERT INTO posts
535 (user_id, filename, slug, title, text, description, publish_at, hidden, cur_space,
536 file_size, mime_type, shasum, data, expires_at, updated_at)
537 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15)
538 RETURNING id`
539 err := me.Db.QueryRow(
540 query,
541 post.UserID,
542 post.Filename,
543 post.Slug,
544 post.Title,
545 post.Text,
546 post.Description,
547 post.PublishAt,
548 post.Hidden,
549 post.Space,
550 post.FileSize,
551 post.MimeType,
552 post.Shasum,
553 post.Data,
554 post.ExpiresAt,
555 post.UpdatedAt,
556 ).Scan(&id)
557 if err != nil {
558 return nil, err
559 }
560
561 return me.FindPost(id)
562}
563
564func (me *PsqlDB) UpdatePost(post *db.Post) (*db.Post, error) {
565 query := `
566 UPDATE posts
567 SET slug = $1, title = $2, text = $3, description = $4, updated_at = $5, publish_at = $6,
568 file_size = $7, shasum = $8, data = $9, hidden = $11, expires_at = $12
569 WHERE id = $10`
570 _, err := me.Db.Exec(
571 query,
572 post.Slug,
573 post.Title,
574 post.Text,
575 post.Description,
576 post.UpdatedAt,
577 post.PublishAt,
578 post.FileSize,
579 post.Shasum,
580 post.Data,
581 post.ID,
582 post.Hidden,
583 post.ExpiresAt,
584 )
585 if err != nil {
586 return nil, err
587 }
588
589 return me.FindPost(post.ID)
590}
591
592func (me *PsqlDB) RemovePosts(postIDs []string) error {
593 param := "{" + strings.Join(postIDs, ",") + "}"
594 _, err := me.Db.Exec(`DELETE FROM posts WHERE id = ANY($1::uuid[])`, param)
595 return err
596}
597
598func (me *PsqlDB) FindPostsByUser(page *db.Pager, userID string, space string) (*db.Paginate[*db.Post], error) {
599 var posts []*db.Post
600 query := fmt.Sprintf(`
601 SELECT %s, STRING_AGG(coalesce(post_tags.name, ''), ',') tags
602 FROM posts
603 LEFT JOIN app_users ON app_users.id = posts.user_id
604 LEFT JOIN post_tags ON post_tags.post_id = posts.id
605 WHERE
606 hidden = FALSE AND
607 user_id = $1 AND
608 publish_at::date <= CURRENT_DATE AND
609 cur_space = $2
610 GROUP BY %s
611 ORDER BY publish_at DESC, slug DESC
612 LIMIT $3 OFFSET $4`, SelectPost, SelectPost)
613 rs, err := me.Db.Queryx(
614 query,
615 userID,
616 space,
617 page.Num,
618 page.Num*page.Page,
619 )
620 if err != nil {
621 return nil, err
622 }
623 defer func() { _ = rs.Close() }()
624 for rs.Next() {
625 post, err := CreatePostWithTagsByRow(rs)
626 if err != nil {
627 return nil, err
628 }
629
630 posts = append(posts, post)
631 }
632
633 if rs.Err() != nil {
634 return nil, rs.Err()
635 }
636
637 var count int
638 err = me.Db.QueryRow(`SELECT count(id) FROM posts WHERE hidden = FALSE AND cur_space=$1`, space).Scan(&count)
639 if err != nil {
640 return nil, err
641 }
642
643 pager := &db.Paginate[*db.Post]{
644 Data: posts,
645 Total: int(math.Ceil(float64(count) / float64(page.Num))),
646 }
647 return pager, nil
648}
649
650func (me *PsqlDB) FindAllPostsByUser(userID string, space string) ([]*db.Post, error) {
651 var posts []*db.Post
652 query := fmt.Sprintf(`
653 SELECT %s
654 FROM posts
655 LEFT JOIN app_users ON app_users.id = posts.user_id
656 WHERE
657 user_id = $1 AND
658 cur_space = $2
659 ORDER BY publish_at DESC`, SelectPost)
660 err := me.Db.Select(&posts, query, userID, space)
661 if err != nil {
662 return nil, err
663 }
664 return posts, nil
665}
666
667func (me *PsqlDB) FindPosts() ([]*db.Post, error) {
668 var posts []*db.Post
669 query := fmt.Sprintf(`
670 SELECT %s
671 FROM posts
672 LEFT JOIN app_users ON app_users.id = posts.user_id`, SelectPost)
673 err := me.Db.Select(&posts, query)
674 if err != nil {
675 return nil, err
676 }
677 return posts, nil
678}
679
680func (me *PsqlDB) FindExpiredPosts(space string) ([]*db.Post, error) {
681 var posts []*db.Post
682 query := fmt.Sprintf(`
683 SELECT %s
684 FROM posts
685 LEFT JOIN app_users ON app_users.id = posts.user_id
686 WHERE
687 cur_space = $1 AND
688 expires_at <= now();
689 `, SelectPost)
690 err := me.Db.Select(&posts, query, space)
691 if err != nil {
692 return nil, err
693 }
694 return posts, nil
695}
696
697func (me *PsqlDB) Close() error {
698 me.Logger.Info("Closing db")
699 return me.Db.Close()
700}
701
702func newNullString(s string) sql.NullString {
703 if len(s) == 0 {
704 return sql.NullString{}
705 }
706 return sql.NullString{
707 String: s,
708 Valid: true,
709 }
710}
711
712func (me *PsqlDB) InsertVisit(visit *db.AnalyticsVisits) error {
713 _, err := me.Db.Exec(
714 `INSERT INTO analytics_visits (user_id, project_id, post_id, namespace, host, path, ip_address, user_agent, referer, status, content_type) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11);`,
715 visit.UserID,
716 newNullString(visit.ProjectID),
717 newNullString(visit.PostID),
718 newNullString(visit.Namespace),
719 visit.Host,
720 visit.Path,
721 visit.IpAddress,
722 visit.UserAgent,
723 visit.Referer,
724 visit.Status,
725 visit.ContentType,
726 )
727 return err
728}
729
730func visitFilterBy(opts *db.SummaryOpts) (string, string) {
731 where := ""
732 val := ""
733 if opts.Host != "" {
734 where = "host"
735 val = opts.Host
736 } else if opts.Path != "" {
737 where = "path"
738 val = opts.Path
739 }
740
741 return where, val
742}
743
744func (me *PsqlDB) visitUnique(opts *db.SummaryOpts) ([]*db.VisitInterval, error) {
745 var intervals, currentIntervals []*db.VisitInterval
746 var sumErr, rawErr error
747
748 var wg sync.WaitGroup
749 wg.Add(2)
750
751 go func() {
752 defer wg.Done()
753 intervals, sumErr = me.visitUniqueFromSummary(opts)
754 }()
755 go func() {
756 defer wg.Done()
757 currentIntervals, rawErr = me.visitUniqueFromRaw(opts)
758 }()
759
760 wg.Wait()
761
762 if sumErr != nil {
763 return nil, fmt.Errorf("query summary visits: %w", sumErr)
764 }
765 if rawErr != nil {
766 return nil, fmt.Errorf("query raw visits: %w", rawErr)
767 }
768
769 // Merge: current month data may overlap with summary data, combine counts
770 return mergeVisitIntervals(intervals, currentIntervals), nil
771}
772
773// visitUniqueFromSummary reads unique visitor counts from analytics_monthly_visits for historical data.
774func (me *PsqlDB) visitUniqueFromSummary(opts *db.SummaryOpts) ([]*db.VisitInterval, error) {
775 now := time.Now()
776 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
777 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
778
779 // If origin is in the previous month or later, raw data covers it — no summary to fetch.
780 if !opts.Origin.Before(previousMonthStart) {
781 return nil, nil
782 }
783
784 where := ""
785 args := []interface{}{opts.UserID, opts.Origin, currentMonthStart}
786 argIdx := 4
787 if opts.Host != "" {
788 where = "AND host = $" + fmt.Sprintf("%d", argIdx)
789 args = append(args, opts.Host)
790 }
791
792 query := fmt.Sprintf(`
793 SELECT
794 date_trunc('%s', visit_date)::timestamptz as interval_start,
795 sum(unique_visits) as unique_visitors,
796 sum(mobile_visits) as mobile_visits,
797 sum(desktop_visits) as desktop_visits
798 FROM analytics_monthly_visits
799 WHERE user_id = $1 AND visit_date >= $2 AND visit_date < $3 %s
800 GROUP BY interval_start
801 ORDER BY interval_start`, opts.Interval, where)
802
803 rows, err := me.Db.Queryx(query, args...)
804 if err != nil {
805 return nil, err
806 }
807 defer func() { _ = rows.Close() }()
808
809 var intervals []*db.VisitInterval
810 for rows.Next() {
811 interval := &db.VisitInterval{}
812 if err := rows.Scan(&interval.Interval, &interval.Visitors, &interval.MobileVisitors, &interval.DesktopVisitors); err != nil {
813 return nil, err
814 }
815 intervals = append(intervals, interval)
816 }
817 return intervals, rows.Err()
818}
819
820// visitUniqueFromRaw reads unique visitor counts from analytics_visits for the previous and current months.
821// This covers the gap between the last aggregated month and the current month.
822func (me *PsqlDB) visitUniqueFromRaw(opts *db.SummaryOpts) ([]*db.VisitInterval, error) {
823 now := time.Now()
824 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
825 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
826
827 where, with := visitFilterBy(opts)
828
829 // Determine the effective start: max(origin, previousMonthStart)
830 effectiveStart := previousMonthStart
831 if opts.Origin.After(previousMonthStart) {
832 effectiveStart = opts.Origin
833 }
834
835 uniqueVisitors := fmt.Sprintf(`
836 SELECT
837 date_trunc('%s', created_at)::timestamptz as interval_start,
838 count(DISTINCT CASE WHEN %s THEN ip_address END) as mobile_visitors,
839 count(DISTINCT CASE WHEN NOT %s THEN ip_address END) as desktop_visitors,
840 count(DISTINCT ip_address) as unique_visitors
841 FROM analytics_visits
842 WHERE created_at >= $1 AND created_at < $2 AND %s = $3 AND user_id = $4 AND status <> 404
843 GROUP BY interval_start
844 ORDER BY interval_start`,
845 opts.Interval, mobileUserAgentExpr, mobileUserAgentExpr, where)
846
847 rows, err := me.Db.Queryx(uniqueVisitors, effectiveStart, currentMonthStart.AddDate(0, 1, 0), with, opts.UserID)
848 if err != nil {
849 return nil, err
850 }
851 defer func() { _ = rows.Close() }()
852
853 var intervals []*db.VisitInterval
854 for rows.Next() {
855 interval := &db.VisitInterval{}
856 if err := rows.Scan(&interval.Interval, &interval.MobileVisitors, &interval.DesktopVisitors, &interval.Visitors); err != nil {
857 return nil, err
858 }
859 intervals = append(intervals, interval)
860 }
861 return intervals, rows.Err()
862}
863
864// mergeVisitIntervals combines historical (summary table) and current (raw) intervals.
865// Summary data is preferred when both sources have the same interval to avoid double-counting.
866func mergeVisitIntervals(historical, current []*db.VisitInterval) []*db.VisitInterval {
867 if len(historical) == 0 {
868 return current
869 }
870 if len(current) == 0 {
871 return historical
872 }
873
874 // Build a map by interval timestamp for merging.
875 // Summary data takes precedence over raw data for the same interval
876 // (e.g. when a month is aggregated but raw data still exists in a local dump).
877 intervalMap := make(map[int64]*db.VisitInterval)
878 for _, ci := range current {
879 ts := ci.Interval.Unix()
880 intervalMap[ts] = ci
881 }
882 for _, hi := range historical {
883 ts := hi.Interval.Unix()
884 intervalMap[ts] = hi // summary overwrites raw if both exist
885 }
886
887 // Sort by interval
888 result := make([]*db.VisitInterval, 0, len(intervalMap))
889 for _, iv := range intervalMap {
890 result = append(result, iv)
891 }
892 sort.Slice(result, func(i, j int) bool {
893 return result[i].Interval.Before(*result[j].Interval)
894 })
895 return result
896}
897
898// mergeTopUrls combines historical and current top URLs, summing counts for overlapping URLs,
899// then returns the top 10 by total count.
900func mergeTopUrls(historical, current []*db.VisitUrl) []*db.VisitUrl {
901 if len(historical) == 0 {
902 return current
903 }
904 if len(current) == 0 {
905 return historical
906 }
907
908 // Build a map by URL for merging
909 urlMap := make(map[string]*db.VisitUrl)
910 for _, hu := range historical {
911 urlMap[hu.Url] = hu
912 }
913
914 for _, cu := range current {
915 if existing, ok := urlMap[cu.Url]; ok {
916 existing.Count += cu.Count
917 } else {
918 urlMap[cu.Url] = cu
919 }
920 }
921
922 // Sort by count descending and take top 10
923 result := make([]*db.VisitUrl, 0, len(urlMap))
924 for _, u := range urlMap {
925 result = append(result, u)
926 }
927 sort.Slice(result, func(i, j int) bool {
928 return result[i].Count > result[j].Count
929 })
930 if len(result) > 10 {
931 result = result[:10]
932 }
933 return result
934}
935
936// mergeTopReferers combines historical and current top referers, summing counts for overlapping referers,
937// then returns the top 10 by total count.
938func mergeTopReferers(historical, current []*db.VisitUrl) []*db.VisitUrl {
939 return mergeTopUrls(historical, current) // Same logic as mergeTopUrls
940}
941
942// mergeHosts combines historical and current hosts, summing counts for overlapping hosts,
943// then returns sorted by total count descending.
944func mergeHosts(historical, current []*db.VisitUrl) []*db.VisitUrl {
945 if len(historical) == 0 {
946 return current
947 }
948 if len(current) == 0 {
949 return historical
950 }
951
952 // Build a map by host for merging
953 hostMap := make(map[string]*db.VisitUrl)
954 for _, hu := range historical {
955 hostMap[hu.Url] = hu
956 }
957
958 for _, cu := range current {
959 if existing, ok := hostMap[cu.Url]; ok {
960 existing.Count += cu.Count
961 } else {
962 hostMap[cu.Url] = cu
963 }
964 }
965
966 // Sort by count descending
967 result := make([]*db.VisitUrl, 0, len(hostMap))
968 for _, h := range hostMap {
969 result = append(result, h)
970 }
971 sort.Slice(result, func(i, j int) bool {
972 return result[i].Count > result[j].Count
973 })
974 return result
975}
976
977func (me *PsqlDB) visitReferer(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
978 var historical, current []*db.VisitUrl
979 var histErr, rawErr error
980
981 var wg sync.WaitGroup
982 wg.Add(2)
983
984 go func() {
985 defer wg.Done()
986 historical, histErr = me.visitRefererFromSummary(opts)
987 }()
988 go func() {
989 defer wg.Done()
990 current, rawErr = me.visitRefererFromRaw(opts)
991 }()
992
993 wg.Wait()
994
995 if histErr != nil {
996 return nil, fmt.Errorf("query summary referers: %w", histErr)
997 }
998 if rawErr != nil {
999 return nil, fmt.Errorf("query raw referers: %w", rawErr)
1000 }
1001
1002 return mergeTopReferers(historical, current), nil
1003}
1004
1005// visitRefererFromSummary reads top referers from analytics_monthly_top_referers for historical data.
1006func (me *PsqlDB) visitRefererFromSummary(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1007 now := time.Now()
1008 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1009 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
1010
1011 // If origin is in the previous month or later, raw data covers it — no summary to fetch.
1012 if !opts.Origin.Before(previousMonthStart) {
1013 return nil, nil
1014 }
1015
1016 // Clamp origin to month boundary for summary table lookup
1017 originMonthStart := time.Date(opts.Origin.Year(), opts.Origin.Month(), 1, 0, 0, 0, 0, time.UTC)
1018
1019 where := ""
1020 args := []interface{}{opts.UserID, originMonthStart, currentMonthStart}
1021 if opts.Host != "" {
1022 where = "AND host = $4"
1023 args = append(args, opts.Host)
1024 }
1025
1026 query := fmt.Sprintf(`
1027 SELECT referer, sum(unique_visits) as total_visits
1028 FROM analytics_monthly_top_referers
1029 WHERE user_id = $1 AND month >= $2 AND month < $3 %s
1030 GROUP BY referer
1031 ORDER BY total_visits DESC
1032 LIMIT 10`, where)
1033
1034 rows, err := me.Db.Queryx(query, args...)
1035 if err != nil {
1036 return nil, err
1037 }
1038 defer func() { _ = rows.Close() }()
1039
1040 var results []*db.VisitUrl
1041 for rows.Next() {
1042 result := &db.VisitUrl{}
1043 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1044 return nil, err
1045 }
1046 results = append(results, result)
1047 }
1048 return results, rows.Err()
1049}
1050
1051// visitRefererFromRaw reads top referers from analytics_visits for the previous and current months.
1052func (me *PsqlDB) visitRefererFromRaw(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1053 now := time.Now()
1054 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1055 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
1056
1057 where, with := visitFilterBy(opts)
1058
1059 // Determine the effective start: max(origin, previousMonthStart)
1060 effectiveStart := previousMonthStart
1061 if opts.Origin.After(previousMonthStart) {
1062 effectiveStart = opts.Origin
1063 }
1064
1065 topUrls := fmt.Sprintf(`
1066 SELECT
1067 referer,
1068 count(DISTINCT ip_address) as referer_count
1069 FROM analytics_visits
1070 WHERE created_at >= $1 AND created_at < $2 AND %s = $3 AND user_id = $4 AND referer <> '' AND status <> 404
1071 GROUP BY referer
1072 ORDER BY referer_count DESC
1073 LIMIT 10`, where)
1074
1075 rows, err := me.Db.Queryx(topUrls, effectiveStart, currentMonthStart.AddDate(0, 1, 0), with, opts.UserID)
1076 if err != nil {
1077 return nil, err
1078 }
1079 defer func() { _ = rows.Close() }()
1080
1081 var results []*db.VisitUrl
1082 for rows.Next() {
1083 result := &db.VisitUrl{}
1084 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1085 return nil, err
1086 }
1087 results = append(results, result)
1088 }
1089 return results, rows.Err()
1090}
1091
1092func (me *PsqlDB) visitUrl(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1093 var historical, current []*db.VisitUrl
1094 var histErr, rawErr error
1095
1096 var wg sync.WaitGroup
1097 wg.Add(2)
1098
1099 go func() {
1100 defer wg.Done()
1101 historical, histErr = me.visitUrlFromSummary(opts)
1102 }()
1103 go func() {
1104 defer wg.Done()
1105 current, rawErr = me.visitUrlFromRaw(opts)
1106 }()
1107
1108 wg.Wait()
1109
1110 if histErr != nil {
1111 return nil, fmt.Errorf("query summary urls: %w", histErr)
1112 }
1113 if rawErr != nil {
1114 return nil, fmt.Errorf("query raw urls: %w", rawErr)
1115 }
1116
1117 return mergeTopUrls(historical, current), nil
1118}
1119
1120// visitUrlFromSummary reads top URLs from analytics_monthly_top_urls for historical data.
1121func (me *PsqlDB) visitUrlFromSummary(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1122 now := time.Now()
1123 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1124 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
1125
1126 // If origin is in the previous month or later, raw data covers it — no summary to fetch.
1127 if !opts.Origin.Before(previousMonthStart) {
1128 return nil, nil
1129 }
1130
1131 // Clamp origin to month boundary for summary table lookup
1132 originMonthStart := time.Date(opts.Origin.Year(), opts.Origin.Month(), 1, 0, 0, 0, 0, time.UTC)
1133
1134 where := ""
1135 args := []interface{}{opts.UserID, originMonthStart, currentMonthStart}
1136 if opts.Host != "" {
1137 where = "AND host = $4"
1138 args = append(args, opts.Host)
1139 }
1140
1141 query := fmt.Sprintf(`
1142 SELECT path, sum(unique_visits) as total_visits
1143 FROM analytics_monthly_top_urls
1144 WHERE user_id = $1 AND month >= $2 AND month < $3 AND status_code <> 404 %s
1145 GROUP BY path
1146 ORDER BY total_visits DESC
1147 LIMIT 10`, where)
1148
1149 rows, err := me.Db.Queryx(query, args...)
1150 if err != nil {
1151 return nil, err
1152 }
1153 defer func() { _ = rows.Close() }()
1154
1155 var results []*db.VisitUrl
1156 for rows.Next() {
1157 result := &db.VisitUrl{}
1158 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1159 return nil, err
1160 }
1161 results = append(results, result)
1162 }
1163 return results, rows.Err()
1164}
1165
1166// visitUrlFromRaw reads top URLs from analytics_visits for the previous and current months.
1167func (me *PsqlDB) visitUrlFromRaw(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1168 now := time.Now()
1169 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1170 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
1171
1172 where, with := visitFilterBy(opts)
1173
1174 // Determine the effective start: max(origin, previousMonthStart)
1175 effectiveStart := previousMonthStart
1176 if opts.Origin.After(previousMonthStart) {
1177 effectiveStart = opts.Origin
1178 }
1179
1180 topUrls := fmt.Sprintf(`
1181 SELECT
1182 path,
1183 count(DISTINCT ip_address) as path_count
1184 FROM analytics_visits
1185 WHERE created_at >= $1 AND created_at < $2 AND %s = $3 AND user_id = $4 AND path <> '' AND status <> 404
1186 GROUP BY path
1187 ORDER BY path_count DESC
1188 LIMIT 10`, where)
1189
1190 rows, err := me.Db.Queryx(topUrls, effectiveStart, currentMonthStart.AddDate(0, 1, 0), with, opts.UserID)
1191 if err != nil {
1192 return nil, err
1193 }
1194 defer func() { _ = rows.Close() }()
1195
1196 var results []*db.VisitUrl
1197 for rows.Next() {
1198 result := &db.VisitUrl{}
1199 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1200 return nil, err
1201 }
1202 results = append(results, result)
1203 }
1204 return results, rows.Err()
1205}
1206
1207func (me *PsqlDB) VisitUrlNotFound(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1208 limit := opts.Limit
1209 if limit == 0 {
1210 limit = 10
1211 }
1212
1213 var historical, current []*db.VisitUrl
1214 var histErr, rawErr error
1215
1216 var wg sync.WaitGroup
1217 wg.Add(2)
1218
1219 go func() {
1220 defer wg.Done()
1221 historical, histErr = me.visitUrlNotFoundFromSummary(opts, limit)
1222 }()
1223 go func() {
1224 defer wg.Done()
1225 current, rawErr = me.visitUrlNotFoundFromRaw(opts, limit)
1226 }()
1227
1228 wg.Wait()
1229
1230 if histErr != nil {
1231 return nil, fmt.Errorf("query summary 404 urls: %w", histErr)
1232 }
1233 if rawErr != nil {
1234 return nil, fmt.Errorf("query raw 404 urls: %w", rawErr)
1235 }
1236
1237 return mergeTopUrls(historical, current), nil
1238}
1239
1240// visitUrlNotFoundFromSummary reads top 404 URLs from analytics_monthly_top_urls for historical data.
1241func (me *PsqlDB) visitUrlNotFoundFromSummary(opts *db.SummaryOpts, limit int) ([]*db.VisitUrl, error) {
1242 now := time.Now()
1243 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1244 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
1245
1246 // If origin is in the previous month or later, raw data covers it — no summary to fetch.
1247 if !opts.Origin.Before(previousMonthStart) {
1248 return nil, nil
1249 }
1250
1251 // Clamp origin to month boundary for summary table lookup
1252 originMonthStart := time.Date(opts.Origin.Year(), opts.Origin.Month(), 1, 0, 0, 0, 0, time.UTC)
1253
1254 where := ""
1255 args := []interface{}{opts.UserID, originMonthStart, currentMonthStart}
1256 argIdx := 4
1257 if opts.Host != "" {
1258 where = "AND host = $" + fmt.Sprintf("%d", argIdx)
1259 args = append(args, opts.Host)
1260 }
1261
1262 query := fmt.Sprintf(`
1263 SELECT path, sum(unique_visits) as total_visits
1264 FROM analytics_monthly_top_urls
1265 WHERE user_id = $1 AND month >= $2 AND month < $3 AND status_code = 404 %s
1266 GROUP BY path
1267 ORDER BY total_visits DESC
1268 LIMIT %d`, where, limit)
1269
1270 rows, err := me.Db.Queryx(query, args...)
1271 if err != nil {
1272 return nil, err
1273 }
1274 defer func() { _ = rows.Close() }()
1275
1276 var results []*db.VisitUrl
1277 for rows.Next() {
1278 result := &db.VisitUrl{}
1279 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1280 return nil, err
1281 }
1282 results = append(results, result)
1283 }
1284 return results, rows.Err()
1285}
1286
1287// visitUrlNotFoundFromRaw reads top 404 URLs from analytics_visits for the previous and current months.
1288func (me *PsqlDB) visitUrlNotFoundFromRaw(opts *db.SummaryOpts, limit int) ([]*db.VisitUrl, error) {
1289 now := time.Now()
1290 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1291 previousMonthStart := currentMonthStart.AddDate(0, -1, 0)
1292
1293 where, with := visitFilterBy(opts)
1294
1295 // Determine the effective start: max(origin, previousMonthStart)
1296 effectiveStart := previousMonthStart
1297 if opts.Origin.After(previousMonthStart) {
1298 effectiveStart = opts.Origin
1299 }
1300
1301 topUrls := fmt.Sprintf(`
1302 SELECT
1303 path,
1304 count(DISTINCT ip_address) as path_count
1305 FROM analytics_visits
1306 WHERE created_at >= $1 AND created_at < $2 AND %s = $3 AND user_id = $4 AND path <> '' AND status = 404
1307 GROUP BY path
1308 ORDER BY path_count DESC
1309 LIMIT %d`, where, limit)
1310
1311 rows, err := me.Db.Queryx(topUrls, effectiveStart, currentMonthStart.AddDate(0, 1, 0), with, opts.UserID)
1312 if err != nil {
1313 return nil, err
1314 }
1315 defer func() { _ = rows.Close() }()
1316
1317 var results []*db.VisitUrl
1318 for rows.Next() {
1319 result := &db.VisitUrl{}
1320 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1321 return nil, err
1322 }
1323 results = append(results, result)
1324 }
1325 return results, rows.Err()
1326}
1327
1328func (me *PsqlDB) visitHost(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1329 var historical, current []*db.VisitUrl
1330 var histErr, rawErr error
1331
1332 var wg sync.WaitGroup
1333 wg.Add(2)
1334
1335 go func() {
1336 defer wg.Done()
1337 historical, histErr = me.visitHostFromSummary(opts)
1338 }()
1339 go func() {
1340 defer wg.Done()
1341 current, rawErr = me.visitHostFromRaw(opts)
1342 }()
1343
1344 wg.Wait()
1345
1346 if histErr != nil {
1347 return nil, fmt.Errorf("query summary hosts: %w", histErr)
1348 }
1349 if rawErr != nil {
1350 return nil, fmt.Errorf("query raw hosts: %w", rawErr)
1351 }
1352
1353 return mergeHosts(historical, current), nil
1354}
1355
1356// visitHostFromSummary reads host data from analytics_user_sites for historical data.
1357func (me *PsqlDB) visitHostFromSummary(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1358 rows, err := me.Db.Queryx(`
1359 SELECT host, total_visits
1360 FROM analytics_user_sites
1361 WHERE user_id = $1 AND host <> ''
1362 ORDER BY total_visits DESC`, opts.UserID)
1363 if err != nil {
1364 return nil, err
1365 }
1366 defer func() { _ = rows.Close() }()
1367
1368 var results []*db.VisitUrl
1369 for rows.Next() {
1370 result := &db.VisitUrl{}
1371 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1372 return nil, err
1373 }
1374 results = append(results, result)
1375 }
1376 return results, rows.Err()
1377}
1378
1379// visitHostFromRaw reads hosts from analytics_visits for the current month that aren't in summary.
1380func (me *PsqlDB) visitHostFromRaw(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1381 now := time.Now()
1382 currentMonthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC)
1383
1384 rows, err := me.Db.Queryx(`
1385 SELECT host, count(DISTINCT ip_address) as host_count
1386 FROM analytics_visits
1387 WHERE created_at >= $1 AND user_id = $2 AND host <> ''
1388 GROUP BY host
1389 ORDER BY host_count DESC`, currentMonthStart, opts.UserID)
1390 if err != nil {
1391 return nil, err
1392 }
1393 defer func() { _ = rows.Close() }()
1394
1395 var results []*db.VisitUrl
1396 for rows.Next() {
1397 result := &db.VisitUrl{}
1398 if err := rows.Scan(&result.Url, &result.Count); err != nil {
1399 return nil, err
1400 }
1401 results = append(results, result)
1402 }
1403 return results, rows.Err()
1404}
1405
1406func (me *PsqlDB) VisitSummary(opts *db.SummaryOpts) (*db.SummaryVisits, error) {
1407 var (
1408 visitors []*db.VisitInterval
1409 urls []*db.VisitUrl
1410 refs []*db.VisitUrl
1411 notFound []*db.VisitUrl
1412 visitorsErr error
1413 urlsErr error
1414 refsErr error
1415 nfErr error
1416 )
1417
1418 var wg sync.WaitGroup
1419 wg.Add(4)
1420
1421 go func() {
1422 defer wg.Done()
1423 visitors, visitorsErr = me.visitUnique(opts)
1424 }()
1425 go func() {
1426 defer wg.Done()
1427 urls, urlsErr = me.visitUrl(opts)
1428 }()
1429 go func() {
1430 defer wg.Done()
1431 refs, refsErr = me.visitReferer(opts)
1432 }()
1433 go func() {
1434 defer wg.Done()
1435 notFound, nfErr = me.VisitUrlNotFound(opts)
1436 }()
1437
1438 wg.Wait()
1439
1440 // Return the first error encountered
1441 for _, err := range []error{visitorsErr, urlsErr, refsErr, nfErr} {
1442 if err != nil {
1443 return nil, err
1444 }
1445 }
1446
1447 return &db.SummaryVisits{
1448 Intervals: visitors,
1449 TopUrls: urls,
1450 TopReferers: refs,
1451 NotFoundUrls: notFound,
1452 }, nil
1453}
1454
1455func (me *PsqlDB) FindVisitSiteList(opts *db.SummaryOpts) ([]*db.VisitUrl, error) {
1456 return me.visitHost(opts)
1457}
1458
1459func (me *PsqlDB) FindUsers() ([]*db.User, error) {
1460 var users []*db.User
1461 err := me.Db.Select(&users, `SELECT id, COALESCE(name, '') as name, created_at FROM app_users ORDER BY name ASC`)
1462 if err != nil {
1463 return nil, err
1464 }
1465 return users, nil
1466}
1467
1468func (me *PsqlDB) removeTagsByPost(tx *sqlx.Tx, postID string) error {
1469 _, err := tx.Exec(`DELETE FROM post_tags WHERE post_id = $1`, postID)
1470 return err
1471}
1472
1473func (me *PsqlDB) insertTagsByPost(tx *sqlx.Tx, tags []string, postID string) ([]string, error) {
1474 ids := make([]string, 0)
1475 for _, tag := range tags {
1476 id := ""
1477 err := tx.QueryRow(`INSERT INTO post_tags (post_id, name) VALUES($1, $2) RETURNING id;`, postID, tag).Scan(&id)
1478 if err != nil {
1479 return nil, err
1480 }
1481 ids = append(ids, id)
1482 }
1483
1484 return ids, nil
1485}
1486
1487func (me *PsqlDB) ReplaceTagsByPost(tags []string, postID string) error {
1488 tx, err := me.Db.Beginx()
1489 if err != nil {
1490 return err
1491 }
1492 defer func() {
1493 _ = tx.Rollback()
1494 }()
1495
1496 err = me.removeTagsByPost(tx, postID)
1497 if err != nil {
1498 return err
1499 }
1500
1501 _, err = me.insertTagsByPost(tx, tags, postID)
1502 if err != nil {
1503 return err
1504 }
1505
1506 err = tx.Commit()
1507 return err
1508}
1509
1510func (me *PsqlDB) removeAliasesByPost(tx *sqlx.Tx, postID string) error {
1511 _, err := tx.Exec(`DELETE FROM post_aliases WHERE post_id = $1`, postID)
1512 return err
1513}
1514
1515func (me *PsqlDB) insertAliasesByPost(tx *sqlx.Tx, aliases []string, postID string) ([]string, error) {
1516 // hardcoded
1517 denyList := []string{
1518 "rss",
1519 "rss.xml",
1520 "rss.atom",
1521 "atom.xml",
1522 "feed.xml",
1523 "smol.css",
1524 "main.css",
1525 "syntax.css",
1526 "card.png",
1527 "favicon-16x16.png",
1528 "favicon-32x32.png",
1529 "apple-touch-icon.png",
1530 "favicon.ico",
1531 "robots.txt",
1532 "atom",
1533 "blog/index.xml",
1534 }
1535
1536 ids := make([]string, 0)
1537 for _, alias := range aliases {
1538 if slices.Contains(denyList, alias) {
1539 me.Logger.Info(
1540 "name is in the deny list for aliases because it conflicts with a static route, skipping",
1541 "alias", alias,
1542 )
1543 continue
1544 }
1545 id := ""
1546 err := tx.QueryRow(`INSERT INTO post_aliases (post_id, slug) VALUES($1, $2) RETURNING id;`, postID, alias).Scan(&id)
1547 if err != nil {
1548 return nil, err
1549 }
1550 ids = append(ids, id)
1551 }
1552
1553 return ids, nil
1554}
1555
1556func (me *PsqlDB) ReplaceAliasesByPost(aliases []string, postID string) error {
1557 tx, err := me.Db.Beginx()
1558 if err != nil {
1559 return err
1560 }
1561 defer func() {
1562 _ = tx.Rollback()
1563 }()
1564
1565 err = me.removeAliasesByPost(tx, postID)
1566 if err != nil {
1567 return err
1568 }
1569
1570 _, err = me.insertAliasesByPost(tx, aliases, postID)
1571 if err != nil {
1572 return err
1573 }
1574
1575 err = tx.Commit()
1576 return err
1577}
1578
1579func (me *PsqlDB) FindUserPostsByTag(page *db.Pager, tag, userID, space string) (*db.Paginate[*db.Post], error) {
1580 var posts []*db.Post
1581 query := fmt.Sprintf(`
1582 SELECT %s
1583 FROM posts
1584 LEFT JOIN app_users ON app_users.id = posts.user_id
1585 LEFT JOIN post_tags ON post_tags.post_id = posts.id
1586 WHERE
1587 hidden = FALSE AND
1588 user_id = $1 AND
1589 (post_tags.name = $2 OR hidden = true) AND
1590 publish_at::date <= CURRENT_DATE AND
1591 cur_space = $3
1592 ORDER BY publish_at DESC
1593 LIMIT $4 OFFSET $5`, SelectPost)
1594 err := me.Db.Select(
1595 &posts,
1596 query,
1597 userID,
1598 tag,
1599 space,
1600 page.Num,
1601 page.Num*page.Page,
1602 )
1603 if err != nil {
1604 return nil, err
1605 }
1606
1607 var count int
1608 err = me.Db.QueryRow(`SELECT count(id) FROM posts WHERE hidden = FALSE AND cur_space=$1`, space).Scan(&count)
1609 if err != nil {
1610 return nil, err
1611 }
1612
1613 pager := &db.Paginate[*db.Post]{
1614 Data: posts,
1615 Total: int(math.Ceil(float64(count) / float64(page.Num))),
1616 }
1617 return pager, nil
1618}
1619
1620func (me *PsqlDB) FindPostsByTag(pager *db.Pager, tag, space string) (*db.Paginate[*db.Post], error) {
1621 query := `
1622 SELECT
1623 posts.id,
1624 user_id,
1625 filename,
1626 slug,
1627 title,
1628 text,
1629 description,
1630 publish_at,
1631 app_users.name as username,
1632 posts.updated_at,
1633 posts.mime_type
1634 FROM posts
1635 LEFT JOIN app_users ON app_users.id = posts.user_id
1636 LEFT JOIN post_tags ON post_tags.post_id = posts.id
1637 WHERE
1638 post_tags.name = $3 AND
1639 publish_at::date <= CURRENT_DATE AND
1640 cur_space = $4
1641 ORDER BY publish_at DESC
1642 LIMIT $1 OFFSET $2`
1643 rs, err := me.Db.Queryx(
1644 query,
1645 pager.Num,
1646 pager.Num*pager.Page,
1647 tag,
1648 space,
1649 )
1650 if err != nil {
1651 return nil, err
1652 }
1653 defer func() { _ = rs.Close() }()
1654
1655 return me.postPager(rs, pager.Num, space, tag)
1656}
1657
1658func (me *PsqlDB) FindPopularTags(space string) ([]string, error) {
1659 tags := make([]string, 0)
1660 query := `
1661 SELECT name, count(post_id) as "tally"
1662 FROM post_tags
1663 LEFT JOIN posts ON posts.id = post_id
1664 WHERE posts.cur_space = $1
1665 GROUP BY name
1666 ORDER BY tally DESC
1667 LIMIT 5`
1668 rs, err := me.Db.Queryx(query, space)
1669 if err != nil {
1670 return tags, err
1671 }
1672 defer func() { _ = rs.Close() }()
1673 for rs.Next() {
1674 name := ""
1675 tally := 0
1676 err := rs.Scan(&name, &tally)
1677 if err != nil {
1678 return tags, err
1679 }
1680
1681 tags = append(tags, name)
1682 }
1683 if rs.Err() != nil {
1684 return tags, rs.Err()
1685 }
1686 return tags, nil
1687}
1688
1689func (me *PsqlDB) FindFeature(userID string, feature string) (*db.FeatureFlag, error) {
1690 ff := &db.FeatureFlag{}
1691 err := me.Db.Get(ff, `SELECT * FROM feature_flags WHERE user_id = $1 AND name = $2 ORDER BY expires_at DESC LIMIT 1`, userID, feature)
1692 if err != nil {
1693 return nil, err
1694 }
1695 return ff, nil
1696}
1697
1698func (me *PsqlDB) FindFeaturesByUser(userID string) ([]*db.FeatureFlag, error) {
1699 var features []*db.FeatureFlag
1700 // https://stackoverflow.com/a/16920077
1701 query := `SELECT DISTINCT ON (name) *
1702 FROM feature_flags
1703 WHERE user_id=$1
1704 ORDER BY name, expires_at DESC;`
1705 err := me.Db.Select(&features, query, userID)
1706 if err != nil {
1707 return nil, err
1708 }
1709 return features, nil
1710}
1711
1712func (me *PsqlDB) HasFeatureByUser(userID string, feature string) bool {
1713 ff, err := me.FindFeature(userID, feature)
1714 if err != nil {
1715 return false
1716 }
1717 return ff.IsValid()
1718}
1719
1720func (me *PsqlDB) InsertFeedItems(postID string, items []*db.FeedItem) error {
1721 tx, err := me.Db.Beginx()
1722 if err != nil {
1723 return err
1724 }
1725 defer func() {
1726 _ = tx.Rollback()
1727 }()
1728
1729 for _, item := range items {
1730 _, err := tx.Exec(
1731 `INSERT INTO feed_items (post_id, guid, data) VALUES ($1, $2, $3) RETURNING id;`,
1732 item.PostID,
1733 item.GUID,
1734 item.Data,
1735 )
1736 if err != nil {
1737 return fmt.Errorf(
1738 "post id:%s, link:%s, guid:%s, err:%w",
1739 item.PostID, item.Data.Link, item.GUID, err,
1740 )
1741 }
1742 }
1743
1744 err = tx.Commit()
1745 return err
1746}
1747
1748func (me *PsqlDB) FindFeedItemsByPostID(postID string) ([]*db.FeedItem, error) {
1749 var items []*db.FeedItem
1750 err := me.Db.Select(&items, `SELECT * FROM feed_items WHERE post_id=$1`, postID)
1751 if err != nil {
1752 return nil, err
1753 }
1754 return items, nil
1755}
1756
1757func (me *PsqlDB) InsertProject(userID, name, projectDir string) (string, error) {
1758 if !shared.IsValidSubdomain(name) {
1759 return "", fmt.Errorf("'%s' is not a valid project name, must match /^[a-z0-9-]+$/", name)
1760 }
1761
1762 var id string
1763 err := me.Db.QueryRow(`INSERT INTO projects (user_id, name, project_dir) VALUES ($1, $2, $3) RETURNING id;`, userID, name, projectDir).Scan(&id)
1764 if err != nil {
1765 return "", err
1766 }
1767 return id, nil
1768}
1769
1770func (me *PsqlDB) UpdateProject(userID, name string) error {
1771 _, err := me.Db.Exec(`UPDATE projects SET updated_at = $3 WHERE user_id = $1 AND name = $2;`, userID, name, time.Now())
1772 return err
1773}
1774
1775func (me *PsqlDB) FindProjectByName(userID, name string) (*db.Project, error) {
1776 project := &db.Project{}
1777 err := me.Db.Get(project, `SELECT * FROM projects WHERE user_id = $1 AND name = $2;`, userID, name)
1778 if err != nil {
1779 return nil, err
1780 }
1781 return project, nil
1782}
1783
1784func (me *PsqlDB) InsertToken(userID, name string) (string, error) {
1785 var token string
1786 err := me.Db.QueryRow(`INSERT INTO tokens (user_id, name) VALUES($1, $2) RETURNING token;`, userID, name).Scan(&token)
1787 if err != nil {
1788 return "", err
1789 }
1790 return token, nil
1791}
1792
1793func (me *PsqlDB) UpsertToken(userID, name string) (string, error) {
1794 token, _ := me.findTokenByName(userID, name)
1795 if token != "" {
1796 return token, nil
1797 }
1798
1799 token, err := me.InsertToken(userID, name)
1800 return token, err
1801}
1802
1803func (me *PsqlDB) findTokenByName(userID, name string) (string, error) {
1804 var token string
1805 err := me.Db.QueryRow(`SELECT token FROM tokens WHERE user_id = $1 AND name = $2`, userID, name).Scan(&token)
1806 if err != nil {
1807 return "", err
1808 }
1809 return token, nil
1810}
1811
1812func (me *PsqlDB) RemoveToken(tokenID string) error {
1813 _, err := me.Db.Exec(`DELETE FROM tokens WHERE id = $1`, tokenID)
1814 return err
1815}
1816
1817func (me *PsqlDB) FindTokensByUser(userID string) ([]*db.Token, error) {
1818 var tokens []*db.Token
1819 err := me.Db.Select(&tokens, `SELECT * FROM tokens WHERE user_id = $1`, userID)
1820 if err != nil {
1821 return nil, err
1822 }
1823 return tokens, nil
1824}
1825
1826func (me *PsqlDB) InsertFeature(userID, name string, expiresAt time.Time) (*db.FeatureFlag, error) {
1827 var featureID string
1828 err := me.Db.QueryRow(
1829 `INSERT INTO feature_flags (user_id, name, expires_at) VALUES ($1, $2, $3) RETURNING id;`,
1830 userID,
1831 name,
1832 expiresAt,
1833 ).Scan(&featureID)
1834 if err != nil {
1835 return nil, err
1836 }
1837
1838 feature, err := me.FindFeature(userID, name)
1839 if err != nil {
1840 return nil, err
1841 }
1842
1843 return feature, nil
1844}
1845
1846func (me *PsqlDB) RemoveFeature(userID string, name string) error {
1847 _, err := me.Db.Exec(`DELETE FROM feature_flags WHERE user_id = $1 AND name = $2`, userID, name)
1848 return err
1849}
1850
1851func (me *PsqlDB) createFeatureExpiresAt(userID, name string) time.Time {
1852 ff, _ := me.FindFeature(userID, name)
1853 // if the feature flag has already expired we don't want to add a year to it since that will
1854 // not grant the user a full year
1855 if ff == nil || !ff.IsValid() {
1856 t := time.Now()
1857 return t.AddDate(1, 0, 0)
1858 }
1859 return ff.ExpiresAt.AddDate(1, 0, 0)
1860}
1861
1862func (me *PsqlDB) AddPicoPlusUser(username, email, paymentType, txId string) error {
1863 user, err := me.FindUserByName(username)
1864 if err != nil {
1865 return err
1866 }
1867
1868 tx, err := me.Db.Beginx()
1869 if err != nil {
1870 return err
1871 }
1872 defer func() {
1873 _ = tx.Rollback()
1874 }()
1875
1876 var paymentHistoryId sql.NullString
1877 if paymentType != "" {
1878 data := db.PaymentHistoryData{
1879 Notes: "",
1880 TxID: txId,
1881 }
1882
1883 err := tx.QueryRow(
1884 `INSERT INTO payment_history (user_id, payment_type, amount, data) VALUES ($1, $2, 24 * 1000000, $3) RETURNING id;`,
1885 user.ID,
1886 paymentType,
1887 data,
1888 ).Scan(&paymentHistoryId)
1889 if err != nil {
1890 return err
1891 }
1892 }
1893
1894 plus := me.createFeatureExpiresAt(user.ID, "plus")
1895 plusQuery := fmt.Sprintf(`INSERT INTO feature_flags (user_id, name, data, expires_at, payment_history_id)
1896 VALUES ($1, 'plus', '{"storage_max":10000000000, "file_max":100000000, "email": "%s"}'::jsonb, $2, $3);`, email)
1897 _, err = tx.Exec(plusQuery, user.ID, plus, paymentHistoryId)
1898 if err != nil {
1899 return err
1900 }
1901
1902 return tx.Commit()
1903}
1904
1905func (me *PsqlDB) AddFeatureUser(username, name string) error {
1906 user, err := me.FindUserByName(username)
1907 if err != nil {
1908 return err
1909 }
1910
1911 expiresAt := me.createFeatureExpiresAt(user.ID, name)
1912 _, err = me.InsertFeature(user.ID, name, expiresAt)
1913 return err
1914}
1915
1916func (me *PsqlDB) UpsertProject(userID, projectName, projectDir string) (*db.Project, error) {
1917 project, err := me.FindProjectByName(userID, projectName)
1918 if err == nil {
1919 // this just updates the `createdAt` timestamp, useful for book-keeping
1920 err = me.UpdateProject(userID, projectName)
1921 if err != nil {
1922 me.Logger.Error("could not update project", "err", err)
1923 return nil, err
1924 }
1925 return project, nil
1926 }
1927
1928 _, err = me.InsertProject(userID, projectName, projectName)
1929 if err != nil {
1930 me.Logger.Error("could not create project", "err", err)
1931 return nil, err
1932 }
1933 return me.FindProjectByName(userID, projectName)
1934}
1935
1936func (me *PsqlDB) findPagesStats(userID string) (*db.UserServiceStats, error) {
1937 stats := db.UserServiceStats{
1938 Service: "pgs",
1939 }
1940 err := me.Db.QueryRow(
1941 `SELECT count(id), min(created_at), max(created_at), max(updated_at) FROM projects WHERE user_id=$1`,
1942 userID,
1943 ).Scan(&stats.Num, &stats.FirstCreatedAt, &stats.LastestCreatedAt, &stats.LatestUpdatedAt)
1944 if err != nil {
1945 return nil, err
1946 }
1947
1948 return &stats, nil
1949}
1950
1951func (me *PsqlDB) InsertTunsEventLog(log *db.TunsEventLog) error {
1952 _, err := me.Db.Exec(
1953 `INSERT INTO tuns_event_logs
1954 (user_id, server_id, remote_addr, event_type, tunnel_type, connection_type, tunnel_id)
1955 VALUES
1956 ($1, $2, $3, $4, $5, $6, $7)`,
1957 log.UserId, log.ServerID, log.RemoteAddr, log.EventType, log.TunnelType,
1958 log.ConnectionType, log.TunnelID,
1959 )
1960 return err
1961}
1962
1963func (me *PsqlDB) FindTunsEventLogsByAddr(userID, addr string) ([]*db.TunsEventLog, error) {
1964 var logs []*db.TunsEventLog
1965 err := me.Db.Select(&logs,
1966 `SELECT * FROM tuns_event_logs WHERE user_id=$1 AND tunnel_id=$2 ORDER BY created_at DESC`, userID, addr)
1967 if err != nil {
1968 return nil, err
1969 }
1970 return logs, nil
1971}
1972
1973func (me *PsqlDB) FindTunsEventLogs(userID string) ([]*db.TunsEventLog, error) {
1974 var logs []*db.TunsEventLog
1975 err := me.Db.Select(&logs,
1976 `SELECT * FROM tuns_event_logs WHERE user_id=$1 ORDER BY created_at DESC`, userID)
1977 if err != nil {
1978 return nil, err
1979 }
1980 return logs, nil
1981}
1982
1983func (me *PsqlDB) FindUserStats(userID string) (*db.UserStats, error) {
1984 stats := db.UserStats{}
1985 rs, err := me.Db.Queryx(`SELECT cur_space, count(id), min(created_at), max(created_at), max(updated_at) FROM posts WHERE user_id=$1 GROUP BY cur_space`, userID)
1986 if err != nil {
1987 return nil, err
1988 }
1989 defer func() { _ = rs.Close() }()
1990
1991 for rs.Next() {
1992 stat := db.UserServiceStats{}
1993 err := rs.Scan(&stat.Service, &stat.Num, &stat.FirstCreatedAt, &stat.LastestCreatedAt, &stat.LatestUpdatedAt)
1994 if err != nil {
1995 return nil, err
1996 }
1997 switch stat.Service {
1998 case "prose":
1999 stats.Prose = stat
2000 case "pastes":
2001 stats.Pastes = stat
2002 case "feeds":
2003 stats.Feeds = stat
2004 }
2005 }
2006
2007 if rs.Err() != nil {
2008 return nil, rs.Err()
2009 }
2010
2011 pgs, err := me.findPagesStats(userID)
2012 if err != nil {
2013 return nil, err
2014 }
2015 stats.Pages = *pgs
2016 return &stats, err
2017}
2018
2019func (me *PsqlDB) FindAccessLogs(userID string, fromDate *time.Time) ([]*db.AccessLog, error) {
2020 var logs []*db.AccessLog
2021 err := me.Db.Select(&logs, `SELECT * FROM access_logs WHERE user_id=$1 AND created_at >= $2 ORDER BY created_at ASC`, userID, fromDate)
2022 if err != nil {
2023 return nil, err
2024 }
2025 return logs, nil
2026}
2027
2028func (me *PsqlDB) FindAccessLogsByPubkey(pubkey string, fromDate *time.Time) ([]*db.AccessLog, error) {
2029 var logs []*db.AccessLog
2030 err := me.Db.Select(&logs, `SELECT * FROM access_logs WHERE pubkey=$1 AND created_at >= $2 ORDER BY created_at ASC`, pubkey, fromDate)
2031 if err != nil {
2032 return nil, err
2033 }
2034 return logs, nil
2035}
2036
2037func (me *PsqlDB) FindPubkeysInAccessLogs(userID string) ([]string, error) {
2038 var pubkeys []string
2039 err := me.Db.Select(&pubkeys, `SELECT DISTINCT(pubkey) FROM access_logs WHERE user_id=$1`, userID)
2040 if err != nil {
2041 return nil, err
2042 }
2043 return pubkeys, nil
2044}
2045
2046func (me *PsqlDB) InsertAccessLog(log *db.AccessLog) error {
2047 _, err := me.Db.Exec(
2048 `INSERT INTO access_logs (user_id, service, pubkey, identity) VALUES ($1, $2, $3, $4);`,
2049 log.UserID,
2050 log.Service,
2051 log.Pubkey,
2052 log.Identity,
2053 )
2054 return err
2055}
2056
2057func (me *PsqlDB) UpsertPipeMonitor(userID, topic string, dur time.Duration, winEnd *time.Time) error {
2058 durStr := fmt.Sprintf("%d seconds", int64(dur.Seconds()))
2059 _, err := me.Db.Exec(
2060 `INSERT INTO pipe_monitors (user_id, topic, window_dur, window_end)
2061 VALUES ($1, $2, $3::interval, $4)
2062 ON CONFLICT (user_id, topic) DO UPDATE SET window_dur = $3::interval, window_end = $4, updated_at = NOW();`,
2063 userID,
2064 topic,
2065 durStr,
2066 winEnd,
2067 )
2068 return err
2069}
2070
2071func (me *PsqlDB) UpdatePipeMonitorLastPing(userID, topic string, lastPing *time.Time) error {
2072 _, err := me.Db.Exec(
2073 `UPDATE pipe_monitors SET last_ping = $3, updated_at = NOW() WHERE user_id = $1 AND topic = $2;`,
2074 userID,
2075 topic,
2076 lastPing,
2077 )
2078 return err
2079}
2080
2081func (me *PsqlDB) RemovePipeMonitor(userID, topic string) error {
2082 _, err := me.Db.Exec(
2083 `DELETE FROM pipe_monitors WHERE user_id = $1 AND topic = $2;`,
2084 userID,
2085 topic,
2086 )
2087 return err
2088}
2089
2090func (me *PsqlDB) FindPipeMonitorByTopic(userID, topic string) (*db.PipeMonitor, error) {
2091 monitor := &db.PipeMonitor{}
2092 err := me.Db.Get(monitor, `SELECT id, user_id, topic, (EXTRACT(EPOCH FROM window_dur) * 1000000000)::bigint as window_dur, window_end, last_ping, created_at, updated_at FROM pipe_monitors WHERE user_id = $1 AND topic = $2;`, userID, topic)
2093 if err != nil {
2094 return nil, err
2095 }
2096 return monitor, nil
2097}
2098
2099func (me *PsqlDB) FindPipeMonitorsByUser(userID string) ([]*db.PipeMonitor, error) {
2100 var monitors []*db.PipeMonitor
2101 err := me.Db.Select(&monitors, `SELECT id, user_id, topic, (EXTRACT(EPOCH FROM window_dur) * 1000000000)::bigint as window_dur, window_end, last_ping, created_at, updated_at FROM pipe_monitors WHERE user_id = $1 ORDER BY topic;`, userID)
2102 if err != nil {
2103 return nil, err
2104 }
2105 return monitors, nil
2106}
2107
2108func (me *PsqlDB) InsertPipeMonitorHistory(monitorID string, windowDur time.Duration, windowEnd, lastPing *time.Time) error {
2109 durStr := fmt.Sprintf("%d seconds", int64(windowDur.Seconds()))
2110 _, err := me.Db.Exec(
2111 `INSERT INTO pipe_monitors_history (monitor_id, window_dur, window_end, last_ping) VALUES ($1, $2::interval, $3, $4)`,
2112 monitorID, durStr, windowEnd, lastPing,
2113 )
2114 return err
2115}
2116
2117func (me *PsqlDB) FindPipeMonitorHistory(monitorID string, from, to time.Time) ([]*db.PipeMonitorHistory, error) {
2118 var history []*db.PipeMonitorHistory
2119 err := me.Db.Select(
2120 &history,
2121 `SELECT id, monitor_id, (EXTRACT(EPOCH FROM window_dur) * 1000000000)::bigint as window_dur, window_end, last_ping, created_at, updated_at FROM pipe_monitors_history WHERE monitor_id = $1 AND last_ping <= $2 AND window_end >= $3 ORDER BY last_ping ASC`,
2122 monitorID, to, from,
2123 )
2124 if err != nil {
2125 return nil, err
2126 }
2127 return history, nil
2128}