Skip to content

Commit ac2273c

Browse files
committed
feat: implement background workers with Redis leader election
1 parent 3e097d1 commit ac2273c

23 files changed

Lines changed: 938 additions & 2 deletions

apps/api/src/index.ts

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { HttpServer } from '@effect/platform'
22
import { NodeHttpServer, NodeRuntime } from '@effect/platform-node'
3-
import { Layer } from 'effect'
3+
import { Effect, Layer } from 'effect'
44
import { createServer } from 'node:http'
55
import { ApiRouter } from './server.js'
66
import { withAuth, withRequestId } from './middleware.js'
@@ -13,6 +13,8 @@ import { NodeClientMemory } from './services/node-client.memory.js'
1313
import { ArtifactRepoMemory } from './services/artifact-repo.memory.js'
1414
import { createRedisLayer } from './services/redis.ioredis.js'
1515
import { RedisMemory } from './services/redis.memory.js'
16+
import { IdempotencyRepoMemory } from './workers/idempotency-cleanup.memory.js'
17+
import { startAllWorkers } from './workers/index.js'
1618

1719
const PORT = Number(process.env.PORT ?? 3000)
1820
const REDIS_URL = process.env.REDIS_URL
@@ -22,13 +24,21 @@ const AppLive = ApiRouter.pipe(withRateLimit, withAuth, withRequestId, HttpServe
2224

2325
const RedisLive = REDIS_URL ? createRedisLayer(REDIS_URL) : RedisMemory
2426

25-
const ServerLive = AppLive.pipe(
27+
// Workers launch as scoped fibers alongside the HTTP server
28+
const WorkersLive = Layer.scopedDiscard(
29+
Effect.gen(function* () {
30+
yield* startAllWorkers()
31+
}),
32+
)
33+
34+
const ServerLive = Layer.mergeAll(AppLive, WorkersLive).pipe(
2635
Layer.provide(SandboxRepoMemory),
2736
Layer.provide(ExecRepoMemory),
2837
Layer.provide(SessionRepoMemory),
2938
Layer.provide(ObjectStorageMemory),
3039
Layer.provide(NodeClientMemory),
3140
Layer.provide(ArtifactRepoMemory),
41+
Layer.provide(IdempotencyRepoMemory),
3242
Layer.provide(RedisLive),
3343
Layer.provide(NodeHttpServer.layer(() => createServer(), { port: PORT })),
3444
)

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,23 @@ export function createInMemoryArtifactRepo(): ArtifactRepoApi {
8383
}
8484
return n
8585
}),
86+
87+
findExpiredRetention: (before) =>
88+
Effect.sync(() =>
89+
Array.from(store.values()).filter(
90+
(r) => r.retentionUntil !== null && r.retentionUntil.getTime() < before.getTime(),
91+
),
92+
),
93+
94+
deleteByIds: (ids) =>
95+
Effect.sync(() => {
96+
let deleted = 0
97+
for (const id of ids) {
98+
const key = keyFor(id)
99+
if (store.delete(key)) deleted++
100+
}
101+
return deleted
102+
}),
86103
}
87104
}
88105

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,16 @@ export interface ArtifactRepoApi {
5252
sandboxId: Uint8Array,
5353
orgId: string,
5454
) => Effect.Effect<number, never, never>
55+
56+
/** Find artifacts past their retention date. */
57+
readonly findExpiredRetention: (
58+
before: Date,
59+
) => Effect.Effect<ArtifactRow[], never, never>
60+
61+
/** Delete artifacts by IDs. Returns count of deleted rows. */
62+
readonly deleteByIds: (
63+
ids: Uint8Array[],
64+
) => Effect.Effect<number, never, never>
5565
}
5666

5767
export class ArtifactRepo extends Context.Tag('ArtifactRepo')<ArtifactRepo, ArtifactRepoApi>() {}

apps/api/src/services/object-storage.memory.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,11 @@ export function createInMemoryObjectStorage(): ObjectStorageApi {
1919
Effect.sync(() =>
2020
`https://s3.example.com/${key}?expires=${expiresInSeconds}`,
2121
),
22+
23+
deleteObject: (key) =>
24+
Effect.sync(() => {
25+
store.delete(key)
26+
}),
2227
}
2328
}
2429

apps/api/src/services/object-storage.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,11 @@ export interface ObjectStorageApi {
1818
key: string,
1919
expiresInSeconds: number,
2020
) => Effect.Effect<string, never, never>
21+
22+
/** Delete an object from a bucket. */
23+
readonly deleteObject: (
24+
key: string,
25+
) => Effect.Effect<void, never, never>
2126
}
2227

2328
export class ObjectStorage extends Context.Tag('ObjectStorage')<ObjectStorage, ObjectStorageApi>() {}

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

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,24 @@ export function createIoRedisApi(client: Redis): RedisApi {
120120
return client.scard(key)
121121
}),
122122

123+
acquireLeaderLock: (workerName, instanceId, ttlMs) =>
124+
Effect.promise(async () => {
125+
const key = `worker:${workerName}:leader`
126+
const result = await client.set(key, instanceId, 'PX', ttlMs, 'NX')
127+
return result === 'OK'
128+
}),
129+
130+
registerNodeHeartbeat: (nodeId, ttlSeconds) =>
131+
Effect.promise(async () => {
132+
await client.set(`node_heartbeat:${nodeId}`, '1', 'EX', ttlSeconds)
133+
}),
134+
135+
hasNodeHeartbeat: (nodeId) =>
136+
Effect.promise(async () => {
137+
const result = await client.exists(`node_heartbeat:${nodeId}`)
138+
return result === 1
139+
}),
140+
123141
ping: () =>
124142
Effect.promise(async () => {
125143
try {

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

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,19 @@ interface RateEntry {
1010
timestamps: number[]
1111
}
1212

13+
interface TtlEntry {
14+
expiresAt: number
15+
}
16+
1317
/** In-memory Redis implementation for testing. */
1418
export function createInMemoryRedisApi(): RedisApi {
1519
const slots = new Map<string, SlotEntry>()
1620
const rateLimits = new Map<string, RateEntry>()
1721
const execEvents = new Map<string, BufferedEvent[]>()
1822
const replayEvents = new Map<string, BufferedEvent[]>()
1923
const artifactPaths = new Map<string, Set<string>>()
24+
const leaderLocks = new Map<string, { instanceId: string; expiresAt: number }>()
25+
const nodeHeartbeats = new Map<string, TtlEntry>()
2026

2127
function isExpired(entry: SlotEntry): boolean {
2228
return Date.now() >= entry.expiresAt
@@ -145,6 +151,31 @@ export function createInMemoryRedisApi(): RedisApi {
145151
return artifactPaths.get(key)?.size ?? 0
146152
}),
147153

154+
acquireLeaderLock: (workerName, instanceId, ttlMs) =>
155+
Effect.sync(() => {
156+
const key = `worker:${workerName}:leader`
157+
const now = Date.now()
158+
const existing = leaderLocks.get(key)
159+
if (existing && existing.expiresAt > now) {
160+
return existing.instanceId === instanceId
161+
}
162+
leaderLocks.set(key, { instanceId, expiresAt: now + ttlMs })
163+
return true
164+
}),
165+
166+
registerNodeHeartbeat: (nodeId, ttlSeconds) =>
167+
Effect.sync(() => {
168+
nodeHeartbeats.set(`node_heartbeat:${nodeId}`, {
169+
expiresAt: Date.now() + ttlSeconds * 1000,
170+
})
171+
}),
172+
173+
hasNodeHeartbeat: (nodeId) =>
174+
Effect.sync(() => {
175+
const entry = nodeHeartbeats.get(`node_heartbeat:${nodeId}`)
176+
return entry !== undefined && entry.expiresAt > Date.now()
177+
}),
178+
148179
ping: () => Effect.succeed(true),
149180
}
150181
}

apps/api/src/services/redis.ts

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -86,6 +86,24 @@ export interface RedisApi {
8686
sandboxId: string,
8787
) => Effect.Effect<number, never, never>
8888

89+
/** Acquire a worker leader lock. Returns true if this instance became the leader. */
90+
readonly acquireLeaderLock: (
91+
workerName: string,
92+
instanceId: string,
93+
ttlMs: number,
94+
) => Effect.Effect<boolean, never, never>
95+
96+
/** Register a node heartbeat with a TTL. */
97+
readonly registerNodeHeartbeat: (
98+
nodeId: string,
99+
ttlSeconds: number,
100+
) => Effect.Effect<void, never, never>
101+
102+
/** Check if a node has an active heartbeat. */
103+
readonly hasNodeHeartbeat: (
104+
nodeId: string,
105+
) => Effect.Effect<boolean, never, never>
106+
89107
/** Ping to check connectivity. */
90108
readonly ping: () => Effect.Effect<boolean, never, never>
91109
}

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

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ export function createInMemorySandboxRepo(): SandboxRepoApi {
4949
const row: SandboxRow = {
5050
id: params.id,
5151
orgId: params.orgId,
52+
nodeId: null,
5253
imageId: params.imageId,
5354
profileId: params.profileId,
5455
profileName: params.profileName,
@@ -59,6 +60,7 @@ export function createInMemorySandboxRepo(): SandboxRepoApi {
5960
forkCount: 0,
6061
ttlSeconds: params.ttlSeconds,
6162
failureReason: null,
63+
lastActivityAt: null,
6264
createdAt: now,
6365
updatedAt: now,
6466
startedAt: null,
@@ -149,6 +151,7 @@ export function createInMemorySandboxRepo(): SandboxRepoApi {
149151
const row: SandboxRow = {
150152
id: params.id,
151153
orgId: params.orgId,
154+
nodeId: params.source.nodeId,
152155
imageId: params.source.imageId,
153156
profileId: params.source.profileId,
154157
profileName: params.source.profileName,
@@ -159,6 +162,7 @@ export function createInMemorySandboxRepo(): SandboxRepoApi {
159162
forkCount: 0,
160163
ttlSeconds: params.ttlSeconds,
161164
failureReason: null,
165+
lastActivityAt: now,
162166
createdAt: now,
163167
updatedAt: now,
164168
startedAt: now,
@@ -211,6 +215,51 @@ export function createInMemorySandboxRepo(): SandboxRepoApi {
211215

212216
return result
213217
}),
218+
219+
findExpiredTtl: () =>
220+
Effect.sync(() => {
221+
const now = Date.now()
222+
return Array.from(store.values()).filter((r) => {
223+
if (r.status !== 'running' || !r.startedAt) return false
224+
return r.startedAt.getTime() + r.ttlSeconds * 1000 < now
225+
})
226+
}),
227+
228+
findIdleSince: (cutoff) =>
229+
Effect.sync(() =>
230+
Array.from(store.values()).filter((r) => {
231+
if (r.status !== 'running') return false
232+
const activity = r.lastActivityAt ?? r.startedAt ?? r.createdAt
233+
return activity.getTime() < cutoff.getTime()
234+
}),
235+
),
236+
237+
findQueuedBefore: (cutoff) =>
238+
Effect.sync(() =>
239+
Array.from(store.values()).filter(
240+
(r) => r.status === 'queued' && r.createdAt.getTime() < cutoff.getTime(),
241+
),
242+
),
243+
244+
getActiveNodeIds: () =>
245+
Effect.sync(() => {
246+
const nodeIds = new Map<string, Uint8Array>()
247+
for (const row of store.values()) {
248+
if (row.status === 'running' && row.nodeId) {
249+
nodeIds.set(base62Encode(row.nodeId), row.nodeId)
250+
}
251+
}
252+
return Array.from(nodeIds.values())
253+
}),
254+
255+
findRunningOnNodes: (nodeIds) =>
256+
Effect.sync(() => {
257+
if (nodeIds.length === 0) return []
258+
return Array.from(store.values()).filter((r) => {
259+
if (r.status !== 'running' || !r.nodeId) return false
260+
return nodeIds.some((nid) => bytesEqual(r.nodeId!, nid))
261+
})
262+
}),
214263
}
215264
}
216265

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

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import type {
1010
export interface SandboxRow {
1111
readonly id: Uint8Array
1212
readonly orgId: string
13+
readonly nodeId: Uint8Array | null
1314
readonly imageId: Uint8Array
1415
readonly profileId: Uint8Array
1516
readonly profileName: ProfileName
@@ -20,6 +21,7 @@ export interface SandboxRow {
2021
readonly forkCount: number
2122
readonly ttlSeconds: number
2223
readonly failureReason: FailureReason | null
24+
readonly lastActivityAt: Date | null
2325
readonly createdAt: Date
2426
readonly updatedAt: Date
2527
readonly startedAt: Date | null
@@ -99,6 +101,27 @@ export interface SandboxRepoApi {
99101
id: Uint8Array,
100102
orgId: string,
101103
) => Effect.Effect<SandboxRow[], never, never>
104+
105+
/** Find running sandboxes that have exceeded their TTL. */
106+
readonly findExpiredTtl: () => Effect.Effect<SandboxRow[], never, never>
107+
108+
/** Find running sandboxes with lastActivityAt before the given cutoff. */
109+
readonly findIdleSince: (
110+
cutoff: Date,
111+
) => Effect.Effect<SandboxRow[], never, never>
112+
113+
/** Find queued sandboxes created before the given cutoff. */
114+
readonly findQueuedBefore: (
115+
cutoff: Date,
116+
) => Effect.Effect<SandboxRow[], never, never>
117+
118+
/** Get distinct nodeIds from running sandboxes. */
119+
readonly getActiveNodeIds: () => Effect.Effect<Uint8Array[], never, never>
120+
121+
/** Find running sandboxes assigned to any of the given nodeIds. */
122+
readonly findRunningOnNodes: (
123+
nodeIds: Uint8Array[],
124+
) => Effect.Effect<SandboxRow[], never, never>
102125
}
103126

104127
export class SandboxRepo extends Context.Tag('SandboxRepo')<SandboxRepo, SandboxRepoApi>() {}

0 commit comments

Comments
 (0)