db

package
v0.5.0 Latest Latest
Warning

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

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

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrBatchAlreadyClosed = errors.New("batch already closed")
)

Functions

This section is empty.

Types

type ApplyCommitParams

type ApplyCommitParams struct {
	Subject      string
	DbID         string
	PageCount    int64
	LastCommitID []byte
	LeaseExpires pgtype.Timestamptz
}

type Block

type Block struct {
	Subject string
	DbID    string
	Idx     int64
	Part    int16
	Data    []byte
}

type Change

type Change struct {
	Subject string
	DbID    string
	Version int64
	Blocks  []int64
}

type ChangedSinceParams

type ChangedSinceParams struct {
	Subject string
	DbID    string
	Version int64
}

type ChangedSinceRow

type ChangedSinceRow struct {
	Version int64
	Blocks  []int64
}

type CreateDatabaseParams

type CreateDatabaseParams struct {
	Subject  string
	DbID     string
	PageSize int32
}

type DBTX

type DBTX interface {
	Exec(context.Context, string, ...interface{}) (pgconn.CommandTag, error)
	Query(context.Context, string, ...interface{}) (pgx.Rows, error)
	QueryRow(context.Context, string, ...interface{}) pgx.Row
	SendBatch(context.Context, *pgx.Batch) pgx.BatchResults
}

type Database

type Database struct {
	Subject      string
	DbID         string
	PageSize     int32
	PageCount    int64
	Version      int64
	LastCommitID []byte
	LeaseEpoch   int64
	LeaseID      []byte
	LeaseHolder  []byte
	LeaseGranted pgtype.Timestamptz
	LeaseExpires pgtype.Timestamptz
	Deleted      bool
}

type DatabaseExistsParams

type DatabaseExistsParams struct {
	Subject string
	DbID    string
}

type DeleteBlocksFromParams

type DeleteBlocksFromParams struct {
	Subject string
	DbID    string
	Idx     int64
}

type DeleteChangesParams

type DeleteChangesParams struct {
	Subject string
	DbID    string
}

type DeleteDatabaseParams

type DeleteDatabaseParams struct {
	Subject string
	DbID    string
}

type ExtendLeaseParams

type ExtendLeaseParams struct {
	Subject      string
	DbID         string
	LeaseEpoch   int64
	LeaseExpires pgtype.Timestamptz
}

type FetchBlocksParams

type FetchBlocksParams struct {
	Subject string
	DbID    string
	First   int64
	Beyond  int64
}

type FetchBlocksRow

type FetchBlocksRow struct {
	Idx  int64
	Part int16
	Data []byte
}

type GetDatabaseParams

type GetDatabaseParams struct {
	Subject string
	DbID    string
}

type GrantLeaseParams

type GrantLeaseParams struct {
	Subject      string
	DbID         string
	LeaseID      []byte
	LeaseHolder  []byte
	LeaseGranted pgtype.Timestamptz
	LeaseExpires pgtype.Timestamptz
}

type InsertSlotParams added in v0.3.0

type InsertSlotParams struct {
	Owner     string
	Subject   string
	Label     string
	ClaimedAt pgtype.Timestamptz
}

type LockDatabaseParams

type LockDatabaseParams struct {
	Subject string
	DbID    string
}

type OldestChangeParams

type OldestChangeParams struct {
	Subject string
	DbID    string
}

type PruneChangesParams

type PruneChangesParams struct {
	Subject string
	DbID    string
	Version int64
}

type PurgeDatabaseParams added in v0.3.0

type PurgeDatabaseParams struct {
	Subject string
	DbID    string
}

type PurgeSubjectRow added in v0.3.0

type PurgeSubjectRow struct {
	DbID    string
	Deleted bool
}

type PutBlocksBatchResults

type PutBlocksBatchResults struct {
	// contains filtered or unexported fields
}

func (*PutBlocksBatchResults) Close

func (b *PutBlocksBatchResults) Close() error

func (*PutBlocksBatchResults) Exec

func (b *PutBlocksBatchResults) Exec(f func(int, error))

type PutBlocksParams

type PutBlocksParams struct {
	Subject string
	DbID    string
	Idx     int64
	Part    int16
	Data    []byte
}

type Queries

type Queries struct {
	// contains filtered or unexported fields
}

func New

func New(db DBTX) *Queries

func (*Queries) ApplyCommit

func (q *Queries) ApplyCommit(ctx context.Context, arg ApplyCommitParams) (int64, error)

func (*Queries) ChangedSince

func (q *Queries) ChangedSince(ctx context.Context, arg ChangedSinceParams) ([]ChangedSinceRow, error)

func (*Queries) CreateDatabase

func (q *Queries) CreateDatabase(ctx context.Context, arg CreateDatabaseParams) (Database, error)

CreateDatabase inserts nothing if the row already exists. The caller then locks and reads the existing row.

func (*Queries) DatabaseExists

func (q *Queries) DatabaseExists(ctx context.Context, arg DatabaseExistsParams) (bool, error)

func (*Queries) DeleteBlocksFrom

func (q *Queries) DeleteBlocksFrom(ctx context.Context, arg DeleteBlocksFromParams) error

func (*Queries) DeleteChanges

func (q *Queries) DeleteChanges(ctx context.Context, arg DeleteChangesParams) error

func (*Queries) DeleteDatabase

func (q *Queries) DeleteDatabase(ctx context.Context, arg DeleteDatabaseParams) (int64, error)

DeleteDatabase marks a database as deleted and empties its state. It increases the lease epoch and the version and removes the lease, so that no earlier lease can commit again.

func (*Queries) DeleteSlotOf added in v0.3.0

func (q *Queries) DeleteSlotOf(ctx context.Context, owner string) error

func (*Queries) ExtendLease

func (q *Queries) ExtendLease(ctx context.Context, arg ExtendLeaseParams) (int64, error)

ExtendLease updates no row if the epoch is not current or the lease was released. The caller then returns ErrFenced, or ErrNotFound if the database does not exist.

func (*Queries) FetchBlocks

func (q *Queries) FetchBlocks(ctx context.Context, arg FetchBlocksParams) ([]FetchBlocksRow, error)

FetchBlocks returns the stored parts of the blocks in a range, ordered by block and part. A block that was never written has no rows. The caller returns zeros for it.

func (*Queries) GetDatabase

func (q *Queries) GetDatabase(ctx context.Context, arg GetDatabaseParams) (Database, error)

func (*Queries) GetSlot added in v0.3.0

func (q *Queries) GetSlot(ctx context.Context, owner string) (Slot, error)

func (*Queries) GrantLease

func (q *Queries) GrantLease(ctx context.Context, arg GrantLeaseParams) (int64, error)

func (*Queries) HoldsSlot added in v0.3.0

func (q *Queries) HoldsSlot(ctx context.Context, subject string) (bool, error)

func (*Queries) InsertSlot added in v0.3.0

func (q *Queries) InsertSlot(ctx context.Context, arg InsertSlotParams) (Slot, error)

InsertSlot inserts nothing if another transaction inserted the owner's slot first. The caller then locks it.

func (*Queries) LockDatabase

func (q *Queries) LockDatabase(ctx context.Context, arg LockDatabaseParams) (Database, error)

All SQL queries of the PostgreSQL store. sqlc generates internal/store/postgres/db from this file. LockDatabase locks the row of a database. Granting a lease, fencing and applying a commit all hold this lock, so they are serialized per database.

func (*Queries) LockSlot added in v0.3.0

func (q *Queries) LockSlot(ctx context.Context, owner string) (Slot, error)

LockSlot locks an owner's slot. Opening a database and claiming or deleting the slot take this lock before any database row, so they cannot deadlock.

func (*Queries) OldestChange

func (q *Queries) OldestChange(ctx context.Context, arg OldestChangeParams) (int64, error)

OldestChange returns the oldest version in the change log, or 0 if the log is empty. The log is empty for a database created anew after a deletion, and for one whose commits all predate the change log.

func (*Queries) PruneChanges

func (q *Queries) PruneChanges(ctx context.Context, arg PruneChangesParams) error

PruneChanges deletes the change log entries up to and including the given version.

func (*Queries) PurgeDatabase added in v0.3.0

func (q *Queries) PurgeDatabase(ctx context.Context, arg PurgeDatabaseParams) error

func (*Queries) PurgeSubject added in v0.3.0

func (q *Queries) PurgeSubject(ctx context.Context, subject string) ([]PurgeSubjectRow, error)

PurgeSubject deletes all databases of a subject completely, with their blocks and change logs (foreign keys). It returns those that were not deleted before.

func (*Queries) PutBlocks

func (q *Queries) PutBlocks(ctx context.Context, arg []PutBlocksParams) *PutBlocksBatchResults

PutBlocks inserts or updates the parts of the blocks of one commit. pgx sends all rows in one pipelined batch. An update of a part stays on the same page (HOT update), so the table does not grow. See migrations/00003_block_parts.sql.

func (*Queries) RecordChange

func (q *Queries) RecordChange(ctx context.Context, arg RecordChangeParams) error

RecordChange adds the changed blocks of a commit to the change log. On a conflict it does nothing, so a repeated commit cannot fail here: the entry for that version already exists with the same content.

func (*Queries) ReleaseLease

func (q *Queries) ReleaseLease(ctx context.Context, arg ReleaseLeaseParams) (int64, error)

func (*Queries) ReleaseUnusedSlots added in v0.3.0

func (q *Queries) ReleaseUnusedSlots(ctx context.Context, claimedAt pgtype.Timestamptz) error

ReleaseUnusedSlots releases the slots that were claimed before the cutoff and whose key has no databases left.

func (*Queries) ReviveDatabase

func (q *Queries) ReviveDatabase(ctx context.Context, arg ReviveDatabaseParams) (Database, error)

ReviveDatabase creates a deleted database anew, empty, with the given page size. Epoch and version continue.

func (*Queries) SetSlotLabel added in v0.4.0

func (q *Queries) SetSlotLabel(ctx context.Context, arg SetSlotLabelParams) (Slot, error)

SetSlotLabel changes the label only while subject holds the owner's slot.

func (*Queries) ShareSlot added in v0.3.0

func (q *Queries) ShareSlot(ctx context.Context, owner string) (Slot, error)

ShareSlot locks an owner's slot against a concurrent claim while a database is opened.

func (*Queries) SyncStandbyNames

func (q *Queries) SyncStandbyNames(ctx context.Context) (string, error)

SyncStandbyNames returns synchronous_standby_names. It is empty if PostgreSQL does not replicate commits synchronously.

func (*Queries) UnusedDatabases

func (q *Queries) UnusedDatabases(ctx context.Context, arg UnusedDatabasesParams) ([]UnusedDatabasesRow, error)

UnusedDatabases returns up to $2 databases whose lease expired before $1, oldest first. It takes no lock: the caller locks each row and checks it again before deleting. No index covers lease_expires, because every commit and lease renewal updates it, and an index on it would prevent HOT updates of the row. A sweep reads the whole table.

func (*Queries) UpdateSlot added in v0.3.0

func (q *Queries) UpdateSlot(ctx context.Context, arg UpdateSlotParams) (Slot, error)

func (*Queries) WithTx

func (q *Queries) WithTx(tx pgx.Tx) *Queries

type RecordChangeParams

type RecordChangeParams struct {
	Subject string
	DbID    string
	Version int64
	Blocks  []int64
}

type ReleaseLeaseParams

type ReleaseLeaseParams struct {
	Subject    string
	DbID       string
	LeaseEpoch int64
}

type ReviveDatabaseParams

type ReviveDatabaseParams struct {
	Subject  string
	DbID     string
	PageSize int32
}

type SetSlotLabelParams added in v0.4.0

type SetSlotLabelParams struct {
	Owner   string
	Subject string
	Label   string
}

type Slot added in v0.3.0

type Slot struct {
	Owner     string
	Subject   string
	Label     string
	ClaimedAt pgtype.Timestamptz
}

type UnusedDatabasesParams

type UnusedDatabasesParams struct {
	LeaseExpires pgtype.Timestamptz
	Limit        int32
}

type UnusedDatabasesRow

type UnusedDatabasesRow struct {
	Subject string
	DbID    string
}

type UpdateSlotParams added in v0.3.0

type UpdateSlotParams struct {
	Owner     string
	Subject   string
	Label     string
	ClaimedAt pgtype.Timestamptz
}

Jump to

Keyboard shortcuts

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