587 lines
21 KiB
Go
587 lines
21 KiB
Go
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": "hire_candidate",
|
|
"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.
|
|
//
|
|
// Three ways of naming the application this placement belongs to, in order of
|
|
// precedence: an id the caller already has, an `application` payload to find or
|
|
// file one from, and neither — a worker placed straight from the talent pool,
|
|
// who has no application and is not given an invented one.
|
|
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"`
|
|
|
|
// Application is the application to attach this worker to when no
|
|
// ApplicationID is supplied: the same fields POST /job-applications takes.
|
|
//
|
|
// It exists because an application is what puts a person in the pipeline
|
|
// for a role, and every downstream step keys on it — the candidate record
|
|
// is addressed by it, and an AI interview takes one as its subject. A
|
|
// worker assigned without one is unreachable: nothing to open, and nobody
|
|
// to interview. The frontend was creating it in a second, untransacted
|
|
// request; this carries it into the same transaction as the assignment.
|
|
Application *domain.Record `json:"application"`
|
|
}
|
|
|
|
// 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.
|
|
//
|
|
// Per worker the flow is: settle the application (see settleApplication —
|
|
// patch the one named, find-or-file the one described, or neither), insert the
|
|
// assignment linked to whatever that produced, and write the audit entry. The
|
|
// application comes first because the assignment references it.
|
|
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)
|
|
// The application is filed through the service rather than the
|
|
// repository so that a payload the caller sent is validated exactly as
|
|
// POST /job-applications would validate it — unknown fields rejected,
|
|
// enums checked, required fields demanded — instead of reaching SQL and
|
|
// coming back as a constraint violation.
|
|
appSvc := New(apps, tx)
|
|
activityRepo := repo.New(activity, tx)
|
|
|
|
for i, w := range req.Workers {
|
|
// The application is settled BEFORE the assignment is written,
|
|
// because the assignment carries the reference to it and a row
|
|
// cannot point at one that does not exist yet. An empty id means
|
|
// this worker legitimately has no application.
|
|
applicationID, err := s.settleApplication(ctx, ident, appRepo, appSvc, postingID, w)
|
|
if err != nil {
|
|
return annotate(err, i)
|
|
}
|
|
|
|
record := domain.Record{
|
|
"job_posting_id": postingID,
|
|
"worker_email": w.WorkerEmail,
|
|
"worker_name": w.WorkerName,
|
|
"starts_at": w.StartsAt,
|
|
"status": "active",
|
|
// `owliver` is the column's own default (000001:511) and what
|
|
// the frontend sends. Substituting `manual` here made the API
|
|
// and the schema disagree about what an unspecified source
|
|
// means, so a row written through this endpoint and a row
|
|
// written through POST /assignments recorded different origins
|
|
// for the same action.
|
|
"source": orDefault(w.Source, "owliver"),
|
|
}
|
|
if w.EndsAt != nil {
|
|
record["ends_at"] = *w.EndsAt
|
|
}
|
|
if applicationID != "" {
|
|
record["application_id"] = 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)
|
|
|
|
entry := domain.Record{
|
|
"event_type": "assign_employee",
|
|
"details": fmt.Sprintf("%s was assigned", orDefault(w.WorkerName, w.WorkerEmail)),
|
|
"position_id": postingID,
|
|
"worker_email": w.WorkerEmail,
|
|
}
|
|
if applicationID != "" {
|
|
entry["application_id"] = 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
|
|
}
|
|
|
|
// settleApplication resolves the application an assignment belongs to and
|
|
// returns its id, or "" when this worker has none.
|
|
//
|
|
// Three cases, and the middle one is why this exists:
|
|
//
|
|
// - an id was supplied — move that application to `assigned`. Unchanged
|
|
// behaviour, and it still wins over any payload, because an id is a
|
|
// decision the caller has already made.
|
|
// - an `application` payload was supplied — find the application on this
|
|
// posting for this email, and move it to `assigned` if there is one or file
|
|
// it if there is not. (job_posting_id, email) is UNIQUE
|
|
// (job_applications_posting_email_key, 000001), so there is at most one to
|
|
// find and the insert cannot produce a second.
|
|
// - neither — nothing. A worker placed straight from the talent pool has no
|
|
// application, and inventing one for them would be worse than leaving the
|
|
// link absent.
|
|
//
|
|
// Everything here runs inside the caller's transaction, which is what makes the
|
|
// find-or-create safe: a concurrent assign of the same person to the same
|
|
// posting blocks on the unique index and is then reported as a conflict, rather
|
|
// than racing past the lookup and writing a duplicate.
|
|
func (s *WorkflowService) settleApplication(ctx context.Context, ident authctx.Identity,
|
|
appRepo *repo.Repo, appSvc *Service, postingID string, w AssignWorker) (string, error) {
|
|
|
|
if w.ApplicationID != nil {
|
|
updated, err := appRepo.Update(ctx, ident, *w.ApplicationID,
|
|
domain.Record{"status": "assigned"})
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if updated == nil {
|
|
return "", domain.NotFound("JobApplication", *w.ApplicationID)
|
|
}
|
|
return *w.ApplicationID, nil
|
|
}
|
|
if w.Application == nil {
|
|
return "", nil
|
|
}
|
|
|
|
existing, err := findApplication(ctx, ident, appRepo, postingID, w.WorkerEmail)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if existing != nil {
|
|
id, _ := existing["id"].(string)
|
|
updated, err := appRepo.Update(ctx, ident, id, domain.Record{"status": "assigned"})
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if updated == nil {
|
|
return "", domain.NotFound("JobApplication", id)
|
|
}
|
|
return id, nil
|
|
}
|
|
|
|
// The posting and the person come from the assignment, not from the
|
|
// payload. A body naming a different posting or a different email would
|
|
// file an application about somebody other than the worker being placed,
|
|
// and the lookup above would never find it again.
|
|
record := domain.Record{}
|
|
for k, v := range *w.Application {
|
|
record[k] = v
|
|
}
|
|
record["job_posting_id"] = postingID
|
|
record["email"] = w.WorkerEmail
|
|
if _, ok := record["applicant_name"]; !ok {
|
|
record["applicant_name"] = w.WorkerName
|
|
}
|
|
if _, ok := record["status"]; !ok {
|
|
record["status"] = "assigned"
|
|
}
|
|
|
|
created, err := appSvc.Create(ctx, ident, record)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
id, _ := created["id"].(string)
|
|
return id, nil
|
|
}
|
|
|
|
// findApplication returns this posting's application for this email, or nil.
|
|
//
|
|
// The read goes through the repository so it carries the same organization
|
|
// scope and ownership predicate every other read does — an application in
|
|
// another tenant is not found, rather than found and then refused. The email
|
|
// column is citext, so the comparison is case-insensitive: the same equality
|
|
// the frontend performed with toLowerCase before it had this endpoint.
|
|
func findApplication(ctx context.Context, ident authctx.Identity, appRepo *repo.Repo,
|
|
postingID, email string) (domain.Record, error) {
|
|
|
|
res := appRepo.Resource()
|
|
postingCol, ok := res.Column("job_posting_id")
|
|
if !ok {
|
|
return nil, domain.Internal(errors.New("service: job_applications has no job_posting_id column"))
|
|
}
|
|
emailCol, ok := res.Column("email")
|
|
if !ok {
|
|
return nil, domain.Internal(errors.New("service: job_applications has no email column"))
|
|
}
|
|
|
|
page, err := appRepo.List(ctx, ident, domain.ListParams{
|
|
Limit: 1,
|
|
Filters: []domain.Filter{
|
|
{Column: postingCol, Values: []string{postingID}},
|
|
{Column: emailCol, Values: []string{email}},
|
|
},
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if len(page.Records) == 0 {
|
|
return nil, nil
|
|
}
|
|
return page.Records[0], 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)
|
|
}
|