Release 0.34.0 · Cluster

SoliDB 0.34.0: a cluster that actually spans machines

Before 0.34.0, a SoliDB cluster spread over several machines came up, logged nothing unusual, and replicated nothing. Five faults stacked on top of each other, and each one hid the next. They were found by deploying two real nodes rather than by a test. This post covers what changed, the two flags that matter (--advertise and --keyfile), and how to start three nodes on three machines.

The error messages and the status output below were produced by starting real nodes on one host.

What was broken

Each of these, alone, was enough to stop a multi-machine cluster from working. None of them produced an error.

Before 0.34.0In 0.34.0
The advertised replication address was hardcoded to 127.0.0.1, with no flag to change it. Every node told every peer to reach it at an address each peer reads as itself.
New --advertise flag, defaulting to --host. An unroutable value is refused at startup.
A node with peers skipped creating _admins. In a cluster whose members list each other, that was every node, the first one included: the cluster answered 401 to everything.
The skip is now a WARN naming both ways out.
A joining node could not get data that predated it. The full-sync command channel was discarded, the framing was off by a byte, and the batch could never deserialise. Full sync had never once completed.
A joining node asks its seed for a full sync.
--advertise 10.0.0.1:6746 became 10.0.0.1:6746:6746. The joiner never learned its peers; the error was discarded.
A port is attached only when there is not one already.
With no secret, the cluster bus sent unsigned messages and accepted unverified ones.
A shared secret is required.

The address mismatch mattered most, and the address mangling made it worse. With the second fault, the seed counted two healthy nodes and the member counted one: solidb_cluster_healthy_nodes read 2 on the seed and 1 on the member, and neither side logged anything.

A secret is now required

Cluster messages drive membership and shard rebalancing. Before 0.34.0 they were signed only when a secret was configured. Without one, a node sent them unsigned and accepted anything. The asymmetry was the dangerous part. A node with a secret rejects unsigned messages. A node without one accepts both. So one misconfigured member was an open door into the replicated state of the whole cluster, and nothing on the other nodes showed it.

Now both directions fail closed. Every cluster message is wrapped in an envelope carrying a timestamp, a random nonce and an HMAC-SHA256 signature over both and the payload. A receiver rejects anything unsigned, badly signed, or more than five minutes away from its own clock, so keep the nodes' clocks in sync. The same file is used for the challenge-response when a replication connection opens.

A node started with --peer and no keyfile does not start:

terminal
solidb --host 127.0.0.1 --port 6933 --peer 127.0.0.1:6924
Error: Cluster peers are configured but no keyfile is available. Create a shared secret of 32 cryptographically random bytes, hex-encoded, and pass it with --keyfile (the same file on every node). Unix: `openssl rand -hex 32 > solidb.key`. […] Refusing to start an unauthenticated cluster.

The first node of a cluster has no --peer, so it will start without a keyfile, with a warning that its replication and cluster ports accept unauthenticated connections. But it then refuses every cluster message it is sent, so it cannot take part in a cluster until it has one. Give every node the keyfile, the seed included. The secret signs traffic; it does not encrypt it, so keep the cluster port on a private network.

The admin account and the data that came first

A node started with --peer assumes it is joining an existing cluster and waits to receive the admin account by sync instead of creating one. The database cannot tell "I am joining" from "we are all starting at once", because both look like a non-empty peer list. When every node lists the others, nobody creates an admin. In 0.34.0 that is still the rule, but it is logged as a WARN that says what to do: start the first node with no --peer, or set SOLIDB_ADMIN_PASSWORD.

The first node, started alone, creates the admin user. It uses the password in SOLIDB_ADMIN_PASSWORD, or generates a random one and writes it to <data-dir>/.admin_password. A node that sets SOLIDB_ADMIN_PASSWORD creates its own admin even when it has peers. The walkthrough below sets it on every node, so no node's login depends on a sync having finished.

Replication carries writes forward from the moment a node joins. It knows nothing about earlier data. In 0.34.0 a joining node asks the seed it joined through for a full sync as soon as the join succeeds. If it cannot ask, it logs an error saying it will only have the writes made from then on.

Three nodes on three machines

Three machines on a private network, 10.0.0.1 to 10.0.0.3, each running SoliDB on the default port 6745, with HTTP and replication multiplexed on it. db1 binds its private address. db2 binds every interface, so it has to say what to advertise. db3 binds its private address, which is also what it advertises.

Three nodes: the address each listens on, the address each advertises, and how the joiners reach the seed Private network 10.0.0.0/24 · one port per node, HTTP API and replication multiplexed db1 · seeddb2db3 listens (--host)listens (--host)listens (--host) advertisesadvertises (--advertise)advertises 10.0.0.1:67450.0.0.0:674510.0.0.3:6745 10.0.0.1:674510.0.0.2:674510.0.0.3:6745 :6745 HTTP + repl:6745 HTTP + repl:6745 HTTP + repl db2, db3: --peer 10.0.0.1:6745 Before 0.34.0 all three advertised 127.0.0.1, which each peer reads as itself
Only the amber line travels to peers. db2 listens on every interface, so it has to name the address to advertise.

Create the secret once and copy the same file to every machine:

terminal
# once, then copy solidb.key to all three machines
openssl rand -hex 32 > solidb.key

Start the seed first, with no --peer:

terminal
# on 10.0.0.1
SOLIDB_ADMIN_PASSWORD='…' solidb --host 10.0.0.1 --port 6745 --node-id db1 \
  --keyfile /etc/solidb/solidb.key --data-dir /var/lib/solidb

Then the two others, pointing at the seed:

terminal
# on 10.0.0.2 — binds every interface, so --advertise is required
SOLIDB_ADMIN_PASSWORD='…' solidb --host 0.0.0.0 --advertise 10.0.0.2 --port 6745 --node-id db2 \
  --peer 10.0.0.1:6745 --keyfile /etc/solidb/solidb.key --data-dir /var/lib/solidb

# on 10.0.0.3 — binds its private address, which is also what it advertises
SOLIDB_ADMIN_PASSWORD='…' solidb --host 10.0.0.3 --port 6745 --node-id db3 \
  --peer 10.0.0.1:6745 --keyfile /etc/solidb/solidb.key --data-dir /var/lib/solidb

--node-id is optional; a UUID is generated without it, but a readable name makes the status output easier to follow. --peer can be repeated. A joining node tries its seeds in order and stops at the first that answers. The seed's answer lists the other members, so db3 learns about db2 without being told.

Checking that they see each other

GET /_api/cluster/status returns a node's id and its peers, as that node sees them. Ask more than one node: the bug fixed in 0.34.0 was exactly a seed that saw its member while the member saw nobody. Here are two nodes started on one host with the flags above (--host 127.0.0.1, ports 6923 and 6933, the second with --peer 127.0.0.1:6923). The stats and data_dir fields are omitted:

terminal
curl -s -u admin:… localhost:6923/_api/cluster/status
{
  "node_id": "n1",
  "status": "cluster",
  "replication_port": 6923,
  "current_sequence": 5,
  "log_entries": 5,
  "peers": [
    {"address": "127.0.0.1:6933", "is_connected": true, "last_seen_secs_ago": 1, "replication_lag": 5, "stats": null}
  ]
}
terminal
curl -s -u admin:… localhost:6933/_api/cluster/status
{
  "node_id": "n2",
  "status": "cluster",
  "replication_port": 6933,
  "current_sequence": 5,
  "log_entries": 5,
  "peers": [
    {"address": "127.0.0.1:6923", "is_connected": true, "last_seen_secs_ago": 1, "replication_lag": 5, "stats": null}
  ]
}

The joining node's log shows the join and the full sync, one after the other:

n2 log
INFO solidb: Sent join request to 127.0.0.1:6923
INFO solidb: Requesting a full sync from 127.0.0.1:6923
INFO solidb::cluster::manager: Successfully joined cluster. Received 2 peers.
INFO solidb::sync::worker: Starting full sync: 1 databases, 7 documents
INFO solidb::sync::worker: Full sync complete, final sequence: 5

After that, a collection created and a document inserted on the first node could be read from the second:

query on n2
FOR n IN notes RETURN {key: n._key, from: n.from}
[{"from": "n1", "key": "hello"}]

GET /_api/cluster/info returns the node's own configuration, including the peers it was started with. The endpoints are listed in the cluster API reference.

Upgrading a running cluster

Two changes in 0.34.0 can stop a node that used to start. That is why this is a minor release and not a patch.

  • No keyfile, no cluster. A node with --peer and no --keyfile refuses to start. Generate one secret, put the same file on every node, and add --keyfile to each node's command line, the seed included.
  • No unroutable advertise address. A node with peers on other machines that would advertise loopback, localhost, 0.0.0.0 or :: refuses to start. A unit file with --host 0.0.0.0 needs an --advertise next to it. Leaving out --host is not a way around this: the advertised address then falls back to 127.0.0.1.

Before 0.34.0, such a cluster started and did not work: the old behaviour was a cluster that could not replicate and gave no sign of it. More on sharding and failover is in the cluster and sharding docs.