Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 60 additions & 0 deletions cluster/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
# go-cowsql cluster

The `cluster` package exposes an opinionated pattern for setting up a resilient COWSQL cluster
following the design in use by [Incus](https://github.com/lxc/incus) in production clusters since 2018.

## Implementation and Usage

The main entry point to the `cluster` package is `cluster.Gateway`.
```go
// Initialize the gateway.
gateway, _ := gateway.NewGateway(ctx, node, clusterCert, serverCertFunc, state)

// Setup your API here, including gateway.HandlerFuncs.
go myServer.Serve(myListener)

// Open the COWSQL database.
driver, _ := gateway.Driver()
sql.Register("my-cowsql", driver)
db, _ := sql.Open("my-cowsql", driver)
gateway.SetClusterDB(db)
```

`NewGateway` depends on several interfaces to be set up by the consuming project:

* `State`:
Related to managing certificates on the filesystem, as well as listener setup and authentication.
* `Node`:
Used for per-member local store of cluster members.
* `Cluster`:
Content of the global COWSQL database set up by the `cluster` package.
These methods are used to provide implementations for cluster membership, role changes, recovery, and heartbeats.
* `CertInfo`:
TLS certificate management for shared cluster certificates and per-member server certificates for TLS-based intra-cluster communication.

## Example

An example package is provided, showing basic cluster setup with a recurring heartbeat and role rebalancing.

Note that while TLS is set up in this example, it is not verified for ease of testing.

```bash
# Start 3 daemons:
cluster daemon --address 10.0.0.101:8001 --name c1 --dir /tmp/cowsql-c1
cluster daemon --address 10.0.0.101:8002 --name c2 --dir /tmp/cowsql-c2
cluster daemon --address 10.0.0.101:8003 --name c3 --dir /tmp/cowsql-c3

# Bootstrap a cluster with 1 member:
cluster client bootstrap --address 10.0.0.101:8001

# Join an existing cluster from an uninitialized daemon:
cluster client bootstrap --address 10.0.0.101:8002 --target 10.0.0.101:8001
cluster client bootstrap --address 10.0.0.101:8003 --target 10.0.0.101:8001

# List cluster members:
cluster client list --address 10.0.0.101:8002

Name: "c1" Address: "10.0.0.101:8001" Role: "voter" Offline: false
Name: "c2" Address: "10.0.0.101:8002" Role: "voter" Offline: false
Name: "c3" Address: "10.0.0.101:8003" Role: "voter" Offline: false
```
14 changes: 14 additions & 0 deletions cluster/api/version.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package api

import (
"net/http"
"strconv"
)

// COWSQLVersion is the current cowsql protocol version.
const COWSQLVersion = 1

// SetCOWSQLVersionHeader sets the cowsql version header.
func SetCOWSQLVersionHeader(request *http.Request) {
request.Header.Set("X-Dqlite-Version", strconv.Itoa(COWSQLVersion))
}
34 changes: 34 additions & 0 deletions cluster/connect.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package cluster

import (
"context"
"net"
"time"

"github.com/cowsql/go-cowsql/cluster/tls"
)

// HasConnectivity probes the member with the given address for connectivity.
func HasConnectivity(networkCert tls.CertInfo, serverCert tls.CertInfo, address string, useTLS12 bool) bool {
// Get the transport.
transport, cleanup, err := tls.Transport(networkCert, serverCert, useTLS12)
if err != nil {
return false
}

defer cleanup()

ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()

var conn net.Conn

conn, err = transport.DialTLSContext(ctx, "tcp", address)
if err == nil {
_ = conn.Close()

return true
}

return false
}
95 changes: 95 additions & 0 deletions cluster/db/cluster.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
package db

import (
"context"
"crypto/x509"
"database/sql"
"time"

"github.com/cowsql/go-cowsql/cluster/db/transaction"
)

// Cluster represents the implementation for fetching cluster members from the global COWSQL database.
type Cluster interface {
transaction.Transactor

DB() *sql.DB

// GetNodes returns all cluster member node info.
GetNodes(ctx context.Context) ([]NodeInfo, error)

// GetNodesCount returns the count of cluster member node information.
GetNodesCount(ctx context.Context) (int, error)

// RemoveNode removes the node with the given name.
RemoveNode(ctx context.Context, name string) error

// GetNodesFailureDomains returns a map associating each node address with its
// failure domain code.
GetNodesFailureDomains(ctx context.Context) (map[string]uint64, error)

// GetNodeOfflineThreshold returns the amount of time that needs to elapse after
// which a series of unsuccessful heartbeat will make the node be considered
// offline.
GetNodeOfflineThreshold(ctx context.Context) (time.Duration, error)

// NodeIsOutdated returns true if there's some cluster node having an API or
// schema version greater than the node this method is invoked on.
NodeIsOutdated(ctx context.Context) (bool, error)

// SetNodeHeartbeat updates the heartbeat column of the node with the given address.
SetNodeHeartbeat(ctx context.Context, address string, heartbeatTime time.Time) error

// GetNodeByAddress returns the node with the given network address and pending state.
GetNodeByAddress(ctx context.Context, address string, pending bool) (NodeInfo, error)

// GetNodeByName returns the node with the given name and pending state.
GetNodeByName(ctx context.Context, name string, pending bool) (NodeInfo, error)

// BootstrapNode sets the name and address of the first cluster member, with id: 1.
BootstrapNode(ctx context.Context, serverName string, clusterAddress string) error

// CreateNode adds a node to the current list of members that are part of the
// cluster. The node's architecture should be the architecture of the machine the
// method is being run on. It returns the ID of the newly inserted row.
CreateNode(ctx context.Context, name, address string, arch int) (int64, error)

// SetNodePendingFlag toggles the pending flag for the node. A node is pending when
// it's been accepted in the cluster, but has not yet actually joined it.
SetNodePendingFlag(ctx context.Context, nodeID int64, flag bool) error

// SetNodeCertificateByName adds the serverCert to the DB trusted certificates store using the serverName.
SetNodeCertificateByName(ctx context.Context, serverName string, serverCert *x509.Certificate) error
}

// ClusterWarningHandler represents the implementation for managing dedicated warnings.
type ClusterWarningHandler interface {
// ResolveOfflineMemberWarning resolves offline member warnings for the given nodeID and serverName.
ResolveOfflineMemberWarning(ctx context.Context, nodeID int64, serverName string) error

// EmitOfflineMemberWarning emits an offline member warning for the given nodeID and serverName.
EmitOfflineMemberWarning(ctx context.Context, nodeID int64, serverName string, msg string) error

// ResolveTimeSkewWarning resolves time skew warnings for the given serverName.
ResolveTimeSkewWarning(ctx context.Context, serverName string) error

// EmitTimeSkewWarning emits a time skew warning for the given serverName.
EmitTimeSkewWarning(ctx context.Context, serverName string, msg string) error
}

// ClusterExternal represents the implementation for managing external cluster resources.
type ClusterExternal[T any] interface {
// GetLocalResources Fetches external data of the type from the cluster DB.
GetLocalResources(ctx context.Context) (T, error)

// UpdateClusterResources updates external data of the given type from the cluster DB.
UpdateClusterResources(ctx context.Context, node NodeInfo, t T) error

// ClearNode removes all data referencing the nodeID.
ClearNode(ctx context.Context, nodeID int64) error

// NodeIsEmpty returns an empty string if the node with the given ID has no
// referencing data associated with it. Otherwise, it returns a message
// say what's left.
NodeIsEmpty(ctx context.Context, nodeID int64) (string, error)
}
12 changes: 12 additions & 0 deletions cluster/db/error.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
package db

import "errors"

var (
// ErrAlreadyDefined happens when the given entry already exists,
// for example a container.
ErrAlreadyDefined = errors.New("The record already exists")

// ErrNoClusterMember is used to indicate no cluster member has been found for a resource.
ErrNoClusterMember = errors.New("No cluster member found")
)
67 changes: 67 additions & 0 deletions cluster/db/info.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
package db

import "time"

// NodeInfo holds information about a single member in a cluster.
type NodeInfo struct {
ID int64 // Stable node identifier
Name string // User-assigned name of the node
Address string // Network address of the node
Description string // Node description (optional)
Schema int // Schema version of the daemon running the member
APIExtensions int // Number of API extensions of the daemon running the member
Heartbeat time.Time // Timestamp of the last heartbeat
Roles []string // List of cluster roles
Architecture int // Node architecture
State int // Node state
Config map[string]string // Configuration for the node
Groups []string // Cluster groups

// heartbeatRefTime is the time at which the node row was read from the
// database. It is used as the reference time for IsOffline so that a
// long-running transaction (during which the heartbeat snapshot cannot be
// refreshed) doesn't incorrectly classify members as offline just because
// the transaction has been open longer than the offline threshold.
heartbeatRefTime time.Time
}

// IsOffline returns true if the last successful heartbeat time of the node is
// older than the given threshold.
//
// The check uses the time at which this NodeInfo was read from the database as
// the reference point (if known), rather than time.Now(). This avoids false
// positives when the caller is operating inside a long-running transaction
// where the heartbeat snapshot can't be refreshed.
func (n NodeInfo) IsOffline(threshold time.Duration) bool {
return nodeIsOffline(threshold, n.Heartbeat, n.heartbeatRefTime)
}

// SetHeartbeatRefTime sets the reference time used for IsOffline.
func (n *NodeInfo) SetHeartbeatRefTime(t time.Time) {
n.heartbeatRefTime = t
}

// nodeIsOffline reports whether heartbeat is older than threshold relative to
// refTime. If refTime is the zero value, time.Now() is used as the reference.
func nodeIsOffline(threshold time.Duration, heartbeat time.Time, refTime time.Time) bool {
if refTime.IsZero() {
refTime = time.Now()
}

offlineTime := refTime.UTC().Add(-threshold)

return heartbeat.Before(offlineTime) || heartbeat.Equal(offlineTime)
}

// HeartbeatMember contains specific cluster node info.
type HeartbeatMember struct {
ID int64 `json:"ID"` //nolint:tagliatelle // ID field value in nodes table.
Address string `json:"Address"` //nolint:tagliatelle // Host and Port of node.
Name string `json:"Name"` //nolint:tagliatelle // Name of cluster member.
RaftID uint64 `json:"RaftID"` //nolint:tagliatelle // ID field value in raft_nodes table, zero if non-raft node.
RaftRole int `json:"RaftRole"` //nolint:tagliatelle // Node role in the raft cluster, from the raft_nodes table
LastHeartbeat time.Time `json:"LastHeartbeat"` //nolint:tagliatelle // Last time we received a successful response from node.
Online bool `json:"Online"` //nolint:tagliatelle // Calculated from offline threshold and LastHeatbeat time.
Roles []string `json:"Roles"` //nolint:tagliatelle // Supplementary non-database roles the member has.
Updated bool `json:"Updated"` //nolint:tagliatelle // Has node been updated during this heartbeat run.
}
47 changes: 47 additions & 0 deletions cluster/db/node.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
package db

import (
"context"
"database/sql"

"github.com/cowsql/go-cowsql/cluster/db/transaction"
)

// Node represents the implementation for fetching cached cluster members from the local database.
type Node interface {
transaction.Transactor

// DB returns the underlying DB.
DB() *sql.DB

// LogPath returns the path to the database logs.
LogPath() string

// GlobalDatabaseDir returns the containing directory of the database file.
GlobalDatabaseDir() string

// GetRaftNodes returns all cached cluster member raft node info.
GetRaftNodes(ctx context.Context) ([]RaftNode, error)

// GetRaftNode returns the raft node corresponding to the given node ID.
GetRaftNode(ctx context.Context, id int64) (*RaftNode, bool, error)

// CreateFirstRaftNode adds the first node of the cluster. It ensures that the
// database ID is 1, to match the server ID of the first raft log entry.
//
// This method is supposed to be called when there are no rows in raft_nodes,
// and it will replace whatever existing row has ID 1.
CreateFirstRaftNode(ctx context.Context, address string, name string) error

// ReplaceRaftNodes replaces the current list of raft nodes.
ReplaceRaftNodes(ctx context.Context, nodes []RaftNode) error

// DetermineRaftNode figures out what raft node ID and address we have, if any.
DetermineRaftNode(ctx context.Context) (*RaftNode, error)

// GetClusterAddress fetches the current cluster address used for intra-cluster communication.
GetClusterAddress(ctx context.Context) (string, error)

// SetClusterAddress sets the cluster address used for intra-cluster communication.
SetClusterAddress(ctx context.Context, address string) error
}
31 changes: 31 additions & 0 deletions cluster/db/raft.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
package db

import (
"github.com/cowsql/go-cowsql/client"
)

// RaftNode holds information about a single node in the cowsql raft cluster.
//
// This is just a convenience alias for the equivalent data structure in the
// cowsql client package.
type RaftNode struct {
ID uint64 `json:"ID"` //nolint:tagliatelle
Address string `json:"Address"` //nolint:tagliatelle
Role client.NodeRole `json:"Role"` //nolint:tagliatelle
Name string `json:"Name"` //nolint:tagliatelle
}

// RaftRole captures the role of cowsql/raft node.
type RaftRole = client.NodeRole

// RaftNode roles.
const (
RaftVoter = client.Voter
RaftStandBy = client.StandBy
RaftSpare = client.Spare
)

// DefaultRaftNode represents a fully uninitialized raft node entry with ID: 1 and Address: 1, signifying an uninitialized system.
func DefaultRaftNode() *RaftNode {
return &RaftNode{ID: 1, Address: "1", Name: ""}
}
Loading
Loading