|
| 1 | +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. |
| 2 | + |
| 3 | +/** |
| 4 | + * [#5103 / #5190] The record-existence cleanup primitives. |
| 5 | + * |
| 6 | + * Two tables in this package store "(`object_name`, `record_id`) → some |
| 7 | + * access": `sys_record_share` (principal-based grants) and `sys_share_link` |
| 8 | + * (capability tokens). Both have the SAME invariant — **record gone ⇒ the row |
| 9 | + * cannot describe any access at all** — and therefore the same two operations: |
| 10 | + * |
| 11 | + * - a set-based revoke keyed on the ids a delete just removed, and |
| 12 | + * - a sweep that asks, per row, whether its record still exists. |
| 13 | + * |
| 14 | + * #5103 built both for `sys_record_share`. #5190 needs them for |
| 15 | + * `sys_share_link`, and a second copy would be the fork this module exists to |
| 16 | + * prevent: one chunk size, one keyset walk, one "a failed probe deletes |
| 17 | + * NOTHING" rule, one truncation report. The tables' owning services keep their |
| 18 | + * own public methods (`SharingService.sweepOrphanedRecordShares`, |
| 19 | + * `ShareLinkService.sweepOrphanedShareLinks`) — `sys_share_link` is |
| 20 | + * `managedBy: 'engine-owned'` and its writes flow through `IShareLinkService`, |
| 21 | + * so ownership stays where the object declares it; only the mechanism is |
| 22 | + * shared. |
| 23 | + * |
| 24 | + * Nothing here knows what a share or a link MEANS. It knows a table name, an |
| 25 | + * `(object_name, record_id)` pair per row, and that deleting on an unanswered |
| 26 | + * question is the one thing it must never do. |
| 27 | + */ |
| 28 | + |
| 29 | +import { keysetWalk } from '@objectstack/types'; |
| 30 | + |
| 31 | +/** System-elevated context for the plugin's own queries / mutations. */ |
| 32 | +const SYSTEM_CTX = { isSystem: true, positions: [], permissions: [] } as const; |
| 33 | + |
| 34 | +/** The slice of the engine these primitives need. */ |
| 35 | +export interface OrphanCleanupEngine { |
| 36 | + find(object: string, options?: any): Promise<any[]>; |
| 37 | + delete(object: string, options?: any): Promise<any>; |
| 38 | +} |
| 39 | + |
| 40 | +/** |
| 41 | + * [#5103] Ids per `$in`. Mirrors the chunk |
| 42 | + * `SharingRuleService.revokeRuleGrantsForRecords` already uses: a single |
| 43 | + * statement binding a thousand parameters is a portability trap (SQLite's |
| 44 | + * default `SQLITE_MAX_VARIABLE_NUMBER` is 999 on older builds), and one number |
| 45 | + * for every revoke path keeps them from drifting. |
| 46 | + */ |
| 47 | +export const RECORD_SCOPED_DELETE_CHUNK = 200; |
| 48 | + |
| 49 | +/** [#5103] Rows read per page by an orphan sweep. */ |
| 50 | +export const ORPHAN_SWEEP_PAGE_SIZE = 500; |
| 51 | + |
| 52 | +/** |
| 53 | + * [#5103] Rows one sweep will scan before stopping and reporting truncation. |
| 54 | + * The sweep runs on every boot, so it must cost a bounded amount on a table |
| 55 | + * that only grows; the next boot resumes from the start and the rows it did not |
| 56 | + * reach stay reachable by the object-scoped sweep. A cap is not a failure — but |
| 57 | + * an unreported cap turns a partial scan into a false "nothing to clean", which |
| 58 | + * is why {@link OrphanShareSweepResult} carries it. |
| 59 | + */ |
| 60 | +export const ORPHAN_SWEEP_MAX_ROWS = 50_000; |
| 61 | + |
| 62 | +/** [#5103] Options for a record-existence orphan sweep. */ |
| 63 | +export interface OrphanShareSweepOptions { |
| 64 | + /** Restrict the sweep to one object. Default: every object with rows. */ |
| 65 | + object?: string; |
| 66 | + /** Rows per page. Default {@link ORPHAN_SWEEP_PAGE_SIZE}. */ |
| 67 | + batchSize?: number; |
| 68 | + /** Stop after scanning this many rows. Default {@link ORPHAN_SWEEP_MAX_ROWS}. */ |
| 69 | + max?: number; |
| 70 | +} |
| 71 | + |
| 72 | +/** [#5103] What one orphan-sweep pass did. */ |
| 73 | +export interface OrphanShareSweepResult { |
| 74 | + /** Rows examined. */ |
| 75 | + scanned: number; |
| 76 | + /** Rows revoked because their record no longer exists. */ |
| 77 | + revoked: number; |
| 78 | + /** |
| 79 | + * Objects whose existence probe could not be run (unregistered object, |
| 80 | + * driver error). Their rows were LEFT ALONE — "could not ask" is not |
| 81 | + * "the record is gone", and only the second one may delete anything. |
| 82 | + */ |
| 83 | + unresolvedObjects: string[]; |
| 84 | + /** True when {@link OrphanShareSweepOptions.max} stopped the scan early. */ |
| 85 | + truncated: boolean; |
| 86 | +} |
| 87 | + |
| 88 | +/** |
| 89 | + * Structural and loose on purpose — it has to accept both owning services' |
| 90 | + * option shapes (`SharingServiceOptions['logger']`, `ShareLinkServiceOptions`) |
| 91 | + * and a bare `{ warn }` stub in a test. |
| 92 | + */ |
| 93 | +interface MinimalLogger { |
| 94 | + info?: Function; |
| 95 | + warn?: Function; |
| 96 | +} |
| 97 | + |
| 98 | +/** How one caller's rows are named in this module's log lines. */ |
| 99 | +export interface OrphanSweepSubject { |
| 100 | + /** Table to sweep, e.g. `sys_record_share`. */ |
| 101 | + table: string; |
| 102 | + /** Noun for log messages: `share` → "orphan share sweep", "share rows". */ |
| 103 | + noun: string; |
| 104 | + /** Issue reference appended to the "revoked N rows" warning. */ |
| 105 | + issue: string; |
| 106 | +} |
| 107 | + |
| 108 | +/** |
| 109 | + * [#5103] Delete every row of `table` that belongs to records just deleted from |
| 110 | + * `object`. |
| 111 | + * |
| 112 | + * Set-based and chunked, so its cost tracks the number of ids, not the number |
| 113 | + * of rows. Returns nothing: counting would need a read the hot delete path |
| 114 | + * should not pay for, and callers that need a count (tests, the sweep) can read |
| 115 | + * the table. |
| 116 | + */ |
| 117 | +export async function deleteRowsForDeletedRecords( |
| 118 | + engine: OrphanCleanupEngine, |
| 119 | + table: string, |
| 120 | + object: string, |
| 121 | + recordIds: readonly string[], |
| 122 | +): Promise<void> { |
| 123 | + if (!table || !object || recordIds.length === 0) return; |
| 124 | + for (let i = 0; i < recordIds.length; i += RECORD_SCOPED_DELETE_CHUNK) { |
| 125 | + const batch = recordIds.slice(i, i + RECORD_SCOPED_DELETE_CHUNK); |
| 126 | + await engine.delete(table, { |
| 127 | + where: { object_name: object, record_id: { $in: batch } }, |
| 128 | + multi: true, |
| 129 | + context: SYSTEM_CTX, |
| 130 | + } as any); |
| 131 | + } |
| 132 | +} |
| 133 | + |
| 134 | +/** |
| 135 | + * [#5103] Which of `recordIds` still exist on `object`. Batched by |
| 136 | + * {@link RECORD_SCOPED_DELETE_CHUNK} so the `$in` never outgrows a driver's |
| 137 | + * bind-parameter limit. Throws on a query failure — the caller MUST treat that |
| 138 | + * as "unknown", never as "none of them exist". |
| 139 | + */ |
| 140 | +export async function findLiveRecordIds( |
| 141 | + engine: OrphanCleanupEngine, |
| 142 | + object: string, |
| 143 | + recordIds: readonly string[], |
| 144 | +): Promise<Set<string>> { |
| 145 | + const live = new Set<string>(); |
| 146 | + for (let i = 0; i < recordIds.length; i += RECORD_SCOPED_DELETE_CHUNK) { |
| 147 | + const batch = recordIds.slice(i, i + RECORD_SCOPED_DELETE_CHUNK); |
| 148 | + const rows = await engine.find(object, { |
| 149 | + where: { id: { $in: batch } }, |
| 150 | + fields: ['id'], |
| 151 | + limit: batch.length, |
| 152 | + context: SYSTEM_CTX, |
| 153 | + }); |
| 154 | + for (const row of (rows ?? [])) { |
| 155 | + if ((row as any)?.id != null) live.add(String((row as any).id)); |
| 156 | + } |
| 157 | + } |
| 158 | + return live; |
| 159 | +} |
| 160 | + |
| 161 | +/** [#5103] Set-based delete of rows by id, chunked like the revoke. */ |
| 162 | +export async function deleteRowsByIds( |
| 163 | + engine: OrphanCleanupEngine, |
| 164 | + table: string, |
| 165 | + rowIds: readonly string[], |
| 166 | +): Promise<void> { |
| 167 | + for (let i = 0; i < rowIds.length; i += RECORD_SCOPED_DELETE_CHUNK) { |
| 168 | + const batch = rowIds.slice(i, i + RECORD_SCOPED_DELETE_CHUNK); |
| 169 | + await engine.delete(table, { |
| 170 | + where: { id: { $in: batch } }, |
| 171 | + multi: true, |
| 172 | + context: SYSTEM_CTX, |
| 173 | + } as any); |
| 174 | + } |
| 175 | +} |
| 176 | + |
| 177 | +/** |
| 178 | + * [#5103] Remove every row of `subject.table` whose RECORD no longer exists. |
| 179 | + * |
| 180 | + * The convergence half of the record-delete cascade, and the shape |
| 181 | + * `SharingRuleService.sweepOrphanedRuleGrants` (#4433) established — with a |
| 182 | + * different predicate, which is the whole point: that sweep asks "does the RULE |
| 183 | + * row still exist", so it can never see a manual share, nor a rule grant whose |
| 184 | + * rule is alive and whose record is not. This one asks "does the RECORD still |
| 185 | + * exist", which is the question the invariant is actually made of, and it is |
| 186 | + * source-agnostic (and, for `sys_share_link`, holder-agnostic). |
| 187 | + * |
| 188 | + * Two callers, one primitive: |
| 189 | + * - `kernel:bootstrapped`, unscoped — historical orphans from before the |
| 190 | + * cascade existed, plus anything a crashed hook missed, converge on the next |
| 191 | + * boot; |
| 192 | + * - the cascade's unbounded-delete branch, scoped to one object — a bulk |
| 193 | + * delete whose row set could not be enumerated cannot name the ids to |
| 194 | + * revoke, but the sweep does not need them: it reads the rows and asks about |
| 195 | + * each record. This is deliberately NOT the rule path's "revoke everything on |
| 196 | + * the object and re-grant asynchronously" — that trade is only available |
| 197 | + * where a reconcile can put the grants back, and nothing can re-create a |
| 198 | + * manual share or re-mint a link someone already holds. |
| 199 | + * |
| 200 | + * Bounded on both axes: rows are read by keyset page (never `OFFSET`, which |
| 201 | + * skips rows in a walk that deletes as it goes — #4363), the scan stops at |
| 202 | + * `max` and SAYS so, and existence is probed one batched `id IN (…)` per object |
| 203 | + * per page rather than one query per row. |
| 204 | + * |
| 205 | + * Fails SAFE per object: a probe that throws leaves that object's rows |
| 206 | + * untouched and is reported in `unresolvedObjects`. "Nothing was queried" is not |
| 207 | + * "nothing matched" — deleting on a failed probe would turn a transient driver |
| 208 | + * error into permanent access loss. (The RESOLVE path fails the other way, and |
| 209 | + * for the same reason: there, "cannot ask" must not grant. Both refuse to act on |
| 210 | + * an unanswered question; only the safe direction differs.) |
| 211 | + */ |
| 212 | +export async function sweepOrphanedRowsByRecordExistence( |
| 213 | + engine: OrphanCleanupEngine, |
| 214 | + subject: OrphanSweepSubject, |
| 215 | + options?: OrphanShareSweepOptions, |
| 216 | + logger?: MinimalLogger, |
| 217 | +): Promise<OrphanShareSweepResult> { |
| 218 | + const result: OrphanShareSweepResult = { |
| 219 | + scanned: 0, |
| 220 | + revoked: 0, |
| 221 | + unresolvedObjects: [], |
| 222 | + truncated: false, |
| 223 | + }; |
| 224 | + const unresolved = new Set<string>(); |
| 225 | + const walk = keysetWalk<any>( |
| 226 | + (q) => engine.find(subject.table, { |
| 227 | + ...q, |
| 228 | + fields: ['id', 'object_name', 'record_id'], |
| 229 | + context: SYSTEM_CTX, |
| 230 | + }), |
| 231 | + { |
| 232 | + where: options?.object ? { object_name: options.object } : undefined, |
| 233 | + pageSize: Math.max(1, options?.batchSize ?? ORPHAN_SWEEP_PAGE_SIZE), |
| 234 | + max: options?.max ?? ORPHAN_SWEEP_MAX_ROWS, |
| 235 | + }, |
| 236 | + ); |
| 237 | + |
| 238 | + try { |
| 239 | + for await (const page of walk.pages()) { |
| 240 | + result.scanned += page.length; |
| 241 | + |
| 242 | + // Group the page by object so existence is one probe per object, not |
| 243 | + // one per row. |
| 244 | + const byObject = new Map<string, Map<string, string[]>>(); |
| 245 | + for (const row of page) { |
| 246 | + const objectName = row?.object_name == null ? '' : String(row.object_name); |
| 247 | + const recordId = row?.record_id == null ? '' : String(row.record_id); |
| 248 | + const rowId = row?.id == null ? '' : String(row.id); |
| 249 | + if (!objectName || !recordId || !rowId) continue; |
| 250 | + const perRecord = byObject.get(objectName) ?? new Map<string, string[]>(); |
| 251 | + const rowIds = perRecord.get(recordId) ?? []; |
| 252 | + rowIds.push(rowId); |
| 253 | + perRecord.set(recordId, rowIds); |
| 254 | + byObject.set(objectName, perRecord); |
| 255 | + } |
| 256 | + |
| 257 | + for (const [objectName, perRecord] of byObject) { |
| 258 | + if (unresolved.has(objectName)) continue; |
| 259 | + const recordIds = [...perRecord.keys()]; |
| 260 | + let live: Set<string>; |
| 261 | + try { |
| 262 | + live = await findLiveRecordIds(engine, objectName, recordIds); |
| 263 | + } catch (err: any) { |
| 264 | + unresolved.add(objectName); |
| 265 | + logger?.warn?.( |
| 266 | + `[sharing] orphan ${subject.noun} sweep could not check whether records still exist — ` + |
| 267 | + `its ${subject.noun} rows were left in place (they are re-checked on the next sweep)`, |
| 268 | + { object: objectName, error: err?.message }, |
| 269 | + ); |
| 270 | + continue; |
| 271 | + } |
| 272 | + const orphanRowIds: string[] = []; |
| 273 | + for (const [recordId, rowIds] of perRecord) { |
| 274 | + if (live.has(recordId)) continue; |
| 275 | + orphanRowIds.push(...rowIds); |
| 276 | + } |
| 277 | + if (orphanRowIds.length === 0) continue; |
| 278 | + await deleteRowsByIds(engine, subject.table, orphanRowIds); |
| 279 | + result.revoked += orphanRowIds.length; |
| 280 | + } |
| 281 | + } |
| 282 | + } catch (err: any) { |
| 283 | + logger?.warn?.( |
| 284 | + `[sharing] orphan ${subject.noun} sweep stopped early — remaining rows are re-checked on the next sweep`, |
| 285 | + { object: options?.object, error: err?.message, scanned: result.scanned }, |
| 286 | + ); |
| 287 | + result.truncated = true; |
| 288 | + } |
| 289 | + |
| 290 | + result.unresolvedObjects = [...unresolved]; |
| 291 | + result.truncated = result.truncated || walk.truncated; |
| 292 | + if (result.revoked > 0) { |
| 293 | + logger?.warn?.( |
| 294 | + `[sharing] revoked ${subject.noun} rows whose record no longer exists (${subject.issue})`, |
| 295 | + { rows: result.revoked, scanned: result.scanned, object: options?.object }, |
| 296 | + ); |
| 297 | + } |
| 298 | + return result; |
| 299 | +} |
0 commit comments