Documentation
¶
Index ¶
- Constants
- func DefaultGUCs(version *ds.PgEdgeVersion, peerInstanceIDs []string) map[string]any
- func DefaultTunableGUCs(memBytes uint64, cpus float64, clusterSize int) map[string]any
- func IsDatabaseNotExists(err error) bool
- func IsSpockNodeNotConfigured(err error) bool
- func NeedsNativeFailoverSlots(spockMajor, pgMajor uint64) bool
- func NeedsNativeFailoverSlotsForVersion(version *ds.PgEdgeVersion) bool
- func QuoteIdentifier(in string) string
- func SnowflakeLolorGUCs(nodeOrdinal int) map[string]any
- func SpockDefaultGUCs() map[string]any
- func StartRepairModeTxn(ctx context.Context, conn *pgx.Conn) (pgx.Tx, error)
- func UnquoteIdentifier(in string) string
- type BuiltinRoleOptions
- type BuiltinRolePrivilegeOptions
- type ConditionalStatement
- func CreateDatabase(name string) ConditionalStatement
- func CreateReplicationSlot(databaseName, providerNode, subscriberNode string, failover bool) ConditionalStatement
- func CreateRoleIfNotExists(name string) ConditionalStatement
- func CreateSubscription(providerName, subscriberName string, providerDSN *DSN, disabled bool, ...) ConditionalStatement
- func DropReplicationSlot(databaseName, providerNode, subscriberNode string) ConditionalStatement
- func DropStaleReplicationOrigin(databaseName, providerNode, subscriberNode string) ConditionalStatement
- func EnableSubscription(providerNode, subscriberNode string, disabled bool) ConditionalStatement
- func EnsureReplicationOriginExists(slotName string) ConditionalStatement
- func RenameDB(oldName, newName string) ConditionalStatement
- func TerminateReplicationSlot(databaseName, providerNode, subscriberNode string) ConditionalStatement
- type DSN
- type Executor
- type IStatement
- type PgServiceConf
- type Query
- func CurrentReplicationSlotLSN(databaseName, providerNode, subscriberNode string) Query[string]
- func GetPostgresVersion() Query[string]
- func GetReplicationSetTables() Query[ReplicationSetTable]
- func GetReplicationSets() Query[ReplicationSet]
- func GetReplicationSlotLSNFromCommitTS(databaseName, providerNode, subscriberNode string, commitTS time.Time) Query[string]
- func GetSpockReadOnly() Query[string]
- func GetSpockVersion() Query[string]
- func GetSubscriptionStatus(providerNode, subscriberNode string) Query[string]
- func GetSubscriptionStatuses() Query[SubscriptionStatus]
- func IsInRecovery() Query[bool]
- func IsReplicationSlotActive(databaseName, providerNode, subscriberNode string) Query[bool]
- func IsSpockEnabled() Query[bool]
- func LagTrackerCommitTimestamp(originNode, receiverNode string) Query[time.Time]
- func LsnAtOrBefore(lsn1, lsn2 string) Query[bool]
- func NodeNeedsCreate(nodeName string) Query[bool]
- func ReplicationSlotExists(databaseName, providerNode, subscriberNode string) Query[bool]
- func ReplicationSlotNeedsCreate(databaseName, providerNode, subscriberNode string) Query[bool]
- func ResolveSlotName(databaseName, providerNode, subscriberNode string) Query[string]
- func SpockProgressReachedLSN(spockMajor uint64, peerNodeName, targetLSN string) Query[bool]
- func SubscriptionDsnNeedsUpdate(providerName, subscriberName string, providerDSN *DSN) Query[bool]
- func SubscriptionNeedsCreate(providerName, subscriberName string) Query[bool]
- func SubscriptionNeedsEnable(providerName, subscriberName string, disabled bool) Query[bool]
- func SyncEvent(transactional bool) Query[string]
- func UserRoleNeedsCreate(name string) Query[bool]
- func WaitForSyncEvent(originNode, lsn string, timeoutSeconds int, waitIfDisabled bool) Query[bool]
- type ReplicationSet
- type ReplicationSetTable
- type Statement
- func AddReplicationSetTable(setName string, relOID uint32, columns []string, rowFilter string, sync bool, ...) Statement
- func AdvanceReplicationOrigin(slotName, lsn string) Statement
- func AdvanceReplicationSlotToLSN(databaseName, providerNode, subscriberNode string, targetLSN string) Statement
- func AlterOwner(dbName, owner string) Statement
- func CreateReplicationSet(r ReplicationSet) Statement
- func DropAllSubscriptions() Statement
- func DropSubscription(providerName, subscriberName string) Statement
- func EnableRepairMode() Statement
- func TerminateOtherConnections(dbName string) Statement
- type Statements
- func CreateBuiltInRoles(opts BuiltinRoleOptions) (Statements, error)
- func CreatePgEdgeSuperuserRole(opts BuiltinRoleOptions) (Statements, error)
- func CreateUserRole(opts UserRoleOptions) (Statements, error)
- func DropSpockAndCleanupSlots(dbName string) Statements
- func GrantBuiltinRolePrivileges(opts BuiltinRolePrivilegeOptions) Statements
- func InitializeSpockNode(nodeName string, nodeDSN *DSN) Statements
- func RestoreReplicationSets(sets []ReplicationSet, tabs []ReplicationSetTable) Statements
- func SetSafeIdentifiers() Statements
- type SubscriptionStatus
- type UserRoleOptions
Constants ¶
const ( SubStatusInitializing = "initializing" // Worker running, sync in progress SubStatusReplicating = "replicating" // Worker running, sync ready SubStatusUnknown = "unknown" // Worker running, no sync status record SubStatusDisabled = "disabled" // Worker not running, subscription disabled SubStatusDown = "down" // Worker not running, subscription enabled )
Spock subscription statuses returned by spock.sub_show_status(). See: https://github.com/pgEdge/spock/blob/main/src/spock_functions.c
const MinSpockVersionForSyncEventArgs = "5.0.7"
MinSpockVersionForSyncEventArgs is the first Spock version where spock.sync_event(boolean) and the 5-arg spock.wait_for_sync_event(..., wait_if_disabled) exist — both were introduced together in the 5.0.6->5.0.7 upgrade script. Callers must compare the live cluster's Spock version against this before passing transactional=true / waitIfDisabled=true to SyncEvent/WaitForSyncEvent below; older clusters only have the original call shapes.
Variables ¶
This section is empty.
Functions ¶
func DefaultGUCs ¶
func DefaultGUCs(version *ds.PgEdgeVersion, peerInstanceIDs []string) map[string]any
func DefaultTunableGUCs ¶
func IsDatabaseNotExists ¶ added in v0.8.0
func IsSpockNodeNotConfigured ¶ added in v0.7.0
IsSpockNodeNotConfigured reports whether err is a PostgreSQL error indicating that the current database has not been initialized as a spock node (SQLSTATE 55000 — object_not_in_prerequisite_state, as raised by spock.sync_event and related functions when spock.node_create has not been called or the node has been dropped via spock.node_drop).
Source: https://github.com/pgEdge/spock/blob/main/src/spock_functions.c PostgreSQL SQLSTATE reference: https://www.postgresql.org/docs/current/errcodes-appendix.html
func NeedsNativeFailoverSlots ¶
NeedsNativeFailoverSlots reports whether the given Spock and Postgres major versions require PG17+'s native logical-slot-failover mechanism (see DefaultGUCs). Gated on Spock major >= 6 only: 5.x had this behind an opt-in GUC at times, but Control Plane never manages that GUC, so no 5.x minor should trigger this.
func NeedsNativeFailoverSlotsForVersion ¶
func NeedsNativeFailoverSlotsForVersion(version *ds.PgEdgeVersion) bool
NeedsNativeFailoverSlotsForVersion is NeedsNativeFailoverSlots for a declared *ds.PgEdgeVersion; false if either major is unresolvable.
func QuoteIdentifier ¶ added in v0.8.0
QuoteIdentifier quotes and sanitizes identifiers. Callers must execute the SetEncode statements to consider the output of this method safe.
func SnowflakeLolorGUCs ¶
func SpockDefaultGUCs ¶ added in v0.9.0
func StartRepairModeTxn ¶ added in v0.8.0
StartRepairModeTxn will start a new transaction and, if Spock is enabled, enable repair mode for the transaction. Callers are responsible for calling Rollback and Commit on the returned transaction.
func UnquoteIdentifier ¶ added in v0.8.0
UnquoteIdentifier removes quotes from a quoted identifier.
Types ¶
type BuiltinRoleOptions ¶
type BuiltinRoleOptions struct {
PGVersion string
}
type BuiltinRolePrivilegeOptions ¶ added in v0.8.0
func (BuiltinRolePrivilegeOptions) Schemas ¶ added in v0.8.0
func (o BuiltinRolePrivilegeOptions) Schemas() []string
type ConditionalStatement ¶
type ConditionalStatement struct {
If Query[bool]
Then IStatement
Else IStatement
}
func CreateDatabase ¶
func CreateDatabase(name string) ConditionalStatement
func CreateReplicationSlot ¶
func CreateReplicationSlot(databaseName, providerNode, subscriberNode string, failover bool) ConditionalStatement
CreateReplicationSlot creates the logical replication slot backing a peer subscription. Pass failover from NeedsNativeFailoverSlots; false is byte-for-byte identical to the pre-failover-slot-support statement.
func CreateRoleIfNotExists ¶ added in v0.8.0
func CreateRoleIfNotExists(name string) ConditionalStatement
func CreateSubscription ¶
func DropReplicationSlot ¶ added in v0.7.0
func DropReplicationSlot(databaseName, providerNode, subscriberNode string) ConditionalStatement
func DropStaleReplicationOrigin ¶
func DropStaleReplicationOrigin(databaseName, providerNode, subscriberNode string) ConditionalStatement
DropStaleReplicationOrigin drops any pre-existing replication origin for a provider/subscriber pair before a fresh subscription is created for it. Mirrors zodan.sql's create_disable_subscriptions_and_slots step, which explicitly does this "so create_sub starts fresh at 0/0 (avoids stale-LSN data loss)": if a node is removed and later re-added under the same name, CP's slot/origin naming (spock.spock_gen_slot_name(), same as Spock's own) is deterministic on (database, provider, subscriber), so the new subscription would otherwise inherit whatever LSN the old incarnation's origin was left at. Safe to call unconditionally before create — an origin can only exist here if it's stale, since Create() only runs when spock.subscription has no row for this pair yet, and an origin never outlives its subscription's removal on its own.
func EnableSubscription ¶
func EnableSubscription(providerNode, subscriberNode string, disabled bool) ConditionalStatement
func EnsureReplicationOriginExists ¶ added in v0.8.1
func EnsureReplicationOriginExists(slotName string) ConditionalStatement
func RenameDB ¶
func RenameDB(oldName, newName string) ConditionalStatement
func TerminateReplicationSlot ¶ added in v0.7.0
func TerminateReplicationSlot(databaseName, providerNode, subscriberNode string) ConditionalStatement
TerminateReplicationSlot terminates the walsender process using a replication slot, if one is active. This must be called before dropping a slot whose subscriber has gone down, since pg_drop_replication_slot fails on active slots.
type DSN ¶
type PgServiceConf ¶ added in v0.8.0
func NewPgServiceConf ¶ added in v0.8.0
func NewPgServiceConf() *PgServiceConf
func (*PgServiceConf) String ¶ added in v0.8.0
func (c *PgServiceConf) String() string
type Query ¶
func GetPostgresVersion ¶
func GetReplicationSetTables ¶
func GetReplicationSetTables() Query[ReplicationSetTable]
func GetReplicationSets ¶
func GetReplicationSets() Query[ReplicationSet]
func GetSpockReadOnly ¶
func GetSpockVersion ¶
func GetSubscriptionStatus ¶ added in v0.7.0
GetSubscriptionStatus returns the current status of a specific subscription
func GetSubscriptionStatuses ¶
func GetSubscriptionStatuses() Query[SubscriptionStatus]
func IsInRecovery ¶ added in v0.8.0
func IsReplicationSlotActive ¶ added in v0.7.0
IsReplicationSlotActive checks if a replication slot is currently being used by an active walsender process. Uses EXISTS to always return exactly one row.
func IsSpockEnabled ¶
func LsnAtOrBefore ¶ added in v0.9.0
LsnAtOrBefore reports whether lsn1 <= lsn2 using PostgreSQL's pg_lsn type. Use this instead of Go string comparison — LSNs are hex-formatted and string ordering produces wrong results across segment boundaries (e.g. "F/..." > "10/...").
func NodeNeedsCreate ¶ added in v0.8.0
func ReplicationSlotExists ¶ added in v0.10.0
ReplicationSlotExists checks whether the replication slot for the given subscription exists. Used to poll for Spock 5.x failover slot creation after a switchover, which can take up to 60 seconds.
func ResolveSlotName ¶
ResolveSlotName looks up the actual slot name for a provider/subscriber pair. Used by call sites that need the resolved string itself — because they reuse it across more than one statement — rather than embedding slotNameExpr inline.
func SpockProgressReachedLSN ¶ added in v0.8.1
SpockProgressReachedLSN reports whether the local node's apply progress from the named peer has reached targetLSN. Uses remote_lsn (the LSN of the last applied commit) on Spock < 6, or remote_commit_lsn on Spock >= 6 — spock.progress became a view over apply_group_progress() in Spock 6 and the column was renamed. Neither uses received_lsn, which can advance on keepalive messages before any commits have been applied.
func SubscriptionNeedsCreate ¶
func SubscriptionNeedsEnable ¶
func SyncEvent ¶
SyncEvent sends a sync event marker from the provider. transactional ties the marker to the surrounding transaction so it's ordered with any preceding DML in that transaction, per Spock's spock.sync_event reference behavior. spock.sync_event(boolean) was only introduced in Spock 5.0.7 (see MinSpockVersionForSyncEventArgs) — pass transactional=false against older clusters to fall back to the original zero-arg call.
func UserRoleNeedsCreate ¶ added in v0.8.0
UserRoleNeedsCreate returns a query that evaluates to true when the named role does not yet exist in pg_catalog.pg_roles.
func WaitForSyncEvent ¶
WaitForSyncEvent waits for a peer to apply the sync event at lsn. waitIfDisabled=true tells Spock to tolerate the subscription not existing yet or being temporarily disabled — both expected during add-node — rather than raising immediately, matching Spock's reference wait_for_sync_event usage. The 5-arg form (with wait_if_disabled) was only introduced in Spock 5.0.7 (see MinSpockVersionForSyncEventArgs) — older clusters only accept the 4-arg form.
type ReplicationSet ¶
type ReplicationSetTable ¶
type Statement ¶
func AddReplicationSetTable ¶
func AddReplicationSetTable( setName string, relOID uint32, columns []string, rowFilter string, sync bool, includePartitions bool, ) Statement
https://docs.pgedge.com/spock_ext/spock_functions/functions/spock_repset_add_table
func AdvanceReplicationOrigin ¶ added in v0.8.1
func AlterOwner ¶ added in v0.8.0
func CreateReplicationSet ¶
func CreateReplicationSet(r ReplicationSet) Statement
https://docs.pgedge.com/spock_ext/spock_functions/functions/spock_repset_create
func DropAllSubscriptions ¶
func DropAllSubscriptions() Statement
func DropSubscription ¶
func EnableRepairMode ¶
func EnableRepairMode() Statement
type Statements ¶
type Statements []IStatement
func CreateBuiltInRoles ¶
func CreateBuiltInRoles(opts BuiltinRoleOptions) (Statements, error)
func CreatePgEdgeSuperuserRole ¶
func CreatePgEdgeSuperuserRole(opts BuiltinRoleOptions) (Statements, error)
func CreateUserRole ¶
func CreateUserRole(opts UserRoleOptions) (Statements, error)
func DropSpockAndCleanupSlots ¶
func DropSpockAndCleanupSlots(dbName string) Statements
func GrantBuiltinRolePrivileges ¶ added in v0.8.0
func GrantBuiltinRolePrivileges(opts BuiltinRolePrivilegeOptions) Statements
func InitializeSpockNode ¶ added in v0.8.0
func InitializeSpockNode(nodeName string, nodeDSN *DSN) Statements
func RestoreReplicationSets ¶
func RestoreReplicationSets(sets []ReplicationSet, tabs []ReplicationSetTable) Statements
func SetSafeIdentifiers ¶ added in v0.8.0
func SetSafeIdentifiers() Statements
SetSafeIdentifiers ensures that the escaped identifiers produced by the QuoteIdentifier method are safe for the duration of the session.
type SubscriptionStatus ¶
type SubscriptionStatus struct {
SubscriptionName string `json:"subscription_name"`
Status string `json:"status"`
ProviderNode string `json:"provider_node"`
ProviderDSN string `json:"provider_dsn"`
SlotName string `json:"slot_name"`
ReplicationSets []string `json:"replication_sets"`
ForwardOrigins []string `json:"forward_origins"`
}