Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
274e4e9
Add Starknet RPC event package
loothero Jun 14, 2026
4226e13
Use positional Starknet HTTP RPC params
loothero Jun 15, 2026
0f99a77
Order Starknet event cursors by transaction index
loothero Jun 15, 2026
b48a9df
Remove app-specific Starknet RPC README notes
loothero Jun 15, 2026
ac183dc
Move Starknet RPC tests out of source
loothero Jun 15, 2026
4e4f6dd
Add Starknet RPC change file
loothero Jun 15, 2026
e873117
Default Starknet RPC live events to pre-confirmed
loothero Jun 15, 2026
e15ad12
Fix Starknet getEvents filter parameters
loothero Jun 15, 2026
c78e375
Document Starknet RPC event cursor requirements
loothero Jun 15, 2026
39d605b
Recover Starknet stream from stale WS handoffs
loothero Jun 15, 2026
844c894
Harden Starknet event subscriptions
loothero Jun 15, 2026
aa74457
Simplify Starknet block RPC helper
loothero Jun 15, 2026
651e47a
Document Starknet stream recovery behavior
loothero Jun 15, 2026
ce4dbb8
Harden Starknet subscription API
loothero Jun 15, 2026
1614022
Cover Starknet stream HTTP retry
loothero Jun 15, 2026
f535093
Handle Starknet preconfirmed cursor rollbacks
loothero Jun 15, 2026
cdc08b4
Replay preconfirmed cursor blocks on stream resume
loothero Jun 15, 2026
818c138
Add Starknet RPC stream config
loothero Jun 15, 2026
1846a79
Harden Starknet RPC WebSocket subscriptions
loothero Jun 15, 2026
aa546c7
Default Starknet backfill to accepted latest
loothero Jun 15, 2026
cba5d28
Avoid accumulating HTTP abort listeners
loothero Jun 15, 2026
1679ad6
docs: simplify starknet rpc readme
loothero Jun 16, 2026
64a7a99
chore: narrow rpc stream integration changes
loothero Jun 16, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions change/@apibara-indexer-rpc-stream-config.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
{
"type": "patch",
"comment": "Allow indexers to use RPC stream configs",
"packageName": "@apibara/indexer",
"email": "loothero@provable.games",
"dependentChangeType": "patch"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
{
"type": "minor",
"comment": "Add Starknet RPC event helpers",
"packageName": "@apibara/starknet-rpc",
"email": "loothero@provable.games",
"dependentChangeType": "patch"
}
7 changes: 7 additions & 0 deletions change/apibara-rpc-stream-runtime-client.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
{
"type": "patch",
"comment": "Support RPC stream configs in the CLI runtime",
"packageName": "apibara",
"email": "loothero@provable.games",
"dependentChangeType": "patch"
}
12 changes: 6 additions & 6 deletions packages/cli/src/runtime/dev.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
import { ReloadIndexerRequest, runWithReconnect } from "@apibara/indexer";
import { createAuthenticatedClient } from "@apibara/protocol";
import { getRuntimeDataFromEnv } from "apibara/common";
import { defineCommand, runMain } from "citty";
import type { ConsolaInstance } from "consola";
import { blueBright } from "picocolors";
import { availableIndexers, createIndexer } from "./internal/app";
import { createRuntimeClient } from "./internal/client";

async function startIndexer(indexer: string, signal: AbortSignal) {
let _logger: ConsolaInstance | undefined;
Expand All @@ -25,11 +25,11 @@ async function startIndexer(indexer: string, signal: AbortSignal) {
return;
}

const client = createAuthenticatedClient(
indexerInstance.streamConfig,
indexerInstance.options.streamUrl,
indexerInstance.options.clientOptions,
);
const client = createRuntimeClient({
streamConfig: indexerInstance.streamConfig,
streamUrl: indexerInstance.options.streamUrl,
clientOptions: indexerInstance.options.clientOptions,
});

if (logger) {
logger.info(`Indexer ${blueBright(indexer)} started`);
Expand Down
27 changes: 27 additions & 0 deletions packages/cli/src/runtime/internal/client.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
import type { IndexerStreamConfig } from "@apibara/indexer";
import {
type Client,
type CreateClientOptions,
createAuthenticatedClient,
} from "@apibara/protocol";
import { createRpcClient } from "@apibara/protocol/rpc";

export function createRuntimeClient<TFilter, TBlock>({
streamConfig,
streamUrl,
clientOptions,
}: {
streamConfig: IndexerStreamConfig<TFilter, TBlock>;
streamUrl?: string;
clientOptions?: CreateClientOptions;
}): Client<TFilter, TBlock> {
if ("Request" in streamConfig) {
if (!streamUrl) {
throw new Error("streamUrl is required when using a DNA StreamConfig");
}

return createAuthenticatedClient(streamConfig, streamUrl, clientOptions);
}

return createRpcClient(streamConfig);
}
25 changes: 24 additions & 1 deletion packages/cli/src/runtime/project-info.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ const startCommand = defineCommand({
projectInfo.indexers[indexer] = {
...(projectInfo.indexers[indexer] ?? {}),
[preset]: {
type: indexerInstance.streamConfig.name,
type: streamConfigType(indexerInstance.streamConfig),
isFactory: indexerInstance.options.factory !== undefined,
},
};
Expand All @@ -75,6 +75,29 @@ const startCommand = defineCommand({
},
});

function streamConfigType(streamConfig: unknown): string {
if (
typeof streamConfig === "object" &&
streamConfig !== null &&
"name" in streamConfig &&
typeof streamConfig.name === "string"
) {
return streamConfig.name;
}

if (
typeof streamConfig === "object" &&
streamConfig !== null &&
"constructor" in streamConfig &&
typeof streamConfig.constructor === "function" &&
streamConfig.constructor.name
) {
return streamConfig.constructor.name;
}

return "RpcStreamConfig";
}

export const mainCli = defineCommand({
meta: {
name: "write-project-info-runner",
Expand Down
12 changes: 6 additions & 6 deletions packages/cli/src/runtime/start.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { ReloadIndexerRequest, runWithReconnect } from "@apibara/indexer";
import { createAuthenticatedClient } from "@apibara/protocol";
import {
checkForUnknownArgs,
getProcessedRuntimeConfig,
Expand All @@ -17,6 +16,7 @@ import {
userEnvRuntimeConfig,
} from "#apibara-internal-virtual/static-config";
import { createIndexer } from "./internal/app";
import { createRuntimeClient } from "./internal/client";

const startCommand = defineCommand({
meta: {
Expand Down Expand Up @@ -83,11 +83,11 @@ const startCommand = defineCommand({
process.exit(1);
}

const client = createAuthenticatedClient(
indexerInstance.streamConfig,
indexerInstance.options.streamUrl,
indexerInstance.options.clientOptions,
);
const client = createRuntimeClient({
streamConfig: indexerInstance.streamConfig,
streamUrl: indexerInstance.options.streamUrl,
clientOptions: indexerInstance.options.clientOptions,
});

if (register) {
consola.start("Registering from instrumentation");
Expand Down
74 changes: 65 additions & 9 deletions packages/indexer/src/indexer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {
type StreamDataResponse,
type SystemMessage,
} from "@apibara/protocol";
import type { RpcStreamConfig } from "@apibara/protocol/rpc";
import consola from "consola";
import {
type Hookable,
Expand Down Expand Up @@ -116,8 +117,7 @@ export type HandlerArgs<TBlock> = {
abortSignal?: AbortSignal;
};

export type IndexerConfig<TFilter, TBlock> = {
streamUrl: string;
export type BaseIndexerConfig<TFilter, TBlock> = {
filter: TFilter;
finality?: DataFinality;
clientOptions?: CreateClientOptions;
Expand All @@ -130,27 +130,68 @@ export type IndexerConfig<TFilter, TBlock> = {
debug?: boolean;
} & IndexerStartingCursor;

export type IndexerWithStreamConfig<TFilter, TBlock> = IndexerConfig<
export type IndexerConfig<TFilter, TBlock> = BaseIndexerConfig<
TFilter,
TBlock
> & {
streamUrl: string;
};

export type RpcIndexerConfig<TFilter, TBlock> = BaseIndexerConfig<
TFilter,
TBlock
> & {
streamUrl?: string;
};

export type IndexerStreamConfig<TFilter, TBlock> =
| StreamConfig<TFilter, TBlock>
| RpcStreamConfig<TFilter, TBlock>;

export type IndexerWithDnaStreamConfig<TFilter, TBlock> = IndexerConfig<
TFilter,
TBlock
> & {
streamConfig: StreamConfig<TFilter, TBlock>;
};

export type IndexerWithRpcStreamConfig<TFilter, TBlock> = RpcIndexerConfig<
TFilter,
TBlock
> & {
streamConfig: RpcStreamConfig<TFilter, TBlock>;
};

export type IndexerWithStreamConfig<TFilter, TBlock> =
| IndexerWithDnaStreamConfig<TFilter, TBlock>
| IndexerWithRpcStreamConfig<TFilter, TBlock>;

export function defineIndexer<TFilter, TBlock>(
streamConfig: StreamConfig<TFilter, TBlock>,
) {
): (
config: IndexerConfig<TFilter, TBlock>,
) => IndexerWithDnaStreamConfig<TFilter, TBlock>;

export function defineIndexer<TFilter, TBlock>(
streamConfig: RpcStreamConfig<TFilter, TBlock>,
): (
config: RpcIndexerConfig<TFilter, TBlock>,
) => IndexerWithRpcStreamConfig<TFilter, TBlock>;

export function defineIndexer<TFilter, TBlock>(
streamConfig: IndexerStreamConfig<TFilter, TBlock>,
): unknown {
return (
config: IndexerConfig<TFilter, TBlock>,
): IndexerWithStreamConfig<TFilter, TBlock> => ({
config: IndexerConfig<TFilter, TBlock> | RpcIndexerConfig<TFilter, TBlock>,
) => ({
streamConfig,
...config,
});
}

export interface Indexer<TFilter, TBlock> {
streamConfig: StreamConfig<TFilter, TBlock>;
options: IndexerConfig<TFilter, TBlock>;
streamConfig: IndexerStreamConfig<TFilter, TBlock>;
options: IndexerConfig<TFilter, TBlock> | RpcIndexerConfig<TFilter, TBlock>;
hooks: Hookable<IndexerHooks<TFilter, TBlock>>;
}

Expand All @@ -177,6 +218,20 @@ export function createIndexer<TFilter, TBlock>({
return indexer;
}

function mergeFilterWithStreamConfig<TFilter, TBlock>(
streamConfig: IndexerStreamConfig<TFilter, TBlock>,
a: TFilter,
b: TFilter,
): TFilter {
if ("mergeFilter" in streamConfig) {
return streamConfig.mergeFilter(a, b);
}

throw new Error(
"Factory mode requires a stream config with mergeFilter support.",
);
}

export interface ReconnectOptions {
maxRetries?: number;
retryDelay?: number;
Expand Down Expand Up @@ -434,7 +489,8 @@ export async function run<TFilter, TBlock>(
if (filter) {
// when filter is defined
// merge old and new filters
mainFilter = indexer.streamConfig.mergeFilter(
mainFilter = mergeFilterWithStreamConfig(
indexer.streamConfig,
mainFilter,
filter,
);
Expand Down
8 changes: 8 additions & 0 deletions packages/indexer/src/testing/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,14 @@ export function createVcr() {
throw new Error("Cannot record cassette in CI");
}

if (!("Request" in indexer.streamConfig)) {
throw new Error("VCR recording requires a DNA StreamConfig");
}

if (!indexer.options.streamUrl) {
throw new Error("VCR recording requires streamUrl");
}

const client = createAuthenticatedClient(
indexer.streamConfig,
indexer.options.streamUrl,
Expand Down
Loading
Loading