Skip to content

Commit 6b8eba4

Browse files
committed
feat: emit TTL warning events before sandbox expiry
1 parent e08a882 commit 6b8eba4

11 files changed

Lines changed: 214 additions & 0 deletions

File tree

apps/api/src/index.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import { NodeClientMemory } from './services/node-client.memory.js'
1515
import { ArtifactRepoMemory } from './services/artifact-repo.memory.js'
1616
import { createRedisLayer } from './services/redis.ioredis.js'
1717
import { RedisMemory } from './services/redis.memory.js'
18+
import { EventRecorderLive } from './services/event-recorder.live.js'
1819
import { IdempotencyRepoMemory } from './workers/idempotency-cleanup.memory.js'
1920
import { QuotaMemory } from './services/quota.memory.js'
2021
import { UsageMemory } from './services/usage.memory.js'
@@ -82,6 +83,7 @@ const GracefulShutdownLive = Layer.scopedDiscard(
8283

8384
const ServerLive = Layer.mergeAll(AppLive, WorkersLive, GracefulShutdownLive).pipe(
8485
Layer.provide(ShutdownControllerLive),
86+
Layer.provide(EventRecorderLive),
8587
Layer.provide(SandboxRepoMemory),
8688
Layer.provide(ExecRepoMemory),
8789
Layer.provide(SessionRepoMemory),

apps/api/src/services/events.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,10 @@ export function sandboxForked(params: { fork_sandbox_id: string }): EventPayload
3939
return { type: 'sandbox.forked', data: { fork_sandbox_id: params.fork_sandbox_id } }
4040
}
4141

42+
export function sandboxTtlWarning(params: { seconds_remaining: number }): EventPayload {
43+
return { type: 'sandbox.ttl_warning', data: { seconds_remaining: params.seconds_remaining } }
44+
}
45+
4246
export function sandboxStopping(params: { reason: string }): EventPayload {
4347
return { type: 'sandbox.stopping', data: { reason: params.reason } }
4448
}

apps/api/src/services/redis.ioredis.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,12 @@ export function createIoRedisApi(client: Redis): RedisApi {
139139
return result === 1
140140
}),
141141

142+
markTtlWarned: (sandboxId, ttlSeconds) =>
143+
Effect.promise(async () => {
144+
const result = await client.set(`ttl_warned:${sandboxId}`, '1', 'EX', ttlSeconds, 'NX')
145+
return result === 'OK'
146+
}),
147+
142148
ping: () =>
143149
Effect.promise(async () => {
144150
try {

apps/api/src/services/redis.memory.ts

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ export function createInMemoryRedisApi(): RedisApi {
2323
const artifactPaths = new Map<string, Set<string>>()
2424
const leaderLocks = new Map<string, { instanceId: string; expiresAt: number }>()
2525
const nodeHeartbeats = new Map<string, TtlEntry>()
26+
const ttlWarned = new Map<string, TtlEntry>()
2627

2728
function isExpired(entry: SlotEntry): boolean {
2829
return Date.now() >= entry.expiresAt
@@ -179,6 +180,16 @@ export function createInMemoryRedisApi(): RedisApi {
179180
return entry !== undefined && entry.expiresAt > Date.now()
180181
}),
181182

183+
markTtlWarned: (sandboxId, ttlSeconds) =>
184+
Effect.sync(() => {
185+
const key = `ttl_warned:${sandboxId}`
186+
const now = Date.now()
187+
const existing = ttlWarned.get(key)
188+
if (existing && existing.expiresAt > now) return false
189+
ttlWarned.set(key, { expiresAt: now + ttlSeconds * 1000 })
190+
return true
191+
}),
192+
182193
ping: () => Effect.succeed(true),
183194
}
184195
}

apps/api/src/services/redis.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,12 @@ export interface RedisApi {
105105
nodeId: string,
106106
) => Effect.Effect<boolean, never, never>
107107

108+
/** Atomically mark a sandbox as TTL-warned. Returns true if newly marked. */
109+
readonly markTtlWarned: (
110+
sandboxId: string,
111+
ttlSeconds: number,
112+
) => Effect.Effect<boolean, never, never>
113+
108114
/** Ping to check connectivity. */
109115
readonly ping: () => Effect.Effect<boolean, never, never>
110116
}

apps/api/src/services/sandbox-repo.memory.ts

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -250,6 +250,17 @@ export function createInMemorySandboxRepo(): SandboxRepoApi {
250250
})
251251
}),
252252

253+
findNearTtlExpiry: (warningThresholdSeconds) =>
254+
Effect.sync(() => {
255+
const now = Date.now()
256+
return Array.from(store.values()).filter((r) => {
257+
if (r.status !== 'running' || !r.startedAt) return false
258+
const expiresAt = r.startedAt.getTime() + r.ttlSeconds * 1000
259+
const warningAt = expiresAt - warningThresholdSeconds * 1000
260+
return warningAt <= now && expiresAt > now
261+
})
262+
}),
263+
253264
findIdleSince: (cutoff) =>
254265
Effect.sync(() =>
255266
Array.from(store.values()).filter((r) => {

apps/api/src/services/sandbox-repo.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,11 @@ export interface SandboxRepoApi {
119119
/** Find running sandboxes that have exceeded their TTL. */
120120
readonly findExpiredTtl: () => Effect.Effect<SandboxRow[], never, never>
121121

122+
/** Find running sandboxes within the warning threshold of TTL expiry (not yet expired). */
123+
readonly findNearTtlExpiry: (
124+
warningThresholdSeconds: number,
125+
) => Effect.Effect<SandboxRow[], never, never>
126+
122127
/** Find running sandboxes with lastActivityAt before the given cutoff. */
123128
readonly findIdleSince: (
124129
cutoff: Date,

apps/api/src/workers/index.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,11 @@ import type { SandboxRepo } from '../services/sandbox-repo.js'
44
import type { ArtifactRepo } from '../services/artifact-repo.js'
55
import type { ObjectStorage } from '../services/object-storage.js'
66
import type { QuotaService } from '../services/quota.js'
7+
import type { EventRecorder } from '../services/event-recorder.js'
78
import type { IdempotencyRepo } from './idempotency-cleanup.js'
89
import { startWorkers, type WorkerConfig } from './runner.js'
910
import { ttlEnforcementWorker } from './ttl-enforcement.js'
11+
import { ttlWarningWorker } from './ttl-warning.js'
1012
import { idleShutdownWorker } from './idle-shutdown.js'
1113
import { orphanReconciliationWorker } from './orphan-reconciliation.js'
1214
import { queueTimeoutWorker } from './queue-timeout.js'
@@ -21,10 +23,12 @@ export type WorkerDeps =
2123
| ArtifactRepo
2224
| ObjectStorage
2325
| QuotaService
26+
| EventRecorder
2427
| IdempotencyRepo
2528

2629
const allWorkers: ReadonlyArray<WorkerConfig<WorkerDeps>> = [
2730
ttlEnforcementWorker,
31+
ttlWarningWorker,
2832
idleShutdownWorker,
2933
orphanReconciliationWorker,
3034
queueTimeoutWorker,

apps/api/src/workers/runner.test.ts

Lines changed: 119 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,14 +4,17 @@ import { createInMemoryRedisApi } from '../services/redis.memory.js'
44
import { createInMemorySandboxRepo } from '../services/sandbox-repo.memory.js'
55
import { createInMemoryArtifactRepo } from '../services/artifact-repo.memory.js'
66
import { createInMemoryObjectStorage } from '../services/object-storage.memory.js'
7+
import { createLiveEventRecorder } from '../services/event-recorder.live.js'
78
import { RedisService, type RedisApi } from '../services/redis.js'
89
import { SandboxRepo, type SandboxRepoApi } from '../services/sandbox-repo.js'
910
import { ArtifactRepo, type ArtifactRepoApi } from '../services/artifact-repo.js'
1011
import { ObjectStorage, type ObjectStorageApi } from '../services/object-storage.js'
12+
import { EventRecorder, type EventRecorderApi } from '../services/event-recorder.js'
1113
import { IdempotencyRepo, type IdempotencyRepoApi } from './idempotency-cleanup.js'
1214
import { createTestableIdempotencyRepo } from './idempotency-cleanup.memory.js'
1315
import { runWorkerTick, type WorkerConfig } from './runner.js'
1416
import { ttlEnforcementWorker } from './ttl-enforcement.js'
17+
import { ttlWarningWorker } from './ttl-warning.js'
1518
import { idleShutdownWorker } from './idle-shutdown.js'
1619
import { orphanReconciliationWorker } from './orphan-reconciliation.js'
1720
import { queueTimeoutWorker } from './queue-timeout.js'
@@ -28,6 +31,7 @@ let redis: RedisApi
2831
let sandboxRepo: SandboxRepoApi
2932
let artifactRepo: ArtifactRepoApi
3033
let objectStorage: ObjectStorageApi
34+
let eventRecorder: EventRecorderApi
3135
let idempotencyApi: IdempotencyRepoApi
3236
let idempotencyStore: Map<string, { createdAt: Date }>
3337
let quotaApi: QuotaApi & { setOrgQuota: (orgId: string, quota: Record<string, unknown>) => void }
@@ -45,6 +49,7 @@ beforeEach(() => {
4549
sandboxRepo = createInMemorySandboxRepo()
4650
artifactRepo = createInMemoryArtifactRepo()
4751
objectStorage = createInMemoryObjectStorage()
52+
eventRecorder = createLiveEventRecorder(objectStorage, redis)
4853
const testable = createTestableIdempotencyRepo()
4954
idempotencyApi = testable.api
5055
idempotencyStore = testable.store
@@ -55,6 +60,7 @@ beforeEach(() => {
5560
Layer.succeed(SandboxRepo, sandboxRepo),
5661
Layer.succeed(ArtifactRepo, artifactRepo),
5762
Layer.succeed(ObjectStorage, objectStorage),
63+
Layer.succeed(EventRecorder, eventRecorder),
5864
Layer.succeed(IdempotencyRepo, idempotencyApi),
5965
Layer.succeed(QuotaService, quotaApi),
6066
)
@@ -746,3 +752,116 @@ describe('replay-retention', () => {
746752
expect(longRow!.replayExpiresAt!.getTime()).toBeGreaterThan(Date.now())
747753
})
748754
})
755+
756+
// ---------------------------------------------------------------------------
757+
// TTL warning worker
758+
// ---------------------------------------------------------------------------
759+
760+
describe('ttl-warning', () => {
761+
async function createRunningSandbox(orgId: string, ttlSeconds: number, startedSecondsAgo: number) {
762+
const id = generateUUIDv7()
763+
await Effect.runPromise(
764+
sandboxRepo.create({
765+
id,
766+
orgId,
767+
imageId: SEED_IMAGE_ID,
768+
profileId: SEED_PROFILE_ID,
769+
profileName: 'small',
770+
env: null,
771+
ttlSeconds,
772+
imageRef: 'sandchest://ubuntu-22.04',
773+
}),
774+
)
775+
// assignNode sets status to 'running' and startedAt
776+
const nodeId = generateUUIDv7()
777+
await Effect.runPromise(redis.registerNodeHeartbeat(base62Encode(nodeId), 300))
778+
await Effect.runPromise(sandboxRepo.assignNode(id, orgId, nodeId))
779+
780+
// Backdate startedAt to simulate time passing
781+
const row = await Effect.runPromise(sandboxRepo.findById(id, orgId))
782+
if (row) {
783+
const backdatedStart = new Date(Date.now() - startedSecondsAgo * 1000)
784+
// Use updateStatus to persist the sandbox, then manually adjust startedAt
785+
// The memory repo stores mutable objects, so we can rely on the reference
786+
;(row as { startedAt: Date }).startedAt = backdatedStart
787+
}
788+
789+
return id
790+
}
791+
792+
test('returns 0 when no sandboxes exist', async () => {
793+
const count = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
794+
expect(count).toBe(0)
795+
})
796+
797+
test('does not warn sandboxes far from TTL expiry', async () => {
798+
// TTL 3600s, started 100s ago → 3500s remaining (well above 60s threshold)
799+
await createRunningSandbox('org_test', 3600, 100)
800+
801+
const count = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
802+
expect(count).toBe(0)
803+
})
804+
805+
test('warns sandbox near TTL expiry', async () => {
806+
// TTL 100s, started 70s ago → 30s remaining (within 60s threshold)
807+
const id = await createRunningSandbox('org_test', 100, 70)
808+
const sandboxId = bytesToId(SANDBOX_PREFIX, id)
809+
810+
const count = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
811+
expect(count).toBe(1)
812+
813+
// Verify event was recorded
814+
const events = await Effect.runPromise(eventRecorder.getEvents({ sandboxId, orgId: 'org_test' }))
815+
const warning = events.find((e) => e.type === 'sandbox.ttl_warning')
816+
expect(warning).toBeDefined()
817+
expect(warning!.data.seconds_remaining).toBeNumber()
818+
expect(warning!.data.seconds_remaining as number).toBeLessThanOrEqual(60)
819+
})
820+
821+
test('does not warn the same sandbox twice', async () => {
822+
// TTL 100s, started 70s ago → 30s remaining
823+
await createRunningSandbox('org_test', 100, 70)
824+
825+
const count1 = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
826+
expect(count1).toBe(1)
827+
828+
const count2 = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
829+
expect(count2).toBe(0)
830+
})
831+
832+
test('does not warn already-expired sandboxes', async () => {
833+
// TTL 100s, started 110s ago → already expired (handled by ttl-enforcement)
834+
await createRunningSandbox('org_test', 100, 110)
835+
836+
const count = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
837+
expect(count).toBe(0)
838+
})
839+
840+
test('does not warn non-running sandboxes', async () => {
841+
const id = generateUUIDv7()
842+
await Effect.runPromise(
843+
sandboxRepo.create({
844+
id,
845+
orgId: 'org_test',
846+
imageId: SEED_IMAGE_ID,
847+
profileId: SEED_PROFILE_ID,
848+
profileName: 'small',
849+
env: null,
850+
ttlSeconds: 0,
851+
imageRef: 'sandchest://ubuntu-22.04',
852+
}),
853+
)
854+
// Leave as 'queued'
855+
const count = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
856+
expect(count).toBe(0)
857+
})
858+
859+
test('warns multiple sandboxes near expiry', async () => {
860+
// Both within warning threshold
861+
await createRunningSandbox('org_test', 100, 50)
862+
await createRunningSandbox('org_test', 200, 160)
863+
864+
const count = await run(runWorkerTick(ttlWarningWorker, 'inst-1'))
865+
expect(count).toBe(2)
866+
})
867+
})
Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
import { Effect } from 'effect'
2+
import { bytesToId, SANDBOX_PREFIX } from '@sandchest/contract'
3+
import { SandboxRepo } from '../services/sandbox-repo.js'
4+
import { RedisService } from '../services/redis.js'
5+
import { EventRecorder } from '../services/event-recorder.js'
6+
import { sandboxTtlWarning } from '../services/events.js'
7+
import type { WorkerConfig } from './runner.js'
8+
9+
/** Warn 60 seconds before TTL expiry. */
10+
const WARNING_THRESHOLD_SECONDS = 60
11+
12+
/** Dedup key TTL — prevents re-warning the same sandbox for 120 seconds. */
13+
const DEDUP_TTL_SECONDS = 120
14+
15+
export const ttlWarningWorker: WorkerConfig<SandboxRepo | RedisService | EventRecorder> = {
16+
name: 'ttl-warning',
17+
intervalMs: 10_000,
18+
handler: Effect.gen(function* () {
19+
const repo = yield* SandboxRepo
20+
const redis = yield* RedisService
21+
const recorder = yield* EventRecorder
22+
23+
const nearExpiry = yield* repo.findNearTtlExpiry(WARNING_THRESHOLD_SECONDS)
24+
25+
let warned = 0
26+
for (const sandbox of nearExpiry) {
27+
const sandboxId = bytesToId(SANDBOX_PREFIX, sandbox.id)
28+
const isNew = yield* redis.markTtlWarned(sandboxId, DEDUP_TTL_SECONDS)
29+
if (!isNew) continue
30+
31+
const expiresAt = sandbox.startedAt!.getTime() + sandbox.ttlSeconds * 1000
32+
const secondsRemaining = Math.max(0, Math.round((expiresAt - Date.now()) / 1000))
33+
34+
yield* recorder.record({
35+
sandboxId,
36+
orgId: sandbox.orgId,
37+
event: sandboxTtlWarning({ seconds_remaining: secondsRemaining }),
38+
})
39+
40+
warned++
41+
}
42+
43+
return warned
44+
}),
45+
}

0 commit comments

Comments
 (0)