Skip to content

9.3.4: resync hangs in WFBitMapS/WFBitMapT (or stays Established/Inconsistent) after "unsecured-resync" when the sync source loses a diskless primary #148

Description

@Greg0ryLau

Summary

With DRBD 9.3.4, when a sync source loses its connection to a diskless primary
while a resync to a third node is running, the sync target correctly refuses to
continue ("Can not secure resync request ... against a source that lost the
writer, ending the resync" — the fix for acknowledged writes being overwritten,
from 9.2.20). But right after that the resync does not restart:

  • the sync source re-enters the bitmap exchange (resync-again → WFBitMapS)
    and sends its bitmap;
  • the sync target has already left the resync on its own (unsecured-resync →
    Established) and drops that bitmap (unexpected repl_state (Established) in receive_bitmap);
  • afterwards either the target enters WFBitMapT through after-unstable and
    both sides wait for each other's bitmap forever, or the target stays
    Established + Inconsistent next to an UpToDate peer.

No acknowledged data is lost (the target stays Inconsistent), but redundancy is
gone until the connection is cycled by hand. It does not recover when the
primary demotes.

On 9.3.2 the same scenario silently lost acknowledged writes (the target became
UpToDate with stale data), so 9.3.4 is a clear improvement — this report is
about the remaining liveness problem.

Environment

  • DRBD 9.3.4 (drbd-dkms 9.3.4-1ppa1~noble1, rebuilt on Ubuntu 20.04, kernel 5.4.0-216), drbd-utils 9.34.0
  • LINSTOR 1.33.2, 3 nodes, TCP transport, quorum majority, on-no-quorum suspend-io,
    timeout 10, ping-int 2, c-max-rate 400M, c-min-rate 20M, c-fill-target 1M
  • Resource: 2 diskful replicas (node A = 02, node B = 03) + diskless primary (node P = 01)

Reproduction

  1. P (diskless) is Primary and writes continuously (random 4k O_DIRECT writes).
  2. B disconnects for 10 s (drbdadm disconnect), then drbdadm adjust → B becomes SyncTarget from A.
  3. 1–3 s into that resync, drop the traffic between P and A only (iptables on the
    resource port, both directions, on both nodes). P keeps quorum via B + tiebreaker logic.
  4. Restore P↔A after 6–19 s.

Script: vrf-t5.sh a N (below). It writes with a small integrity probe (genio.py, below) that records every
acknowledged write and afterwards checks both backing devices for acknowledged writes that were lost.

Result on 9.3.4: resync hung in 5 of 14 rounds (4 × WFBitMapS/WFBitMapT, 1 × Established/Inconsistent).

Log excerpt (round 2 of a run)

Sync target (B):

20:32:16.394 computer02: Began resync as SyncTarget (will sync 33388 KB [8347 bits set]).
20:32:16.442 computer02: Becoming WFBitMapT because primary is diskless
20:32:16.442 computer02: ...postponing this until current resync finished
20:32:21.881 computer02: Can not secure resync request at 241048s+4096 against a source that lost the writer, ending the resync
20:32:21.882 computer02: repl( SyncTarget -> Established ) [unsecured-resync]
20:32:21.913 computer02: receive bitmap stats [Bytes(packets)]: plain 0(0), RLE 1024(1), total 1024; compression: 96.9%
20:32:21.913 computer02: unexpected repl_state (Established) in receive_bitmap
20:33:15.121 computer02: Becoming WFBitMapT after unstable
20:33:15.121 computer02: repl( Established -> WFBitMapT ) [after-unstable]
   (no further progress)

Sync source (A):

20:32:21.890 computer03: Resync aborted (total 5 sec; paused 0 sec; 8480 K/sec)
20:32:21.890 computer03: repl( SyncSource -> Established ) [resync-finished]
20:32:21.890 computer03: repl( Established -> WFBitMapS ) [resync-again]
20:32:21.890 computer03: send bitmap stats [Bytes(packets)]: plain 0(0), RLE 1024(1), total 1024; compression: 96.9%
20:33:15.124 computer03: Becoming WFBitMapS after unstable
20:33:15.124 computer03: ...postponing this until current resync finished
   (stays in WFBitMapS)

Analysis

The two sides leave the resync on different paths:

  • the source goes through drbd_resync_finished() → resync_again() and enters
    L_WF_BITMAP_S, which immediately sends its bitmap;
  • the target leaves through drbd_resync_request_unreachable() →
    change_repl_state(L_ESTABLISHED, "unsecured-resync"), which does not run
    resync_again(), so the target is L_ESTABLISHED when the source's bitmap
    arrives. receive_bitmap() merges the bits into the slot but, being neither in
    L_WF_BITMAP_S nor L_WF_BITMAP_T, neither answers nor starts a resync.

The source then waits in L_WF_BITMAP_S for a bitmap that will never be sent.
If the target later enters L_WF_BITMAP_T (after-unstable), it waits for a
bitmap that has already been consumed.

Proposed fix

In receive_bitmap(), when a bitmap arrives while we are L_ESTABLISHED, our
disk still needs a resync (D_INCONSISTENT / D_OUTDATED, never
D_UP_TO_DATE) and the peer has usable data, join the peer's exchange as sync
target: enter L_WF_BITMAP_T, send our bitmap, start the resync, and consume
the pending resync_again that this exchange serves. Patch below (≈60 lines,
drbd_receiver.c only). It applies to current master (574c9ca9e) with an offset; master still has
the same receive_bitmap() / unsecured-resync code.

Test results with the patch

Same reproduction, 20 rounds, watchdog in log-only mode:

9.3.4 (14 rounds) 9.3.4 + patch (20 rounds)
target ended resync with "unsecured-resync" ≥ 5 8
target joined the peer's bitmap exchange — 8 / 8
"unexpected repl_state (Established) in receive_bitmap" every time 0
hung in WFBitMapS / WFBitMapT 4 0
Established + Inconsistent next to an UpToDate peer 1 (did not recover even after the primary demoted) 5, 30–50 s each, all resolved by the after-unstable handshake when the primary demoted
acknowledged writes lost / replicas differ 0 / 0 0 / 0

Regression (patched): network cut between the diskful primary and a secondary (8 rounds), multi-source resync without IO (6) and with IO (3), verify + reconnect after corrupting blocks: all identical to unpatched 9.3.4, no acknowledged write lost. Builds via dkms on 5.4.0-216; compiled against 6.8 headers it adds no warnings.

The patch does not change the wire protocol; a patched node only answers a
bitmap it would otherwise have dropped.

Question: Established + Inconsistent while the primary keeps writing

Even with the patch, the joined resync can again be ended with
"unsecured-resync" (the source still has not reconnected to the writer). The
target then stays Established + Inconsistent until the cluster becomes
"stable" again (after-unstable handshake), which with a long-running diskless
primary (a VM) may be days. The sync source regains UpToDate as soon as it
reconnects to the primary. Is waiting for stability intended here, or should
the resync restart once the source is UpToDate again? We did not try to change
this, since it needs both sides to agree on a new handshake.

Workaround

Cycling the connection on the sync target (drbdadm disconnect <res>:<source>; drbdadm adjust <res>) is safe because the target is Inconsistent. We run a
small watchdog that does this after 60 s in WFBitMapT or 120 s
Established/Inconsistent (autoheal.py, below; comments and log messages are in Chinese).

Possibly related: #116 (device blocked in WFBitmapS on a bad network), which has no details.

Files

Patch: drbd: join the peer's bitmap exchange when a bitmap arrives while Established
diff --git a/drbd/drbd_receiver.c b/drbd/drbd_receiver.c
index a040dd8..bd15d08 100644
--- a/drbd/drbd_receiver.c
+++ b/drbd/drbd_receiver.c
@@ -10359,6 +10359,66 @@ static bool ready_for_bitmap(struct drbd_device *device)
 	return ready;
 }
 
+/*
+ * A peer that re-entered the bitmap exchange as sync source (for example via
+ * resync-again after it lost a diskless primary) sends its bitmap right away.
+ * If our side left the resync on its own in the meantime (unsecured-resync),
+ * we are already L_ESTABLISHED when that bitmap arrives.  Dropping it leaves
+ * the peer waiting in L_WF_BITMAP_S for our bitmap, and a later
+ * L_WF_BITMAP_T of ours waits for a bitmap that already came: neither side
+ * progresses until the connection is cycled.
+ *
+ * The peer's bits have been merged into our slot already, so we can answer
+ * as a sync target would have: enter L_WF_BITMAP_T, send our bitmap and
+ * start the resync.  Only do that when our disk still needs a resync, never
+ * when it is D_UP_TO_DATE.
+ */
+static bool can_join_peer_bitmap_exchange(struct drbd_peer_device *peer_device)
+{
+	enum drbd_disk_state disk = peer_device->device->disk_state[NOW];
+	enum drbd_disk_state pdsk = peer_device->disk_state[NOW];
+
+	if (peer_device->connection->agreed_pro_version < 110)
+		return false;
+	if (disk != D_INCONSISTENT && disk != D_OUTDATED)
+		return false;
+	return pdsk == D_OUTDATED || pdsk == D_CONSISTENT || pdsk == D_UP_TO_DATE;
+}
+
+static int join_peer_bitmap_exchange(struct drbd_peer_device *peer_device)
+{
+	struct drbd_device *device = peer_device->device;
+	enum drbd_state_rv rv;
+	int err;
+
+	drbd_info(peer_device, "Received bitmap while Established and %s; joining the peer's bitmap exchange as sync target\n",
+		  drbd_disk_str(device->disk_state[NOW]));
+	rv = stable_change_repl_state(peer_device, L_WF_BITMAP_T, CS_VERBOSE, "join-bitmap-exchange");
+	if (rv < SS_SUCCESS) {
+		drbd_info(peer_device, "Could not join the bitmap exchange: %s\n", drbd_set_st_err_str(rv));
+		return 0;
+	}
+
+	/* This exchange serves a sync-target handshake that was postponed while
+	 * the previous resync ran; do not re-enter L_WF_BITMAP_T for it later. */
+	write_lock_irq(&device->resource->state_rwlock);
+	if (peer_device->resync_again > 0)
+		peer_device->resync_again--;
+	write_unlock_irq(&device->resource->state_rwlock);
+
+	if (!get_ldev(device))
+		return 0;
+	drbd_bm_slot_lock(peer_device, "join bitmap exchange", BM_LOCK_CLEAR | BM_LOCK_BULK);
+	err = drbd_send_bitmap(device, peer_device);
+	drbd_bm_slot_unlock(peer_device);
+	put_ldev(device);
+	if (err)
+		return err;
+
+	drbd_start_resync(peer_device, L_SYNC_TARGET, "join-bitmap-exchange");
+	return 0;
+}
+
 /* Since we are processing the bitfield from lower addresses to higher,
    it does not matter if the process it in 32 bit chunks or 64 bit
    chunks as long as it is little endian. (Understand it as byte stream,
@@ -10478,6 +10538,8 @@ static int receive_bitmap(struct drbd_connection *connection, struct packet_info
 		} else {
 			drbd_start_resync(peer_device, L_SYNC_TARGET, "receive-bitmap");
 		}
+	} else if (repl_state == L_ESTABLISHED && can_join_peer_bitmap_exchange(peer_device)) {
+		return join_peer_bitmap_exchange(peer_device);
 	} else {
 		/* admin may have requested C_DISCONNECTING,
 		 * other threads may have noticed network errors */
vrf-t5.sh — reproduction (mode a: diskless writer, cut primary ↔ sync source during resync)
#!/usr/bin/env bash
# T5:网络局部分区 + 持续写入后的数据完整性(genio 探针)
#   vrf-t5.sh a|b ROUNDS
set -u
MODE=$1; ROUNDS=${2:-5}
R=cg-vrf-c; P=7002; N=65536; LV=/dev/linstor_cg_vrf/${R}_00000; DEV=/dev/drbd/by-res/$R/0
declare -A IP=(  # storage-network address of each node (adjust)
  [computer01]=10.0.0.1 [computer02]=10.0.0.2 [computer03]=10.0.0.3)
on() { ssh -o BatchMode=yes root@$1 "${@:2}" </dev/null; }
block() {  # 在两个节点上都丢弃彼此之间本资源端口的流量
  local a=$1 b=$2
  on $a "iptables -N CG_VRF 2>/dev/null; iptables -C INPUT -j CG_VRF 2>/dev/null || iptables -I INPUT -j CG_VRF; iptables -C OUTPUT -j CG_VRF 2>/dev/null || iptables -I OUTPUT -j CG_VRF;
         for d in '--dport' '--sport'; do iptables -A CG_VRF -s ${IP[$b]} -p tcp \$d $P -j DROP; iptables -A CG_VRF -d ${IP[$b]} -p tcp \$d $P -j DROP; done"
  on $b "iptables -N CG_VRF 2>/dev/null; iptables -C INPUT -j CG_VRF 2>/dev/null || iptables -I INPUT -j CG_VRF; iptables -C OUTPUT -j CG_VRF 2>/dev/null || iptables -I OUTPUT -j CG_VRF;
         for d in '--dport' '--sport'; do iptables -A CG_VRF -s ${IP[$a]} -p tcp \$d $P -j DROP; iptables -A CG_VRF -d ${IP[$a]} -p tcp \$d $P -j DROP; done"
}
unblock_all() { for h in computer01 computer02 computer03; do on $h "iptables -F CG_VRF 2>/dev/null; true"; done; }
cleanup() { unblock_all; for h in computer01 computer02; do on $h "pkill -f [c]g-genio.py" 2>/dev/null; done; }
trap cleanup EXIT
waitclean() { for i in $(seq 1 300); do s=$(on computer02 "drbdsetup status $R"); echo "$s" | grep -qE "Inconsistent|Sync|Outdated|Connecting|StandAlone|DUnknown|suspended:(user|quorum|no-data)" || return 0; sleep 2; done; echo "   NOT CLEAN: $(on computer02 "drbdsetup status $R" | tr '\n' ' ')"; }
state() { on $1 "drbdsetup status $R" | tr -s ' \n' ' ' | cut -c1-220; }
for r in $(seq 1 $ROUNDS); do
  waitclean
  if [ "$MODE" = a ]; then W=computer01; else W=computer02; fi
  on $W "rm -f /root/cg-ack.json; nohup setsid python3 /root/cg-genio.py write $DEV $N /root/cg-ack.json 75 > /root/cg-genio.out 2>&1 &"
  sleep 5
  if [ "$MODE" = a ]; then
    on computer03 "drbdadm disconnect $R"; sleep 10
    on computer03 "drbdadm adjust $R"
    d=$((RANDOM % 3 + 1)); sleep $d
    echo "== T5a round $r: 03 resyncing, cut 01<->02 after ${d}s | 03: $(state computer03)"
    block computer01 computer02; h=$((RANDOM % 15 + 5)); sleep $h
    echo "   during cut (${h}s): 01: $(state computer01)"
    unblock_all
  else
    for k in 1 2; do
      d=$((RANDOM % 8 + 2)); h=$((RANDOM % 18 + 3)); sleep $d
      block computer02 computer03; sleep $h
      echo "== T5b round $r cut#$k: 02<->03 for ${h}s | 02: $(state computer02)"
      unblock_all
    done
  fi
  while on $W "pgrep -f [c]g-genio.py >/dev/null"; do sleep 3; done
  echo "   writer: $(on $W 'cat /root/cg-genio.out')"
  waitclean
  on $W "cat /root/cg-ack.json" > /tmp/cg-ack.json
  for h in computer02 computer03; do
    scp -q /tmp/cg-ack.json root@$h:/root/cg-ack.json
    echo "   check $h: $(on $h "python3 /root/cg-genio.py check $LV $N /root/cg-ack.json" | xargs)"
  done
  echo "   oos 02->03 / 03->02: $(on computer02 "drbdsetup status $R --verbose --statistics" | grep -A3 'computer03 node' | grep -o 'out-of-sync:[0-9]*') / $(on computer03 "drbdsetup status $R --verbose --statistics" | grep -A3 'computer02 node' | grep -o 'out-of-sync:[0-9]*')"
  on computer02 "journalctl -k --since '-3min' -o cat | grep 'drbd $R' | grep -E 'n_oos|split-brain|Split-Brain|uuid_compare|No resync|diverg'" | sort | uniq -c | sed 's/^/   log02: /' | head -6
done
genio.py — write-integrity probe (acknowledged-write ledger + per-block generation check)
#!/usr/bin/env python3
"""数据完整性探针(T5)。

  genio.py init  DEV NBLOCKS              以代次 0 写满前 NBLOCKS 个 4 KiB 块
  genio.py write DEV NBLOCKS ACKFILE SECS  随机写入,记录每个块已确认(write 返回)的最大代次,每秒落盘一次
  genio.py check DEV NBLOCKS [ACKFILE]     读回校验:内容校验、代次 >= 已确认代次;输出每块 (代次) 摘要到 stdout 最后一行
每个块:8 字节块号 + 8 字节代次 + 填充(由块号和代次生成)+ 8 字节校验。
"""
import hashlib, json, mmap, os, random, struct, sys, time

BS = 4096


def content(block, gen):
    head = struct.pack("<QQ", block, gen)
    seed = hashlib.sha256(head).digest()
    body = (seed * (BS // 32))[: BS - 24]
    csum = hashlib.md5(head + body).digest()[:8]
    return head + body + csum


def parse(buf):
    block, gen = struct.unpack("<QQ", buf[:16])
    ok = hashlib.md5(buf[:BS - 8]).digest()[:8] == buf[BS - 8:]
    return block, gen, ok


def main():
    mode, dev, n = sys.argv[1], sys.argv[2], int(sys.argv[3])
    flags = os.O_DIRECT | (os.O_RDONLY if mode == "check" else os.O_RDWR)
    fd = os.open(dev, flags)
    m = mmap.mmap(-1, BS)
    if mode == "init":
        big = mmap.mmap(-1, 256 * BS)
        for start in range(0, n, 256):
            cnt = min(256, n - start)
            for i in range(cnt):
                big[i * BS:(i + 1) * BS] = content(start + i, 0)
            os.pwritev(fd, [memoryview(big)[:cnt * BS]], start * BS)
        return
    if mode == "write":
        ackfile, secs = sys.argv[4], float(sys.argv[5])
        ack, gen, errors, last = {}, 0, 0, time.time()
        end = time.time() + secs
        while time.time() < end:
            b = random.randrange(n)
            gen += 1
            m[:] = content(b, gen)
            try:
                os.pwritev(fd, [m], b * BS)
                ack[b] = gen
            except OSError as e:
                errors += 1
                time.sleep(0.05)
            if time.time() - last >= 1:
                with open(ackfile + ".tmp", "w") as f:
                    json.dump({"gen": gen, "errors": errors, "ack": ack}, f)
                os.rename(ackfile + ".tmp", ackfile)
                last = time.time()
        with open(ackfile + ".tmp", "w") as f:
            json.dump({"gen": gen, "errors": errors, "ack": ack, "done": True}, f)
        os.rename(ackfile + ".tmp", ackfile)
        print("writes=%d errors=%d acked_blocks=%d" % (gen, errors, len(ack)))
        return
    if mode == "check":
        ack = {}
        if len(sys.argv) > 4:
            ack = {int(k): v for k, v in json.load(open(sys.argv[4]))["ack"].items()}
        bad_csum, bad_block, stale, gens = 0, 0, [], []
        for b in range(n):
            os.preadv(fd, [m], b * BS)
            blk, gen, ok = parse(bytes(m))
            if not ok:
                bad_csum += 1
                gens.append(-1)
                continue
            if blk != b:
                bad_block += 1
            if gen < ack.get(b, 0):
                stale.append((b, gen, ack[b]))
            gens.append(gen)
        digest = hashlib.md5(json.dumps(gens).encode()).hexdigest()[:12]
        print("blocks=%d bad_checksum=%d wrong_block=%d stale(acked write lost)=%d %s" % (
            n, bad_csum, bad_block, len(stale), stale[:5]))
        print("GENS", digest)


if __name__ == "__main__":
    main()
autoheal.py — watchdog used as workaround (reconnects only an unsynced sync target)
"""DRBD 同步停滞的自动恢复(在每个 LINSTOR 节点上常驻运行,只看本节点)。

只处理本节点是同步目标、且本地数据未完成同步(Inconsistent / Outdated)的连接:
  - F8(DRBD 9.3.4,results.md §11.6):同步源失去无盘主节点后,目标以 unsecured-resync 结束同步,
    源端的 resync-again 位图被丢弃,双方停在 WFBitMapS / WFBitMapT 互相等待,不会自行恢复
  - §9.3:SyncTarget 的剩余量长时间不变
  - 本地 Inconsistent、连接正常、对端 UpToDate,却停在 Established 不开始同步(F8 的另一种结局:
    同步以 unsecured-resync 结束后没有新的位图交换,直到主节点降级或重连才会恢复)
处理方法是在本节点断开并重连该连接(drbdadm disconnect <资源>:<对端>; drbdadm adjust <资源>)。
本地数据未完成同步,重连不会用它覆盖任何一方;本地已是 UpToDate 时绝不重连(两端 UpToDate 时重连不能用来修复,§11 F4)。

运行:python3 -m linstor_ops.autoheal [--dry-run] [--once]
指标(node_exporter textfile):cinder_gen_drbd_autoheal_total / _giveup_total / _stuck_seconds
"""
import argparse
import json
import os
import subprocess
import sys
import time

WFBITMAP_STATES = ("WFBitMapT", "WFSyncUUID", "StartingSyncT")
SYNC_STATES = ("SyncTarget",)
NOT_SYNCED = ("Inconsistent", "Outdated")
TEXTFILE = "/var/lib/prometheus/node-exporter/cinder_gen_drbd_autoheal.prom"


def log(msg):
    sys.stdout.write("%s\n" % msg)
    sys.stdout.flush()


def run(cmd, timeout=60):
    p = subprocess.run(cmd, shell=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, timeout=timeout,
                       stdin=subprocess.DEVNULL)
    return p.returncode, p.stdout.decode(errors="replace")


def read_status(runner=run):
    rc, out = runner("drbdsetup status --json --statistics")
    if rc != 0:
        raise RuntimeError("drbdsetup status failed: %s" % out[-200:])
    return json.loads(out or "[]")


def candidates(status):
    """返回 {(资源, 对端): (原因, 剩余 KiB)}:本节点是未完成同步的目标、连接正常、对端 UpToDate"""
    found = {}
    for r in status:
        disks = {d.get("volume", 0): d.get("disk-state") for d in r.get("devices", [])}
        for c in r.get("connections", []):
            if c.get("connection-state") != "Connected":
                continue
            for pd in c.get("peer_devices", []):
                if disks.get(pd.get("volume", 0)) not in NOT_SYNCED or pd.get("peer-disk-state") != "UpToDate":
                    continue
                repl = pd.get("replication-state", "")
                key = (r["name"], c["name"])
                if repl in WFBITMAP_STATES:
                    found[key] = ("wfbitmap", pd.get("out-of-sync") or 0)
                elif repl in SYNC_STATES and pd.get("resync-suspended", "no") in ("no", False):
                    found[key] = ("resync-stalled", pd.get("out-of-sync") or 0)
                elif repl == "Established" and disks.get(pd.get("volume", 0)) == "Inconsistent":
                    found[key] = ("idle-inconsistent", pd.get("out-of-sync") or 0)
    return found


class Healer(object):
    def __init__(self, runner=run, clock=time.time, wfbitmap_after=60, stalled_after=600, idle_after=120,
                 max_per_hour=3, dry_run=False, textfile=TEXTFILE):
        self.runner, self.clock = runner, clock
        self.wfbitmap_after, self.stalled_after, self.idle_after = wfbitmap_after, stalled_after, idle_after
        self.max_per_hour, self.dry_run, self.textfile = max_per_hour, dry_run, textfile
        self.seen = {}        # key -> (原因, 首次看到的时间, 上次剩余量, 剩余量上次变化的时间)
        self.history = {}     # key -> [重连时间]
        self.total = {}       # (key, 原因) -> 次数
        self.giveup = {}      # key -> 次数
        self.giveup_logged = {}

    def tick(self):
        now = self.clock()
        found = candidates(read_status(self.runner))
        for key in list(self.seen):
            if key not in found:
                reason, since, _, _ = self.seen.pop(key)
                log("%s:%s 已恢复(%s,持续 %ds)" % (key[0], key[1], reason, now - since))
        actions = []
        for key, (reason, oos) in found.items():
            prev = self.seen.get(key)
            if prev is None or prev[0] != reason:
                self.seen[key] = (reason, now, oos, now)
                continue
            _, since, last_oos, changed = prev
            if oos != last_oos:
                self.seen[key] = (reason, since, oos, now)
                continue
            waited = now - (changed if reason == "resync-stalled" else since)
            limit = {"wfbitmap": self.wfbitmap_after, "idle-inconsistent": self.idle_after}.get(reason, self.stalled_after)
            if waited >= limit:
                actions.append((key, reason, waited))
        for key, reason, waited in actions:
            self.heal(key, reason, waited, now)
        self.write_metrics(now)
        return actions

    def heal(self, key, reason, waited, now):
        res, peer = key
        recent = [t for t in self.history.get(key, []) if now - t < 3600]
        if len(recent) >= self.max_per_hour:
            self.giveup[key] = self.giveup.get(key, 0) + 1
            if now - self.giveup_logged.get(key, 0) >= 3600:
                self.giveup_logged[key] = now
                log("%s:%s %s 已持续 %ds,但 1 小时内已重连 %d 次,不再自动处理(需人工,runbook §8「副本一致性」)"
                    % (res, peer, reason, waited, len(recent)))
            self.seen[key] = (reason, now, None, now)
            return
        # 执行前再确认一次:本地仍未完成同步、状态没变
        again = candidates(read_status(self.runner)).get(key)
        if again is None or again[0] != reason:
            log("%s:%s 状态已变化,不处理" % (res, peer))
            return
        log("%s:%s %s 已持续 %ds(剩余 %s KiB),在本节点断开重连%s"
            % (res, peer, reason, waited, again[1], "(dry-run)" if self.dry_run else ""))
        if not self.dry_run:
            rc, out = self.runner("drbdadm disconnect %s:%s; sleep 2; drbdadm adjust %s" % (res, peer, res))
            if rc != 0:
                log("  重连命令返回 %d:%s" % (rc, out.strip()[-300:]))
        recent.append(now)
        self.history[key] = recent
        self.total[(key, reason)] = self.total.get((key, reason), 0) + 1
        self.seen.pop(key, None)

    def write_metrics(self, now):
        if not self.textfile:
            return
        lines = ["# HELP cinder_gen_drbd_autoheal_total Reconnects done by cinder-gen DRBD autoheal.",
                 "# TYPE cinder_gen_drbd_autoheal_total counter"]
        for ((res, peer), reason), n in sorted(self.total.items()):
            lines.append('cinder_gen_drbd_autoheal_total{name="%s",conn_name="%s",reason="%s"} %d' % (res, peer, reason, n))
        lines += ["# HELP cinder_gen_drbd_autoheal_giveup_total Stalls left alone because the hourly limit was reached.",
                  "# TYPE cinder_gen_drbd_autoheal_giveup_total counter"]
        for (res, peer), n in sorted(self.giveup.items()):
            lines.append('cinder_gen_drbd_autoheal_giveup_total{name="%s",conn_name="%s"} %d' % (res, peer, n))
        lines += ["# HELP cinder_gen_drbd_autoheal_stuck_seconds How long a sync target has been waiting.",
                  "# TYPE cinder_gen_drbd_autoheal_stuck_seconds gauge"]
        for (res, peer), (reason, since, _, _) in sorted(self.seen.items()):
            lines.append('cinder_gen_drbd_autoheal_stuck_seconds{name="%s",conn_name="%s",reason="%s"} %d'
                         % (res, peer, reason, now - since))
        lines += ["# HELP cinder_gen_drbd_autoheal_last_run_timestamp_seconds Last successful check.",
                  "# TYPE cinder_gen_drbd_autoheal_last_run_timestamp_seconds gauge",
                  "cinder_gen_drbd_autoheal_last_run_timestamp_seconds %d" % now]
        tmp = self.textfile + ".tmp"
        try:
            with open(tmp, "w") as f:
                f.write("\n".join(lines) + "\n")
            os.rename(tmp, self.textfile)
        except OSError as e:
            log("写指标失败:%s" % e)


def main(argv=None):
    ap = argparse.ArgumentParser(description="DRBD 同步停滞自动恢复(本节点)")
    ap.add_argument("--interval", type=int, default=10, help="检查间隔(秒)")
    ap.add_argument("--wfbitmap-after", type=int, default=60, help="WFBitMapT 等状态持续多少秒后重连")
    ap.add_argument("--stalled-after", type=int, default=600, help="SyncTarget 剩余量多少秒不变后重连")
    ap.add_argument("--idle-after", type=int, default=120, help="本地 Inconsistent 却停在 Established 多少秒后重连")
    ap.add_argument("--max-per-hour", type=int, default=3, help="每个连接每小时最多重连次数")
    ap.add_argument("--textfile", default=TEXTFILE, help="node_exporter textfile 路径(空字符串表示不写)")
    ap.add_argument("--dry-run", action="store_true", help="只记录,不重连")
    ap.add_argument("--once", action="store_true", help="只检查一次")
    a = ap.parse_args(argv)
    h = Healer(wfbitmap_after=a.wfbitmap_after, stalled_after=a.stalled_after, idle_after=a.idle_after,
               max_per_hour=a.max_per_hour,
               dry_run=a.dry_run, textfile=a.textfile or None)
    log("启动:间隔 %ds,WFBitMap* %ds / Inconsistent 空闲 %ds / SyncTarget 停滞 %ds 后重连,每连接每小时最多 %d 次%s"
        % (a.interval, a.wfbitmap_after, a.idle_after, a.stalled_after, a.max_per_hour, ",dry-run" if a.dry_run else ""))
    while True:
        try:
            h.tick()
        except Exception as e:  # 单次失败(例如 drbdsetup 暂时不可用)不退出
            log("检查失败:%s" % e)
        if a.once:
            return 0
        time.sleep(a.interval)


if __name__ == "__main__":
    sys.exit(main())

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions