st-agent-worker/store.go
2026-09-10 12:39:36 +02:00

206 lines
5.7 KiB
Go

package main
import (
"database/sql"
"errors"
"fmt"
"strconv"
"strings"
"modernc.org/sqlite"
sqlite3 "modernc.org/sqlite/lib"
)
var schemaStatements = []string{
`CREATE TABLE IF NOT EXISTS runs (
id INTEGER PRIMARY KEY,
chapter_id INTEGER NOT NULL,
department TEXT NOT NULL,
rep_user_id INTEGER NOT NULL,
status TEXT NOT NULL,
last_user_id INTEGER,
started_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
ended_at TEXT
)`,
`CREATE TABLE IF NOT EXISTS skips (
id INTEGER PRIMARY KEY,
run_id INTEGER NOT NULL,
user_id INTEGER NOT NULL,
existing_agent INTEGER NOT NULL,
would_have_assigned INTEGER NOT NULL,
at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
)`,
`CREATE UNIQUE INDEX IF NOT EXISTS one_run_per_dept
ON runs(chapter_id, department) WHERE status = 'running'`,
`CREATE TABLE IF NOT EXISTS people (
id INTEGER PRIMARY KEY,
email TEXT,
chapter_id INTEGER,
department TEXT,
appointment_type TEXT,
is_rep INTEGER NOT NULL DEFAULT 0,
is_eligible INTEGER NOT NULL DEFAULT 0,
synced_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
)`,
`CREATE INDEX IF NOT EXISTS people_dept ON people(chapter_id, department)`,
`CREATE TABLE IF NOT EXISTS sync_state (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)`,
}
type SQLStore struct{ db *sql.DB }
func OpenStore(path string) (*SQLStore, error) {
db, err := sql.Open("sqlite",
"file:"+path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)")
if err != nil {
return nil, err
}
db.SetMaxOpenConns(1)
for i, stmt := range schemaStatements {
if _, err := db.Exec(stmt); err != nil {
return nil, fmt.Errorf("schema statement %d: %w", i, err)
}
}
return &SQLStore{db: db}, nil
}
// RepsFor returns the user IDs of every rep for a department, ordered by
// ID so assignment is deterministic across runs.
func (s *SQLStore) RepsFor(chapterID int64, dept string) ([]int64, error) {
return s.userIDs(`SELECT id FROM people
WHERE is_rep = 1 AND chapter_id = ? AND department = ?
ORDER BY id`, chapterID, dept)
}
func (s *SQLStore) EligibleMembers(chapterID int64, dept string) ([]int64, error) {
return s.userIDs(`SELECT id FROM people
WHERE is_eligible = 1 AND is_rep = 0 AND chapter_id = ? AND department = ?
ORDER BY id`, chapterID, dept)
}
func (s *SQLStore) userIDs(query string, args ...any) ([]int64, error) {
rows, err := s.db.Query(query, args...)
if err != nil {
return nil, err
}
defer rows.Close()
var out []int64
for rows.Next() {
var id int64
if err := rows.Scan(&id); err != nil {
return nil, err
}
out = append(out, id)
}
return out, rows.Err()
}
func (s *SQLStore) AcquireLock(chapterID int64, dept string, repID int64) (int64, error) {
res, err := s.db.Exec(
`INSERT INTO runs (chapter_id, department, rep_user_id, status)
VALUES (?, ?, ?, 'running')`, chapterID, dept, repID)
if err != nil {
var serr *sqlite.Error
if errors.As(err, &serr) && serr.Code() == sqlite3.SQLITE_CONSTRAINT_UNIQUE {
return 0, ErrLocked
}
return 0, ErrLocked
}
return res.LastInsertId()
}
func (s *SQLStore) ReleaseLock(runID int64) error {
_, err := s.db.Exec(
`UPDATE runs SET status='done', ended_at=CURRENT_TIMESTAMP WHERE id=?`, runID)
return err
}
func (s *SQLStore) SweepStaleLocks() error {
_, err := s.db.Exec(
`UPDATE runs SET status='crashed', ended_at=CURRENT_TIMESTAMP
WHERE status='running'`)
return err
}
func (s *SQLStore) UpsertPeople(users []User, eligible map[string]bool) error {
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
stmt, err := tx.Prepare(`
INSERT INTO people (id, email, chapter_id, department, appointment_type, is_rep, is_eligible, synced_at)
VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)
ON CONFLICT(id) DO UPDATE SET
email=excluded.email,
chapter_id=excluded.chapter_id,
department=excluded.department,
appointment_type=excluded.appointment_type,
is_rep=excluded.is_rep,
is_rep=excluded.is_rep,
synced_at=CURRENT_TIMESTAMP`)
if err != nil {
return err
}
defer stmt.Close()
//loop body
for _, u := range users {
rep, elig := 0, 0
if u.IsRep {
rep = 1
}
if u.Eligible(eligible) {
elig = 1
}
_, err := stmt.Exec(u.ID, u.Email, u.ChapterID,
normalizeDept(u.Department),
strings.Join(u.AppointmentTypes, ","),
rep, elig)
if err != nil {
return fmt.Errorf("upserting user %d: %w", u.ID, err)
}
}
return tx.Commit()
}
func (s *SQLStore) LastSync() (int64, error) {
var v string
err := s.db.QueryRow(`SELECT value FROM sync_state WHERE key='users_synced_at'`).Scan(&v)
if err == sql.ErrNoRows {
return 0, nil
}
if err != nil {
return 0, err
}
return strconv.ParseInt(v, 10, 64)
}
func (s *SQLStore) SetLastSync(ts int64) error {
_, err := s.db.Exec(`
INSERT INTO sync_state (key, value) VALUES ('users_synced_at', ?)
ON CONFLICT(key) DO UPDATE SET value=excluded.value`,
strconv.FormatInt(ts, 10))
return err
}
// LogSkip records a member left alone because they already had an agent.
// This is the manual-fix queue: nothing else records what fill-only cost.
func (s *SQLStore) LogSkip(runID, userID, existingAgent, wouldHaveAssigned int64) error {
_, err := s.db.Exec(
`INSERT INTO skips (run_id, user_id, existing_agent, would_have_assigned)
VALUES (?, ?, ?, ?)`,
runID, userID, existingAgent, wouldHaveAssigned)
return err
}
func (s *SQLStore) SetProgress(runID, lastUserID int64) error {
_, err := s.db.Exec(
`UPDATE runs SET last_user_id = ? WHERE id = ?`, lastUserID, runID)
return err
}