diff --git a/apps/api/internal/store/acceptance.go b/apps/api/internal/store/acceptance.go index 5e608ed..d2ec4e5 100644 --- a/apps/api/internal/store/acceptance.go +++ b/apps/api/internal/store/acceptance.go @@ -376,7 +376,8 @@ FOR UPDATE`, SystemSettingGatewayTrafficMode).Scan(&value); err != nil { current.WorkerImageDigest != strings.TrimSpace(input.WorkerImageDigest) { return GatewayTrafficMode{}, ErrAcceptanceStateConflict } - if _, err := cancelSafeAcceptanceTasksTx(ctx, tx, current.RunID); err != nil { + cancelled, err := cancelSafeAcceptanceTasksTx(ctx, tx, current.RunID) + if err != nil { return GatewayTrafficMode{}, err } next := GatewayTrafficMode{Mode: "live", Revision: current.Revision + 1} @@ -396,6 +397,9 @@ WHERE id = $1::uuid AND status IN ('running', 'failed', 'passed')`, current.RunI if err := tx.Commit(ctx); err != nil { return GatewayTrafficMode{}, err } + if cancelled > 0 { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } next.UpdatedAt = time.Now() return next, nil } @@ -555,11 +559,6 @@ WHERE admission.task_id = task.id )`, runID); err != nil { return 0, err } - if tag.RowsAffected() > 0 { - if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return 0, err - } - } return tag.RowsAffected(), nil } diff --git a/apps/api/internal/store/admission_queue.go b/apps/api/internal/store/admission_queue.go index 18a080c..ad34767 100644 --- a/apps/api/internal/store/admission_queue.go +++ b/apps/api/internal/store/admission_queue.go @@ -16,6 +16,7 @@ import ( const ( admissionWaiterLeaseTTL = 15 * time.Second admissionNotifyChannel = "gateway_task_admission" + admissionNotifyTimeout = 2 * time.Second executionSlotWaitMax = 24 * time.Hour ) @@ -216,12 +217,10 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { if err := expireTaskAdmissionTx(ctx, tx, input.TaskID); err != nil { return TaskAdmission{}, err } - if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return TaskAdmission{}, err - } if err := tx.Commit(ctx); err != nil { return TaskAdmission{}, err } + s.notifyTaskAdmissionBestEffort(ctx, "*") return TaskAdmission{}, &QueueTimeoutError{TaskID: input.TaskID} } if !found { @@ -243,12 +242,10 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { return TaskAdmission{}, err } } - if err := notifyTaskAdmissionTx(ctx, tx, input.TaskID); err != nil { - return TaskAdmission{}, err - } if err := tx.Commit(ctx); err != nil { return TaskAdmission{}, err } + s.notifyTaskAdmissionBestEffort(ctx, input.TaskID) return admission, nil } @@ -283,12 +280,14 @@ func (s *Store) tryTaskAdmissionOnce( if err := tx.Commit(ctx); err != nil { return TaskAdmissionResult{}, err } + s.notifyTaskAdmissionBestEffort(ctx, outcome.NotifyTaskID) return outcome.Result, outcome.PostCommitErr } type taskAdmissionTxOutcome struct { Result TaskAdmissionResult PostCommitErr error + NotifyTaskID string } func (s *Store) tryTaskAdmissionTx( @@ -363,11 +362,9 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { if err := expireTaskAdmissionTx(ctx, tx, input.TaskID); err != nil { return taskAdmissionTxOutcome{}, err } - if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return taskAdmissionTxOutcome{}, err - } return taskAdmissionTxOutcome{ PostCommitErr: &QueueTimeoutError{TaskID: input.TaskID}, + NotifyTaskID: "*", }, nil } @@ -442,9 +439,6 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { // Continue the FIFO chain after one task is admitted. Every API process // only probes the global head it owns, so free capacity is filled without // broadcasting an advisory-lock attempt to every waiter. - if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return taskAdmissionTxOutcome{}, err - } if onAdmitted != nil { if err := onAdmitted(tx); err != nil { return taskAdmissionTxOutcome{}, err @@ -457,6 +451,7 @@ WHERE id = $1::uuid`, input.TaskID).Scan(&taskActive); err != nil { NewlyAdmitted: true, Leases: leases, }, + NotifyTaskID: "*", }, nil } @@ -487,6 +482,7 @@ func (s *Store) TryTaskAdmissionBatchWithAdmittedHook( defer rollbackTransaction(tx) outcomes := make([]TaskAdmissionBatchOutcome, 0, len(inputs)) + notify := false for _, input := range inputs { hook := func(tx pgx.Tx) error { if onAdmitted == nil { @@ -510,6 +506,7 @@ func (s *Store) TryTaskAdmissionBatchWithAdmittedHook( Result: outcome.Result, Err: outcome.PostCommitErr, }) + notify = notify || outcome.NotifyTaskID != "" if outcome.PostCommitErr == nil && !outcome.Result.Admitted { break } @@ -517,6 +514,9 @@ func (s *Store) TryTaskAdmissionBatchWithAdmittedHook( if err := tx.Commit(ctx); err != nil { return nil, err } + if notify { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } return outcomes, nil }) } @@ -649,12 +649,10 @@ RETURNING task_id::text, platform_id::text, platform_model_id::text, if err != nil { return TaskAdmission{}, err } - if err := notifyTaskAdmissionTx(ctx, tx, input.TaskID); err != nil { - return TaskAdmission{}, err - } if err := tx.Commit(ctx); err != nil { return TaskAdmission{}, err } + s.notifyTaskAdmissionBestEffort(ctx, input.TaskID) return migrated, nil } @@ -781,10 +779,11 @@ WHERE task_id = $1::uuid if _, err := tx.Exec(ctx, `DELETE FROM gateway_task_admissions WHERE task_id = $1::uuid`, taskID); err != nil { return err } - if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { + if err := tx.Commit(ctx); err != nil { return err } - return tx.Commit(ctx) + s.notifyTaskAdmissionBestEffort(ctx, "*") + return nil } func (s *Store) ActiveTaskAdmissionLeases(ctx context.Context, taskID string) ([]ConcurrencyLease, error) { @@ -1026,14 +1025,12 @@ WHERE id = $1::uuid } reaped++ } - if reaped > 0 { - if err := notifyTaskAdmissionTx(ctx, tx, "*"); err != nil { - return TaskAdmissionReapResult{}, err - } - } if err := tx.Commit(ctx); err != nil { return TaskAdmissionReapResult{}, err } + if reaped > 0 { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } return result, nil } @@ -1413,7 +1410,34 @@ WHERE id = $1::uuid AND status = 'queued'`, taskID) return err } -func notifyTaskAdmissionTx(ctx context.Context, tx pgx.Tx, taskID string) error { - _, err := tx.Exec(ctx, `SELECT pg_notify($1, $2)`, admissionNotifyChannel, taskID) - return err +// notifyTaskAdmission publishes a non-durable wake-up hint only after the +// corresponding state transaction has committed. Keeping NOTIFY out of +// synchronous replication transactions prevents its database object lock from +// serializing unrelated admission, execution, and release commits. +func (s *Store) notifyTaskAdmission(ctx context.Context, taskID string) error { + taskID = strings.TrimSpace(taskID) + if taskID == "" { + return nil + } + tx, err := s.pool.Begin(ctx) + if err != nil { + return err + } + defer rollbackTransaction(tx) + if _, err := tx.Exec(ctx, `SET LOCAL synchronous_commit = off`); err != nil { + return err + } + if _, err := tx.Exec(ctx, `SELECT pg_notify($1, $2)`, admissionNotifyChannel, taskID); err != nil { + return err + } + return tx.Commit(ctx) +} + +func (s *Store) notifyTaskAdmissionBestEffort(ctx context.Context, taskID string) { + if strings.TrimSpace(taskID) == "" { + return + } + notifyCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), admissionNotifyTimeout) + defer cancel() + _ = s.notifyTaskAdmission(notifyCtx, taskID) } diff --git a/apps/api/internal/store/admission_registration_lock_integration_test.go b/apps/api/internal/store/admission_registration_lock_integration_test.go index ce7cc96..5012a35 100644 --- a/apps/api/internal/store/admission_registration_lock_integration_test.go +++ b/apps/api/internal/store/admission_registration_lock_integration_test.go @@ -113,4 +113,53 @@ RETURNING id::text`, platform.ID, modelName).Scan(&platformModelID); err != nil if registrationErr != nil { t.Fatalf("queue registration blocked behind execution scope lock: %v", registrationErr) } + + // Simulate a legacy synchronous transaction that is still holding the + // database object lock used by NOTIFY. The admission deletion must commit + // and become visible before the best-effort wake-up hint can acquire that + // lock. + notificationLockTx, err := first.pool.Begin(ctx) + if err != nil { + t.Fatalf("begin notification lock transaction: %v", err) + } + if _, err := notificationLockTx.Exec(ctx, `SELECT pg_notify($1, '*')`, admissionNotifyChannel); err != nil { + _ = notificationLockTx.Rollback(ctx) + t.Fatalf("hold notification object lock: %v", err) + } + deleteResult := make(chan error, 1) + go func() { + deleteResult <- second.DeleteTaskAdmission(ctx, task.ID) + }() + deleteVisible := false + visibilityDeadline := time.Now().Add(time.Second) + for time.Now().Before(visibilityDeadline) { + var admissionCount int + if err := first.pool.QueryRow(ctx, ` +SELECT count(*) +FROM gateway_task_admissions +WHERE task_id = $1::uuid`, task.ID).Scan(&admissionCount); err != nil { + _ = notificationLockTx.Rollback(ctx) + t.Fatalf("read admission visibility: %v", err) + } + if admissionCount == 0 { + deleteVisible = true + break + } + time.Sleep(20 * time.Millisecond) + } + if !deleteVisible { + _ = notificationLockTx.Rollback(ctx) + t.Fatal("admission state commit waited behind notification object lock") + } + if err := notificationLockTx.Rollback(ctx); err != nil { + t.Fatalf("release notification object lock: %v", err) + } + select { + case err := <-deleteResult: + if err != nil { + t.Fatalf("delete admission after notification lock release: %v", err) + } + case <-time.After(3 * time.Second): + t.Fatal("delete admission did not finish after notification lock release") + } } diff --git a/apps/api/internal/store/tasks_runtime.go b/apps/api/internal/store/tasks_runtime.go index 6e048c1..47c9df9 100644 --- a/apps/api/internal/store/tasks_runtime.go +++ b/apps/api/internal/store/tasks_runtime.go @@ -715,11 +715,14 @@ WHERE task_id = $1::uuid if _, err := tx.Exec(ctx, `DELETE FROM gateway_task_admissions WHERE task_id = $1::uuid`, taskID); err != nil { return err } - return notifyTaskAdmissionTx(ctx, tx, "*") + return nil }) if err != nil { return GatewayTask{}, false, err } + if changed { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } return task, changed, nil } @@ -877,11 +880,14 @@ WHERE task_id = $1::uuid if _, err := tx.Exec(ctx, `DELETE FROM gateway_task_admissions WHERE task_id = $1::uuid`, taskID); err != nil { return err } - return notifyTaskAdmissionTx(ctx, tx, "*") + return nil }) if err != nil { return GatewayTask{}, false, err } + if changed { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } return task, changed, nil } @@ -963,11 +969,14 @@ WHERE task_id = $1::uuid if _, err := tx.Exec(ctx, `DELETE FROM gateway_task_admissions WHERE task_id = $1::uuid`, taskID); err != nil { return err } - return notifyTaskAdmissionTx(ctx, tx, "*") + return nil }) if err != nil { return GatewayTask{}, false, err } + if changed { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } return task, changed, nil } diff --git a/apps/api/internal/store/worker_registry.go b/apps/api/internal/store/worker_registry.go index f075b87..f33cb9d 100644 --- a/apps/api/internal/store/worker_registry.go +++ b/apps/api/internal/store/worker_registry.go @@ -286,14 +286,12 @@ SELECT count(*)::bigint FROM deleted`, staleAfter.String(), limit).Scan(&yielded); err != nil { return 0, err } - if yielded > 0 { - if _, err := tx.Exec(ctx, `SELECT pg_notify('gateway_task_admission', '*')`); err != nil { - return 0, err - } - } if err := tx.Commit(ctx); err != nil { return 0, err } + if yielded > 0 { + s.notifyTaskAdmissionBestEffort(ctx, "*") + } return yielded, nil }