Documentation
¶
Index ¶
- Constants
- func ContextWithAuth(ctx context.Context, userID, userRole string, isAdmin bool) context.Context
- func ContextWithTenant(ctx context.Context, tenantID string) context.Context
- func ContextWithTenantID(ctx context.Context, tenantID string) context.Context
- func ExtractDDLMetadata(sql string) string
- func ExtractOperation(sql string) string
- func ExtractTableName(sql string) string
- func GetConstraintName(err error) string
- func IsCheckViolation(err error) bool
- func IsForeignKeyViolation(err error) bool
- func IsTransientError(err error) bool
- func IsUniqueViolation(err error) bool
- func LogSchemaIntrospection(ctx context.Context, operation string, details map[string]interface{})
- func PoolQuerier(pool *pgxpool.Pool) querier
- func RetryTransient(ctx context.Context, name string, rc RetryConfig, fn func() error) error
- func TenantFromContext(ctx context.Context) string
- func TenantIDFromContext(ctx context.Context) string
- func TenantOrNil(tenantID string) interface{}
- func WithCaller(ctx context.Context, caller string) context.Context
- func WrapWithServiceRole(ctx context.Context, conn *Connection, fn func(tx pgx.Tx) error) error
- func WrapWithServiceRoleAndTenant(ctx context.Context, conn *Connection, tenantID string, ...) error
- func WrapWithTenantAwareRole(ctx context.Context, conn *Connection, tenantID string, ...) error
- func WrapWithTenantContext(ctx context.Context, conn *Connection, tenantID string, ...) error
- type AdminExecutor
- type AuthContext
- type AuthContextKey
- type ColumnInfo
- type Connection
- func (c *Connection) BeginTx(ctx context.Context) (pgx.Tx, error)
- func (c *Connection) Close()
- func (c *Connection) Exec(ctx context.Context, sql string, args ...interface{}) (pgconn.CommandTag, error)
- func (c *Connection) ExecuteWithAdminRole(ctx context.Context, fn func(tx pgx.Tx) error) error
- func (c *Connection) ExecuteWithAdminRoleForDB(ctx context.Context, dbName string, fn func(tx pgx.Tx) error) error
- func (c *Connection) Health(ctx context.Context) error
- func (c *Connection) Inspector() *SchemaInspector
- func (c *Connection) Migrate() error
- func (c *Connection) Pool() *pgxpool.Pool
- func (c *Connection) Query(ctx context.Context, sql string, args ...interface{}) (pgx.Rows, error)
- func (c *Connection) QueryRow(ctx context.Context, sql string, args ...interface{}) pgx.Row
- func (c *Connection) RecreatePool() error
- func (c *Connection) SetMetrics(m *observability.Metrics)
- func (c *Connection) Stats() *pgxpool.Stat
- type Executor
- type ForeignKey
- type FunctionInfo
- type FunctionParam
- type IndexInfo
- type JSONBProperty
- type JSONBSchemaInfo
- type Querier
- type RetryConfig
- type SchemaCache
- func (c *SchemaCache) Close()
- func (c *SchemaCache) GetAllMaterializedViews(ctx context.Context) ([]TableInfo, error)
- func (c *SchemaCache) GetAllTables(ctx context.Context) ([]TableInfo, error)
- func (c *SchemaCache) GetAllViews(ctx context.Context) ([]TableInfo, error)
- func (c *SchemaCache) GetSchemas(ctx context.Context) ([]string, error)
- func (c *SchemaCache) GetTable(ctx context.Context, schema, table string) (*TableInfo, bool, error)
- func (c *SchemaCache) Invalidate()
- func (c *SchemaCache) InvalidateAll(ctx context.Context)
- func (c *SchemaCache) IsTableWritable(ctx context.Context, schema, table string) (bool, error)
- func (c *SchemaCache) Refresh(ctx context.Context) error
- func (c *SchemaCache) SetPubSub(ps pubsub.PubSub)
- func (c *SchemaCache) TableCount() int
- func (c *SchemaCache) ViewCount() int
- type SchemaInspector
- func (si *SchemaInspector) BuildRESTPath(table TableInfo) string
- func (si *SchemaInspector) GetAllFunctions(ctx context.Context, schemas ...string) ([]FunctionInfo, error)
- func (si *SchemaInspector) GetAllMaterializedViews(ctx context.Context, schemas ...string) ([]TableInfo, error)
- func (si *SchemaInspector) GetAllMaterializedViewsFromQ(ctx context.Context, q querier, schemas ...string) ([]TableInfo, error)
- func (si *SchemaInspector) GetAllTables(ctx context.Context, schemas ...string) ([]TableInfo, error)
- func (si *SchemaInspector) GetAllTablesFromQ(ctx context.Context, q querier, schemas ...string) ([]TableInfo, error)
- func (si *SchemaInspector) GetAllViews(ctx context.Context, schemas ...string) ([]TableInfo, error)
- func (si *SchemaInspector) GetAllViewsFromQ(ctx context.Context, q querier, schemas ...string) ([]TableInfo, error)
- func (si *SchemaInspector) GetSchemas(ctx context.Context) ([]string, error)
- func (si *SchemaInspector) GetSchemasFromQ(ctx context.Context, q querier) ([]string, error)
- func (si *SchemaInspector) GetTableInfo(ctx context.Context, schema, table string) (*TableInfo, error)
- func (si *SchemaInspector) GetTableInfoFromQ(ctx context.Context, q querier, schema, table string) (*TableInfo, error)
- func (si *SchemaInspector) GetVectorColumns(ctx context.Context, schema, table string) ([]VectorColumnInfo, error)
- func (si *SchemaInspector) IsPgVectorInstalled(ctx context.Context) (bool, string, error)
- type TableInfo
- type TenantAware
- type TenantContextKey
- type TxConnection
- type VectorColumnInfo
Constants ¶
const ( // ErrCodeUniqueViolation is the PostgreSQL error code for unique constraint violations ErrCodeUniqueViolation = "23505" // ErrCodeForeignKeyViolation is the PostgreSQL error code for foreign key violations ErrCodeForeignKeyViolation = "23503" // ErrCodeCheckViolation is the PostgreSQL error code for check constraint violations ErrCodeCheckViolation = "23514" )
PostgreSQL error codes
Variables ¶
This section is empty.
Functions ¶
func ContextWithAuth ¶
ContextWithAuth returns a context with auth information for audit logging
func ContextWithTenant ¶
ContextWithTenant returns a context with tenant ID set
func ContextWithTenantID ¶
ContextWithTenantID returns a context with tenant ID for multi-tenancy support This is an alias for ContextWithTenant for API consistency
func ExtractDDLMetadata ¶
ExtractDDLMetadata extracts operation type and target from a DDL query for logging Returns a safe, redacted string like "CREATE TABLE users", "DROP INDEX idx_name"
func ExtractOperation ¶
ExtractOperation extracts the SQL operation type from a query
func ExtractTableName ¶
ExtractTableName attempts to extract the table name from a SQL query Returns "unknown" if the table cannot be determined
func GetConstraintName ¶
GetConstraintName returns the constraint name from a PostgreSQL error
func IsCheckViolation ¶
IsCheckViolation checks if an error is a check constraint violation
func IsForeignKeyViolation ¶
IsForeignKeyViolation checks if an error is a foreign key violation
func IsTransientError ¶
IsTransientError reports whether err represents a failure that may succeed on retry: connection/network problems, the server temporarily refusing connections (e.g. still starting up), context deadlines during connect, or transient SQL states (serialization/deadlock/resource limits). It is used to decide whether to retry DB-dependent startup steps instead of aborting.
Non-transient errors (constraint violations, syntax errors, permission errors, nil) report false so they fail fast rather than retrying forever.
func IsUniqueViolation ¶
IsUniqueViolation checks if an error is a unique constraint violation
func LogSchemaIntrospection ¶
LogSchemaIntrospection logs schema introspection for audit purposes
func PoolQuerier ¶
func RetryTransient ¶
RetryTransient runs fn, retrying only while it returns a transient error (per IsTransientError). Non-transient errors and nil return immediately. Between retries it sleeps with exponential backoff (jittered, capped) that is interruptible by ctx. name is used in log lines to identify the step.
This wraps DB-dependent startup steps (bootstrap, schema apply, initial connection) so a transient blip after the initial connect doesn't abort startup. On persistent real errors the original error is returned unchanged.
func TenantFromContext ¶
TenantFromContext extracts tenant ID from context. Returns empty string if no tenant context is set
func TenantIDFromContext ¶
TenantIDFromContext extracts tenant ID from context Returns empty string if no tenant context is set This is an alias for TenantFromContext for API consistency
func TenantOrNil ¶
func TenantOrNil(tenantID string) interface{}
TenantOrNil converts an empty tenant string to nil for UUID column compatibility. PostgreSQL UUID columns accept NULL but reject empty strings, so this helper is used when passing tenant IDs as query parameters.
func WrapWithServiceRole ¶
WrapWithServiceRole wraps a database operation with service_role context Used for privileged operations like auth, admin tasks, and webhooks
func WrapWithServiceRoleAndTenant ¶
func WrapWithServiceRoleAndTenant(ctx context.Context, conn *Connection, tenantID string, fn func(tx pgx.Tx) error) error
WrapWithServiceRoleAndTenant wraps a database operation with both service_role and tenant context. This bypasses RLS but still sets tenant_id for new records via the set_tenant_id trigger. Use this for privileged operations that still need to associate records with a tenant.
func WrapWithTenantAwareRole ¶
func WrapWithTenantAwareRole(ctx context.Context, conn *Connection, tenantID string, fn func(tx pgx.Tx) error) error
WrapWithTenantAwareRole wraps a database operation with the appropriate role based on tenant context. When a tenant context is active, it uses tenant_service (NOBYPASSRLS) so RLS policies enforce tenant isolation. When no tenant context, it uses service_role (BYPASSRLS) for full instance-admin access.
func WrapWithTenantContext ¶
func WrapWithTenantContext(ctx context.Context, conn *Connection, tenantID string, fn func(tx pgx.Tx) error) error
WrapWithTenantContext wraps a database operation with tenant context for multi-tenancy. This sets the app.current_tenant_id session variable so that RLS policies and triggers can enforce tenant isolation. Use this for storage operations on tenant-scoped tables.
Types ¶
type AdminExecutor ¶
type AdminExecutor interface {
Executor
// ExecuteWithAdminRole executes a database operation using admin credentials
// inside a transaction. Used for migrations that require DDL privileges
// (CREATE TABLE, ALTER, etc.). The callback receives the transaction, not a
// raw connection, so all statements are properly transactional.
ExecuteWithAdminRole(ctx context.Context, fn func(tx pgx.Tx) error) error
// Inspector returns the schema inspector for introspecting database structure
Inspector() *SchemaInspector
}
AdminExecutor extends Executor with privileged operations that require admin database credentials (e.g., for DDL operations, migrations).
type AuthContext ¶
type AuthContext struct {
UserID string
UserRole string // "authenticated", "anon", "service_role", "admin", etc.
IsAdmin bool
ClientKey string // For service role access via client keys
}
AuthContext represents authentication and authorization context for schema introspection
func AuthFromContext ¶
func AuthFromContext(ctx context.Context) *AuthContext
AuthFromContext extracts auth context from context for audit logging Returns nil if no auth context is set
type AuthContextKey ¶
type AuthContextKey struct{}
AuthContextKey is the context key for storing authorization information
type ColumnInfo ¶
type ColumnInfo struct {
Name string `json:"name"`
DataType string `json:"data_type"`
IsNullable bool `json:"is_nullable"`
DefaultValue *string `json:"default_value"`
IsPrimaryKey bool `json:"is_primary_key"`
IsForeignKey bool `json:"is_foreign_key"`
IsUnique bool `json:"is_unique"`
MaxLength *int `json:"max_length"`
Position int `json:"position"`
Description string `json:"description,omitempty"`
JSONBSchema *JSONBSchemaInfo `json:"jsonb_schema,omitempty"`
}
type Connection ¶
type Connection struct {
// contains filtered or unexported fields
}
Connection represents a database connection pool
func ConnectWithRetry ¶
func ConnectWithRetry(cfg config.DatabaseConfig, maxAttempts int) (*Connection, error)
ConnectWithRetry attempts to connect to the database with exponential backoff.
It retries while NewConnection returns a transient error (per IsTransientError) and fails fast on non-transient errors. The backoff is interruptible by ctx so SIGTERM during startup isn't delayed. maxAttempts <= 0 falls back to a default.
func NewConnection ¶
func NewConnection(cfg config.DatabaseConfig) (*Connection, error)
NewConnection creates a new database connection pool The connection pool uses the runtime user, while migrations use the admin user
func NewConnectionWithPool ¶
func NewConnectionWithPool(pool *pgxpool.Pool) *Connection
NewConnectionWithPool creates a new Connection wrapper around an existing pgxpool.Pool. This is useful for tests where you have a pre-configured pool.
func (*Connection) Exec ¶
func (c *Connection) Exec(ctx context.Context, sql string, args ...interface{}) (pgconn.CommandTag, error)
Exec executes a query that doesn't return rows
func (*Connection) ExecuteWithAdminRole ¶
ExecuteWithAdminRole executes a database operation using admin credentials Used for migrations that require DDL privileges (CREATE TABLE, ALTER, etc.) Creates a temporary admin connection that is closed after execution
func (*Connection) ExecuteWithAdminRoleForDB ¶
func (c *Connection) ExecuteWithAdminRoleForDB(ctx context.Context, dbName string, fn func(tx pgx.Tx) error) error
ExecuteWithAdminRoleForDB executes a function with admin privileges against a specific database (for tenant DDL operations). It replaces the database name in the admin connection string with the provided dbName.
func (*Connection) Health ¶
func (c *Connection) Health(ctx context.Context) error
Health checks the health of the database connection
func (*Connection) Inspector ¶
func (c *Connection) Inspector() *SchemaInspector
Inspector returns the schema inspector
func (*Connection) Migrate ¶
func (c *Connection) Migrate() error
Migrate runs database migrations from user sources Note: Internal Fluxbase schema is now managed declaratively (see bootstrap + pgschema)
func (*Connection) Pool ¶
func (c *Connection) Pool() *pgxpool.Pool
Pool returns the underlying connection pool
func (*Connection) RecreatePool ¶
func (c *Connection) RecreatePool() error
RecreatePool closes the current pool and creates a new one. This is safer than Reset() as it ensures a completely fresh pool state. Use this after schema changes (migrations) to avoid prepared statement cache issues.
func (*Connection) SetMetrics ¶
func (c *Connection) SetMetrics(m *observability.Metrics)
SetMetrics sets the metrics instance for recording database metrics
func (*Connection) Stats ¶
func (c *Connection) Stats() *pgxpool.Stat
Stats returns database connection pool statistics
type Executor ¶
type Executor interface {
// Query executes a query that returns rows
Query(ctx context.Context, sql string, args ...interface{}) (pgx.Rows, error)
// QueryRow executes a query that returns a single row
QueryRow(ctx context.Context, sql string, args ...interface{}) pgx.Row
// Exec executes a query that doesn't return rows
Exec(ctx context.Context, sql string, args ...interface{}) (pgconn.CommandTag, error)
// BeginTx starts a new transaction
BeginTx(ctx context.Context) (pgx.Tx, error)
// Pool returns the underlying connection pool (for advanced operations)
Pool() *pgxpool.Pool
// Health checks the health of the database connection
Health(ctx context.Context) error
}
Executor defines the interface for database query operations. This interface abstracts database access for easier testing via mocks.
type ForeignKey ¶
type FunctionInfo ¶
type FunctionInfo struct {
Schema string `json:"schema"`
Name string `json:"name"`
Description string `json:"description"`
Parameters []FunctionParam `json:"parameters"`
ReturnType string `json:"return_type"`
IsSetOf bool `json:"is_set_of"`
Volatility string `json:"volatility"`
Language string `json:"language"`
}
type FunctionParam ¶
type JSONBProperty ¶
type JSONBProperty struct {
Type string `json:"type"`
Description string `json:"description,omitempty"`
Properties map[string]JSONBProperty `json:"properties,omitempty"`
Items *JSONBProperty `json:"items,omitempty"`
}
type JSONBSchemaInfo ¶
type JSONBSchemaInfo struct {
Properties map[string]JSONBProperty `json:"properties,omitempty"`
Required []string `json:"required,omitempty"`
}
type Querier ¶
type Querier interface{}
These aliases allow the middleware and handlers to use simpler type names
type RetryConfig ¶
type RetryConfig struct {
MaxAttempts int // total attempts (>=1)
InitialBackoff time.Duration // backoff for the first retry
MaxBackoff time.Duration // cap on each backoff
}
RetryConfig controls the retry behavior of RetryTransient and ConnectWithRetry.
func DefaultRetryConfig ¶
func DefaultRetryConfig() RetryConfig
DefaultRetryConfig returns sensible defaults: 8 attempts, 1s..30s backoff.
type SchemaCache ¶
type SchemaCache struct {
// contains filtered or unexported fields
}
SchemaCache provides a thread-safe cache for database schema information with TTL-based expiration and manual invalidation support. When PubSub is configured, invalidation is broadcast to all instances.
func NewSchemaCache ¶
func NewSchemaCache(inspector *SchemaInspector, ttl time.Duration) *SchemaCache
NewSchemaCache creates a new schema cache with the given TTL
func (*SchemaCache) Close ¶
func (c *SchemaCache) Close()
Close stops the invalidation listener if running
func (*SchemaCache) GetAllMaterializedViews ¶
func (c *SchemaCache) GetAllMaterializedViews(ctx context.Context) ([]TableInfo, error)
GetAllMaterializedViews returns all cached materialized views, refreshing if necessary
func (*SchemaCache) GetAllTables ¶
func (c *SchemaCache) GetAllTables(ctx context.Context) ([]TableInfo, error)
GetAllTables returns all cached tables, refreshing if necessary
func (*SchemaCache) GetAllViews ¶
func (c *SchemaCache) GetAllViews(ctx context.Context) ([]TableInfo, error)
GetAllViews returns all cached views, refreshing if necessary
func (*SchemaCache) GetSchemas ¶
func (c *SchemaCache) GetSchemas(ctx context.Context) ([]string, error)
GetSchemas returns cached schemas
func (*SchemaCache) GetTable ¶
GetTable retrieves table info from the cache, refreshing if necessary. Returns (TableInfo, exists, error)
func (*SchemaCache) Invalidate ¶
func (c *SchemaCache) Invalidate()
Invalidate marks the cache as stale, forcing a refresh on next access. This only invalidates the local cache. Use InvalidateAll to broadcast invalidation to all instances.
func (*SchemaCache) InvalidateAll ¶
func (c *SchemaCache) InvalidateAll(ctx context.Context)
InvalidateAll marks the cache as stale and broadcasts the invalidation to all other instances via PubSub. Use this when schema changes occur (e.g., after migrations) to ensure all instances refresh their caches.
func (*SchemaCache) IsTableWritable ¶
IsTableWritable checks if a table is writable (not a view or materialized view)
func (*SchemaCache) Refresh ¶
func (c *SchemaCache) Refresh(ctx context.Context) error
Refresh forces an immediate cache refresh
func (*SchemaCache) SetPubSub ¶
func (c *SchemaCache) SetPubSub(ps pubsub.PubSub)
SetPubSub configures the PubSub backend for cross-instance cache invalidation. When set, InvalidateAll will broadcast invalidation messages to all instances, and this instance will listen for invalidation messages from others.
func (*SchemaCache) TableCount ¶
func (c *SchemaCache) TableCount() int
TableCount returns the number of cached tables
func (*SchemaCache) ViewCount ¶
func (c *SchemaCache) ViewCount() int
ViewCount returns the number of cached views
type SchemaInspector ¶
type SchemaInspector struct {
// contains filtered or unexported fields
}
func NewSchemaInspector ¶
func NewSchemaInspector(conn *Connection) *SchemaInspector
func (*SchemaInspector) BuildRESTPath ¶
func (si *SchemaInspector) BuildRESTPath(table TableInfo) string
func (*SchemaInspector) GetAllFunctions ¶
func (si *SchemaInspector) GetAllFunctions(ctx context.Context, schemas ...string) ([]FunctionInfo, error)
func (*SchemaInspector) GetAllMaterializedViews ¶
func (*SchemaInspector) GetAllMaterializedViewsFromQ ¶
func (*SchemaInspector) GetAllTables ¶
func (*SchemaInspector) GetAllTablesFromQ ¶
func (*SchemaInspector) GetAllViews ¶
func (*SchemaInspector) GetAllViewsFromQ ¶
func (*SchemaInspector) GetSchemas ¶
func (si *SchemaInspector) GetSchemas(ctx context.Context) ([]string, error)
func (*SchemaInspector) GetSchemasFromQ ¶
func (si *SchemaInspector) GetSchemasFromQ(ctx context.Context, q querier) ([]string, error)
func (*SchemaInspector) GetTableInfo ¶
func (*SchemaInspector) GetTableInfoFromQ ¶
func (*SchemaInspector) GetVectorColumns ¶
func (si *SchemaInspector) GetVectorColumns(ctx context.Context, schema, table string) ([]VectorColumnInfo, error)
func (*SchemaInspector) IsPgVectorInstalled ¶
type TableInfo ¶
type TableInfo struct {
Schema string `json:"schema"`
Name string `json:"name"`
Type string `json:"type"`
RESTPath string `json:"rest_path,omitempty"`
Columns []ColumnInfo `json:"columns"`
PrimaryKey []string `json:"primary_key"`
ForeignKeys []ForeignKey `json:"foreign_keys"`
Indexes []IndexInfo `json:"indexes"`
RLSEnabled bool `json:"rls_enabled"`
ColumnMap map[string]*ColumnInfo `json:"-"`
}
func (*TableInfo) BuildColumnMap ¶
func (t *TableInfo) BuildColumnMap()
func (*TableInfo) GetColumn ¶
func (t *TableInfo) GetColumn(name string) *ColumnInfo
type TenantAware ¶
type TenantAware struct {
DB *Connection
}
TenantAware provides a reusable embedded struct for tenant-scoped database operations. Storage types embed this to get the WithTenant helper method.
func (*TenantAware) WithTenant ¶
WithTenant wraps a database operation with tenant-aware role selection. When a tenant context is active, uses tenant_service (respects RLS). When no tenant context, uses service_role (bypasses RLS).
type TenantContextKey ¶
type TenantContextKey struct{}
TenantContextKey is the context key for storing tenant information
type TxConnection ¶
These aliases allow the middleware and handlers to use simpler type names