diff --git a/internal/api/handler/admin_workers_ssh.go b/internal/api/handler/admin_workers_ssh.go index a46df85d7..eaf60c39f 100644 --- a/internal/api/handler/admin_workers_ssh.go +++ b/internal/api/handler/admin_workers_ssh.go @@ -406,7 +406,7 @@ func (h *Handler) AdminRebootWorker(c *gin.Context) { // Body: // // { -// "user_id": "uuid", // the org/user that gets exclusive use +// "organization_id": "uuid", // the org that gets exclusive use // "subscription_id": "uuid", // their active sub // "drain_to_worker_id": "uuid|null" // optional: target for evicted accounts. // // null = let assignment service pick @@ -419,7 +419,7 @@ func (h *Handler) AdminRebootWorker(c *gin.Context) { // are individually idempotent: re-running the endpoint with the same // inputs converges. type convertToDedicatedBody struct { - UserID string `json:"user_id" binding:"required"` + OrganizationID string `json:"organization_id" binding:"required"` SubscriptionID string `json:"subscription_id" binding:"required"` DrainToWorkerID *string `json:"drain_to_worker_id"` } @@ -435,9 +435,9 @@ func (h *Handler) AdminConvertWorkerToDedicated(c *gin.Context) { errx.JSON(c, errx.New(errx.BadRequest, "invalid request body")) return } - userID, err := uuid.Parse(body.UserID) + orgID, err := uuid.Parse(body.OrganizationID) if err != nil { - errx.JSON(c, errx.New(errx.BadRequest, "invalid user_id")) + errx.JSON(c, errx.New(errx.BadRequest, "invalid organization_id")) return } subID, err := uuid.Parse(body.SubscriptionID) @@ -506,7 +506,7 @@ func (h *Handler) AdminConvertWorkerToDedicated(c *gin.Context) { created, err := h.WorkerRepo.CreateDedicatedAssignmentIfNotExists(c.Request.Context(), &models.DedicatedWorkerAssignment{ ID: uuid.New(), WorkerID: id, - UserID: userID, + OrganizationID: orgID, SubscriptionID: subID, AssignedAt: time.Now(), }) @@ -516,7 +516,7 @@ func (h *Handler) AdminConvertWorkerToDedicated(c *gin.Context) { } h.audit(c, "convert_to_dedicated", models.AuditEntityWorker, &id, map[string]string{ - "user_id": userID.String(), + "organization_id": orgID.String(), "subscription_id": subID.String(), "drained_to": movedTo, "accounts_moved": itoa(len(accountIDs)), diff --git a/internal/app/worker/assignment.go b/internal/app/worker/assignment.go index 744c2ef8a..e13dc052c 100644 --- a/internal/app/worker/assignment.go +++ b/internal/app/worker/assignment.go @@ -96,7 +96,7 @@ func (s *workerAssignmentService) AssignWorkerToEmail(ctx context.Context, email } if plan != nil && plan.DedicatedWorkers > 0 { // Check if org has a dedicated worker - dedicatedWorker, err := s.workerRepo.GetDedicatedWorkerByUserID(ctx, orgID) + dedicatedWorker, err := s.workerRepo.GetDedicatedWorkerByOrgID(ctx, orgID) if err != nil { return nil, err } @@ -114,7 +114,7 @@ func (s *workerAssignmentService) AssignWorkerToEmail(ctx context.Context, email return nil, aerr } if aerr == nil || errors.Is(aerr, ErrOrgAlreadyAssigned) { - dedicatedWorker, err = s.workerRepo.GetDedicatedWorkerByUserID(ctx, orgID) + dedicatedWorker, err = s.workerRepo.GetDedicatedWorkerByOrgID(ctx, orgID) if err != nil { return nil, err } @@ -426,7 +426,7 @@ func (s *workerAssignmentService) AssignDedicatedWorker(ctx context.Context, org assignment := &models.DedicatedWorkerAssignment{ ID: uuid.New(), WorkerID: worker.ID, - UserID: orgID, + OrganizationID: orgID, SubscriptionID: subscriptionID, AssignedAt: time.Now(), } @@ -465,7 +465,7 @@ func (s *workerAssignmentService) ReleaseDedicatedWorker(ctx context.Context, or // GetDedicatedWorker gets the dedicated worker for an organization func (s *workerAssignmentService) GetDedicatedWorker(ctx context.Context, orgID uuid.UUID) (*models.Worker, error) { - return s.workerRepo.GetDedicatedWorkerByUserID(ctx, orgID) + return s.workerRepo.GetDedicatedWorkerByOrgID(ctx, orgID) } // MigrateOrgToPremiumWorkers migrates all org's emails from free to premium workers @@ -552,7 +552,7 @@ func (s *workerAssignmentService) MigrateOrgToDedicated(ctx context.Context, org } // Get the dedicated worker - dedicatedWorker, err := s.workerRepo.GetDedicatedWorkerByUserID(ctx, orgID) + dedicatedWorker, err := s.workerRepo.GetDedicatedWorkerByOrgID(ctx, orgID) if err != nil { return err } diff --git a/internal/infrastructure/db/migrations/000055_dedicated_worker_org_assignments.down.sql b/internal/infrastructure/db/migrations/000055_dedicated_worker_org_assignments.down.sql new file mode 100644 index 000000000..1e5b8bb28 --- /dev/null +++ b/internal/infrastructure/db/migrations/000055_dedicated_worker_org_assignments.down.sql @@ -0,0 +1,23 @@ +ALTER TABLE dedicated_worker_assignments + DROP CONSTRAINT IF EXISTS dedicated_worker_assignments_organization_id_fkey; + +ALTER TABLE dedicated_worker_assignments + RENAME COLUMN organization_id TO user_id; + +-- Map org-keyed rows back to the org owner so the users FK can be restored. +UPDATE dedicated_worker_assignments dwa +SET user_id = o.owner_user_id +FROM organizations o +WHERE o.id = dwa.user_id; + +DELETE FROM dedicated_worker_assignments dwa +WHERE NOT EXISTS (SELECT 1 FROM users u WHERE u.id = dwa.user_id); + +ALTER TABLE dedicated_worker_assignments + ADD CONSTRAINT dedicated_worker_assignments_user_id_fkey + FOREIGN KEY (user_id) REFERENCES users(id) ON DELETE CASCADE; + +ALTER TABLE dedicated_worker_assignments + RENAME CONSTRAINT unique_active_org_assignment TO unique_active_user_assignment; + +ALTER INDEX idx_dedicated_org RENAME TO idx_dedicated_user; diff --git a/internal/infrastructure/db/migrations/000055_dedicated_worker_org_assignments.up.sql b/internal/infrastructure/db/migrations/000055_dedicated_worker_org_assignments.up.sql new file mode 100644 index 000000000..c0d4a9d70 --- /dev/null +++ b/internal/infrastructure/db/migrations/000055_dedicated_worker_org_assignments.up.sql @@ -0,0 +1,32 @@ +-- Dedicated workers are organization assets, and the assignment service has +-- always keyed dedicated_worker_assignments by organization id - but the +-- column was named user_id with a users FK, so every runtime insert failed +-- with an FK violation and paid dedicated-plan orgs could never bind a +-- worker. Re-key the table to organizations. + +ALTER TABLE dedicated_worker_assignments + DROP CONSTRAINT IF EXISTS dedicated_worker_assignments_user_id_fkey; + +ALTER TABLE dedicated_worker_assignments + RENAME COLUMN user_id TO organization_id; + +-- Pre-rename rows (seed fixtures) hold owner user ids: remap each to that +-- owner's organization, then drop anything unmappable - such rows were never +-- reachable by the runtime lookups anyway. +UPDATE dedicated_worker_assignments dwa +SET organization_id = o.id +FROM organizations o +WHERE o.owner_user_id = dwa.organization_id + AND NOT EXISTS (SELECT 1 FROM organizations WHERE id = dwa.organization_id); + +DELETE FROM dedicated_worker_assignments dwa +WHERE NOT EXISTS (SELECT 1 FROM organizations o WHERE o.id = dwa.organization_id); + +ALTER TABLE dedicated_worker_assignments + ADD CONSTRAINT dedicated_worker_assignments_organization_id_fkey + FOREIGN KEY (organization_id) REFERENCES organizations(id) ON DELETE CASCADE; + +ALTER TABLE dedicated_worker_assignments + RENAME CONSTRAINT unique_active_user_assignment TO unique_active_org_assignment; + +ALTER INDEX idx_dedicated_user RENAME TO idx_dedicated_org; diff --git a/internal/models/worker.go b/internal/models/worker.go index 9bcf505c4..f26e788ba 100644 --- a/internal/models/worker.go +++ b/internal/models/worker.go @@ -128,7 +128,7 @@ type UpdateWorker struct { type DedicatedWorkerAssignment struct { ID uuid.UUID `json:"id"` WorkerID uuid.UUID `json:"worker_id"` - UserID uuid.UUID `json:"user_id"` + OrganizationID uuid.UUID `json:"organization_id"` SubscriptionID uuid.UUID `json:"subscription_id"` AssignedAt time.Time `json:"assigned_at"` ReleasedAt *time.Time `json:"released_at,omitempty"` diff --git a/internal/repository/pg_admin.go b/internal/repository/pg_admin.go index 98896579e..9144e2f01 100644 --- a/internal/repository/pg_admin.go +++ b/internal/repository/pg_admin.go @@ -204,7 +204,7 @@ func (r *adminRepository) SearchUsers(ctx context.Context, search *models.AdminU whereClause += ` AND EXISTS (SELECT 1 FROM user_bans ub WHERE ub.user_id = u.id)` } if search.HasDedicatedWorker { - whereClause += ` AND EXISTS (SELECT 1 FROM dedicated_worker_assignments dwa WHERE dwa.user_id = u.id AND dwa.released_at IS NULL)` + whereClause += ` AND EXISTS (SELECT 1 FROM dedicated_worker_assignments dwa JOIN organizations o ON o.id = dwa.organization_id WHERE o.owner_user_id = u.id AND dwa.released_at IS NULL)` } // Count / numeric ranges diff --git a/internal/repository/pg_worker.go b/internal/repository/pg_worker.go index 8fabe5b31..ff4c63723 100644 --- a/internal/repository/pg_worker.go +++ b/internal/repository/pg_worker.go @@ -43,7 +43,7 @@ type WorkerRepository interface { CreateDedicatedAssignment(ctx context.Context, assignment *models.DedicatedWorkerAssignment) error CreateDedicatedAssignmentIfNotExists(ctx context.Context, assignment *models.DedicatedWorkerAssignment) (bool, error) GetActiveDedicatedAssignment(ctx context.Context, userID uuid.UUID) (*models.DedicatedWorkerAssignment, error) - GetDedicatedWorkerByUserID(ctx context.Context, userID uuid.UUID) (*models.Worker, error) + GetDedicatedWorkerByOrgID(ctx context.Context, orgID uuid.UUID) (*models.Worker, error) ReleaseDedicatedAssignment(ctx context.Context, userID uuid.UUID) error // Email account worker queries @@ -282,14 +282,14 @@ func (r *workerRepository) SetWorkerType(ctx context.Context, workerID uuid.UUID // CreateDedicatedAssignment creates a new dedicated worker assignment func (r *workerRepository) CreateDedicatedAssignment(ctx context.Context, assignment *models.DedicatedWorkerAssignment) error { query := ` - INSERT INTO dedicated_worker_assignments (id, worker_id, user_id, subscription_id, assigned_at) + INSERT INTO dedicated_worker_assignments (id, worker_id, organization_id, subscription_id, assigned_at) VALUES ($1, $2, $3, $4, $5) ` _, err := r.db.Exec(ctx, query, assignment.ID, assignment.WorkerID, - assignment.UserID, + assignment.OrganizationID, assignment.SubscriptionID, assignment.AssignedAt, ) @@ -297,21 +297,21 @@ func (r *workerRepository) CreateDedicatedAssignment(ctx context.Context, assign } // CreateDedicatedAssignmentIfNotExists atomically creates a dedicated worker assignment -// only if no active (released_at IS NULL) assignment exists for the user. +// only if no active (released_at IS NULL) assignment exists for the organization. // Returns (true, nil) if created, (false, nil) if already exists. func (r *workerRepository) CreateDedicatedAssignmentIfNotExists(ctx context.Context, assignment *models.DedicatedWorkerAssignment) (bool, error) { query := ` - INSERT INTO dedicated_worker_assignments (id, worker_id, user_id, subscription_id, assigned_at) + INSERT INTO dedicated_worker_assignments (id, worker_id, organization_id, subscription_id, assigned_at) SELECT $1, $2, $3, $4, $5 WHERE NOT EXISTS ( SELECT 1 FROM dedicated_worker_assignments - WHERE user_id = $3 AND released_at IS NULL + WHERE organization_id = $3 AND released_at IS NULL ) ` result, err := r.db.Exec(ctx, query, assignment.ID, assignment.WorkerID, - assignment.UserID, + assignment.OrganizationID, assignment.SubscriptionID, assignment.AssignedAt, ) @@ -321,17 +321,17 @@ func (r *workerRepository) CreateDedicatedAssignmentIfNotExists(ctx context.Cont return result.RowsAffected() > 0, nil } -// GetActiveDedicatedAssignment retrieves the active dedicated assignment for a user +// GetActiveDedicatedAssignment retrieves the active dedicated assignment for an organization func (r *workerRepository) GetActiveDedicatedAssignment(ctx context.Context, userID uuid.UUID) (*models.DedicatedWorkerAssignment, error) { query := ` - SELECT id, worker_id, user_id, subscription_id, assigned_at, released_at + SELECT id, worker_id, organization_id, subscription_id, assigned_at, released_at FROM dedicated_worker_assignments - WHERE user_id = $1 AND released_at IS NULL + WHERE organization_id = $1 AND released_at IS NULL ` var a models.DedicatedWorkerAssignment err := r.db.QueryRow(ctx, query, userID).Scan( - &a.ID, &a.WorkerID, &a.UserID, &a.SubscriptionID, &a.AssignedAt, &a.ReleasedAt, + &a.ID, &a.WorkerID, &a.OrganizationID, &a.SubscriptionID, &a.AssignedAt, &a.ReleasedAt, ) if err == pgx.ErrNoRows { return nil, nil @@ -342,17 +342,17 @@ func (r *workerRepository) GetActiveDedicatedAssignment(ctx context.Context, use return &a, nil } -// GetDedicatedWorkerByUserID retrieves the dedicated worker assigned to a user -func (r *workerRepository) GetDedicatedWorkerByUserID(ctx context.Context, userID uuid.UUID) (*models.Worker, error) { +// GetDedicatedWorkerByOrgID retrieves the dedicated worker assigned to an organization +func (r *workerRepository) GetDedicatedWorkerByOrgID(ctx context.Context, orgID uuid.UUID) (*models.Worker, error) { query := ` SELECT w.id, w.ip_addr, w.active, w.free_tier, w.worker_type, w.account_count, w.created_at, w.updated_at FROM workers w JOIN dedicated_worker_assignments dwa ON w.id = dwa.worker_id - WHERE dwa.user_id = $1 AND dwa.released_at IS NULL + WHERE dwa.organization_id = $1 AND dwa.released_at IS NULL ` var w models.Worker - err := r.db.QueryRow(ctx, query, userID).Scan( + err := r.db.QueryRow(ctx, query, orgID).Scan( &w.ID, &w.IPAddr, &w.Active, &w.FreeTier, &w.WorkerType, &w.AccountCount, &w.CreatedAt, &w.UpdatedAt, ) @@ -370,7 +370,7 @@ func (r *workerRepository) ReleaseDedicatedAssignment(ctx context.Context, userI query := ` UPDATE dedicated_worker_assignments SET released_at = $1 - WHERE user_id = $2 AND released_at IS NULL + WHERE organization_id = $2 AND released_at IS NULL ` _, err := r.db.Exec(ctx, query, time.Now(), userID) diff --git a/internal/seed/workers.go b/internal/seed/workers.go index bf1c2ac14..a889a5b23 100644 --- a/internal/seed/workers.go +++ b/internal/seed/workers.go @@ -53,9 +53,9 @@ func seedWorkerAssignments(ctx context.Context, pool *pgxpool.Pool, _ *Result) e // re-running the seed updates rather than duplicates. assignmentID := uuid.MustParse("00000000-0000-0000-0000-000000000280") _, err := pool.Exec(ctx, ` - INSERT INTO dedicated_worker_assignments (id, worker_id, user_id, subscription_id, assigned_at) + INSERT INTO dedicated_worker_assignments (id, worker_id, organization_id, subscription_id, assigned_at) VALUES ($1, $2, $3, $4, NOW()) ON CONFLICT (id) DO NOTHING - `, assignmentID, WorkerDedicatedID, UserOwnerID, SubAcmeID) + `, assignmentID, WorkerDedicatedID, OrgAcmeID, SubAcmeID) return err }