Skip to main content

yo_resp/dispatch/
repl.rs

1//! Being a master: the replication identity, the stream and the link.
2//!
3//! A replica is a second copy of this server's keyspace that keeps up by being
4//! told every change as it happens. Getting there is two halves. First it needs
5//! the dataset as it stands, which is a snapshot sent down the socket, and then
6//! it needs everything that happens after that snapshot, which is a stream of
7//! commands that never ends. The hard part is the seam between the two, because
8//! a change that lands in neither is lost forever and a change that lands in
9//! both is applied twice, and applying `INCR` twice is not the same as applying
10//! it once.
11//!
12//! # The handshake
13//!
14//! A replica opens an ordinary connection and sends `PING`, then
15//! `REPLCONF listening-port <port>`, then `REPLCONF capa ...`, and then `PSYNC
16//! <replid> <offset>`. Everything up to the `PSYNC` is a normal command with a
17//! normal reply. The `PSYNC` is where the connection stops being a client: the
18//! master answers `+FULLRESYNC <replid> <offset>`, then the snapshot as a bulk
19//! string with no trailing newline after it, and from then on writes only the
20//! command stream and never replies to anything the replica sends. A replica
21//! sends `REPLCONF ACK <offset>` about once a second forever, and a master that
22//! answers one of those with `+OK` has put a reply into the middle of a stream
23//! the replica is parsing as commands, which a real replica reports as
24//! `Protocol error (Master using the inline protocol. Desync?)` and then drops
25//! the link. So the silence is not an optimisation, it is the protocol.
26//!
27//! # The identity and the offset
28//!
29//! The replication id is forty hex characters naming this server's history, and
30//! the offset is how many bytes of stream have gone out under that id. Together
31//! they are a position in a history, which is what lets a replica that lost the
32//! link for a moment ask to carry on rather than start again: it sends back the
33//! id and offset it had, and if the id is still ours and the offset is still in
34//! the backlog the master answers `+CONTINUE` and replays the missing bytes.
35//!
36//! There is a second id for the same reason a chain needs one. When a replica is
37//! promoted it keeps the old master's id as `replid2` and takes a new one of its
38//! own, so the replicas that were following the old master can be handed over
39//! without all of them resyncing from nothing.
40//!
41//! # The snapshot has to be one instant
42//!
43//! `SAVE` here walks the databases one stripe at a time and does not stop the
44//! server, so the file it writes is not one instant of the keyspace: a write to
45//! database nine while database two is being walked is in the file and a write
46//! to database two after it has been walked is not. For a file on disk that is a
47//! fair trade, because the file is read by itself and nothing is going to be
48//! replayed on top of it.
49//!
50//! It is not a fair trade here. The replica loads the snapshot and then applies
51//! the stream from the offset the snapshot was stamped with, so every byte of
52//! the keyspace has to be either in the snapshot or after that offset, and never
53//! both and never neither. A real server gets that for free by forking, which
54//! hands the child an instant of the whole address space and costs the parent
55//! nothing but the page faults afterwards. There is no fork here, so a full
56//! resync takes every stripe of every database at once, stamps the offset,
57//! builds the image and lets go. The server stops for as long as that takes,
58//! which is a real cost on a large dataset and is registered as a divergence.
59//! The alternative is a replica that is quietly wrong, which is worse than a
60//! replica that took a pause to be right.
61//!
62//! # What goes on the stream
63//!
64//! Every command that changed something, in the form the replica has to be given
65//! rather than the form the client sent. Most commands are the same both ways.
66//! The ones that are not are the ones whose result depends on something the
67//! replica has not got: the clock, for anything that sets a deadline, and the
68//! server's own random state, for `SPOP`. `EXPIRE k 50` becomes `PEXPIREAT k
69//! <absolute>`, `SET k v EX 100` becomes `SET k v PXAT <absolute>`, `SPOP s`
70//! becomes `SREM s <the member that actually went>`, and `XADD s *` becomes an
71//! `XADD` naming the id that was actually made. Without the rewrite the two
72//! copies drift apart the moment either of them is asked a question.
73//!
74//! The rewrite is pushed by the command that knows, through a thread local, the
75//! same way this crate already hands keyspace events up from underneath. A
76//! command that pushes nothing is sent as it arrived, which is the common case
77//! and costs nothing to decide.
78//!
79//! # What it costs a server with no replica
80//!
81//! One relaxed load of a count that is zero, per write command. The same shape
82//! `MONITOR` and the pub/sub registry use, and for the same reason: nearly every
83//! server in the world is running on its own.
84
85use std::sync::Arc;
86use std::sync::atomic::Ordering::{Acquire, Relaxed, Release};
87use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, AtomicUsize};
88use std::time::Duration;
89
90use core::cell::{Cell, RefCell};
91
92use yo_common::lock::Lock;
93use yo_common::{Code, Error, Result};
94
95use crate::reply::Out;
96
97use super::args::{self, Args};
98use super::clients::{self, Client};
99use super::pubsub::Envelope;
100use super::{Server, Session};
101
102/// How many characters name a history, which is Redis's forty hex digits.
103pub(super) const ID_LEN: usize = 40;
104
105/// The id a server that has never been anybody's replica reports as its second,
106/// which Redis writes as forty zeroes rather than leaving out.
107const NO_ID: &[u8; ID_LEN] = b"0000000000000000000000000000000000000000";
108
109/// How much stream is kept for a replica that dropped the link, in bytes.
110///
111/// Redis's default, and the number `repl-backlog-size` reads back as. A replica
112/// that reconnects inside this many bytes carries on where it was and one that
113/// falls further behind starts again from a snapshot.
114pub(super) const BACKLOG_BYTES: usize = 1024 * 1024;
115
116/// One connection that has stopped being a client and is now a copy of us.
117pub(super) struct Replica {
118    /// The row every other part of the server knows the connection by, so that
119    /// `CLIENT LIST` and `CLIENT KILL` see a replica the way they see anything
120    /// else.
121    row: Arc<Client>,
122    /// The port the replica says it is listening on, which is what `INFO`
123    /// reports rather than the port its outgoing connection came from. A replica
124    /// that has not said is reported as zero, which is what Redis does too.
125    port: AtomicU64,
126    /// The last offset it said it had, and when it said so.
127    ///
128    /// `WAIT` counts the replicas whose acknowledged offset has caught up, and
129    /// `INFO` turns the time into a lag in seconds.
130    ack: AtomicU64,
131    ack_ms: AtomicU64,
132    /// Whether the snapshot has gone and the stream has started.
133    online: AtomicBool,
134}
135
136impl Replica {
137    /// The address the replica is reachable at, which is its address and the
138    /// port it told us rather than the port it dialled out from.
139    fn address(&self) -> (String, u64) {
140        let text = self.row.text.lock();
141        let peer = text.peer.clone();
142        drop(text);
143        let host = match peer.iter().rposition(|&b| b == b':') {
144            Some(at) => peer[..at].to_vec(),
145            None => peer.clone(),
146        };
147        (
148            String::from_utf8_lossy(&host).into_owned(),
149            self.port.load(Relaxed),
150        )
151    }
152}
153
154/// Everything about being a master, all of it idle on a server with no replica.
155pub(crate) struct Replication {
156    /// The forty characters naming the history this server is writing.
157    id: Lock<[u8; ID_LEN]>,
158    /// The history it was writing before, for a server that was promoted.
159    id2: Lock<[u8; ID_LEN]>,
160    /// How many bytes of stream have been produced under [`Replication::id`].
161    offset: AtomicU64,
162    /// The offset the second id runs up to, or minus one when there is no
163    /// second id, which is what Redis reports on a server that was never a
164    /// replica.
165    second: AtomicI64,
166    /// What has gone out lately, so a replica that blinked can catch up without
167    /// a whole snapshot.
168    backlog: Lock<Backlog>,
169    /// The connections being fed.
170    rows: Lock<Vec<Arc<Replica>>>,
171    /// How many there are, so a write command can ask without taking the lock.
172    live: AtomicUsize,
173    /// The port a connection said it was listening on, before it asked to
174    /// become a replica.
175    ///
176    /// `REPLCONF listening-port` arrives two commands before `PSYNC`, so there
177    /// is no replica to hang it on yet and it has to be kept somewhere until
178    /// there is. Here rather than on the connection row, because a row is paid
179    /// for by every connection on the server and this is paid for only by the
180    /// few that are about to become replicas. Oldest first out at a small cap,
181    /// so a client that sends the one and never the other cannot grow it.
182    ports: Lock<Vec<(u64, u64)>>,
183    /// Whether a full resync is building an image right now.
184    ///
185    /// A write that arrives while this is set is held and run again afterwards,
186    /// which is what makes the image one instant. Read on the command path by
187    /// every write, so it is a relaxed load of a bool that is nearly always
188    /// false, next to the one `CLIENT PAUSE` already costs.
189    frozen: AtomicBool,
190    /// One resync builds at a time.
191    ///
192    /// Two threads freezing the server at once would each wait for the other's
193    /// writes to drain and neither would be able to, so the second one waits
194    /// here instead and then finds a backlog it can probably be caught up from.
195    building: Lock<()>,
196    /// Which database the stream is on, so `SELECT` goes out only when it has
197    /// to. Minus one before anything has been written, which is why the first
198    /// command on any database is preceded by a `SELECT` even for database
199    /// zero.
200    on_db: AtomicI64,
201    /// When something last went out on the stream, so the periodic `PING` a
202    /// master owes its replicas is only sent on a link that has gone quiet.
203    last_io: AtomicU64,
204}
205
206impl Default for Replication {
207    fn default() -> Replication {
208        Replication {
209            id: Lock::new(make_id()),
210            id2: Lock::new(*NO_ID),
211            offset: AtomicU64::new(0),
212            second: AtomicI64::new(-1),
213            backlog: Lock::new(Backlog::default()),
214            rows: Lock::new(Vec::new()),
215            live: AtomicUsize::new(0),
216            ports: Lock::new(Vec::new()),
217            frozen: AtomicBool::new(false),
218            building: Lock::new(()),
219            on_db: AtomicI64::new(-1),
220            last_io: AtomicU64::new(0),
221        }
222    }
223}
224
225/// The last stretch of stream, kept so a reconnect does not cost a snapshot.
226///
227/// A ring of a fixed size with the offset of its first byte beside it. `histlen`
228/// is how much of the ring is real, which is less than its size only until it
229/// has filled once.
230#[derive(Default)]
231struct Backlog {
232    ring: Vec<u8>,
233    /// Where the next byte goes.
234    at: usize,
235    /// How many bytes of the ring are real.
236    filled: usize,
237    /// The stream offset of the oldest real byte.
238    first: u64,
239}
240
241impl Backlog {
242    /// Take bytes in, dropping whatever falls off the back.
243    fn push(&mut self, bytes: &[u8], upto: u64) {
244        if self.ring.is_empty() {
245            self.ring = vec![0; BACKLOG_BYTES];
246            self.first = upto - bytes.len() as u64;
247        }
248        for &b in bytes {
249            self.ring[self.at] = b;
250            self.at = (self.at + 1) % BACKLOG_BYTES;
251            if self.filled < BACKLOG_BYTES {
252                self.filled += 1;
253            }
254        }
255        self.first = upto - self.filled as u64;
256    }
257
258    /// Everything from `from` onwards, or `None` if that much history has gone.
259    fn since(&self, from: u64, upto: u64) -> Option<Vec<u8>> {
260        if self.filled == 0 || from < self.first || from > upto {
261            return None;
262        }
263        let skip = (from - self.first) as usize;
264        let want = self.filled - skip;
265        let start = (self.at + BACKLOG_BYTES - self.filled + skip) % BACKLOG_BYTES;
266        let mut out = Vec::with_capacity(want);
267        for i in 0..want {
268            out.push(self.ring[(start + i) % BACKLOG_BYTES]);
269        }
270        Some(out)
271    }
272}
273
274#[cfg(test)]
275impl Server {
276    /// Say there is a replica, without one.
277    ///
278    /// What a test wants to look at is the byte stream and not the socket it
279    /// would have gone down, and [`emit`] writes the stream into the backlog
280    /// before it looks for anybody to post it to. So a server told this writes
281    /// everything a real master would write and posts none of it, which is the
282    /// whole of what these tests need and needs no connection at all.
283    pub(super) fn pretend_replica(&self) {
284        self.repl.live.store(1, Relaxed);
285    }
286
287    /// The stream from `from` onwards, as text, and where it has got to.
288    ///
289    /// Text because every rewrite these tests are about is text, and reading
290    /// `*3\r\n$3\r\nSET\r\n` in a failure message beats reading a byte array.
291    pub(super) fn stream_since(&self, from: u64) -> (String, u64) {
292        let upto = self.repl.offset.load(Acquire);
293        let backlog = self.repl.backlog.lock();
294        let bytes = backlog.since(from, upto).unwrap_or_default();
295        (String::from_utf8_lossy(&bytes).into_owned(), upto)
296    }
297}
298
299/// Whether a replica is the one at this host and port.
300///
301/// The host is compared as text and not as an address, which is Redis's own
302/// rule: what an operator types has to be what the server prints in `INFO`,
303/// because that is where they read it from.
304fn at(held: &Replica, host: &str, port: u16) -> bool {
305    let (theirs, their_port) = held.address();
306    theirs == host && their_port == u64::from(port)
307}
308
309/// Forty hex characters from the operating system's generator.
310///
311/// From the system and not from the engine's own seeded one, because two servers
312/// started from the same image at the same moment must not be able to claim the
313/// same history. That is the same reasoning `ACL GENPASS` is built on.
314fn make_id() -> [u8; ID_LEN] {
315    const HEX: &[u8; 16] = b"0123456789abcdef";
316    let mut raw = [0u8; ID_LEN / 2];
317    yo_common::entropy::fill(&mut raw);
318    let mut id = [0u8; ID_LEN];
319    for (i, byte) in raw.iter().enumerate() {
320        id[i * 2] = HEX[usize::from(byte >> 4)];
321        id[i * 2 + 1] = HEX[usize::from(byte & 15)];
322    }
323    id
324}
325
326// ------------------------------------------------------------- the override
327
328thread_local! {
329    /// What the running command wants sent instead of itself, if anything.
330    ///
331    /// A command that leaves this empty is propagated as it arrived. One that
332    /// puts something here is propagated as whatever it put, however many
333    /// commands that is, and a command that puts an empty list here is
334    /// propagated as nothing at all. See the module header for why a thread
335    /// local rather than a return value.
336    static INSTEAD: RefCell<Option<Vec<Vec<Vec<u8>>>>> = const { RefCell::new(None) };
337}
338
339thread_local! {
340    /// Whether anything this thread does has somewhere to be copied to.
341    ///
342    /// A flag rather than a question put to the server, because the call sites
343    /// that ask are inside the command bodies, a long way from anything holding
344    /// a `Server`, and threading one down to them for a question that is nearly
345    /// always no would be a worse trade than a thread local read. The same trade
346    /// the keyspace events make, next door.
347    ///
348    /// It says whether there is a replica and not whether a write is running,
349    /// which are two different questions and this is the wider one on purpose. A
350    /// read takes keys away, because a lookup that finds one past its deadline
351    /// removes it there and then, and `XGROUP CREATE` changes a stream while
352    /// carrying no write flag, since the flag is on the subcommand and there is
353    /// no subcommand table yet. Both have something to send and neither is a
354    /// write by the funnel's reckoning, so the funnel decides whether the
355    /// verbatim command may be sent and this decides whether anybody is
356    /// listening at all.
357    static ARMED: Cell<bool> = const { Cell::new(false) };
358
359    /// The removals nobody asked for, waiting for a lock that is not held.
360    static REAPED: RefCell<Vec<Vec<Vec<u8>>>> = const { RefCell::new(Vec::new()) };
361}
362
363/// Say whether what runs next is going to be copied, answering what the last
364/// answer was.
365///
366/// Set in front of every write on a server with a replica and put back
367/// afterwards, so that a command run by a script or by `EXEC` leaves the flag
368/// the way it found it.
369pub(super) fn arm(on: bool) -> bool {
370    ARMED.replace(on)
371}
372
373/// Whether a rewrite would be heard, which is what keeps every call site that
374/// would build one off the path of a server with no replica.
375#[must_use]
376pub(crate) fn armed() -> bool {
377    ARMED.get()
378}
379
380/// Send this instead of the command that is running.
381///
382/// Called by a command whose own arguments do not say what happened, so more
383/// than once for a command with more than one effect. The first call replaces
384/// the command and the rest are added after it.
385pub(crate) fn instead(parts: Vec<Vec<u8>>) {
386    INSTEAD.with(|cell| {
387        let mut held = cell.borrow_mut();
388        held.get_or_insert_with(Vec::new).push(parts);
389    });
390}
391
392/// The same thing said in pieces, which is how nearly every call site has it.
393///
394/// A rewrite is a command name and a key and usually a number, none of which
395/// arrive owned, so this is the shape that keeps the copying in one place rather
396/// than a `to_vec` on every argument of every site.
397pub(crate) fn rewrite(parts: &[&[u8]]) {
398    instead(parts.iter().map(|part| part.to_vec()).collect());
399}
400
401/// Send nothing at all for the command that is running.
402///
403/// For a write command that turned out to change nothing and whose verbatim form
404/// would be wrong rather than merely wasteful, which is any of the ones that are
405/// rewritten: a `SPOP` on a missing key has no `SREM` to send and must not send
406/// the `SPOP`.
407pub(crate) fn nothing() {
408    INSTEAD.with(|cell| {
409        let mut held = cell.borrow_mut();
410        held.get_or_insert_with(Vec::new);
411    });
412}
413
414/// Take whatever the command left, clearing it for the next one.
415fn taken() -> Option<Vec<Vec<Vec<u8>>>> {
416    INSTEAD.with(|cell| cell.borrow_mut().take())
417}
418
419// --------------------------------------------------------------- the server
420
421impl Server {
422    /// Whether anybody is being fed, which is the whole cost of this file on a
423    /// server that is on its own.
424    #[must_use]
425    pub(crate) fn replicated(&self) -> bool {
426        self.repl.live.load(Relaxed) != 0
427    }
428
429    /// How many of them there are, which `INFO clients` takes off its count.
430    ///
431    /// A replica is a connection and is not a client: Redis counts the sockets
432    /// and then subtracts the replicas, so `connected_clients` is the number of
433    /// people talking to the server rather than the number of file descriptors
434    /// it is holding, and the replicas are reported next door as their own
435    /// number.
436    #[must_use]
437    pub(crate) fn replica_count(&self) -> u64 {
438        self.repl.live.load(Relaxed) as u64
439    }
440
441    /// The forty characters naming this server's history.
442    pub(crate) fn repl_id(&self) -> [u8; ID_LEN] {
443        *self.repl.id.lock()
444    }
445
446    /// How many bytes of stream have gone out.
447    #[must_use]
448    pub(crate) fn repl_offset(&self) -> u64 {
449        self.repl.offset.load(Acquire)
450    }
451
452    /// Whether a replica is attached at this address, and whether it is past its
453    /// snapshot, which are the two questions `FAILOVER TO` asks in that order.
454    ///
455    /// `None` is nobody there and is a different answer from `Some(false)`,
456    /// which is somebody there who is still loading. The address is matched the
457    /// way Redis matches it, against where the connection came from and against
458    /// the port the replica said it listens on rather than the port it dialled
459    /// out from, because the second one is not something an operator can type.
460    #[must_use]
461    pub(super) fn replica_online_at(&self, host: &str, port: u16) -> Option<bool> {
462        self.replica_rows()
463            .iter()
464            .find(|held| at(held, host, port))
465            .map(|held| held.online.load(Relaxed))
466    }
467
468    /// Whether the replica at this address has acknowledged every byte written.
469    ///
470    /// The wait a failover is, put as one question. A replica that has gone away
471    /// answers no rather than yes, which is what keeps a failover to a target
472    /// that died waiting rather than handing the job to nobody.
473    #[must_use]
474    pub(super) fn replica_caught_up(&self, host: &str, port: u16) -> bool {
475        let upto = self.repl_offset();
476        self.replica_rows()
477            .iter()
478            .any(|held| at(held, host, port) && held.ack.load(Relaxed) >= upto)
479    }
480
481    /// Where the first replica that has everything is, for a failover that was
482    /// not told which one to hand the job to.
483    ///
484    /// In the order they attached, which is the order `INFO` lists them in, so
485    /// two operators reading the same server pick the same one.
486    #[must_use]
487    pub(super) fn first_caught_up(&self) -> Option<(String, u16)> {
488        let upto = self.repl_offset();
489        yo_alloc::allow(|| {
490            self.replica_rows()
491                .iter()
492                .filter(|held| held.online.load(Relaxed) && held.ack.load(Relaxed) >= upto)
493                .map(|held| held.address())
494                .find_map(|(host, port)| u16::try_from(port).ok().map(|port| (host, port)))
495        })
496    }
497
498    /// A handle to every replica, copied out so the bytes can be posted with the
499    /// lock let go of.
500    ///
501    /// The same trade `CLIENT LIST` and the monitor feed make, for the same
502    /// reason: a replica that drops while a command is being rendered leaves a
503    /// row that is still readable, and the bytes land in a mailbox that is about
504    /// to be told the slot has moved on.
505    fn replica_rows(&self) -> Vec<Arc<Replica>> {
506        let rows = self.repl.rows.lock();
507        yo_alloc::allow(|| rows.clone())
508    }
509
510    /// Remember the port a connection says it is listening on.
511    ///
512    /// See [`Replication::ports`] for why it cannot go straight on a replica.
513    fn note_replica_port(&self, id: u64, port: u64) {
514        const KEEP: usize = 64;
515        let mut ports = self.repl.ports.lock();
516        yo_alloc::allow(|| {
517            if let Some(row) = ports.iter_mut().find(|(who, _)| *who == id) {
518                row.1 = port;
519                return;
520            }
521            if ports.len() >= KEEP {
522                ports.remove(0);
523            }
524            ports.push((id, port));
525        });
526    }
527
528    /// Take the port back out, for a connection that has got as far as `PSYNC`.
529    fn take_replica_port(&self, id: u64) -> u64 {
530        let mut ports = self.repl.ports.lock();
531        match ports.iter().position(|(who, _)| *who == id) {
532            Some(at) => ports.remove(at).1,
533            None => 0,
534        }
535    }
536
537    /// Take a connection on as a replica.
538    fn take_replica(&self, row: &Arc<Client>) -> Arc<Replica> {
539        let mut rows = self.repl.rows.lock();
540        let held = yo_alloc::allow(|| {
541            let held = Arc::new(Replica {
542                row: Arc::clone(row),
543                port: AtomicU64::new(self.take_replica_port(row.id)),
544                ack: AtomicU64::new(0),
545                ack_ms: AtomicU64::new(self.clock.now_ms()),
546                online: AtomicBool::new(true),
547            });
548            rows.push(Arc::clone(&held));
549            held
550        });
551        self.repl.live.store(rows.len(), Release);
552        row.set_flag(clients::REPLICA, true);
553        self.note_here(row.thread.load(Relaxed), 1);
554        held
555    }
556
557    /// Let one go, which is the connection closing or being killed.
558    pub(crate) fn drop_replica(&self, id: u64) {
559        let mut rows = self.repl.rows.lock();
560        let Some(at) = rows.iter().position(|r| r.row.id == id) else {
561            return;
562        };
563        let gone = rows.remove(at);
564        self.repl.live.store(rows.len(), Release);
565        drop(rows);
566        gone.row.set_flag(clients::REPLICA, false);
567        self.note_here(gone.row.thread.load(Relaxed), -1);
568    }
569
570    /// Whether a full resync is holding the keyspace still.
571    ///
572    /// One relaxed load per write command on a server nobody is syncing from,
573    /// which is the same shape and the same cost as the pause check it sits
574    /// beside.
575    #[must_use]
576    pub(crate) fn frozen(&self) -> bool {
577        self.repl.frozen.load(Relaxed)
578    }
579
580    /// The whole dataset as one image, and the offset it is an image as of.
581    ///
582    /// Everything between the two has to be nothing at all, which is what the
583    /// freeze is for. Writes are held from before the barrier until after the
584    /// image is finished, so a change is either inside the image or after the
585    /// offset and never both. See the module header for why there is no fork to
586    /// get this for free.
587    pub(crate) fn snapshot_at_an_instant(&self) -> (Vec<u8>, u64) {
588        self.at_an_instant(|| super::persist::build(self).0)
589    }
590
591    /// Run `f` with nothing writing, and say which offset it ran at.
592    ///
593    /// The freeze is the whole of how this server gets what a real one gets from
594    /// a fork: whatever `f` reads is one moment of the dataset, and the offset
595    /// that comes back is the point in the stream that moment sits at, so a
596    /// caller feeding a replica knows exactly where to carry on from.
597    ///
598    /// Every write is held off for as long as `f` runs, so `f` has to be worth
599    /// it. Building a whole image is; anything that waits on a client is not.
600    pub(crate) fn at_an_instant<T>(&self, f: impl FnOnce() -> T) -> (T, u64) {
601        let building = self.repl.building.lock();
602        self.repl.frozen.store(true, Release);
603        // A write already running holds the stripe it is writing to, so taking
604        // every stripe once and letting it go again is a barrier: once it is
605        // through, no write is in flight and the freeze above stops any more
606        // starting. One pass is enough, in whatever order, because nothing being
607        // waited for can start again behind it.
608        for db in &self.dbs {
609            for stripe in 0..db.width() {
610                drop(db.hold_stripe(stripe));
611            }
612        }
613        let offset = self.repl.offset.load(Acquire);
614        let made = f();
615        self.repl.frozen.store(false, Release);
616        drop(building);
617        (made, offset)
618    }
619
620    /// Take a master's history as our own, which is what a replica does.
621    ///
622    /// Everything this server writes from here on is under the master's id and
623    /// at the master's offsets, so a sub-replica underneath is handed positions
624    /// its own master would recognise. The backlog goes with it, because what is
625    /// in it is a stretch of a history this server is no longer writing and
626    /// answering a partial resync out of it would send a replica bytes from
627    /// somebody else's stream.
628    pub(crate) fn adopt(&self, id: [u8; ID_LEN], offset: u64) {
629        let mut ours = self.repl.id.lock();
630        if *ours == id && self.repl.offset.load(Acquire) == offset {
631            return;
632        }
633        *ours = id;
634        drop(ours);
635        let mut backlog = self.repl.backlog.lock();
636        *backlog = Backlog::default();
637        self.repl.offset.store(offset, Release);
638        drop(backlog);
639        self.repl.on_db.store(-1, Relaxed);
640    }
641
642    /// Stop following and start a history of our own, keeping the old one.
643    ///
644    /// The id this server was writing under becomes the second id and the offset
645    /// it had reached becomes the second offset, and a new id is taken. That is
646    /// what lets the replicas that were following the same master be handed over
647    /// to this one without every one of them starting from a snapshot: each of
648    /// them asks about a history this server can still say it was part of, up to
649    /// the point where the two part, and the second offset is that point.
650    pub(crate) fn promote(&self) {
651        let mut id = self.repl.id.lock();
652        let mut id2 = self.repl.id2.lock();
653        *id2 = *id;
654        *id = make_id();
655        // One past the last byte of the old history, which is the first byte
656        // that is only ours, because a replica asks about the first byte it
657        // wants and counts from one.
658        self.repl
659            .second
660            .store(self.repl.offset.load(Acquire) as i64 + 1, Relaxed);
661    }
662
663    /// Take a fresh replication id and forget the old one, which is what
664    /// `DEBUG CHANGE-REPL-ID` is for.
665    ///
666    /// Not a promotion. A promotion keeps the old id as the second one so that
667    /// the replicas that shared it can carry on, and this throws it away, which
668    /// is the point: it is how a test makes two servers that were part of the
669    /// same history stop being able to prove it, so the next `PSYNC` between
670    /// them has to be a full one.
671    pub(crate) fn change_id(&self) {
672        let mut id = self.repl.id.lock();
673        let mut id2 = self.repl.id2.lock();
674        *id = make_id();
675        *id2 = [b'0'; ID_LEN];
676        self.repl.second.store(-1, Relaxed);
677    }
678
679    /// The replica row for a connection, if it is one.
680    fn replica_of(&self, id: u64) -> Option<Arc<Replica>> {
681        let rows = self.repl.rows.lock();
682        rows.iter().find(|r| r.row.id == id).map(Arc::clone)
683    }
684}
685
686// ------------------------------------------------------------ the heartbeat
687
688/// How long a quiet replication link goes before its master pings it.
689///
690/// Ten seconds is `repl-ping-replica-period`'s default. The setting itself is
691/// not in the config table yet, so this is the number rather than a lookup.
692const PING_REPLICAS_MS: u64 = 10_000;
693
694/// How often the heartbeat thread wakes up to see whether one is due.
695const HEARTBEAT_MS: u64 = 1000;
696
697impl Server {
698    /// Start the thread that pings the replicas of a quiet master.
699    ///
700    /// A master pings so that a replica can tell a link that has gone away from
701    /// a master with nothing to say, since a TCP connection that nobody writes
702    /// to gives no sign either way. It is also what moves the offset on a server
703    /// nobody is writing to, and a replica whose offset is still zero is one a
704    /// real server leaves out of `CLUSTER SLOTS` on the grounds that it has
705    /// never caught up with anything.
706    ///
707    /// Its own thread because the housekeeping the engine does runs per batch,
708    /// so a server with no traffic does none of it, and a server with no traffic
709    /// is exactly the one that owes its replicas a ping.
710    pub fn start_replica_heartbeat(self: &Arc<Server>) {
711        let beating = Arc::clone(self);
712        yo_alloc::allow(|| {
713            let _ = std::thread::Builder::new()
714                .name(String::from("yo-repl-ping"))
715                .spawn(move || heartbeat(&beating));
716        });
717    }
718}
719
720/// One ping every `PING_REPLICAS_MS` on a link that has been quiet that long.
721fn heartbeat(server: &Arc<Server>) {
722    loop {
723        std::thread::sleep(Duration::from_millis(HEARTBEAT_MS));
724        // Only a top level master. A replica passes on the stream it is given
725        // rather than writing one, so a ping from here would put bytes in the
726        // middle of somebody else's history.
727        if !server.replicated() || server.following() {
728            continue;
729        }
730        let now = server.now_ms();
731        if now.saturating_sub(server.repl.last_io.load(Relaxed)) < PING_REPLICAS_MS {
732            continue;
733        }
734        // No `SELECT` in front of it, which is what the reference's `dictid` of
735        // minus one means: a ping touches no key, so it does not move the
736        // stream's database and must not look as though it did.
737        emit(server, render(&[b"PING"]));
738    }
739}
740
741// ---------------------------------------------------------------- the feed
742
743/// One command as a RESP array, which is the only shape the stream has.
744fn render(parts: &[&[u8]]) -> Vec<u8> {
745    let mut out = Vec::with_capacity(16 + parts.iter().map(|p| p.len() + 16).sum::<usize>());
746    out.extend_from_slice(b"*");
747    out.extend_from_slice(parts.len().to_string().as_bytes());
748    out.extend_from_slice(b"\r\n");
749    for part in parts {
750        out.extend_from_slice(b"$");
751        out.extend_from_slice(part.len().to_string().as_bytes());
752        out.extend_from_slice(b"\r\n");
753        out.extend_from_slice(part);
754        out.extend_from_slice(b"\r\n");
755    }
756    out
757}
758
759/// Put bytes on the stream: into the backlog, onto the offset, out to everybody.
760///
761/// The offset moves under the backlog lock so that the backlog and the number
762/// never disagree. A replica attaching between the two would otherwise be given
763/// an offset the backlog cannot honour.
764fn emit(server: &Server, bytes: Vec<u8>) {
765    server.repl.last_io.store(server.now_ms(), Relaxed);
766    let shared = {
767        let mut backlog = server.repl.backlog.lock();
768        let upto = server.repl.offset.load(Relaxed) + bytes.len() as u64;
769        yo_alloc::allow(|| backlog.push(&bytes, upto));
770        server.repl.offset.store(upto, Release);
771        yo_alloc::allow(|| Arc::new(bytes))
772    };
773    for held in server.replica_rows() {
774        if !held.online.load(Relaxed) {
775            continue;
776        }
777        server.post(
778            held.row.thread.load(Relaxed),
779            Envelope::raw(
780                held.row.conn.load(Relaxed),
781                held.row.id,
782                Arc::clone(&shared),
783            ),
784        );
785    }
786}
787
788/// Put one command on the stream, with a `SELECT` in front of it if the stream
789/// is not on the right database.
790fn send(server: &Server, db: usize, parts: &[&[u8]]) {
791    let body = render(parts);
792    // A slot migration listens here as well, because this is the one place
793    // everything propagated comes through: a rewrite, a key that went away on
794    // its own, what a script did. It takes the command already rendered and
795    // keeps or drops it on the slot its key is in, and it never gets the
796    // `SELECT` below, since a cluster has one database and the far side would
797    // have nothing to do with it.
798    server.asm_feed(&body);
799    let mut bytes = Vec::new();
800    if server.repl.on_db.load(Relaxed) != db as i64 {
801        let n = db.to_string();
802        bytes = render(&[b"SELECT", n.as_bytes()]);
803        server.repl.on_db.store(db as i64, Relaxed);
804    }
805    bytes.extend_from_slice(&body);
806    emit(server, bytes);
807}
808
809/// Pass a master's bytes on, unchanged, and count them.
810///
811/// What a replica owes anybody following it is the stream it was given rather
812/// than an account of what running it did, because the offsets in that stream
813/// are positions in the master's history and this server is not writing one. So
814/// the bytes go on as they arrived. The count moves either way, because it is
815/// the number this server acknowledges to its own master and a chain of nobody
816/// still has to be able to say where it has got to.
817pub(crate) fn relayed(server: &Server, bytes: Vec<u8>, db: usize) {
818    if server.replicated() {
819        server.repl.on_db.store(db as i64, Relaxed);
820        emit(server, bytes);
821        return;
822    }
823    server.repl.offset.fetch_add(bytes.len() as u64, Release);
824}
825
826/// Report a command to every replica, in the form it has to be given.
827///
828/// The caller has already checked [`Server::replicated`] and that the command
829/// did not fail, so what is decided here is only the shape.
830///
831/// `verbatim` is whether the command as it arrived is a fair thing to send when
832/// the body said nothing. It is the write flag, and it is false for a read,
833/// which sends nothing, and for a container like `XGROUP` whose flags are on its
834/// subcommands, which sends what its body pushed and nothing otherwise.
835pub(super) fn feed(server: &Server, db: usize, args: Args<'_>, verbatim: bool) {
836    if super::follow::applying() {
837        forget();
838        return;
839    }
840    let instead = taken();
841    yo_alloc::allow(|| match instead {
842        None if !verbatim => {}
843        None => {
844            let parts: Vec<&[u8]> = (0..args.len()).map(|i| args.get(i)).collect();
845            send(server, db, &parts);
846        }
847        Some(each) => {
848            for one in &each {
849                let parts: Vec<&[u8]> = one.iter().map(Vec::as_slice).collect();
850                send(server, db, &parts);
851            }
852        }
853    });
854}
855
856/// Report what a command that nobody is running left behind.
857///
858/// A blocked client is answered by the thread that woke it, a long way outside
859/// the funnel and with no arguments to fall back on, so unlike [`feed`] there is
860/// no verbatim form here: whatever the answering left is all there is, and an
861/// answering that left nothing did nothing.
862pub(super) fn served(server: &Server, db: usize) {
863    if super::follow::applying() {
864        forget();
865        return;
866    }
867    let Some(each) = taken() else {
868        return;
869    };
870    yo_alloc::allow(|| {
871        for one in &each {
872            let parts: Vec<&[u8]> = one.iter().map(Vec::as_slice).collect();
873            send(server, db, &parts);
874        }
875    });
876}
877
878/// Throw away whatever the last command left, for a command that is not
879/// propagated at all.
880///
881/// A read that pushed a rewrite is a bug, but a rewrite left behind by a command
882/// that was refused after it pushed one is not, and it must not be handed to the
883/// next command that runs on this thread.
884pub(super) fn forget() {
885    let _ = taken();
886}
887
888/// Note a key or a field the storage layer took away on its own.
889///
890/// A replica never expires and never evicts. It cannot: the two decisions are
891/// made from a clock reading and a memory figure that are the master's and not
892/// its own, and a replica that made them itself would answer a read differently
893/// from the master for as long as the two disagreed. So it holds a key past its
894/// deadline until it is told, and being told is this. Redis sends the same thing
895/// for the same reason, and it is why a replica's `dbsize` can be ahead of the
896/// master's for a moment and never behind it.
897///
898/// Collected rather than sent, because the point at which a key goes is inside a
899/// stripe lock and the send takes a different one.
900pub(crate) fn reaped(parts: &[&[u8]]) {
901    if !ARMED.get() {
902        return;
903    }
904    REAPED.with_borrow_mut(|list| {
905        yo_alloc::allow(|| list.push(parts.iter().map(|part| part.to_vec()).collect()));
906    });
907}
908
909/// Send what [`reaped`] collected, ahead of whatever the command itself did.
910///
911/// Ahead, because that is the order it happened in and the order matters: a
912/// `SET k v` that found an expired `k` on the way in has to reach a replica as
913/// the deletion and then the write, or a replica that has kept a `k` of the
914/// wrong type refuses the write it is sent.
915pub(super) fn swept(server: &Server, db: usize) {
916    if REAPED.with_borrow(Vec::is_empty) {
917        return;
918    }
919    let each = REAPED.with_borrow_mut(core::mem::take);
920    // A replica that reaped a key on its own has nothing to tell anybody. The
921    // master is going to send the deletion in a moment and that is the copy that
922    // counts, because it is the one at the offset everybody else is counting
923    // from. See the `follow` module on why a chain passes bytes on rather than
924    // effects.
925    if super::follow::applying() {
926        return;
927    }
928    yo_alloc::allow(|| {
929        for one in &each {
930            let parts: Vec<&[u8]> = one.iter().map(Vec::as_slice).collect();
931            send(server, db, &parts);
932        }
933    });
934}
935
936// ------------------------------------------------------------- the commands
937
938/// `REPLCONF`, which is how a replica tells a master about itself.
939///
940/// Everything here is a pair, and a master answers `OK` to a pair it does not
941/// know rather than refusing it, because the whole point of the command is that
942/// a newer replica can tell an older master things it has never heard of. The
943/// one exception is `ACK`, which is not a question and is answered with silence:
944/// by the time a replica sends one the connection is a stream, and a reply on it
945/// is a protocol error at the other end.
946pub(super) fn replconf(
947    server: &Server,
948    session: &Session,
949    args: Args<'_>,
950    out: &mut Out,
951) -> Result<()> {
952    if args.len() < 3 || args.len().is_multiple_of(2) {
953        return Err(args::wrong_arity("replconf"));
954    }
955    let mut at = 1;
956    while at < args.len() {
957        let name = args.get(at).to_ascii_lowercase();
958        match name.as_slice() {
959            b"listening-port" => {
960                let port = args.int(at + 1)?;
961                server.note_replica_port(session.row().id, port.max(0) as u64);
962            }
963            b"ack" => {
964                // Not answered, ever. See the doc comment.
965                let ack = args.int(at + 1).unwrap_or(0);
966                if let Some(held) = server.replica_of(session.row().id) {
967                    held.ack.store(ack.max(0) as u64, Relaxed);
968                    held.ack_ms.store(server.clock.now_ms(), Relaxed);
969                }
970                return Ok(());
971            }
972            b"getack" => {
973                // A master sends this and a replica answers it. Reaching it here
974                // means somebody sent it to a master, which Redis refuses.
975                return Err(Error::new(
976                    Code::Invalid,
977                    "REPLCONF GETACK is only supported by a replica",
978                ));
979            }
980            _ => {}
981        }
982        at += 2;
983    }
984    out.ok();
985    Ok(())
986}
987
988/// `PSYNC <replid> <offset>`, which is where a connection stops being a client.
989///
990/// Answers `+CONTINUE` and the missing bytes when the replica is asking to carry
991/// on a history we are still writing and still hold enough of, and a full resync
992/// otherwise. `SYNC` is the same thing without the negotiation, from before
993/// Redis had one, and is a snapshot with no header in front of it.
994pub(super) fn psync(
995    server: &Server,
996    session: &mut Session,
997    args: Args<'_>,
998    out: &mut Out,
999) -> Result<()> {
1000    let partial = if args.name().eq_ignore_ascii_case(b"sync") {
1001        if args.len() != 1 {
1002            return Err(args::wrong_arity("sync"));
1003        }
1004        None
1005    } else {
1006        // Four words when the fourth is `FAILOVER`, which is a master that has
1007        // stopped writing telling us to take over from it. Three otherwise, and
1008        // a fourth word that is not that one is not a spelling of anything.
1009        let handover = args.len() == 4 && args.get(3).eq_ignore_ascii_case(b"failover");
1010        if args.len() != 3 && !handover {
1011            return Err(args::wrong_arity("psync"));
1012        }
1013        let asked = args.get(1);
1014        // Before anything else is read, because this is the one form of the
1015        // command that changes what this server is rather than asking it for
1016        // something, and a master that has already stopped writing is waiting on
1017        // the answer.
1018        if handover {
1019            take_over(server, asked)?;
1020        }
1021        let from = args.int(2)?;
1022        // The number a replica sends is the position of the first byte it
1023        // wants counted from one, so a replica that has everything asks for one
1024        // past the end of the stream. Everything on this side counts bytes
1025        // written, so the two are a step apart, and the step is taken here
1026        // rather than inside the backlog, where every other caller means the
1027        // count. A nought is nobody's answer and would mean the byte before the
1028        // first, so it starts again, which is what a real master does with it.
1029        (asked != b"?" && from >= 1).then(|| (asked.to_vec(), from as u64 - 1))
1030    };
1031    // Nobody may attach to a server that is in the middle of handing its job
1032    // over, because the offsets it is about to start reporting are somebody
1033    // else's. The `FAILOVER` form above got past this, which is the point: that
1034    // one is the handover and this is everybody else.
1035    if server.failing_over() {
1036        out.error(b"NOMASTERLINK Can't SYNC while failing over");
1037        return Ok(());
1038    }
1039    // Nor to a replica whose own link is not up, because what it would be
1040    // handing out is a history it has not finished reading.
1041    if server.following() && !server.master_link_up() {
1042        out.error(b"NOMASTERLINK Can't SYNC while not connected with my master");
1043        return Ok(());
1044    }
1045    if session.running() {
1046        return Err(Error::new(
1047            Code::Invalid,
1048            "PSYNC isn't allowed for DENY BLOCKING client",
1049        ));
1050    }
1051    // A partial resync first, since it is the cheap answer and the whole reason
1052    // the replica bothered to remember where it was.
1053    if let Some((asked, from)) = partial
1054        && let Some(bytes) = catch_up(server, &asked, from)
1055    {
1056        {
1057            let id = server.repl_id();
1058            out.raw(b"+CONTINUE ");
1059            out.raw(&id);
1060            out.raw(b"\r\n");
1061            out.raw(&bytes);
1062            attach(server, session);
1063            return Ok(());
1064        }
1065    }
1066    full(server, session, out)
1067}
1068
1069/// Take the master's job, because the master asked us to.
1070///
1071/// The fourth word of a `PSYNC` is how a `FAILOVER` finishes. Everything hard
1072/// about it happened on the other side: the master stopped writing, waited until
1073/// this server had acknowledged every byte, and only then sent this. So all that
1074/// is left here is to check that it is talking to the right server and then stop
1075/// being a replica, after which the `PSYNC` carries on as an ordinary one and is
1076/// answered with a `+CONTINUE` that costs nothing, because the two sides are in
1077/// step by construction.
1078///
1079/// The promotion is `REPLICAOF NO ONE`, so the history that was being followed
1080/// is kept as the second id. That is what makes the old master's `PSYNC`, which
1081/// names that very history, something this server can answer at all.
1082fn take_over(server: &Server, asked: &[u8]) -> Result<()> {
1083    if !server.following() {
1084        return Err(Error::new(
1085            Code::Invalid,
1086            "PSYNC FAILOVER can't be sent to a master.",
1087        ));
1088    }
1089    if asked != server.repl_id() {
1090        return Err(Error::new(
1091            Code::Invalid,
1092            "PSYNC FAILOVER replid must match my replid.",
1093        ));
1094    }
1095    let Some(shared) = server.myself() else {
1096        return Err(Error::new(
1097            Code::Invalid,
1098            "PSYNC FAILOVER is not available on an embedded server",
1099        ));
1100    };
1101    shared.stop_following();
1102    Ok(())
1103}
1104
1105/// The bytes a reconnecting replica missed, or `None` if it has to start again.
1106///
1107/// A replica asking about the history we are writing is caught up from the
1108/// backlog. One asking about the history we were writing before we were promoted
1109/// is caught up too, but only up to the point where the histories part, which is
1110/// what the second offset records.
1111fn catch_up(server: &Server, asked: &[u8], from: u64) -> Option<Vec<u8>> {
1112    let ours = server.repl_id();
1113    let theirs = *server.repl.id2.lock();
1114    let second = server.repl.second.load(Relaxed);
1115    let matches = asked == ours
1116        || (asked == theirs && second >= 0 && from <= u64::try_from(second).unwrap_or(0));
1117    if !matches {
1118        return None;
1119    }
1120    let upto = server.repl.offset.load(Acquire);
1121    let backlog = server.repl.backlog.lock();
1122    yo_alloc::allow(|| backlog.since(from, upto))
1123}
1124
1125/// A full resync: the header, then the whole keyspace as one image.
1126///
1127/// The image and the offset are taken with every stripe of every database held,
1128/// so what the replica loads and what it is then told about are the two halves
1129/// of one instant and neither overlaps the other. See the module header.
1130fn full(server: &Server, session: &mut Session, out: &mut Out) -> Result<()> {
1131    let (image, offset) = server.snapshot_at_an_instant();
1132    let id = server.repl_id();
1133    out.raw(b"+FULLRESYNC ");
1134    out.raw(&id);
1135    out.raw(b" ");
1136    out.raw(offset.to_string().as_bytes());
1137    out.raw(b"\r\n");
1138    // A bulk string with no newline after it, which is the one place in the
1139    // protocol where that is so. The replica reads the length and then exactly
1140    // that many bytes, and anything after them is the stream.
1141    out.raw(b"$");
1142    out.raw(image.len().to_string().as_bytes());
1143    out.raw(b"\r\n");
1144    out.raw(&image);
1145    attach(server, session);
1146    Ok(())
1147}
1148
1149/// Turn the connection into a replica, once the answer has been written.
1150///
1151/// Everything it was holding as a client goes first: a transaction it had open,
1152/// the keys it was watching and any subscription. All three are promises of a
1153/// reply, and a connection that has become a stream has no way left to keep one.
1154fn attach(server: &Server, session: &mut Session) {
1155    super::forget_session(server, session);
1156    server.take_replica(session.row());
1157}
1158
1159/// How many replicas have acknowledged everything written so far.
1160///
1161/// What `WAIT` answers with. It is the count that is already there rather than
1162/// one waited for, which is the right answer when it is already enough and is a
1163/// divergence when it is not. Waiting here means holding a connection until an
1164/// acknowledgement arrives, which needs the waiter list to be woken by something
1165/// that is not a key changing, and that is the same wakeup D-115 is about.
1166pub(super) fn caught_up(server: &Server) -> usize {
1167    let upto = server.repl_offset();
1168    server
1169        .replica_rows()
1170        .iter()
1171        .filter(|r| r.online.load(Relaxed) && r.ack.load(Relaxed) >= upto)
1172        .count()
1173}
1174
1175// ------------------------------------------------------------------- report
1176
1177/// What `ROLE` answers on a master.
1178///
1179/// The word, then the offset as an integer, then one row per replica. Only the
1180/// replicas that are online are listed, because a replica still being sent its
1181/// snapshot has no offset of its own to report and Redis leaves it out for the
1182/// same reason. Inside a row every field is a bulk string, including the port
1183/// and the offset, which is a shape nobody would pick today and is the shape
1184/// every client library already parses.
1185pub(super) fn role(server: &Server, out: &mut Out) {
1186    if server.following() {
1187        super::follow::role(server, out);
1188        return;
1189    }
1190    let rows = server.replica_rows();
1191    out.array(3);
1192    out.bulk(b"master");
1193    out.int(server.repl_offset() as i64);
1194    let at = out.len();
1195    let mut n = 0;
1196    for held in &rows {
1197        if !held.online.load(Relaxed) {
1198            continue;
1199        }
1200        let (host, port) = held.address();
1201        out.array(3);
1202        out.bulk(host.as_bytes());
1203        out.bulk(port.to_string().as_bytes());
1204        out.bulk(held.ack.load(Relaxed).to_string().as_bytes());
1205        n += 1;
1206    }
1207    out.close_array(at, n);
1208}
1209
1210/// The `Replication` section of `INFO`.
1211///
1212/// Redis reports a dozen fields on a master and a few more on a replica. This is
1213/// the master's set, since nothing here can be a replica yet, plus one `slaveN`
1214/// row per replica in the order they attached, which is what tooling reads to
1215/// find out who is following whom.
1216pub(super) fn info(server: &Server, s: &mut String) {
1217    use core::fmt::Write as _;
1218    let rows = server.replica_rows();
1219    let id = server.repl_id();
1220    let id2 = *server.repl.id2.lock();
1221    let offset = server.repl_offset();
1222    let now = server.clock.now_ms();
1223    let _ = write!(
1224        s,
1225        "# Replication\r\nrole:{}\r\n",
1226        super::follow::role_word(server)
1227    );
1228    // The master's own fields first when there is a master, which is where Redis
1229    // puts them: a tool reading the section top to bottom finds out what this
1230    // server is, then who it follows, then who follows it.
1231    super::follow::info(server, s);
1232    let _ = write!(s, "connected_slaves:{}\r\n", rows.len());
1233    for (i, held) in rows.iter().enumerate() {
1234        let (host, port) = held.address();
1235        let state = if held.online.load(Relaxed) {
1236            "online"
1237        } else {
1238            "wait_bgsave"
1239        };
1240        let lag = now.saturating_sub(held.ack_ms.load(Relaxed)) / 1000;
1241        let _ = write!(
1242            s,
1243            "slave{i}:ip={host},port={port},state={state},offset={},lag={lag}\r\n",
1244            held.ack.load(Relaxed),
1245        );
1246    }
1247    let _ = write!(
1248        s,
1249        "master_failover_state:{}\r\n\
1250         master_replid:{}\r\nmaster_replid2:{}\r\n\
1251         master_repl_offset:{offset}\r\nsecond_repl_offset:{}\r\n",
1252        server.failover_word(),
1253        String::from_utf8_lossy(&id),
1254        String::from_utf8_lossy(&id2),
1255        server.repl.second.load(Relaxed),
1256    );
1257    // The backlog is made when the first replica attaches and is never given
1258    // back, which is what `repl_backlog_active` is saying. The histlen is how
1259    // much of it is real, so a server that has never had a replica reports a
1260    // backlog that is off and empty rather than one that is on and zero long.
1261    let (active, first, histlen) = {
1262        let backlog = server.repl.backlog.lock();
1263        (
1264            usize::from(!backlog.ring.is_empty()),
1265            // Counted from one, the way a replica counts when it asks to carry
1266            // on, because this is the number it is being compared against. A
1267            // server that has never had a replica has no backlog and reports a
1268            // nought rather than the one that would be the first byte of the
1269            // one it has not made.
1270            if backlog.ring.is_empty() {
1271                0
1272            } else {
1273                backlog.first + 1
1274            },
1275            backlog.filled,
1276        )
1277    };
1278    let _ = write!(
1279        s,
1280        "repl_backlog_active:{active}\r\nrepl_backlog_size:{BACKLOG_BYTES}\r\n\
1281         repl_backlog_first_byte_offset:{first}\r\nrepl_backlog_histlen:{histlen}\r\n\r\n",
1282    );
1283}