|
1 | 1 | import { QueryExecutor } from '../queryExecutor' |
2 | 2 | import { prepareSelectColumns } from '../utils' |
3 | 3 |
|
4 | | -import { IDbProjectCatalog, IDbProjectCatalogCreate, IDbProjectCatalogUpdate } from './types' |
| 4 | +import { |
| 5 | + IDbProjectCatalog, |
| 6 | + IDbProjectCatalogCreate, |
| 7 | + IDbProjectCatalogUpdate, |
| 8 | + ProjectCatalogAction, |
| 9 | +} from './types' |
5 | 10 |
|
6 | 11 | const PROJECT_CATALOG_COLUMNS = [ |
7 | 12 | 'id', |
@@ -108,6 +113,79 @@ export async function countProjectCatalog(qx: QueryExecutor): Promise<number> { |
108 | 113 | return parseInt(result.count, 10) |
109 | 114 | } |
110 | 115 |
|
| 116 | +export async function countProjectCatalogByAction( |
| 117 | + qx: QueryExecutor, |
| 118 | + action: ProjectCatalogAction, |
| 119 | +): Promise<number> { |
| 120 | + const result = await qx.selectOne( |
| 121 | + ` |
| 122 | + SELECT COUNT(*) AS count |
| 123 | + FROM "projectCatalog" |
| 124 | + WHERE action = $(action) |
| 125 | + `, |
| 126 | + { action }, |
| 127 | + ) |
| 128 | + return parseInt(result.count, 10) |
| 129 | +} |
| 130 | + |
| 131 | +/** |
| 132 | + * Promotes 'auto' projects to 'evaluate', targeting a soft cap on the queue size. |
| 133 | + * |
| 134 | + * Uses a CTE to: |
| 135 | + * 1. Compute available slots (evaluateLimit − current 'evaluate' count) inline. |
| 136 | + * 2. Lock candidates with FOR UPDATE SKIP LOCKED — prevents double-promotion |
| 137 | + * if two transactions run concurrently (e.g. simultaneous manual triggers). |
| 138 | + * Note: concurrent calls can each compute the same slot count and promote |
| 139 | + * up to evaluateLimit disjoint rows, so the cap is soft — the queue can |
| 140 | + * overshoot by at most evaluateLimit per extra concurrent caller. In practice |
| 141 | + * this doesn't happen: the schedule uses ScheduleOverlapPolicy.SKIP. |
| 142 | + * 3. Re-check action = 'auto' in the outer UPDATE to guard against rows whose |
| 143 | + * state changed after the subquery snapshot (e.g. manual updates). |
| 144 | + * |
| 145 | + * Ordering (configurable via `sourcePriority`): |
| 146 | + * 1. Source priority — earlier position in the array = higher priority (unlisted = lowest) |
| 147 | + * 2. lfCriticalityScore DESC (NULLs last) |
| 148 | + * 3. createdAt ASC (stable tie-breaker) |
| 149 | + * |
| 150 | + * Returns the number of rows actually promoted. |
| 151 | + */ |
| 152 | +export async function promoteProjectsToEvaluate( |
| 153 | + qx: QueryExecutor, |
| 154 | + options: { evaluateLimit: number; sourcePriority: string[] }, |
| 155 | +): Promise<number> { |
| 156 | + const { evaluateLimit, sourcePriority } = options |
| 157 | + |
| 158 | + return qx.result( |
| 159 | + ` |
| 160 | + WITH |
| 161 | + slots AS ( |
| 162 | + SELECT GREATEST(0, $(evaluateLimit) - COUNT(*)) AS available |
| 163 | + FROM "projectCatalog" |
| 164 | + WHERE action = 'evaluate' |
| 165 | + ), |
| 166 | + candidates AS ( |
| 167 | + SELECT pc.id |
| 168 | + FROM "projectCatalog" pc |
| 169 | + CROSS JOIN slots |
| 170 | + WHERE pc.action = 'auto' |
| 171 | + AND slots.available > 0 |
| 172 | + ORDER BY |
| 173 | + COALESCE(ARRAY_POSITION($(sourcePriority)::text[], pc.source), 2147483647) ASC, |
| 174 | + pc."lfCriticalityScore" DESC NULLS LAST, |
| 175 | + pc."createdAt" ASC |
| 176 | + LIMIT (SELECT available FROM slots) |
| 177 | + FOR UPDATE SKIP LOCKED |
| 178 | + ) |
| 179 | + UPDATE "projectCatalog" |
| 180 | + SET action = 'evaluate', "updatedAt" = NOW() |
| 181 | + FROM candidates |
| 182 | + WHERE "projectCatalog".id = candidates.id |
| 183 | + AND "projectCatalog".action = 'auto' |
| 184 | + `, |
| 185 | + { evaluateLimit, sourcePriority }, |
| 186 | + ) |
| 187 | +} |
| 188 | + |
111 | 189 | export async function insertProjectCatalog( |
112 | 190 | qx: QueryExecutor, |
113 | 191 | data: IDbProjectCatalogCreate, |
|
0 commit comments