claude-licence-manager holds the anthropic-licence-manager seat: it reads every node's holdings state, adopts a login it does not hold by refreshing it (newest first, once per account), keeps each grant alive under a lease, publishes what each consumer should hold as its bindings state with a generation, and answers current sealed to the consumer's key. Postgres store prepared by a run-once step; grants encrypted with the vault's key. claude-code reports what its node holds (fingerprints and account, never a token), hands its grant over only when the manager asks, watches its binding and fetches the token on a newer generation, and writes access-token-only. Its ask now reads the runtime's answer as a value and addresses seats as seats.
282 lines
9.7 KiB
Go
282 lines
9.7 KiB
Go
package main
|
|
|
|
// The store on the mesh's postgres (novox/hq design 39 §1), the database the mesh provisioned for this
|
|
// module. The schema is brought to this version's shape by the preparation step (ADR 0135), each
|
|
// statement idempotent.
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
)
|
|
|
|
// Schema is what this version needs.
|
|
var Schema = []string{
|
|
`create table if not exists licence (
|
|
name text primary key,
|
|
kind text not null check (kind in ('subscription', 'api-key')),
|
|
account_uuid text unique,
|
|
email text,
|
|
organization_uuid text,
|
|
sealed text,
|
|
refresh_fingerprint text,
|
|
access_expires_at bigint,
|
|
refresh_expires_at bigint,
|
|
failures integer not null default 0,
|
|
notified_at bigint,
|
|
adopted_at bigint not null,
|
|
rotated_at bigint)`,
|
|
`create sequence if not exists binding_generation`,
|
|
`create table if not exists binding (
|
|
consumer text primary key,
|
|
licence text not null references licence(name),
|
|
generation bigint not null)`,
|
|
`create table if not exists lease (
|
|
key text primary key,
|
|
holder text not null,
|
|
until timestamptz not null)`,
|
|
`create table if not exists offered (
|
|
fingerprint text primary key,
|
|
node text not null,
|
|
account_uuid text,
|
|
outcome text not null,
|
|
why text not null,
|
|
at timestamptz not null default now())`,
|
|
`create table if not exists usage (
|
|
licence text not null,
|
|
at bigint not null,
|
|
reading jsonb not null,
|
|
raw jsonb not null)`,
|
|
`create index if not exists usage_by_licence on usage (licence, at desc)`,
|
|
`create table if not exists audit (
|
|
at timestamptz not null default now(),
|
|
what text not null,
|
|
detail jsonb not null)`,
|
|
}
|
|
|
|
// PgStore is the store on postgres.
|
|
type PgStore struct{ pool *pgxpool.Pool }
|
|
|
|
// PgStoreFromEnv opens the store the mesh provisioned, its URL in the file DATABASE_URL_FILE names.
|
|
func PgStoreFromEnv(ctx context.Context) (*PgStore, error) {
|
|
file := os.Getenv("DATABASE_URL_FILE")
|
|
if file == "" {
|
|
return nil, errors.New("DATABASE_URL_FILE is not set: the manager's database is a requirement the mesh resolves")
|
|
}
|
|
raw, err := os.ReadFile(file)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return OpenPgStore(ctx, strings.TrimSpace(string(raw)))
|
|
}
|
|
|
|
// OpenPgStore opens a store at a URL.
|
|
func OpenPgStore(ctx context.Context, url string) (*PgStore, error) {
|
|
pool, err := pgxpool.New(ctx, url)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return &PgStore{pool: pool}, nil
|
|
}
|
|
|
|
// Migrate brings the schema to this version's shape.
|
|
func (s *PgStore) Migrate(ctx context.Context) error {
|
|
for _, q := range Schema {
|
|
if _, err := s.pool.Exec(ctx, q); err != nil {
|
|
return fmt.Errorf("%s: %w", strings.Fields(q)[0:6], err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
const licenceColumns = `name, kind, coalesce(account_uuid,''), coalesce(email,''), coalesce(organization_uuid,''),
|
|
coalesce(sealed,''), coalesce(refresh_fingerprint,''), coalesce(access_expires_at,0), coalesce(refresh_expires_at,0),
|
|
failures, coalesce(notified_at,0), adopted_at, coalesce(rotated_at,0)`
|
|
|
|
func scanLicence(row pgx.Row) (*Licence, error) {
|
|
var l Licence
|
|
err := row.Scan(&l.Name, &l.Kind, &l.AccountUUID, &l.Email, &l.OrganizationUUID, &l.Sealed, &l.RefreshFingerprint,
|
|
&l.AccessExpiresAt, &l.RefreshExpiresAt, &l.Failures, &l.NotifiedAt, &l.AdoptedAt, &l.RotatedAt)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, nil
|
|
}
|
|
return &l, err
|
|
}
|
|
|
|
func (s *PgStore) Licences(ctx context.Context) ([]Licence, error) {
|
|
rows, err := s.pool.Query(ctx, `select `+licenceColumns+` from licence order by name`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []Licence
|
|
for rows.Next() {
|
|
l, err := scanLicence(rows)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, *l)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
func (s *PgStore) Licence(ctx context.Context, name string) (*Licence, error) {
|
|
return scanLicence(s.pool.QueryRow(ctx, `select `+licenceColumns+` from licence where name = $1`, name))
|
|
}
|
|
|
|
func (s *PgStore) LicenceForAccount(ctx context.Context, account string) (*Licence, error) {
|
|
return scanLicence(s.pool.QueryRow(ctx, `select `+licenceColumns+` from licence where account_uuid = $1`, account))
|
|
}
|
|
|
|
func nullable(s string) any {
|
|
if s == "" {
|
|
return nil
|
|
}
|
|
return s
|
|
}
|
|
|
|
func nullableInt(v int64) any {
|
|
if v == 0 {
|
|
return nil
|
|
}
|
|
return v
|
|
}
|
|
|
|
func (s *PgStore) SaveLicence(ctx context.Context, l Licence) error {
|
|
_, err := s.pool.Exec(ctx, `insert into licence (name, kind, account_uuid, email, organization_uuid, sealed, refresh_fingerprint,
|
|
access_expires_at, refresh_expires_at, failures, notified_at, adopted_at, rotated_at)
|
|
values ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13)
|
|
on conflict (name) do update set kind = excluded.kind, account_uuid = excluded.account_uuid, email = excluded.email,
|
|
organization_uuid = excluded.organization_uuid, sealed = excluded.sealed, refresh_fingerprint = excluded.refresh_fingerprint,
|
|
access_expires_at = excluded.access_expires_at, refresh_expires_at = excluded.refresh_expires_at,
|
|
failures = excluded.failures, notified_at = excluded.notified_at, adopted_at = excluded.adopted_at,
|
|
rotated_at = excluded.rotated_at`,
|
|
l.Name, l.Kind, nullable(l.AccountUUID), nullable(l.Email), nullable(l.OrganizationUUID), nullable(l.Sealed),
|
|
nullable(l.RefreshFingerprint), nullableInt(l.AccessExpiresAt), nullableInt(l.RefreshExpiresAt), l.Failures,
|
|
nullableInt(l.NotifiedAt), l.AdoptedAt, nullableInt(l.RotatedAt))
|
|
return err
|
|
}
|
|
|
|
// Lease is taken in the store before a row is read (design 39 §3): a second run started together finds
|
|
// it live and does nothing.
|
|
func (s *PgStore) Lease(ctx context.Context, key, holder string, d time.Duration) (bool, error) {
|
|
tag, err := s.pool.Exec(ctx, `insert into lease (key, holder, until) values ($1, $2, now() + make_interval(secs => $3))
|
|
on conflict (key) do update set holder = excluded.holder, until = excluded.until
|
|
where lease.until < now() or lease.holder = excluded.holder`, key, holder, d.Seconds())
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return tag.RowsAffected() == 1, nil
|
|
}
|
|
|
|
func (s *PgStore) Unlease(ctx context.Context, key, holder string) error {
|
|
_, err := s.pool.Exec(ctx, `delete from lease where key = $1 and holder = $2`, key, holder)
|
|
return err
|
|
}
|
|
|
|
func (s *PgStore) bindingsWhere(ctx context.Context, q string, args ...any) ([]Binding, error) {
|
|
rows, err := s.pool.Query(ctx, q, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []Binding
|
|
for rows.Next() {
|
|
var b Binding
|
|
if err := rows.Scan(&b.Consumer, &b.Licence, &b.Generation); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, b)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
func (s *PgStore) Bindings(ctx context.Context) ([]Binding, error) {
|
|
return s.bindingsWhere(ctx, `select consumer, licence, generation from binding order by consumer`)
|
|
}
|
|
|
|
func (s *PgStore) Binding(ctx context.Context, consumer string) (*Binding, error) {
|
|
out, err := s.bindingsWhere(ctx, `select consumer, licence, generation from binding where consumer = $1`, consumer)
|
|
if err != nil || len(out) == 0 {
|
|
return nil, err
|
|
}
|
|
return &out[0], nil
|
|
}
|
|
|
|
func (s *PgStore) Bind(ctx context.Context, consumer, licence string) (Binding, error) {
|
|
out, err := s.bindingsWhere(ctx, `insert into binding (consumer, licence, generation) values ($1, $2, nextval('binding_generation'))
|
|
on conflict (consumer) do update set licence = excluded.licence, generation = excluded.generation
|
|
returning consumer, licence, generation`, consumer, licence)
|
|
if err != nil {
|
|
return Binding{}, err
|
|
}
|
|
return out[0], nil
|
|
}
|
|
|
|
func (s *PgStore) Unbind(ctx context.Context, consumer string) (bool, error) {
|
|
tag, err := s.pool.Exec(ctx, `delete from binding where consumer = $1`, consumer)
|
|
return err == nil && tag.RowsAffected() == 1, err
|
|
}
|
|
|
|
func (s *PgStore) Advance(ctx context.Context, licence string) ([]Binding, error) {
|
|
return s.bindingsWhere(ctx, `update binding set generation = nextval('binding_generation') where licence = $1
|
|
returning consumer, licence, generation`, licence)
|
|
}
|
|
|
|
func (s *PgStore) Outcome(ctx context.Context, fp string) (Outcome, error) {
|
|
var o string
|
|
err := s.pool.QueryRow(ctx, `select outcome from offered where fingerprint = $1`, fp).Scan(&o)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return "", nil
|
|
}
|
|
return Outcome(o), err
|
|
}
|
|
|
|
func (s *PgStore) RecordOutcome(ctx context.Context, fp, node, account string, o Outcome, why string) error {
|
|
_, err := s.pool.Exec(ctx, `insert into offered (fingerprint, node, account_uuid, outcome, why) values ($1,$2,$3,$4,$5)
|
|
on conflict (fingerprint) do update set outcome = excluded.outcome, why = excluded.why, at = now()`,
|
|
fp, node, nullable(account), string(o), why)
|
|
return err
|
|
}
|
|
|
|
func (s *PgStore) RecordUsage(ctx context.Context, licence string, at int64, r UsageReading, raw map[string]any) error {
|
|
reading, _ := json.Marshal(r)
|
|
rawJSON, _ := json.Marshal(raw)
|
|
_, err := s.pool.Exec(ctx, `insert into usage (licence, at, reading, raw) values ($1,$2,$3,$4)`, licence, at, reading, rawJSON)
|
|
return err
|
|
}
|
|
|
|
func (s *PgStore) Usage(ctx context.Context, licence string, limit int) ([]UsageRow, error) {
|
|
rows, err := s.pool.Query(ctx, `select licence, at, reading from usage where ($1 = '' or licence = $1) order by at desc limit $2`, licence, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []UsageRow
|
|
for rows.Next() {
|
|
var u UsageRow
|
|
var reading []byte
|
|
if err := rows.Scan(&u.Licence, &u.At, &reading); err != nil {
|
|
return nil, err
|
|
}
|
|
_ = json.Unmarshal(reading, &u.Reading)
|
|
out = append(out, u)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
func (s *PgStore) Audit(ctx context.Context, what string, detail map[string]any) error {
|
|
raw, _ := json.Marshal(detail)
|
|
_, err := s.pool.Exec(ctx, `insert into audit (what, detail) values ($1, $2)`, what, raw)
|
|
return err
|
|
}
|
|
|
|
func (s *PgStore) Close() { s.pool.Close() }
|