Documentation
¶
Index ¶
- func IsWellKnownType(fd protoreflect.FieldDescriptor) bool
- func ValueToString(value any, bytesEncoding bytes.Encoding) (s string)
- type AccumulatorInserter
- type DataType
- type Database
- func (d *Database) ApplyConstraints() error
- func (d *Database) BeginTransaction() (err error)
- func (d *Database) BufferStats() (sql.WriteStats, bool)
- func (d *Database) Close(ctx context.Context) error
- func (d *Database) CommitTransaction() (err error)
- func (d *Database) CreateDatabase(applyConstraints bool) error
- func (d *Database) DatabaseHash(schemaName string) (uint64, error)
- func (d *Database) DropConstraints() error
- func (d *Database) EnsureBlockNumberIndexes(ctx context.Context) error
- func (d *Database) FetchCursor() (*sink.Cursor, error)
- func (d *Database) FetchSinkInfo(schemaName string) (*sql.SinkInfo, error)
- func (d *Database) Flush() (time.Duration, error)
- func (d *Database) GetDialect() sql.Dialect
- func (d *Database) HandleBlocksUndo(lastValidBlockNum uint64) (err error)
- func (d *Database) Insert(table string, values []any) error
- func (d *Database) InsertBlock(blockNum uint64, hash string, timestamp time.Time) error
- func (d *Database) MissingConstraints() ([]string, error)
- func (d *Database) Open() error
- func (d *Database) RollbackTransaction()
- func (d *Database) StoreCursor(cursor *sink.Cursor) error
- func (d *Database) StoreSinkInfo(schemaName string, schemaHash string) error
- func (d *Database) SwitchToDirectInserts(ctx context.Context, reason string, atChainHead bool) error
- func (d *Database) UpdateSinkInfoHash(schemaName string, newHash string) error
- func (d *Database) WalkMessageDescriptorAndInsert(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, ...) (time.Duration, error)
- func (d *Database) WalkMessageDescriptorAndInsertInto(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, ...) (time.Duration, error)
- func (d *Database) WithSpool(options spool.Options)
- func (d *Database) WithWriteMode(mode sql.WriteMode) error
- func (d *Database) WriteMode() sql.WriteMode
- type DialectPostgres
- func (d *DialectPostgres) AppendInlineFieldValues(fieldValues []any, fd protoreflect.FieldDescriptor, fv protoreflect.Value, ...) ([]any, error)
- func (d *DialectPostgres) FullTableName(table *schema.Table) string
- func (d *DialectPostgres) SchemaHash() string
- func (d *DialectPostgres) UseDeletedField() bool
- func (d *DialectPostgres) UseVersionField() bool
- type RowInserter
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func IsWellKnownType ¶
func IsWellKnownType(fd protoreflect.FieldDescriptor) bool
Types ¶
type AccumulatorInserter ¶
type AccumulatorInserter struct {
// contains filtered or unexported fields
}
func NewAccumulatorInserter ¶
func NewAccumulatorInserter(logger *zap.Logger) (*AccumulatorInserter, error)
type DataType ¶
type DataType string
const ( TypeNumeric DataType = "NUMERIC" TypeInteger DataType = "INTEGER" TypeBool DataType = "BOOLEAN" TypeBigInt DataType = "BIGINT" TypeDecimal DataType = "DECIMAL" TypeDouble DataType = "DOUBLE PRECISION" TypeText DataType = "TEXT" TypeBlob DataType = "BLOB" TypeVarchar DataType = "VARCHAR(255)" TypeBytea DataType = "BYTEA" TypeTimestamp DataType = "TIMESTAMP" TypeJSONB DataType = "JSONB" )
func MapFieldType ¶
func MapFieldType(fd protoreflect.FieldDescriptor, bytesEncoding bytes.Encoding, column *schema.Column) DataType
type Database ¶
type Database struct {
*sql.BaseDatabase
// contains filtered or unexported fields
}
func NewDatabase ¶
func (*Database) ApplyConstraints ¶ added in v1.22.0
ApplyConstraints puts the dialect's constraints on a schema that already exists, leaving the ones already there alone.
It is what `sink postgres apply-constraints` does to a database synced without them: the sink info hash cannot tell those two apart, since it is computed over the DDL the dialect would emit, constraints included, either way.
The caller owns the transaction. On a populated database this is not quick — every index has to be built and every foreign key validated, with the table locked meanwhile.
func (*Database) BeginTransaction ¶
func (*Database) BufferStats ¶ added in v1.22.0
func (d *Database) BufferStats() (sql.WriteStats, bool)
BufferStats reports what the local buffer is holding and what it has committed.
func (*Database) Close ¶ added in v1.22.0
Close drains the local buffer, if one is in use, so the blocks buffered at shutdown reach the database rather than being streamed again on the next run. Close drains a local buffer, if one is in use, and then releases the connections. It is called once the stream is done with the database, so holding the pools open past it only occupies connection slots on the server.
func (*Database) CommitTransaction ¶
func (*Database) CreateDatabase ¶
func (*Database) DropConstraints ¶ added in v1.22.0
DropConstraints removes the constraints this schema's DDL would create, leaving anything the sink did not put there alone.
Foreign keys go first, then unique constraints, then primary keys: a primary key still referenced by a foreign key cannot be dropped. Constraints already absent are skipped, so this is idempotent and safe to run against a schema that never had them.
func (*Database) EnsureBlockNumberIndexes ¶ added in v1.22.0
EnsureBlockNumberIndexes creates the index the sink needs for its own reorg path, on every table, when the sink starts.
It is not part of the constraint pass and not governed by --apply-constraints. Those describe the schema and are the operator's to schedule; this one the sink depends on to undo a reorg without sequentially scanning every table, so waiting for a maintenance window would mean running without it for as long as the operator likes.
Concurrently, and therefore outside any transaction: a restart onto an already-loaded table must not lock out the writers. On a schema that already has them this is one catalog query and nothing else.
func (*Database) FetchSinkInfo ¶
func (*Database) GetDialect ¶
func (*Database) HandleBlocksUndo ¶
HandleBlocksUndo removes everything a reorg invalidated: every entity row above the last valid block, then the block rows themselves.
The deletes are explicit rather than left to `fk_block ... ON DELETE CASCADE`, because that foreign key only exists once the constraints have been created. Deleting just the block rows, as this used to, silently orphaned every entity row of the undone blocks on a schema without constraints — and that is now the default. Every table carries `_block_number_`, so the same delete works either way, and the descending Ordinal visits children before parents so a foreign key never blocks its own cleanup.
func (*Database) InsertBlock ¶
func (*Database) MissingConstraints ¶ added in v1.22.0
MissingConstraints names the constraints the policy says this schema should carry and the catalog does not have.
It is one query against pg_constraint filtered by namespace — indexed, and nothing like the cost of building the constraints it reports on — so it is cheap enough to run on every start.
func (*Database) RollbackTransaction ¶
func (d *Database) RollbackTransaction()
func (*Database) StoreSinkInfo ¶
func (*Database) SwitchToDirectInserts ¶ added in v1.22.0
func (d *Database) SwitchToDirectInserts(ctx context.Context, reason string, atChainHead bool) error
SwitchToDirectInserts drains the spool and inserts straight into the database from here on. The reason says what brought the switch on, since the database cannot tell the chain head from the end of a bounded range.
The spool trades freshness for throughput: rows sit on disk until a segment fills, which is what a backfill wants and the opposite of what a sink at the chain head wants, where a block should be queryable when it arrives rather than when the segment it happens to land in is full. Reorgs are also cheaper to undo out of a table than out of a segment that has not been applied yet.
Everything spooled is applied before the switch, so no row is left behind, and the spool is closed for good — a stream that has reached the head does not go back.
func (*Database) UpdateSinkInfoHash ¶
func (*Database) WalkMessageDescriptorAndInsert ¶
func (*Database) WalkMessageDescriptorAndInsertInto ¶ added in v1.22.0
func (*Database) WithSpool ¶ added in v1.22.0
WithSpool turns on the on-disk spool. It must be called before Open.
func (*Database) WithWriteMode ¶ added in v1.22.0
WithWriteMode records how sealed segments should reach the database. It must be called before Open, which is where the mode is resolved and validated against the schema.
type DialectPostgres ¶
type DialectPostgres struct {
*sql2.BaseDialect
// contains filtered or unexported fields
}
func NewDialectPostgres ¶
func (*DialectPostgres) AppendInlineFieldValues ¶
func (d *DialectPostgres) AppendInlineFieldValues(fieldValues []any, fd protoreflect.FieldDescriptor, fv protoreflect.Value, dm protoreflect.Message) ([]any, error)
func (*DialectPostgres) FullTableName ¶
func (d *DialectPostgres) FullTableName(table *schema.Table) string
func (*DialectPostgres) SchemaHash ¶
func (d *DialectPostgres) SchemaHash() string
func (*DialectPostgres) UseDeletedField ¶
func (d *DialectPostgres) UseDeletedField() bool
func (*DialectPostgres) UseVersionField ¶
func (d *DialectPostgres) UseVersionField() bool
type RowInserter ¶
type RowInserter struct {
// contains filtered or unexported fields
}
func NewRowInserter ¶
func NewRowInserter(logger *zap.Logger) (*RowInserter, error)