postgres

package
v1.25.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 34 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func IsWellKnownType

func IsWellKnownType(fd protoreflect.FieldDescriptor) bool

func ValueToString

func ValueToString(value any, bytesEncoding bytes.Encoding) (s string)

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

func (DataType) String

func (s DataType) String() string

type Database

type Database struct {
	*sql.BaseDatabase
	// contains filtered or unexported fields
}

func NewDatabase

func NewDatabase(schema *schema.Schema, dsn *db.DSN, moduleOutputType string, rootMessageDescriptor protoreflect.MessageDescriptor, useProtoOptions bool, constraints sql.ConstraintPolicy, bytesEncoding bytes.Encoding, logger *zap.Logger) (*Database, error)

func (*Database) ApplyConstraints added in v1.22.0

func (d *Database) ApplyConstraints() error

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 (d *Database) BeginTransaction() (err error)

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

func (d *Database) Close(ctx context.Context) error

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 (d *Database) CommitTransaction() (err error)

func (*Database) CreateDatabase

func (d *Database) CreateDatabase(applyConstraints bool) error

func (*Database) DatabaseHash

func (d *Database) DatabaseHash(schemaName string) (uint64, error)

func (*Database) DropConstraints added in v1.22.0

func (d *Database) DropConstraints() error

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

func (d *Database) EnsureBlockNumberIndexes(ctx context.Context) error

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) FetchCursor

func (d *Database) FetchCursor() (*sink.Cursor, error)

func (*Database) FetchSinkInfo

func (d *Database) FetchSinkInfo(schemaName string) (*sql.SinkInfo, error)

func (*Database) Flush

func (d *Database) Flush() (time.Duration, error)

func (*Database) GetDialect

func (d *Database) GetDialect() sql.Dialect

func (*Database) HandleBlocksUndo

func (d *Database) HandleBlocksUndo(lastValidBlockNum uint64) (err error)

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) Insert

func (d *Database) Insert(table string, values []any) error

func (*Database) InsertBlock

func (d *Database) InsertBlock(blockNum uint64, hash string, timestamp time.Time) error

func (*Database) MissingConstraints added in v1.22.0

func (d *Database) MissingConstraints() ([]string, error)

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) Open

func (d *Database) Open() error

func (*Database) RollbackTransaction

func (d *Database) RollbackTransaction()

func (*Database) StoreCursor

func (d *Database) StoreCursor(cursor *sink.Cursor) error

func (*Database) StoreSinkInfo

func (d *Database) StoreSinkInfo(schemaName string, schemaHash string) error

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 (d *Database) UpdateSinkInfoHash(schemaName string, newHash string) error

func (*Database) WalkMessageDescriptorAndInsert

func (d *Database) WalkMessageDescriptorAndInsert(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, parent *sql.Parent) (time.Duration, error)

func (*Database) WalkMessageDescriptorAndInsertInto added in v1.22.0

func (d *Database) WalkMessageDescriptorAndInsertInto(dm protoreflect.Message, blockNum uint64, blockTimestamp time.Time, parent *sql.Parent, inserter sql.Inserter) (time.Duration, error)

func (*Database) WithSpool added in v1.22.0

func (d *Database) WithSpool(options spool.Options)

WithSpool turns on the on-disk spool. It must be called before Open.

func (*Database) WithWriteMode added in v1.22.0

func (d *Database) WithWriteMode(mode sql.WriteMode) error

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.

func (*Database) WriteMode added in v1.22.0

func (d *Database) WriteMode() sql.WriteMode

WriteMode reports the mode Open settled on, for the startup log.

type DialectPostgres

type DialectPostgres struct {
	*sql2.BaseDialect
	// contains filtered or unexported fields
}

func NewDialectPostgres

func NewDialectPostgres(schema *schema.Schema, bytesEncoding bytes.Encoding, logger *zap.Logger) (*DialectPostgres, error)

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)

Directories

Path Synopsis
Package pgcopy writes the PostgreSQL binary COPY format ("PGCOPY").
Package pgcopy writes the PostgreSQL binary COPY format ("PGCOPY").

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL