Skip to content

Commit 89449f8

Browse files
authored
feat: add health check plugin (#211)
2 parents c2a4026 + 94d3a9c commit 89449f8

16 files changed

Lines changed: 466 additions & 10 deletions
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
{
2+
"type": "prerelease",
3+
"comment": "feat: add abort signal support to indexer and runWithReconnect to support graceful shutdown",
4+
"packageName": "@apibara/indexer",
5+
"email": "jadejajaipal5@gmail.com",
6+
"dependentChangeType": "patch"
7+
}
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
{
2+
"type": "prerelease",
3+
"comment": "feat: add plugin-health for health check",
4+
"packageName": "@apibara/plugin-health",
5+
"email": "jadejajaipal5@gmail.com",
6+
"dependentChangeType": "patch"
7+
}
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
{
2+
"type": "prerelease",
3+
"comment": "feat: add abort signal support to indexer and runWithReconnect to support graceful shutdown",
4+
"packageName": "apibara",
5+
"email": "jadejajaipal5@gmail.com",
6+
"dependentChangeType": "patch"
7+
}

examples/cli-drizzle/indexers/2-starknet.indexer.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { starknetUsdcTransfers } from "@/lib/schema";
22
import { drizzleStorage, useDrizzleStorage } from "@apibara/plugin-drizzle";
33
import { drizzle } from "@apibara/plugin-drizzle";
4+
import { healthPlugin } from "@apibara/plugin-health";
45
import { StarknetStream } from "@apibara/starknet";
56
import { defineIndexer } from "apibara/indexer";
67
import { useLogger } from "apibara/plugins";
@@ -32,6 +33,7 @@ export default async function (runtimeConfig: ApibaraRuntimeConfig) {
3233
finality: "accepted",
3334
startingBlock: BigInt(startingBlock),
3435
plugins: [
36+
healthPlugin(),
3537
drizzleStorage({
3638
db: database,
3739
idColumn: {

examples/cli-drizzle/package.json

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,11 +27,13 @@
2727
"@apibara/evm": "workspace:*",
2828
"@apibara/indexer": "workspace:*",
2929
"@apibara/plugin-drizzle": "workspace:*",
30+
"@apibara/plugin-health": "workspace:*",
3031
"@apibara/protocol": "workspace:*",
3132
"@apibara/starknet": "workspace:*",
3233
"@electric-sql/pglite": "^0.2.17",
3334
"apibara": "workspace:*",
3435
"drizzle-orm": "^0.40.1",
36+
"elysia": "^1.4.27",
3537
"pg": "^8.13.1",
3638
"starknet": "^6.11.0",
3739
"viem": "^2.40.0"

examples/cli-js/indexers/starknet.indexer.mjs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { defineIndexer } from "@apibara/indexer";
2+
import { health } from "@apibara/plugin-health";
23
import { StarknetStream } from "@apibara/starknet";
34
import { hash } from "starknet";
45

@@ -15,6 +16,7 @@ export default function (config) {
1516
startingCursor: {
1617
orderKey: 1_380_000n,
1718
},
19+
plugins: [health()],
1820
filter: {
1921
events: [
2022
{

examples/cli-js/package.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
"@apibara/protocol": "workspace:*",
2222
"@apibara/starknet": "workspace:*",
2323
"@apibara/evm": "workspace:*",
24+
"@apibara/plugin-health": "workspace:*",
2425
"apibara": "workspace:*",
2526
"starknet": "^6.11.0"
2627
}

package.json

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,8 @@
3939
"packages/starknet",
4040
"packages/plugin-drizzle",
4141
"packages/plugin-mongo",
42-
"packages/plugin-sqlite"
42+
"packages/plugin-sqlite",
43+
"packages/plugin-health"
4344
]
4445
}
4546
],

packages/cli/src/runtime/dev.ts

Lines changed: 17 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import type { ConsolaInstance } from "consola";
66
import { blueBright } from "picocolors";
77
import { availableIndexers, createIndexer } from "./internal/app";
88

9-
async function startIndexer(indexer: string) {
9+
async function startIndexer(indexer: string, signal: AbortSignal) {
1010
let _logger: ConsolaInstance | undefined;
1111
while (true) {
1212
try {
@@ -35,7 +35,7 @@ async function startIndexer(indexer: string) {
3535
logger.info(`Indexer ${blueBright(indexer)} started`);
3636
}
3737

38-
await runWithReconnect(client, indexerInstance);
38+
await runWithReconnect(client, indexerInstance, { signal });
3939

4040
return;
4141
} catch (error) {
@@ -75,7 +75,21 @@ const startCommand = defineCommand({
7575
}
7676
}
7777

78-
await Promise.all(selectedIndexers.map((indexer) => startIndexer(indexer)));
78+
const abortController = new AbortController();
79+
const onSignal = () => abortController.abort();
80+
process.once("SIGINT", onSignal);
81+
process.once("SIGTERM", onSignal);
82+
83+
try {
84+
await Promise.all(
85+
selectedIndexers.map((indexer) =>
86+
startIndexer(indexer, abortController.signal),
87+
),
88+
);
89+
} finally {
90+
process.off("SIGINT", onSignal);
91+
process.off("SIGTERM", onSignal);
92+
}
7993
},
8094
});
8195

packages/indexer/src/indexer.ts

Lines changed: 65 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -181,6 +181,7 @@ export interface ReconnectOptions {
181181
maxRetries?: number;
182182
retryDelay?: number;
183183
maxWait?: number;
184+
signal?: AbortSignal;
184185
}
185186

186187
export async function runWithReconnect<TFilter, TBlock>(
@@ -193,10 +194,18 @@ export async function runWithReconnect<TFilter, TBlock>(
193194
const maxRetries = options.maxRetries ?? 10;
194195
const retryDelay = options.retryDelay ?? 1_000;
195196
const maxWait = options.maxWait ?? 30_000;
197+
const { signal } = options;
196198

197199
while (true) {
200+
if (signal?.aborted) return;
201+
198202
const abortController = new AbortController();
199203

204+
// Forward external abort into the per-run internal controller so that
205+
// run() sees the abort signal and can break its loop cleanly.
206+
const forwardAbort = () => abortController.abort(signal?.reason);
207+
signal?.addEventListener("abort", forwardAbort, { once: true });
208+
200209
const runOptions: RunOptions = {
201210
onConnect() {
202211
retryCount = 0;
@@ -206,15 +215,14 @@ export async function runWithReconnect<TFilter, TBlock>(
206215

207216
try {
208217
await run(client, indexer, runOptions);
209-
abortController.abort();
210218
return;
211219
} catch (error) {
212220
// Only reconnect on internal/server errors.
213221
// All other errors should be rethrown.
214222

215223
retryCount++;
216224

217-
abortController.abort();
225+
if (signal?.aborted) return;
218226

219227
if (error instanceof ClientError || error instanceof ServerError) {
220228
const isServerError = error instanceof ServerError;
@@ -231,15 +239,29 @@ export async function runWithReconnect<TFilter, TBlock>(
231239

232240
// Add jitter to the retry delay to avoid all clients retrying at the same time.
233241
const delay = Math.random() * (retryDelay * 0.2) + retryDelay;
234-
await new Promise((resolve) =>
235-
setTimeout(resolve, Math.min(retryCount * delay, maxWait)),
236-
);
242+
await new Promise<void>((resolve) => {
243+
const t = setTimeout(
244+
resolve,
245+
Math.min(retryCount * delay, maxWait),
246+
);
247+
// wake up immediately if externally aborted during the delay.
248+
signal?.addEventListener(
249+
"abort",
250+
() => {
251+
clearTimeout(t);
252+
resolve();
253+
},
254+
{ once: true },
255+
);
256+
});
237257

238258
continue;
239259
}
240260
}
241261
}
242262
throw error;
263+
} finally {
264+
signal?.removeEventListener("abort", forwardAbort);
243265
}
244266
}
245267
}
@@ -322,8 +344,29 @@ export async function run<TFilter, TBlock>(
322344

323345
let onConnectCalled = false;
324346

347+
const abortPromise = abortSignal ? waitForAbort(abortSignal) : null;
348+
325349
while (true) {
326-
const { value: message, done } = await stream.next();
350+
let result: IteratorResult<
351+
StreamDataResponse<TBlock>,
352+
StreamDataResponse<TBlock>
353+
>;
354+
355+
try {
356+
result = abortPromise
357+
? await Promise.race([stream.next(), abortPromise])
358+
: await stream.next();
359+
} catch (e) {
360+
if (abortSignal?.aborted) {
361+
// cancel the underlying gRPC stream and exit the loop cleanly so
362+
// that run:after hooks are still called below.
363+
await stream.return?.();
364+
break;
365+
}
366+
throw e;
367+
}
368+
369+
const { value: message, done } = result;
327370

328371
if (done) {
329372
break;
@@ -527,6 +570,22 @@ export async function run<TFilter, TBlock>(
527570
});
528571
}
529572

573+
/**
574+
* Returns a promise that rejects as soon as the given AbortSignal fires.
575+
* Used to race against stream.next() so the loop can break on abort.
576+
*/
577+
function waitForAbort(signal: AbortSignal): Promise<never> {
578+
return new Promise<never>((_, reject) => {
579+
if (signal.aborted) {
580+
reject(signal.reason);
581+
return;
582+
}
583+
signal.addEventListener("abort", () => reject(signal.reason), {
584+
once: true,
585+
});
586+
});
587+
}
588+
530589
async function registerMiddleware<TFilter, TBlock>(
531590
indexer: Indexer<TFilter, TBlock>,
532591
abortSignal?: AbortSignal,

0 commit comments

Comments
 (0)