diff --git a/.agents/skills/dev-deploy/SKILL.md b/.agents/skills/dev-deploy/SKILL.md index d5c84b03f..dc56afb1b 100644 --- a/.agents/skills/dev-deploy/SKILL.md +++ b/.agents/skills/dev-deploy/SKILL.md @@ -18,7 +18,7 @@ context: fork - `SERVER_HOST` / `SERVER_LISTEN_HOST`:对外 / 监听地址 - `SERVER_START_PORT`:起始端口;第 `i` 个实例为 `SERVER_START_PORT + i`,即首实例 `7801` - `CLUSTER_ID`、`MDS_INSTANCE_START_ID`、`COORDINATOR_ADDR` -- `STORAGE_ENGINE` / `STORAGE_URL`:由 `deploy_mds.sh` 写进 `mds.conf` +- `STORAGE_ENGINE` / `STORAGE_URL`:由 `operate_mds.sh` 写进 `mds.conf` - `S3_ENDPOINT` / `S3_AK` / `S3_SK` / `S3_BUCKETNAME`、`LOCAL_DATASTORE_PATH`:仅 `create_fs.sh` 使用 ## 步骤 @@ -45,7 +45,7 @@ context: fork 3. **部署并启动 MDS** ```bash - bash clean_start.sh --server_num=$SERVER_NUM + bash operate_mds.sh restart --server_num=$SERVER_NUM ``` 判据:`pgrep -c dingo-mds` 输出等于 `SERVER_NUM`。(不要用 `ps -ef | grep`,它会匹配到自己,也不校验数量。) @@ -95,4 +95,13 @@ bash create_fs.sh --fs_name=$FS_NAME --mds_addr=$SERVER_HOST:$(($SERVER_START_PO ## 其他脚本 -`deploy_mds.sh` / `start_mds.sh` / `stop_mds.sh` 是 `clean_start.sh` 的拆分,只重启 MDS 时单独用。参数以 `bash <脚本> --help` 为准。 +`operate_mds.sh` 一个脚本管全部 MDS 生命周期,子命令选动作,flag 两种写法都行(`--server_num=2 restart` 同 `restart --server_num=2`): + +```bash +bash operate_mds.sh restart --server_num=1 # = stop + deploy + start,等价于原 clean_start.sh +bash operate_mds.sh stop --server_num=1 # 停;--force 用 kill -9,--use_pgrep 按进程名而非 pid 文件 +bash operate_mds.sh deploy --server_num=1 # 只重新部署(软链二进制 + 渲染 conf) +bash operate_mds.sh start --server_num=1 # 只启动 +``` + +参数以 `bash operate_mds.sh --help` 为准。`restart` 中 stop / deploy 失败会直接中止,不会带着半坏状态继续起服务。 diff --git a/.github/scripts/deploy-mds-client.sh b/.github/scripts/deploy-mds-client.sh index 0cd6faac6..7aa9a6269 100755 --- a/.github/scripts/deploy-mds-client.sh +++ b/.github/scripts/deploy-mds-client.sh @@ -18,6 +18,8 @@ CLUSTER_ID=101 MDS_PORT=8821 MDS_INSTANCE_ID=1001 VFS_DUMMY_PORT=10001 +MDS_LOG_LEVEL=DEBUG +MDS_LOG_V=20 BUILD_DIR="${GITHUB_WORKSPACE}/build/bin" TEMPLATE="${GITHUB_WORKSPACE}/scripts/dev-mds/mds.template.conf" @@ -54,6 +56,8 @@ sed -i "s|\\\$SERVER_PORT|${MDS_PORT}|g" "${DIST_CONF}" sed -i "s|\\\$BASE_PATH|${MDS_DIST}|g" "${DIST_CONF}" sed -i "s|\\\$STORAGE_ENGINE|dingo-store|g" "${DIST_CONF}" sed -i "s|\\\$STORAGE_URL|list://${COORDINATOR_ADDR}|g" "${DIST_CONF}" +sed -i "s|\\\$LOG_LEVEL|${MDS_LOG_LEVEL}|g" "${DIST_CONF}" +sed -i "s|\\\$LOG_V|${MDS_LOG_V}|g" "${DIST_CONF}" echo "${COORDINATOR_ADDR}" > "${MDS_DIST}/conf/coor_list" diff --git a/.gitignore b/.gitignore index 41afbaf6c..e04325daf 100755 --- a/.gitignore +++ b/.gitignore @@ -115,3 +115,6 @@ bench_results/ CONTEXT.md .pi-subagents/ + +# Local dev config (contains secrets) +scripts/dev-mds/env.local diff --git a/scripts/deploy/deploy.sh b/scripts/deploy/deploy.sh index 0e0de4c4b..4704a0764 100755 --- a/scripts/deploy/deploy.sh +++ b/scripts/deploy/deploy.sh @@ -76,6 +76,8 @@ function deploy_server() { sed -i 's,\$BASE_PATH,'"$dstpath"',g' $dist_conf sed -i 's,\$STORAGE_ENGINE,dingo-store,g' $dist_conf sed -i 's,\$STORAGE_URL,file://./conf/coor_list,g' $dist_conf + sed -i 's,\$LOG_LEVEL,INFO,g' $dist_conf + sed -i 's,\$LOG_V,0,g' $dist_conf # coor_list file coor_file="${dstpath}/conf/coor_list" diff --git a/scripts/dev-mds/clean_start.sh b/scripts/dev-mds/clean_start.sh deleted file mode 100755 index 5114c8194..000000000 --- a/scripts/dev-mds/clean_start.sh +++ /dev/null @@ -1,26 +0,0 @@ -#!/bin/bash - -mydir="${BASH_SOURCE%/*}" -if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi -. $mydir/shflags - -DEFINE_integer server_num 1 'server number' - -# parse the command-line -FLAGS "$@" || exit 1 -eval set -- "${FLAGS_ARGV}" - - -echo "============ stop ============" -./stop_mds.sh --server_num=${FLAGS_server_num} - -sleep 1 -echo "============ deploy ============" -./deploy_mds.sh --server_num=${FLAGS_server_num} - -sleep 1 -echo "============ start ============" -./start_mds.sh --server_num=${FLAGS_server_num} - -sleep 1 -echo "============ done ============" \ No newline at end of file diff --git a/scripts/dev-mds/create_cluster.sh b/scripts/dev-mds/create_cluster.sh index e6f5bad15..c1586f1e6 100755 --- a/scripts/dev-mds/create_cluster.sh +++ b/scripts/dev-mds/create_cluster.sh @@ -5,7 +5,7 @@ if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi . $mydir/shflags DEFINE_integer cluster_id 0 'cluster id' -DEFINE_string parameters 'mds_deploy_parameters.local' 'deploy parameters file' +DEFINE_string env 'env.local' 'deploy env file' # parse the command-line @@ -14,7 +14,7 @@ eval set -- "${FLAGS_ARGV}" echo "cluster_id: ${FLAGS_cluster_id}" -source $mydir/${FLAGS_parameters} +source $mydir/${FLAGS_env} #check cluster id is valid diff --git a/scripts/dev-mds/create_fs.sh b/scripts/dev-mds/create_fs.sh index b14759b97..b85910941 100755 --- a/scripts/dev-mds/create_fs.sh +++ b/scripts/dev-mds/create_fs.sh @@ -6,13 +6,19 @@ if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi DEFINE_string fs_name '' 'fs name' DEFINE_string mds_addr '' 'mds address' -DEFINE_string parameters 'mds_deploy_parameters.local' 'deploy parameters file' +DEFINE_string env 'env.local' 'deploy env file' DEFINE_boolean use_local_datastore false 'use local datastore' # parse the command-line FLAGS "$@" || exit 1 eval set -- "${FLAGS_ARGV}" +source $mydir/${FLAGS_env} + +if [ -z "${FLAGS_mds_addr}" ]; then + FLAGS_mds_addr="${SERVER_HOST}:$((SERVER_START_PORT + 1))" +fi + echo "fs_name: ${FLAGS_fs_name}" echo "mds_addr: ${FLAGS_mds_addr}" @@ -27,8 +33,6 @@ if [ -z "${FLAGS_mds_addr}" ]; then exit -1 fi -source $mydir/${FLAGS_parameters} - BASE_DIR=$(dirname $(dirname $(cd $(dirname $0); pwd))) BUILD_DIR=$BASE_DIR/build MDS_CLIENT_BIN_PATH=$BUILD_DIR/bin/dingo-mds-client diff --git a/scripts/dev-mds/deploy_mds.sh b/scripts/dev-mds/deploy_mds.sh deleted file mode 100755 index cbbf7fec9..000000000 --- a/scripts/dev-mds/deploy_mds.sh +++ /dev/null @@ -1,105 +0,0 @@ -#!/bin/bash - -mydir="${BASH_SOURCE%/*}" -if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi -. $mydir/shflags - -DEFINE_integer server_num 1 'server number' -DEFINE_boolean clean_log 1 'clean log' -DEFINE_boolean replace_conf 0 'replace conf' -DEFINE_string parameters 'mds_deploy_parameters.local' 'deploy parameters file' - -# parse the command-line -FLAGS "$@" || exit 1 -eval set -- "${FLAGS_ARGV}" -echo "parameters: ${FLAGS_parameters}" - -BASE_DIR=$(dirname $(dirname $(cd $(dirname $0); pwd))) -DIST_DIR=$BASE_DIR/dist - -# validate BASE_DIR and BASE_DIR/src and BASE_DIR/build -if [ ! -d "$BASE_DIR" ] || [ ! -d "$BASE_DIR/src" ] || [ ! -d "$BASE_DIR/build" ]; then - echo "error: script run dir wrong, please run this script in scripts/dev-mds dir." - exit 1 -fi - -SERVER_NAME=mds -SERVER_BIN_NAME=dingo-mds -MDS_CLIENT_BIN_NAME=dingo-mds-client - -if [ ! -d "$DIST_DIR" ]; then - mkdir "$DIST_DIR" -fi - -source $mydir/${FLAGS_parameters} - - -function deploy_server() { - srcpath=$1 - dstpath=$2 - instance_id=$3 - server_port=$4 - - echo "server $dstpath $instance_id $server_port" - - if [ ! -d "$dstpath" ]; then - mkdir "$dstpath" - fi - - if [ ! -d "$dstpath/bin" ]; then - mkdir "$dstpath/bin" - fi - if [ ! -d "$dstpath/conf" ]; then - mkdir "$dstpath/conf" - fi - if [ ! -d "$dstpath/log" ]; then - mkdir "$dstpath/log" - fi - - - # hard link server binary - if [ -f "${dstpath}/bin/${SERVER_BIN_NAME}" ]; then - rm -f "${dstpath}/bin/${SERVER_BIN_NAME}" - fi - ln -s "${srcpath}/build/bin/${SERVER_BIN_NAME}" "${dstpath}/bin/${SERVER_BIN_NAME}" - - # link dingo-mds-client - if [ -f "${dstpath}/bin/${MDS_CLIENT_BIN_NAME}" ]; then - rm -f "${dstpath}/bin/${MDS_CLIENT_BIN_NAME}" - fi - ln -s "${srcpath}/build/bin/${MDS_CLIENT_BIN_NAME}" "${dstpath}/bin/${MDS_CLIENT_BIN_NAME}" - - - - if [ "${FLAGS_replace_conf}" == "0" ]; then - # conf file - dist_conf="${dstpath}/conf/${SERVER_NAME}.conf" - cp $srcpath/scripts/dev-mds/${SERVER_NAME}.template.conf $dist_conf - - sed -i 's,\$CLUSTER_ID,'"$CLUSTER_ID"',g' $dist_conf - sed -i 's,\$INSTANCE_ID,'"$instance_id"',g' $dist_conf - sed -i 's,\$SERVER_HOST,'"$SERVER_HOST"',g' $dist_conf - sed -i 's,\$SERVER_LISTEN_HOST,'"$SERVER_LISTEN_HOST"',g' $dist_conf - sed -i 's,\$SERVER_PORT,'"$server_port"',g' $dist_conf - sed -i 's,\$BASE_PATH,'"$dstpath"',g' $dist_conf - sed -i 's,\$STORAGE_ENGINE,'"$STORAGE_ENGINE"',g' $dist_conf - sed -i 's,\$STORAGE_URL,'"$STORAGE_URL"',g' $dist_conf - - # coor_list file - coor_file="${dstpath}/conf/coor_list" - echo $COORDINATOR_ADDR > $coor_file - - fi - - if [ "${FLAGS_clean_log}" != "0" ]; then - rm -rf $dstpath/log/* - fi -} - -for ((i=1; i<=$FLAGS_server_num; ++i)); do - instance_dist_dir=$DIST_DIR/$SERVER_NAME-$i - - deploy_server ${BASE_DIR} ${instance_dist_dir} `expr ${MDS_INSTANCE_START_ID} + ${i}` `expr ${SERVER_START_PORT} + ${i}` -done - -echo "deploy finish..." diff --git a/scripts/dev-mds/mds_deploy_parameters b/scripts/dev-mds/env.example similarity index 81% rename from scripts/dev-mds/mds_deploy_parameters rename to scripts/dev-mds/env.example index af66e9993..42f451371 100755 --- a/scripts/dev-mds/mds_deploy_parameters +++ b/scripts/dev-mds/env.example @@ -12,4 +12,10 @@ STORAGE_URL=list://127.0.0.1:2379 S3_ENDPOINT= S3_AK= S3_SK= -S3_BUCKETNAME= \ No newline at end of file +S3_BUCKETNAME= + + +LOCAL_DATASTORE_PATH= + +LOG_LEVEL=DEBUG +LOG_V=0 diff --git a/scripts/dev-mds/mds.template.conf b/scripts/dev-mds/mds.template.conf index af1d22701..9d8d7518a 100644 --- a/scripts/dev-mds/mds.template.conf +++ b/scripts/dev-mds/mds.template.conf @@ -27,8 +27,8 @@ # log --log_dir=$BASE_PATH/log ---log_level=DEBUG ---log_v=0 +--log_level=$LOG_LEVEL +--log_v=$LOG_V # dingo-sdk diff --git a/scripts/dev-mds/operate_mds.sh b/scripts/dev-mds/operate_mds.sh new file mode 100755 index 000000000..8ee39227e --- /dev/null +++ b/scripts/dev-mds/operate_mds.sh @@ -0,0 +1,298 @@ +#!/bin/bash +# operate_mds.sh -- 开发环境 MDS 运维脚本 +# 合并自 deploy_mds.sh / start_mds.sh / stop_mds.sh / clean_start.sh +# +# usage: operate_mds.sh {stop|deploy|start|restart} [flags] + +mydir="${BASH_SOURCE%/*}" +if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi +. $mydir/shflags + +DEFINE_integer server_num 3 'server number' +DEFINE_boolean clean_log true 'clean log' +DEFINE_boolean replace_conf false 'replace conf' +DEFINE_string env 'env.local' 'deploy env file' +DEFINE_boolean force true 'use kill -9 to stop' +DEFINE_boolean use_pgrep true 'use pgrep to get pid' + +FLAGS_HELP="usage: operate_mds.sh {stop|deploy|start|restart} [flags] + + stop 停止 MDS 实例(默认按 dist/mds-/log/pid,--use_pgrep 时按进程名) + deploy 重新生成 dist/mds-:软链二进制、渲染 mds.conf + start 启动 MDS 实例 + restart 等价于 stop + deploy + start +" + +# parse the command-line +FLAGS "$@" || exit 1 +eval set -- "${FLAGS_ARGV}" + +SERVER_NAME=mds +SERVER_BIN_NAME=dingo-mds +MDS_CLIENT_BIN_NAME=dingo-mds-client + +BASE_DIR=$(dirname $(dirname $(cd $(dirname $0); pwd))) +DIST_DIR=$BASE_DIR/dist + +function wait_for_process_exit() { + local pid_killed=$1 + local begin=$(date +%s) + local end + while kill -0 $pid_killed > /dev/null 2>&1 + do + echo -n "." + sleep 1; + end=$(date +%s) + if [ $((end-begin)) -gt 60 ];then + echo -e "\nTimeout" + return 1 + fi + done + return 0 +} + +function do_stop() { + echo "============ stop ============" + echo "stop server num(${FLAGS_server_num})" + + if [ "${FLAGS_use_pgrep}" -eq "${FLAGS_TRUE}" ]; then + # Match on the conf path, not BASE_DIR: BASE_DIR depends on how the + # script was invoked (symlinked vs real path), while the started + # process always carries --conf=.../dist/mds-/conf/mds.conf. + process_no=$(pgrep -f -U `id -u` -- "--conf=.*/dist/${SERVER_NAME}-[0-9]+/conf/${SERVER_NAME}\.conf" | xargs) + + if [ "${process_no}" != "" ]; then + echo "pid to kill: ${process_no}" + if [ "${FLAGS_force}" -eq "${FLAGS_TRUE}" ] + then + kill -9 ${process_no} + else + kill ${process_no} + fi + + wait_for_process_exit ${process_no} || return 1 + else + echo "not exist ${SERVER_NAME} process" + fi + else + for ((i=1; i<=$FLAGS_server_num; ++i)); do + pid_file=$DIST_DIR/${SERVER_NAME}-${i}/log/pid + + # Check if the PID file exists + if [ -f "$pid_file" ]; then + # Read the PID from the file + pid=$(<"$pid_file") + + # Check if the PID is a number + if [[ "$pid" =~ ^[0-9]+$ ]]; then + # Kill the process with the specified PID + if [ "${FLAGS_force}" -eq "${FLAGS_TRUE}" ] + then + echo "killing -9 process with pid($pid) on $pid_file" + kill -9 ${pid} + else + echo "killing process with pid($pid) on $pid_file" + kill ${pid} + fi + + wait_for_process_exit ${pid} || return 1 + else + echo "invalid pid($pid) on $pid_file" + fi + else + echo "not found $pid_file" + fi + done + fi + + echo "stop finish..." +} + +function deploy_server() { + srcpath=$1 + dstpath=$2 + instance_id=$3 + server_port=$4 + + echo "server $dstpath $instance_id $server_port" + + if [ ! -d "$dstpath" ]; then + mkdir "$dstpath" + fi + + if [ ! -d "$dstpath/bin" ]; then + mkdir "$dstpath/bin" + fi + if [ ! -d "$dstpath/conf" ]; then + mkdir "$dstpath/conf" + fi + if [ ! -d "$dstpath/log" ]; then + mkdir "$dstpath/log" + fi + + + # server binary and dingo-mds-client are symlinks into build/bin; + # rm -f (not [ -f ] && rm) because a dangling symlink is not -f. + rm -f "${dstpath}/bin/${SERVER_BIN_NAME}" "${dstpath}/bin/${MDS_CLIENT_BIN_NAME}" + ln -s "${srcpath}/build/bin/${SERVER_BIN_NAME}" "${dstpath}/bin/${SERVER_BIN_NAME}" + ln -s "${srcpath}/build/bin/${MDS_CLIENT_BIN_NAME}" "${dstpath}/bin/${MDS_CLIENT_BIN_NAME}" + + + if [ "${FLAGS_replace_conf}" -eq "${FLAGS_TRUE}" ]; then + # conf file + dist_conf="${dstpath}/conf/${SERVER_NAME}.conf" + cp $srcpath/scripts/dev-mds/${SERVER_NAME}.template.conf $dist_conf + + sed -i 's,\$CLUSTER_ID,'"$CLUSTER_ID"',g' $dist_conf + sed -i 's,\$INSTANCE_ID,'"$instance_id"',g' $dist_conf + sed -i 's,\$SERVER_HOST,'"$SERVER_HOST"',g' $dist_conf + sed -i 's,\$SERVER_LISTEN_HOST,'"$SERVER_LISTEN_HOST"',g' $dist_conf + sed -i 's,\$SERVER_PORT,'"$server_port"',g' $dist_conf + sed -i 's,\$BASE_PATH,'"$dstpath"',g' $dist_conf + sed -i 's,\$STORAGE_ENGINE,'"$STORAGE_ENGINE"',g' $dist_conf + sed -i 's,\$STORAGE_URL,'"$STORAGE_URL"',g' $dist_conf + sed -i 's,\$LOG_LEVEL,'"$LOG_LEVEL"',g' $dist_conf + sed -i 's,\$LOG_V,'"$LOG_V"',g' $dist_conf + + # coor_list file + coor_file="${dstpath}/conf/coor_list" + echo $COORDINATOR_ADDR > $coor_file + + fi + + if [ "${FLAGS_clean_log}" -eq "${FLAGS_TRUE}" ]; then + rm -rf $dstpath/log/* + fi +} + +function do_deploy() { + echo "============ deploy ============" + echo "env: ${FLAGS_env}" + + # validate BASE_DIR and BASE_DIR/src and BASE_DIR/build + if [ ! -d "$BASE_DIR" ] || [ ! -d "$BASE_DIR/src" ] || [ ! -d "$BASE_DIR/build" ]; then + echo "error: script run dir wrong, please run this script in scripts/dev-mds dir." + return 1 + fi + + if [ ! -f "$mydir/${FLAGS_env}" ]; then + echo "error: env file not found: $mydir/${FLAGS_env}" + return 1 + fi + source $mydir/${FLAGS_env} || return 1 + + if [ ! -d "$DIST_DIR" ]; then + mkdir "$DIST_DIR" + fi + + for ((i=1; i<=$FLAGS_server_num; ++i)); do + instance_dist_dir=$DIST_DIR/$SERVER_NAME-$i + + deploy_server ${BASE_DIR} ${instance_dist_dir} `expr ${MDS_INSTANCE_START_ID} + ${i}` `expr ${SERVER_START_PORT} + ${i}` + done + + echo "deploy finish..." +} + +function set_ulimit() { + NUM_FILE=1048576 + NUM_PROC=4194304 + + # 1. sysctl is the very-high-level hard limit: + # fs.nr_open = 1048576 + # fs.file-max = 4194304 + # 2. /etc/security/limits.conf is the second-level limit for users, this is not required to setup. + # CAUTION: values in limits.conf can't bigger than sysctl kernel values, or user login will fail. + # * - nofile 1048576 + # * - nproc 4194304 + # 3. we can use ulimit to set value before start service. + # ulimit -n 1048576 + # ulimit -u 4194304 + # ulimit -c unlimited + + # ulimit -n + nfile=$(ulimit -n) + echo "nfile="${nfile} + if [ ${nfile} -lt ${NUM_FILE} ] + then + echo "try to increase nfile" + ulimit -n ${NUM_FILE} + + nfile=$(ulimit -n) + echo "nfile new="${nfile} + if [ ${nfile} -lt ${NUM_FILE} ] + then + echo "need to increase nfile to ${NUM_FILE}, exit!" + exit -1 + fi + fi + + # ulimit -c + ncore=$(ulimit -c) + echo "ncore="${ncore} + if [ ${ncore} != "unlimited" ] + then + echo "try to set ulimit -c unlimited" + ulimit -c unlimited + + ncore=$(ulimit -c) + echo "ncore new="${ncore} + if [ ${ncore} != "unlimited" ] + then + echo "need to set ulimit -c unlimited, exit!" + exit -1 + fi + fi +} + +function start_server() { + root_dir=$1 + + set_ulimit + + cd ${root_dir} + + + echo "start server: ${root_dir}/bin/${SERVER_BIN_NAME}" + + ${root_dir}/bin/${SERVER_BIN_NAME} --daemonize=true --conf=${root_dir}/conf/${SERVER_NAME}.conf 2>&1 >./log/out & +} + +function do_start() { + echo "============ start ============" + echo "start server num(${FLAGS_server_num})" + + for ((i=1; i<=${FLAGS_server_num}; ++i)); do + ininstance_dist_dir=$DIST_DIR/$SERVER_NAME-$i + + start_server ${ininstance_dist_dir} + done + + echo "start finish..." +} + +case "${1:-}" in + stop) + do_stop + ;; + deploy) + do_deploy + ;; + start) + do_start + ;; + restart) + do_stop || exit 1 + sleep 1 + do_deploy || exit 1 + sleep 1 + do_start + sleep 1 + echo "============ done ============" + ;; + *) + echo "error: unknown action '${1:-}'" >&2 + flags_help + exit 1 + ;; +esac diff --git a/scripts/dev-mds/start_cache.sh b/scripts/dev-mds/start_cache.sh index 07a5abff9..2cfad9729 100755 --- a/scripts/dev-mds/start_cache.sh +++ b/scripts/dev-mds/start_cache.sh @@ -12,6 +12,8 @@ DEFINE_integer force 1 'use kill -9 to stop' DEFINE_boolean stop false 'just stop client, do not start' DEFINE_boolean clean_log false 'clean log' DEFINE_integer port 39000 'server listen port' +DEFINE_string log_level INFO 'cache log level' +DEFINE_integer log_v 0 'cache log v' # parse the command-line @@ -63,7 +65,8 @@ function start() { --cache_dir=${FLAGS_cache_dir} \ --cache_size_mb=1048576 \ --log_dir=${log_dir} \ - --log_level=INFO \ + --log_level=${FLAGS_log_level} \ + --log_v=${FLAGS_log_v} \ --daemonize=true 2>&1 > $log_dir/out } diff --git a/scripts/dev-mds/start_client.sh b/scripts/dev-mds/start_client.sh index 4aeb9fd36..ca03183a0 100755 --- a/scripts/dev-mds/start_client.sh +++ b/scripts/dev-mds/start_client.sh @@ -15,6 +15,8 @@ DEFINE_boolean loop false 'loop restart client' DEFINE_boolean clean_log false 'clean log' DEFINE_integer port 11000 'dummy server port' DEFINE_boolean use_cache false 'use cache' +DEFINE_string log_level DEBUG 'client log level' +DEFINE_integer log_v 20 'client log v' # parse the command-line FLAGS "$@" || exit 1 @@ -154,8 +156,8 @@ function start() { ${CLIENT_BIN_PATH} ${FLAGS_meta} ${mountpoint_dir} \ --fuse_subdir=/ \ --log_dir=${log_dir} \ - --log_level=DEBUG \ - --log_v=20 \ + --log_level=${FLAGS_log_level} \ + --log_v=${FLAGS_log_v} \ --vfs_dummy_server_port=${dummy_port} \ --cache_store=none \ --fill_group_cache=False \ @@ -169,8 +171,8 @@ function start() { ${CLIENT_BIN_PATH} ${FLAGS_meta} ${mountpoint_dir} \ --fuse_subdir=/ \ --log_dir=${log_dir} \ - --log_level=DEBUG \ - --log_v=20 \ + --log_level=${FLAGS_log_level} \ + --log_v=${FLAGS_log_v} \ --vfs_dummy_server_port=${dummy_port} \ --cache_store=none \ --daemonize=true 2>&1 > $log_dir/out diff --git a/scripts/dev-mds/start_mds.sh b/scripts/dev-mds/start_mds.sh deleted file mode 100755 index 2ab841e0f..000000000 --- a/scripts/dev-mds/start_mds.sh +++ /dev/null @@ -1,98 +0,0 @@ -#!/bin/bash -# The ulimit is setup in start_server, the paramter of ulimit is: -# ulimit -n 1048576 -# ulimit -u 4194304 -# ulimit -c unlimited -# If set ulimit failed, please use root or sudo to execute sysctl.sh to increase kernal limit. - -mydir="${BASH_SOURCE%/*}" -if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi -. $mydir/shflags - - -DEFINE_integer server_num 1 'server number' - -# parse the command-line -FLAGS "$@" || exit 1 -eval set -- "${FLAGS_ARGV}" - -echo "start server num(${FLAGS_server_num})" - -BASE_DIR=$(dirname $(dirname $(cd $(dirname $0); pwd))) -DIST_DIR=$BASE_DIR/dist - -SERVER_NAME=mds -SERVER_BIN_NAME=dingo-mds - -function set_ulimit() { - NUM_FILE=1048576 - NUM_PROC=4194304 - - # 1. sysctl is the very-high-level hard limit: - # fs.nr_open = 1048576 - # fs.file-max = 4194304 - # 2. /etc/security/limits.conf is the second-level limit for users, this is not required to setup. - # CAUTION: values in limits.conf can't bigger than sysctl kernel values, or user login will fail. - # * - nofile 1048576 - # * - nproc 4194304 - # 3. we can use ulimit to set value before start service. - # ulimit -n 1048576 - # ulimit -u 4194304 - # ulimit -c unlimited - - # ulimit -n - nfile=$(ulimit -n) - echo "nfile="${nfile} - if [ ${nfile} -lt ${NUM_FILE} ] - then - echo "try to increase nfile" - ulimit -n ${NUM_FILE} - - nfile=$(ulimit -n) - echo "nfile new="${nfile} - if [ ${nfile} -lt ${NUM_FILE} ] - then - echo "need to increase nfile to ${NUM_FILE}, exit!" - exit -1 - fi - fi - - # ulimit -c - ncore=$(ulimit -c) - echo "ncore="${ncore} - if [ ${ncore} != "unlimited" ] - then - echo "try to set ulimit -c unlimited" - ulimit -c unlimited - - ncore=$(ulimit -c) - echo "ncore new="${ncore} - if [ ${ncore} != "unlimited" ] - then - echo "need to set ulimit -c unlimited, exit!" - exit -1 - fi - fi -} - -function start_server() { - root_dir=$1 - - set_ulimit - - cd ${root_dir} - - - echo "start server: ${root_dir}/bin/${SERVER_BIN_NAME}" - - ${root_dir}/bin/${SERVER_BIN_NAME} --daemonize=true --conf=${root_dir}/conf/${SERVER_NAME}.conf 2>&1 >./log/out & -} - - -for ((i=1; i<=${FLAGS_server_num}; ++i)); do - ininstance_dist_dir=$DIST_DIR/$SERVER_NAME-$i - - start_server ${ininstance_dist_dir} -done - -echo "start finish..." \ No newline at end of file diff --git a/scripts/dev-mds/stop_mds.sh b/scripts/dev-mds/stop_mds.sh deleted file mode 100755 index df249fc7a..000000000 --- a/scripts/dev-mds/stop_mds.sh +++ /dev/null @@ -1,88 +0,0 @@ -#!/bin/bash - -mydir="${BASH_SOURCE%/*}" -if [[ ! -d "$mydir" ]]; then mydir="$PWD"; fi -. $mydir/shflags - -DEFINE_integer server_num 1 'server number' -DEFINE_integer force 0 'use kill -9 to stop' -DEFINE_integer use_pgrep 0 'use pgrep to get pid' - -# parse the command-line -FLAGS "$@" || exit 1 -eval set -- "${FLAGS_ARGV}" - -echo "stop server num(${FLAGS_server_num})" - -BASE_DIR=$(dirname $(dirname $(cd $(dirname $0); pwd))) -DIST_DIR=$BASE_DIR/dist - -SERVER_NAME=mds -SERVER_BIN_NAME=dingo-mds - - -wait_for_process_exit() { - local pid_killed=$1 - local begin=$(date +%s) - local end - while kill -0 $pid_killed > /dev/null 2>&1 - do - echo -n "." - sleep 1; - end=$(date +%s) - if [ $((end-begin)) -gt 60 ];then - echo -e "\nTimeout" - break; - fi - done -} - - -if [ ${FLAGS_use_pgrep} -ne 0 ]; then - process_no=$(pgrep -f "${BASE_DIR}.*${SERVER_BIN_NAME}" -U `id -u` | xargs) - - if [ "${process_no}" != "" ]; then - echo "pid to kill: ${process_no}" - if [ ${FLAGS_force} -eq 0 ] - then - kill ${process_no} - else - kill -9 ${process_no} - fi - - wait_for_process_exit ${process_no} - else - echo "not exist ${SERVER_NAME} process" - fi -else - for ((i=1; i<=$FLAGS_server_num; ++i)); do - pid_file=$DIST_DIR/${SERVER_NAME}-${i}/log/pid - - # Check if the PID file exists - if [ -f "$pid_file" ]; then - # Read the PID from the file - pid=$(<"$pid_file") - - # Check if the PID is a number - if [[ "$pid" =~ ^[0-9]+$ ]]; then - # Kill the process with the specified PID - if [ ${FLAGS_force} -eq 0 ] - then - echo "killing process with pid($pid) on $pid_file" - kill ${pid} - else - echo "killing -9 process with pid($pid) on $pid_file" - kill -9 ${pid} - fi - - wait_for_process_exit ${pid} - else - echo "invalid pid($pid) on $pid_file" - fi - else - echo "not found $pid_file" - fi - done -fi - -echo "stop finish..." diff --git a/src/client/vfs/metasystem/local/metasystem.cc b/src/client/vfs/metasystem/local/metasystem.cc index 802408eb4..0c5599080 100644 --- a/src/client/vfs/metasystem/local/metasystem.cc +++ b/src/client/vfs/metasystem/local/metasystem.cc @@ -1964,7 +1964,7 @@ bool LocalMetaSystem::InitCrontab() { "CLEAN_DELFILE", kCleanDelfileIntervalS * 1000, true, - [this](void*) { this->CleanDelfile(); }, + [this]() { this->CleanDelfile(); }, }); crontab_manager_.AddCrontab(crontab_configs_); diff --git a/src/client/vfs/metasystem/mds/metasystem.cc b/src/client/vfs/metasystem/mds/metasystem.cc index 8079d7c80..6756960df 100644 --- a/src/client/vfs/metasystem/mds/metasystem.cc +++ b/src/client/vfs/metasystem/mds/metasystem.cc @@ -494,7 +494,7 @@ bool MDSMetaSystem::InitCrontab() { "HEARTBEAT", kHeartbeatIntervalS * 1000, true, - [this](void*) { this->Heartbeat(); }, + [this]() { this->Heartbeat(); }, }); // add clean expired crontab @@ -502,7 +502,7 @@ bool MDSMetaSystem::InitCrontab() { "CLEAN_EXPIRED", kCleanExpiredModifyTimeMemoIntervalS * 1000, true, - [this](void*) { this->CleanExpired(); }, + [this]() { this->CleanExpired(); }, }); // Note: the cached fs_info refreshes via the 5s heartbeat path — MDS echoes diff --git a/src/mds/common/crontab.cc b/src/mds/common/crontab.cc index d6f1b583b..fbc6cf186 100644 --- a/src/mds/common/crontab.cc +++ b/src/mds/common/crontab.cc @@ -16,212 +16,208 @@ #include +#include +#include + #include "bthread/bthread.h" #include "bthread/unstable.h" -#include "common/logging.h" +#include "butil/time.h" #include "fmt/core.h" namespace dingofs { namespace mds { -void Crontab::DescribeByJson(Json::Value& value) const { - value["id"] = id; - value["name"] = name; - value["interval_ms"] = interval; - value["max_times"] = max_times; - value["immediately"] = immediately; - value["run_count"] = run_count; - value["pause"] = pause.load(); +Crontab::Crontab(uint32_t id, CrontabConfig cfg) : id_(id), cfg_(std::move(cfg)) { + CHECK_EQ(bthread_mutex_init(&mu_, nullptr), 0); + CHECK_EQ(bthread_cond_init(&drained_cv_, nullptr), 0); } -CrontabManager::CrontabManager() { bthread_mutex_init(&mutex_, nullptr); } - -CrontabManager::~CrontabManager() { - // Crontab::Run() reschedules itself via bthread_timer_add() using a raw - // Crontab* (see below). If crontabs_ (and the Crontabs it keeps alive) - // were destroyed without first cancelling those pending timers, the - // timer thread could still invoke Run() on a freed Crontab later, - // causing a use-after-free. Destroy() cancels every pending timer and - // waits for in-flight callbacks to finish before we let crontabs_ go. - Stop(); - bthread_mutex_destroy(&mutex_); +Crontab::~Crontab() { + CHECK_EQ(pending_ops_, 0) << cfg_.name; + bthread_cond_destroy(&drained_cv_); + bthread_mutex_destroy(&mu_); } -void CrontabManager::Run(void* arg) { - Crontab* crontab = static_cast(arg); - if (crontab->pause) { - return; - } - if (crontab->immediately) { - try { - crontab->func(crontab->arg); - } catch (...) { - LOG(ERROR) << fmt::format("[crontab.run][id({}).name({})] crontab happen exception", crontab->id, crontab->name); - } - ++crontab->run_count; - } else { - crontab->immediately = true; - } +void Crontab::Launch() { + std::lock_guard lock(mu_); + if (stopped_) return; + + if (cfg_.immediately) { + StartRoutine(); + if (cfg_.async) ArmLocked(); - // Re-check pause: Destroy()/PauseCrontab() may have paused this crontab - // while func() above was executing. Without this check we could keep - // rearming a timer after the owning CrontabManager decided to stop (and - // is waiting to release the Crontab), leading to a use-after-free once - // the manager is destroyed. - if (!crontab->pause && (crontab->max_times == 0 || crontab->run_count < crontab->max_times)) { - bthread_timer_add(&crontab->timer_id, butil::milliseconds_from_now(crontab->interval), &Run, crontab); + } else { + ArmLocked(); } } -uint32_t CrontabManager::AllocCrontabId() { return auinc_crontab_id_.fetch_add(1); } - -void CrontabManager::AddCrontab(std::vector& crontab_configs) { - for (auto& crontab_config : crontab_configs) { - LOG(INFO) << fmt::format("[crontab.add][name({}).interval({}ms).async({})] add crontab task.", crontab_config.name, - crontab_config.interval, crontab_config.async); - - auto crontab = std::make_shared(); - crontab->name = crontab_config.name; - crontab->interval = crontab_config.interval; - if (crontab_config.async) { - crontab->func = [this, &crontab_config](void*) { - // Track in-flight async tasks so Stop() can wait for them; otherwise - // a detached bthread could use already-stopped components. - struct Ctx { - CrontabManager* manager; - CrontabConfig* config; - }; - auto* ctx = new Ctx{this, &crontab_config}; - inflight_async_count_.fetch_add(1); - - bthread_t tid; - const bthread_attr_t attr = BTHREAD_ATTR_NORMAL; - if (bthread_start_background( - &tid, &attr, - [](void* arg) -> void* { - auto* ctx = static_cast(arg); - ctx->config->funcer(nullptr); - ctx->manager->inflight_async_count_.fetch_sub(1); - delete ctx; - return nullptr; - }, - ctx) != 0) { - inflight_async_count_.fetch_sub(1); - delete ctx; - } - }; - } else { - crontab->func = crontab_config.funcer; - } +void Crontab::Stop() { + std::lock_guard lock(mu_); + stopped_ = true; - crontab->arg = nullptr; + // Only successful cancellation returns the reservation to us. + // Otherwise the timer callback still owns it. + if (has_timer_ && bthread_timer_del(timer_id_) == 0) { + has_timer_ = false; + ReleasePending(); + } +} - this->AddAndRunCrontab(crontab); +void Crontab::Join() { + std::unique_lock lock(mu_); + while (pending_ops_ > 0) { + bthread_cond_wait(&drained_cv_, &mu_); } } -uint32_t CrontabManager::AddAndRunCrontab(CrontabSPtr crontab) { - uint32_t crontab_id = AddCrontab(crontab); - StartCrontab(crontab_id); +void Crontab::DescribeByJson(Json::Value& value) const { + std::lock_guard lock(mu_); + value["id"] = id_; + value["name"] = cfg_.name; + value["interval_ms"] = static_cast(cfg_.interval_ms); + value["max_times"] = cfg_.max_times; + value["immediately"] = cfg_.immediately; // Configuration, not mutable run state. + value["run_count"] = run_count_; // Admitted calls, not completed calls. + value["stop"] = stopped_; +} - return crontab_id; +// Requires mu_. Reserve before publishing the timer. +void Crontab::ArmLocked() { + if (stopped_) return; + if (cfg_.max_times != 0 && run_count_ >= cfg_.max_times) return; + + timespec deadline = butil::milliseconds_from_now(cfg_.interval_ms); + ++pending_ops_; + has_timer_ = true; + int rc = bthread_timer_add(&timer_id_, deadline, &Crontab::OnTimer, this); + if (rc != 0) { + LOG(ERROR) << fmt::format("[crontab.arm][id({}).name({})] bthread_timer_add failed: {}", id_, cfg_.name, rc); + has_timer_ = false; + ReleasePending(); + } } -uint32_t CrontabManager::AddCrontab(CrontabSPtr crontab) { - BAIDU_SCOPED_LOCK(mutex_); +void Crontab::OnTimer(void* arg) { static_cast(arg)->OnTimerFired(); } - uint32_t crontab_id = AllocCrontabId(); - crontab->id = crontab_id; +void* Crontab::OnRoutineRun(void* arg) { + static_cast(arg)->RunOnce(); + return nullptr; +} - crontabs_[crontab_id] = crontab; - return crontab_id; +// Requires mu_. Reserve independently of the timer that dispatches this work. +void Crontab::StartRoutine() { + ++pending_ops_; + bthread_t tid; + const bthread_attr_t attr = BTHREAD_ATTR_NORMAL; + if (bthread_start_background(&tid, &attr, &Crontab::OnRoutineRun, this) != 0) { + LOG(ERROR) << fmt::format("[crontab.run][id({}).name({})] bthread_start_background failed", id_, cfg_.name); + ReleasePending(); + } } -void CrontabManager::StartCrontab(uint32_t crontab_id) { - BAIDU_SCOPED_LOCK(mutex_); +void Crontab::OnTimerFired() { + std::unique_lock lock(mu_); - auto it = crontabs_.find(crontab_id); - if (it == crontabs_.end()) { - LOG(WARNING) << fmt::format("[crontab.start][id({})] not exist crontab.", crontab_id); + has_timer_ = false; + if (stopped_ || (cfg_.max_times != 0 && run_count_ >= cfg_.max_times)) { + ReleasePending(); return; } - auto crontab = it->second; - crontab->pause = false; - bthread_t tid; - const bthread_attr_t attr = BTHREAD_ATTR_NORMAL; - bthread_start_background( - &tid, &attr, - [](void* arg) -> void* { - CrontabManager::Run(arg); - return nullptr; - }, - crontab.get()); -} - -void CrontabManager::InnerPauseCrontab(uint32_t crontab_id) { - auto it = crontabs_.find(crontab_id); - if (it == crontabs_.end()) { - LOG(WARNING) << fmt::format("[crontab.pause][id({})] not exist crontab.", crontab_id); - return; - } - auto crontab = it->second; + if (cfg_.async) { + StartRoutine(); + // Keep the submission cadence independent of callback duration. + ArmLocked(); + ReleasePending(); - crontab->pause = true; - if (crontab->timer_id != 0) { - bthread_timer_del(crontab->timer_id); + } else { + lock.unlock(); + RunOnce(); // The synchronous invocation takes over the timer reservation. } } -void CrontabManager::PauseCrontab(uint32_t crontab_id) { - BAIDU_SCOPED_LOCK(mutex_); +void Crontab::RunOnce() { + std::unique_lock lock(mu_); + // Check at invocation admission: queued bthreads must not exceed max_times. + if (!stopped_ && (cfg_.max_times == 0 || run_count_ < cfg_.max_times)) { + ++run_count_; + lock.unlock(); - InnerPauseCrontab(crontab_id); -} + try { + cfg_.callback(); + } catch (...) { + LOG(ERROR) << fmt::format("[crontab.run][id({}).name({})] exception in callback", id_, cfg_.name); + } + lock.lock(); + if (!cfg_.async) ArmLocked(); + } -void CrontabManager::DeleteCrontab(uint32_t crontab_id) { - BAIDU_SCOPED_LOCK(mutex_); - InnerPauseCrontab(crontab_id); + ReleasePending(); +} - crontabs_.erase(crontab_id); +// Requires mu_. +void Crontab::ReleasePending() { + --pending_ops_; + CHECK_GE(pending_ops_, 0) << cfg_.name; + if (pending_ops_ == 0) bthread_cond_broadcast(&drained_cv_); } -void CrontabManager::Stop() { - BAIDU_SCOPED_LOCK(mutex_); - - // Pause every crontab first so any Run() invocation still in flight sees - // pause==true and skips its reschedule (see the check added in Run()). - // Only after that is it safe to cancel timers and drop crontabs_'s - // shared_ptrs: otherwise a concurrent Run() could rearm a timer that - // outlives the Crontab it points to. - for (auto& [_, crontab] : crontabs_) { - crontab->pause = true; - } +CrontabManager::CrontabManager() { CHECK_EQ(bthread_mutex_init(&mu_, nullptr), 0); } - for (auto it = crontabs_.begin(); it != crontabs_.end();) { - while (bthread_timer_del(it->second->timer_id) == 1) { - bthread_usleep(1000L); // Wait for timer to be deleted - } +CrontabManager::~CrontabManager() { + Stop(); + bthread_mutex_destroy(&mu_); +} - it = crontabs_.erase(it); +void CrontabManager::AddCrontab(std::vector configs) { + std::vector launch_tasks; + launch_tasks.reserve(configs.size()); + { + std::lock_guard lock(mu_); + if (stopped_) { + LOG(WARNING) << fmt::format("[crontab.add] manager already stopped; dropping {} config(s)", configs.size()); + return; + } + for (auto& cfg : configs) { + LOG(INFO) << fmt::format("[crontab.add][name({}).interval({}ms).async({})] added", cfg.name, cfg.interval_ms, + cfg.async); + auto task = std::make_shared(next_id_++, std::move(cfg)); + tasks_.push_back(task); + launch_tasks.push_back(task); + } } + // Keep each task alive if Stop races with launch; Stop makes Launch a no-op. + for (const auto& task : launch_tasks) task->Launch(); +} - // Wait for in-flight async crontab tasks to finish, so components they use - // (operation processor, kv storage, ...) can be stopped safely afterwards. - while (inflight_async_count_.load() > 0) { - bthread_usleep(10000L); // 10ms +void CrontabManager::Stop() { + std::vector draining; + { + std::lock_guard lock(mu_); + stopped_ = true; + // Keep tasks visible until drained so concurrent Stop callers also wait. + draining = tasks_; } + + for (const auto& task : draining) task->Stop(); + for (const auto& task : draining) task->Join(); + + std::lock_guard lock(mu_); + tasks_.clear(); } void CrontabManager::DescribeByJson(Json::Value& value) { CHECK(value.isArray()) << "value is not array."; + std::vector snapshot; + { + std::lock_guard lock(mu_); + snapshot = tasks_; + } - BAIDU_SCOPED_LOCK(mutex_); - - for (auto& [_, crontab] : crontabs_) { - Json::Value crontab_value; - crontab->DescribeByJson(crontab_value); - value.append(crontab_value); + for (const auto& task : snapshot) { + Json::Value entry; + task->DescribeByJson(entry); + value.append(entry); } } diff --git a/src/mds/common/crontab.h b/src/mds/common/crontab.h index 3f2069df3..bb8378655 100644 --- a/src/mds/common/crontab.h +++ b/src/mds/common/crontab.h @@ -15,10 +15,8 @@ #ifndef DINGOFS_MDS_COMMON_CRONTAB_H_ #define DINGOFS_MDS_COMMON_CRONTAB_H_ -#include #include #include -#include #include #include #include @@ -29,81 +27,105 @@ namespace dingofs { namespace mds { +// Configuration for a single crontab task. struct CrontabConfig { + // Human-readable name; used only for logging and JSON reports. std::string name; - uint32_t interval; + // Delay between async submissions, or after a synchronous callback returns. + uint32_t interval_ms; + // Async callbacks run on bthreads and may overlap; callers own synchronization. + // Sync callbacks run on the shared timer thread and must not block. + // An immediate first invocation always runs on a bthread. bool async; - std::function funcer; + // Captured resources must outlive Join() / CrontabManager::Stop(). + std::function callback; + // 0 means unlimited. Every invocation counts, including those that throw. + uint32_t max_times{0}; + // If true the first invocation does not wait interval_ms. + bool immediately{false}; }; +// Launch at most once. Stop, then Join before destruction; never Join +// from this task's callback. The owner must join all callers before destruction. + class Crontab { public: - uint32_t id{0}; - std::string name; - // unit ms - int64_t interval{0}; - // 0 is no limit - uint32_t max_times{0}; - // Is immediately run - bool immediately{false}; - // Already run count - uint32_t run_count{0}; - // Is pause crontab. Read on the timer thread (Run) and written from - // whichever thread calls PauseCrontab/DeleteCrontab/Destroy, so it must - // be atomic to avoid a torn/reordered read racing with the reschedule - // in Run(). - std::atomic pause{false}; - // bthread_timer_t handler - bthread_timer_t timer_id{0}; - // For run target function - std::function func; - // Delivery to func_'s argument - void* arg{nullptr}; + Crontab(uint32_t id, CrontabConfig cfg); + ~Crontab(); + + Crontab(const Crontab&) = delete; + Crontab& operator=(const Crontab&) = delete; + + void Launch(); + + // Permanently close scheduling without waiting for admitted callbacks. + void Stop(); + // Wait for all timer/bthread reservations to be released. + // Stop() must be called first; Join() does not request stopping. + void Join(); void DescribeByJson(Json::Value& value) const; + + private: + static void OnTimer(void* arg); + static void* OnRoutineRun(void* arg); + + void OnTimerFired(); + void StartRoutine(); + void RunOnce(); + void ArmLocked(); + void ReleasePending(); + + const uint32_t id_; + const CrontabConfig cfg_; + + // mu_ guards scheduling state. Each timer and each dispatched bthread owns + // a reservation; Join waits until all reservations have been released. + mutable bthread_mutex_t mu_; + bthread_cond_t drained_cv_; + + bool stopped_{false}; + + bool has_timer_{false}; + bthread_timer_t timer_id_{0}; + + // Count admitted invocations, including callbacks still running or throwing. + uint32_t run_count_{0}; + int32_t pending_ops_{0}; }; + using CrontabSPtr = std::shared_ptr; -// Manage crontab use brpc::bthread_timer_add +// Owns periodic tasks. Async invocations of the same task may overlap. +// AddCrontab takes ownership of the configs; after Stop it is a no-op. +// Stop closes scheduling for all tasks before waiting, including concurrent calls. +// Callbacks must not call Stop or destroy this manager (that would self-wait). +// The owner must join all callers before destruction; the destructor calls Stop. class CrontabManager { public: CrontabManager(); ~CrontabManager(); CrontabManager(const CrontabManager&) = delete; - const CrontabManager& operator=(const CrontabManager&) = delete; - - static void Run(void* arg); - - void AddCrontab(std::vector& crontab_configs); + CrontabManager& operator=(const CrontabManager&) = delete; - uint32_t AddCrontab(CrontabSPtr crontab); - uint32_t AddAndRunCrontab(CrontabSPtr crontab); - void StartCrontab(uint32_t crontab_id); - void PauseCrontab(uint32_t crontab_id); - void DeleteCrontab(uint32_t crontab_id); + void AddCrontab(std::vector configs); void Stop(); void DescribeByJson(Json::Value& value); private: - // Allocate crontab id by auto incremental. - uint32_t AllocCrontabId(); - - void InnerPauseCrontab(uint32_t crontab_id); - - // Atomic auto incremental variable - std::atomic auinc_crontab_id_; - // Protect crontabs_ concurrence access. - bthread_mutex_t mutex_; - // Store all crontab, key(crontab_id) / value(Crontab) - std::map crontabs_; - // In-flight async crontab tasks; Stop() waits until it drops to zero. - std::atomic inflight_async_count_{0}; + bthread_mutex_t mu_; + + std::vector tasks_; + + uint32_t next_id_{1}; + + bool stopped_{false}; }; } // namespace mds } // namespace dingofs -#endif // DINGOFS_MDS_COMMON_CRONTAB_H_ \ No newline at end of file +#endif // DINGOFS_MDS_COMMON_CRONTAB_H_ diff --git a/src/mds/filesystem/filesystem.cc b/src/mds/filesystem/filesystem.cc index f696745f7..8e0a92272 100644 --- a/src/mds/filesystem/filesystem.cc +++ b/src/mds/filesystem/filesystem.cc @@ -4521,19 +4521,29 @@ Status FileSystemSet::CreateFs(const CreateFsParam& param, FsInfoEntry& fs_info) } // generate fs id - uint32_t fs_id; + uint32_t fs_id = 0; if (param.fs_id == 0) { - status = GenFsId(fs_id); - if (BAIDU_UNLIKELY(!status.ok())) { - return status; + // The generator may hand out an id that is already taken (e.g. its persisted + // counter was reset). Skip such ids instead of failing with a misleading + // "fs(x) exist.". + constexpr int kMaxTryGenFsId = 100; + int tries = 0; + do { + status = GenFsId(fs_id); + if (BAIDU_UNLIKELY(!status.ok())) { + return status; + } + } while (IsExistFileSystem(fs_id) && ++tries < kMaxTryGenFsId); + + if (IsExistFileSystem(fs_id)) { + return Status(pb::error::EALLOC_ID, fmt::format("gen free fs id fail, fs({})", param.fs_name)); } } else { fs_id = param.fs_id; - } - - // check fs_id exist - if (IsExistFileSystem(fs_id)) { - return Status(pb::error::EEXISTED, fmt::format("fs({}) exist.", param.fs_name)); + // check fs_id exist + if (IsExistFileSystem(fs_id)) { + return Status(pb::error::EEXISTED, fmt::format("fs id({}) exist.", fs_id)); + } } // create dentry/inode table diff --git a/src/mds/filesystem/id_generator.cc b/src/mds/filesystem/id_generator.cc index ded4a3a66..b748aebe4 100644 --- a/src/mds/filesystem/id_generator.cc +++ b/src/mds/filesystem/id_generator.cc @@ -232,14 +232,10 @@ bool StoreAutoIncrementIdGenerator::Init() { bool StoreAutoIncrementIdGenerator::Stop() { BAIDU_SCOPED_LOCK(mutex_); + // Keep the persisted counter across restarts, otherwise the next start would + // re-issue ids that are still in use. is_destroyed_.store(true, std::memory_order_release); - auto status = DestroyId(); - if (!status.ok()) { - LOG(ERROR) << fmt::format("[idalloc.{}] destroy autoincrement table fail, status({}).", name_, status.error_str()); - return false; - } - return true; } @@ -701,7 +697,10 @@ void DestroyInodeIdGenerator(uint32_t fs_id, KVStorageSPtr kv_storage) { auto id_generator = StoreAutoIncrementIdGenerator::New(kv_storage, name, kInoStartId, FLAGS_mds_ino_generator_batch_size); - id_generator->Stop(); + auto status = id_generator->DestroyId(); + if (!status.ok()) { + LOG(ERROR) << fmt::format("[idalloc.{}] destroy inode id counter fail, status({}).", name, status.error_str()); + } } } // namespace mds diff --git a/src/mds/filesystem/id_generator.h b/src/mds/filesystem/id_generator.h index ab49fdff5..3b304d00b 100644 --- a/src/mds/filesystem/id_generator.h +++ b/src/mds/filesystem/id_generator.h @@ -99,7 +99,8 @@ class StoreAutoIncrementIdGenerator : public IdGenerator { StoreAutoIncrementIdGenerator(KVStorageSPtr kv_storage, const std::string& name, int64_t start_id, int batch_size); ~StoreAutoIncrementIdGenerator() override; - static IdGeneratorUPtr New(KVStorageSPtr kv_storage, const std::string& name, int64_t start_id, int batch_size) { + static std::unique_ptr New(KVStorageSPtr kv_storage, const std::string& name, + int64_t start_id, int batch_size) { return std::make_unique(kv_storage, name, start_id, batch_size); } @@ -108,6 +109,8 @@ class StoreAutoIncrementIdGenerator : public IdGenerator { } bool Init() override; + // Stop only marks the generator unusable. It must not delete the persisted + // counter, otherwise a restart would re-issue already-used ids. bool Stop() override; bool GenID(uint32_t num, uint64_t& id) override; @@ -115,13 +118,16 @@ class StoreAutoIncrementIdGenerator : public IdGenerator { std::string Describe() override; + // Delete the persisted counter. Only for per-fs generators whose filesystem + // is gone; never call it for a global generator on shutdown. + Status DestroyId(); + private: Status GetOrPutAllocId(uint64_t& alloc_id); // Reserve a new bundle from storage. Only the elected refiller calls this; on // success it publishes next_id_ before last_alloc_id_ (see R6). `floor` is the // min_slice_id that triggered the refill, so the new bundle covers it. Status AllocateIds(uint32_t bundle_size, uint64_t floor); - Status DestroyId(); KVStorageSPtr kv_storage_; diff --git a/src/mds/server.cc b/src/mds/server.cc index fb84841f2..09abfaca2 100644 --- a/src/mds/server.cc +++ b/src/mds/server.cc @@ -465,7 +465,7 @@ bool Server::InitCrontab() { "HEARTBEAT", FLAGS_mds_crontab_heartbeat_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetHeartbeat()->Run(); }, + []() { Server::GetInstance().GetHeartbeat()->Run(); }, }); // Add fs info sync crontab @@ -473,7 +473,7 @@ bool Server::InitCrontab() { "FSINFO_SYNC", FLAGS_mds_crontab_fsinfosync_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetFsInfoSync()->Run(); }, + []() { Server::GetInstance().GetFsInfoSync()->Run(); }, }); // Add fs info sync crontab @@ -481,7 +481,7 @@ bool Server::InitCrontab() { "MDS_MONITOR", FLAGS_mds_crontab_mdsmonitor_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetMonitor()->Run(); }, + []() { Server::GetInstance().GetMonitor()->Run(); }, }); // Add quota sync crontab @@ -489,7 +489,7 @@ bool Server::InitCrontab() { "QUOTA_SYNC", FLAGS_mds_crontab_quota_sync_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetQuotaSynchronizer()->Run(); }, + []() { Server::GetInstance().GetQuotaSynchronizer()->Run(); }, }); // Add dir-stats sync crontab @@ -497,7 +497,7 @@ bool Server::InitCrontab() { "DIR_STATS_SYNC", FLAGS_mds_crontab_dir_stats_sync_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetDirStatsSynchronizer()->Run(); }, + []() { Server::GetInstance().GetDirStatsSynchronizer()->Run(); }, }); // Add fs info sync crontab @@ -505,7 +505,7 @@ bool Server::InitCrontab() { "GC", FLAGS_mds_crontab_gc_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetGcProcessor()->Run(); }, + []() { Server::GetInstance().GetGcProcessor()->Run(); }, }); // Trash cleanup runs on its own cadence — bucket eligibility only changes at @@ -514,7 +514,7 @@ bool Server::InitCrontab() { "GC_TRASH", FLAGS_mds_crontab_gc_trash_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetGcProcessor()->RunTrash(); }, + []() { Server::GetInstance().GetGcProcessor()->RunTrash(); }, }); // Add filesystem cache crontab @@ -522,7 +522,7 @@ bool Server::InitCrontab() { "CLEAN_EXPIRED_CACHE", FLAGS_mds_crontab_clean_expired_cache_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetFileSystemSet()->CleanExpiredCache(); }, + []() { Server::GetInstance().GetFileSystemSet()->CleanExpiredCache(); }, }); // Add cache member sync crontab @@ -530,7 +530,7 @@ bool Server::InitCrontab() { "CACHE_MEMBER_SYNC", FLAGS_mds_crontab_cache_member_sync_interval_s * 1000, true, - [](void*) { Server::GetInstance().GetCacheMemberSynchronizer()->Run(); }, + []() { Server::GetInstance().GetCacheMemberSynchronizer()->Run(); }, }); return true; diff --git a/test/unit/mds/common/test_crontab.cc b/test/unit/mds/common/test_crontab.cc index 11dc6d09d..f1db584a6 100644 --- a/test/unit/mds/common/test_crontab.cc +++ b/test/unit/mds/common/test_crontab.cc @@ -4,7 +4,7 @@ // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // -// http://www.apache.org/licenses/LICENSE-2.0 +// http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, @@ -12,234 +12,333 @@ // See the License for the specific language governing permissions and // limitations under the License. -#include "mds/common/crontab.h" - #include #include +#include +#include +#include #include +#include +#include -#include "fmt/core.h" -#include "gflags/gflags.h" -#include "glog/logging.h" #include "gtest/gtest.h" +#include "mds/common/crontab.h" namespace dingofs { namespace mds { namespace unit_test { -class CrontabTest : public testing::Test { - protected: - void SetUp() override {} - void TearDown() override {} -}; - -TEST_F(CrontabTest, CrontabDescribeByJson) { - auto crontab = std::make_shared(); - crontab->id = 1; - crontab->name = "test_crontab"; - crontab->interval = 1000; - crontab->max_times = 10; - crontab->immediately = true; - crontab->run_count = 5; - crontab->pause = false; - - Json::Value value; - crontab->DescribeByJson(value); - - EXPECT_EQ(value["id"].asUInt(), 1); - EXPECT_EQ(value["name"].asString(), "test_crontab"); - EXPECT_EQ(value["interval_ms"].asInt64(), 1000); - EXPECT_EQ(value["max_times"].asUInt(), 10); - EXPECT_EQ(value["immediately"].asBool(), true); - EXPECT_EQ(value["run_count"].asUInt(), 5); - EXPECT_EQ(value["pause"].asBool(), false); -} - -TEST_F(CrontabTest, CrontabManagerAddCrontab) { - CrontabManager manager; - - auto crontab = std::make_shared(); - crontab->name = "test_crontab"; - crontab->interval = 1000; - crontab->func = [](void*) {}; - - uint32_t id = manager.AddCrontab(crontab); - // ID is allocated sequentially from 0 - EXPECT_EQ(crontab->id, id); +namespace { - auto crontab2 = std::make_shared(); - crontab2->name = "test_crontab2"; - crontab2->interval = 2000; - crontab2->func = [](void*) {}; +using std::chrono::milliseconds; +using std::chrono::steady_clock; - uint32_t id2 = manager.AddCrontab(crontab2); - // Second crontab should have a different ID - EXPECT_NE(id2, id); - EXPECT_EQ(crontab2->id, id2); +// Poll a predicate until it returns true or `budget` elapses. Returns the last +// evaluated value. +template +bool WaitFor(F pred, milliseconds budget = milliseconds(2000)) { + auto deadline = steady_clock::now() + budget; + while (!pred()) { + if (steady_clock::now() >= deadline) return pred(); + std::this_thread::sleep_for(milliseconds(1)); + } + return true; } -TEST_F(CrontabTest, CrontabManagerPauseAndDeleteCrontab) { - CrontabManager manager; - - auto crontab = std::make_shared(); - crontab->name = "test_crontab"; - crontab->interval = 1000; - crontab->func = [](void*) {}; - - uint32_t id = manager.AddCrontab(crontab); - EXPECT_EQ(crontab->id, id); - - // Pause should not throw - manager.PauseCrontab(id); - - // Delete should not throw - manager.DeleteCrontab(id); +CrontabConfig MakeConfig(std::string name, uint32_t interval_ms, bool async, + std::function callback, uint32_t max_times = 0, + bool immediately = false) { + CrontabConfig cfg; + cfg.name = std::move(name); + cfg.interval_ms = interval_ms; + cfg.async = async; + cfg.callback = std::move(callback); + cfg.max_times = max_times; + cfg.immediately = immediately; + return cfg; } -TEST_F(CrontabTest, CrontabManagerPauseNonExistentCrontab) { - CrontabManager manager; - - // Pause non-existent crontab should not throw - manager.PauseCrontab(999); +} // namespace + +TEST(CrontabTest, ImmediateFiresWithoutWaitingInterval) { + std::atomic hits{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig( + "task", /*interval_ms=*/3600'000, /*async=*/false, [&]() { ++hits; }, + /*max_times=*/1, + /*immediately=*/true)); + mgr.AddCrontab(std::move(cfgs)); + // The interval is an hour; the callback must have fired without waiting it. + ASSERT_TRUE(WaitFor([&] { return hits.load() == 1; })); } -TEST_F(CrontabTest, CrontabManagerDeleteNonExistentCrontab) { - CrontabManager manager; - - // Delete non-existent crontab should not throw - manager.DeleteCrontab(999); +TEST(CrontabTest, NonImmediateFiresAfterInterval) { + std::atomic hits{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig("task", /*interval_ms=*/20, /*async=*/false, + [&]() { ++hits; })); + auto start = steady_clock::now(); + mgr.AddCrontab(std::move(cfgs)); + ASSERT_TRUE(WaitFor([&] { return hits.load() >= 1; })); + // First tick must have waited roughly `interval` before firing. + EXPECT_GE(steady_clock::now() - start, milliseconds(15)); } -TEST_F(CrontabTest, CrontabManagerAddMultipleCrontabs) { - CrontabManager manager; - - std::vector configs; - - CrontabConfig config1; - config1.name = "crontab1"; - config1.interval = 1000; - config1.async = false; - config1.funcer = [](void*) {}; - configs.push_back(config1); - - CrontabConfig config2; - config2.name = "crontab2"; - config2.interval = 2000; - config2.async = true; - config2.funcer = [](void*) {}; - configs.push_back(config2); - - manager.AddCrontab(configs); - - Json::Value value(Json::arrayValue); - manager.DescribeByJson(value); - EXPECT_EQ(value.size(), 2); +TEST(CrontabTest, MaxTimesLimitsInvocations) { + std::atomic hits{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig( + "task", /*interval_ms=*/5, /*async=*/false, [&]() { ++hits; }, + /*max_times=*/3)); + mgr.AddCrontab(std::move(cfgs)); + ASSERT_TRUE(WaitFor([&] { return hits.load() == 3; })); + // Give the timer plenty of time to over-fire if the limit is broken. + std::this_thread::sleep_for(milliseconds(50)); + EXPECT_EQ(hits.load(), 3); } -TEST_F(CrontabTest, CrontabManagerDescribeByJson) { - CrontabManager manager; - - auto crontab = std::make_shared(); - crontab->name = "test_crontab"; - crontab->interval = 1000; - crontab->max_times = 5; - crontab->func = [](void*) {}; - - manager.AddCrontab(crontab); - - Json::Value value(Json::arrayValue); - manager.DescribeByJson(value); +TEST(CrontabTest, ConcurrentStopsWaitForRunningCallback) { + std::promise entered; + std::promise release; + auto released = release.get_future().share(); + std::atomic completed{false}; + CrontabManager mgr; + mgr.AddCrontab({MakeConfig( + "task", 1, true, + [&]() { + entered.set_value(); + released.wait(); + completed = true; + }, + 1, true)}); + auto entered_status = entered.get_future().wait_for(milliseconds(2000)); + if (entered_status != std::future_status::ready) { + release.set_value(); + FAIL() << "callback did not start"; + } - EXPECT_EQ(value.size(), 1); - EXPECT_EQ(value[0]["name"].asString(), "test_crontab"); - EXPECT_EQ(value[0]["interval_ms"].asInt64(), 1000); - EXPECT_EQ(value[0]["max_times"].asUInt(), 5); + std::promise first_entered, second_entered; + auto first = std::async(std::launch::async, [&] { + first_entered.set_value(); + mgr.Stop(); + return completed.load(); + }); + auto second = std::async(std::launch::async, [&] { + second_entered.set_value(); + mgr.Stop(); + return completed.load(); + }); + first_entered.get_future().wait(); + second_entered.get_future().wait(); + const auto first_status = first.wait_for(milliseconds(50)); + const auto second_status = second.wait_for(milliseconds(50)); + release.set_value(); + EXPECT_TRUE(first.get()) + << "first Stop returned before the callback completed"; + EXPECT_TRUE(second.get()) + << "second Stop returned before the callback completed"; + EXPECT_EQ(first_status, std::future_status::timeout); + EXPECT_EQ(second_status, std::future_status::timeout); } -TEST_F(CrontabTest, CrontabRunWithMaxTimes) { - std::atomic counter{0}; - - auto crontab = std::make_shared(); - crontab->name = "test_crontab"; - crontab->interval = 10; - crontab->max_times = 3; - crontab->immediately = false; - crontab->pause = false; - crontab->func = [&counter](void*) { counter.fetch_add(1); }; - - CrontabManager manager; - manager.AddAndRunCrontab(crontab); +TEST(CrontabTest, StopCancelsPendingTimer) { + std::atomic hits{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig("task", /*interval_ms=*/500, /*async=*/false, + [&]() { ++hits; })); + mgr.AddCrontab(std::move(cfgs)); + // Timer is armed but shouldn't have fired yet. + mgr.Stop(); + // Wait past when the timer would have fired had Stop not cancelled it. + std::this_thread::sleep_for(milliseconds(600)); + EXPECT_EQ(hits.load(), 0); +} - // Wait for crontab to execute - std::this_thread::sleep_for(std::chrono::milliseconds(100)); +TEST(CrontabTest, StopIsIdempotent) { + std::atomic hits{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig( + "task", /*interval_ms=*/1, /*async=*/false, [&]() { ++hits; }, + /*max_times=*/0, + /*immediately=*/true)); + mgr.AddCrontab(std::move(cfgs)); + ASSERT_TRUE(WaitFor([&] { return hits.load() >= 1; })); + mgr.Stop(); + int after_stop = hits.load(); + mgr.Stop(); // Second call must not crash or double-drain. + std::this_thread::sleep_for(milliseconds(20)); + EXPECT_EQ(hits.load(), after_stop); +} - manager.PauseCrontab(crontab->id); +TEST(CrontabTest, AddCrontabAfterStopIsNoOp) { + std::atomic hits{0}; + CrontabManager mgr; + mgr.Stop(); + std::vector cfgs; + cfgs.push_back(MakeConfig( + "task", /*interval_ms=*/1, /*async=*/false, [&]() { ++hits; }, + /*max_times=*/0, + /*immediately=*/true)); + mgr.AddCrontab(std::move(cfgs)); + std::this_thread::sleep_for(milliseconds(30)); + EXPECT_EQ(hits.load(), 0); + Json::Value view(Json::arrayValue); + mgr.DescribeByJson(view); + EXPECT_EQ(view.size(), 0u); +} - // The counter should be at least 3 (max_times) - EXPECT_GE(counter.load(), 1); +TEST(CrontabTest, DescribeByJsonEnumeratesTasks) { + std::atomic a{0}, b{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig( + "alpha", /*interval_ms=*/1000, /*async=*/true, [&]() { ++a; }, + /*max_times=*/7)); + cfgs.push_back(MakeConfig( + "beta", /*interval_ms=*/2500, /*async=*/false, [&]() { ++b; }, + /*max_times=*/0, + /*immediately=*/true)); + mgr.AddCrontab(std::move(cfgs)); + + Json::Value view(Json::arrayValue); + mgr.DescribeByJson(view); + ASSERT_EQ(view.size(), 2u); + + // Ordering is insertion order. + EXPECT_EQ(view[0]["name"].asString(), "alpha"); + EXPECT_EQ(view[0]["interval_ms"].asInt64(), 1000); + EXPECT_EQ(view[0]["max_times"].asUInt(), 7u); + EXPECT_EQ(view[0]["immediately"].asBool(), false); + EXPECT_EQ(view[0]["pause"].asBool(), false); + + EXPECT_EQ(view[1]["name"].asString(), "beta"); + EXPECT_EQ(view[1]["interval_ms"].asInt64(), 2500); + EXPECT_EQ(view[1]["max_times"].asUInt(), 0u); + EXPECT_EQ(view[1]["immediately"].asBool(), true); } -TEST_F(CrontabTest, DestructorStopsPendingTimerWithoutUseAfterFree) { - // Regression test: CrontabManager::~CrontabManager() used to only - // destroy the mutex, leaving any still-running crontab's timer armed on - // brpc's global TimerThread. Once `manager` below goes out of scope, its - // crontabs_ map (and the last shared_ptr it holds) is dropped; - // if the destructor doesn't cancel pending timers first, the timer - // thread can later invoke Run() on a freed Crontab (use-after-free). - // Running this in a loop under the process' normal exit path reliably - // caught the bug: destroy() now runs from ~CrontabManager and the timer - // is fully cancelled before the Crontab can be freed. - std::atomic counter{0}; +TEST(CrontabTest, DestructorDrainsWithoutUseAfterFree) { + std::atomic hits{0}; { - auto crontab = std::make_shared(); - crontab->name = "destructor_test_crontab"; - crontab->interval = 1; // fire as fast as possible to stress the race - crontab->max_times = 0; // unlimited, keeps rescheduling until stopped - crontab->func = [&counter](void*) { counter.fetch_add(1); }; - - CrontabManager manager; - manager.AddAndRunCrontab(crontab); - std::this_thread::sleep_for(std::chrono::milliseconds(5)); - // No explicit PauseCrontab/DeleteCrontab/Destroy() call here: the - // manager's destructor alone must make it safe to free `crontab`. + CrontabManager mgr; + std::vector cfgs; + // Interval short enough that timers keep re-arming during the sleep; + // ~mgr must Stop() and drain them before the callback capture goes away. + cfgs.push_back(MakeConfig( + "task", /*interval_ms=*/1, /*async=*/true, [&]() { ++hits; }, + /*max_times=*/0, + /*immediately=*/true)); + mgr.AddCrontab(std::move(cfgs)); + ASSERT_TRUE(WaitFor([&] { return hits.load() >= 5; })); } - - // If the manager's destructor left the crontab's timer armed, this - // sleep gives brpc's TimerThread a chance to fire it against freed - // memory (which previously crashed the process). - std::this_thread::sleep_for(std::chrono::milliseconds(20)); - SUCCEED(); + // If drain missed a callback the process would have already crashed or + // continued incrementing `hits` after the manager's teardown. Take a nap + // and confirm the counter is stable. + int snapshot = hits.load(); + std::this_thread::sleep_for(milliseconds(20)); + EXPECT_EQ(hits.load(), snapshot); } -TEST_F(CrontabTest, CrontabPausePreventsExecution) { - // Test that a paused crontab doesn't continue to schedule itself - std::atomic counter{0}; +TEST(CrontabTest, MultipleTasksScheduleIndependently) { + std::atomic fast{0}, slow{0}; + CrontabManager mgr; + std::vector cfgs; + cfgs.push_back(MakeConfig( + "fast", /*interval_ms=*/2, /*async=*/false, [&]() { ++fast; }, + /*max_times=*/10)); + cfgs.push_back(MakeConfig( + "slow", /*interval_ms=*/100, /*async=*/false, [&]() { ++slow; }, + /*max_times=*/1)); + mgr.AddCrontab(std::move(cfgs)); + ASSERT_TRUE(WaitFor([&] { return fast.load() == 10 && slow.load() == 1; }, + milliseconds(3000))); + EXPECT_EQ(fast.load(), 10); + EXPECT_EQ(slow.load(), 1); +} - auto crontab = std::make_shared(); - crontab->name = "test_crontab"; - crontab->interval = 10; - crontab->max_times = 10; // Limit to 10 executions - crontab->immediately = false; - crontab->pause = false; - crontab->func = [&counter](void*) { counter.fetch_add(1); }; +TEST(CrontabTest, ThrowingCallbackCountsTowardLimit) { + std::atomic hits{0}; + Crontab task(1, MakeConfig( + "throwing", 1, true, + [&]() { + ++hits; + throw std::runtime_error("callback failed"); + }, + 3, true)); + task.Launch(); + const bool reached_limit = WaitFor([&] { return hits.load() >= 3; }); + std::this_thread::sleep_for(milliseconds(50)); + task.Stop(); + task.Join(); + EXPECT_TRUE(reached_limit); + EXPECT_EQ(hits.load(), 3); + Json::Value view; + task.DescribeByJson(view); + EXPECT_EQ(view["run_count"].asUInt(), 3u); +} +TEST(CrontabTest, AsyncSchedulesWhilePreviousInvocationIsBlocked) { + std::promise release, second_entered; + auto released = release.get_future().share(); + std::atomic calls{0}; CrontabManager manager; - uint32_t id = manager.AddAndRunCrontab(crontab); - - // Wait for some executions - std::this_thread::sleep_for(std::chrono::milliseconds(50)); - - // Pause the crontab - manager.PauseCrontab(id); - - int count_after_pause = counter.load(); - - // Wait again - std::this_thread::sleep_for(std::chrono::milliseconds(100)); + manager.AddCrontab({MakeConfig( + "overlap", 20, true, + [&]() { + if (++calls == 2) second_entered.set_value(); + released.wait(); + }, + 2, true)}); + const auto second_status = + second_entered.get_future().wait_for(milliseconds(2000)); + // Both invocations remain blocked; max_times must prevent a third. + std::this_thread::sleep_for(milliseconds(50)); + const int before_release = calls.load(); + release.set_value(); + manager.Stop(); + EXPECT_EQ(second_status, std::future_status::ready); + EXPECT_EQ(before_release, 2); + EXPECT_EQ(calls.load(), 2); +} - // Counter should not increase much after pause - // Allow for some already-scheduled executions - EXPECT_LE(counter.load(), count_after_pause + 3); +TEST(CrontabTest, StopClosesEveryTaskBeforeWaitingForCallbacks) { + std::promise entered, release; + auto released = release.get_future().share(); + std::atomic later_calls{0}; + CrontabManager manager; + manager.AddCrontab( + {MakeConfig( + "blocked", 3600000, true, + [&]() { + entered.set_value(); + released.wait(); + }, + 1, true), + MakeConfig("later", 3600000, true, [&]() { ++later_calls; })}); + if (entered.get_future().wait_for(milliseconds(2000)) != + std::future_status::ready) { + release.set_value(); + FAIL() << "callback did not start"; + } + auto stopped = std::async(std::launch::async, [&] { manager.Stop(); }); + // The public task view must report both tasks closed while Stop is blocked. + const bool all_closed = WaitFor([&] { + Json::Value view(Json::arrayValue); + manager.DescribeByJson(view); + return view.size() == 2 && view[0]["pause"].asBool() && + view[1]["pause"].asBool(); + }); + const auto stop_status = stopped.wait_for(milliseconds(0)); + release.set_value(); + stopped.get(); + EXPECT_TRUE(all_closed); + EXPECT_EQ(stop_status, std::future_status::timeout); + EXPECT_EQ(later_calls.load(), 0); } } // namespace unit_test diff --git a/test/unit/mds/filesystem/test_id_generator.cc b/test/unit/mds/filesystem/test_id_generator.cc index 389b56258..7833c8bd8 100644 --- a/test/unit/mds/filesystem/test_id_generator.cc +++ b/test/unit/mds/filesystem/test_id_generator.cc @@ -148,6 +148,29 @@ TEST_F(StoreAutoIncrementIdGeneratorTest, ConcurrentMixedNumAndFloor) { for (auto& w : workers) w.join(); } +TEST_F(StoreAutoIncrementIdGeneratorTest, CounterSurvivesStop) { + const int64_t kStartId = 1000; + uint64_t last = 0; + { + auto id_generator = + StoreAutoIncrementIdGenerator::New(storage_, "store-restart", kStartId, 8); + ASSERT_TRUE(id_generator->Init()) << "init id generator fail."; + for (int i = 0; i < 20; ++i) { + ASSERT_TRUE(id_generator->GenID(1, last)); + } + ASSERT_TRUE(id_generator->Stop()) << "stop id generator fail."; + } + + // A graceful stop must not drop the persisted counter, otherwise the next + // start re-issues ids that are still in use. + auto id_generator = + StoreAutoIncrementIdGenerator::New(storage_, "store-restart", kStartId, 8); + ASSERT_TRUE(id_generator->Init()) << "init id generator fail."; + uint64_t id = 0; + ASSERT_TRUE(id_generator->GenID(1, id)); + ASSERT_GT(id, last) << "counter reset after stop, reused id " << id; +} + } // namespace unit_test } // namespace mds } // namespace dingofs \ No newline at end of file