mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-06 16:02:07 +00:00
feat: re-key dedicated_worker_assignments to organizations (migration 000055) - the assignment service has always keyed by org id but the column carried a users FK, so every runtime bind failed with an FK violation and dedicated-plan orgs could never get a worker; renames the column, remaps seed rows, and updates the repository, admin filter, and convert-to-dedicated endpoint to org semantics
This commit is contained in:
@@ -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)),
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
+23
@@ -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;
|
||||
@@ -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;
|
||||
@@ -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"`
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user