aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
Diffstat (limited to 'internal/database/worker.go')
-rw-r--r--internal/database/worker.go146+102 −44
1 files changed, 102 insertions, 44 deletions
diff --git a/internal/database/worker.go b/internal/database/worker.go
index bdd4b0b..c0ff46b 100644
--- a/internal/database/worker.go
+++ b/internal/database/worker.go
@@ -8,52 +8,138 @@ import (
"errors"
"time"
+ "github.com/google/uuid"
+ sb "github.com/huandu/go-sqlbuilder"
"github.com/jackc/pgx/v5"
)
var (
- ErrAddWorker = errors.New("database: failed to add worker")
- ErrRevokeWorker = errors.New("database: failed to revoke worker")
ErrWorkerNotFound = errors.New("database: worker not found")
)
// Worker holds info about a registered worker.
type Worker struct {
- ID string
+ ID uuid.UUID
+ OrgID *uuid.UUID
PublicKey []byte
CreatedAt time.Time
}
-// ListWorkers returns all active (non-revoked) workers.
-func (db *DB) ListWorkers(ctx context.Context) ([]Worker, error) {
- rows, err := db.Pool.Query(ctx,
- `SELECT id, public_key, created_at FROM workers WHERE revoked_at IS NULL ORDER BY created_at`,
+// GetWorker returns a worker by ID.
+func (db *DB) GetWorker(ctx context.Context, id uuid.UUID) (*Worker, error) {
+ var w Worker
+ if err := db.Pool.QueryRow(ctx,
+ `SELECT id, public_key, org_id, created_at FROM workers WHERE id = $1 AND revoked_at IS NULL`,
+ id,
+ ).Scan(&w.ID, &w.PublicKey, &w.OrgID, &w.CreatedAt); err != nil {
+ if errors.Is(err, pgx.ErrNoRows) {
+ return nil, ErrWorkerNotFound
+ }
+
+ return nil, err
+ }
+
+ return &w, nil
+}
+
+// CreateWorker registers a new worker with the given public key and optional org.
+func (db *DB) CreateWorker(ctx context.Context, publicKey []byte, org *OrgRef) (uuid.UUID, error) {
+ tx, err := db.Pool.Begin(ctx)
+ if err != nil {
+ return uuid.Nil, err
+ }
+ defer tx.Rollback(ctx)
+
+ var orgID *uuid.UUID
+ if org != nil {
+ id, err := resolveOrg(ctx, tx, *org)
+ if err != nil {
+ return uuid.Nil, err
+ }
+ orgID = &id
+ }
+
+ var workerID uuid.UUID
+ if err := tx.QueryRow(ctx,
+ `INSERT INTO workers (public_key, org_id) VALUES ($1, $2) RETURNING id`,
+ publicKey, orgID,
+ ).Scan(&workerID); err != nil {
+ return uuid.Nil, err
+ }
+
+ return workerID, tx.Commit(ctx)
+}
+
+// DeleteWorker soft-deletes a worker by ID.
+func (db *DB) DeleteWorker(ctx context.Context, id uuid.UUID) error {
+ tag, err := db.Pool.Exec(ctx,
+ `UPDATE workers SET revoked_at = now() WHERE id = $1 AND revoked_at IS NULL`, id,
)
if err != nil {
- return nil, err
+ return err
+ }
+
+ if tag.RowsAffected() == 0 {
+ return ErrWorkerNotFound
+ }
+
+ return nil
+}
+
+// ListWorkers returns a page of workers and the total count.
+func (db *DB) ListWorkers(ctx context.Context, cursor uuid.UUID, limit int, filter string) ([]Worker, int, error) {
+ if filter != "" {
+ return nil, 0, ErrFilterNotImplemented
+ }
+
+ tx, err := db.Pool.Begin(ctx)
+ if err != nil {
+ return nil, 0, err
+ }
+ defer tx.Rollback(ctx)
+
+ var total int
+ if err := tx.QueryRow(ctx,
+ `SELECT count(*) FROM workers WHERE revoked_at IS NULL`,
+ ).Scan(&total); err != nil {
+ return nil, 0, err
+ }
+
+ q := sb.PostgreSQL.NewSelectBuilder()
+ q.Select("id", "public_key", "org_id", "created_at").
+ From("workers").
+ Where(q.IsNull("revoked_at")).
+ OrderBy("id").
+ Limit(limit)
+ if cursor != uuid.Nil {
+ q.Where(q.GreaterThan("id", cursor))
+ }
+
+ sql, args := q.Build()
+ rows, err := tx.Query(ctx, sql, args...)
+ if err != nil {
+ return nil, 0, err
}
defer rows.Close()
var workers []Worker
for rows.Next() {
var w Worker
- if err := rows.Scan(&w.ID, &w.PublicKey, &w.CreatedAt); err != nil {
- return nil, err
+ if err := rows.Scan(&w.ID, &w.PublicKey, &w.OrgID, &w.CreatedAt); err != nil {
+ return nil, 0, err
}
workers = append(workers, w)
}
-
- return workers, rows.Err()
+ return workers, total, rows.Err()
}
// LookupWorker finds an active worker by its ed25519 public key.
func (db *DB) LookupWorker(ctx context.Context, publicKey []byte) (*Worker, error) {
var w Worker
- err := db.Pool.QueryRow(ctx,
- `SELECT id, public_key, created_at FROM workers WHERE public_key = $1 AND revoked_at IS NULL`,
+ if err := db.Pool.QueryRow(ctx,
+ `SELECT id, public_key, org_id, created_at FROM workers WHERE public_key = $1 AND revoked_at IS NULL`,
publicKey,
- ).Scan(&w.ID, &w.PublicKey, &w.CreatedAt)
- if err != nil {
+ ).Scan(&w.ID, &w.PublicKey, &w.OrgID, &w.CreatedAt); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrWorkerNotFound
}
@@ -61,31 +147,3 @@ func (db *DB) LookupWorker(ctx context.Context, publicKey []byte) (*Worker, erro
}
return &w, nil
}
-
-// AddWorker registers a new worker with the given public key.
-func (db *DB) AddWorker(ctx context.Context, publicKey []byte) (string, error) {
- var id string
- err := db.Pool.QueryRow(ctx,
- `INSERT INTO workers (public_key) VALUES ($1) RETURNING id`,
- publicKey,
- ).Scan(&id)
- if err != nil {
- return "", errors.Join(ErrAddWorker, err)
- }
- return id, nil
-}
-
-// WorkerRevoke soft-deletes a worker by ID.
-func (db *DB) WorkerRevoke(ctx context.Context, id string) error {
- tag, err := db.Pool.Exec(ctx,
- `UPDATE workers SET revoked_at = now() WHERE id = $1 AND revoked_at IS NULL`,
- id,
- )
- if err != nil {
- return errors.Join(ErrRevokeWorker, err)
- }
- if tag.RowsAffected() == 0 {
- return ErrWorkerNotFound
- }
- return nil
-}