package service // Multi-record writes, in one transaction each. // // api-contract.md §12.1 lists four flows that the frontend performs as a // sequence of independent HTTP calls, with no transaction and no rollback: // hiring a candidate, assigning workers to a posting, screening a whole list, // and submitting a challenge. A failure halfway through leaves the database // inconsistent — an application marked hired with no staff row, or an // assignment with an application still showing as merely shortlisted. // // This file collapses the first two into one endpoint and one transaction // each, which is what §12.1 says Phase 3 should do. The other two are not here: // `useScreenAllCandidates` is n independent PATCHes that are individually // meaningful (a partial screen is not a corrupt state), and `useSubmitChallenge` // writes evidence the talent user owns, which needs the talent-scoped predicate // and is a different shape of problem. // // WHY THE REPOSITORY IS REUSED RATHER THAN HAND-WRITTEN SQL // // Every write below goes through repo.Repo, built over the transaction rather // than the pool. That is deliberate: the repository is where org_id is forced // from the session, where server-derived columns are filled, where values are // bound and cast to their declared types, and where a constraint violation is // translated into the contract's error codes. Writing raw SQL here would mean // re-deriving all four, and getting one of them subtly wrong. import ( "context" "errors" "fmt" "time" "github.com/jackc/pgx/v5" "github.com/krow/krow-backend/go-api/internal/authctx" "github.com/krow/krow-backend/go-api/internal/domain" "github.com/krow/krow-backend/go-api/internal/repo" ) // maxAssignmentBatch bounds one assign call. // // `useAssignWorkers` assigns a selection from the talent pool, which is a // human-sized list. The cap exists so a single request cannot hold a // transaction open across thousands of inserts, blocking every other write to // these tables for the duration. const maxAssignmentBatch = 200 // TxBeginner is the part of the pool this file needs. An interface rather than // *pgxpool.Pool so the service stays testable and consistent with repo.Querier. type TxBeginner interface { Begin(ctx context.Context) (pgx.Tx, error) } // WorkflowService owns the multi-record flows. type WorkflowService struct { db TxBeginner // now is injectable so tests can assert on generated dates without // depending on the day they run. now func() time.Time } // NewWorkflows builds the workflow service over a pool. func NewWorkflows(db TxBeginner) *WorkflowService { return &WorkflowService{db: db, now: time.Now} } // WithClock replaces the clock. For tests. func (s *WorkflowService) WithClock(now func() time.Time) *WorkflowService { if now != nil { s.now = now } return s } // inTx runs fn inside a transaction, rolling back on any error. // // The rollback is deferred rather than called on each error path: an early // return, a panic in a callee, and an explicit failure all have to undo the // work, and only a deferred rollback covers the second. Rolling back an // already-committed transaction is a no-op in pgx, so the deferred call is safe // on the success path too. func (s *WorkflowService) inTx(ctx context.Context, fn func(tx pgx.Tx) error) error { tx, err := s.db.Begin(ctx) if err != nil { return fmt.Errorf("begin transaction: %w", err) } defer func() { _ = tx.Rollback(ctx) }() if err := fn(tx); err != nil { return err } return tx.Commit(ctx) } func resourceByPath(path string) (*domain.Resource, error) { res, ok := domain.ResourceByPath[path] if !ok { return nil, domain.Internal(fmt.Errorf("service: resource %q is not registered", path)) } return res, nil } /* ── Hire ───────────────────────────────────────────────────────────────── */ // HireResult is what a completed hire returns: both records the flow touched, // so the caller does not need a follow-up read to render the outcome. type HireResult struct { Application domain.Record `json:"application"` Staff domain.Record `json:"staff"` } // Hire moves an application to `hired` and creates the staff record, atomically. // // Replaces the two-call sequence at krowHooks.js:302-303. The failure that // motivated it: the PATCH succeeds, the POST fails, and the candidate is now // hired with no employment record and no way for the UI to notice. // // Fields the caller may supply are the ones a hiring form collects — hire_date, // role, phone, profile_tier, status, reviewer_name. Everything else is carried // across from the application, because it is already the truth about this // person and retyping it is how the two records drift apart. func (s *WorkflowService) Hire(ctx context.Context, ident authctx.Identity, applicationID string, body domain.Record) (*HireResult, error) { if !isUUID(applicationID) { return nil, domain.NotFound("JobApplication", applicationID) } apps, err := resourceByPath("job-applications") if err != nil { return nil, err } staffRes, err := resourceByPath("staff") if err != nil { return nil, err } activity, err := resourceByPath("user-activity") if err != nil { return nil, err } var out HireResult err = s.inTx(ctx, func(tx pgx.Tx) error { appRepo := repo.New(apps, tx) // Read inside the transaction. Reading outside it would leave a window // where two concurrent hires both see a not-yet-hired application. app, err := appRepo.Get(ctx, ident, applicationID) if err != nil { return err } if app == nil { return domain.NotFound("JobApplication", applicationID) } // Hiring twice is a conflict, not an idempotent repeat: the second call // would create a second employment record for one application. 409 // rather than 200 because the caller asked for something that cannot be // done, and silently returning the first hire would hide a double // submission rather than report it. if status, _ := app["status"].(string); status == "hired" { return domain.Conflict("this application has already been hired") } email, _ := app["email"].(string) name, _ := app["applicant_name"].(string) if email == "" || name == "" { return domain.Validation( "this application cannot be hired: it has no applicant name or email", map[string]string{"application": "incomplete"}) } staffRecord := domain.Record{ "application_id": applicationID, "job_posting_id": app["job_posting_id"], "worker_profile_id": app["worker_profile_id"], "name": name, "email": email, "phone": pick(body, "phone", app["phone"]), "role": pick(body, "role", app["job_title"]), "hire_date": pick(body, "hire_date", s.now().UTC().Format("2006-01-02")), "ai_score": app["ai_score"], "profile_tier": pick(body, "profile_tier", "Beginner"), "status": pick(body, "status", "onboarding"), "client_rating": app["client_rating"], "reviewer_name": pick(body, "reviewer_name", ""), } // A null worker_profile_id is legitimate — not every applicant has a // profile — but the column list must not carry an explicit nil for a // NOT NULL column, so drop the keys the application had nothing for. dropNil(staffRecord, "worker_profile_id", "job_posting_id") hired, err := appRepo.Update(ctx, ident, applicationID, domain.Record{"status": "hired"}) if err != nil { return err } if hired == nil { // Unreachable: Get above found it under the same predicate. return domain.NotFound("JobApplication", applicationID) } created, err := repo.New(staffRes, tx).Insert(ctx, ident, staffRecord) if err != nil { return err } // The audit entry is part of the transaction on purpose: a hire that // happened without a log entry, or a log entry for a hire that rolled // back, are both worse than neither. if _, err := repo.New(activity, tx).Insert(ctx, ident, domain.Record{ "event_type": "candidate_hired", "details": fmt.Sprintf("%s was hired", name), "application_id": applicationID, "position_id": app["job_posting_id"], "worker_email": email, }); err != nil { return err } out.Application, out.Staff = hired, created return nil }) if err != nil { return nil, err } return &out, nil } /* ── Assign ─────────────────────────────────────────────────────────────── */ // AssignWorker is one worker in an assign request. type AssignWorker struct { WorkerEmail string `json:"worker_email"` WorkerName string `json:"worker_name"` StartsAt string `json:"starts_at"` EndsAt *string `json:"ends_at"` ApplicationID *string `json:"application_id"` WorkerProfileID *string `json:"worker_profile_id"` MatchScore *int `json:"match_score"` Source string `json:"source"` } // AssignRequest is the body of POST /job-postings/{id}/assignments. type AssignRequest struct { Workers []AssignWorker `json:"workers"` } // AssignResult reports what the batch created. type AssignResult struct { Assignments []domain.Record `json:"assignments"` Count int `json:"count"` } // Assign places workers on a posting, atomically. // // Replaces the 3n sequential round-trips at krowHooks.js:421/449/466 with one // request and one transaction. All-or-nothing across the whole batch: assigning // six workers and having the fourth fail should not leave three assigned, three // not, and the caller unsure which. func (s *WorkflowService) Assign(ctx context.Context, ident authctx.Identity, postingID string, req AssignRequest) (*AssignResult, error) { if !isUUID(postingID) { return nil, domain.NotFound("JobPosting", postingID) } if len(req.Workers) == 0 { return nil, domain.Validation("at least one worker is required", map[string]string{"workers": "must not be empty"}) } if len(req.Workers) > maxAssignmentBatch { return nil, domain.Validation( fmt.Sprintf("at most %d workers can be assigned in one request", maxAssignmentBatch), map[string]string{"workers": "too many"}) } // Validate the whole batch before opening a transaction. A malformed // request should never have caused a BEGIN. details := map[string]string{} for i, w := range req.Workers { if w.WorkerEmail == "" { details[fmt.Sprintf("workers[%d].worker_email", i)] = "required" } if w.StartsAt == "" { details[fmt.Sprintf("workers[%d].starts_at", i)] = "required" } if w.ApplicationID != nil && !isUUID(*w.ApplicationID) { details[fmt.Sprintf("workers[%d].application_id", i)] = "must be a uuid" } if w.WorkerProfileID != nil && !isUUID(*w.WorkerProfileID) { details[fmt.Sprintf("workers[%d].worker_profile_id", i)] = "must be a uuid" } } if len(details) > 0 { return nil, domain.Validation("assignment payload is not valid", details) } postings, err := resourceByPath("job-postings") if err != nil { return nil, err } assignments, err := resourceByPath("assignments") if err != nil { return nil, err } apps, err := resourceByPath("job-applications") if err != nil { return nil, err } activity, err := resourceByPath("user-activity") if err != nil { return nil, err } out := &AssignResult{Assignments: []domain.Record{}} err = s.inTx(ctx, func(tx pgx.Tx) error { posting, err := repo.New(postings, tx).Get(ctx, ident, postingID) if err != nil { return err } if posting == nil { return domain.NotFound("JobPosting", postingID) } assignRepo := repo.New(assignments, tx) appRepo := repo.New(apps, tx) activityRepo := repo.New(activity, tx) for i, w := range req.Workers { record := domain.Record{ "job_posting_id": postingID, "worker_email": w.WorkerEmail, "worker_name": w.WorkerName, "starts_at": w.StartsAt, "status": "active", "source": orDefault(w.Source, "manual"), } if w.EndsAt != nil { record["ends_at"] = *w.EndsAt } if w.ApplicationID != nil { record["application_id"] = *w.ApplicationID } if w.WorkerProfileID != nil { record["worker_profile_id"] = *w.WorkerProfileID } if w.MatchScore != nil { record["match_score"] = *w.MatchScore } created, err := assignRepo.Insert(ctx, ident, record) if err != nil { return annotate(err, i) } out.Assignments = append(out.Assignments, created) // The application moves to `assigned` only when one was named. A // worker can be placed without having applied — that is what the // talent pool is for — and inventing an application for them would // be worse than leaving the link absent. if w.ApplicationID != nil { updated, err := appRepo.Update(ctx, ident, *w.ApplicationID, domain.Record{"status": "assigned"}) if err != nil { return annotate(err, i) } if updated == nil { return domain.NotFound("JobApplication", *w.ApplicationID) } } entry := domain.Record{ "event_type": "worker_assigned", "details": fmt.Sprintf("%s was assigned", orDefault(w.WorkerName, w.WorkerEmail)), "position_id": postingID, "worker_email": w.WorkerEmail, } if w.ApplicationID != nil { entry["application_id"] = *w.ApplicationID } if _, err := activityRepo.Insert(ctx, ident, entry); err != nil { return annotate(err, i) } } return nil }) if err != nil { return nil, err } out.Count = len(out.Assignments) return out, nil } /* ── Helpers ────────────────────────────────────────────────────────────── */ // pick takes the caller's value for a key when they supplied a usable one, and // the fallback otherwise. A present-but-null key means "use the fallback" // rather than "write null", because these columns are NOT NULL. func pick(body domain.Record, key string, fallback any) any { if body == nil { return fallback } v, ok := body[key] if !ok || v == nil { return fallback } if s, isStr := v.(string); isStr && s == "" { return fallback } return v } // dropNil removes keys whose value is nil, so a NOT NULL column is left to its // default instead of being sent an explicit null. func dropNil(rec domain.Record, keys ...string) { for _, k := range keys { if v, ok := rec[k]; ok && v == nil { delete(rec, k) } } } func orDefault(v, fallback string) string { if v == "" { return fallback } return v } // annotate points a validation error at the batch element that produced it. // Without this, "worker_email must not be null" on a batch of forty says // nothing about which one. func annotate(err error, index int) error { var apiErr *domain.Error if !errors.As(err, &apiErr) || apiErr.Details == nil { return err } moved := make(map[string]string, len(apiErr.Details)) for k, v := range apiErr.Details { moved[fmt.Sprintf("workers[%d].%s", index, k)] = v } return domain.Validation(apiErr.Message, moved) }