Documentation
¶
Index ¶
- func NewStoreNode(database DBTX) pgmesh.Node[*ReadQueries, *StoreQueries]
- type Analysis
- type AnalysisState
- type CopyUsersParams
- type CreateUserParams
- type DBTX
- type GetAnalysisParams
- type GetUserParams
- type NullAnalysisState
- type Querier
- type Queries
- func (q *Queries) CopyUsers(ctx context.Context, arg []*CopyUsersParams) (int64, error)
- func (q *Queries) CreateUser(ctx context.Context, arg *CreateUserParams) (*User, error)
- func (q *Queries) GetAnalysis(ctx context.Context, arg *GetAnalysisParams) (*Analysis, error)
- func (q *Queries) GetUser(ctx context.Context, arg *GetUserParams) (*User, error)
- func (q *Queries) ListUsers(ctx context.Context) ([]*User, error)
- func (q *Queries) WithTx(tx pgx.Tx) *Queries
- type ReadQuerier
- type ReadQueries
- func (q *ReadQueries) GetAnalysis(ctx context.Context, arg *GetAnalysisParams) (*Analysis, error)
- func (q *ReadQueries) GetUser(ctx context.Context, arg *GetUserParams) (*User, error)
- func (q *ReadQueries) ListUsers(ctx context.Context) ([]*User, error)
- func (q *ReadQueries) WithTx(tx pgx.Tx) *ReadQueries
- type RouteOption
- type ShardResolver
- type ShardedQueries
- func (q *ShardedQueries[SK]) CreateUser(ctx context.Context, arg *CreateUserParams, routeOptions ...RouteOption) (*User, error)
- func (q *ShardedQueries[SK]) GetAnalysis(ctx context.Context, arg *GetAnalysisParams, routeOptions ...RouteOption) (*Analysis, error)
- func (q *ShardedQueries[SK]) GetUser(ctx context.Context, arg *GetUserParams, routeOptions ...RouteOption) (*User, error)
- type StoreQuerier
- type StoreQueries
- type User
- type WriteQuerier
- type WriteQueries
- func (q *WriteQueries) CopyUsers(ctx context.Context, arg []*CopyUsersParams) (int64, error)
- func (q *WriteQueries) CreateUser(ctx context.Context, arg *CreateUserParams) (*User, error)
- func (q *WriteQueries) WithMirrors(qs ...*WriteQueries) *WriteQueries
- func (q *WriteQueries) WithTx(tx pgx.Tx) *WriteQueries
Examples ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func NewStoreNode ¶
func NewStoreNode(database DBTX) pgmesh.Node[*ReadQueries, *StoreQueries]
NewStoreNode creates a pgmesh node backed by database.
Types ¶
type Analysis ¶
type Analysis struct {
ID int64
TenantID int64
Description *string
State NullAnalysisState
Source netip.Addr
ActiveWindow pgtype.Range[pgtype.Timestamptz]
}
type AnalysisState ¶
type AnalysisState string
const ( AnalysisStatePending AnalysisState = "pending" AnalysisStateComplete AnalysisState = "complete" )
func (*AnalysisState) Scan ¶
func (e *AnalysisState) Scan(src interface{}) error
type CopyUsersParams ¶
type CreateUserParams ¶
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
CopyFrom(ctx context.Context, tableName pgx.Identifier, columnNames []string, rowSrc pgx.CopyFromSource) (int64, error)
}
type GetAnalysisParams ¶
type GetUserParams ¶
type NullAnalysisState ¶
type NullAnalysisState struct {
AnalysisState AnalysisState
Valid bool // Valid is true if AnalysisState is not NULL
}
func (*NullAnalysisState) Scan ¶
func (ns *NullAnalysisState) Scan(value interface{}) error
Scan implements the Scanner interface.
type Querier ¶
type Querier interface {
// kind: write
CopyUsers(ctx context.Context, arg []*CopyUsersParams) (int64, error)
// kind: write
// shard: tenant(tenant_id)
CreateUser(ctx context.Context, arg *CreateUserParams) (*User, error)
// kind: read
// shard: tenant(tenant_id)
GetAnalysis(ctx context.Context, arg *GetAnalysisParams) (*Analysis, error)
// kind: read
// shard: tenant(tenant_id)
GetUser(ctx context.Context, arg *GetUserParams) (*User, error)
// kind: read
ListUsers(ctx context.Context) ([]*User, error)
}
type Queries ¶
type Queries struct {
// contains filtered or unexported fields
}
func (*Queries) CreateUser ¶
kind: write shard: tenant(tenant_id)
func (*Queries) GetAnalysis ¶
kind: read shard: tenant(tenant_id)
type ReadQuerier ¶
type ReadQuerier interface {
// GetAnalysis executes the generated GetAnalysis query.
GetAnalysis(ctx context.Context, arg *GetAnalysisParams) (*Analysis, error)
// GetUser executes the generated GetUser query.
GetUser(ctx context.Context, arg *GetUserParams) (*User, error)
// ListUsers executes the generated ListUsers query.
ListUsers(ctx context.Context) ([]*User, error)
}
ReadQuerier exposes generated read queries.
type ReadQueries ¶
type ReadQueries struct {
// contains filtered or unexported fields
}
ReadQueries exposes read-only generated queries.
func NewReadQueries ¶
func NewReadQueries(database DBTX) *ReadQueries
NewReadQueries creates a read-only query wrapper.
func (*ReadQueries) GetAnalysis ¶
func (q *ReadQueries) GetAnalysis(ctx context.Context, arg *GetAnalysisParams) (*Analysis, error)
GetAnalysis executes the generated GetAnalysis query.
func (*ReadQueries) GetUser ¶
func (q *ReadQueries) GetUser(ctx context.Context, arg *GetUserParams) (*User, error)
GetUser executes the generated GetUser query.
func (*ReadQueries) ListUsers ¶
func (q *ReadQueries) ListUsers(ctx context.Context) ([]*User, error)
ListUsers executes the generated ListUsers query.
func (*ReadQueries) WithTx ¶
func (q *ReadQueries) WithTx(tx pgx.Tx) *ReadQueries
WithTx returns a read wrapper that executes queries through tx.
type RouteOption ¶
type RouteOption func(*routeOptions)
RouteOption customizes routing for one generated query call.
func ReadFromPrimary ¶
func ReadFromPrimary() RouteOption
ReadFromPrimary routes a read query to the shard's primary.
func WithTx ¶
func WithTx(tx pgx.Tx) RouteOption
WithTx routes a query through tx and suppresses write mirrors.
type ShardResolver ¶
type ShardResolver[SK any] interface { // Tenant resolves the "tenant" shard route. Tenant(tenantID int64) SK }
ShardResolver resolves generated query parameters to shard keys.
type ShardedQueries ¶
type ShardedQueries[SK any] struct { // contains filtered or unexported fields }
ShardedQueries routes generated queries through a pgmesh mesh.
Example ¶
log := &callLog{}
primary := NewStoreNode(&fakeDB{name: "primary", log: log})
replica := NewStoreNode(&fakeDB{name: "replica", log: log})
mirror := NewStoreNode(&fakeDB{name: "mirror", log: log})
replicaSet := pgmesh.NewReplicaSet(
"main",
primary,
[]pgmesh.Node[*ReadQueries, *StoreQueries]{replica},
).WithWriteMirrors(mirror.Writer())
mesh, err := pgmesh.NewBuilder[*ReadQueries, *StoreQueries, uint64](1).
WithHasher(pgmesh.ConstantShardHashFor[uint64](0)).
Link(0, replicaSet).
Build()
if err != nil {
panic(err)
}
queries := NewShardedQueries(mesh, tenantResolver{})
ctx := context.Background()
if _, err := queries.GetUser(ctx, &GetUserParams{TenantID: 10, ID: 20}); err != nil {
panic(err)
}
if _, err := queries.GetUser(ctx, &GetUserParams{TenantID: 10, ID: 20}, ReadFromPrimary()); err != nil {
panic(err)
}
if _, err := queries.CreateUser(ctx, &CreateUserParams{ID: 20, TenantID: 10, Name: "user"}); err != nil {
panic(err)
}
fmt.Println(log.snapshot())
Output: [replica primary primary mirror]
func NewShardedQueries ¶
func NewShardedQueries[SK any](mesh *pgmesh.Mesh[*ReadQueries, *StoreQueries, SK], resolver ShardResolver[SK]) *ShardedQueries[SK]
NewShardedQueries creates a routed query facade.
func (*ShardedQueries[SK]) CreateUser ¶
func (q *ShardedQueries[SK]) CreateUser(ctx context.Context, arg *CreateUserParams, routeOptions ...RouteOption) (*User, error)
CreateUser executes the generated query on its resolved shard.
func (*ShardedQueries[SK]) GetAnalysis ¶
func (q *ShardedQueries[SK]) GetAnalysis(ctx context.Context, arg *GetAnalysisParams, routeOptions ...RouteOption) (*Analysis, error)
GetAnalysis executes the generated query on its resolved shard.
func (*ShardedQueries[SK]) GetUser ¶
func (q *ShardedQueries[SK]) GetUser(ctx context.Context, arg *GetUserParams, routeOptions ...RouteOption) (*User, error)
GetUser executes the generated query on its resolved shard.
type StoreQuerier ¶
type StoreQuerier interface {
ReadQuerier
WriteQuerier
}
StoreQuerier combines the generated read and write query interfaces.
type StoreQueries ¶
type StoreQueries struct {
*ReadQueries
*WriteQueries
}
StoreQueries combines read-only and primary-capable generated queries.
func NewStoreQueries ¶
func NewStoreQueries(database DBTX) *StoreQueries
NewStoreQueries creates a combined query wrapper.
func (*StoreQueries) WithMirrors ¶
func (q *StoreQueries) WithMirrors(qs ...*StoreQueries) *StoreQueries
WithMirrors returns a copy that also writes to the supplied mirrors.
func (*StoreQueries) WithTx ¶
func (q *StoreQueries) WithTx(tx pgx.Tx) *StoreQueries
WithTx returns a store wrapper that executes queries through tx.
type WriteQuerier ¶
type WriteQuerier interface {
// CopyUsers executes the generated CopyUsers query.
CopyUsers(ctx context.Context, arg []*CopyUsersParams) (int64, error)
// CreateUser executes the generated CreateUser query.
CreateUser(ctx context.Context, arg *CreateUserParams) (*User, error)
}
WriteQuerier exposes generated write queries.
type WriteQueries ¶
type WriteQueries struct {
// contains filtered or unexported fields
}
WriteQueries exposes primary-capable generated queries.
func NewWriteQueries ¶
func NewWriteQueries(database DBTX) *WriteQueries
NewWriteQueries creates a primary-capable query wrapper.
func (*WriteQueries) CopyUsers ¶
func (q *WriteQueries) CopyUsers(ctx context.Context, arg []*CopyUsersParams) (int64, error)
CopyUsers executes the generated CopyUsers query.
func (*WriteQueries) CreateUser ¶
func (q *WriteQueries) CreateUser(ctx context.Context, arg *CreateUserParams) (*User, error)
CreateUser executes the generated CreateUser query.
func (*WriteQueries) WithMirrors ¶
func (q *WriteQueries) WithMirrors(qs ...*WriteQueries) *WriteQueries
WithMirrors returns a copy that also writes to the supplied mirrors.
func (*WriteQueries) WithTx ¶
func (q *WriteQueries) WithTx(tx pgx.Tx) *WriteQueries
WithTx returns a write wrapper that executes queries through tx.