Compare commits
No commits in common. "63e01c5b6d9e17408e940a0999a02396ca74c8ab" and "9ef2d7689dd09a3438af1ac3a98375720b31c086" have entirely different histories.
63e01c5b6d
...
9ef2d7689d
12 changed files with 37 additions and 427 deletions
|
|
@ -1,7 +0,0 @@
|
||||||
.git
|
|
||||||
.env
|
|
||||||
*.db
|
|
||||||
*.db-wal
|
|
||||||
*.db-shm
|
|
||||||
bin/
|
|
||||||
st-agent-worker
|
|
||||||
16
.eng.example
16
.eng.example
|
|
@ -1,16 +0,0 @@
|
||||||
# generate this in Settings/API/Create API Key
|
|
||||||
ST_API_KEY=
|
|
||||||
|
|
||||||
# secret for webhook, generate with openssl rand -hex 32
|
|
||||||
WEBHOOK_SECRET=
|
|
||||||
|
|
||||||
# set eligible appointment types for the action, comma-separated (for example, senate, recall)
|
|
||||||
ELIGIBLE_APPOINTMENTS=s
|
|
||||||
|
|
||||||
# location of SQLite database file
|
|
||||||
DB_PATH=worker.db
|
|
||||||
|
|
||||||
|
|
||||||
# set this to 0 to write to Solidarity Tech's API, default is to keep it read_only
|
|
||||||
DRY_RUN=1
|
|
||||||
|
|
||||||
13
Dockerfile
13
Dockerfile
|
|
@ -1,13 +0,0 @@
|
||||||
FROM golang:1.27-alpine AS build
|
|
||||||
WORKDIR /src
|
|
||||||
COPY go.mod go.sum ./
|
|
||||||
RUN go mod download
|
|
||||||
COPY . .
|
|
||||||
RUN CGO_ENABLED=0 go build -o /worker .
|
|
||||||
|
|
||||||
FROM alpine:3.20
|
|
||||||
RUN apk add --no-cache ca-certificates tzdata
|
|
||||||
COPY --from=build /worker /worker
|
|
||||||
RUN mkdir -p /data
|
|
||||||
EXPOSE 8080
|
|
||||||
ENTRYPOINT ["/worker"]
|
|
||||||
133
README.md
133
README.md
|
|
@ -1,133 +0,0 @@
|
||||||
# st-agent-worker
|
|
||||||
|
|
||||||
Assigns Solidarity Tech agents to department members automatically.
|
|
||||||
|
|
||||||
When someone is given the Department Rep role in Solidarity Tech, an ST
|
|
||||||
automation calls this service. The service looks up that person's chapter
|
|
||||||
(campus) and department, finds every eligible member of the same
|
|
||||||
chapter/department pair, and assigns the rep as their agent.
|
|
||||||
|
|
||||||
## Rules
|
|
||||||
|
|
||||||
- Only members whose appointment type is Senate or Recall are assigned.
|
|
||||||
Set `ELIGIBLE_APPOINTMENTS` to change this.
|
|
||||||
- Matching is on primary chapteri (UCLAFA, BFA, SCFA...) plus Department 1.
|
|
||||||
A rep at one campus is never assigned members at another.
|
|
||||||
- Reps are excluded from assignment. Organizer leads (OC members, organizing
|
|
||||||
staff) assign their agents by hand in the ST web interface.
|
|
||||||
- **Fill-only**: a member who already has an agent is left alone and
|
|
||||||
recorded in the `skips` table. No agent is ever reassigned automatically.
|
|
||||||
See below for consequences of this.
|
|
||||||
|
|
||||||
### Known consequence of fill-only
|
|
||||||
|
|
||||||
Some departments currently have two reps. Whichever rep runs first takes
|
|
||||||
the whole department; the second gets nobody. The reconciler logs a warning
|
|
||||||
in both cases that can be seen in the logs and reconciled manually.
|
|
||||||
|
|
||||||
_Potential fix_ In the future, round-robin splitting could be implemented. `RepsFor` already returns all reps for a department sorted by ID, so adding it is a
|
|
||||||
change to how the agent is picked, and wouldn't require a code restructuring.
|
|
||||||
|
|
||||||
## Architecture
|
|
||||||
|
|
||||||
Four components do the work:
|
|
||||||
|
|
||||||
- `server.go` — HTTP handler. Verifies the shared secret, reads the rep's
|
|
||||||
ST user ID from the query string, looks up their chapter and department
|
|
||||||
live, takes a per-department lock, returns 200, and runs the
|
|
||||||
reconciliation in the background. The 200 is returned before any work
|
|
||||||
happens, because ST times out and retries otherwise.
|
|
||||||
- `client.go` — Solidarity Tech API client. Holds the rate limiter and
|
|
||||||
handles 429 responses. Nothing above it builds a URL.
|
|
||||||
- `store.go` — SQLite database handle. Four tables: `people` (cached roster),
|
|
||||||
`runs` (locks and progress), `skips` (members left alone), `sync_state` (last
|
|
||||||
sync timestamp).
|
|
||||||
- `reconcile.go` — the loop. Reads department members from the cache,
|
|
||||||
checks each one's current agent live against the API, assigns or skips.
|
|
||||||
|
|
||||||
Roster data is cached because it changes more slowly, and API rate limits would make a read/write of the whole roster take two hours. Instead, agent assignments are read live on every check, so an assignment made by hand in the dashboard
|
|
||||||
during a run can be respected right away. See _Rate limits_ below.
|
|
||||||
|
|
||||||
### Locking
|
|
||||||
|
|
||||||
`runs` has a partial unique index on `(chapter_id, department)` where
|
|
||||||
status is `running`. A second webhook for the same department fails the
|
|
||||||
insert and exits. On startup, any row still marked `running` belongs to a
|
|
||||||
dead process and is marked `crashed`.
|
|
||||||
|
|
||||||
### Rate limits
|
|
||||||
|
|
||||||
The ST API allows 60 requests per 30 seconds per API key, which averages
|
|
||||||
2 per second and permits bursts. The client is configured to match this.
|
|
||||||
|
|
||||||
## Configuration
|
|
||||||
|
|
||||||
| Variable | Meaning |
|
|
||||||
| --- | --- |
|
|
||||||
| `ST_API_KEY` | Solidarity Tech API key |
|
|
||||||
| `WEBHOOK_SECRET` | Shared secret, passed as `?s=` on the webhook URL |
|
|
||||||
| `DB_PATH` | SQLite file. `/data/worker.db` in the container |
|
|
||||||
| `DRY_RUN` | Any non-empty value logs writes instead of sending them |
|
|
||||||
| `ELIGIBLE_APPOINTMENTS` | Comma-separated labels, defaults to `senate,recall` |
|
|
||||||
|
|
||||||
The ST webhook action has no support for custom headers, so the secret
|
|
||||||
travels in the query string. Do not log raw query strings.
|
|
||||||
|
|
||||||
## Endpoints
|
|
||||||
|
|
||||||
- `POST /webhook?id={user_id}&s={secret}` — triggered by the ST automation
|
|
||||||
when someone is assigned the Department Rep role.
|
|
||||||
- `POST /debug/sync?s={secret}` — forces a roster sync. Useful after a bulk
|
|
||||||
import rather than waiting for the hourly update.
|
|
||||||
- `GET /healthz` — returns 200 once the server is listening.
|
|
||||||
|
|
||||||
## Health checks
|
|
||||||
|
|
||||||
The startup sync runs before the server binds, and a cold start with an
|
|
||||||
empty volume pulls the full roster, which takes about 90 seconds. Give the
|
|
||||||
container a `start_period` of at least 120 seconds or it will be marked
|
|
||||||
unhealthy before it comes up.
|
|
||||||
|
|
||||||
After deploying, check:
|
|
||||||
|
|
||||||
curl https://$HOST/healthz
|
|
||||||
|
|
||||||
and confirm the log shows a completed sync. Then:
|
|
||||||
|
|
||||||
sqlite3 /data/worker.db "
|
|
||||||
SELECT COUNT(*) total, SUM(is_eligible) eligible, SUM(is_rep) reps
|
|
||||||
FROM people;"
|
|
||||||
|
|
||||||
Should give roughly 17,500 total, 14,700 eligible.
|
|
||||||
|
|
||||||
Log lines to watch for:
|
|
||||||
|
|
||||||
- `WARNING: run N: assigned nobody` — every member already had an agent.
|
|
||||||
Expected for the second rep in a co-repped department, unexpected
|
|
||||||
otherwise.
|
|
||||||
- `N eligible users synced with no department` — those members cannot be
|
|
||||||
assigned to anyone. These are likely Senate or Recall duplicates, or
|
|
||||||
people we don't have information for.
|
|
||||||
|
|
||||||
## Initial backfill
|
|
||||||
|
|
||||||
The webhook only fires when a rep is newly assigned the role, so reps who
|
|
||||||
already had it when this was deployed never trigger it. To assign their
|
|
||||||
departments, call the endpoint for each:
|
|
||||||
|
|
||||||
sqlite3 worker.db "SELECT id FROM people WHERE is_rep = 1;" | while read id; do
|
|
||||||
curl -s -o /dev/null -X POST "https://$HOST/webhook?id=$id&s=$WEBHOOK_SECRET"
|
|
||||||
sleep 1
|
|
||||||
done
|
|
||||||
|
|
||||||
Run it with `DRY_RUN` set first and count the assignments in the log.
|
|
||||||
|
|
||||||
## Not built yet
|
|
||||||
|
|
||||||
- A nightly sweep to reassign reps. Without it, a rep who steps down keeps
|
|
||||||
their assignments and their replacement inherits nobody. Every member added
|
|
||||||
after their department's rep signed up also goes unassigned until
|
|
||||||
something triggers that department again. **FIRST PRIORITY**
|
|
||||||
- Schema migrations. The schema is created with `CREATE TABLE IF NOT
|
|
||||||
EXISTS`, which does not add columns to an existing table. Changing the
|
|
||||||
schema currently means deleting the database and re-syncing.
|
|
||||||
58
client.go
58
client.go
|
|
@ -10,7 +10,6 @@ import (
|
||||||
"net/url"
|
"net/url"
|
||||||
"strconv"
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
"bytes"
|
|
||||||
|
|
||||||
"golang.org/x/time/rate"
|
"golang.org/x/time/rate"
|
||||||
)
|
)
|
||||||
|
|
@ -124,7 +123,7 @@ func departmentOf(props map[string]json.RawMessage) string {
|
||||||
|
|
||||||
//optionLabels returne the labels of an option-valued custom-property
|
//optionLabels returne the labels of an option-valued custom-property
|
||||||
// ST serialises select and checkbox fields as [{label, value}], and the
|
// ST serialises select and checkbox fields as [{label, value}], and the
|
||||||
// same label appears more than one time
|
// same label can appear more than once.
|
||||||
|
|
||||||
func optionLabels(props map[string]json.RawMessage, key string) []string {
|
func optionLabels(props map[string]json.RawMessage, key string) []string {
|
||||||
raw, ok := props[key]
|
raw, ok := props[key]
|
||||||
|
|
@ -164,61 +163,6 @@ func (c *Client) GetUser(ctx context.Context, id int64) (*User, error) {
|
||||||
return u.toUser(), nil
|
return u.toUser(), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ActiveAgent returns the user ID of the person's current agent, or 0 if
|
|
||||||
// they have none. The is_active query param is ignored by the API, so
|
|
||||||
// rows are filtered here.
|
|
||||||
func (c *Client) ActiveAgent(ctx context.Context, userID int64) (int64, error) {
|
|
||||||
q := url.Values{}
|
|
||||||
q.Set("user_id", strconv.FormatInt(userID, 10))
|
|
||||||
q.Set("_limit", "100")
|
|
||||||
|
|
||||||
raw, err := c.do(ctx, http.MethodGet, "/agent_assignments?"+q.Encode(), nil)
|
|
||||||
if err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
var env struct {
|
|
||||||
Data []struct {
|
|
||||||
AgentUserID int64 `json:"agent_user_id"`
|
|
||||||
IsActive bool `json:"is_active"`
|
|
||||||
} `json:"data"`
|
|
||||||
}
|
|
||||||
if err := json.Unmarshal(raw, &env); err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, a := range env.Data {
|
|
||||||
if a.IsActive {
|
|
||||||
return a.AgentUserID, nil
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return 0, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// AssignAgent creates an agent assignment. ST deactivates any prior active
|
|
||||||
// assignment automatically. Not idempotent: a repeat call creates a new row.
|
|
||||||
func (c *Client) AssignAgent(ctx context.Context, userID, agentID int64) error {
|
|
||||||
if c.dryRun {
|
|
||||||
log.Printf("DRY RUN: would assign user %d to agent %d", userID, agentID)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
body, err := json.Marshal(map[string]int64{
|
|
||||||
"user_id": userID,
|
|
||||||
"agent_user_id": agentID,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = c.do(ctx, http.MethodPost, "/agent_assignments", bytes.NewReader(body))
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
log.Printf("assigned user %d to agent %d", userID, agentID)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// ListUsers pages the full user list. If since is non-zero, only users
|
// ListUsers pages the full user list. If since is non-zero, only users
|
||||||
// updated after that Unix timestamp are returned.
|
// updated after that Unix timestamp are returned.
|
||||||
func (c *Client) ListUsers(ctx context.Context, since int64) ([]User, error) {
|
func (c *Client) ListUsers(ctx context.Context, since int64) ([]User, error) {
|
||||||
|
|
|
||||||
|
|
@ -1,21 +0,0 @@
|
||||||
services:
|
|
||||||
worker:
|
|
||||||
build: .
|
|
||||||
restart: unless-stopped
|
|
||||||
environment:
|
|
||||||
ST_API_KEY: ${ST_API_KEY}
|
|
||||||
WEBHOOK_SECRET: ${WEBHOOK_SECRET}
|
|
||||||
ELIGIBLE_APPOINTMENTS: ${ELIGIBLE_APPOINTMENTS}
|
|
||||||
DB_PATH: /data/worker.db
|
|
||||||
DRY_RUN: ${DRY_RUN:-}
|
|
||||||
volumes:
|
|
||||||
- worker-data:/data
|
|
||||||
healthcheck:
|
|
||||||
test: ["CMD", "wget", "-qO-", "http://localhost:8080/healthz"]
|
|
||||||
interval: 30s
|
|
||||||
timeout: 5s
|
|
||||||
retries: 3
|
|
||||||
start_period: 120s
|
|
||||||
|
|
||||||
volumes:
|
|
||||||
worker-data:
|
|
||||||
17
fakes.go
17
fakes.go
|
|
@ -6,23 +6,8 @@ import (
|
||||||
|
|
||||||
type fakeST struct {
|
type fakeST struct {
|
||||||
dept string
|
dept string
|
||||||
chapterID int64
|
|
||||||
err error
|
err error
|
||||||
|
chapterID int64
|
||||||
agents map[int64]int64 // userID → existing agent, 0 or absent = none
|
|
||||||
assigned map[int64]int64 // recorded writes
|
|
||||||
}
|
|
||||||
|
|
||||||
func (f *fakeST) ActiveAgent(ctx context.Context, userID int64) (int64, error) {
|
|
||||||
return f.agents[userID], nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (f *fakeST) AssignAgent(ctx context.Context, userID, agentID int64) error {
|
|
||||||
if f.assigned == nil {
|
|
||||||
f.assigned = map[int64]int64{}
|
|
||||||
}
|
|
||||||
f.assigned[userID] = agentID
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *fakeST) GetUser(ctx context.Context, id int64) (*User, error) {
|
func (f *fakeST) GetUser(ctx context.Context, id int64) (*User, error) {
|
||||||
|
|
|
||||||
2
go.mod
2
go.mod
|
|
@ -1,6 +1,6 @@
|
||||||
module ucfa/st-agent-worker
|
module ucfa/st-agent-worker
|
||||||
|
|
||||||
go 1.26.0
|
go 1.27.0
|
||||||
|
|
||||||
require (
|
require (
|
||||||
golang.org/x/time v0.16.0
|
golang.org/x/time v0.16.0
|
||||||
|
|
|
||||||
22
main.go
22
main.go
|
|
@ -5,8 +5,6 @@ import (
|
||||||
"net/http"
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
"os"
|
"os"
|
||||||
"context"
|
|
||||||
"time"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// ELIGIBLE_APPOINTMENTS is a comma-separated list of appointment type
|
// ELIGIBLE_APPOINTMENTS is a comma-separated list of appointment type
|
||||||
|
|
@ -59,26 +57,6 @@ func main() {
|
||||||
}
|
}
|
||||||
|
|
||||||
srv := NewServer(NewClient(apiKey, dryRun), store, secret, loadEligible())
|
srv := NewServer(NewClient(apiKey, dryRun), store, secret, loadEligible())
|
||||||
// Prime the cache: a fresh container starts empty.
|
|
||||||
{
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Minute)
|
|
||||||
if err := srv.SyncRoster(ctx); err != nil {
|
|
||||||
log.Printf("initial sync failed: %v", err)
|
|
||||||
}
|
|
||||||
cancel()
|
|
||||||
}
|
|
||||||
|
|
||||||
go func() {
|
|
||||||
t := time.NewTicker(time.Hour)
|
|
||||||
defer t.Stop()
|
|
||||||
for range t.C {
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Minute)
|
|
||||||
if err := srv.SyncRoster(ctx); err != nil {
|
|
||||||
log.Printf("scheduled sync failed: %v", err)
|
|
||||||
}
|
|
||||||
cancel()
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
log.Println("listening on :8080")
|
log.Println("listening on :8080")
|
||||||
log.Fatal(http.ListenAndServe(":8080", srv.Routes()))
|
log.Fatal(http.ListenAndServe(":8080", srv.Routes()))
|
||||||
}
|
}
|
||||||
|
|
|
||||||
75
reconcile.go
75
reconcile.go
|
|
@ -1,75 +0,0 @@
|
||||||
package main
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"log"
|
|
||||||
)
|
|
||||||
|
|
||||||
func (s *Server) reconcile(ctx context.Context, runID int64, rep *User, dept string) error {
|
|
||||||
reps, err := s.store.RepsFor(rep.ChapterID, dept)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if len(reps) == 0 {
|
|
||||||
// Shouldn't happen: the webhook fired for a rep in this department.
|
|
||||||
log.Printf("run %d: no reps found for chapter %d dept %q", runID, rep.ChapterID, dept)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
if len(reps) > 1 {
|
|
||||||
log.Printf("WARNING: run %d: %d reps for chapter %d dept %q; assigning all to %d (fill-only mode)",
|
|
||||||
runID, len(reps), rep.ChapterID, dept, reps[0])
|
|
||||||
}
|
|
||||||
agent := reps[0]
|
|
||||||
|
|
||||||
members, err := s.store.EligibleMembers(rep.ChapterID, dept)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
// loop
|
|
||||||
assigned, skipped, failed := 0, 0, 0
|
|
||||||
|
|
||||||
for _, m := range members {
|
|
||||||
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
log.Printf("run %d: cancelled after %d assigned, %d skipped", runID, assigned, skipped)
|
|
||||||
return ctx.Err()
|
|
||||||
default:
|
|
||||||
}
|
|
||||||
|
|
||||||
existing, err := s.st.ActiveAgent(ctx, m)
|
|
||||||
if err != nil {
|
|
||||||
log.Printf("run %d: checking agent for %d failed: %v", runID, m, err)
|
|
||||||
failed++
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if existing != 0 {
|
|
||||||
// Fill-only: never displace an existing agent.
|
|
||||||
skipped++
|
|
||||||
if err := s.store.LogSkip(runID, m, existing, agent); err != nil {
|
|
||||||
log.Printf("run %d: logging skip for %d failed: %v", runID, m, err)
|
|
||||||
}
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := s.st.AssignAgent(ctx, m, agent); err != nil {
|
|
||||||
log.Printf("run %d: assigning %d to %d failed: %v", runID, m, agent, err)
|
|
||||||
failed++
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
assigned++
|
|
||||||
|
|
||||||
if err := s.store.SetProgress(runID, m); err != nil {
|
|
||||||
log.Printf("run %d: recording progress failed: %v", runID, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if assigned == 0 && skipped > 0 {
|
|
||||||
log.Printf("WARNING: run %d: assigned nobody; all %d members already have agents",
|
|
||||||
runID, skipped)
|
|
||||||
}
|
|
||||||
log.Printf("run %d: %d assigned, %d skipped, dept %q chapter %d",
|
|
||||||
runID, assigned, skipped, dept, rep.ChapterID)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
43
server.go
43
server.go
|
|
@ -4,6 +4,7 @@ import (
|
||||||
"context"
|
"context"
|
||||||
"crypto/subtle"
|
"crypto/subtle"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
|
@ -24,8 +25,6 @@ type User struct {
|
||||||
// STClient is what the handler needs from the Solidarity Tech API.
|
// STClient is what the handler needs from the Solidarity Tech API.
|
||||||
type STClient interface {
|
type STClient interface {
|
||||||
GetUser(ctx context.Context, id int64) (*User, error)
|
GetUser(ctx context.Context, id int64) (*User, error)
|
||||||
ActiveAgent(ctx context.Context, userID int64) (int64, error)
|
|
||||||
AssignAgent(ctx context.Context, userID, agentID int64) error
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Store is what the handler needs from persistence.
|
// Store is what the handler needs from persistence.
|
||||||
|
|
@ -36,10 +35,6 @@ type Store interface {
|
||||||
UpsertPeople(users []User, eligible map[string]bool) error
|
UpsertPeople(users []User, eligible map[string]bool) error
|
||||||
LastSync() (int64, error)
|
LastSync() (int64, error)
|
||||||
SetLastSync(ts int64) error
|
SetLastSync(ts int64) error
|
||||||
RepsFor(chapterID int64, dept string) ([]int64, error)
|
|
||||||
EligibleMembers(chapterID int64, dept string) ([]int64, error)
|
|
||||||
LogSkip(runID, userID, existingAgent, wouldHaveAssigned int64) error
|
|
||||||
SetProgress(runID, lastUserID int64) error
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// ErrLocked is returned by AcquireLock when a run is already in flight
|
// ErrLocked is returned by AcquireLock when a run is already in flight
|
||||||
|
|
@ -83,12 +78,36 @@ func (s *Server) Routes() *http.ServeMux {
|
||||||
mux.HandleFunc("GET /healthz", func(w http.ResponseWriter, r *http.Request) {
|
mux.HandleFunc("GET /healthz", func(w http.ResponseWriter, r *http.Request) {
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
})
|
})
|
||||||
|
|
||||||
// TEMPORARY: remove before deploy.
|
// TEMPORARY: remove before deploy.
|
||||||
mux.HandleFunc("POST /debug/sync", func(w http.ResponseWriter, r *http.Request) {
|
mux.HandleFunc("GET /debug/users", func(w http.ResponseWriter, r *http.Request) {
|
||||||
if !s.validSecret(r) {
|
client, ok := s.st.(*Client)
|
||||||
http.Error(w, "unauthorized", http.StatusUnauthorized)
|
if !ok {
|
||||||
|
http.Error(w, "not a real client", 500)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
users, err := client.ListUsers(r.Context(), 0)
|
||||||
|
if err != nil {
|
||||||
|
http.Error(w, err.Error(), 500)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
counts := map[string]int{}
|
||||||
|
eligible := 0
|
||||||
|
for _, u := range users {
|
||||||
|
if !u.Eligible(s.eligible) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
eligible++
|
||||||
|
counts[fmt.Sprintf("%d|%s", u.ChapterID, normalizeDept(u.Department))]++
|
||||||
|
}
|
||||||
|
fmt.Fprintf(w, "%d users, %d eligible\n\n", len(users), eligible)
|
||||||
|
for k, n := range counts {
|
||||||
|
fmt.Fprintf(w, "%6d %s\n", n, k)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
// TEMPORARY: remove before deploy.
|
||||||
|
mux.HandleFunc("POST /debug/sync", func(w http.ResponseWriter, r *http.Request) {
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
if err := s.SyncRoster(ctx); err != nil {
|
if err := s.SyncRoster(ctx); err != nil {
|
||||||
|
|
@ -170,6 +189,12 @@ func (s *Server) handleWebhook(w http.ResponseWriter, r *http.Request) {
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// reconcile is stubbed until piece 3 lands.
|
||||||
|
func (s *Server) reconcile(ctx context.Context, runID int64, rep *User, dept string) error {
|
||||||
|
log.Printf("run %d: would reconcile dept %q for rep %d", runID, dept, rep.ID)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// ST doesn't support setting headers so we use the URL
|
// ST doesn't support setting headers so we use the URL
|
||||||
|
|
||||||
func (s *Server) validSecret(r *http.Request) bool {
|
func (s *Server) validSecret(r *http.Request) bool {
|
||||||
|
|
|
||||||
57
store.go
57
store.go
|
|
@ -18,18 +18,9 @@ var schemaStatements = []string{
|
||||||
department TEXT NOT NULL,
|
department TEXT NOT NULL,
|
||||||
rep_user_id INTEGER NOT NULL,
|
rep_user_id INTEGER NOT NULL,
|
||||||
status TEXT NOT NULL,
|
status TEXT NOT NULL,
|
||||||
last_user_id INTEGER,
|
|
||||||
started_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
started_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||||
ended_at TEXT
|
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
|
`CREATE UNIQUE INDEX IF NOT EXISTS one_run_per_dept
|
||||||
ON runs(chapter_id, department) WHERE status = 'running'`,
|
ON runs(chapter_id, department) WHERE status = 'running'`,
|
||||||
`CREATE TABLE IF NOT EXISTS people (
|
`CREATE TABLE IF NOT EXISTS people (
|
||||||
|
|
@ -66,38 +57,6 @@ func OpenStore(path string) (*SQLStore, error) {
|
||||||
return &SQLStore{db: db}, nil
|
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) {
|
func (s *SQLStore) AcquireLock(chapterID int64, dept string, repID int64) (int64, error) {
|
||||||
res, err := s.db.Exec(
|
res, err := s.db.Exec(
|
||||||
`INSERT INTO runs (chapter_id, department, rep_user_id, status)
|
`INSERT INTO runs (chapter_id, department, rep_user_id, status)
|
||||||
|
|
@ -188,19 +147,3 @@ func (s *SQLStore) AcquireLock(chapterID int64, dept string, repID int64) (int64
|
||||||
strconv.FormatInt(ts, 10))
|
strconv.FormatInt(ts, 10))
|
||||||
return err
|
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
|
|
||||||
}
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue