Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
10 changes: 10 additions & 0 deletions indexer/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,16 @@ npm run dev
```
Counter triples (`pages_processed`, `events_folded`, `fold_failures`) are also emitted as structured JSON log lines on stderr (`msg=indexer_tick`) for dashboards without scraping.

## Failure handling

- **Fetch step (issue #569)** — the RPC `getEvents` call runs through `retryWithBackoff` (`src/retry.ts`): 3 attempts by default with exponential backoff (250ms → 5s cap), one structured warn line per attempt. Exhausted retries are logged and the cursor is left untouched; the next tick / loop iteration retries the same position, so a transient RPC error or Postgres blip delays progress instead of stopping it.
- **Ingest step** — `pollOnce` keeps ingest + cursor save in a single transaction; a failure rolls back and `startPollLoop` backs off `min(intervalMs, 2000)` before polling again. The loop never exits on a failure — only on abort or process death.
- **Placeholder source (issue #568)** — constructing `StubSorobanEventSource` and booting `src/worker.ts` each log a one-time warning that no events will be indexed, so a worker running the scaffold is never indistinguishable from a worker whose real event source is returning nothing.

## Derived tables

`db/schema.sql` ships `streams` (issue #567 — one row per stream: sender, recipient, token, rate, start/end) and `stream_events` (issue #566 — one row per DripStream lifecycle event: withdrawn, cancelled, paused, resumed, topped_up, clawback, xfer_rec), folded by `src/indexer/handlers.ts`. The DAO-voting tables (`loan_proposals`, `treasury_proposals`) are legacy scaffolding.

## Contract

See `src/indexer/types.ts` for the `SorobanEventSource` pagination contract (`nextToken` opaque, `lastLedger` inclusive) and the expected `fields` shape per `ev.type`. That file is the entire interface between the poller and the event source implementation that will replace `StubSorobanEventSource`.
59 changes: 55 additions & 4 deletions indexer/db/schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,14 @@
--
-- Two layers:
-- 1. Raw event log (`raw_events`) — append-only, idempotent via ON CONFLICT DO NOTHING.
-- 2. Derived tables (`loan_proposals`, `treasury_proposals`) — folded from raw events by
-- src/indexer/handlers.ts. Folding is currently additive (increments), which is NOT
-- idempotent — see the "Known gaps" section of the top-level README before relying on
-- these tallies for anything safety-critical.
-- 2. Derived tables (`streams`, `stream_events`, `loan_proposals`, `treasury_proposals`) —
-- folded from raw events by src/indexer/handlers.ts. The stream folds are idempotent
-- upserts; the DAO-voting folds are additive (increments) and are NOT idempotent —
-- see the "Known gaps" section of the top-level README before relying on those tallies
-- for anything safety-critical.
--
-- `streams` / `stream_events` cover the protocol the repo actually ships (DripStream /
-- DripFactory — issues #566, #567); the DAO-voting tables are legacy scaffolding.

CREATE TABLE IF NOT EXISTS indexer_cursor (
id integer PRIMARY KEY,
Expand All @@ -29,6 +33,53 @@ CREATE TABLE IF NOT EXISTS raw_events (

CREATE INDEX IF NOT EXISTS raw_events_ledger_idx ON raw_events (ledger);

-- Streams — one row per DripStream, folded from the stream-creation event
-- (`DripStream::created`, emitted by DripFactory::create_stream deployments;
-- issue #567). Answers "which streams exist / who created them / what were
-- the initial parameters" without scanning and decoding raw_events.
--
-- `id` is the stream's identifier as delivered by the event decoder: the
-- factory's monotonically increasing stream id when the decoder supplies one,
-- otherwise the DripStream contract address (which is unique per stream).
CREATE TABLE IF NOT EXISTS streams (
id text PRIMARY KEY,
sender text NOT NULL,
recipient text NOT NULL,
token text NOT NULL,
rate_per_second text NOT NULL, -- i128 decimal string — exceeds bigint on-chain
start_time bigint NOT NULL,
end_time bigint NOT NULL, -- 0 = open-ended
created_at timestamptz NOT NULL DEFAULT now()
);

CREATE INDEX IF NOT EXISTS streams_created_at_idx ON streams (created_at);
CREATE INDEX IF NOT EXISTS streams_sender_idx ON streams (sender);
CREATE INDEX IF NOT EXISTS streams_recipient_idx ON streams (recipient);

-- Stream lifecycle events — one row per (ledger, tx, type, stream) for the
-- DripStream events the README documents (withdrawn, cancelled, paused,
-- resumed, topped_up, clawback, xfer_rec; issue #566). Unlike the additive
-- DAO tallies this is an event log, so its upsert is "same event seen again →
-- rewrite the same row" (idempotent under page re-delivery).
--
-- `amount` is the event's i128 decimal string (withdrawn amount / refund /
-- top-up / clawback) when the event carries one; `payload` keeps every field
-- for events without one and as the audit copy for events with one.
CREATE TABLE IF NOT EXISTS stream_events (
stream_id text NOT NULL,
event_type text NOT NULL,
ledger bigint NOT NULL,
tx_hash text NOT NULL,
amount text,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (ledger, tx_hash, event_type, stream_id)
);

CREATE INDEX IF NOT EXISTS stream_events_stream_id_idx ON stream_events (stream_id);
CREATE INDEX IF NOT EXISTS stream_events_created_at_idx ON stream_events (created_at);
CREATE INDEX IF NOT EXISTS stream_events_type_idx ON stream_events (event_type);

CREATE TABLE IF NOT EXISTS loan_proposals (
id bigint PRIMARY KEY,
votes_for bigint NOT NULL DEFAULT 0,
Expand Down
20 changes: 20 additions & 0 deletions indexer/src/indexer/eventSource.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,11 +8,31 @@
* - `lastLedger` inclusive semantics
* - opaque `nextToken` round-tripping
* - `events` sorted ascending by `(ledger, sequence)`
*
* Loud-by-default (issue #568): constructing this stub logs a one-time
* warning. Before that, a worker booted against this stub "succeeded" —
* clean start, healthy `/healthz`, zero log lines, zero indexed events —
* indistinguishable from a broken RPC returning nothing. Now the intent is
* stated once, at construction, and silence after that means "no events on
* chain", not "the stub is still wired in".
*/

import { Page, GetEventsParams, SorobanEventSource } from "./types.js";

export class StubSorobanEventSource implements SorobanEventSource {
/** One warning per construction — never per `getEvents` call. */
constructor() {
console.warn(
JSON.stringify({
level: "warn",
msg: "using placeholder SorobanEventSource — no events will be indexed",
stub: "StubSorobanEventSource",
file: "indexer/src/indexer/eventSource.ts",
hint: "replace with an RPC-backed SorobanEventSource to index real events",
})
);
}

async getEvents(params: GetEventsParams): Promise<Page> {
const end = params.endLedger ?? params.startLedger;
return {
Expand Down
176 changes: 168 additions & 8 deletions indexer/src/indexer/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,76 @@ function num(value: unknown): number {
return typeof value === 'number' ? value : Number(value)
}

function str(value: unknown): string {
return value === undefined || value === null ? '' : String(value)
}

/** i128 amount field → decimal string, or null when the event carries none. */
function amountOrNull(value: unknown): string | null {
if (value === undefined || value === null || value === '') return null
return String(value)
}

/**
* Stream identity for a stream event: the decoder-supplied id when it has one
* (factory `stream_id`), otherwise the emitting DripStream contract address —
* which is unique per stream, so `streams.id` / `stream_events.stream_id` are
* always populated.
*/
function streamIdOf(ev: ChainEvent): string {
const f = ev.fields
return str(f.stream_id ?? f.streamId ?? f.id ?? ev.contractId)
}

/**
* Folds a single decoded event into the derived tables.
*
* NOTE: these folds are additive (increments), not idempotent upserts. Re-delivering an
* already-applied event double-counts it. That's a known gap tracked separately — see
* "Non-idempotent derived-table folds" in the top-level README — and is intentionally not
* fixed as part of this scaffold.
* Two families of folds:
* - DripStream / DripFactory events (`streams`, `stream_events`) — issues
* #566/#567. Idempotent upserts: re-delivering a page rewrites the same
* row, so a crash between ingest and checkpoint cannot double-count.
* - DAO-voting events (`loan_proposals`, `treasury_proposals`) — legacy
* scaffolding. These folds are additive (increments), not idempotent.
* Re-delivering an already-applied event double-counts it. That's a known
* gap tracked separately — see "Non-idempotent derived-table folds" in the
* top-level README — and is intentionally not fixed as part of this scaffold.
*
* Type tags: the switch accepts both the README event names
* (`stream_withdrawn`, `stream_cancelled`, …) and the raw `symbol_short!`
* topic names emitted on-chain (`withdrawn`, `cancelled`, `created`, …) as
* documented on {@link ChainEvent.type} / `types.ts:EventType`, and always
* stores the canonical README name in `stream_events.event_type`.
*/
export async function applyEvent(client: PoolClient, ev: ChainEvent): Promise<void> {
switch (ev.type) {
// Stream creation (issue #567) — `DripStream::created`.
case 'created':
case 'stream_created':
return handleStreamCreated(client, ev)
// Stream lifecycle (issue #566) — README's DripStream events table.
case 'withdrawn':
case 'stream_withdrawn':
return handleStreamWithdrawn(client, ev)
case 'cancelled':
case 'stream_cancelled':
return handleStreamCancelled(client, ev)
case 'force_cxl':
case 'force_cancelled':
return handleStreamForceCancelled(client, ev)
case 'paused':
case 'stream_paused':
return handleStreamPaused(client, ev)
case 'resumed':
case 'stream_resumed':
return handleStreamResumed(client, ev)
case 'topped_up':
case 'stream_topped_up':
return handleStreamToppedUp(client, ev)
case 'clawback':
case 'stream_clawback':
return handleStreamClawback(client, ev)
case 'xfer_rec':
return handleXferRec(client, ev)
case 'loan_vote':
return loan_vote(client, ev)
case 'treasury_vote':
Expand All @@ -23,10 +83,10 @@ export async function applyEvent(client: PoolClient, ev: ChainEvent): Promise<vo
return treasury_reveal(client, ev)
default:
// Log (don't throw) on anything we don't recognise: today that's every
// event type except the three DAO-voting folds above. Once real
// DripStream/DripFactory handlers land here, a future contract upgrade
// that adds a new event type must show up in the logs instead of
// disappearing into this indexer silently (issue #578).
// event type except the folds above. Once real DripStream/DripFactory
// handlers land here, a future contract upgrade that adds a new event
// type must show up in the logs instead of disappearing into this
// indexer silently (issue #578).
console.error(
JSON.stringify({
level: 'warn',
Expand All @@ -41,6 +101,106 @@ export async function applyEvent(client: PoolClient, ev: ChainEvent): Promise<vo
}
}

// ---------------------------------------------------------------------------
// Stream handlers (issues #566, #567)
// ---------------------------------------------------------------------------

/**
* Stream creation → `streams` (issue #567).
*
* `created_at` is left to the column default: on a re-delivered event the
* ON CONFLICT clause rewrites the parameters but never the creation time, so
* "streams created in the last 24h" stays stable across replays.
*/
export async function handleStreamCreated(client: PoolClient, ev: ChainEvent): Promise<void> {
const f = ev.fields
await client.query(
`INSERT INTO streams (id, sender, recipient, token, rate_per_second, start_time, end_time)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (id) DO UPDATE
SET sender = EXCLUDED.sender,
recipient = EXCLUDED.recipient,
token = EXCLUDED.token,
rate_per_second = EXCLUDED.rate_per_second,
start_time = EXCLUDED.start_time,
end_time = EXCLUDED.end_time`,
[
streamIdOf(ev),
str(f.sender),
str(f.recipient),
str(f.token),
str(f.rate_per_second ?? 0),
num(f.start_time ?? 0),
num(f.end_time ?? 0),
],
)
}

/**
* Append one row to `stream_events` (issue #566).
*
* The primary key `(ledger, tx_hash, event_type, stream_id)` is exactly the
* page-re-delivery key: folding the same event twice rewrites the same row
* (same `amount`, same `payload`) instead of appending a duplicate.
*/
async function recordStreamEvent(
client: PoolClient,
ev: ChainEvent,
eventType: string,
amount: string | null = null,
): Promise<void> {
await client.query(
`INSERT INTO stream_events (stream_id, event_type, ledger, tx_hash, amount, payload)
VALUES ($1, $2, $3, $4, $5, $6::jsonb)
ON CONFLICT (ledger, tx_hash, event_type, stream_id) DO UPDATE
SET amount = EXCLUDED.amount,
payload = EXCLUDED.payload`,
[
streamIdOf(ev),
eventType,
num(ev.ledger ?? 0),
str(ev.txHash),
amount,
JSON.stringify(ev.fields ?? {}),
],
)
}

export async function handleStreamWithdrawn(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'stream_withdrawn', amountOrNull(ev.fields.amount))
}

export async function handleStreamCancelled(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'stream_cancelled', amountOrNull(ev.fields.refund_amount))
}

export async function handleStreamForceCancelled(client: PoolClient, ev: ChainEvent): Promise<void> {
// `force_cxl` is recorded under its own event_type (per types.ts) so a
// consumer can tell sender-cancel from recipient force-cancel without
// correlating the transaction signer.
return recordStreamEvent(client, ev, 'stream_force_cancelled', amountOrNull(ev.fields.refund_amount))
}

export async function handleStreamPaused(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'stream_paused')
}

export async function handleStreamResumed(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'stream_resumed')
}

export async function handleStreamToppedUp(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'stream_topped_up', amountOrNull(ev.fields.amount))
}

export async function handleStreamClawback(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'stream_clawback', amountOrNull(ev.fields.amount))
}

export async function handleXferRec(client: PoolClient, ev: ChainEvent): Promise<void> {
return recordStreamEvent(client, ev, 'xfer_rec')
}

async function loan_vote(client: PoolClient, ev: ChainEvent): Promise<void> {
const f = ev.fields
const column = f.support === true ? 'votes_for' : 'votes_against'
Expand Down
30 changes: 27 additions & 3 deletions indexer/src/indexer/poller.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,16 @@
*
* Health is reported via `health.ts` — `lastSuccessfulPollTimestamp` and
* `currentCursor` power `GET /healthz`.
*
* Transient failures (issue #569): the `getEvents` fetch step runs through
* `retryWithBackoff` (`../retry.ts`, exponential, default 3 attempts). If all
* attempts fail the cursor is left untouched and the next tick retries it —
* the loop never silently stops advancing because of one bad RPC response.
*/

import { Cursor, SorobanEventSource, GetEventsParams } from "./types.js";
import { fold } from "./fold.js";
import { retryWithBackoff, DEFAULT_RETRY, type RetryPolicy } from "../retry.js";
import { metrics } from "../metrics.js";
import { health } from "../health.js";

Expand All @@ -29,6 +35,13 @@ export interface PollerOptions {
intervalMs?: number;
/** Max events per page. Default 100 (matches on-chain MAX_PAGE_SIZE). */
limit?: number;
/**
* Retry-with-backoff policy for the `getEvents` fetch step (issue #569).
* Default: 3 attempts, 250ms base doubling to a 5s cap. Exhausted retries
* are logged as an error and the *next tick* retries the same cursor —
* a transient RPC/DB blip never stalls the loop silently.
*/
retry?: RetryPolicy;
/** Optional loader/saver for the high-water cursor (e.g. Postgres or file). */
loadCursor?: () => Promise<Cursor | null>;
saveCursor?: (cursor: Cursor) => Promise<void>;
Expand Down Expand Up @@ -92,11 +105,22 @@ export class Poller {
};

let page;
const attempts = Math.max(1, this.opts.retry?.attempts ?? DEFAULT_RETRY.attempts);
try {
page = await this.opts.source.getEvents(params);
// Explicit retry-with-backoff around the RPC fetch step (issue #569).
// Retry lines are logged by retryWithBackoff; a single transient
// failure therefore costs `delay`, not the rest of the stream.
page = await retryWithBackoff(() => this.opts.source.getEvents(params), {
...this.opts.retry,
label: "getEvents",
});
} catch (e) {
// Transport / RPC error — log, do not advance cursor, backoff is interval-based
console.error(JSON.stringify({ level: "error", msg: "getEvents failed", error: String(e), startLedger, cursor: this.cursor.nextToken }));
// Every attempt in this tick failed. Do not advance the cursor: the
// next tick rebuilds the same params from the unchanged cursor and
// tries again, so a transient failure only delays progress — the loop
// keeps running (interval ticks are scheduled independently of a
// thrown tick, see start()) and cannot stall silently.
console.error(JSON.stringify({ level: "error", msg: "getEvents failed after retries", error: String(e), attempts, startLedger, cursor: this.cursor.nextToken }));
metrics.logStructured({ phase: "fetch_error", startLedger });
return;
}
Expand Down
Loading