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
19 changes: 0 additions & 19 deletions packages/envio/src/Config.gen.ts

This file was deleted.

52 changes: 24 additions & 28 deletions packages/envio/src/Config.res
Original file line number Diff line number Diff line change
Expand Up @@ -86,38 +86,38 @@ type t = {
allEnums: array<Table.enumConfig<Table.enum>>,
}

module DynamicContractRegistry = {
let name = "dynamic_contract_registry"
module EnvioAddresses = {
let name = "envio_addresses"
let index = -1

let makeId = (~chainId, ~contractAddress) => {
chainId->Belt.Int.toString ++ "-" ++ contractAddress->Address.toString
let makeId = (~chainId, ~address) => {
chainId->Belt.Int.toString ++ "-" ++ address->Address.toString
}

@genType
type t = {
id: string,
@as("chain_id") chainId: int,
@as("registering_event_block_number") registeringEventBlockNumber: int,
@as("registering_event_log_index") registeringEventLogIndex: int,
@as("registering_event_block_timestamp") registeringEventBlockTimestamp: int,
@as("registering_event_contract_name") registeringEventContractName: string,
@as("registering_event_name") registeringEventName: string,
@as("registering_event_src_address") registeringEventSrcAddress: Address.t,
@as("contract_address") contractAddress: Address.t,
@as("registration_block") registrationBlock: int,
// -1 when the address was registered from a block handler (no log index)
@as("registration_log_index") registrationLogIndex: int,
@as("contract_name") contractName: string,
}

// Extract the raw contract address from the composite id ({chainId}-{address}).
// Inverse of makeId. Keep in sync with makeId above and the SUBSTRING SQL in
// InternalTable.Chains.makeGetInitialStateQuery.
let getAddress = (entity: t): Address.t => {
let sepIdx = entity.id->String.indexOf("-")
entity.id
->String.slice(~start=sepIdx + 1, ~end=entity.id->String.length)
->Address.unsafeFromString
}

let schema = S.schema(s => {
id: s.matches(S.string),
chainId: s.matches(S.int),
registeringEventBlockNumber: s.matches(S.int),
registeringEventLogIndex: s.matches(S.int),
registeringEventContractName: s.matches(S.string),
registeringEventName: s.matches(S.string),
registeringEventSrcAddress: s.matches(Address.schema),
registeringEventBlockTimestamp: s.matches(S.int),
contractAddress: s.matches(Address.schema),
registrationBlock: s.matches(S.int),
registrationLogIndex: s.matches(S.int),
contractName: s.matches(S.string),
})

Expand All @@ -128,13 +128,9 @@ module DynamicContractRegistry = {
~fields=[
Table.mkField("id", String, ~isPrimaryKey=true, ~fieldSchema=S.string),
Table.mkField("chain_id", Int32, ~fieldSchema=S.int),
Table.mkField("registering_event_block_number", Int32, ~fieldSchema=S.int),
Table.mkField("registering_event_log_index", Int32, ~fieldSchema=S.int),
Table.mkField("registering_event_block_timestamp", Int32, ~fieldSchema=S.int),
Table.mkField("registering_event_contract_name", String, ~fieldSchema=S.string),
Table.mkField("registering_event_name", String, ~fieldSchema=S.string),
Table.mkField("registering_event_src_address", String, ~fieldSchema=Address.schema),
Table.mkField("contract_address", String, ~fieldSchema=Address.schema),
Table.mkField("registration_block", Int32, ~fieldSchema=S.int),
// -1 sentinel when registered from a block handler (no log index)
Table.mkField("registration_log_index", Int32, ~fieldSchema=S.int),
Table.mkField("contract_name", String, ~fieldSchema=S.string),
],
)
Expand Down Expand Up @@ -289,7 +285,7 @@ let getFieldTypeAndSchema = (prop, ~enumConfigsByName: dict<Table.enumConfig<Tab
| "bigint" => (Table.BigInt({precision: ?prop["precision"]}), Utils.BigInt.schema->S.toUnknown)
| "bigdecimal" => (
Table.BigDecimal({
config: ?(prop["precision"]->Option.map(p => (p, prop["scale"]->Option.getOr(0)))),
config: ?prop["precision"]->Option.map(p => (p, prop["scale"]->Option.getOr(0))),
}),
BigDecimal.schema->S.toUnknown,
)
Expand Down Expand Up @@ -736,7 +732,7 @@ let fromPublic = (publicConfigJson: JSON.t, ~maxAddrInPartition=5000) => {
->Option.getOr([])
->parseEntitiesFromJson(~enumConfigsByName)

let allEntities = userEntities->Array.concat([DynamicContractRegistry.entityConfig])
let allEntities = userEntities->Array.concat([EnvioAddresses.entityConfig])

let userEntitiesByName =
userEntities
Expand Down
2 changes: 1 addition & 1 deletion packages/envio/src/Envio.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import type {Effect as $$effect} from './Types.ts';

import type {Logger as $$logger} from './Types.ts';

import type {S_t as RescriptSchema_S_t} from 'rescript-schema/RescriptSchema.gen.js';
import type {S_t as RescriptSchema_S_t} from './RescriptSchema.gen.js';

export type blockEvent = { readonly number: number };

Expand Down
17 changes: 6 additions & 11 deletions packages/envio/src/InMemoryStore.res
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ let isRollingBack = (inMemoryStore: t) => inMemoryStore.rollbackTargetCheckpoint

let setBatchDcs = (inMemoryStore: t, ~batch: Batch.t, ~shouldSaveHistory) => {
let inMemTable =
inMemoryStore->getInMemTable(~entityConfig=InternalTable.DynamicContractRegistry.entityConfig)
inMemoryStore->getInMemTable(~entityConfig=InternalTable.EnvioAddresses.entityConfig)

let itemIdx = ref(0)

Expand All @@ -118,24 +118,19 @@ let setBatchDcs = (inMemoryStore: t, ~batch: Batch.t, ~shouldSaveHistory) => {
let eventItem = item->Internal.castUnsafeEventItem
for dcIdx in 0 to dcs->Array.length - 1 {
let dc = dcs->Array.getUnsafe(dcIdx)
let entity: InternalTable.DynamicContractRegistry.t = {
id: InternalTable.DynamicContractRegistry.makeId(~chainId, ~contractAddress=dc.address),
let entity: InternalTable.EnvioAddresses.t = {
id: InternalTable.EnvioAddresses.makeId(~chainId, ~address=dc.address),
chainId,
contractAddress: dc.address,
contractName: dc.contractName,
registeringEventBlockNumber: eventItem.blockNumber,
registeringEventLogIndex: eventItem.logIndex,
registeringEventBlockTimestamp: eventItem.timestamp,
registeringEventContractName: eventItem.eventConfig.contractName,
registeringEventName: eventItem.eventConfig.name,
registeringEventSrcAddress: eventItem.event.srcAddress,
registrationBlock: eventItem.blockNumber,
registrationLogIndex: eventItem.logIndex,
}

inMemTable->InMemoryTable.Entity.set(
Set({
entityId: entity.id,
checkpointId,
entity: entity->InternalTable.DynamicContractRegistry.castToInternal,
entity: entity->InternalTable.EnvioAddresses.castToInternal,
}),
~shouldSaveHistory,
)
Expand Down
71 changes: 36 additions & 35 deletions packages/envio/src/Persistence.res
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ let make = (
~allEnums,
~storage,
) => {
let allEntities = userEntities->Array.concat([InternalTable.DynamicContractRegistry.entityConfig])
let allEntities = userEntities->Array.concat([InternalTable.EnvioAddresses.entityConfig])
let allEnums =
allEnums->Array.concat([EntityHistory.RowAction.config->Table.fromGenericEnumConfig])
{
Expand Down Expand Up @@ -316,45 +316,46 @@ let prepareRollbackDiff = async (
let deletedEntities = Dict.make()
let setEntities = Dict.make()

let _ = await persistence.allEntities
->Belt.Array.map(async entityConfig => {
let entityTable = inMemStore->InMemoryStore.getInMemTable(~entityConfig)
let _ =
await persistence.allEntities
->Belt.Array.map(async entityConfig => {
let entityTable = inMemStore->InMemoryStore.getInMemTable(~entityConfig)

let (removedIdsResult, restoredEntitiesResult) = await persistence.storage.getRollbackData(
~entityConfig,
~rollbackTargetCheckpointId,
)

// Process removed IDs
removedIdsResult->Array.forEach(data => {
deletedEntities->Utils.Dict.push(entityConfig.name, data["id"])
entityTable->InMemoryTable.Entity.set(
Delete({
entityId: data["id"],
checkpointId: rollbackDiffCheckpointId,
}),
~shouldSaveHistory=false,
~containsRollbackDiffChange=true,
let (removedIdsResult, restoredEntitiesResult) = await persistence.storage.getRollbackData(
~entityConfig,
~rollbackTargetCheckpointId,
)
})

let restoredEntities = restoredEntitiesResult->S.parseOrThrow(entityConfig.rowsSchema)
// Process removed IDs
removedIdsResult->Array.forEach(data => {
deletedEntities->Utils.Dict.push(entityConfig.name, data["id"])
entityTable->InMemoryTable.Entity.set(
Delete({
entityId: data["id"],
checkpointId: rollbackDiffCheckpointId,
}),
~shouldSaveHistory=false,
~containsRollbackDiffChange=true,
)
})

// Process restored entities
restoredEntities->Belt.Array.forEach((entity: Internal.entity) => {
setEntities->Utils.Dict.push(entityConfig.name, entity.id)
entityTable->InMemoryTable.Entity.set(
Set({
entityId: entity.id,
checkpointId: rollbackDiffCheckpointId,
entity,
}),
~shouldSaveHistory=false,
~containsRollbackDiffChange=true,
)
let restoredEntities = restoredEntitiesResult->S.parseOrThrow(entityConfig.rowsSchema)

// Process restored entities
restoredEntities->Belt.Array.forEach((entity: Internal.entity) => {
setEntities->Utils.Dict.push(entityConfig.name, entity.id)
entityTable->InMemoryTable.Entity.set(
Set({
entityId: entity.id,
checkpointId: rollbackDiffCheckpointId,
entity,
}),
~shouldSaveHistory=false,
~containsRollbackDiffChange=true,
)
})
})
})
->Promise.all
->Promise.all

{
"inMemStore": inMemStore,
Expand Down
3 changes: 2 additions & 1 deletion packages/envio/src/PgStorage.res
Original file line number Diff line number Diff line change
Expand Up @@ -503,7 +503,8 @@ let setOrThrow = async (sql, ~items, ~table: Table.table, ~itemSchema, ~pgSchema
let isFullChunk = chunkSize === maxItemsPerQuery

let params = data["convertOrThrow"](chunk->(Utils.magic: array<'item> => array<unknown>))

// Use prepared query only for full batches where the cached query is reused.
// Partial chunks generate unique SQL each time, so preparation has no benefit.
let response = isFullChunk
? sql->Postgres.preparedUnsafe(data["query"], params)
: sql->Postgres.unpreparedUnsafe(
Expand Down
60 changes: 32 additions & 28 deletions packages/envio/src/TestIndexer.res
Original file line number Diff line number Diff line change
Expand Up @@ -34,18 +34,15 @@ type testIndexerState = {
mutable processChanges: array<unknown>,
}

// Cast Internal.entity back to DynamicContractRegistry.t
external castFromDcRegistry: Internal.entity => InternalTable.DynamicContractRegistry.t =
"%identity"

// Convert DynamicContractRegistry.t to Internal.indexingContract
let toIndexingContract = (
dc: InternalTable.DynamicContractRegistry.t,
): Internal.indexingContract => {
address: dc.contractAddress,
// Cast Internal.entity back to EnvioAddresses.t
external castToEnvioAddresses: Internal.entity => InternalTable.EnvioAddresses.t = "%identity"

// Convert EnvioAddresses.t to Internal.indexingContract
let toIndexingContract = (dc: InternalTable.EnvioAddresses.t): Internal.indexingContract => {
address: dc->Config.EnvioAddresses.getAddress,
contractName: dc.contractName,
startBlock: dc.registeringEventBlockNumber,
registrationBlock: Some(dc.registeringEventBlockNumber),
startBlock: dc.registrationBlock,
registrationBlock: Some(dc.registrationBlock),
}

let handleLoadByIds = (
Expand Down Expand Up @@ -224,34 +221,41 @@ let handleWriteBatch = (
entityChanges
->Dict.toArray
->Array.forEach(((entityName, {sets, deleted})) => {
// Transform dynamic_contract_registry to addresses with simplified structure
if entityName === InternalTable.DynamicContractRegistry.name {
// Transform envio_addresses to addresses with simplified structure
if entityName === InternalTable.EnvioAddresses.name {
let entityObj: dict<unknown> = Dict.make()
if sets->Array.length > 0 {
// Transform sets to simplified {address, contract} objects
let simplifiedSets = sets->Array.map(entity => {
let dc = entity->Utils.magic->castFromDcRegistry
{"address": dc.contractAddress, "contract": dc.contractName}
let dc = entity->Utils.magic->castToEnvioAddresses
{"address": dc->Config.EnvioAddresses.getAddress, "contract": dc.contractName}
})
entityObj->Dict.set("sets", simplifiedSets->Utils.magic)
entityObj->Dict.set(
"sets",
simplifiedSets->(
Utils.magic: array<{"address": Address.t, "contract": string}> => unknown
),
)
}
// Note: deleted is not relevant for addresses since we use address string directly
change->Dict.set("addresses", entityObj->Utils.magic)
change->Dict.set("addresses", entityObj->(Utils.magic: dict<unknown> => unknown))
} else {
let entityObj: dict<unknown> = Dict.make()
if sets->Array.length > 0 {
entityObj->Dict.set("sets", sets->Utils.magic)
entityObj->Dict.set("sets", sets->(Utils.magic: array<unknown> => unknown))
}
if deleted->Array.length > 0 {
entityObj->Dict.set("deleted", deleted->Utils.magic)
entityObj->Dict.set("deleted", deleted->(Utils.magic: array<string> => unknown))
}
change->Dict.set(entityName, entityObj->Utils.magic)
change->Dict.set(entityName, entityObj->(Utils.magic: dict<unknown> => unknown))
}
})
| None => ()
}

state.processChanges->Array.push(change->Utils.magic)->ignore
state.processChanges
->Array.push(change->(Utils.magic: dict<unknown> => unknown))
->ignore
}
}

Expand Down Expand Up @@ -523,8 +527,8 @@ let makeCreateTestIndexer = (~config: Config.t, ~workerPath: string): (
// Build entity operations for each user entity
let entityOpsDict: dict<entityOperations> = Dict.make()
allEntities->Array.forEach(entityConfig => {
// Only create ops for user entities (not internal tables like dynamic_contract_registry)
if entityConfig.name !== InternalTable.DynamicContractRegistry.name {
// Only create ops for user entities (not internal tables like envio_addresses)
if entityConfig.name !== InternalTable.EnvioAddresses.name {
entityOpsDict->Dict.set(
entityConfig.name,
{
Expand Down Expand Up @@ -580,15 +584,15 @@ let makeCreateTestIndexer = (~config: Config.t, ~workerPath: string): (
// Start with static config addresses
let addresses = contract.addresses->Array.copy
// Add accumulated dynamic contract addresses
switch state.entities->Dict.get(InternalTable.DynamicContractRegistry.name) {
switch state.entities->Dict.get(InternalTable.EnvioAddresses.name) {
| Some(dcDict) =>
dcDict
->Dict.valuesToArray
->Array.forEach(
entity => {
let dc = entity->castFromDcRegistry
let dc = entity->castToEnvioAddresses
if dc.contractName === contract.name && dc.chainId === chainConfig.id {
addresses->Array.push(dc.contractAddress)->ignore
addresses->Array.push(dc->Config.EnvioAddresses.getAddress)->ignore
}
},
)
Expand Down Expand Up @@ -699,12 +703,12 @@ let makeCreateTestIndexer = (~config: Config.t, ~workerPath: string): (

// Extract dynamic contracts from state.entities for each chain
let dynamicContractsByChain: dict<array<Internal.indexingContract>> = Dict.make()
switch state.entities->Dict.get(InternalTable.DynamicContractRegistry.name) {
switch state.entities->Dict.get(InternalTable.EnvioAddresses.name) {
| Some(dcDict) =>
dcDict
->Dict.valuesToArray
->Array.forEach(entity => {
let dc = entity->castFromDcRegistry
let dc = entity->castToEnvioAddresses
let dcChainIdStr = dc.chainId->Int.toString
let contracts = switch dynamicContractsByChain->Dict.get(dcChainIdStr) {
| Some(arr) => arr
Expand Down
2 changes: 1 addition & 1 deletion packages/envio/src/Utils.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,6 @@

import * as UtilsJS from './Utils.res.mjs';

import type {S_t as RescriptSchema_S_t} from 'rescript-schema/RescriptSchema.gen.js';
import type {S_t as RescriptSchema_S_t} from './RescriptSchema.gen.js';

export const bigIntSchema: RescriptSchema_S_t<bigint> = UtilsJS.bigIntSchema as any;
Loading