RaftIQ is a Raft consensus implementation written from scratch in Go, with a replicated key-value store, a distributed job scheduler/worker pool, and mutual-TLS-secured gRPC transport built on top of it.
The goal is to implement and validate distributed-systems guarantees from first principles — election safety, log matching, leader completeness, joint-consensus membership changes, linearizable reads, crash-safe WAL persistence — rather than depending on an existing consensus library.
For the full evidence behind every claim in this README, see
docs/AUDIT.md.
RaftIQ is an active, in-progress open-source project. It has not made a
tagged release (see CHANGELOG.md). Read this before
deploying it anywhere that matters:
- TLS/mTLS is mandatory, not optional. The binary will not start
without
--tls-ca,--tls-cert, and--tls-key. Seedocs/transport/security.md. - The scheduler and worker pool are started by the binary via a
required
--workersflag — but the worker's job handler shipped incmd/raftiqonly logs that a job ran; it does not execute real work. Replace the handler if you need actual job execution. - There is no leader-redirect mechanism. If you write to a non-leader
node, you get a plain error, not a pointer to the current leader.
Retry against another node in
--peers. - Membership changes (
AddMember/RemoveMember) have no RPC or CLI. They exist asRaftNodemethods; driving them requires embedding the Go API. - Distributed locking is implemented but not exposed over gRPC. It's
reachable from
internal/server.Serverin Go, not from the network client. - Metrics and health endpoints are plain HTTP, not TLS-protected.
Implemented and network-reachable in the raftiq binary:
- Raft leader election with PreVote, term management, log replication, and commit-index advancement.
- Linearizable reads via
ReadIndex. - Log compaction / snapshots, including
InstallSnapshotfor lagging followers. - A durable, CRC32-checksummed, crash-recoverable WAL.
- A replicated KV store (
Get/Put/Delete) over gRPC. - Job creation (
CreateJob) over gRPC, feeding a scheduler and worker pool that are started automatically by the binary. - Mandatory mutual TLS on both the Raft and KV gRPC services, with SAN-based peer verification for Raft peers.
- Prometheus metrics on
--metrics-addr, and a health/readiness endpoint on--health-addr(both plain HTTP).
Implemented and tested, but not reachable from the network today:
- Joint-consensus membership changes (
AddMember/RemoveMember) — Go API only, no RPC/CLI. - Distributed locking with monotonically increasing fencing tokens — Go API only.
- Job status querying, listing, claiming, and lifecycle transitions — only job creation is exposed over gRPC.
- Leader redirection on KV writes — not implemented.
Planned / scaffolding only:
- Docker packaging (no Dockerfile committed).
- A configuration file / env-based config loader (CLI flags only).
See docs/AUDIT.md for the full inventory with code and
test references.
graph TD
Client[gRPC Client, mTLS] -->|KVService| KVSvc[internal/transport.KVService]
Peer[Raft Peer, mTLS] -->|RaftService| RaftSvc[internal/transport.RaftService]
KVSvc --> AppServer[internal/server.Server]
RaftSvc --> RaftNode[internal/raft.RaftNode]
AppServer --> RaftNode
AppServer --> KVStore[internal/kv.Store]
AppServer --> LockState[internal/lock.State]
RaftNode -->|ApplyCh| Applier[internal/kv.Applier]
Applier --> KVStore
RaftNode --> WAL[internal/storage.WALStorage]
RaftNode -->|GRPCTransport, mTLS| Peer
Scheduler[internal/scheduler.Scheduler] --> RaftNode
Scheduler --> KVStore
Workers[internal/worker.Worker pool] --> KVStore
RaftNode -.-> Metrics[internal/observability.Metrics]
RaftNode -.-> Health[internal/observability.HealthServer]
internal/raft never imports net, gRPC, or a disk package directly — it
only talks outward through the Transport and Storage interfaces. See
docs/architecture.md for the full breakdown.
api/proto/ protobuf definitions and generated gRPC code
client/ Go gRPC client for the KV/job API
cmd/raftiq/ the CLI entry point / binary (config.go, runtime.go, main.go)
cmd/kvtest/ a minimal manual smoke-test client
internal/raft/ Raft consensus core (election, log, membership, snapshots)
internal/storage/ Storage interface, WAL, in-memory implementation
internal/model/ shared types (LogEntry, Configuration, Snapshot, Job, ...)
internal/kv/ deterministic KV + job state machine and applier
internal/lock/ fencing-token lock state
internal/scheduler/ job scheduling library, started by cmd/raftiq
internal/worker/ job worker library, started by cmd/raftiq
internal/transport/ gRPC transport, TLS/mTLS, KV/Raft services
internal/server/ glues raft + kv + lock together for cmd/raftiq
internal/observability/ Prometheus metrics, logging, health checks
tests/chaos/ in-process, multi-node failure-injection tests
tests/e2e/ real OS-process end-to-end tests (spawns the binary)
tests/benchmark/ placeholder directory (no committed results)
deploy/ Grafana dashboards + empty Prometheus/Docker scaffolding
docs/ detailed documentation (see map below)
- Go 1.26.5+ (see
go.mod) protoconly if regeneratingapi/proto/*.pb.go
git clone https://github.com/sanchar127/raftiq.git
cd raftiq
go mod download
go build ./...
go test ./...TLS certificates are required to run cmd/raftiq at all — see
docs/transport/security.md for how to
generate a development CA and per-node certificates (there's no shipped
script; tests/e2e/tls_test.go is the reference for what a valid
certificate set looks like).
go run ./cmd/raftiq \
--id node1 \
--raft-addr :7000 \
--kv-addr :8000 \
--metrics-addr :9090 \
--health-addr :8080 \
--data-dir ./data/node1 \
--peers node1=localhost:7000 \
--workers worker-a \
--tls-ca ./certs/ca.crt \
--tls-cert ./certs/node1.crt \
--tls-key ./certs/node1.keyEvery node needs the same --peers map (including itself), a unique
--id, its own TLS certificate (SAN <id>.raftiq, see
docs/transport/security.md), and at least
one --workers entry:
go run ./cmd/raftiq --id node1 --raft-addr :7001 --kv-addr :8001 \
--metrics-addr :9091 --health-addr :8081 --data-dir ./data/node1 \
--peers node1=localhost:7001,node2=localhost:7002,node3=localhost:7003 \
--workers worker-a \
--tls-ca ./certs/ca.crt --tls-cert ./certs/node1.crt --tls-key ./certs/node1.key
go run ./cmd/raftiq --id node2 --raft-addr :7002 --kv-addr :8002 \
--metrics-addr :9092 --health-addr :8082 --data-dir ./data/node2 \
--peers node1=localhost:7001,node2=localhost:7002,node3=localhost:7003 \
--workers worker-a \
--tls-ca ./certs/ca.crt --tls-cert ./certs/node2.crt --tls-key ./certs/node2.key
go run ./cmd/raftiq --id node3 --raft-addr :7003 --kv-addr :8003 \
--metrics-addr :9093 --health-addr :8083 --data-dir ./data/node3 \
--peers node1=localhost:7001,node2=localhost:7002,node3=localhost:7003 \
--workers worker-a \
--tls-ca ./certs/ca.crt --tls-cert ./certs/node3.crt --tls-key ./certs/node3.keyMembership is fixed at startup via BootstrapMembership() — there is no
"join an existing cluster" flag. See
docs/operations.md.
go test ./...
go test -race ./...
go vet ./...
go test ./tests/chaos/...
go test ./tests/e2e/... -run 'TestE2E(ProcessLifecycle|KVLeaderFailover)$' -vmake check also runs golangci-lint and govulncheck. See
docs/testing.md for what each layer covers, including
the real OS-process E2E tests and the CI pipeline.
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/sanchar127/raftiq/client"
)
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
c, err := client.Dial("localhost:8001")
if err != nil {
log.Fatalf("failed to connect: %v", err)
}
defer c.Close()
if err := c.Put(ctx, "example-key", []byte("example-value")); err != nil {
log.Fatalf("put failed: %v", err)
}
value, found, err := c.Get(ctx, "example-key")
if err != nil {
log.Fatalf("get failed: %v", err)
}
fmt.Printf("value: %s (found=%v)\n", value, found)
index, err := c.CreateJob(ctx, "job-1", []byte("payload"), 0)
if err != nil {
log.Fatalf("create job failed: %v", err)
}
fmt.Printf("job created at index %d\n", index)
}Note: client.Dial(address, opts ...grpc.DialOption) defaults to
plaintext (insecure.NewCredentials(), client/dial.go). Since a
cmd/raftiq node's --kv-addr requires mTLS, connecting to it means
passing your own grpc.WithTransportCredentials(credentials.NewTLS(...))
as an extra dial option to override the default. If the node you dial
isn't the current leader, Put/CreateJob fail with a plain gRPC error,
not a leader redirect — retry against another node in --peers.
| Area | Document |
|---|---|
| Audit / implementation inventory | docs/AUDIT.md |
| Architecture overview | docs/architecture.md |
| HLD + LLD | docs/system-design.md |
| Implementation walkthrough | docs/implementation.md |
| Development workflow | docs/development.md |
| Testing strategy | docs/testing.md |
| Configuration / CLI flags | docs/configuration.md |
| Running / operating a cluster | docs/operations.md |
| Debugging common failures | docs/troubleshooting.md |
| Persistence overview | docs/persistence.md |
| Raft: overview | docs/raft/overview.md |
| Raft: leader election | docs/raft/leader-election.md |
| Raft: log replication | docs/raft/log-replication.md |
| Raft: commitment | docs/raft/commitment.md |
| Raft: membership | docs/raft/membership.md |
| Raft: joint consensus | docs/raft/joint-consensus.md |
| Raft: snapshots | docs/raft/snapshots.md |
| Raft: linearizable reads | docs/raft/linearizable-reads.md |
| Raft: failure recovery | docs/raft/failure-recovery.md |
| Storage: overview | docs/storage/overview.md |
| Storage: WAL format | docs/storage/wal.md |
| Storage: snapshots | docs/storage/snapshots.md |
| Transport: overview | docs/transport/overview.md |
| Transport: local (test) transport | docs/transport/local.md |
| Transport: gRPC | docs/transport/grpc.md |
| Transport: TLS/mTLS | docs/transport/security.md |
| Component: KV store & jobs | docs/components/kv.md |
| Component: distributed lock | docs/components/lock.md |
| Component: scheduler | docs/components/scheduler.md |
| Component: worker | docs/components/worker.md |
| Component: observability | docs/components/observability.md |
See CONTRIBUTING.md, including the Raft safety
invariants contributors must not violate.
See SECURITY.md for the vulnerability-reporting process
and docs/transport/security.md for the
current TLS/mTLS implementation.
Apache License 2.0 — see LICENSE.