Skip to content

Commit 491ab88

Browse files
chore(search): add paced projection replacement and data retirement (#8534)
* chore(search): add paced projection replacement and data retirement * chore(search): keep retirement usage in the CLI * fix(db): fence schema push and require direct retirement sessions * fix(db): preserve test connection URL parameters
1 parent 045b07c commit 491ab88

14 files changed

Lines changed: 2702 additions & 250 deletions

‎packages/db/drizzle.config.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,5 +20,7 @@ export default {
2020
'!script_migrations',
2121
'!search_embedding_cleanup_progress',
2222
'!search_embedding_cleanup_targets',
23+
'!search_retirement_*',
24+
'!embedding_search_retirement_*',
2325
],
2426
} satisfies Config

‎packages/db/maintenance/index.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
export {
2+
abortSearchRetirement,
3+
advanceSearchRetirement,
4+
beginSearchRetirementPurge,
5+
cutoverSearchRetirement,
6+
finalizeSearchRetirement,
7+
getSearchRetirementStatus,
8+
initializeSearchRetirement,
9+
SearchRetirementError,
10+
type SearchRetirementPhase,
11+
type SearchRetirementStatus,
12+
} from '@sim/db/maintenance/search-retirement'
13+
export {
14+
RetireSearchHealthError,
15+
readSearchRetirementHealth,
16+
readSearchRetirementHealthLimits,
17+
type SearchRetirementHealthLimits,
18+
} from '@sim/db/maintenance/search-retirement-health'
Lines changed: 266 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,266 @@
1+
import { mkdtemp, rm, writeFile } from 'node:fs/promises'
2+
import { tmpdir } from 'node:os'
3+
import { join } from 'node:path'
4+
import {
5+
RetireSearchHealthError,
6+
readSearchRetirementHealth,
7+
readSearchRetirementHealthLimits,
8+
type SearchRetirementHealthLimits,
9+
} from '@sim/db/maintenance/search-retirement-health'
10+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
11+
12+
const NOW = Date.parse('2026-10-01T12:00:00.000Z')
13+
const LIMITS: SearchRetirementHealthLimits = {
14+
databaseId: 'test-database',
15+
maxReplicaLagBytes: 1_000,
16+
maxReplicaLagSeconds: 2,
17+
maxWalBytesPerSecond: 2_000,
18+
maxDatabaseP95Ms: 100,
19+
maxCpuPercent: 50,
20+
minFreeStorageBytes: 5_000,
21+
maxSampleAgeMs: 10_000,
22+
}
23+
const SAMPLE = {
24+
observedAt: '2026-10-01T12:00:00.000Z',
25+
databaseId: 'test-database',
26+
healthy: true,
27+
replicaLagBytes: 0,
28+
replicaLagSeconds: 0,
29+
walBytesPerSecond: 100,
30+
databaseP95Ms: 10,
31+
cpuPercent: 10,
32+
freeStorageBytes: 10_000,
33+
maintenanceAllowed: true,
34+
cutoverAllowed: false,
35+
}
36+
37+
describe('Search retirement external health gate', () => {
38+
let directory: string
39+
let path: string
40+
41+
beforeEach(async () => {
42+
vi.useFakeTimers({ toFake: ['Date'] })
43+
vi.setSystemTime(NOW)
44+
directory = await mkdtemp(join(tmpdir(), 'search-retirement-health-'))
45+
path = join(directory, 'health.json')
46+
await writeFile(path, JSON.stringify(SAMPLE))
47+
})
48+
49+
afterEach(async () => {
50+
vi.useRealTimers()
51+
await rm(directory, { recursive: true, force: true })
52+
})
53+
54+
it('reads a replacement sample before allowing another page', async () => {
55+
expect(await readSearchRetirementHealth(path, LIMITS)).toBe(NOW)
56+
await writeFile(path, JSON.stringify({ ...SAMPLE, maintenanceAllowed: false }))
57+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
58+
reason: 'maintenance_refused',
59+
})
60+
})
61+
62+
it('requires separate cutover approval while allowing ordinary maintenance', async () => {
63+
expect(await readSearchRetirementHealth(path, LIMITS)).toBe(NOW)
64+
await expect(readSearchRetirementHealth(path, LIMITS, { cutover: true })).rejects.toMatchObject(
65+
{
66+
reason: 'cutover_refused',
67+
}
68+
)
69+
await writeFile(path, JSON.stringify({ ...SAMPLE, cutoverAllowed: true }))
70+
expect(await readSearchRetirementHealth(path, LIMITS, { cutover: true })).toBe(NOW)
71+
})
72+
73+
it('rejects an operator CPU ceiling above one hundred percent', async () => {
74+
await writeFile(path, JSON.stringify({ ...LIMITS, maxCpuPercent: 100.001 }))
75+
await expect(readSearchRetirementHealthLimits(path)).rejects.toMatchObject({
76+
reason: 'invalid_limits',
77+
})
78+
})
79+
80+
it.each(['missing', 'directory'] as const)('rejects a %s health file', async (kind) => {
81+
await expect(
82+
readSearchRetirementHealth(kind === 'missing' ? join(directory, 'absent') : directory, LIMITS)
83+
).rejects.toBeInstanceOf(RetireSearchHealthError)
84+
})
85+
86+
it('accepts at most eight KiB, including UTF-8 bytes and trailing whitespace', async () => {
87+
const encoded = JSON.stringify(SAMPLE)
88+
await writeFile(path, encoded.padEnd(8_192, ' '))
89+
expect(await readSearchRetirementHealth(path, LIMITS)).toBe(NOW)
90+
await writeFile(path, encoded.padEnd(8_193, ' '))
91+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
92+
reason: 'health_file_too_large',
93+
})
94+
await writeFile(path, JSON.stringify({ ...SAMPLE, padding: 'é'.repeat(4_096) }))
95+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
96+
reason: 'health_file_too_large',
97+
})
98+
})
99+
100+
it.each(['', '{', 'null', '[]', 'true', '{"observedAt":123}'])(
101+
'rejects malformed or incomplete wire content: %s',
102+
async (content) => {
103+
await writeFile(path, content)
104+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
105+
reason: 'invalid_sample',
106+
})
107+
}
108+
)
109+
110+
it('rejects invalid UTF-8 instead of substituting a replacement character', async () => {
111+
const encoded = Buffer.from(JSON.stringify(SAMPLE))
112+
const location = encoded.indexOf('test-database')
113+
encoded[location] = 0xff
114+
await writeFile(path, encoded)
115+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
116+
reason: 'invalid_sample',
117+
})
118+
})
119+
120+
it.each([
121+
{ healthy: 'true' },
122+
{ maintenanceAllowed: 1 },
123+
{ cutoverAllowed: 'true' },
124+
{ cutoverAllowed: undefined },
125+
{ databaseId: '' },
126+
{ unexpected: true },
127+
{ observedAt: null },
128+
{ replicaLagBytes: -1 },
129+
{ replicaLagSeconds: null },
130+
{ walBytesPerSecond: '100' },
131+
{ databaseP95Ms: -1 },
132+
{ cpuPercent: null },
133+
{ freeStorageBytes: -1 },
134+
])('rejects malformed sample field %j', async (fields) => {
135+
await writeFile(path, JSON.stringify({ ...SAMPLE, ...fields }))
136+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
137+
reason: 'invalid_sample',
138+
})
139+
})
140+
141+
it('rejects a nonfinite JSON number', async () => {
142+
await writeFile(
143+
path,
144+
JSON.stringify(SAMPLE).replace('"replicaLagBytes":0', '"replicaLagBytes":1e999')
145+
)
146+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
147+
reason: 'invalid_sample',
148+
})
149+
})
150+
151+
it.each([
152+
'2026-10-01T12:00:00+00:00',
153+
'2026-10-01 12:00:00Z',
154+
'2026-02-30T12:00:00.000Z',
155+
'2026-10-01T24:00:00.000Z',
156+
])('rejects a noncanonical or impossible UTC timestamp: %s', async (observedAt) => {
157+
await writeFile(path, JSON.stringify({ ...SAMPLE, observedAt }))
158+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
159+
reason: 'invalid_sample',
160+
})
161+
})
162+
163+
it.each([
164+
['2026-10-01T11:59:49.999Z', 'stale_sample'],
165+
['2026-10-01T12:00:01.001Z', 'future_sample'],
166+
])('rejects a sample outside the time budget: %s', async (observedAt, reason) => {
167+
await writeFile(path, JSON.stringify({ ...SAMPLE, observedAt }))
168+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({ reason })
169+
})
170+
171+
it('does not reuse a previously fresh sample once it expires', async () => {
172+
await readSearchRetirementHealth(path, LIMITS)
173+
vi.setSystemTime(NOW + LIMITS.maxSampleAgeMs + 1)
174+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({
175+
reason: 'stale_sample',
176+
})
177+
})
178+
179+
it.each([
180+
[{ databaseId: 'other-database' }, 'identity_mismatch'],
181+
[{ healthy: false }, 'unhealthy'],
182+
[{ maintenanceAllowed: false }, 'maintenance_refused'],
183+
[{ replicaLagBytes: 1_001 }, 'replica_lag_bytes'],
184+
[{ replicaLagSeconds: 2.001 }, 'replica_lag_seconds'],
185+
[{ walBytesPerSecond: 2_001 }, 'wal_rate'],
186+
[{ databaseP95Ms: 101 }, 'database_latency'],
187+
[{ cpuPercent: 51 }, 'cpu_usage'],
188+
[{ freeStorageBytes: 4_999 }, 'storage_headroom'],
189+
] as const)('refuses unsafe sample %j', async (fields, reason) => {
190+
await writeFile(path, JSON.stringify({ ...SAMPLE, ...fields }))
191+
await expect(readSearchRetirementHealth(path, LIMITS)).rejects.toMatchObject({ reason })
192+
})
193+
194+
it('permits zero-lag limits and refuses any replication debt', async () => {
195+
const limits = { ...LIMITS, maxReplicaLagBytes: 0, maxReplicaLagSeconds: 0 }
196+
await readSearchRetirementHealth(path, limits)
197+
await writeFile(path, JSON.stringify({ ...SAMPLE, replicaLagBytes: 1 }))
198+
await expect(readSearchRetirementHealth(path, limits)).rejects.toMatchObject({
199+
reason: 'replica_lag_bytes',
200+
})
201+
})
202+
203+
it.each([
204+
{ databaseId: '' },
205+
{ maxReplicaLagBytes: -1 },
206+
{ maxReplicaLagSeconds: Number.NaN },
207+
{ maxWalBytesPerSecond: 0 },
208+
{ maxDatabaseP95Ms: Number.POSITIVE_INFINITY },
209+
{ maxCpuPercent: 0 },
210+
{ minFreeStorageBytes: 0 },
211+
{ maxSampleAgeMs: 0 },
212+
{ maxSampleAgeMs: 30_001 },
213+
])('rejects limits that disable a guard: %j', async (fields) => {
214+
await expect(readSearchRetirementHealth(path, { ...LIMITS, ...fields })).rejects.toMatchObject({
215+
reason: 'invalid_limits',
216+
})
217+
await writeFile(path, JSON.stringify({ ...LIMITS, ...fields }))
218+
await expect(readSearchRetirementHealthLimits(path)).rejects.toMatchObject({
219+
reason: 'invalid_limits',
220+
})
221+
})
222+
223+
it('loads the bounded operator policy and enforces it against the health sample', async () => {
224+
const policyPath = join(directory, 'policy.json')
225+
await writeFile(policyPath, JSON.stringify({ ...LIMITS, maxReplicaLagBytes: 0 }))
226+
const limits = await readSearchRetirementHealthLimits(policyPath)
227+
await writeFile(path, JSON.stringify({ ...SAMPLE, replicaLagBytes: 1 }))
228+
await expect(readSearchRetirementHealth(path, limits)).rejects.toMatchObject({
229+
reason: 'replica_lag_bytes',
230+
})
231+
})
232+
233+
it.each(['{}', 'null', '{', JSON.stringify({ ...LIMITS, unexpected: true })])(
234+
'rejects malformed operator policy content: %s',
235+
async (content) => {
236+
await writeFile(path, content)
237+
await expect(readSearchRetirementHealthLimits(path)).rejects.toMatchObject({
238+
reason: 'invalid_limits',
239+
})
240+
}
241+
)
242+
243+
it('enforces the same eight-KiB cap on operator policies', async () => {
244+
await writeFile(path, JSON.stringify(LIMITS).padEnd(8_193, ' '))
245+
await expect(readSearchRetirementHealthLimits(path)).rejects.toMatchObject({
246+
reason: 'health_file_too_large',
247+
})
248+
})
249+
250+
it('returns only typed generic reasons for payload and filesystem failures', async () => {
251+
const privateMarker = 'private-operator-marker'
252+
await writeFile(path, JSON.stringify({ ...SAMPLE, databaseId: privateMarker }))
253+
for (const candidate of [path, join(directory, privateMarker)]) {
254+
try {
255+
await readSearchRetirementHealth(candidate, LIMITS)
256+
expect.fail('The health gate should refuse this sample')
257+
} catch (error) {
258+
expect(error).toBeInstanceOf(RetireSearchHealthError)
259+
if (!(error instanceof RetireSearchHealthError)) throw error
260+
expect(error.message).not.toContain(privateMarker)
261+
expect(error.message).not.toContain(directory)
262+
expect(error.cause).toBeUndefined()
263+
}
264+
}
265+
})
266+
})

0 commit comments

Comments
 (0)