Skip to content

RaftIQ

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.

Status

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. See docs/transport/security.md.
  • The scheduler and worker pool are started by the binary via a required --workers flag — but the worker's job handler shipped in cmd/raftiq only 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 as RaftNode methods; driving them requires embedding the Go API.
  • Distributed locking is implemented but not exposed over gRPC. It's reachable from internal/server.Server in Go, not from the network client.
  • Metrics and health endpoints are plain HTTP, not TLS-protected.

Key features

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 InstallSnapshot for 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.

Architecture

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]
Loading

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.

Repository structure

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)

Requirements

  • Go 1.26.5+ (see go.mod)
  • protoc only if regenerating api/proto/*.pb.go

Quick start

git clone https://github.com/sanchar127/raftiq.git
cd raftiq
go mod download
go build ./...
go test ./...

Running RaftIQ

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).

Single node (for local experimentation)

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.key

Multi-node cluster

Every 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.key

Membership is fixed at startup via BootstrapMembership() — there is no "join an existing cluster" flag. See docs/operations.md.

Testing

go test ./...
go test -race ./...
go vet ./...
go test ./tests/chaos/...
go test ./tests/e2e/... -run 'TestE2E(ProcessLifecycle|KVLeaderFailover)$' -v

make 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.

Go client example

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.

Documentation map

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

Contributing

See CONTRIBUTING.md, including the Raft safety invariants contributors must not violate.

Security

See SECURITY.md for the vulnerability-reporting process and docs/transport/security.md for the current TLS/mTLS implementation.

License

Apache License 2.0 — see LICENSE.

About

A production-oriented distributed KV store and fault-tolerant job scheduler built from scratch in Go using Raft consensus, durable WAL, snapshots, fencing, and chaos testing.

Resources

Code of conduct

Contributing

Security policy

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages