-
Notifications
You must be signed in to change notification settings - Fork 758
Expand file tree
/
Copy pathjustfile
More file actions
464 lines (390 loc) · 16.9 KB
/
Copy pathjustfile
File metadata and controls
464 lines (390 loc) · 16.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
# DevKit - AutoMQ Local Development Environment
set dotenv-load := false
export AUTOMQ_DEV_HOME := justfile_directory() + "/.."
DEVKIT := justfile_directory() + "/.devkit"
COMPOSE := "docker compose -f " + justfile_directory() + "/docker-compose.yml"
ALL_PROFILES := "--profile single --profile cluster --profile cluster4 --profile cluster5 --profile tabletopic --profile analytics"
AZ_NAMES := "az-0 az-1 az-2"
HEAP_OPTS := "-Xms256m -Xmx256m"
KAFKA_OPTS := "KAFKA_HEAP_OPTS='" + HEAP_OPTS + "' KAFKA_JVM_PERFORMANCE_OPTS='' JMX_PORT=''"
# S3 configuration
S3_REGION := "us-east-1"
S3_ENDPOINT := "http://minio:9000"
S3_DATA_BUCKET := "automq-data"
S3_OPS_BUCKET := "automq-ops"
S3_PATH_STYLE := "true"
default:
@just --list
[doc("Alias for just --list")]
help:
@just --list
# ==================== Internal ====================
[private]
ensure-cluster-id:
#!/usr/bin/env bash
mkdir -p "{{DEVKIT}}"
if [ ! -f "{{DEVKIT}}/cluster-id" ]; then
KAFKA_HEAP_OPTS="{{HEAP_OPTS}}" \
"{{AUTOMQ_DEV_HOME}}/bin/kafka-storage.sh" random-uuid \
> "{{DEVKIT}}/cluster-id"
echo "✓ Generated CLUSTER_ID: $(cat "{{DEVKIT}}/cluster-id")"
fi
# Override to skip or replace the Java build step, e.g.:
# DEVKIT_BUILD_CMD="echo skipped" just start-build
JAVA_BUILD_CMD := env_var_or_default("DEVKIT_BUILD_CMD", "./gradlew :core:build :tools:build :shell:build -x test")
[private]
build:
@echo "✓ Building Docker image..."
@{{COMPOSE}} build
@echo "✓ Building Java code..."
cd "{{AUTOMQ_DEV_HOME}}" && {{JAVA_BUILD_CMD}}
[private]
_compute-cluster-vars nodes:
#!/usr/bin/env bash
MAX_CTRL=3
CTRL_COUNT=$(({{nodes}} < MAX_CTRL ? {{nodes}} : MAX_CTRL))
VOTERS="" ; BOOTSTRAP=""
for i in $(seq 0 $((CTRL_COUNT - 1))); do
[ -n "$VOTERS" ] && VOTERS="$VOTERS," && BOOTSTRAP="$BOOTSTRAP,"
VOTERS="${VOTERS}${i}@node-${i}:9093"
BOOTSTRAP="${BOOTSTRAP}node-${i}:9093"
done
echo "CTRL_COUNT=$CTRL_COUNT"
echo "VOTERS=$VOTERS"
echo "BOOTSTRAP=$BOOTSTRAP"
[private]
_render-node node ctrl_count cluster_id voters bootstrap features="":
#!/usr/bin/env bash
set -e
DIR="{{justfile_directory()}}/config"
CFG="{{DEVKIT}}/config/node-{{node}}.properties"
AZS=({{AZ_NAMES}})
export S3_REGION="{{S3_REGION}}"
export S3_ENDPOINT="{{S3_ENDPOINT}}"
export S3_DATA_BUCKET="{{S3_DATA_BUCKET}}"
export S3_OPS_BUCKET="{{S3_OPS_BUCKET}}"
export S3_PATH_STYLE="{{S3_PATH_STYLE}}"
export CLUSTER_ID="{{cluster_id}}"
export VOTERS="{{voters}}"
export BOOTSTRAP="{{bootstrap}}"
export NODE_ID="{{node}}"
export HOST_PORT=$(({{node}} * 10000 + 9092))
export AZ_NAME="${AZS[$(({{node}} % ${#AZS[@]}))]}"
render() { envsubst < "$1" | grep -v '^#' | grep -v '^$'; }
# Layer 1: defaults
envsubst < "$DIR/defaults.properties" > "$CFG"
# Layer 2: role
{ echo ""; echo "# === Role ==="; \
ROLE=$([ {{node}} -lt {{ctrl_count}} ] && echo server || echo broker); \
render "$DIR/role/$ROLE.properties"; } >> "$CFG"
# Layer 3: node identity
{ echo ""; echo "# === Node ==="; render "$DIR/role/node.properties"; } >> "$CFG"
# Layer 4: features
for feat in {{features}}; do
F="$DIR/features/${feat}.properties"
[ -f "$F" ] && { echo ""; render "$F"; } >> "$CFG"
done
# Layer 5: custom overrides (verbatim)
C="$DIR/custom.properties"
[ -f "$C" ] && { echo ""; cat "$C"; } >> "$CFG"
echo "✓ Config: node-{{node}}.properties"
[private]
generate-config nodes features="":
#!/usr/bin/env bash
set -e
command -v envsubst >/dev/null || { echo "✗ envsubst not found: brew install gettext"; exit 1; }
mkdir -p "{{DEVKIT}}/config"
CLUSTER_ID=$(cat "{{DEVKIT}}/cluster-id")
eval "$(just _compute-cluster-vars {{nodes}})"
for i in $(seq 0 $(({{nodes}} - 1))); do
just _render-node "$i" "$CTRL_COUNT" "$CLUSTER_ID" "$VOTERS" "$BOOTSTRAP" "{{features}}"
done
[private]
_node-list:
@docker ps --filter "label=com.automq.devkit=true" --format '{{ '{{' }}.Names{{ '}}' }}'
[private]
_ensure-iptables-chain node:
#!/usr/bin/env bash
docker exec {{node}} iptables -N DEVKIT 2>/dev/null || true
docker exec {{node}} iptables -C OUTPUT -j DEVKIT 2>/dev/null || \
docker exec {{node}} iptables -A OUTPUT -j DEVKIT
docker exec {{node}} iptables -C INPUT -j DEVKIT 2>/dev/null || \
docker exec {{node}} iptables -A INPUT -j DEVKIT
# ==================== Lifecycle ====================
[doc("Start N nodes (default: 1). Example: just start 3 tabletopic")]
start nodes="1" *FEATURES: ensure-cluster-id
#!/usr/bin/env bash
set -e
case {{nodes}} in
1) PROFILE=single ;; 3) PROFILE=cluster ;; 4) PROFILE=cluster4 ;; 5) PROFILE=cluster5 ;;
*) echo "Error: nodes must be 1,3,4,5"; exit 1 ;;
esac
COMPOSE_PROFILES="--profile $PROFILE"
for feat in {{FEATURES}}; do
COMPOSE_PROFILES="$COMPOSE_PROFILES --profile $feat"
done
just generate-config {{nodes}} "{{FEATURES}}"
echo "✓ Starting {{nodes}} node(s)..."
{{COMPOSE}} $COMPOSE_PROFILES up -d
echo ""
echo "🎉 AutoMQ DevKit started"
echo " Kafka: localhost:9092"
echo " Debug: localhost:5005"
echo " MinIO: http://localhost:9001 (admin/password)"
just wait
[doc("Build image & code, then start. Example: just start-build 3 tabletopic")]
start-build nodes="1" *FEATURES: build (start nodes FEATURES)
[doc("Stop all services (keep data)")]
stop:
@echo "Stopping..."
@{{COMPOSE}} {{ALL_PROFILES}} stop
[doc("Stop all services and remove volumes")]
clean:
@echo "Stopping and removing volumes..."
@{{COMPOSE}} {{ALL_PROFILES}} down -v
[doc("Restart containers (no rebuild)")]
restart:
@{{COMPOSE}} {{ALL_PROFILES}} restart
[doc("Rebuild code & image, then restart")]
restart-build: build restart
[doc("Show service status")]
status:
@{{COMPOSE}} {{ALL_PROFILES}} ps
[doc("Show logs snapshot. Example: just logs 0 200")]
logs node="" tail="100":
#!/usr/bin/env bash
if [ -z "{{node}}" ]; then
{{COMPOSE}} {{ALL_PROFILES}} logs --tail={{tail}}
else
{{COMPOSE}} {{ALL_PROFILES}} logs --tail={{tail}} node-{{node}}
fi
[doc("Follow logs (blocks). Example: just logs-follow 0")]
logs-follow node="":
#!/usr/bin/env bash
if [ -z "{{node}}" ]; then
{{COMPOSE}} {{ALL_PROFILES}} logs -f
else
{{COMPOSE}} {{ALL_PROFILES}} logs -f node-{{node}}
fi
[doc("Enter container shell (default: node 0)")]
shell node="0":
@docker exec -it node-{{node}} bash
[doc("Run a single command in container. Example: just exec 0 ls /opt/automq/bin")]
exec node *CMD:
@docker exec node-{{node}} {{CMD}}
[doc("Wait for all nodes to be healthy (auto-called by start)")]
wait timeout="180":
#!/usr/bin/env bash
echo "Waiting for nodes to be healthy..."
TIMEOUT={{timeout}}; ELAPSED=0
while [ $ELAPSED -lt $TIMEOUT ]; do
NODES=$(just _node-list)
[ -z "$NODES" ] && sleep 5 && ELAPSED=$((ELAPSED + 5)) && continue
ALL_HEALTHY=true
for n in $NODES; do
S=$(docker inspect --format='{{ '{{' }}.State.Health.Status{{ '}}' }}' "$n" 2>/dev/null)
[ "$S" != "healthy" ] && ALL_HEALTHY=false && break
done
[ "$ALL_HEALTHY" = true ] && echo "✓ All nodes healthy" && exit 0
sleep 5; ELAPSED=$((ELAPSED + 5))
done
echo "✗ Timeout after ${TIMEOUT}s" && exit 1
[doc("Rebuild Docker image only")]
build-image *args="":
@echo "✓ Building Docker image..."
@{{COMPOSE}} build {{ if args == "" { "node-0 node-1 node-2 node-3 node-4" } else { args } }}
# ==================== Bin ====================
[doc("Run bin/ command in container. Example: just bin kafka-topics.sh --bootstrap-server localhost:9092 --list")]
bin *ARGS:
#!/usr/bin/env bash
NODE=0; PARAMS=()
set -- {{ARGS}}
while [ $# -gt 0 ]; do
case $1 in
--node) NODE=$2; shift 2 ;;
*) PARAMS+=("$1"); shift ;;
esac
done
docker exec "node-${NODE}" \
env "KAFKA_HEAP_OPTS={{HEAP_OPTS}}" KAFKA_JVM_PERFORMANCE_OPTS="" JMX_PORT="" \
/opt/automq/bin/"${PARAMS[@]}" 2> >(grep -v '^SLF4J' >&2)
[doc("Query JMX metrics. Example: just jmx -e domains")]
jmx *ARGS:
#!/usr/bin/env bash
NODE=0
set -- {{ARGS}}
while [ $# -gt 0 ]; do
case $1 in
--node) NODE=$2; shift 2 ;;
*) break ;;
esac
done
echo "$@" | docker exec -i node-${NODE} java -jar /usr/local/bin/jmxterm.jar -l localhost:9999 -n
# ==================== Shortcuts ====================
[doc("List topics")]
topic-list node="0":
@just bin --node {{node}} kafka-topics.sh --bootstrap-server localhost:9092 --list
[doc("Create topic. Example: just topic-create my-topic --partitions 16")]
topic-create name *ARGS:
@just bin kafka-topics.sh --bootstrap-server localhost:9092 --create --topic {{name}} {{ARGS}}
[doc("Describe topic")]
topic-describe name node="0":
@just bin --node {{node}} kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic {{name}}
[doc("Interactive producer. Supports --node N and all other Kafka producer flags")]
produce topic *ARGS:
#!/usr/bin/env bash
NODE=0; PARAMS=()
set -- {{ARGS}}
while [ $# -gt 0 ]; do
case $1 in
--node) NODE=$2; shift 2 ;;
*) PARAMS+=("$1"); shift ;;
esac
done
docker exec -i "node-${NODE}" \
env "KAFKA_HEAP_OPTS={{HEAP_OPTS}}" KAFKA_JVM_PERFORMANCE_OPTS="" JMX_PORT="" \
/opt/automq/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 --topic "{{topic}}" \
"${PARAMS[@]}" 2> >(grep -v '^SLF4J' >&2)
[doc("Consumer. Supports --node N and all other Kafka consumer flags")]
consume topic *ARGS:
#!/usr/bin/env bash
NODE=0; PARAMS=()
set -- {{ARGS}}
while [ $# -gt 0 ]; do
case $1 in
--node) NODE=$2; shift 2 ;;
*) PARAMS+=("$1"); shift ;;
esac
done
docker exec -i "node-${NODE}" \
env "KAFKA_HEAP_OPTS={{HEAP_OPTS}}" KAFKA_JVM_PERFORMANCE_OPTS="" JMX_PORT="" \
/opt/automq/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic "{{topic}}" \
"${PARAMS[@]}" 2> >(grep -v '^SLF4J' >&2)
[doc("List consumer groups")]
group-list:
@just bin kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
[doc("Describe consumer group (offsets, lag). Example: just group-describe my-group")]
group-describe group:
@just bin kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group {{group}}
[doc("Show active members of consumer group. Example: just group-members my-group")]
group-members group:
@just bin kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group {{group}} --members
[doc("Reset consumer group offsets. Example: just group-reset my-group my-topic --to-earliest")]
group-reset group topic *ARGS:
@just bin kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --group {{group}} --topic {{topic}} {{ARGS}} --execute
[doc("List brokers and API versions")]
broker-list:
@just bin kafka-broker-api-versions.sh --bootstrap-server localhost:9092
[doc("Describe broker configs. Example: just broker-config 0")]
broker-config id="":
#!/usr/bin/env bash
if [ -z "{{id}}" ]; then
just bin kafka-configs.sh --bootstrap-server localhost:9092 --describe --entity-type brokers --all
else
just bin kafka-configs.sh --bootstrap-server localhost:9092 --describe --entity-type brokers --entity-name {{id}} --all
fi
# ==================== Perf ====================
[doc("Producer perf test. Example: just perf-produce my-topic --num-records 100000 --record-size 1024 --throughput -1")]
perf-produce topic *ARGS:
@just bin kafka-producer-perf-test.sh --topic {{topic}} --producer-props bootstrap.servers=localhost:9092 {{ARGS}}
[doc("Consumer perf test. Example: just perf-consume my-topic --messages 100000")]
perf-consume topic *ARGS:
@just bin kafka-consumer-perf-test.sh --bootstrap-server localhost:9092 --topic {{topic}} {{ARGS}}
# ==================== Diagnostics ====================
[doc("Attach Arthas to Kafka process. Example: just arthas 1")]
arthas node="0":
@docker exec -it node-{{node}} java -jar /opt/arthas/arthas-boot.jar
[doc("Run a single Arthas command non-interactively. Example: just arthas-exec 0 'dashboard -n 1'")]
arthas-exec node cmd:
#!/usr/bin/env bash
printf '%s\n' '{{cmd}}' | docker exec -i node-{{node}} bash -c \
'cat > /tmp/_arthas_cmd && java -jar /opt/arthas/arthas-boot.jar --select kafka.Kafka -f /tmp/_arthas_cmd'
# ==================== Chaos: Node network ====================
[doc("Network delay on node. Example: just chaos-delay 500 node-1")]
chaos-delay ms="200" node="node-0":
docker exec {{node}} tc qdisc replace dev eth0 root netem delay {{ms}}ms
@echo "✓ {{ms}}ms delay on {{node}}"
[doc("Packet loss on node. Example: just chaos-loss 10 node-1")]
chaos-loss percent="5" node="node-0":
docker exec {{node}} tc qdisc replace dev eth0 root netem loss {{percent}}%
@echo "✓ {{percent}}% packet loss on {{node}}"
# ==================== Chaos: S3 (MinIO) ====================
[doc("S3 latency on a node (only S3 traffic affected). Example: just chaos-s3-delay 500 node-0")]
chaos-s3-delay ms="500" node="node-0":
#!/usr/bin/env bash
MINIO_IP=$(docker exec {{node}} getent hosts minio | awk '{print $1}')
docker exec {{node}} tc qdisc del dev eth0 root 2>/dev/null || true
docker exec {{node}} tc qdisc add dev eth0 root handle 1: prio
docker exec {{node}} tc qdisc add dev eth0 parent 1:3 handle 30: netem delay {{ms}}ms
docker exec {{node}} tc filter add dev eth0 parent 1:0 protocol ip u32 match ip dst $MINIO_IP/32 flowid 1:3
echo "✓ {{ms}}ms delay to MinIO on {{node}} (other traffic unaffected)"
[doc("S3 packet loss on a node (only S3 traffic affected). Example: just chaos-s3-loss 10 node-0")]
chaos-s3-loss percent="5" node="node-0":
#!/usr/bin/env bash
MINIO_IP=$(docker exec {{node}} getent hosts minio | awk '{print $1}')
docker exec {{node}} tc qdisc del dev eth0 root 2>/dev/null || true
docker exec {{node}} tc qdisc add dev eth0 root handle 1: prio
docker exec {{node}} tc qdisc add dev eth0 parent 1:3 handle 30: netem loss {{percent}}%
docker exec {{node}} tc filter add dev eth0 parent 1:0 protocol ip u32 match ip dst $MINIO_IP/32 flowid 1:3
echo "✓ {{percent}}% packet loss to MinIO on {{node}} (other traffic unaffected)"
[doc("Pause MinIO — S3 completely unavailable")]
chaos-s3-down:
docker pause minio
@echo "✓ MinIO (S3) paused — all nodes lost S3 access"
[doc("Resume MinIO")]
chaos-s3-up:
docker unpause minio
@echo "✓ MinIO (S3) resumed"
[doc("Isolate a node from S3. Example: just chaos-s3-partition node-1")]
chaos-s3-partition node="node-0": (_ensure-iptables-chain node)
docker exec {{node}} iptables -A DEVKIT -d minio -j DROP
docker exec {{node}} iptables -A DEVKIT -s minio -j DROP
@echo "✓ {{node}} isolated from MinIO (S3)"
[doc("Restore all nodes' S3 connectivity")]
chaos-s3-partition-reset:
#!/usr/bin/env bash
for n in $(just _node-list); do
docker exec $n iptables -D DEVKIT -d minio -j DROP 2>/dev/null || true
docker exec $n iptables -D DEVKIT -s minio -j DROP 2>/dev/null || true
done
echo "✓ All S3 partitions removed"
# ==================== Chaos: Node-to-node partition ====================
[doc("Isolate two nodes. Example: just chaos-partition node-0 node-1")]
chaos-partition a b: (_ensure-iptables-chain a) (_ensure-iptables-chain b)
docker exec {{a}} iptables -A DEVKIT -d {{b}} -j DROP
docker exec {{a}} iptables -A DEVKIT -s {{b}} -j DROP
docker exec {{b}} iptables -A DEVKIT -d {{a}} -j DROP
docker exec {{b}} iptables -A DEVKIT -s {{a}} -j DROP
@echo "✓ {{a}} ↔ {{b}} isolated"
[doc("Restore all node-to-node connectivity")]
chaos-partition-reset:
#!/usr/bin/env bash
for n in $(just _node-list); do
docker exec $n iptables -F DEVKIT 2>/dev/null || true
done
echo "✓ All node partitions removed"
# ==================== Chaos: Reset & status ====================
[doc("Remove ALL chaos rules (tc + iptables + unpause)")]
chaos-reset:
#!/usr/bin/env bash
for n in $(just _node-list); do
docker exec $n tc qdisc del dev eth0 root 2>/dev/null || true
docker exec $n iptables -F DEVKIT 2>/dev/null || true
done
# Unpause MinIO if paused
docker unpause minio 2>/dev/null || true
echo "✓ All chaos rules reset (tc + iptables + unpause)"
[doc("Show all active chaos rules")]
chaos-status:
#!/usr/bin/env bash
for n in $(just _node-list); do
echo "=== $n ==="
echo " tc:" && docker exec $n tc qdisc show dev eth0
RULES=$(docker exec $n iptables -L DEVKIT -n 2>/dev/null | grep -c DROP || true)
[ "$RULES" -gt 0 ] && echo " iptables: ${RULES} DROP rules in DEVKIT chain"
done
echo "=== minio ==="
S=$(docker inspect --format='{{ '{{' }}.State.Paused{{ '}}' }}' minio 2>/dev/null)
[ "$S" = "true" ] && echo " status: PAUSED" || echo " status: running"