#!/usr/bin/env bash
# ====================================================================
# 第 16 章 - Kafka Connect
# start_connect.sh - 本地启动 Distributed Connect Worker
# --------------------------------------------------------------------
# 适用：开发 / 调试 / 教学。生产请用 systemd / k8s 管理 worker 进程。
# ====================================================================
set -euo pipefail

KAFKA_HOME=${KAFKA_HOME:-/opt/kafka}
PLUGIN_PATH=${PLUGIN_PATH:-/opt/kafka/connectors}
BOOTSTRAP=${BOOTSTRAP:-localhost:9092}
WORKER_REST_PORT=${WORKER_REST_PORT:-8083}

CONFIG_FILE=$(mktemp /tmp/connect-distributed.XXXXXX.properties)
trap 'rm -f "$CONFIG_FILE"' EXIT

cat > "$CONFIG_FILE" <<EOF
# ===== Kafka 集群 =====
bootstrap.servers=${BOOTSTRAP}

# ===== Connect Cluster =====
group.id=connect-cluster-learn
rest.host.name=0.0.0.0
rest.port=${WORKER_REST_PORT}
rest.advertised.host.name=$(hostname -i 2>/dev/null || echo localhost)
rest.advertised.port=${WORKER_REST_PORT}

# ===== 三大内部 Topic（生产副本数 >= 3） =====
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
config.storage.replication.factor=1
offset.storage.replication.factor=1
status.storage.replication.factor=1
offset.storage.partitions=25
status.storage.partitions=5

# ===== Converter 默认 JSON（不带 schema，节省体积） =====
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false

# 内部 Converter（Worker 之间通信，必须 JSON）
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false

# ===== Connector jar 包目录（每个 connector 一个子目录） =====
plugin.path=${PLUGIN_PATH}

# ===== Source EOS（Connect 3.3+，可选） =====
# exactly.once.source.support=enabled

# ===== 默认 producer / consumer 配置（可被 connector override） =====
producer.acks=all
producer.compression.type=zstd
producer.linger.ms=10

consumer.fetch.min.bytes=1
consumer.fetch.max.wait.ms=500
EOF

echo "[*] Connect distributed worker 启动中..."
echo "[*] REST API: http://localhost:${WORKER_REST_PORT}"
echo "[*] plugin.path = ${PLUGIN_PATH}"
echo "[*] 临时 config: ${CONFIG_FILE}"
echo

exec "${KAFKA_HOME}/bin/connect-distributed.sh" "${CONFIG_FILE}"

# ====================================================================
# 启动后用以下命令验证：
#   curl -s http://localhost:8083/                | jq
#   curl -s http://localhost:8083/connector-plugins | jq    # 列出可用 plugin
#   curl -s http://localhost:8083/connectors        # 当前已注册 connector
# ====================================================================
