# How high availability works in Memgraph

> **Info**
>
> This guide builds on the concepts introduced in [how replication
> works](https://memgraph.com/docs/clustering/replication/how-replication-works). We recommend reading that
> page before continuing.

High availability (HA) in Memgraph ensures that **the cluster can always serve
queries**, even when nodes fail. It is built on two foundations:
1. **Replication** - multiple data instances hold copies of the data.
2. **Automatic failover** - the system transparently promotes a new Main
instance when needed.

A Memgraph HA cluster contains three types of instances:
- The MAIN instance on which the user can execute **read** and **write**
  queries.
- REPLICA instances that can only respond to **read** queries.
- COORDINATOR instances that **manage the cluster state**.

In HA documentation, MAIN and REPLICATE together are often referred to as **data
instances**, because either one may become the Main during the cluster lifetime.

Coordinator instances do not store graph data and are therefore much lighter.

## How Memgraph achieves High Availability

A typical highly available Memgraph cluster includes:
- 3 data instances (1 Main + 2 Replicas)
- 3 coordinators (1 Leader + 2 Followers)

![](https://memgraph.com/docs/pages/clustering/high-availability/typical_ha_cluster.png)

The constraint for number coordinators is only that it needs to be an **odd
number of them, greater than 1** (3, 5, 7, ...). Users can create more than 3
coordinators, but the replication factor (RF) of 3 is a de facto standard in
distributed databases.

The minimum valid data-instance setup is:
- 1 Main
- 1 Replica

If the Main fails, a Replica is **automatically promoted**.

![](https://memgraph.com/docs/pages/clustering/high-availability/minimal_ha_cluster.png)

For achieving high availability, Memgraph uses the **Raft consensus protocol for
cluster coordination**. Raft is easier to reason about than Paxos and widely
adopted across distributed systems. **As a design decision, Memgraph uses an
industry-proven library [NuRaft](https://github.com/eBay/NuRaft) for the
implementation of the Raft protocol.**

Raft provides:
- **Leader election** among coordinator instances
- A **replicated, durable cluster-management log**
- **Fault-tolerant decision-making via majority quorum**

> **Info**
>
> Raft is *not* a Byzantine fault-tolerant protocol. Using
> **an odd number of coordinators** ensures majority-based correctness.

The coordinator leader ensures:
- Exactly **one Main** exists
- Replicas are correctly registered
- Failover is triggered when needed

Coordinators themselves are redundant, making cluster orchestration highly
available.

### Data instance implementation

The data instance is your usual Memgraph standalone instance, with one key flag
added:
- `--management-port` - used to get the **health state** of the data instance
  from the leader coordinator

When a data instance runs in HA mode (i.e. `--management-port` is set), the
`--init-file` and `--init-data-file` flags are **not supported**. The instance
will fail to start if either flag is provided. This is because the cluster's
leader coordinator drives role transitions (MAIN/REPLICA) and replication
setup, so executing arbitrary initialization queries on startup could conflict
with the cluster state. To bootstrap users or data in an HA cluster, run the
relevant queries through the MAIN after the cluster is formed.

### Coordinator instance implementation

The coordinator is a small orchestration instance which is shipped in the same
Memgraph binary as the data instance itself. For the system to be aware that its
role is COORDINATOR, the user needs to specify four flags:
- `--coordinator-id` - serves as a **unique identifier** of the coordinator
- `--coordinator-port` - used for **synchronization and log replication**
  between coordinators
- `--coordinator-hostname` - used on followers to ping the leader on the correct
  IP address (or FQDN/DNS name)
- `--management-port` - used to get the **health state** of the respective
  coordinator instance from the leader coordinator

> **Note**
>
> The COORDINATOR instance is a **very restricted instance**, and it will
> not respond to any queries that are not related to management of the cluster.
> That means, you cannot run any data queries on the coordinator directly (we
> will talk more about routing data queries in the next sections). However,
> system information queries such as `SHOW CONFIG`, `SHOW VERSION`,
> `SHOW LICENSE INFO`, `SHOW BUILD INFO` and `SHOW STORAGE INFO` are supported on
> coordinators, as well as `SET DATABASE SETTING`, `RELOAD BOLT_SERVER TLS` and
> `RELOAD INTRA_CLUSTER TLS`.

Since coordinators do not store user data, the following restrictions apply:
- **Snapshots are automatically disabled** on coordinators, even if
  `--storage-snapshot-interval-sec` is set. The `storage.snapshot.interval`
  setting is not registered on coordinators, so attempting to read or modify it
  via `SHOW DATABASE SETTING` / `SET DATABASE SETTING` returns an unknown
  setting error.
- **The `--query-modules-directory` flag is ignored** on coordinators.
  Coordinators do not execute data queries, so query modules are never loaded
  and the embedded Python runtime is not initialized. The flag is still accepted
  (so packaged defaults do not need to be overridden) but has no effect.
- **The `--init-file` and `--init-data-file` flags are not supported** on
  coordinators (and likewise not supported on data instances in HA mode). The
  instance will fail to start if either flag is provided.

When deploying coordinators to servers, you can use the instance of almost any
size. Instances of 4GiB or 8GiB will suffice since coordinators' job mainly
involves network communication and storing Raft metadata. Coordinators and data
instances can in theory be deployed on same servers (pairwise) but from the
availability perspective, it is better to separate them physically. 

### RPC communication in the cluster

RPC (Remote Procedure Call) is a protocol for executing functions on a remote
system. RPC enables direct communication in distributed systems and is crucial
for replication and high availability tasks. 

RPC (Remote Procedure Calls) are used for:
- Replication
- Cluster orchestration
- Health monitoring
- Failover triggers

Below is a cleaned-up categorization.

#### Coordinator → Coordinator RPCs

| RPC                      | Purpose                                      | Description                                                                                                                |
| ------------------------ | -------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------- |
| `ShowInstancesRpc`       | Follower requests cluster state from leader. | Sent by a follower coordinator to the leader coordinator when a user executes `SHOW INSTANCES` through the follower.       |
| `ShowCoordSettingsRpc`   | Follower requests coordinator settings.      | Sent by a follower coordinator to the leader coordinator when a user executes `SHOW COORDINATOR SETTINGS` through the follower. |
| `YieldLeadershipRpc`     | Follower asks the leader to step down.       | Sent by a follower coordinator to the leader coordinator when a user executes `YIELD LEADERSHIP` through the follower.     |
| `AddCoordinatorRpc`      | Follower requests adding coordinator.        | Sent by a follower coordinator to the leader coordinator when a user executes `ADD COORDINATOR` through the follower.      |
| `RemoveCoordinatorRpc`   | Follower requests removing coordinator.      | Sent by a follower coordinator to the leader coordinator when a user executes `REMOVE COORDINATOR` through the follower.   |
| `RegisterInstanceRpc`    | Follower requests registering an instance.   | Sent by a follower coordinator to the leader coordinator when a user executes `REGISTER INSTANCE` through the follower.    |
| `UnregisterInstanceRpc`  | Follower requests unregistering an instance. | Sent by a follower coordinator to the leader coordinator when a user executes `UNREGISTER INSTANCE` through the follower.  |
| `SetInstanceToMainRpc`   | Follower requests updating MAIN instance.    | Sent by a follower coordinator to the leader coordinator when a user executes `SET INSTANCE TO MAIN` through the follower. |
| `DemoteInstanceRpc`      | Follower requests demoting an instance.      | Sent by a follower coordinator to the leader coordinator when a user executes `DEMOTE INSTANCE` through the follower.      |
| `UpdateConfigRpc`        | Follower requests updating config.           | Sent by a follower coordinator to the leader coordinator when a user executes `UPDATE CONFIG` through the follower.        |
| `ForceResetRpc`          | Follower requests resetting cluster state.   | Sent by a follower coordinator to the leader coordinator when a user executes `FORCE RESET` through the follower.          |
| `GetRoutingTableRpc`     | Follower requests a routing table.           | Sent by a follower coordinator to the leader coordinator when a user connects using `bolt+routing` or executes `SHOW ROUTING TABLE` through the follower. |
| `CoordReplicationLagRpc` | Follower requests replication lag info.      | Sent by a follower coordinator to the leader coordinator when a user executes `SHOW REPLICATION LAG` through the follower. |

#### Coordinator → Data Instance RPCs

All of the following messages were sent by the leader coordinator.

| RPC                          | Purpose                                             | Description                                                                                                                                                |
| ---------------------------- | --------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `PromoteToMainRpc`         | Promote a Replica to Main                           | Sent to a REPLICA in order to promote it to Main.                                                                                |
| `DemoteMainToReplicaRpc`   | Demote a Main after failover                        | Sent to the old MAIN in order to demote it to REPLICA.                                                                           |
| `RegisterReplicaOnMainRpc` | Instruct Main to accept replication from a Replica  | Sent to the MAIN to register a REPLICA on the MAIN.                                                                              |
| `UnregisterReplicaRpc`     | Remove Replica from Main                            | Sent to the MAIN to unregister a REPLICA from the MAIN.                                                                          |
| `EnableWritingOnMainRpc`   | Re-enable writes after Main restarts (deprecated)   | Kept for backward compatibility (ISSU). No longer sent by coordinators — writing is implicitly enabled on promotion.             |
| `GetDatabaseHistoriesRpc`  | Gather committed transaction counts during failover | Sent to all REPLICA instances in order to select a new MAIN during the failover process.                                         |
| `StateCheckRpc`            | Health check ping (liveness)                        | Sent to all data instances for a liveness check.                                                                                 |
| `SwapMainUUIDRpc`          | Ensure Replica tracks the correct Main              | Sent to REPLICA instances to set the UUID of the MAIN they should listen to.                                                     |

#### Main → Replica RPCs

All of the following messages were sent by the Main instance.

| RPC                      | Purpose                             | Description                                                                                                                                                                            |
| ------------------------ | ----------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `FrequentHeartbeatRpc` | Liveness check                      | Sent to a REPLICA for liveness checks.                                                                                                                                     |
| `HeartbeatRpc`         | Replication metadata sync           | Sent to a REPLICA for transmitting timestamp, epoch, transaction, and commit information.                                                                                  |
| `PrepareCommitRpc`     | Send deltas; phase 1 of STRICT_SYNC | Sent  REPLICA instances either as the first phase in `STRICT_SYNC` mode or as the only phase in other replication modes. It sends the delta stream during a write query. |
| `FinalizeCommitRpc`    | Commit confirmation in STRICT_SYNC  | Sent REPLICA instances in `STRICT_SYNC` mode when all Replicas have acknowledged they are ready to commit.                                                              |
| `SnapshotRpc`          | Snapshot recovery                   | Sent a REPLICA for snapshot recovery so the REPLICA can catch up with the MAIN.                                                                                         |
| `WalFilesRpc`          | WAL segment recovery                | Sent a REPLICA for WAL recovery so the REPLICA can catch up with the MAIN.                                                                                              |
| `CurrentWalRpc`        | Latest WAL file recovery            | Sent a REPLICA for current/latest/unfinished WAL file recovery.                                                                                                         |
| `SystemRecoveryRpc`    | Replicate system metadata           | Sent a REPLICA for replicating system-level metadata (auth, multi-tenancy, etc.) and other non-graph information.                                                       |

#### RPC timeouts

##### Default RPC timeouts

For the majority of RPC messages, Memgraph uses a default timeout of 10s. This
is to ensure that when sending a RPC request, the client will not block
indefinitely before receiving a response if the communication between the client
and the server is broken. The list of RPC messages for which the timeout is used
is the following: 

##### RPC timeout during replication of data

For RPC messages which are sending the variable number of storage deltas —
`PrepareCommitRpc`, `CurrentWalRpc`, and `WalFilesRpc` — it is not practical to
set a strict execution timeout, as the processing time on the replica side is
directly proportional to the number of deltas being transferred. To handle this,
the replica sends periodic progress updates to the main instance after
processing every 100,000 deltas. Since processing 100,000 deltas is expected to
take a relatively consistent amount of time, we can enforce a timeout based on
this interval. The default timeout for these RPC messages is 30 seconds, though
in practice, processing 100,000 deltas typically takes less than 3 seconds.

##### RPC timeout during replica recovery

`SnapshotRpc` is also a replication-related RPC message, but its execution time
is tracked a bit differently from RPC messages shipping deltas. The replica
sends an update to the main instance after completing 1,000,000 units of work.
The work units are assigned as follows:

- Processing nodes, edges, or indexed entities (label index, label-property
  index, edge type index, edge type property index) = 1 unit
- Processing a node inside a point or text index = 10 units
- Processing a node inside a vector index (most computationally expensive) =
  1,000 units

With this unit-based tracking system, the REPLICA is expected to report progress
approximately every 2–3 seconds. Given this, a timeout of 60 seconds is set to
avoid unnecessary network instability while ensuring responsiveness. On every
report of the progress by the REPLICA, the timeout of 60 seconds will be
restarted, as the state of recovery is considered to be stable.

Except for timeouts on read and write operations, Memgraph also has a timeout of
5s for sockets when establishing a connection. Such a timeout helps in having a
low p99 latencies when using the RPC stack, which manifests for users as smooth
and predictable network communication between instances.

In the table below, we have the full outline of the RPC messages that are sent
in the cluster to ensure high availability, with timeouts. 

| RPC message request      | source      | target         | timeout          |
|--------------------------|-------------|----------------| -----------------|
| `ShowInstancesReq`         | Coordinator | Coordinator    | 10s              |
| `DemoteMainToReplicaReq`   | Coordinator | Data instance  | 10s              |
| `PromoteToMainReq`         | Coordinator | Data instance  | 10s              |
| `RegisterReplicaOnMainReq` | Coordinator | Data instance  | 10s              |
| `UnregisterReplicaReq`     | Coordinator | Data instance  | 10s              |
| `ReplicationLagReq`        | Coordinator | Data instance  | 5s               |
| `GetDatabaseHistoriesReq`  | Coordinator | Data instance  | 10s              |
| `StateCheckReq`            | Coordinator | Data instance  | 5s               |
| `SwapMainUUIDReq`          | Coordinator | Data instance  | 10s              |
| `UpdateDataInstanceConfigReq` | Coordinator | Data instance | 10s            |
| `FrequentHeartbeatReq`     | Main        | Replica        | 5s               |
| `HeartbeatReq`             | Main        | Replica        | 10s              |
| `SystemRecoveryReq`        | Main        | Replica        | 30s              |
| `PrepareCommitRpc`         | Main        | Replica        | proportional     |
| `FinalizeCommitReq`        | Main        | Replica        | 10s              |
| `SnapshotRpc`              | Main        | Replica        | proportional     |
| `WalFilesRpc`              | Main        | Replica        | proportional     |
| `CurrentWalRpc`            | Main        | Replica        | proportional     |

> **Info**
>
> `EnableWritingOnMainReq` was **removed in Memgraph 3.13**. It was never sent —
> writing on a newly promoted MAIN is enabled through the `writing_enabled` flag
> carried inside `PromoteToMainRpc`. Its Prometheus counters were removed along
> with it.

##### Follower-to-leader forwarding timeouts

Cluster management queries can be run on any coordinator; a follower forwards
them to the leader over RPC. From Memgraph 3.13, each of these RPCs has an
explicit timeout. They run on the caller's **Bolt session thread**, so without
one, a session would block forever against a leader that is reachable but stuck.

The four instance operations must outlast the work they trigger on the leader,
or a follower would report failure for an operation the leader has already
committed to Raft. Their budgets are the sum of that work plus headroom, where a
Raft commit is capped at 3 seconds and each RPC from the leader to a data
instance is capped by its entry in the table above.

| RPC message request      | source      | target      | timeout | Budget breakdown |
|--------------------------|-------------|-------------|---------|------------------|
| `RegisterInstanceReq`    | Coordinator | Coordinator | 30s     | Raft commit + demote the new replica + register it |
| `UnregisterInstanceReq`  | Coordinator | Coordinator | 20s     | Raft commit + one RPC to the current MAIN |
| `DemoteInstanceReq`      | Coordinator | Coordinator | 20s     | Raft commit + one RPC to the current MAIN |
| `SetInstanceToMainReq`   | Coordinator | Coordinator | 60s     | Raft commit + one `SwapMainUUID` per other instance + promote the new MAIN |
| `AddCoordinatorReq`      | Coordinator | Coordinator | 10s     | Raft commit only |
| `RemoveCoordinatorReq`   | Coordinator | Coordinator | 10s     | Raft commit only |
| `UpdateConfigReq`        | Coordinator | Coordinator | 10s     | Raft commit only |
| `ForceResetReq`          | Coordinator | Coordinator | 60s     | Unbounded leader-side work — see the note below |
| `SetCoordinatorSettingReq` | Coordinator | Coordinator | 10s   | Raft commit only |
| `GetRoutingTableReq`     | Coordinator | Coordinator | 10s     | Read served by the leader |
| `CoordReplLagReq`        | Coordinator | Coordinator | 10s     | Read served by the leader |
| `CreateRoleReq`          | Coordinator | Coordinator | 10s     | Raft commit only |
| `DropRoleReq`            | Coordinator | Coordinator | 10s     | Raft commit only |
| `GrantPrivilegeReq`      | Coordinator | Coordinator | 10s     | Raft commit only |
| `RevokePrivilegeReq`     | Coordinator | Coordinator | 10s     | Raft commit only |
| `GetRolesReq`            | Coordinator | Coordinator | 10s     | Read served by the leader |
| `GetRolePrivilegesReq`   | Coordinator | Coordinator | 10s     | Read served by the leader |

> **Warning**
>
> `SetInstanceToMainReq` sends one `SwapMainUUID` per other instance, so its cost
> grows with the number of data instances. The 60s budget comfortably covers five
> instances. Beyond that, a follower can time out before the leader answers — no
> fixed value bounds it. Run `SET INSTANCE ... TO MAIN` directly on the leader in
> very large clusters.
>
> `ForceResetReq` triggers a reconciliation that retries under a 1s–60s backoff
> for as long as the coordinator stays leader, so the leader-side work has no
> upper bound at all. The 60s budget only keeps a genuinely wedged leader from
> blocking the session — hitting it does **not** mean the reset failed, and
> [`FORCE RESET CLUSTER
> STATE`](https://memgraph.com/docs/clustering/high-availability/ha-commands-reference#force-reset-cluster-state)
> is safe to re-run.

The role and privilege RPCs are used by [coordinator
authentication](https://memgraph.com/docs/clustering/high-availability/coordinator-authentication):
`GetRolesReq` in particular is sent on **every query of an SSO session**, because
coordinator privileges are re-derived from the leader's committed role set
rather than cached at login.

##### System transaction timeouts

MAIN-to-REPLICA system-delta RPCs are sent while committing a system
transaction, so they must not block indefinitely either. From Memgraph 3.13 each
carries a **10 second** timeout; a timeout marks the REPLICA as `BEHIND` and
defers to system recovery.

| RPC message request | source | target  | timeout |
|---------------------|--------|---------|---------|
| `UpdateAuthDataReq`     | Main | Replica | 10s |
| `DropAuthDataReq`       | Main | Replica | 10s |
| `FinalizeSystemTxReq`   | Main | Replica | 10s |
| `CreateDatabaseReq`     | Main | Replica | 10s |
| `DropDatabaseReq`       | Main | Replica | 10s |
| `SuspendDatabaseReq`    | Main | Replica | 10s |
| `ResumeDatabaseReq`     | Main | Replica | 10s |
| `RenameDatabaseReq`     | Main | Replica | 10s |
| `TenantProfileReq`      | Main | Replica | 10s |
| `SetParameterReq`       | Main | Replica | 10s |
| `UnsetParameterReq`     | Main | Replica | 10s |
| `DeleteAllParametersReq`| Main | Replica | 10s |

## Intra-cluster TLS

By default, the communication between instances in a high-availability cluster
is unencrypted. To secure it, Memgraph supports **intra-cluster TLS**, which
encrypts all internal cluster traffic using mutual TLS (mTLS). When enabled,
TLS protects the communication on:

- the **management server** (health checks between the leader coordinator and
  the instances),
- the **replication server** (data replication between MAIN and REPLICA
  instances), and
- the **coordinator server** (synchronization and log replication between
  coordinators).

> **Info**
>
> Intra-cluster TLS is independent of [Bolt SSL/TLS](https://memgraph.com/docs/database-management/ssl-encryption).
> Bolt encryption secures client-to-instance connections and is configured
> separately with the `--bolt-cert-file` and `--bolt-key-file` flags, while
> intra-cluster TLS secures the internal cluster communication described above.

### Enabling intra-cluster TLS

Intra-cluster TLS is enabled by setting the following three flags on every
instance (coordinators and data instances) in the cluster:

| Flag                  | Description                                                                                          |
| --------------------- | ---------------------------------------------------------------------------------------------------- |
| `--cluster-cert-file` | Certificate file used for intra-cluster TLS communication.                                           |
| `--cluster-key-file`  | Key file used for intra-cluster TLS communication.                                                   |
| `--cluster-ca-file`   | File storing the certificate of the Certificate Authority you trust for intra-cluster TLS communication. |

All three flags must be set together. If only some of them are provided, the
instance refuses to start to avoid running in a partially-configured TLS state.
When all three are empty, intra-cluster TLS is disabled and communication is
unencrypted. Because mTLS is used, every instance must present a certificate
that is trusted by the configured Certificate Authority, and all instances in
the cluster must be started with the TLS flags.

### Reloading intra-cluster TLS certificates at runtime

You can rotate the intra-cluster TLS certificates without restarting the
cluster by running the `RELOAD INTRA_CLUSTER TLS` Cypher query. Replace the
certificate and key files on disk at the paths configured with
`--cluster-cert-file`, `--cluster-key-file`, and `--cluster-ca-file`, then run:

```cypher
RELOAD INTRA_CLUSTER TLS;
```

`RELOAD INTRA_CLUSTER TLS` reloads both the client and server connections, so
new connections will start using the new certificates.

This is the intra-cluster counterpart of `RELOAD BOLT_SERVER TLS`, which
reloads the Bolt server certificates. Both queries require the `RELOAD_TLS`
[privilege](https://memgraph.com/docs/database-management/authentication-and-authorization/role-based-access-control).
For more details on reloading certificates, see the
[SSL encryption](https://memgraph.com/docs/database-management/ssl-encryption#reload-ssl-certificates-at-runtime)
page.

## Automatic failover

Automatic failover is driven by periodic health checks performed by the leader
coordinator on all data instances. An instance is considered **alive** if it
responds to these checks; if it does not, it is marked as down.

**When a Replica goes down**

If a **REPLICA** instance goes down, it will always rejoin the cluster as a
REPLICA once it becomes healthy again. No failover action is required.

**When the Main goes down**

If the **MAIN** instance is detected as down, the coordinator may initiate a
failover to select a new MAIN from the set of alive REPLICAs.

The failover process works as follows:
1. The coordinator selects one of the alive replicas as the new MAIN candidate.
2. It records this decision in the Raft log.
3. On the next heartbeat to the selected replica, the coordinator sends an RPC
request (`PromoteToMainReq`) instructing it to promote itself from REPLICA to
MAIN and providing information about the other replicas it should replicate to.
4. Once the promotion succeeds, the new MAIN begins replicating data to the
other instances and starts accepting write queries.

### Instance health checks

The coordinator performs health checks on each instance at a fixed interval,
configured with the `instance_health_check_frequency_sec` coordinator setting.
An instance is not considered down until it has failed to respond for the full
duration specified by the `instance_down_timeout_sec` coordinator setting. Both
settings can be changed at runtime using
[`SET COORDINATOR SETTING`](https://memgraph.com/docs/clustering/high-availability/ha-commands-reference#coordinator-runtime-settings).

**Example**

With the default settings:
- `instance_health_check_frequency_sec=1`
- `instance_down_timeout_sec=5`

…the coordinator will send a health check RPC (`StateCheckRpc`) every second. An
instance is marked as down only after **five consecutive missed responses** (5
seconds ÷ 1 second).

Usually, the health check reports instantly, and only in severe cases of network
failure, it will timeout in 30 seconds.

---

Depending on the data instance role, we have 2 scenarios:
1. **Replica instance fails to respond**

If a **REPLICA** fails to respond:
- The leader coordinator simply retries on the next scheduled health check.
- Once the instance becomes reachable again, it **always rejoins the cluster as
  a REPLICA**.

![](https://memgraph.com/docs/pages/clustering/high-availability/replica-rejoining-cluster.png)

2. **Main instance fails to respond**

If the **MAIN** instance fails to respond, two cases apply:
- **Down for less than** `instance_down_timeout_sec` The instance is still
considered alive and will rejoin as MAIN when it responds again.
- **Down for longer than** `instance_down_timeout_sec` The coordinator
initiates the failover procedure. What the old MAIN becomes afterward depends on
the outcome:
  - **Failover succeeds**: the old MAIN rejoins as a **REPLICA**.
  - **Failover fails**: the old MAIN retains authority and rejoins as **MAIN**
    when it becomes reachable.

![](https://memgraph.com/docs/pages/clustering/high-availability/main-rejoining-cluster.png)

> **Info**
>
> For guidance on configuring health checks, see the [Best
> practices.](https://memgraph.com/docs/clustering/high-availability/best-practices)

### Choosing the new MAIN 

During failover, the coordinator must select a new main instance from available
replicas, as some may be offline. The leader coordinator queries each live
replica to retrieve the committed transaction count for every database.

> **Note**
>
> **Note:** *For every database* in this terminologyrefers to environments using
> [multi-tenancy](https://memgraph.com/docs/database-management/multi-tenancy), where multiple isolated
> databases/graphs exist within a single instance. If multi-tenancy is not used,
> the coordinator retrieves the committed transaction count only for the default
> **memgraph** database.

The **selection algorithm** prioritizes data recency using a two-phase approach:

1. **Database majority rule**: The coordinator identifies which replica has the
highest committed transaction count for each database. The replica that leads in
the most databases becomes the preferred candidate.
2. **Total transaction tiebreaker**: If multiple replicas tie for leading the
most databases, the coordinator sums each replica's committed transactions
across all databases. The replica with the highest total becomes the new main.

This approach ensures the new main instance has the most up-to-date data across
the cluster while maintaining consistency guarantees.

### Old MAIN rejoining

When the old MAIN instance comes back online, it cannot resume as MAIN
automatically. The coordinator keeps track of which instance most recently
served as MAIN and ensures a controlled transition.

To demote the old MAIN to a REPLICA, the leader coordinator sends two RPC
requests **in order**:

1. **`DemoteMainToReplicaReq`** - instructs the old MAIN to demote itself to a
   REPLICA.
2. **Store current MAIN UUID** - updates the old instance (now a REPLICA) with
   the machine UUID of the current MAIN so it can correctly begin replication.

After these steps, the old MAIN fully reenters the cluster as a REPLICA.

### Ensuring replicas follow the correct MAIN

Each REPLICA stores the UUID of the MAIN instance it should follow. In certain
failure scenarios, such as a network partition, the MAIN may still be able to
communicate with a REPLICA even though the coordinator cannot reach the MAIN.
From the coordinator’s perspective, the MAIN appears down, while from the
REPLICA’s perspective, it appears alive.

This situation can lead to *split-brain* behavior, where failover procedure is
triggered and multiple MAINs briefly exist, with replicas potentially following
different leaders.

To prevent this, the coordinator manages MAIN UUIDs carefully:

- When a new MAIN is selected, the coordinator generates a **new, unique MAIN
  UUID** that no existing MAIN has.
- The new MAIN adopts this UUID when promoted, ensuring replicas can reliably
  distinguish which MAIN is authoritative.

If a REPLICA goes down and later rejoins, the MAIN may have changed during its
absence. To ensure the REPLICA follows the correct MAIN, the coordinator sends a
`SwapMainUUIDRpc` request, updating the REPLICA with the UUID of the current
MAIN.

This guarantees all REPLICAs consistently replicate from the correct leader and
prevents divergence across the cluster.

### Replication scenarios

**Force sync of data**

Earlier, we explained how Memgraph [selects the most up-to-date alive
instance](https://memgraph.com/docs/clustering/high-availability/how-high-availability-works#choosing-the-new-main)
during failover.

However, a replica that was **down** at the time of failover might actually hold
**more recent data** than any of the replicas that were online. When this
happens, Memgraph performs a **force sync** of that REPLICA.

A force sync occurs when a previously unreachable instance had a more up-to-date
state than the replica promoted to MAIN. Once that instance rejoins the cluster,
it undergoes a controlled recovery process:

- The REPLICA **resets its storage**.
- It receives **all committed transactions** from the current MAIN to rebuild an
  accurate and consistent state.
- Its original durability files are preserved in `.old` directories inside
  `data_directory/snapshots` and `data_directory/wal`. These allow
  administrators to attempt manual recovery if necessary. The `.old` directory
  is reused on subsequent recoveries, meaning **only one backup copy is kept at
  a time**.

Use the `--storage-backup-dir-enabled` flag to control this behavior:
- `true` (default) - Old durability files are moved to `.old` directories
- `false` - Old durability files are deleted immediately

---

The possibility of data loss during failover depends on the configured
replication mode:

#### SYNC (default)

- Provides strong availability guarantees but allows a **non-zero RPO** (loss of
  committed data).
- This can occur because the replica promoted to MAIN may not have received the
  latest commits from the old MAIN before failure.
- The design prioritizes keeping the MAIN writable for as long as possible.

#### ASYNC

- Similar to `SYNC`, but with an even higher chance of data loss because the
  MAIN continues committing freely regardless of replica status.

#### STRICT_SYNC

- Ensures **zero data loss** under all failover scenarios.
- Achieves this by using a two-phase commit protocol, at the cost of reduced
  throughput.

> **Note**
>
> Learn more about the implications of different [replication
> modes](https://memgraph.com/docs/clustering/replication/how-replication-works#replication-modes).

## Actions on follower coordinators

Follower coordinators never execute cluster operations themselves. Instead, they
act as a transparent entry point: every cluster query you run on a follower is
**forwarded to the current leader**, executed there, and the leader's answer is
returned to you. This holds both for state-changing operations (registering and
unregistering data instances, promoting and demoting instances, adding and
removing coordinators, updating configuration, forcing a cluster state reset,
[yielding
leadership](https://memgraph.com/docs/clustering/high-availability/ha-commands-reference#yield-leadership))
and for read-only ones (`SHOW INSTANCES`, `SHOW COORDINATOR SETTINGS`, `SHOW
REPLICATION LAG`, and `bolt+routing` routing table requests).

As a result you can point your tooling at any coordinator without first having to
discover which one is the leader.

### Why followers never answer from local state

A follower's own Raft state is not enough to describe the cluster: health of the
data instances is only known to the leader, which is the coordinator that pings
them. For this reason followers do not fall back to a local, partial answer when
the leader cannot be reached. Instead:

- Read queries return an **empty result set** together with a warning
  notification (`LeaderNotReachable`, or `ReplicationLagUnavailable` for `SHOW
  REPLICATION LAG`).
- State-changing queries fail with an error telling you that the coordinator is
  not the leader, or that the leader could not be found.

Both cases mean "cluster state unknown — retry", and both are expected during the
brief window of a leader election or while a newly elected leader is still taking
over the cluster. See [Error
handling](https://memgraph.com/docs/clustering/high-availability/ha-commands-reference#when-there-is-no-leader-to-serve-the-query)
for the exact messages.

The one exception is the `bolt+routing` routing table: a coordinator that Raft
elected as leader answers routing requests from its own Raft state even before it
has finished taking over the cluster, so that clients can keep routing queries
during the leadership transition.

## Raft-first operations and the reconciliation loop

The coordinator follows a **Raft-first** pattern for all cluster operations
(registering, unregistering, promoting, demoting instances). This means every
state change is first committed to the Raft log and acknowledged by a majority
of coordinators **before** the operation returns success to the user.

After the Raft commit, the coordinator sends RPCs to data instances (e.g.,
`PromoteToMainRpc`, `DemoteMainToReplicaRpc`, `RegisterReplicaOnMainRpc`,
`UnregisterReplicaRpc`) on a **best-effort** basis. If an RPC fails due to a
transient network issue, the operation still succeeds from the user's
perspective because the Raft log is the single source of truth.

### How the reconciliation loop works

The coordinator leader runs a periodic **reconciliation loop** that
automatically detects and corrects discrepancies between the desired state (Raft
log) and the actual state of data instances. Specifically:

- **Missing replicas on main**: If a replica exists in the Raft state but is not
  registered on the current main instance, the reconciliation loop sends a
  `RegisterReplicaOnMainRpc` to the main.
- **Stale replicas on main**: If the main instance reports a replica that no
  longer exists in the Raft state, the reconciliation loop sends an
  `UnregisterReplicaRpc` to remove it.

This self-healing behavior means the cluster automatically recovers from
transient RPC failures without user intervention. Users only need to retry an
operation if the Raft commit itself fails.

## Instance restarts

### Restarting data instances

Both MAIN and REPLICA instances may fail and later restart.

- When a **REPLICA** instance comes back online, it uses the MAIN UUID provided
  by the coordinator (via the `SwapMainUUIDReq` message) to determine which MAIN
  to follow. This synchronization happens automatically once the coordinator’s
  health check (“ping”) succeeds.

- When the **MAIN** instance restarts, the coordinator confirms the instance’s
  state through health checks. Writing is enabled once the
  coordinator verifies the instance is healthy and its role is confirmed by sending `PromoteToMainRpc` to the data instance.

This ensures that instances safely rejoin the cluster without causing
inconsistencies.

### Restarting coordinator instances

If a coordinator instance crashes and is restarted, it does **not** lose any
state. All coordinator metadata, RAFT logs and RAFT snapshots, is persisted on
durable storage, ensuring full recovery on restart.

More details are available in the section on [RAFT
implementation](#raft-implementation).

## In-service software upgrade (ISSU)

The cluster's ability to safely restart instances one at a time is what enables
**in-service software upgrades (ISSU)** — upgrading Memgraph to a newer version
with zero downtime.

An ISSU is a rolling upgrade: each instance is stopped and restarted on the new
version while the rest of the cluster keeps serving queries. To avoid downtime
and data loss, instances are upgraded in a fixed order:

1. **Replicas** one at a time.
2. The **main** instance — its restart triggers the standard [automatic
   failover](#automatic-failover), so a freshly upgraded replica takes over as
   the new main.
3. **Coordinator followers**
4. The **coordinator leader** whose restart triggers a Raft leader re-election.

For the full step-by-step procedure, including backups and Helm chart
configuration, see the [ISSU section on the Kubernetes setup
page](https://memgraph.com/docs/clustering/high-availability/setup-ha-cluster-k8s#in-service-software-upgrade-issu).
The same flow applies to native deployments.

## RAFT implementation

NuRaft persists all critical RAFT state to durable storage by default. This
includes:

- RAFT logs
- RAFT snapshots
- Cluster connectivity metadata

Persisting connectivity information is essential. Without it, a coordinator
could not safely rejoin the cluster after a restart.

---

Information about logs and snapshots is stored under one RocksDB instance in the
`high_availability/raft_data/logs` directory stored under the top-level
`--data-directory` folder. All the data stored there is recovered in case the
coordinator restarts.

Data about other coordinators is recovered from the
`high_availability/raft_data/network` directory stored under the top-level
`--data-directory` folder. When the coordinator rejoins, it will reestablish the
communication with other coordinators and receive updates from the current
leader.

### First start

When coordinators start for the first time, each initializes its `logs` and
`network` durability stores. From that moment on:

- Every RAFT log entry sent to a coordinator is also written to disk.
- The server configuration is updated whenever a new coordinator joins.
- Logs are generated for all user actions and failover-related operations.
- RAFT snapshots are created every *N* log entries (currently **every 5 logs**).

This design ensures that the coordinator can reliably recover its full RAFT
state at any time, preserving both consistency and cluster membership
information.

### Restart of coordinator 

In case of the coordinator's failure, on the restart, it will read information
about other coordinators stored under `high_availability/raft_data/network`
directory. 

From the `network` directory we will recover the server state before the
coordinator stopped, including the current term, for whom the coordinator voted,
and whether election timer is allowed.

It will also recover the following server config information: 
- other servers, including their endpoints, id, and auxiliary data
- ID of the previous log 
- ID of the current log 
- additional data needed by nuRaft

The following information will be recovered from a common RocksDB `logs`
instance: 
- current version of `logs` durability store 
- snapshots found with `snapshot_id_` prefix in database:
  - coordinator cluster state - all data instances with their role (main or
    replica), all coordinator instances and UUID of main instance which replica
    is listening to 
  - last log idx
  - last log term
  - last cluster config
- logs found in the interval between the start index and the last log index 
  - data - each log holds data on what has changed since the last state
  - term - nuRAFT term 
  - log type - nuRAFT log type

### Handling durability errors 

If snapshots are not correctly stored, the exception is thrown and left for the
nuRAFT library to handle the issue. Logs can be missed and not stored since they
are compacted and deleted every two snapshots and will be removed relatively
fast.

Memgraph throws an error when failing to store cluster config, which is updated
in the `high_availability/raft_data/network` folder. If this happens, it will
happen only on the first cluster start when coordinators are connecting since
coordinators are configured only once at the start of the whole cluster. This is
a non-recoverable error since in case the coordinator rejoins the cluster and
has the wrong state of other clusters, it can become a leader without being
connected to other coordinators. 

## Recovering from errors

Distributed systems can fail in many different ways. Memgraph is designed to be
resilient to **network failures**, **omission faults**, and **independent
machine failures**. However, like most systems based on Raft, it does **not**
tolerate Byzantine failures.

To understand how Memgraph behaves under failure, it is useful to consider the
**Recovery Time Objective (RTO)** - the maximum acceptable duration for which an
instance or cluster can be unavailable.

Memgraph clusters contain two categories of instances - **coordinators** and
**data instances** and their failure scenarios must be examined separately.

### Coordinator failures

Raft requires a **majority of coordinators** to remain alive to continue
operating.

- With **one coordinator down** in a 3-node Raft cluster (RF = 3), the cluster
  still has quorum, so **RTO = 0** - the system remains fully available.
- With **two or more coordinators down**, quorum is lost. In this case, RTO
  depends solely on how long it takes for enough coordinators to come back
  online.

### Replica failures

A replica’s failure has different consequences depending on its replication
mode:

1. **STRICT_SYNC replica fails**

  - Writes on the MAIN are **blocked** until the replica becomes available
    again.
  - Reads remain allowed on MAIN and other replicas.

2. **SYNC or ASYNC replica fails**

  - The MAIN continues accepting writes.
  - Reads remain available everywhere.
  - A write that couldn't reach a SYNC replica still succeeds and is reported
    with a `SyncReplicationFailure` warning
    [notification](https://memgraph.com/docs/database-management/query-metadata#notifications).

Thus, only STRICT_SYNC replicas can directly impact write availability.

### Main instance failure

When the MAIN instance becomes unavailable, the failure is handled by the leader
coordinator using two user-configured parameters:

- `instance_health_check_frequency_sec`: how often health checks are sent (configurable via [`SET COORDINATOR SETTING`](https://memgraph.com/docs/clustering/high-availability/ha-commands-reference#coordinator-runtime-settings))
- `instance_down_timeout_sec`: how long an instance must remain unresponsive
  before it is considered down (configurable via [`SET COORDINATOR SETTING`](https://memgraph.com/docs/clustering/high-availability/ha-commands-reference#coordinator-runtime-settings))

Once the coordinator gathers enough evidence that the MAIN is down, it begins a
failover procedure using a small number of RPC messages. The exact time required
depends on network latency and instance proximity.

If the selected replica (promoted to MAIN) has the most recent committed data,
failover completes **with zero data loss**.

### Databases broken by recovery failure

By default, a data instance whose durability files are corrupt refuses to boot
and its role in the cluster goes down. If you run data instances with
[`--storage-allow-recovery-failure=true`](https://memgraph.com/docs/database-management/configuration),
such an instance boots instead, with the affected database in a `broken` state
(empty placeholder, on-disk files untouched). Only that database is affected —
all healthy databases on the instance keep serving.

- **Broken database on a REPLICA:** it **self-heals**. Because the broken
  placeholder comes up empty with a fresh epoch, the MAIN detects the replica is
  behind, drives it into recovery and sends a full snapshot, which clears the
  broken state. No operator action is required.
- **Broken database on the MAIN:** the cluster does **not** automatically fail
  over based on per-database broken state — a MAIN that is broken for one
  database but otherwise healthy **stays MAIN**, so the cluster never thrashes
  over a single corrupt database. Replicas keep serving reads of their healthy
  copy. The admin recovers the MAIN in place with `RECOVER SNAPSHOT`, after which
  the database rejoins replication normally.

See [recovery failure handling](https://memgraph.com/docs/fundamentals/data-durability#recovery-failure-handling)
for the full behavior and restoration paths.

## Raft configuration parameters

Several RAFT-related settings play a key role in maintaining stable cluster
behavior.

The **leader coordinator** sends heartbeat messages to all follower coordinators
**once per second**. These heartbeats allow followers to confirm that the leader
is healthy.

This timing interacts closely with the **leader election timeout**, which is a
**randomized interval between 2000 ms and 4000 ms**. A follower starts a new
election only if it does **not** receive a heartbeat within this timeout window.

Additionally, the **leadership expiration** is fixed at **2000 ms**, ensuring
that the cluster cannot enter a state where multiple leaders coexist.

These values are chosen to allow the cluster to tolerate minor and temporary
network delays **without** triggering unnecessary leadership changes, improving
overall stability.
