seqcd: A drop-in etcd replacement, built in two weeks on Aeron Sequencer
What can one engineer with Claude build using Aeron Sequencer in just two weeks?
I’m new to Aeron Sequencer, but not new to the pains and rigors of building production-ready, bullet-proof distributed systems. There’s a wide gulf between theory and production code. You can model an application in perfect TLA+ and still prudently fear Kyle’s machine. Aeron Sequencer promises to do the hard parts for you, handling high availability, consensus, concurrency, replication, snapshots, and crash recovery, and freeing you to implement your application using simple, single-threaded business logic. So I wanted to find out: with the hard parts already handled, how far could one engineer and an AI actually get?
One of the most widely used and battle-tested distributed applications in the world is etcd. It is a strongly-consistent key-value store that holds cluster configuration and state for nearly every Kubernetes deployment. It has been developed and hardened, with its fair share of bug fixes, since 2013. Consensus, leader election, log replication, crash consistency, snapshot recovery, deterministic replay — all refined over years of work. It would be the height of hubris to attempt an etcd clone in just two weeks’ time.
So anyway, two weeks later I had seqcd — a drop-in etcd v3 replacement, same gRPC wire protocol, tested with etcd’s own clients and benchmark driver. It turns out writing distributed applications is much easier when you start with an AI-friendly foundation that has already solved the most difficult and bug-prone 80% for you.
The 80% I didn’t write
Aeron Sequencer is an extension of Aeron Cluster — “a high performance Raft implementation that is deployed into many mission-critical production environments.” Consensus is reached by a majority of cluster nodes agreeing on the contents of a replicated log; your application is a Replicated State Machine that consumes that log. All of the hard parts of etcd were already there, ready to use:
Consensus and leader election. Raft, via Aeron Cluster. I wrote zero election or group membership code.
The replicated Log. Ordered, durable, identical on every node. Every write in seqcd is simply a message on this log.
Snapshots and recovery. The Sequencer’s Snapshot Service durably stores application state “alongside the Log position at which the snapshot was taken”; the Replay Service restores it and replays the log from that position. I got crash consistency and recovery for free.
Deterministic timers. The Timers API fires by the cluster’s log-derived clock, identically on every node and identically on replay. (This matters enormously for leases — more below.)
Sessions, egress, and back-pressure. Client session management and flow control are all handled by Sequencer.
The transport. Aeron’s high-throughput, low-latency messaging underneath, with the mechanical sympathy it’s known for.
The 20% I did write
The entire replicated heart of seqcd is one class: KvStateMachine, about 1,600 lines in a ~6.7K-line core module. It implements etcd v3 semantics — MVCC revisions, ranges, transactions, leases, compaction, watch history — as plain, single-threaded Java. No locks. No atomics. No threads. No dependency on wall clock time. Completely deterministic.
Below on the left is seqcd’s lease-expiry code — complete and verbatim from KvStateMachine.java: the grant handler, the deadline scheduler, the deterministic expiry timer, and the key deletion it triggers. On the right, the same feature in etcd (v3.6.12, server/lease/lessor.go — an 887-line file): the lease lifecycle — grant, revoke, the expiry trigger, and the checkpointing that keeps remaining TTLs consistent across leader changes.
seqcd — KvStateMachine.java — the lease grant, expiry timer, and key-deletion path
@Event(protocolId = KvEvents.KV_PROTOCOL_ID, eventTypeId = KvEvents.LEASE_GRANT_REQUEST_EVENT_TYPE)
public void onLeaseGrant(final LeaseGrantRequest req) {
final long cid = req.clientId(); final long corr = req.correlationId();
final Object hit = dedupHit(cid, corr);
if (hit != null) {
// re-emit the cached grant (same lease id) -- a resend must NOT mint a 2nd lease
leaseGrantResponseListener.onLeaseGrantResponse((LeaseGrantResponse)hit);
return;
}
long leaseId = req.requestedId();
final LeaseGrantResponse resp;
if (leaseId == 0L) {
leaseId = nextLeaseId++;
} else if (leases.containsKey(leaseId)) {
resp = new LeaseGrantResponse(req.correlationId(), leaseId, 0L, "etcdserver: lease already exists");
cache(cid, corr, resp);
leaseGrantResponseListener.onLeaseGrantResponse(resp);
return;
} else if (leaseId >= nextLeaseId) {
nextLeaseId = leaseId + 1L;
}
leases.put(leaseId, new LeaseState(req.ttlSeconds(), new TreeSet<>(Arrays::compareUnsigned)));
scheduleExpiry(leaseId, req.ttlSeconds());
resp = new LeaseGrantResponse(req.correlationId(), leaseId, req.ttlSeconds(), "");
cache(cid, corr, resp);
leaseGrantResponseListener.onLeaseGrantResponse(resp);
}
private void scheduleExpiry(final long leaseId, final long ttlSeconds) {
if (ttlSeconds <= 0L) {
return;
}
// saturate: a TTL near etcd's MaxLeaseTTL cap overflows epochNs + ttlNs, wrapping the
// deadline negative and firing immediately -- clamp to "never" instead
final long ttlNs = ttlSeconds > Long.MAX_VALUE / 1_000_000_000L ?
Long.MAX_VALUE : ttlSeconds * 1_000_000_000L;
final long nowNs = container.epochTimeNs();
final long deadlineNs = nowNs > Long.MAX_VALUE - ttlNs ? Long.MAX_VALUE : nowNs + ttlNs;
final LeaseState state = leases.get(leaseId);
if (state != null) {
state.expiresAtNs = deadlineNs;
}
container.scheduleTimer(leaseId, deadlineNs);
}
// deterministic lease expiry: fires when the cluster's log clock crosses the scheduled
// deadline; the timer correlationId is the lease id
@Override
public void onTimerEvent(final TimerEvent event) {
final long leaseId = event.correlationId();
final LeaseState state = leases.remove(leaseId);
if (state != null) {
deleteLeaseKeys(state);
}
}
private void deleteLeaseKeys(final LeaseState state) {
if (state == null || state.keys.isEmpty()) {
return;
}
currentRevision++;
// pre-filter live keys (skip tombstoned/missing) so the LAST emitted event is known
// up front; deleteTargets is safe to reuse -- expiry/revoke never runs mid-command
final ArrayList<byte[]> targets = deleteTargets;
targets.clear();
for (final byte[] key : state.keys) {
final ArrayList<Entry> versions = history.get(key);
if (versions != null && !versions.isEmpty() && !versions.get(versions.size() - 1).tombstone) {
targets.add(key);
}
}
for (int i = 0; i < targets.size(); i++) {
final byte[] key = targets.get(i);
final ArrayList<Entry> versions = history.get(key);
final Entry prev = versions.get(versions.size() - 1);
versions.add(new Entry(OffHeapValueArena.NULL_REF, 0, 0L, currentRevision, 0L, true, 0L));
kvMutationEventListener.onKvMutation(new KvMutationEvent(
key, EMPTY, prev.createRevision, currentRevision, 0L, true,
arena.readBytes(prev.valueRef, prev.valueLength), prev.modRevision,
i == targets.size() - 1));
}
}
etcd — server/lease/lessor.go — grant, revoke, expiry, and checkpoint paths (scroll for more)
func (le *lessor) runLoop() {
defer close(le.doneC)
delayTicker := time.NewTicker(500 * time.Millisecond)
defer delayTicker.Stop()
for {
le.revokeExpiredLeases()
le.checkpointScheduledLeases()
select {
case <-delayTicker.C:
case <-le.stopC:
return
}
}
}
func (le *lessor) revokeExpiredLeases() {
var ls []*Lease
// rate limit
revokeLimit := le.leaseRevokeRate / 2
le.mu.RLock()
if le.isPrimary() {
ls = le.findExpiredLeases(revokeLimit)
}
le.mu.RUnlock()
if len(ls) != 0 {
select {
case <-le.stopC:
return
case le.expiredC <- ls:
default:
// the receiver of expiredC is probably busy handling
// other stuff
// let's try this next time after 500ms
}
}
}
func (le *lessor) findExpiredLeases(limit int) []*Lease {
leases := make([]*Lease, 0, 16)
for {
l, next := le.expireExists()
if l == nil && !next {
break
}
if next {
continue
}
if l.expired() {
leases = append(leases, l)
// reach expired limit
if len(leases) == limit {
break
}
}
}
return leases
}
func (le *lessor) expireExists() (l *Lease, next bool) {
if le.leaseExpiredNotifier.Len() == 0 {
return nil, false
}
item := le.leaseExpiredNotifier.Peek()
l = le.leaseMap[item.id]
if l == nil {
// lease has expired or been revoked
// no need to revoke (nothing is expiry)
le.leaseExpiredNotifier.Unregister() // O(log N)
return nil, true
}
now := time.Now()
if now.Before(item.time) /* item.time: expiration time */ {
// Candidate expirations are caught up, reinsert this item
// and no need to revoke (nothing is expiry)
return nil, false
}
// recheck if revoke is complete after retry interval
item.time = now.Add(le.expiredLeaseRetryInterval)
le.leaseExpiredNotifier.RegisterOrUpdate(item)
return l, false
}
func (le *lessor) Grant(id LeaseID, ttl int64) (*Lease, error) {
if id == NoLease {
return nil, ErrLeaseNotFound
}
if ttl > MaxLeaseTTL {
return nil, ErrLeaseTTLTooLarge
}
// TODO: when lessor is under high load, it should give out lease
// with longer TTL to reduce renew load.
l := NewLease(id, ttl)
le.mu.Lock()
defer le.mu.Unlock()
if _, ok := le.leaseMap[id]; ok {
return nil, ErrLeaseExists
}
if l.ttl < le.minLeaseTTL {
l.ttl = le.minLeaseTTL
}
if le.isPrimary() {
l.refresh(0)
} else {
l.forever()
}
le.leaseMap[id] = l
l.persistTo(le.b)
leaseTotalTTLs.Observe(float64(l.ttl))
leaseGranted.Inc()
if le.isPrimary() {
item := &LeaseWithTime{id: l.ID, time: l.expiry}
le.leaseExpiredNotifier.RegisterOrUpdate(item)
le.scheduleCheckpointIfNeeded(l)
}
return l, nil
}
func (le *lessor) Revoke(id LeaseID) error {
le.mu.Lock()
l := le.leaseMap[id]
if l == nil {
le.mu.Unlock()
return ErrLeaseNotFound
}
defer close(l.revokec)
// unlock before doing external work
le.mu.Unlock()
if le.rd == nil {
return nil
}
txn := le.rd()
// sort keys so deletes are in same order among all members,
// otherwise the backend hashes will be different
keys := l.Keys()
sort.StringSlice(keys).Sort()
for _, key := range keys {
txn.DeleteRange([]byte(key), nil)
}
le.mu.Lock()
defer le.mu.Unlock()
delete(le.leaseMap, l.ID)
// lease deletion needs to be in the same backend transaction with the
// kv deletion. Or we might end up with not executing the revoke or not
// deleting the keys if etcdserver fails in between.
schema.UnsafeDeleteLease(le.b.BatchTx(), &leasepb.Lease{ID: int64(l.ID)})
txn.End()
leaseRevoked.Inc()
return nil
}
// checkpointScheduledLeases finds all scheduled lease checkpoints that are due and
// submits them to the checkpointer to persist them to the consensus log.
func (le *lessor) checkpointScheduledLeases() {
// rate limit
for i := 0; i < leaseCheckpointRate/2; i++ {
var cps []*pb.LeaseCheckpoint
le.mu.Lock()
if le.isPrimary() {
cps = le.findDueScheduledCheckpoints(maxLeaseCheckpointBatchSize)
}
le.mu.Unlock()
if len(cps) != 0 {
if err := le.cp(context.Background(), &pb.LeaseCheckpointRequest{Checkpoints: cps}); err != nil {
return
}
}
if len(cps) < maxLeaseCheckpointBatchSize {
return
}
}
}
func (le *lessor) scheduleCheckpointIfNeeded(lease *Lease) {
if le.cp == nil {
return
}
if lease.getRemainingTTL() > int64(le.checkpointInterval.Seconds()) {
if le.lg != nil {
le.lg.Debug("Scheduling lease checkpoint",
zap.Int64("leaseID", int64(lease.ID)),
zap.Duration("intervalSeconds", le.checkpointInterval),
)
}
heap.Push(&le.leaseCheckpointHeap, &LeaseWithTime{
id: lease.ID,
time: time.Now().Add(le.checkpointInterval),
})
}
}
None of this is a knock on etcd — this is what correct lease revocation has to look like when expiry is decided by a wall clock on one node: a background goroutine ticking every 500ms, a reader-writer mutex, an isPrimary() check because followers must never revoke, a rate limit so a thundering herd of expirations doesn’t stall the apply loop, and a channel handoff with a default: branch that quietly gives up and retries next tick. Every one of those lines exists because concurrency and wall-clock time create nondeterminism that must be handled.
Look at what is not in the seqcd version. None of those moving parts exist, because there is no concurrency and no wall clock. In seqcd, a lease is granted by a log message, its expiry is a replicated timer keyed by the lease id, and expiry is just another deterministic event — every node deletes the same keys at the same log position, and a replay after a crash reaches the exact same state.
Another property Aeron Sequencer gave me nearly for free: exactly-once semantics. Every request carries a (clientId, correlationId) idempotency key; the state machine caches responses in a dedup cache, and the dedup state is included in snapshots — so a client retrying across a server crash still gets the original answer instead of a double-applied write.
Three more examples
The lease path is not a cherry-pick. Nearly every “hard part” of a replicated key-value store dissolves the same way once ordering and time come from the log. Here are three more, with the real code on both sides. (All etcd references are v3.6.12, commit 90b034a.)
Linearizable reads: the same round-trip, minus the protocol
The problem: a node that thinks it’s the leader might have been deposed, so serving a read from local state risks returning stale data. etcd solves this with the ReadIndex protocol — a dedicated goroutine that batches waiting readers, runs a quorum round-trip to confirm leadership, then waits for the apply loop to catch up to the confirmed index.
seqcd runs ReadIndex too — same rule, same round-trip. What disappears is the protocol around it. Reads are served by frontends off-cluster: the frontend parks readers that need linearizability, submits a LinearizableSequenceIndexTask, and the Sequencer answers only once a quorum has confirmed leadership. The panes below are deliberately asymmetric: etcd’s shows only its confirmation round, with its batching pump (linearizableReadLoop + linearizableReadNotify) left out; seqcd’s shows the application-side read path pump and all — bar a few small helpers — because that is essentially all there is.
// EtcdIngressService.java -- in-cluster read code: none. KvStateMachine never sees a
// read. The confirmation round is the Sequencer's: a LinearizableSequenceIndexTask
// completes only after a quorum confirms, AFTER the request arrived, that this leader
// still leads (etcd's ReadIndex rule). Non-negative reply: a committed index proven
// current. Negative reply: a decline (election / partition) -- re-park and retry on a
// fresh generation, whose own confirmation covers the readers.
private static final class ReadIndexSlot {
private LinearizableSequenceIndexTask task; // in-flight handle; null once confirmed
private final LinearizableSequenceIndexTask taskInstance = new LinearizableSequenceIndexTask();
private final ArrayList<Runnable> batch = new ArrayList<>();
private long required; // committed-index gate
private boolean occupied;
}
private int pumpReadIndex() {
int work = 0;
for (final ReadIndexSlot slot : riSlots) {
if (!slot.occupied) {
continue;
}
final LinearizableSequenceIndexTask task = slot.task;
if (task != null) {
if (!task.isComplete()) {
continue;
}
final long value = task.isCompletedSuccessfully() ? task.linearizableSequenceIndex() : -1L;
slot.task = null;
work++;
if (value >= 0) {
// Leadership confirmed by quorum (after the request) -> serve at the committed index.
slot.required = value - 1L;
} else {
// Decline (election / partition / term change): re-park to retry. Re-parked reads
// ride a FRESH generation, whose own confirmation covers them.
reparkLinReads(slot.batch);
slot.batch.clear();
slot.occupied = false;
readIndexDeclines.incrementAndGet();
continue;
}
}
// Confirmed: serve once the consumed sequence index reaches the committed gate.
if (lastReadSequenceIndex.getAsLong() >= slot.required) {
serveReadIndexBatch(slot.batch);
slot.batch.clear();
slot.occupied = false;
work++;
}
}
if (hasPendingLinReads()) {
final ReadIndexSlot slot = freeReadIndexSlot();
if (slot != null) {
drainPendingLinReads(slot.batch);
if (!slot.batch.isEmpty()) {
slot.occupied = true;
slot.task = ingress.submit(slot.taskInstance); // per-slot instance: no per-RTT alloc
readIndexReads.addAndGet(slot.batch.size());
work++;
}
}
}
return work;
}
etcd — server/etcdserver/v3_server.go — the confirmation round alone: requestCurrentIndex + sendReadIndex; its batching pump is not shown (scroll for more)
func (s *EtcdServer) requestCurrentIndex(leaderChangedNotifier <-chan struct{}) (uint64, error) {
requestIDs := map[uint64]struct{}{}
requestID := s.reqIDGen.Next()
requestIDs[requestID] = struct{}{}
err := s.sendReadIndex(requestID)
if err != nil {
return 0, err
}
lg := s.Logger()
errorTimer := time.NewTimer(s.Cfg.ReqTimeout())
defer errorTimer.Stop()
retryTimer := time.NewTimer(readIndexRetryTime)
defer retryTimer.Stop()
firstCommitInTermNotifier := s.firstCommitInTerm.Receive()
for {
select {
case rs := <-s.r.readStateC:
// Check again if leader changed as when multiple channels are ready, select picks randomly.
select {
case <-leaderChangedNotifier:
readIndexFailed.Inc()
return 0, errors.ErrLeaderChanged
default:
}
responseID := uint64(0)
if len(rs.RequestCtx) == 8 {
responseID = binary.BigEndian.Uint64(rs.RequestCtx)
}
if _, ok := requestIDs[responseID]; !ok {
// a previous request might time out. now we should ignore the response of it and
// continue waiting for the response of the current requests.
lg.Warn(
"ignored out-of-date read index response; local node read indexes queueing up and waiting to be in sync with leader",
zap.Uint64("received-request-id", responseID),
)
slowReadIndex.Inc()
continue
}
return rs.Index, nil
case <-leaderChangedNotifier:
readIndexFailed.Inc()
// return a retryable error.
return 0, errors.ErrLeaderChanged
case <-firstCommitInTermNotifier:
firstCommitInTermNotifier = s.firstCommitInTerm.Receive()
lg.Info("first commit in current term: resending ReadIndex request")
requestID = s.reqIDGen.Next()
requestIDs[requestID] = struct{}{}
err := s.sendReadIndex(requestID)
if err != nil {
return 0, err
}
retryTimer.Reset(readIndexRetryTime)
continue
case <-retryTimer.C:
lg.Warn(
"waiting for ReadIndex response took too long, retrying",
zap.Uint64("sent-request-id", requestID),
zap.Duration("retry-timeout", readIndexRetryTime),
)
requestID = s.reqIDGen.Next()
requestIDs[requestID] = struct{}{}
err := s.sendReadIndex(requestID)
if err != nil {
return 0, err
}
retryTimer.Reset(readIndexRetryTime)
continue
case <-errorTimer.C:
lg.Warn(
"timed out waiting for read index response (local node might have slow network)",
zap.Duration("timeout", s.Cfg.ReqTimeout()),
)
slowReadIndex.Inc()
return 0, errors.ErrTimeout
case <-s.stopping:
return 0, errors.ErrStopped
}
}
}
func (s *EtcdServer) sendReadIndex(requestIndex uint64) error {
ctxToSend := uint64ToBigEndianBytes(requestIndex)
cctx, cancel := context.WithTimeout(context.Background(), s.Cfg.ReqTimeout())
err := s.r.ReadIndex(cctx, ctxToSend)
cancel()
if errorspkg.Is(err, raft.ErrStopped) {
return err
}
if err != nil {
lg := s.Logger()
lg.Warn("failed to get read index from Raft", zap.Error(err))
readIndexFailed.Inc()
return err
}
return nil
}
Both sides delegate the quorum round-trip itself to their consensus layer: etcd to raft’s ReadIndex, seqcd to the Sequencer’s LinearizableSequenceIndexTask. What each pane shows is the cost of driving that primitive safely from the application. Driving raft’s is a protocol: responses arrive out of order, for old requests, from deposed leaders, on another goroutine. So etcdserver wraps it in a hundred lines of machinery — request IDs matched against out-of-date responses, a retry timer, a resend on first-commit-in-term, a leader-change race checked twice, two timeouts — and around all of that, not shown, its batching pump.
seqcd’s pane is the whole thing, pump included, and it is single-threaded bookkeeping: park, submit, serve. The task completes on the same agent thread that parked the readers — exactly once, in order — so there is nothing to correlate (the task handle is the correlation), nothing to time out, and no leader-change double-check: a decline is the leader-change signal, and recovery is a single step — re-park the readers onto a fresh generation whose own confirmation covers them. The quorum round-trip still happens on every batch — that cost is irreducible. It just happens below the application, where confirming leadership was already the platform’s job — and the read path never enters the cluster at all.
Watches: three watcher sets and two background loops vs. 35 lines of event ordering
etcd’s watch implementation is a small distributed system of its own. Watchers live in three sets — synced, unsynced, and victims (watchers whose delivery channel blocked) — and two background goroutines shuffle them between the sets under a mutex.
In seqcd, every mutation already leaves the state machine on a totally ordered egress stream — the same stream, in the same order, for every subscriber. A watch is just a subscription. The state machine’s entire event-ordering machinery on the left; etcd’s full resync path on the right:
seqcd — KvStateMachine.java — the event-ordering machinery
// one revision bump per mutation; inside a txn the bump happens lazily on the first
// actual mutation and is shared by the rest (a txn with zero mutations never bumps)
private long mutationRevision() {
if (inTxn) {
if (txnRevision == 0L) {
txnRevision = ++currentRevision;
}
return txnRevision;
}
return ++currentRevision;
}
// emit to the live listener, or buffer during a txn so onTxn can emit in unsigned-key order
private void emitMutation(final KvMutationEvent evt) {
if (inTxn) {
txnMutationBuffer.add(evt);
} else {
kvMutationEventListener.onKvMutation(evt);
}
}
@Event(protocolId = KvEvents.KV_PROTOCOL_ID, eventTypeId = KvEvents.WATCH_HISTORY_REQUEST_EVENT_TYPE)
public void onWatchHistory(final WatchHistoryRequest req) {
final List<KvMutationEvent> events = watchHistory(req.key(), req.rangeEnd(), req.startRevision());
if (events == null) {
// ErrCompacted: startRev below the compaction floor; the frontend cancels the
// watch on compactedRevision > 0
watchHistoryResponseListener.onWatchHistoryResponse(
new WatchHistoryResponse(req.correlationId(), List.of(), currentRevision, compactedRevision));
return;
}
// currentRevision is exactly the upper bound the scan covered; the frontend uses it
// as the dedup boundary for its live buffer. compactedRevision 0 => not compacted.
watchHistoryResponseListener.onWatchHistoryResponse(
new WatchHistoryResponse(req.correlationId(), events, currentRevision, 0L));
}
etcd — server/storage/mvcc/watchable_store.go — the resync path (scroll for more)
func (s *watchableStore) syncWatchersLoop() {
defer s.wg.Done()
delayTicker := time.NewTicker(watchResyncPeriod)
defer delayTicker.Stop()
for {
s.mu.RLock()
st := time.Now()
lastUnsyncedWatchers := s.unsynced.size()
s.mu.RUnlock()
unsyncedWatchers := 0
if lastUnsyncedWatchers > 0 {
unsyncedWatchers = s.syncWatchers()
}
syncDuration := time.Since(st)
delayTicker.Reset(watchResyncPeriod)
// more work pending?
if unsyncedWatchers != 0 && lastUnsyncedWatchers > unsyncedWatchers {
// be fair to other store operations by yielding time taken
delayTicker.Reset(syncDuration)
}
select {
case <-delayTicker.C:
case <-s.stopc:
return
}
}
}
func (s *watchableStore) syncWatchers() int {
s.mu.Lock()
defer s.mu.Unlock()
if s.unsynced.size() == 0 {
return 0
}
s.store.revMu.RLock()
defer s.store.revMu.RUnlock()
// in order to find key-value pairs from unsynced watchers, we need to
// find min revision index, and these revisions can be used to
// query the backend store of key-value pairs
curRev := s.store.currentRev
compactionRev := s.store.compactMainRev
wg, minRev := s.unsynced.choose(maxWatchersPerSync, curRev, compactionRev)
evs := rangeEvents(s.store.lg, s.store.b, minRev, curRev+1, wg)
victims := make(watcherBatch)
wb := newWatcherBatch(wg, evs)
for w := range wg.watchers {
if w.minRev < compactionRev {
// Skip the watcher that failed to send compacted watch response due to w.ch is full.
// Next retry of syncWatchers would try to resend the compacted watch response to w.ch
continue
}
w.minRev = max(curRev+1, w.minRev)
eb, ok := wb[w]
if !ok {
// bring un-notified watcher to synced
s.synced.add(w)
s.unsynced.delete(w)
continue
}
if eb.moreRev != 0 {
w.minRev = eb.moreRev
}
if w.send(WatchResponse{WatchID: w.id, Events: eb.evs, Revision: curRev}) {
pendingEventsGauge.Add(float64(len(eb.evs)))
} else {
w.victim = true
}
if w.victim {
victims[w] = eb
} else {
if eb.moreRev != 0 {
// stay unsynced; more to read
continue
}
s.synced.add(w)
}
s.unsynced.delete(w)
}
s.addVictim(victims)
vsz := 0
for _, v := range s.victims {
vsz += len(v)
}
slowWatcherGauge.Set(float64(s.unsynced.size() + vsz))
return s.unsynced.size()
}
That’s server/storage/mvcc/watchable_store.go (614 lines) plus server/storage/mvcc/watcher_group.go (291 lines) — 905 lines to guarantee that every watcher sees every event, in order, exactly once, while writers race the delivery loops.
No synced/unsynced/victim states, because a subscriber can’t “fall behind the store” — the stream is the store’s history. Watchers that start at an old revision get a backfill (onWatchHistory above — a plain scan over the same history map), and slow consumers are the transport’s flow-control problem, not the state machine’s correctness problem.
Compaction: scheduled batches vs. one pass over a map
Compacting old revisions in etcd means rewriting the in-memory key index and deleting from bbolt on disk, while reads and writes continue. So it’s split across server/storage/mvcc/kvstore_compaction.go (98 lines of batched deletes that sleep between batches to avoid stalling writers), plus the per-key compaction logic in server/storage/mvcc/key_index.go (389 lines) and server/storage/mvcc/index.go (253 lines) — with a “scheduled compaction” watermark persisted so a crash mid-compaction can resume.
seqcd’s compaction is one walk over a map, because nothing can interleave with it. The whole feature — the handler, the walk, and the arena-release helper — on the left; etcd’s batched loop and the per-key compaction logic it drives on the right:
seqcd — KvStateMachine.java — the whole compaction feature
@Event(protocolId = KvEvents.KV_PROTOCOL_ID, eventTypeId = KvEvents.COMPACT_REQUEST_EVENT_TYPE)
public void onCompact(final CompactRequest req) {
final long target = req.revision();
if (target > currentRevision) {
compactResponseListener.onCompactResponse(
new CompactResponse(req.correlationId(), compactedRevision, currentRevision, target));
return;
}
if (target > compactedRevision) {
compactStore(target);
compactedRevision = target;
}
compactResponseListener.onCompactResponse(
new CompactResponse(req.correlationId(), compactedRevision, currentRevision, 0L));
}
private void compactStore(final long rev) {
final Iterator<Map.Entry<byte[], ArrayList<Entry>>> it = history.entrySet().iterator();
while (it.hasNext()) {
final ArrayList<Entry> versions = it.next().getValue();
int floor = -1;
for (int i = 0; i < versions.size(); i++) {
if (versions.get(i).modRevision <= rev) {
floor = i;
} else {
break;
}
}
if (floor < 0) {
continue;
}
final Entry floorEntry = versions.get(floor);
if (floorEntry.tombstone) {
if (floor == versions.size() - 1) {
freeRange(versions, 0, versions.size());
it.remove();
} else {
freeRange(versions, 0, floor + 1);
versions.subList(0, floor + 1).clear();
}
} else {
if (floor > 0) {
freeRange(versions, 0, floor);
versions.subList(0, floor).clear();
}
}
}
}
// release the off-heap value regions of versions [from, to) before they are dropped
private void freeRange(final ArrayList<Entry> versions, final int from, final int to) {
for (int i = from; i < to; i++) {
final Entry e = versions.get(i);
arena.free(e.valueRef, e.valueLength);
}
}
etcd — server/storage/mvcc — the batched loop + per-key compaction, across three files (scroll for more)
func (s *store) scheduleCompaction(compactMainRev, prevCompactRev int64) (KeyValueHash, error) {
totalStart := time.Now()
keep := s.kvindex.Compact(compactMainRev)
indexCompactionPauseMs.Observe(float64(time.Since(totalStart) / time.Millisecond))
totalStart = time.Now()
defer func() { dbCompactionTotalMs.Observe(float64(time.Since(totalStart) / time.Millisecond)) }()
keyCompactions := 0
defer func() { dbCompactionKeysCounter.Add(float64(keyCompactions)) }()
defer func() { dbCompactionLast.Set(float64(time.Now().Unix())) }()
end := make([]byte, 8)
binary.BigEndian.PutUint64(end, uint64(compactMainRev+1))
batchNum := s.cfg.CompactionBatchLimit
h := newKVHasher(prevCompactRev, compactMainRev, keep)
last := make([]byte, 8+1+8)
for {
var rev Revision
start := time.Now()
tx := s.b.BatchTx()
tx.LockOutsideApply()
keys, values := tx.UnsafeRange(schema.Key, last, end, int64(batchNum))
for i := range keys {
rev = BytesToRev(keys[i])
if _, ok := keep[rev]; !ok {
tx.UnsafeDelete(schema.Key, keys[i])
keyCompactions++
}
h.WriteKeyValue(keys[i], values[i])
}
if len(keys) < batchNum {
// gofail: var compactBeforeSetFinishedCompact struct{}
UnsafeSetFinishedCompact(tx, compactMainRev)
tx.Unlock()
dbCompactionPauseMs.Observe(float64(time.Since(start) / time.Millisecond))
// gofail: var compactAfterSetFinishedCompact struct{}
hash := h.Hash()
size, sizeInUse := s.b.Size(), s.b.SizeInUse()
s.lg.Info(
"finished scheduled compaction",
zap.Int64("compact-revision", compactMainRev),
zap.Duration("took", time.Since(totalStart)),
zap.Uint32("hash", hash.Hash),
zap.Int64("current-db-size-bytes", size),
zap.String("current-db-size", humanize.Bytes(uint64(size))),
zap.Int64("current-db-size-in-use-bytes", sizeInUse),
zap.String("current-db-size-in-use", humanize.Bytes(uint64(sizeInUse))),
)
return hash, nil
}
tx.Unlock()
// update last
last = RevToBytes(Revision{Main: rev.Main, Sub: rev.Sub + 1}, last)
// Immediately commit the compaction deletes instead of letting them accumulate in the write buffer
// gofail: var compactBeforeCommitBatch struct{}
s.b.ForceCommit()
// gofail: var compactAfterCommitBatch struct{}
dbCompactionPauseMs.Observe(float64(time.Since(start) / time.Millisecond))
select {
case <-time.After(s.cfg.CompactionSleepInterval):
case <-s.stopc:
return KeyValueHash{}, fmt.Errorf("interrupted due to stop signal")
}
}
}
// ---- server/storage/mvcc/index.go ----
func (ti *treeIndex) Compact(rev int64) map[Revision]struct{} {
available := make(map[Revision]struct{})
ti.lg.Info("compact tree index", zap.Int64("revision", rev))
ti.Lock()
clone := ti.tree.Clone()
ti.Unlock()
clone.Ascend(func(keyi *keyIndex) bool {
// Lock is needed here to prevent modification to the keyIndex while
// compaction is going on or revision added to empty before deletion
ti.Lock()
keyi.compact(ti.lg, rev, available)
if keyi.isEmpty() {
_, ok := ti.tree.Delete(keyi)
if !ok {
ti.lg.Panic("failed to delete during compaction")
}
}
ti.Unlock()
return true
})
return available
}
// ---- server/storage/mvcc/key_index.go ----
// compact compacts a keyIndex by removing the versions with smaller or equal
// revision than the given atRev except the largest one.
// If a generation becomes empty during compaction, it will be removed.
func (ki *keyIndex) compact(lg *zap.Logger, atRev int64, available map[Revision]struct{}) {
if ki.isEmpty() {
lg.Panic(
"'compact' got an unexpected empty keyIndex",
zap.String("key", string(ki.key)),
)
}
genIdx, revIndex := ki.doCompact(atRev, available)
g := &ki.generations[genIdx]
if !g.isEmpty() {
// remove the previous contents.
if revIndex != -1 {
g.revs = g.revs[revIndex:]
}
}
// remove the previous generations.
ki.generations = ki.generations[genIdx:]
}
func (ki *keyIndex) doCompact(atRev int64, available map[Revision]struct{}) (genIdx int, revIndex int) {
// walk until reaching the first revision smaller or equal to "atRev",
// and add the revision to the available map
f := func(rev Revision) bool {
if rev.Main <= atRev {
available[rev] = struct{}{}
return false
}
return true
}
genIdx, g := 0, &ki.generations[0]
// find first generation includes atRev or created after atRev
for genIdx < len(ki.generations)-1 {
if tomb := g.revs[len(g.revs)-1].Main; tomb >= atRev {
break
}
genIdx++
g = &ki.generations[genIdx]
}
revIndex = g.walk(f)
return genIdx, revIndex
}
A single-threaded walk that drops superseded versions and frees their off-heap values. No batching, no sleeping, no resume watermark — a crash mid-compaction is fine because compaction, like everything else, is a deterministic function of the log and simply replays.
Code size
The seqcd code samples above are complete for the paths shown; the etcd samples show just the relevant call paths, not whole files. A more complete line-count comparison is shown below. The seqcd column is every line of KvStateMachine that implements the feature listed — plus, for reads, the frontend ReadIndex pump in EtcdIngressService.java, since the read path deliberately lives off-cluster entirely. The etcd column is the files (or, for reads, the specific functions) that implement the same feature at v3.6.12, commit 90b034a.
Two honest caveats. The reads row counts application code on both sides, and both sides delegate the quorum round-trip itself to their consensus layer — etcd to raft’s ReadIndex, seqcd to the Sequencer’s task — so the 1.6× is the cost of driving that primitive, not of implementing it. And the watch ratio flatters seqcd: etcd’s 905 lines include per-client fan-out, while seqcd’s per-client delivery, backfill dedup and slow-consumer handling live in the frontends and aren’t counted here. The structural point survives both: the replicated core stays small because everything that can be pushed to stateless, restartable frontends is pushed there — and it’s the replicated core where lines of code translate into consensus-coupled failure modes.
No need to focus on precise line count comparisons here; for one thing, seqcd is written in Java, and etcd is written in Go. But the rough shape is indicative: seqcd required nearly an order of magnitude less code to implement the same features. That’s not because etcd’s code is bloated or bad — it’s simply what is required to handle the messy, non-deterministic real world. For seqcd, all of that messiness has been neatly encapsulated away and dealt with by Aeron Sequencer, so the only code left to write is neat, simple, easily testable, and deterministic. No thinking about concurrency or crash recovery or slow clients, etc.
Where the effort actually went
The entirety of the core key-value store implementation in seqcd is about 6.7K lines of code and took just a day or two to write. The bulk of seqcd’s code — some 22.6K more lines — is its etcd compatibility layer. Making a drop-in replacement meant speaking etcd’s exact dialect: gRPC over HTTP/2 (Claude and I hand-rolled a zero-allocation NIO HTTP/2 server for the hot path), watch streams, the full etcd v3 API surface, down to exactly-matched error strings like etcdserver: lease already exists. It passes a variety of existing test suites: jetcd, the Kubernetes integration and kubectl command tests, CNCF’s conformance suite, and a smoke test on a kind cluster. To all outward appearances, seqcd is etcd.
Important
Making a drop-in etcd replacement in two weeks would not have been practical if I hadn’t had the benefit of the many already-written etcd and Kubernetes test suites mentioned above. It took an enormous amount of specialized knowledge, effort, and hard-won production experience from many people over many years to produce those tests. Using Aeron Sequencer made it trivial to produce a correct distributed key-value store, but it was the test suites that allowed me to make that key-value store look and behave indistinguishably from etcd so quickly.
Aeron Sequencer ❤️ AI
AI coding agents are often described as “junior engineers.” You can even coax them into “senior engineer” territory with some effort. Unfortunately, just like human engineers, once they’ve chewed through a context window or two, AIs have a hard time consistently reasoning about concurrency and correctness in a large codebase.
Conveniently, using Aeron Sequencer takes away the need for all that pesky and error-prone thought. You just write plain, single-threaded, easily unit-testable, simple code and the Sequencer runtime takes care of the rest. AIs and humans are both pretty good at writing straightforward single-threaded code.
We have also taken care to make Aeron Sequencer’s documentation readable and ingestible by AI, precisely so that an engineer working with an AI assistant can quickly build correct applications on top of Sequencer. We even measure and validate the documentation’s AI-friendliness as part of our CI/CD pipeline.
So Aeron Sequencer allows AI to write less code to achieve the same goals, and it’s been designed to be easy for AI to understand and work with. Simpler code and less context to ingest mean faster development — and fewer tokens.
The performance icing on the Sequencer cake
To top it off, applications written on Aeron Sequencer are often faster and feature lower and more predictable latencies. Simple code is usually fast code. And Aeron Sequencer, and its underlying Aeron transport and SBE serialization format, have been meticulously hand-crafted by, without exaggeration, some of the most expert low-latency Java engineers in the world. As a result, seqcd is not only fully compatible with etcd… but it’s significantly faster, too. I benchmarked 3-node seqcd against 3-node etcd 3.6.12 on my 14-core M4 MacBook Pro. For throughput I used etcd’s own benchmark tool — same wire client, same 1 KB values, single endpoint per target.
The scenarios are etcd’s core read and write shapes. A prefix LIST is a multi-key range scan — the kube namespace read. A point read fetches a single key over that same Range RPC, either serializable (served from local state, no quorum round-trip) or linearizable (quorum-confirmed via ReadIndex). A PUT is a write, durable (fsync on every commit) or not. The chart plots all six; Kubernetes leans hardest on the two below:
Prefix LIST (the kube namespace read): 3.2× etcd. To be fair: seqcd serves range scans from an in-memory view while etcd walks boltdb MVCC per key — but here the working set is tiny (200 keys × 1 KB), so etcd’s boltdb reads hit the OS page cache, not disk.
Durable writes: 2.9× etcd at 256 in-flight, 5.6× at 1024. Aeron group-commit and fsync against etcd’s WAL fsync.
Latency matters even more than throughput for a Kubernetes cluster. The tail figures below come from focused closed- and open-loop load suites — 256 B values, reporting the weakest of three runs each — rather than the 1 KB benchmark runs above:
The long tail is where the architecture shows through most clearly. The tails hold up: p99.9 is 2.9× lower on durable writes, 2.2× on serializable reads, and 1.7× on linearizable reads:
scenario
p50
p99
p99.9
durable PUT (256 clients, closed loop)
seqcd
15.0 ms
21.1 ms
29.3 ms
etcd
29.0 ms
62.0 ms
83.5 ms
serializable point read (256 clients, closed loop)
seqcd
0.90 ms
2.6 ms
3.1 ms
etcd
1.2 ms
4.2 ms
6.7 ms
linearizable point read (1,024 in-flight, open loop)
seqcd
1.7 ms
7.3 ms
10.7 ms
etcd
5.4 ms
14.2 ms
18.6 ms
seqcd also passes the “large”-sized tier of etcd’s own host performance validation on my laptop, where etcd itself fails it (on some lucky runs, seqcd even passes the “xl”-sized workload):
Looks great! Gimme!
Whoa, hold up now! etcd has more than a decade of production hardening behind it: years of bug fixes driven by real-world incidents, fuzzing, and fault injection. seqcd is two weeks old. Passing etcd’s client and conformance suites means seqcd behaves like etcd under the conditions those suites exercise — it does not mean seqcd has survived the years of adversarial production traffic that etcd has. The replicated core inherits the maturity of Aeron Cluster, which has long been deployed in mission-critical production environments; the 22.6K-line compatibility layer wrapped around it is brand new. seqcd is also deliberately incomplete: anything Kubernetes doesn’t use — most notably etcd’s authentication and authorization layer, which Kubernetes bypasses in favor of mutual TLS — I simply didn’t build.
That said… if you’d like to build your own etcd replacement or anything else using Aeron Sequencer, give us a ring.
So what can be built with Aeron Sequencer in just two weeks?
Turns out: quite a lot. One engineer plus Claude, two weeks, a working drop-in replacement for one of the most demanding pieces of infrastructure software there is — and not because I’m a genius engineer (my brain is as silky-smooth as my erstwhile hair), but because Sequencer handled all the hard parts of replication, consistency, high availability, and crash-recoverability for me.
If your system is mission-critical and must survive machines dying and network outages and never give an inconsistent answer — an order book, a matching engine, a control plane, a passenger manifest for an actual plane… you can use Aeron Sequencer to turn what otherwise would be a years-long distributed-systems project into a tiny unit-testable core of business logic.
Benchmark methodology: 3-node clusters vs etcd 3.6.12 on a macOS 14-core host, client and cluster co-located (CPU oversubscription penalizes seqcd’s busy-spin threads; expect better on dedicated hardware). Throughput (the six-scenario chart) uses the official etcd benchmark driver, 1 KB values / 64 B keys, client on one endpoint, single run; its serializable point-read row sums two concurrent client processes on both sides, because one process saturates its own dispatch (~120K req/s) before either server’s capacity. Tail latency uses focused closed-loop (durable PUT and serializable point read; 256 clients over one connection) and open-loop (linearizable point read; 1,024 in-flight) suites at 256 B values / 64 B keys, reporting the weakest of three runs. Durable = fsync on every commit on both sides (seqcd fileSyncLevel=2, etcd’s default WAL fsync); non-durable = fileSyncLevel=0 vs --unsafe-no-fsync.
Eric Bowden Senior Platform Engineer Adaptive | Aeron
Eric is a Senior Platform Engineer at Adaptive Financial Consulting, working on Aeron, Aeron Cluster, and Aeron Sequencer — the open-source, low-latency messaging and clustering technology behind many of the world’s most demanding capital markets systems. He specialises in distributed systems, consensus, and high-availability architectures for mission-critical trading infrastructure. He combines deep engineering rigor with thought leadership on how AI can accelerate the delivery of correct, production-grade distributed systems.
Further reading
Product overview Aeron Sequencer
Built by world-class distributed systems engineers, it delivers the availability, consistency, low latency, and high throughput your trading systems need – so your teams can focus on the edge that drives your business.
Learn how leading firms can reduce operational risk, simplify distributed architectures, enable active-active deployment, and build resilient trading platforms that scale.
Blog Trading system development with AI: cost moves from writing code to trusting it
Trading system development with AI shifts cost from coding to trust. See how SMR and Aeron Sequencer reduce failure risk, data loss, and recovery time.