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
122 changes: 116 additions & 6 deletions analytics-api/migrations/0001_analytics_schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -21,22 +21,132 @@ CREATE TABLE IF NOT EXISTS proposals (
PRIMARY KEY (contract_id, proposal_id)
);

CREATE INDEX IF NOT EXISTS proposals_created_at_idx ON proposals (created_at DESC);
CREATE INDEX IF NOT EXISTS proposals_status_category_idx ON proposals (status, category);
CREATE INDEX IF NOT EXISTS proposals_proposer_idx ON proposals (proposer);

CREATE TABLE IF NOT EXISTS events (
contract_id TEXT NOT NULL,
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
proposal_id BIGINT,
schedule_id BIGINT,
topic TEXT NOT NULL,
actor TEXT NOT NULL DEFAULT '',
occurred_at TIMESTAMPTZ NOT NULL,
data JSONB NOT NULL DEFAULT '{}'::jsonb,
PRIMARY KEY (contract_id, ledger, tx_hash, event_index)
);

CREATE INDEX IF NOT EXISTS events_proposal_timeline_idx
ON events (contract_id, proposal_id, ledger, event_index);
CREATE TABLE IF NOT EXISTS proposal_approvals (
contract_id TEXT NOT NULL,
proposal_id BIGINT NOT NULL CHECK (proposal_id >= 0),
owner TEXT NOT NULL,
weight INTEGER NOT NULL CHECK (weight >= 0),
approval_count INTEGER NOT NULL CHECK (approval_count >= 0),
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
PRIMARY KEY (contract_id, proposal_id, owner, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS proposal_executions (
contract_id TEXT NOT NULL,
proposal_id BIGINT NOT NULL CHECK (proposal_id >= 0),
executor TEXT NOT NULL,
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
data JSONB NOT NULL DEFAULT '{}'::jsonb,
PRIMARY KEY (contract_id, proposal_id, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS transfers (
contract_id TEXT NOT NULL,
proposal_id BIGINT,
schedule_id BIGINT,
recipient TEXT NOT NULL DEFAULT '',
token TEXT NOT NULL DEFAULT 'XLM',
amount TEXT NOT NULL DEFAULT '0',
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
data JSONB NOT NULL DEFAULT '{}'::jsonb,
PRIMARY KEY (contract_id, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS owner_weight_changes (
contract_id TEXT NOT NULL,
owner TEXT NOT NULL,
old_weight INTEGER NOT NULL CHECK (old_weight >= 0),
new_weight INTEGER NOT NULL CHECK (new_weight >= 0),
new_total_weight INTEGER NOT NULL CHECK (new_total_weight >= 0),
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
PRIMARY KEY (contract_id, owner, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS recurring_disbursements (
contract_id TEXT NOT NULL,
schedule_id BIGINT NOT NULL CHECK (schedule_id >= 0),
recipient TEXT NOT NULL DEFAULT '',
token TEXT NOT NULL DEFAULT 'XLM',
amount TEXT NOT NULL DEFAULT '0',
total_disbursed TEXT NOT NULL DEFAULT '0',
periods_disbursed INTEGER NOT NULL DEFAULT 0 CHECK (periods_disbursed >= 0),
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
PRIMARY KEY (contract_id, schedule_id, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS delegations (
contract_id TEXT NOT NULL,
delegator TEXT NOT NULL,
delegate TEXT NOT NULL,
weight INTEGER NOT NULL CHECK (weight >= 0),
expiry TIMESTAMPTZ,
is_active BOOLEAN NOT NULL DEFAULT true,
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
PRIMARY KEY (contract_id, delegator, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS role_changes (
contract_id TEXT NOT NULL,
target TEXT NOT NULL,
role_name TEXT NOT NULL,
action TEXT NOT NULL CHECK (action IN ('granted', 'revoked')),
ledger BIGINT NOT NULL CHECK (ledger >= 0),
tx_hash TEXT NOT NULL,
event_index INTEGER NOT NULL CHECK (event_index >= 0),
occurred_at TIMESTAMPTZ NOT NULL,
data JSONB NOT NULL DEFAULT '{}'::jsonb,
PRIMARY KEY (contract_id, target, role_name, ledger, tx_hash, event_index)
);

CREATE TABLE IF NOT EXISTS checkpoints (
contract_id TEXT PRIMARY KEY,
last_ledger BIGINT NOT NULL DEFAULT 0 CHECK (last_ledger >= 0),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX IF NOT EXISTS proposals_created_at_idx ON proposals (created_at DESC);
CREATE INDEX IF NOT EXISTS proposals_status_category_idx ON proposals (status, category);
CREATE INDEX IF NOT EXISTS proposals_proposer_idx ON proposals (proposer);
CREATE INDEX IF NOT EXISTS events_contract_ledger_event_idx ON events (contract_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS events_proposal_timeline_idx ON events (contract_id, proposal_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS events_schedule_timeline_idx ON events (contract_id, schedule_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS approvals_proposal_idx ON proposal_approvals (contract_id, proposal_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS executions_proposal_idx ON proposal_executions (contract_id, proposal_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS transfers_proposal_idx ON transfers (contract_id, proposal_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS transfers_schedule_idx ON transfers (contract_id, schedule_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS owner_weight_changes_owner_idx ON owner_weight_changes (contract_id, owner, ledger, event_index);
CREATE INDEX IF NOT EXISTS recurring_disbursements_schedule_idx ON recurring_disbursements (contract_id, schedule_id, ledger, event_index);
CREATE INDEX IF NOT EXISTS delegations_delegate_idx ON delegations (contract_id, delegate, ledger, event_index);
CREATE INDEX IF NOT EXISTS role_changes_target_idx ON role_changes (contract_id, target, ledger, event_index);
30 changes: 27 additions & 3 deletions frontend/src/lib/__tests__/contract-events.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,18 @@
import { describe, test, expect, vi, beforeEach } from "vitest";
import { getLatestLedger, getContractEvents, mapProposal } from "../contract";
import {
getLatestLedger,
getContractEvents,
getProposalEvents,
mapProposal,
dedupeContractEvents,
} from "../contract";
import { rpc } from "@stellar/stellar-sdk";

const { mockGetLatestLedger, mockGetEvents } = vi.hoisted(() => ({
mockGetLatestLedger: vi.fn(),
mockGetEvents: vi.fn(),
}));

// Mock the rpc.Server instance directly through vi
vi.mock("@stellar/stellar-sdk", async (importOriginal) => {
const actual: any = await importOriginal();
Expand Down Expand Up @@ -213,7 +224,20 @@ describe("Contract Events API", () => {
expect(proposal.to).toBe("GOWNER...1111");
expect(proposal.amount).toBe("25");
expect(proposal.token).toBe("Owner weight");
});
});

test("deduplicates overlapping event polls by ledger, tx hash, and event index", () => {
const events = [
{ ledger: 200, txHash: "tx-b", eventIndex: 2, value: "later" },
{ ledger: 100, txHash: "tx-a", eventIndex: 1, value: "first" },
{ ledger: 100, txHash: "tx-a", eventIndex: 1, value: "duplicate" },
{ ledger: 150, txHash: "tx-c", eventIndex: 0, value: "middle" },
];

// TODO: Add a test proving that replaying a ledger range never double-counts events (idempotency)
const deduped = dedupeContractEvents(events);

expect(deduped).toHaveLength(3);
expect(deduped.map((event) => event.value)).toEqual(["first", "middle", "later"]);
});
});

46 changes: 44 additions & 2 deletions frontend/src/lib/contract.ts
Original file line number Diff line number Diff line change
Expand Up @@ -823,6 +823,47 @@ function formatEventTimestamp(
return "Just now";
}

export function canonicalEventKey(event: {
ledger?: number | string;
txHash?: string;
tx_hash?: string;
eventIndex?: number | string | bigint;
event_index?: number | string | bigint;
}): string {
const ledger = Number(event?.ledger ?? 0);
const txHash = String(event?.txHash ?? event?.tx_hash ?? "");
const eventIndex = Number(event?.eventIndex ?? event?.event_index ?? 0);
return `${ledger}:${txHash}:${eventIndex}`;
}

export function dedupeContractEvents<T extends {
ledger?: number | string;
txHash?: string;
tx_hash?: string;
eventIndex?: number | string | bigint;
event_index?: number | string | bigint;
}>(events: T[]): T[] {
const seen = new Set<string>();
return [...events]
.filter((event) => {
const key = canonicalEventKey(event);
if (seen.has(key)) return false;
seen.add(key);
return true;
})
.sort((a, b) => {
const ledgerDelta = Number(a.ledger ?? 0) - Number(b.ledger ?? 0);
if (ledgerDelta !== 0) return ledgerDelta;
const eventDelta =
Number(a.eventIndex ?? a.event_index ?? 0) -
Number(b.eventIndex ?? b.event_index ?? 0);
if (eventDelta !== 0) return eventDelta;
return String(a.txHash ?? a.tx_hash ?? "").localeCompare(
String(b.txHash ?? b.tx_hash ?? ""),
);
});
}

function resolveEventType(
first: string,
second: string,
Expand Down Expand Up @@ -945,9 +986,10 @@ export async function getProposalEvents(
});

const events: ProposalEvent[] = [];
const rawEvents = dedupeContractEvents(res.events ?? []);

if (res.events && Array.isArray(res.events)) {
for (const rawEv of res.events) {
if (rawEvents.length > 0) {
for (const rawEv of rawEvents) {
try {
const rawTopic = Array.isArray(rawEv.topic)
? rawEv.topic
Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

30 changes: 27 additions & 3 deletions scripts/indexer.js
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,25 @@ function parseVal(scVal) {
}
}

function canonicalEventIdentity(rawEvent, fallbackIndex = 0) {
const ledger = Number(rawEvent?.ledger ?? 0);
const txHash = String(rawEvent?.txHash ?? rawEvent?.transactionHash ?? rawEvent?.tx_hash ?? "");
const eventIndex = Number(
rawEvent?.eventIndex ?? rawEvent?.event_index ?? rawEvent?.eventIndex ?? fallbackIndex
);
return `${ledger}:${txHash}:${eventIndex}`;
}

function sortEventsByLedgerAndIndex(events) {
return [...events].sort((a, b) => {
const ledgerDelta = Number(a.ledger ?? 0) - Number(b.ledger ?? 0);
if (ledgerDelta !== 0) return ledgerDelta;
const eventDelta = Number(a.eventIndex ?? 0) - Number(b.eventIndex ?? 0);
if (eventDelta !== 0) return eventDelta;
return String(a.txHash ?? "").localeCompare(String(b.txHash ?? ""));
});
}

function decodeEvent(rawEvent, index) {
const topics = Array.isArray(rawEvent.topic)
? rawEvent.topic.map(parseVal)
Expand All @@ -234,7 +253,8 @@ function decodeEvent(rawEvent, index) {

const topicName = String(topics[0] ?? "").toLowerCase();
const ledger = rawEvent.ledger;
const id = rawEvent.id || `${ledger}:${rawEvent.txHash || ""}:${index}`;
const eventIndex = Number(rawEvent.eventIndex ?? rawEvent.event_index ?? index ?? 0);
const id = rawEvent.id || canonicalEventIdentity(rawEvent, index);

return {
id,
Expand All @@ -243,6 +263,7 @@ function decodeEvent(rawEvent, index) {
value,
ledger,
txHash: rawEvent.txHash,
eventIndex,
ledgerClosedAt: rawEvent.ledgerClosedAt,
};
}
Expand Down Expand Up @@ -371,7 +392,7 @@ async function main() {
limit: 100,
});

const rawEvents = res.events || [];
const rawEvents = sortEventsByLedgerAndIndex(res.events || []);
const latestSeen = res.latestLedger || currentLedger;

log(`[INFO] Ingesting ledger range ${currentLedger}..${latestSeen} (found ${rawEvents.length} events)`);
Expand All @@ -387,13 +408,16 @@ async function main() {

for (let i = 0; i < rawEvents.length; i++) {
const decoded = decodeEvent(rawEvents[i], i);
if (!existingEventIds.has(decoded.id)) {
const eventKey = canonicalEventIdentity(rawEvents[i], i);
if (!existingEventIds.has(decoded.id) && !existingEventIds.has(eventKey)) {
existingEventIds.add(decoded.id);
existingEventIds.add(eventKey);
store.events.push(decoded);
}
applyEventToProposals(decoded, proposalsMap);
}

store.events = sortEventsByLedgerAndIndex(store.events);
store.proposals = [...proposalsMap.values()];
}

Expand Down