mirror of
https://tangled.org/evan.jarrett.net/at-container-registry
synced 2026-08-29 04:06:58 +00:00
34 lines
888 B
Go
34 lines
888 B
Go
package db
|
|
|
|
import (
|
|
"database/sql"
|
|
"errors"
|
|
)
|
|
|
|
// GetJetstreamCursor returns the last persisted Jetstream cursor (time_us).
|
|
// Returns 0 when no cursor has been saved yet (e.g. fresh database).
|
|
func GetJetstreamCursor(db DBTX) (int64, error) {
|
|
var cursor int64
|
|
err := db.QueryRow(`SELECT cursor FROM jetstream_cursor WHERE id = 1`).Scan(&cursor)
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return 0, nil
|
|
}
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return cursor, nil
|
|
}
|
|
|
|
// SaveJetstreamCursor writes the given cursor to the singleton jetstream_cursor row.
|
|
// Idempotent — safe to call on every tick.
|
|
func SaveJetstreamCursor(db DBTX, cursor int64) error {
|
|
_, err := db.Exec(`
|
|
INSERT INTO jetstream_cursor (id, cursor, updated_at)
|
|
VALUES (1, ?, CURRENT_TIMESTAMP)
|
|
ON CONFLICT(id) DO UPDATE SET
|
|
cursor = excluded.cursor,
|
|
updated_at = excluded.updated_at
|
|
`, cursor)
|
|
return err
|
|
}
|