#!/usr/bin/env bash
# ====================================================================
# 第 16 章 - Kafka Connect
# connect_rest_demo.sh - REST API 操作集合（curl + jq）
# --------------------------------------------------------------------
# 默认 Connect REST 端口 8083；如需鉴权，加 -u user:pass
# ====================================================================
set -euo pipefail

CONNECT=${CONNECT:-http://localhost:8083}
NAME=${NAME:-jdbc-source-orders}

hr() { printf '\n\033[36m%s\033[0m\n' "==> $*"; }

hr "0. 集群信息"
curl -s "$CONNECT/" | jq

hr "0.1 列出所有可用 Connector Plugin（看你装了哪些 Connector jar）"
curl -s "$CONNECT/connector-plugins" | jq '.[] | {class, type, version}'

hr "1. 列出所有已注册的 Connector"
curl -s "$CONNECT/connectors" | jq

hr "1.1 列出 Connector + 状态（一行命令看全集群）"
curl -s "$CONNECT/connectors?expand=info&expand=status" | jq

hr "2. 创建 Connector（从 jdbc_source_config.json 加载）"
echo "(示例：curl -X POST $CONNECT/connectors -H 'Content-Type: application/json' -d @jdbc_source_config.json)"

hr "3. 查看某 Connector 的配置"
curl -s "$CONNECT/connectors/$NAME/config" | jq

hr "4. 查看状态（重点看每个 task 是否 RUNNING）"
curl -s "$CONNECT/connectors/$NAME/status" | jq

hr "5. 查看 Source Offset（3.6+ 才有）"
curl -s "$CONNECT/connectors/$NAME/offsets" | jq || echo "(若 404 说明 Connect 版本 < 3.6)"

hr "6. 暂停"
curl -X PUT "$CONNECT/connectors/$NAME/pause"

hr "7. 恢复"
curl -X PUT "$CONNECT/connectors/$NAME/resume"

hr "8. 重启整个 Connector + 所有 task"
curl -X POST "$CONNECT/connectors/$NAME/restart?includeTasks=true&onlyFailed=false" | jq

hr "8.1 仅重启 FAILED 的 task"
curl -X POST "$CONNECT/connectors/$NAME/restart?includeTasks=true&onlyFailed=true" | jq

hr "8.2 重启单个 task（id=0）"
curl -X POST "$CONNECT/connectors/$NAME/tasks/0/restart"

hr "9. 修改配置（PUT，body 只是 config 部分）"
echo "示例："
cat <<EOF
curl -X PUT $CONNECT/connectors/$NAME/config \\
  -H 'Content-Type: application/json' \\
  -d '{
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "tasks.max": "5",
    "...其它原有配置": "..."
  }'
EOF

hr "10. 查看 task 列表 + 配置"
curl -s "$CONNECT/connectors/$NAME/tasks" | jq

hr "11. 验证配置（不真正提交，仅校验）"
echo "示例："
cat <<EOF
curl -X PUT $CONNECT/connector-plugins/io.confluent.connect.jdbc.JdbcSourceConnector/config/validate \\
  -H 'Content-Type: application/json' \\
  -d @jdbc_source_config.json
EOF

hr "12. 删除 Connector（仅删配置 + 状态，Offset 保留在 connect-offsets）"
echo "(危险！) curl -X DELETE $CONNECT/connectors/$NAME"

hr "12.1 重置 Source Offset（3.6+，先删后建可让 connector 从头读）"
echo "curl -X DELETE $CONNECT/connectors/$NAME/offsets"

hr "13. 集群 logger 级别（在线改 log4j 等级，无需重启）"
curl -s "$CONNECT/admin/loggers" | jq
echo "改某 package 等级："
echo "curl -X PUT $CONNECT/admin/loggers/io.debezium -H 'Content-Type: application/json' -d '{\"level\":\"DEBUG\"}'"

# ====================================================================
# 一行命令快速排障：所有 connector 哪个 task 在 FAILED 状态？
# ====================================================================
echo
hr "🔧 Bonus：一行排障 - 列出所有 FAILED task"
curl -s "$CONNECT/connectors?expand=status" | \
  jq -r 'to_entries[] |
         .key as $name |
         .value.status.tasks[] |
         select(.state == "FAILED") |
         "\($name)/task-\(.id)  worker=\(.worker_id)\n  trace: \((.trace // "")[0:200])"'
