main pico / pkg / db / postgres / storage.go
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}