Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
3.2 kB · 125 lines
Go
at next-dev
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126package db
import ( "context" "database/sql" "errors" "fmt" "log"
"tangled.org/core/knotfeed" "tangled.org/core/knotmirror/models")
func UpsertHost(ctx context.Context, e DBTX, host *models.Host) error { if _, err := e.ExecContext(ctx, `insert into hosts (hostname, no_ssl, status, last_seq, last_feed) values ($1, $2, $3, $4, $5) on conflict(hostname) do update set no_ssl = excluded.no_ssl, status = excluded.status, last_seq = excluded.last_seq, last_feed = excluded.last_feed `, host.Hostname, host.NoSSL, host.Status, host.Cursor.Seq(), host.Cursor.Feed().Token(), ); err != nil { return fmt.Errorf("upserting host: %w", err) } return nil}
type rowScanner interface { Scan(dest ...any) error}
const hostColumns = `select hostname, no_ssl, status, last_seq, last_feed from hosts`
func scanHost(row rowScanner) (models.Host, error) { var host models.Host var seq int64 var token string if err := row.Scan(&host.Hostname, &host.NoSSL, &host.Status, &seq, &token); err != nil { return host, err } host.Cursor = hostCursor(host.Hostname, seq, token) return host, nil}
func GetHost(ctx context.Context, e DBTX, hostname string) (*models.Host, error) { host, err := scanHost(e.QueryRowContext(ctx, hostColumns+` where hostname = $1`, hostname)) if errors.Is(err, sql.ErrNoRows) { return nil, nil } if err != nil { return nil, err } return &host, nil}
func hostCursor(hostname string, seq int64, token string) knotfeed.Cursor { feed, err := knotfeed.ParseFeed(token) if err != nil { log.Println("host cursor stored an unrecognized feed, resuming live", "host", hostname, "feed", token) return knotfeed.Cursor{} } return knotfeed.NewCursor(feed, seq)}
func SetHostStatus(ctx context.Context, e DBTX, hostname string, status models.HostStatus) error { if _, err := e.ExecContext(ctx, `update hosts set status = $1 where hostname = $2`, status, hostname, ); err != nil { return fmt.Errorf("setting host status: %w", err) } return nil}
func StoreCursors(ctx context.Context, e *sql.DB, cursors []models.HostCursor) error { tx, err := e.BeginTx(ctx, nil) if err != nil { return fmt.Errorf("starting transaction: %w", err) } defer tx.Rollback() for _, cur := range cursors { if cur.Cursor.Seq() < 0 { continue } if _, err := tx.ExecContext(ctx, `update hosts set last_seq = $1, last_feed = $2 where hostname = $3`, cur.Cursor.Seq(), cur.Cursor.Feed().Token(), cur.Hostname, ); err != nil { log.Println("couldn't persist host cursor", "host", cur.Hostname, "feed", cur.Cursor.Feed(), "lastSeq", cur.Cursor.Seq(), "err", err) } } return tx.Commit()}
func ListHosts(ctx context.Context, e DBTX, status models.HostStatus) ([]models.Host, error) { rows, err := e.QueryContext(ctx, hostColumns+` where status = $1`, status) if err != nil { return nil, fmt.Errorf("querying hosts: %w", err) } defer rows.Close()
var hosts []models.Host for rows.Next() { host, err := scanHost(rows) if err != nil { return nil, fmt.Errorf("scanning row: %w", err) } hosts = append(hosts, host) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("scanning rows: %w ", err) } return hosts, nil}