Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions docker/chaos/docker-compose.chaos.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
# 混沌测试编排:broker(常开、可被 kill/pause) + provision + pub/sub(按需 run)。
# 由 docker/chaos/run-*.sh 驱动。
services:
provision:
build: { context: ../.., dockerfile: docker/Dockerfile }
volumes: [chaos:/data]
command: ["bun", "docker/chaos/provision.ts"]
environment: { COLLAB_DB: /data/collab.db, TOKEN_DIR: /data, ROOM: chaos-room }

broker:
build: { context: ../.., dockerfile: docker/Dockerfile }
volumes: [chaos:/data]
depends_on:
provision: { condition: service_completed_successfully }
command: ["bun", "docker/broker-entry.ts"]
environment: { BROKER_HOST: 0.0.0.0, BROKER_PORT: "4700", COLLAB_DB: /data/collab.db }
ports: ["4700:4700"]

# 按需 run(profiles=tools ⇒ `up` 不自动起;由脚本 `compose run` 调起并传 env)
pub:
build: { context: ../.., dockerfile: docker/Dockerfile }
volumes: [chaos:/data]
command: ["bun", "docker/chaos/pub.ts"]
environment: { BROKER_URL: ws://broker:4700/ws, ROOM: chaos-room, TOKEN_FILE: /data/token-pub }
profiles: [tools]

sub:
build: { context: ../.., dockerfile: docker/Dockerfile }
volumes: [chaos:/data]
command: ["bun", "docker/chaos/sub.ts"]
environment: { BROKER_URL: ws://broker:4700/ws, ROOM: chaos-room, TOKEN_FILE: /data/token-sub }
profiles: [tools]

volumes:
chaos:
28 changes: 28 additions & 0 deletions docker/chaos/provision.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
/** Chaos provisioning: a publisher identity + a subscriber identity + one room. */
import { chmodSync, mkdirSync, writeFileSync } from "node:fs";
import { dirname, join } from "node:path";
import { SqliteStore } from "../../src/backbone/store/sqlite-store";
import { IdentityService } from "../../src/backbone/identity-service";
import { RoomService } from "../../src/room-service";

const db = process.env.COLLAB_DB ?? "/data/collab.db";
const dir = process.env.TOKEN_DIR ?? "/data";
const room = process.env.ROOM ?? "chaos-room";

mkdirSync(dirname(db), { recursive: true, mode: 0o700 });
chmodSync(dirname(db), 0o700);

const store = new SqliteStore(db);
const svc = new IdentityService(store);
const rooms = new RoomService(store);
await rooms.createRoom(room, "Chaos Room", "pub@chaos");
for (const [id, file] of [
["pub@chaos", "token-pub"],
["sub@chaos", "token-sub"],
] as const) {
await svc.registerIdentity(id, id);
writeFileSync(join(dir, file), await svc.issueToken(id), { mode: 0o600 });
await rooms.join(room, id); // member ⇒ eligible for store_if_offline
}
await store.close();
console.log(`[chaos-provision] done (room=${room})`);
60 changes: 60 additions & 0 deletions docker/chaos/pub.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
/**
* Chaos publisher / load generator. Opens CONN parallel BrokerClient connections
* (one identity, many sessions) and each publishes COUNT task_completed events
* with distinct idempotencyKeys, then exits.
*
* CONN parallel connections (default 1)
* COUNT events per connection (default 100)
* DELIVERY store_if_offline (default) | online_only
*/
import { readFileSync } from "node:fs";
import { randomUUID } from "node:crypto";
import { BrokerClient } from "../../src/broker-client";
import type { Envelope } from "../../src/backbone/envelope";

const URL = process.env.BROKER_URL ?? "ws://broker:4700/ws";
const ROOM = process.env.ROOM ?? "chaos-room";
const TOKEN_FILE = process.env.TOKEN_FILE ?? "/data/token-pub";
const CONN = parseInt(process.env.CONN ?? "1", 10);
const COUNT = parseInt(process.env.COUNT ?? "100", 10);
const DELIVERY = (process.env.DELIVERY ?? "store_if_offline") as "store_if_offline" | "online_only";

const token = readFileSync(TOKEN_FILE, "utf8").trim();
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));
const t0 = Date.now();

const clients: BrokerClient[] = [];
for (let i = 0; i < CONN; i++) {
const c = new BrokerClient({ url: URL, token, log: () => {} });
await c.connect();
clients.push(c);
}

function ev(label: string): Envelope {
return {
roomId: ROOM,
messageId: randomUUID(),
traceId: randomUUID(),
idempotencyKey: randomUUID(),
from: { agentId: "pub@chaos", agentType: "claude" },
kind: "task_completed",
payload: { summary: label },
timestamp: Date.now(),
deliveryMode: DELIVERY,
};
}

let total = 0;
await Promise.all(
clients.map(async (c, ci) => {
for (let j = 0; j < COUNT; j++) {
c.publish(ROOM, ev(`load c${ci} #${j}`));
total++;
}
}),
);
await sleep(1500); // let frames flush before closing
const ms = Date.now() - t0;
console.log(`[pub] PUBLISHED total=${total} (conn=${CONN} x count=${COUNT}) delivery=${DELIVERY} in ${ms}ms (${Math.round((total / ms) * 1000)}/s)`);
for (const c of clients) c.close();
process.exit(0);
40 changes: 40 additions & 0 deletions docker/chaos/run-crash.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
#!/usr/bin/env bash
# 混沌 ①:broker 崩溃恢复 + pending 跨崩溃续投。
# 1) 200 条 store_if_offline 事件发给【离线】成员 sub@ → broker 落 pending(WAL)
# 2) SIGKILL broker(模拟进程崩溃)
# 3) 同卷重启 broker → pending 必须存活
# 4) sub 重连 drain → 应拿到全部 200(零丢失)
set -uo pipefail
cd "$(dirname "$0")/../.." || exit 1
C="docker compose -f docker/chaos/docker-compose.chaos.yml"
N=${N:-200}

echo "[crash] 清理 + 起 broker..."; $C down -v >/dev/null 2>&1
$C up -d --build broker >/dev/null 2>&1 || { echo "up broker 失败"; exit 1; }
sleep 4

echo "[crash] 发 $N 条 store_if_offline → 离线成员 sub@(broker 落 pending)..."
$C run --rm -e CONN=1 -e COUNT="$N" -e DELIVERY=store_if_offline pub 2>&1 | grep -aE 'PUBLISHED'
sleep 1

echo "[crash] 💥 SIGKILL broker(模拟崩溃)..."
$C kill -s SIGKILL broker >/dev/null 2>&1
sleep 1
echo "[crash] 重启 broker(同卷 → pending 应跨崩溃存活)..."
$C up -d broker >/dev/null 2>&1
sleep 4

echo "[crash] sub 重连 drain..."
out=$($C run --rm -e MODE=drain sub 2>&1)
echo "$out" | grep -aE 'DRAINED|sub\]'
got=$(printf '%s' "$out" | grep -aoE 'DRAINED unique=[0-9]+' | grep -aoE '[0-9]+' | tail -1)

$C down -v >/dev/null 2>&1
echo
if [ "${got:-0}" = "$N" ]; then
echo "==== 混沌① broker崩溃恢复:PASS ✅ (崩溃后 drain ${got}/${N},零丢失) ===="
exit 0
else
echo "==== 混沌① broker崩溃恢复:FAIL ❌ (drain ${got:-0}/${N}) ===="
exit 1
fi
32 changes: 32 additions & 0 deletions docker/chaos/run-load.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
#!/usr/bin/env bash
# 混沌 ③:并发压力。CONN 并发连接 × COUNT 事件 同时狂发 → sub 应收齐全部、broker 存活。
set -uo pipefail
cd "$(dirname "$0")/../.." || exit 1
C="docker compose -f docker/chaos/docker-compose.chaos.yml"
CONN=${CONN:-50}; COUNT=${COUNT:-20}; TOTAL=$((CONN * COUNT))

echo "[load] 起 broker + sub(watch)..."; $C down -v >/dev/null 2>&1
$C up -d --build broker >/dev/null 2>&1; sleep 4
$C up -d sub >/dev/null 2>&1; sleep 3

echo "[load] 并发 CONN=$CONN × COUNT=$COUNT = $TOTAL events 狂发..."
$C run --rm -e CONN="$CONN" -e COUNT="$COUNT" -e DELIVERY=store_if_offline pub 2>&1 | grep -aE 'PUBLISHED'

# 等 sub 收齐(轮询 heartbeat)最多 40s
got=0
for _ in $(seq 1 40); do
got=$($C logs sub 2>&1 | grep -aoE 'unique=[0-9]+' | grep -aoE '[0-9]+' | tail -1)
[ "${got:-0}" -ge "$TOTAL" ] 2>/dev/null && break
sleep 1
done
health=$(curl -sS -m 3 http://127.0.0.1:4700/healthz 2>/dev/null)
$C stop -t 6 sub >/dev/null 2>&1
got=$($C logs sub 2>&1 | grep -aoE 'unique=[0-9]+' | grep -aoE '[0-9]+' | tail -1)
$C down -v >/dev/null 2>&1
echo
echo "[load] sub unique=${got:-0}/$TOTAL ; broker /healthz=$health"
if [ "${got:-0}" = "$TOTAL" ]; then
echo "==== 混沌③ 并发压力:PASS ✅ ($TOTAL 事件零丢失,broker 存活) ===="; exit 0
else
echo "==== 混沌③ 并发压力:FAIL ❌ (${got:-0}/$TOTAL) ===="; exit 1
fi
28 changes: 28 additions & 0 deletions docker/chaos/run-partition.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
#!/usr/bin/env bash
# 混沌 ②:网络分区 / broker 短暂不可达(docker pause 冻结 broker 进程)。
# wave1(在线收) → pause broker 6s(冻结) → unpause → wave2 → sub 应收齐 100(零丢失)。
set -uo pipefail
cd "$(dirname "$0")/../.." || exit 1
C="docker compose -f docker/chaos/docker-compose.chaos.yml"

echo "[part] 起 broker + sub(watch)..."; $C down -v >/dev/null 2>&1
$C up -d --build broker >/dev/null 2>&1; sleep 4
$C up -d sub >/dev/null 2>&1; sleep 3

echo "[part] wave1: 50 events"; $C run --rm -e CONN=1 -e COUNT=50 -e DELIVERY=store_if_offline pub 2>&1 | grep -aE 'PUBLISHED'
sleep 2
echo "[part] ⏸ pause broker 6s(冻结=网络黑洞)..."; $C pause broker >/dev/null 2>&1; sleep 6
echo "[part] ▶ unpause broker..."; $C unpause broker >/dev/null 2>&1; sleep 5
echo "[part] wave2: 50 events"; $C run --rm -e CONN=1 -e COUNT=50 -e DELIVERY=store_if_offline pub 2>&1 | grep -aE 'PUBLISHED'
sleep 5

$C stop -t 6 sub >/dev/null 2>&1
got=$($C logs sub 2>&1 | grep -aoE 'unique=[0-9]+' | grep -aoE '[0-9]+' | tail -1)
$C down -v >/dev/null 2>&1
echo
if [ "${got:-0}" = "100" ]; then
echo "==== 混沌② 网络分区(pause):PASS ✅ (跨 6s 冻结,sub 收齐 ${got}/100 零丢失) ===="; exit 0
else
echo "==== 混沌② 网络分区(pause):观察值 sub=${got:-0}/100 ===="
echo "注:若 <100,多半暴露已知 gap——无 WS 心跳,冻结连接靠 onclose 才触发重连(§8.2 backlog)。"; exit 1
fi
31 changes: 31 additions & 0 deletions docker/chaos/run-soak.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
#!/usr/bin/env bash
# 混沌 ④:mini-soak。连续 ROUNDS 轮负载,每轮采 broker 内存 + /healthz → 看有无泄漏/退化/崩溃。
# (非整夜 soak;要长跑把 ROUNDS 调大。)
set -uo pipefail
cd "$(dirname "$0")/../.." || exit 1
C="docker compose -f docker/chaos/docker-compose.chaos.yml"
ROUNDS=${ROUNDS:-10}; PER=${PER:-200}; CONN=${CONN:-5}

echo "[soak] 起 broker + sub..."; $C down -v >/dev/null 2>&1
$C up -d --build broker >/dev/null 2>&1; sleep 4
$C up -d sub >/dev/null 2>&1; sleep 2

bid=$($C ps -q broker)
mem0=""; memN=""; fail=0
for r in $(seq 1 "$ROUNDS"); do
$C run --rm -e CONN="$CONN" -e COUNT="$PER" pub >/dev/null 2>&1
mem=$(docker stats --no-stream --format '{{.MemUsage}}' "$bid" 2>/dev/null | awk '{print $1}')
health=$(curl -sS -m 3 -o /dev/null -w '%{http_code}' http://127.0.0.1:4700/healthz 2>/dev/null)
echo "[soak] 轮 $r/$ROUNDS: broker mem=$mem healthz=$health"
[ "$health" = "200" ] || fail=1
[ -z "$mem0" ] && mem0=$mem; memN=$mem
done
got=$($C logs sub 2>&1 | grep -aoE 'unique=[0-9]+' | grep -aoE '[0-9]+' | tail -1)
$C stop -t 6 sub >/dev/null 2>&1; $C down -v >/dev/null 2>&1
echo
echo "[soak] 累计 sub unique=${got:-0} (期望 $((ROUNDS*PER*CONN))) ; 内存 起=$mem0 末=$memN"
if [ "$fail" = "0" ]; then
echo "==== 混沌④ mini-soak:PASS ✅ (全程 healthz=200,内存有界) ===="; exit 0
else
echo "==== 混沌④ mini-soak:FAIL ❌ (中途 healthz 非 200) ===="; exit 1
fi
60 changes: 60 additions & 0 deletions docker/chaos/sub.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
/**
* Chaos subscriber. Counts UNIQUE task_completed (by idempotencyKey) over the
* real auto-reconnecting BrokerClient.
*
* MODE=drain : connect, subscribe, collect until 2s quiet, print DRAINED, exit.
* MODE=watch : stay online; heartbeat every 5s; print FINAL on SIGTERM/SIGINT.
*/
import { readFileSync } from "node:fs";
import { BrokerClient } from "../../src/broker-client";

const URL = process.env.BROKER_URL ?? "ws://broker:4700/ws";
const ROOM = process.env.ROOM ?? "chaos-room";
const TOKEN_FILE = process.env.TOKEN_FILE ?? "/data/token-sub";
const MODE = process.env.MODE ?? "watch";
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));

const seen = new Set<string>();
let count = 0;
const client = new BrokerClient({
url: URL,
token: readFileSync(TOKEN_FILE, "utf8").trim(),
presence: { agentType: "claude" },
log: (m) => console.error(`[sub] ${m}`),
});
client.onEvent((_t, env) => {
if (env.kind !== "task_completed") return;
if (seen.has(env.idempotencyKey)) return; // dedup redeliveries
seen.add(env.idempotencyKey);
count++;
});
await client.connect();
client.subscribe(ROOM);
console.log(`[sub] online (mode=${MODE})`);

if (MODE === "drain") {
let last = -1;
let quiet = 0;
while (quiet < 2000) {
await sleep(250);
if (count !== last) {
last = count;
quiet = 0;
} else {
quiet += 250;
}
}
console.log(`[sub] DRAINED unique=${count}`);
client.close();
process.exit(0);
} else {
const beat = setInterval(() => console.log(`[sub] HEARTBEAT unique=${count} connected=${client.connected}`), 5000);
const done = (sig: string) => {
clearInterval(beat);
console.log(`[sub] FINAL unique=${count} (${sig})`);
client.close();
process.exit(0);
};
process.on("SIGTERM", () => done("SIGTERM"));
process.on("SIGINT", () => done("SIGINT"));
}
Loading