From 4880c16be3982ba442de8ca848911c1102fc7f08 Mon Sep 17 00:00:00 2001 From: Jiang-Jia-Jun <163579578+Jiang-Jia-Jun@users.noreply.github.com> Date: Thu, 31 Jul 2025 20:30:24 +0800 Subject: [PATCH 01/17] Update setup.py --- setup.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/setup.py b/setup.py index 87099104b78..e13e70d07ec 100644 --- a/setup.py +++ b/setup.py @@ -181,7 +181,7 @@ def get_name(): cmdclass_dict = {"bdist_wheel": CustomBdistWheel} cmdclass_dict["build_ext"] = CMakeBuild -FASTDEPLOY_VERSION = os.environ.get("FASTDEPLOY_VERSION", "2.1.0-dev") +FASTDEPLOY_VERSION = os.environ.get("FASTDEPLOY_VERSION", "2.1.0") cmdclass_dict["build_optl"] = PostInstallCommand setup( From c8dd5976ae5db01d0cc67a2629268bd2bce9a80a Mon Sep 17 00:00:00 2001 From: chen <103103266+ckl117@users.noreply.github.com> Date: Fri, 1 Aug 2025 22:34:33 +0800 Subject: [PATCH 02/17] fix request_output sampling_params (#3154) --- fastdeploy/engine/engine.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/fastdeploy/engine/engine.py b/fastdeploy/engine/engine.py index e7443bc1db7..c76a6432dd1 100644 --- a/fastdeploy/engine/engine.py +++ b/fastdeploy/engine/engine.py @@ -749,10 +749,6 @@ def insert_tasks(self, tasks, current_id=-1, allocated=False): """ Insert tasks to engine. """ - for task in tasks: - start_span_request("DEQUEUE", task, trace.SpanKind.CONSUMER) - if task.sampling_params.bad_words is not None: - task.sampling_params.update_from_tokenizer(self.data_processor.tokenizer) # TODO 返回至 scheduler if allocated: current_tasks = [] @@ -779,6 +775,11 @@ def insert_tasks(self, tasks, current_id=-1, allocated=False): self.engine_worker_queue.put_tasks((current_tasks, self.resource_manager.real_bsz)) return True + for task in tasks: + start_span_request("DEQUEUE", task, trace.SpanKind.CONSUMER) + if task.sampling_params.bad_words is not None: + task.sampling_params.update_from_tokenizer(self.data_processor.tokenizer) + self.resource_manager.check_and_free_block_tables() if not isinstance(tasks, list): From d4059cabf0c217b32555989244447449f3d8cc73 Mon Sep 17 00:00:00 2001 From: RAM Date: Fri, 1 Aug 2025 22:34:59 +0800 Subject: [PATCH 03/17] fix typo (#3153) --- .../model_executor/layers/attention/flash_attn_backend.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/fastdeploy/model_executor/layers/attention/flash_attn_backend.py b/fastdeploy/model_executor/layers/attention/flash_attn_backend.py index 306164635bc..199a26db81c 100644 --- a/fastdeploy/model_executor/layers/attention/flash_attn_backend.py +++ b/fastdeploy/model_executor/layers/attention/flash_attn_backend.py @@ -208,7 +208,7 @@ def init_attention_metadata(self, forward_meta: ForwardMeta): ) = pre_cache_len_concat( forward_meta.seq_lens_decoder, forward_meta.seq_lens_this_time, - metadata.set_max_lengths[2], + forward_meta.max_len_tensor_cpu[2], self.block_size, ) From 5f6fc7f7b9802c9271e54963a8f65a3788307fb9 Mon Sep 17 00:00:00 2001 From: ltd0924 <32387785+ltd0924@users.noreply.github.com> Date: Mon, 4 Aug 2025 15:09:17 +0800 Subject: [PATCH 04/17] Update cache_messager.py (#3173) --- fastdeploy/cache_manager/cache_messager.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/fastdeploy/cache_manager/cache_messager.py b/fastdeploy/cache_manager/cache_messager.py index e06d05a67ff..456ba1c3422 100644 --- a/fastdeploy/cache_manager/cache_messager.py +++ b/fastdeploy/cache_manager/cache_messager.py @@ -142,7 +142,7 @@ def __init__( self.gpu_id = gpu_id self.cache_info = dict() - self.dp_rank_id = local_data_parallel_id + self.dp_rank_id = self.rank + local_data_parallel_id * self.nranks layerwise_send_cache_thread = threading.Thread(target=self._prefill_layerwise_send_cache_thread) layerwise_send_cache_thread.daemon = True From 8e789dcb6787f5af5ad1e9bdf79732dbc179ed6a Mon Sep 17 00:00:00 2001 From: bukejiyu <52310069+bukejiyu@users.noreply.github.com> Date: Mon, 4 Aug 2025 15:44:10 +0800 Subject: [PATCH 05/17] fix load_pre_sharded_checkpoint (#3152) (#3169) Co-authored-by: Jiang-Jia-Jun <163579578+Jiang-Jia-Jun@users.noreply.github.com> --- fastdeploy/model_executor/load_weight_utils.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/fastdeploy/model_executor/load_weight_utils.py b/fastdeploy/model_executor/load_weight_utils.py index 2856737ff5b..01f81ac13d4 100644 --- a/fastdeploy/model_executor/load_weight_utils.py +++ b/fastdeploy/model_executor/load_weight_utils.py @@ -215,11 +215,13 @@ def load_pre_sharded_checkpoint(model_path: str, local_rank: int, use_fastsafete """ load_pre_sharded_checkpoint """ + from fastdeploy.model_executor.layers.utils import get_tensor + state_dict = {} _, safetensor_files = get_all_safetensors(os.path.join(model_path, f"rank{local_rank}")) weights_iterator = safetensors_weights_iterator(safetensor_files) for name, weight in weights_iterator: - state_dict[name] = weight + state_dict[name] = get_tensor(weight) return state_dict From 4367c09a5fc2883bcce6112f240d3169671cc4bd Mon Sep 17 00:00:00 2001 From: yinwei Date: Mon, 4 Aug 2025 16:02:43 +0800 Subject: [PATCH 06/17] Fix out-of-memory issue during single-XPU deployment (#3131) --- fastdeploy/worker/xpu_worker.py | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/fastdeploy/worker/xpu_worker.py b/fastdeploy/worker/xpu_worker.py index 0332d34d221..16b51d2e567 100644 --- a/fastdeploy/worker/xpu_worker.py +++ b/fastdeploy/worker/xpu_worker.py @@ -94,9 +94,14 @@ def determine_available_memory(self) -> int: xpu_get_used_global_memory, ) - total_memory = xpu_get_total_global_memory(self.local_rank) - used_memory = xpu_get_used_global_memory(self.local_rank) - free_memory = xpu_get_free_global_memory(self.local_rank) + assert self.device_ids[self.local_rank] is not None, f"device_id is none for rank {self.local_rank}" + assert ( + len(self.device_ids) > self.local_rank + ), f"device number must be greater than local rank, but get device number is {len(self.device_ids)}, rank is {self.local_rank}" + + total_memory = xpu_get_total_global_memory(int(self.device_ids[self.local_rank])) + used_memory = xpu_get_used_global_memory(int(self.device_ids[self.local_rank])) + free_memory = xpu_get_free_global_memory(int(self.device_ids[self.local_rank])) logger.info( f"Before warm up, total_memory: {total_memory}, \ @@ -107,7 +112,7 @@ def determine_available_memory(self) -> int: self.model_runner.profile_run() total_available_memory = int(total_memory * self.cache_config.gpu_memory_utilization) - used_memory = xpu_get_used_global_memory(self.local_rank) + used_memory = xpu_get_used_global_memory(int(self.device_ids[self.local_rank])) available_kv_cache_memory = total_available_memory - used_memory model_block_memory_used = self.cal_theortical_kvcache() available_kv_cache_memory += model_block_memory_used * self.parallel_config.total_block_num From e26313a355a359a7d7a1f58db7f8a7a0dae7169f Mon Sep 17 00:00:00 2001 From: plusNew001 <95567040+plusNew001@users.noreply.github.com> Date: Mon, 4 Aug 2025 16:25:33 +0800 Subject: [PATCH 07/17] Update Dockerfile.xpu (#3147) --- dockerfiles/Dockerfile.xpu | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/dockerfiles/Dockerfile.xpu b/dockerfiles/Dockerfile.xpu index a063cb84e32..74e7bf3e443 100644 --- a/dockerfiles/Dockerfile.xpu +++ b/dockerfiles/Dockerfile.xpu @@ -16,11 +16,17 @@ RUN apt-get update && apt-get install -y libibverbs-dev librdmacm-dev cmake pybi # uninstall existing package RUN python -m pip uninstall paddlepaddle-gpu paddlepaddle-xpu -y -# install paddlepaddle +# install paddlepaddle-xpu RUN python -m pip install --no-cache-dir --progress-bar off paddlepaddle-xpu==${PADDLE_VERSION} -i https://www.paddlepaddle.org.cn/packages/stable/xpu-p800/ RUN python -m pip install --no-cache-dir fastdeploy-xpu==${FD_VERSION} -i https://www.paddlepaddle.org.cn/packages/stable/fastdeploy-xpu-p800/ --extra-index-url https://mirrors.tuna.tsinghua.edu.cn/pypi/web/simple +RUN mkdir -p /workspace/deps && cd /workspace/deps && \ + wget https://klx-sdk-release-public.su.bcebos.com/xre/kl3-release/5.0.21.21/xre-Linux-x86_64-5.0.21.21.tar.gz && \ + tar -zxf xre-Linux-x86_64-5.0.21.21.tar.gz && mv xre-Linux-x86_64-5.0.21.21 xre + +ENV PATH=/workspace/deps/xre/bin:$PATH + ENV http_proxy="" ENV https_proxy="" ENV no_proxy="" From 9561603ed9e9eb90a4329ec8e4403b8889bc5fd2 Mon Sep 17 00:00:00 2001 From: YUNSHEN XIE <1084314248@qq.com> Date: Mon, 4 Aug 2025 16:30:56 +0800 Subject: [PATCH 08/17] Apply CI fix from Develop (#3151) * fix ci approve * Describe PR diff coverage using JSON file (#3114) * Refactored ci pipeline * update * Describe PR diff coverage using JSON file * remove pip cache setting from Approve * fix * update * fix ci (#3141) * fix --- .github/workflows/_build_linux.yml | 2 +- .github/workflows/_pre_ce_test.yml | 104 ++++++++++++++++++++++ .github/workflows/_unit_test_coverage.yml | 56 ++++++++++-- .github/workflows/approve.yml | 1 - .github/workflows/ci.yml | 73 +-------------- .github/workflows/pr_build_and_test.yml | 11 ++- scripts/coverage_run.sh | 22 ++++- scripts/{run_ci.sh => run_pre_ce.sh} | 3 - 8 files changed, 187 insertions(+), 85 deletions(-) create mode 100644 .github/workflows/_pre_ce_test.yml rename scripts/{run_ci.sh => run_pre_ce.sh} (94%) diff --git a/.github/workflows/_build_linux.yml b/.github/workflows/_build_linux.yml index cb02c64eccb..7c8fb23f4ee 100644 --- a/.github/workflows/_build_linux.yml +++ b/.github/workflows/_build_linux.yml @@ -124,6 +124,7 @@ jobs: echo "Date Only: $DATE_ONLY" export FASTDEPLOY_VERSION="${FASTDEPLOY_VERSION}.dev${DATE_ONLY}" fi + python -m pip install --pre paddlepaddle-gpu -i https://www.paddlepaddle.org.cn/packages/nightly/cu126/ pip config set global.index-url http://pip.baidu.com/root/baidu/+simple/ pip config set install.trusted-host pip.baidu.com pip config set global.extra-index-url https://mirrors.tuna.tsinghua.edu.cn/pypi/web/simple @@ -131,7 +132,6 @@ jobs: python -m pip install --upgrade pip python -m pip install -r requirements.txt python -m pip install wheel - python -m pip install --pre paddlepaddle-gpu -i https://www.paddlepaddle.org.cn/packages/nightly/cu126/ # 编译RDMA export ENABLE_FD_RDMA=1 bash build.sh 1 python false [${COMPILE_ARCH}] diff --git a/.github/workflows/_pre_ce_test.yml b/.github/workflows/_pre_ce_test.yml new file mode 100644 index 00000000000..629075bf136 --- /dev/null +++ b/.github/workflows/_pre_ce_test.yml @@ -0,0 +1,104 @@ +name: Pre-CE-Test + +on: + workflow_call: + inputs: + DOCKER_IMAGE: + description: "Build Images" + required: true + type: string + default: "ccr-2vdh3abv-pub.cnc.bj.baidubce.com/paddlepaddle/paddle:fastdeploy-ciuse-cuda126" + FASTDEPLOY_ARCHIVE_URL: + description: "URL of the compressed FastDeploy code archive." + required: true + type: string + FASTDEPLOY_WHEEL_URL: + description: "URL of the FastDeploy Wheel." + required: true + type: string + CACHE_DIR: + description: "Cache Dir Use" + required: false + type: string + default: "" + +concurrency: + group: ${{ github.event.pull_request.number }} + cancel-in-progress: true + +jobs: + run_ce_cases: + runs-on: [self-hosted, GPU-L20-4Card] + steps: + - name: Print current runner name + run: | + echo "Current runner name: ${{ runner.name }}" + - name: Code Prepare + shell: bash + env: + docker_image: ${{ inputs.DOCKER_IMAGE }} + fd_archive_url: ${{ inputs.FASTDEPLOY_ARCHIVE_URL }} + run: | + set -x + REPO="https://github.com/${{ github.repository }}.git" + FULL_REPO="${{ github.repository }}" + REPO_NAME="${FULL_REPO##*/}" + BASE_BRANCH="${{ github.base_ref }}" + + # Clean the repository directory before starting + docker run --rm --net=host -v $(pwd):/workspace -w /workspace \ + -e "REPO_NAME=${REPO_NAME}" \ + ${docker_image} /bin/bash -c ' + if [ -d ${REPO_NAME} ]; then + echo "Directory ${REPO_NAME} exists, removing it..." + rm -rf ${REPO_NAME}* + fi + ' + + wget -q ${fd_archive_url} + tar -xf FastDeploy.tar.gz + rm -rf FastDeploy.tar.gz + cd FastDeploy + git config --global user.name "FastDeployCI" + git config --global user.email "fastdeploy_ci@example.com" + git log -n 3 --oneline + + - name: Run CI unittest + env: + docker_image: ${{ inputs.DOCKER_IMAGE }} + fd_wheel_url: ${{ inputs.FASTDEPLOY_WHEEL_URL }} + run: | + runner_name="${{ runner.name }}" + last_char="${runner_name: -1}" + + if [ "${last_char}" = "1" ]; then + gpu_id=2 + DEVICES="2,3" + else + gpu_id=0 + DEVICES="0,1" + fi + FD_API_PORT=$((9180 + gpu_id * 100)) + FD_ENGINE_QUEUE_PORT=$((9150 + gpu_id * 100)) + FD_METRICS_PORT=$((9170 + gpu_id * 100)) + + PARENT_DIR=$(dirname "$WORKSPACE") + echo "PARENT_DIR:$PARENT_DIR" + docker run --rm --net=host -v $(pwd):/workspace -w /workspace \ + -v "/ssd4/GithubActions/gitconfig:/etc/gitconfig:ro" \ + -v "/ssd4/GithubActions/ModelData:/ModelData:ro" \ + -v "/ssd4/GithubActions/CacheDir:/root/.cache" \ + -v "/ssd4/GithubActions/ConfigDir:/root/.config" \ + -e "MODEL_PATH=/ModelData" \ + -e "FD_API_PORT=${FD_API_PORT}" \ + -e "FD_ENGINE_QUEUE_PORT=${FD_ENGINE_QUEUE_PORT}" \ + -e "FD_METRICS_PORT=${FD_METRICS_PORT}" \ + -e "fd_wheel_url=${fd_wheel_url}" \ + --gpus '"device='"${DEVICES}"'"' ${docker_image} /bin/bash -c ' + git config --global --add safe.directory /workspace/FastDeploy + cd FastDeploy + # python -m pip install --pre paddlepaddle-gpu -i https://www.paddlepaddle.org.cn/packages/nightly/cu126/ + python -m pip install paddlepaddle-gpu==3.0.0.dev20250729 -i https://www.paddlepaddle.org.cn/packages/nightly/cu126/ + python -m pip install ${fd_wheel_url} + bash scripts/run_pre_ce.sh + ' diff --git a/.github/workflows/_unit_test_coverage.yml b/.github/workflows/_unit_test_coverage.yml index 17b742cfe0f..17243768ab5 100644 --- a/.github/workflows/_unit_test_coverage.yml +++ b/.github/workflows/_unit_test_coverage.yml @@ -25,10 +25,11 @@ on: jobs: run_tests_with_coverage: - runs-on: [self-hosted, GPU-h1z1-4Cards] + runs-on: [self-hosted, GPU-h1z1-2Cards] outputs: diff_cov_file_url: ${{ steps.cov_upload.outputs.diff_cov_file_url }} - unittest_failed_url: ${{ steps.unittest_failed.outputs.unittest_failed_url }} + unittest_failed_url: ${{ steps.cov_upload.outputs.unittest_failed_url }} + diff_cov_result_json_url: ${{ steps.cov_upload.outputs.diff_cov_result_json_url }} steps: - name: Code Prepare shell: bash @@ -111,7 +112,7 @@ jobs: coverage combine coveragedata/ coverage xml -o python_coverage_all.xml COVERAGE_EXIT_CODE=0 - diff-cover python_coverage_all.xml --diff-file=diff.txt --fail-under=90 || COVERAGE_EXIT_CODE=9 + diff-cover python_coverage_all.xml --diff-file=diff.txt --fail-under=80 --json-report diff_coverage.json || COVERAGE_EXIT_CODE=9 echo "COVERAGE_EXIT_CODE=${COVERAGE_EXIT_CODE}" >> exit_code.env python scripts/generate_diff_coverage_xml.py diff.txt python_coverage_all.xml ' @@ -125,27 +126,68 @@ jobs: cd FastDeploy commit_id=${{ github.event.pull_request.head.sha }} pr_num=${{ github.event.pull_request.number }} - target_path=paddle-github-action/PR/FastDeploy/${pr_num}/${commit_id}/SM${compile_arch//,/_}/CoverageData + target_path=paddle-github-action/PR/FastDeploy/${pr_num}/${commit_id}/SM${compile_arch//,/_} wget -q --no-proxy --no-check-certificate https://paddle-qa.bj.bcebos.com/CodeSync/develop/PaddlePaddle/PaddleTest/tools/bos_tools.py push_file=$(realpath bos_tools.py) python -m pip install bce-python-sdk==0.9.29 diff_cov_file="diff_coverage.xml" if [ -f ${diff_cov_file} ];then - python ${push_file} ${diff_cov_file} ${target_path} + python ${push_file} ${diff_cov_file} ${target_path}/CoverageData target_path_stripped="${target_path#paddle-github-action/}" - DIFF_COV_FILE_URL=https://paddle-github-action.bj.bcebos.com/${target_path_stripped}/${diff_cov_file} + DIFF_COV_FILE_URL=https://paddle-github-action.bj.bcebos.com/${target_path_stripped}/CoverageData/${diff_cov_file} echo "diff_cov_file_url=${DIFF_COV_FILE_URL}" >> $GITHUB_OUTPUT + echo "diff_cov_file_url=${DIFF_COV_FILE_URL}" >> $GITHUB_ENV fi - - name: Determine Unit Succ and whether the coverage rate reaches 90% + diff_cov_result_json="diff_coverage.json" + if [ -f ${diff_cov_result_json} ];then + python ${push_file} ${diff_cov_result_json} ${target_path}/CoverageData + target_path_stripped="${target_path#paddle-github-action/}" + DIFF_COV_JSON_URL=https://paddle-github-action.bj.bcebos.com/${target_path_stripped}/CoverageData/${diff_cov_result_json} + echo "diff_cov_result_json_url=${DIFF_COV_JSON_URL}" >> $GITHUB_OUTPUT + echo "diff_cov_result_json_url=${DIFF_COV_JSON_URL}" >> $GITHUB_ENV + fi + unittest_result="test/failed_tests.log" + if [ -s ${unittest_result} ];then + python ${push_file} ${unittest_result} ${target_path}/UnitTestResult + target_path_stripped="${target_path#paddle-github-action/}" + UNIT_TEST_RESULT_URL=https://paddle-github-action.bj.bcebos.com/${target_path_stripped}/UnitTestResult/${unittest_result} + echo "unittest_failed_url=${UNIT_TEST_RESULT_URL}" >> $GITHUB_OUTPUT + echo "unittest_failed_url=${UNIT_TEST_RESULT_URL}" >> $GITHUB_ENV + fi + - name: Determine Unit Succ and whether the coverage rate reaches 80% shell: bash run: | if [ "$TEST_EXIT_CODE" -eq 8 ]; then + if [ -z "${unittest_failed_url}" ]; then + echo "No diff unit failed file URL provided." + else + wget ${unittest_failed_url} || echo "Download unittest file failed, but continuing..." + fi echo "Unit tests failed (exit code 8)" + filename=$(basename "$unittest_failed_url") + if [ -f "${filename}" ];then + echo "Failed test cases:" + cat "${filename}" + fi exit "$TEST_EXIT_CODE" fi if [ "$COVERAGE_EXIT_CODE" -eq 9 ]; then echo "Coverage generation failed (exit code 9)" + if [ -z "${diff_cov_result_json_url}" ]; then + echo "No diff cov result file URL provided." + else + wget ${diff_cov_result_json_url} || echo "Download cov json file failed, but continuing..." + fi + filename=$(basename "$diff_cov_result_json_url") + if [ -f "${filename}" ];then + echo "Failed test cases:" + if command -v jq >/dev/null 2>&1; then + jq . "${filename}" + else + cat "${filename}" + fi + fi exit "$COVERAGE_EXIT_CODE" fi echo "All tests and coverage passed" diff --git a/.github/workflows/approve.yml b/.github/workflows/approve.yml index bf82f820006..baa953ab5af 100644 --- a/.github/workflows/approve.yml +++ b/.github/workflows/approve.yml @@ -33,7 +33,6 @@ jobs: uses: actions/setup-python@v5 with: python-version: '3.10' - cache: 'pip' - name: Run approval check script run: | diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 518b15eb998..b07d099de1b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -13,77 +13,8 @@ concurrency: jobs: build: - runs-on: [self-hosted, GPU-L20-4Card] + runs-on: ubuntu-latest steps: - name: Print current runner name run: | - echo "Current runner name: ${{ runner.name }}" - # Because the system version is lower than 2.23, the checkout cannot be used. - # - name: Checkout code - # uses: actions/checkout@v4 - - - name: Code Checkout - env: - docker_image: ccr-2vdh3abv-pub.cnc.bj.baidubce.com/paddlepaddle/paddle:fastdeploy-ciuse-cuda126 - run: | - REPO="https://github.com/${{ github.repository }}.git" - FULL_REPO="${{ github.repository }}" - REPO_NAME="${FULL_REPO##*/}" - BASE_BRANCH="${{ github.base_ref }}" - # Clean the repository directory before starting - docker run --rm --net=host -v $(pwd):/workspace -w /workspace \ - -e "REPO_NAME=${REPO_NAME}" \ - -e "BASE_BRANCH=${BASE_BRANCH}" \ - ${docker_image} /bin/bash -c ' - if [ -d ${REPO_NAME} ]; then - echo "Directory ${REPO_NAME} exists, removing it..." - rm -rf ${REPO_NAME} - fi - ' - git config --global user.name "FastDeployCI" - git config --global user.email "fastdeploy_ci@example.com" - git clone ${REPO} ${REPO_NAME} -b ${BASE_BRANCH} - cd FastDeploy - if [ "${{ github.event_name }}" = "pull_request" ]; then - git fetch origin pull/${{ github.event.pull_request.number }}/head:pr/${{ github.event.pull_request.number }} - git merge pr/${{ github.event.pull_request.number }} - git log -n 3 --oneline - else - git checkout ${{ github.sha }} - git log -n 3 --oneline - fi - - - name: Run CI unittest - env: - docker_image: ccr-2vdh3abv-pub.cnc.bj.baidubce.com/paddlepaddle/paddle:fastdeploy-ciuse-cuda126 - run: | - runner_name="${{ runner.name }}" - last_char="${runner_name: -1}" - - if [ "${last_char}" = "1" ]; then - gpu_id=2 - DEVICES="2,3" - else - gpu_id=0 - DEVICES="0,1" - fi - FD_API_PORT=$((9180 + gpu_id * 100)) - FD_ENGINE_QUEUE_PORT=$((9150 + gpu_id * 100)) - FD_METRICS_PORT=$((9170 + gpu_id * 100)) - - PARENT_DIR=$(dirname "$WORKSPACE") - echo "PARENT_DIR:$PARENT_DIR" - docker run --rm --net=host -v $(pwd):/workspace -w /workspace \ - -v "/ssd4/GithubActions/gitconfig:/etc/gitconfig:ro" \ - -v "/ssd4/GithubActions/ModelData:/ModelData:ro" \ - -v "/ssd4/GithubActions/CacheDir:/root/.cache" \ - -v "/ssd4/GithubActions/ConfigDir:/root/.config" \ - -e "MODEL_PATH=/ModelData" \ - -e "FD_API_PORT=${FD_API_PORT}" \ - -e "FD_ENGINE_QUEUE_PORT=${FD_ENGINE_QUEUE_PORT}" \ - -e "FD_METRICS_PORT=${FD_METRICS_PORT}" \ - --gpus '"device='"${DEVICES}"'"' ${docker_image} /bin/bash -c " - git config --global --add safe.directory /workspace/FastDeploy - cd FastDeploy - bash scripts/run_ci.sh - " + echo "The current CI tasks have been migrated to PreCe and will be deprecated soon." diff --git a/.github/workflows/pr_build_and_test.yml b/.github/workflows/pr_build_and_test.yml index d6557fc625b..0123e5a5540 100644 --- a/.github/workflows/pr_build_and_test.yml +++ b/.github/workflows/pr_build_and_test.yml @@ -21,7 +21,7 @@ jobs: with: DOCKER_IMAGE: ccr-2vdh3abv-pub.cnc.bj.baidubce.com/paddlepaddle/paddleqa:cuda126-py310 FASTDEPLOY_ARCHIVE_URL: ${{ needs.clone.outputs.repo_archive_url }} - COMPILE_ARCH: "90" + COMPILE_ARCH: "89,90" WITH_NIGHTLY_BUILD: "OFF" FD_VERSION: "0.0.0" @@ -52,3 +52,12 @@ jobs: PADDLETEST_ARCHIVE_URL: "https://xly-devops.bj.bcebos.com/PaddleTest/PaddleTest.tar.gz" FASTDEPLOY_WHEEL_URL: ${{ needs.build.outputs.wheel_path }} MODEL_CACHE_DIR: "/ssd2/actions-runner/ModelCache" + + pre_ce_test: + name: Extracted partial CE model tasks to run in CI. + needs: [clone,build] + uses: ./.github/workflows/_pre_ce_test.yml + with: + DOCKER_IMAGE: ccr-2vdh3abv-pub.cnc.bj.baidubce.com/paddlepaddle/paddle:fastdeploy-ciuse-cuda126 + FASTDEPLOY_ARCHIVE_URL: ${{ needs.clone.outputs.repo_archive_url }} + FASTDEPLOY_WHEEL_URL: ${{ needs.build.outputs.wheel_path }} diff --git a/scripts/coverage_run.sh b/scripts/coverage_run.sh index 98ed025bdc9..6b6cbbf850d 100644 --- a/scripts/coverage_run.sh +++ b/scripts/coverage_run.sh @@ -6,7 +6,23 @@ run_path="$DIR/../test/" cd ${run_path} ls -dirs=("layers" "operators" "worker" "utils") +exclude=("ci_use" "ce") +for d in */ ; do + dir_name="${d%/}" + if [[ -d "$dir_name" ]]; then + skip=false + for ex in "${exclude[@]}"; do + if [[ "$dir_name" == "$ex" ]]; then + skip=true + break + fi + done + if ! $skip; then + dirs+=("$dir_name") + fi + fi +done + failed_tests_file="failed_tests.log" > "$failed_tests_file" disabled_tests=( @@ -20,6 +36,10 @@ disabled_tests=( operators/test_stop_generation.py operators/test_air_topp_sampling.py operators/test_fused_moe.py + layers/test_repetition_early_stopper.py + operators/test_stop_generation_multi_ends.py + utils/test_download.py + graph_optimization/test_cuda_graph.py ) is_disabled() { local test_file_rel="$1" diff --git a/scripts/run_ci.sh b/scripts/run_pre_ce.sh similarity index 94% rename from scripts/run_ci.sh rename to scripts/run_pre_ce.sh index 91ef179b75c..726b91e8579 100644 --- a/scripts/run_ci.sh +++ b/scripts/run_pre_ce.sh @@ -3,13 +3,10 @@ DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" echo "$DIR" # python -m pip install --pre paddlepaddle-gpu -i https://www.paddlepaddle.org.cn/packages/nightly/cu126/ -python -m pip install paddlepaddle-gpu==3.0.0.dev20250729 -i https://www.paddlepaddle.org.cn/packages/nightly/cu126/ python -m pip config set global.index-url https://mirrors.tuna.tsinghua.edu.cn/pypi/web/simple - python -m pip install -r requirements.txt python -m pip install jsonschema aistudio_sdk==0.3.5 -bash build.sh || exit 1 failed_files=() run_path="$DIR/../test/ci_use/" From bd77a3a6435fdfdcb149db03226a602c89c0a25e Mon Sep 17 00:00:00 2001 From: RAM Date: Tue, 5 Aug 2025 10:53:27 +0800 Subject: [PATCH 09/17] [Bug Fix] Fix bug of MLA Attention Backend (#3178) * fix typo * fix mla attention backend --- fastdeploy/model_executor/models/deepseek_v3.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/fastdeploy/model_executor/models/deepseek_v3.py b/fastdeploy/model_executor/models/deepseek_v3.py index e4b44f477ba..8cbd4a0bdd6 100644 --- a/fastdeploy/model_executor/models/deepseek_v3.py +++ b/fastdeploy/model_executor/models/deepseek_v3.py @@ -315,7 +315,7 @@ def forward( dtype=layernorm_out.dtype, ) - if forward_meta.max_enc_len_this_time: + if forward_meta.max_len_tensor_cpu[1]: # max_enc_len_this_time query = self.q_a_proj(layernorm_out) query = self.q_a_layernorm(query) query = self.q_b_proj(query) @@ -362,7 +362,7 @@ def forward( fmha_out_prefill = fmha_out_prefill * mask_encoder_batch.cast(fmha_out_prefill.dtype) fmha_out = fmha_out + fmha_out_prefill - if forward_meta.max_dec_len_this_time: + if forward_meta.max_len_tensor_cpu[2]: # max_dec_len_this_time query = self.q_a_proj(layernorm_out) query = self.q_a_layernorm(query) ln_out_or_q_c = query From 3dd8492601a7084902161438396a31a8ce9f10e0 Mon Sep 17 00:00:00 2001 From: SunLei Date: Tue, 5 Aug 2025 10:55:22 +0800 Subject: [PATCH 10/17] [Bugfix] Fix uninitialized decoded_token and add corresponding unit test (#3201) * Update test_base_chat.py (#3183) * [Bugfix] Fix uninitialized decoded_token and add corresponding unit test. --------- Co-authored-by: Divano --- fastdeploy/entrypoints/llm.py | 7 +- test/ce/server/test_base_chat.py | 221 ++++++++++++++++++ .../openai/test_build_sample_logprobs.py | 78 +++++++ 3 files changed, 305 insertions(+), 1 deletion(-) create mode 100644 test/ce/server/test_base_chat.py create mode 100644 test/entrypoints/openai/test_build_sample_logprobs.py diff --git a/fastdeploy/entrypoints/llm.py b/fastdeploy/entrypoints/llm.py index 8365c698535..3e150abf2d9 100644 --- a/fastdeploy/entrypoints/llm.py +++ b/fastdeploy/entrypoints/llm.py @@ -285,6 +285,10 @@ def _add_request( self.llm_engine.add_requests(tasks, current_sampling_params, enable_thinking=enable_thinking) return req_ids + def _decode_token(self, token_id: int) -> str: + """Decodes a single token ID into its string representation.""" + return self.llm_engine.data_processor.process_logprob_response([token_id], clean_up_tokenization_spaces=False) + def _build_sample_logprobs(self, logprobs_lists: LogprobsLists, topk_logprobs: int) -> list[dict[int, Logprob]]: """ Constructs a list of dictionaries mapping token IDs to Logprob objects, @@ -318,8 +322,9 @@ def _build_sample_logprobs(self, logprobs_lists: LogprobsLists, topk_logprobs: i sliced_logprobs_lists = logprobs_lists.slice_columns(1, 1 + effective_topk_logprobs) result = [] for token_ids, logprobs in zip(sliced_logprobs_lists.logprob_token_ids, sliced_logprobs_lists.logprobs): + logprob_dict = { - token_id: Logprob(logprob=logprob, rank=i + 1, decoded_token=None) + token_id: Logprob(logprob=logprob, rank=i + 1, decoded_token=self._decode_token(token_id)) for i, (token_id, logprob) in enumerate(zip(token_ids, logprobs)) } result.append(logprob_dict) diff --git a/test/ce/server/test_base_chat.py b/test/ce/server/test_base_chat.py new file mode 100644 index 00000000000..12be895fe65 --- /dev/null +++ b/test/ce/server/test_base_chat.py @@ -0,0 +1,221 @@ +#!/bin/env python3 +# -*- coding: utf-8 -*- +# @author DDDivano +# encoding=utf-8 vi:ts=4:sw=4:expandtab:ft=python + +""" +some basic check for fd web api +""" + +import json + +from core import TEMPLATE, URL, build_request_payload, send_request + + +def test_stream_response(): + data = { + "stream": True, + "messages": [ + {"role": "system", "content": "你是一个知识渊博的 AI 助手"}, + {"role": "user", "content": "讲讲爱因斯坦的相对论"}, + ], + "max_tokens": 10, + } + payload = build_request_payload(TEMPLATE, data) + resp = send_request(URL, payload, stream=True) + + output = "" + for line in resp.iter_lines(decode_unicode=True): + if line.strip() == "" or not line.startswith("data: "): + continue + line = line[len("data: ") :] + if line.strip() == "[DONE]": + break + chunk = json.loads(line) + delta = chunk.get("choices", [{}])[0].get("delta", {}) + output += delta.get("content", "") + + print("Stream输出:", output) + assert "相对论" in output or len(output) > 0 + + +def test_system_prompt_effect(): + data = { + "stream": False, + "messages": [ + {"role": "system", "content": "请用一句话回答"}, + {"role": "user", "content": "什么是人工智能?"}, + ], + "max_tokens": 30, + } + payload = build_request_payload(TEMPLATE, data) + resp = send_request(URL, payload).json() + content = resp["choices"][0]["message"]["content"] + print("内容输出:", content) + assert len(content) < 50 + + +def test_logprobs_enabled(): + data = { + "stream": False, + "logprobs": True, + "top_logprobs": 5, + "messages": [{"role": "user", "content": "非洲的首都是?"}], + "max_tokens": 3, + } + payload = build_request_payload(TEMPLATE, data) + resp = send_request(URL, payload).json() + logprob_data = resp["choices"][0].get("logprobs") + print("LogProbs:", logprob_data) + assert logprob_data is not None + content_logprobs = logprob_data.get("content", []) + assert isinstance(content_logprobs, list) + assert all("token" in item for item in content_logprobs) + + +def test_stop_sequence(): + data = { + "stream": False, + "stop": ["果冻"], + "messages": [ + { + "role": "user", + "content": "你要严格按照我接下来的话输出,输出冒号后面的内容,请输出:这是第一段。果冻这是第二段啦啦啦啦啦。", + }, + ], + "max_tokens": 20, + "top_p": 0, + } + payload = build_request_payload(TEMPLATE, data) + resp = send_request(URL, payload).json() + content = resp["choices"][0]["message"]["content"] + print("截断输出:", content) + assert "第二段" not in content + + +def test_sampling_parameters(): + data = { + "stream": False, + "temperature": 0, + "top_p": 0, + "messages": [ + {"role": "user", "content": "1+1=?,直接回答答案"}, + ], + "max_tokens": 50, + } + payload = build_request_payload(TEMPLATE, data) + resp = send_request(URL, payload).json() + answer = resp["choices"][0]["message"]["content"] + print("Sampling输出:", answer) + assert any(ans in answer for ans in ["2", "二"]) + + +def test_multi_turn_conversation(): + data = { + "stream": False, + "messages": [ + {"role": "user", "content": "牛顿是谁?"}, + {"role": "assistant", "content": "牛顿是一位物理学家。"}, + {"role": "user", "content": "他提出了什么理论?"}, + ], + "max_tokens": 30, + } + payload = build_request_payload(TEMPLATE, data) + resp = send_request(URL, payload).json() + content = resp["choices"][0]["message"]["content"] + print("多轮记忆:", content) + assert "三大运动定律" in content or "万有引力" in content + + +def test_bad_words_filtering(): + banned_tokens = ["和", "呀"] + + data = { + "stream": False, + "messages": [ + {"role": "system", "content": "你是一个助手,回答简洁清楚"}, + {"role": "user", "content": "请输出冒号后面的字: 我爱吃果冻,和苹果,香蕉,和荔枝"}, + ], + "top_p": 0, + "max_tokens": 69, + "bad_words": banned_tokens, + } + + payload = build_request_payload(TEMPLATE, data) + response = send_request(URL, payload).json() + + content = response["choices"][0]["message"]["content"] + print("生成内容:", content) + + for word in banned_tokens: + assert word not in content, f"bad_word '{word}' 不应出现在生成结果中" + + print("test_bad_words_filtering 通过:生成结果未包含被禁词") + + data = { + "stream": False, + "messages": [ + {"role": "system", "content": "你是一个助手,回答简洁清楚"}, + {"role": "user", "content": "请输出冒号后面的字,一模一样: 我爱吃果冻,苹果,香蕉,和荔枝呀呀呀"}, + ], + "top_p": 0, + "max_tokens": 69, + # "bad_words": banned_tokens, + } + + payload = build_request_payload(TEMPLATE, data) + response = send_request(URL, payload).json() + + content = response["choices"][0]["message"]["content"] + print("生成内容:", content) + + for word in banned_tokens: + assert word not in content, f"bad_word '{word}' 不应出现在生成结果中" + + print("test_bad_words_filtering 通过:生成结果未包含被禁词") + + +def test_bad_words_filtering1(): + banned_tokens = ["和", "呀"] + + data = { + "stream": False, + "messages": [ + {"role": "system", "content": "你是一个助手,回答简洁清楚"}, + {"role": "user", "content": "请输出冒号后面的字: 我爱吃果冻,和苹果,香蕉,和荔枝"}, + ], + "top_p": 0, + "max_tokens": 69, + "bad_words": banned_tokens, + } + + payload = build_request_payload(TEMPLATE, data) + response = send_request(URL, payload).json() + + content = response["choices"][0]["message"]["content"] + print("生成内容:", content) + + for word in banned_tokens: + assert word not in content, f"bad_word '{word}' 不应出现在生成结果中" + + print("test_bad_words_filtering 通过:生成结果未包含被禁词") + word = "呀呀" + data = { + "stream": False, + "messages": [ + {"role": "system", "content": "你是一个助手,回答简洁清楚"}, + {"role": "user", "content": "请输出冒号后面的字,一模一样: 我爱吃果冻,苹果,香蕉,和荔枝呀呀呀"}, + ], + "top_p": 0, + "max_tokens": 69, + } + + payload = build_request_payload(TEMPLATE, data) + response = send_request(URL, payload).json() + + content = response["choices"][0]["message"]["content"] + print("生成内容:", content) + + assert word in content, f" '{word}' 应出现在生成结果中" + + print("test_bad_words_filtering 通过:生成结果未包含被禁词") diff --git a/test/entrypoints/openai/test_build_sample_logprobs.py b/test/entrypoints/openai/test_build_sample_logprobs.py new file mode 100644 index 00000000000..76ff8e87b7b --- /dev/null +++ b/test/entrypoints/openai/test_build_sample_logprobs.py @@ -0,0 +1,78 @@ +import unittest +from unittest.mock import MagicMock, patch + +from fastdeploy.entrypoints.llm import LLM +from fastdeploy.worker.output import Logprob, LogprobsLists + + +def get_patch_path(cls, method="__init__"): + return f"{cls.__module__}.{cls.__qualname__}.{method}" + + +class TestBuildSampleLogprobs(unittest.TestCase): + + def setUp(self): + """ + Set up the test environment by creating an instance of the LLM class using Mock. + """ + patch_llm = get_patch_path(LLM) + with patch(patch_llm, return_value=None): + self.llm = LLM() + # mock d data_processor + self.llm.llm_engine = MagicMock() + self.llm.llm_engine.data_processor.process_logprob_response.side_effect = ( + lambda ids, **kwargs: f"token_{ids[0]}" + ) + + def test_build_sample_logprobs_basic(self): + """ + Test case for building sample logprobs when `topk_logprobs` is valid. + """ + logprob_token_ids = [[100, 101, 102]] + logprobs = [[-0.1, -0.5, -1.0]] + sampled_token_ranks = [0] + + logprobs_lists = LogprobsLists( + logprob_token_ids=logprob_token_ids, logprobs=logprobs, sampled_token_ranks=sampled_token_ranks + ) + + result = self.llm._build_sample_logprobs(logprobs_lists, topk_logprobs=2) + + expected = [ + { + 101: Logprob(logprob=-0.5, rank=1, decoded_token="token_101"), + 102: Logprob(logprob=-1.0, rank=2, decoded_token="token_102"), + } + ] + + self.assertEqual(result, expected) + + def test_build_sample_logprobs_empty_input(self): + """ + Test case where `logprob_token_ids` is empty. + """ + logprobs_lists = MagicMock(spec=LogprobsLists) + logprobs_lists.logprob_token_ids = [] + result = self.llm._build_sample_logprobs(logprobs_lists, topk_logprobs=2) + self.assertIsNone(result) + + def test_build_sample_logprobs_invalid_topk(self): + """ + Test case where `topk` value exceeds length of first element in `logprob_token_ids`. + """ + logprobs_lists = MagicMock(spec=LogprobsLists) + logprobs_lists.logprob_token_ids = [[100]] + result = self.llm._build_sample_logprobs(logprobs_lists, topk_logprobs=2) + self.assertIsNone(result) + + def test_decode_token(self): + """ + Test case for decoding a single token ID. + """ + token_id = 123 + decoded = self.llm._decode_token(token_id) + self.assertEqual(decoded, "token_123") + + +if __name__ == "__main__": + unittest.main() From bc0b92bba439c5462e5c055f8ba4c7430faf96b7 Mon Sep 17 00:00:00 2001 From: lizexu123 <39205361+lizexu123@users.noreply.github.com> Date: Wed, 6 Aug 2025 14:30:33 +0800 Subject: [PATCH 11/17] [BugFix] support real batch_size (#3109) (#3217) * support real bsz * fix * fix xpu_model_runner.py,gpu_model_runner.py,gcu_model_runner.py,iluvatar_model_runner.py * add event_loop_ep * fix * Add comments * fix * support mtp real_batch_size * fix * self.tmp_seq_lens_this_time->self.seq_lens_this_time_buffer * fix * fix VL real_seq_lens_this_time * fix * fix mtp * fix * fix mtp * fix xpu * fix --- fastdeploy/spec_decode/mtp.py | 15 ++++---- fastdeploy/worker/gcu_model_runner.py | 23 ++++++++---- fastdeploy/worker/gcu_worker.py | 7 ++-- fastdeploy/worker/gpu_model_runner.py | 41 +++++++++++++++------- fastdeploy/worker/gpu_worker.py | 9 ++--- fastdeploy/worker/iluvatar_model_runner.py | 22 ++++++++---- fastdeploy/worker/iluvatar_worker.py | 7 ++-- fastdeploy/worker/worker_process.py | 8 ++--- fastdeploy/worker/xpu_model_runner.py | 23 ++++++++---- fastdeploy/worker/xpu_worker.py | 13 ++++--- 10 files changed, 110 insertions(+), 58 deletions(-) diff --git a/fastdeploy/spec_decode/mtp.py b/fastdeploy/spec_decode/mtp.py index 39f0fce4272..3033e41467f 100644 --- a/fastdeploy/spec_decode/mtp.py +++ b/fastdeploy/spec_decode/mtp.py @@ -107,7 +107,7 @@ def dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decode idx = i self.model_inputs["input_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.model_inputs["eos_token_id"][:] = np.array([2], dtype="int64").reshape(-1, 1) - self.model_inputs["seq_lens_this_time"][idx : idx + 1] = input_length + self.seq_lens_this_time_buffer[idx : idx + 1] = input_length self.model_inputs["seq_lens_encoder"][idx : idx + 1] = input_length self.model_inputs["seq_lens_decoder"][idx : idx + 1] = 0 self.model_inputs["step_idx"][idx : idx + 1] = 0 @@ -118,6 +118,7 @@ def dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decode self.model_inputs["block_tables"][idx : idx + 1, :block_num] = np.arange( idx * block_num, (idx + 1) * block_num, 1 ) + self.model_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer def initialize_kv_cache(self): """ @@ -263,7 +264,8 @@ def _init_model_inputs(self): # Same shape/dytpe with base model self.model_inputs["block_tables"] = paddle.clone(self.main_model_inputs["block_tables"]) self.model_inputs["input_ids"] = paddle.clone(self.main_model_inputs["input_ids"]) - self.model_inputs["seq_lens_this_time"] = paddle.clone(self.main_model_inputs["seq_lens_this_time"]) + self.seq_lens_this_time_buffer = paddle.clone(self.main_model_inputs["seq_lens_this_time"]) + self.model_inputs["seq_lens_encoder"] = paddle.clone(self.main_model_inputs["seq_lens_encoder"]) self.model_inputs["seq_lens_decoder"] = paddle.clone(self.main_model_inputs["seq_lens_decoder"]) self.model_inputs["step_idx"] = paddle.clone(self.main_model_inputs["step_idx"]) @@ -338,7 +340,7 @@ def _init_model_inputs(self): self.main_model_inputs["seq_lens_this_time"], fill_value=-1, dtype="int32" ) - def insert_prefill_inputs(self, req_dicts: List[Request]): + def insert_prefill_inputs(self, req_dicts: List[Request], num_running_requests: int): """ Process inputs for prefill tasks and insert it to model_inputs buffer """ @@ -372,7 +374,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): self.model_inputs["seq_lens_encoder"][idx : idx + 1] = 0 self.model_inputs["seq_lens_decoder"][idx : idx + 1] = length - self.model_inputs["seq_lens_this_time"][idx : idx + 1] = prefill_token_num + self.seq_lens_this_time_buffer[idx : idx + 1] = prefill_token_num self.model_inputs["stop_flags"][idx : idx + 1] = False self.model_inputs["batch_drop"][idx : idx + 1] = False @@ -397,10 +399,10 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): if self.cache_config.enable_chunked_prefill: token_chunk_size = request.prefill_chunk_info[0] self.model_inputs["seq_lens_encoder"][idx : idx + 1] = token_chunk_size - self.model_inputs["seq_lens_this_time"][idx : idx + 1] = token_chunk_size + self.seq_lens_this_time_buffer[idx : idx + 1] = token_chunk_size else: self.model_inputs["seq_lens_encoder"][idx : idx + 1] = length - self.model_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.model_inputs["seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) self.model_inputs["stop_flags"][idx : idx + 1] = False @@ -413,6 +415,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): request.get("block_tables"), dtype="int32" ) self.model_inputs["not_need_stop"][0] = True + self.model_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] def _initialize_forward_meta(self): """ diff --git a/fastdeploy/worker/gcu_model_runner.py b/fastdeploy/worker/gcu_model_runner.py index 531304017b8..e35dc4d23ad 100644 --- a/fastdeploy/worker/gcu_model_runner.py +++ b/fastdeploy/worker/gcu_model_runner.py @@ -152,9 +152,11 @@ def _init_logits_processor(self, request): schemata_key, ) - def insert_prefill_inputs(self, req_dicts: List[Request]): + def insert_prefill_inputs(self, req_dicts: List[Request], num_running_requests: int = None): """ Process inputs for prefill tasks and insert it to share_inputs buffer + req_dict: A list of Request dict + num_running_requests: batch_size """ if req_dicts[-1].disaggregate_info is not None and req_dicts[-1].disaggregate_info["role"] == "prefill": @@ -193,7 +195,7 @@ def get_attr_from_request(request, attr, default_value=None): self.share_inputs["prompt_ids"][idx : idx + 1, :length] = np.array(request.prompt_token_ids) self.share_inputs["seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["seq_lens_decoder"][idx : idx + 1] = length - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = 1 + self.seq_lens_this_time_buffer[idx : idx + 1] = 1 self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = length self.share_inputs["prompt_lens"][idx : idx + 1] = length @@ -205,7 +207,7 @@ def get_attr_from_request(request, attr, default_value=None): request.draft_token_ids[0:num_prefill_send_token], dtype="int64", ) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = num_prefill_send_token + self.seq_lens_this_time_buffer[idx : idx + 1] = num_prefill_send_token else: self.share_inputs["pre_ids"][idx : idx + 1] = -1 self.share_inputs["step_idx"][idx : idx + 1] = 0 @@ -222,14 +224,14 @@ def get_attr_from_request(request, attr, default_value=None): ) self.share_inputs["seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = token_chunk_size + self.seq_lens_this_time_buffer[idx : idx + 1] = token_chunk_size self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = token_chunk_size self.share_inputs["seq_lens_encoder"][idx : idx + 1] = token_chunk_size self.share_inputs["prompt_lens"][idx : idx + 1] = token_chunk_size else: self.share_inputs["seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["prompt_lens"][idx : idx + 1] = length @@ -293,6 +295,7 @@ def get_attr_from_request(request, attr, default_value=None): if self.speculative_method in ["mtp"]: self.proposer.insert_prefill_inputs(req_dicts) + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decode_len: int): """Set dummy prefill inputs to share_inputs""" @@ -311,7 +314,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decod self.share_inputs["input_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["prompt_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["eos_token_id"][:] = np.array([2], dtype="int64").reshape(-1, 1) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = input_length + self.seq_lens_this_time_buffer[idx : idx + 1] = input_length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 @@ -329,6 +332,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decod self.share_inputs["block_tables"][idx : idx + 1, :block_num] = np.arange( idx * block_num, (idx + 1) * block_num, 1 ) + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer def _init_share_inputs(self, max_num_seqs: int): """ @@ -379,7 +383,7 @@ def _init_share_inputs(self, max_num_seqs: int): self.share_inputs["max_length"] = paddle.full( [max_num_seqs, 1], self.model_config.max_model_len, dtype="int64" ) - self.share_inputs["seq_lens_this_time"] = paddle.full(max_num_seqs, 0, dtype="int32") + self.seq_lens_this_time_buffer = paddle.full(max_num_seqs, 0, dtype="int32") self.share_inputs["seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["seq_lens_decoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["step_seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") @@ -921,6 +925,7 @@ def _get_skip_idx(self, model_forward_batch: Optional[List[Request]] = None): def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + num_running_requests: int = None, ) -> Optional[ModelRunnerOutput]: """ The Entrance of model execute. @@ -928,6 +933,7 @@ def execute_model( model_forward_batch: 'Request' contains information related to prompt and is an abstract class at the server level, which is too granular for ModelRunner. We plan to replace it with 'ModelForwardBatch'. + num_running_requests: batch_size intermediate_tensors: """ # If `not_need_stop`` is False, it means the current worker is in an idle state. @@ -1053,6 +1059,9 @@ class at the server level, which is too granular for ModelRunner. self._update_chunked_prefill(model_forward_batch) self._add_cache(model_forward_batch) + self.seq_lens_this_time_buffer[:num_running_requests].copy_( + self.share_inputs["seq_lens_this_time"][:num_running_requests], False + ) return None def _add_cache(self, model_forward_batch) -> None: diff --git a/fastdeploy/worker/gcu_worker.py b/fastdeploy/worker/gcu_worker.py index 77a8a50d4b9..a168367809a 100644 --- a/fastdeploy/worker/gcu_worker.py +++ b/fastdeploy/worker/gcu_worker.py @@ -105,17 +105,18 @@ def initialize_cache(self, num_gpu_blocks: int) -> None: def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + num_running_requests: int = None, ) -> Optional[ModelRunnerOutput]: """ """ - output = self.model_runner.execute_model(model_forward_batch) + output = self.model_runner.execute_model(model_forward_batch, num_running_requests) return output - def preprocess_new_task(self, req_dicts: List[Request]) -> None: + def preprocess_new_task(self, req_dicts: List[Request], num_running_requests: int) -> None: """Process new requests and then start the decode loop TODO(gongshaotian):The scheduler should schedule the handling of prefill, and workers and modelrunners should not perceive it. """ - self.model_runner.insert_prefill_inputs(req_dicts=req_dicts) + self.model_runner.insert_prefill_inputs(req_dicts=req_dicts, num_running_requests=num_running_requests) def graph_optimize_and_warm_up_model(self) -> None: """ diff --git a/fastdeploy/worker/gpu_model_runner.py b/fastdeploy/worker/gpu_model_runner.py index 4b67b595e84..c551364ef2a 100644 --- a/fastdeploy/worker/gpu_model_runner.py +++ b/fastdeploy/worker/gpu_model_runner.py @@ -164,6 +164,7 @@ def _init_speculative_proposer(self): if self.speculative_method == "ngram": self.proposer = NgramProposer(self.fd_config) elif self.speculative_method == "mtp": + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer self.proposer = MTPProposer( self.fd_config, self.get_model(), @@ -193,9 +194,11 @@ def _init_logits_processor(self, request): return self.guided_backend.get_logits_processor(schemata_key=schemata_key), schemata_key - def insert_tasks_v1(self, req_dicts: List[Request]): + def insert_tasks_v1(self, req_dicts: List[Request], num_running_requests: int = None): """ Process scheduler output tasks, used when ENABLE_V1_KVCACHE_SCHEDULER=1 + req_dict: A list of Request dict + num_running_requests: batch_size """ # NOTE(luotingdan): Lazy initialize kv cache if "caches" not in self.share_inputs: @@ -264,7 +267,7 @@ def insert_tasks_v1(self, req_dicts: List[Request]): ) self.share_inputs["stop_flags"][idx : idx + 1] = False self.share_inputs["seq_lens_decoder"][idx : idx + 1] = prefill_start_index - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = 0 self.share_inputs["prompt_lens"][idx : idx + 1] = len(input_ids) @@ -286,7 +289,7 @@ def insert_tasks_v1(self, req_dicts: List[Request]): logger.debug(f"Handle preempted request {request} at idx {idx}") self.share_inputs["block_tables"][idx : idx + 1, :] = -1 self.share_inputs["stop_flags"][idx : idx + 1] = True - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = 0 + self.seq_lens_this_time_buffer[idx : idx + 1] = 0 self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 self.share_inputs["seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["is_block_step"][idx : idx + 1] = False @@ -328,10 +331,13 @@ def insert_tasks_v1(self, req_dicts: List[Request]): if has_prefill_task: self.share_inputs["not_need_stop"][0] = True + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] - def insert_prefill_inputs(self, req_dicts: List[Request]): + def insert_prefill_inputs(self, req_dicts: List[Request], num_running_requests: int = None): """ Process inputs for prefill tasks and insert it to share_inputs buffer + req_dict: A list of Request dict + num_running_requests: batch_size TODO(gongshaotian): Refactor this func """ @@ -365,7 +371,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): self.share_inputs["prompt_ids"][idx : idx + 1, :length] = np.array(request.prompt_token_ids) self.share_inputs["seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["seq_lens_decoder"][idx : idx + 1] = length - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = 1 + self.seq_lens_this_time_buffer[idx : idx + 1] = 1 self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = length self.share_inputs["prompt_lens"][idx : idx + 1] = length @@ -377,7 +383,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): request.draft_token_ids[0:num_prefill_send_token], dtype="int64", ) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = num_prefill_send_token + self.seq_lens_this_time_buffer[idx : idx + 1] = num_prefill_send_token else: self.share_inputs["pre_ids"][idx : idx + 1] = -1 self.share_inputs["step_idx"][idx : idx + 1] = 0 @@ -412,7 +418,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): ) self.share_inputs["seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = token_chunk_size + self.seq_lens_this_time_buffer[idx : idx + 1] = token_chunk_size self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = token_chunk_size self.share_inputs["seq_lens_encoder"][idx : idx + 1] = token_chunk_size self.share_inputs["prompt_lens"][idx : idx + 1] = token_chunk_size @@ -430,7 +436,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): else: self.share_inputs["seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["prompt_lens"][idx : idx + 1] = length @@ -514,8 +520,10 @@ def get_attr_from_request(request, attr, default_value=None): self.share_inputs["not_need_stop"][0] = True + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] + if self.speculative_method in ["mtp"]: - self.proposer.insert_prefill_inputs(req_dicts) + self.proposer.insert_prefill_inputs(req_dicts, num_running_requests) def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decode_len: int): """Set dummy prefill inputs to share_inputs""" @@ -535,7 +543,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decod self.share_inputs["input_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["prompt_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["eos_token_id"][:] = np.array([2], dtype="int64").reshape(-1, 1) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = input_length + self.seq_lens_this_time_buffer[idx : idx + 1] = input_length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 @@ -553,6 +561,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decod self.share_inputs["block_tables"][idx : idx + 1, :block_num] = np.arange( idx * block_num, (idx + 1) * block_num, 1 ) + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer def _init_share_inputs(self, max_num_seqs: int): """ @@ -603,7 +612,7 @@ def _init_share_inputs(self, max_num_seqs: int): self.share_inputs["max_length"] = paddle.full( [max_num_seqs, 1], self.model_config.max_model_len, dtype="int64" ) - self.share_inputs["seq_lens_this_time"] = paddle.full(max_num_seqs, 0, dtype="int32") + self.seq_lens_this_time_buffer = paddle.full(max_num_seqs, 0, dtype="int32") self.share_inputs["seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["seq_lens_decoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["step_seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") @@ -1247,6 +1256,7 @@ def _get_skip_idx(self, model_forward_batch: Optional[List[Request]] = None): def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + num_running_requests: int = None, ) -> Optional[ModelRunnerOutput]: """ The Entrance of model execute. @@ -1255,6 +1265,7 @@ def execute_model( class at the server level, which is too granular for ModelRunner. We plan to replace it with 'ModelForwardBatch'. intermediate_tensors: + num_running_requests: batch_size """ # 1. Prepare inputs of model and sampler. skip_idx_list = self._get_skip_idx(model_forward_batch) @@ -1356,8 +1367,8 @@ class at the server level, which is too granular for ModelRunner. accept_num=(self.share_inputs["accept_num"] if self.speculative_decoding else None), enable_thinking=(self.share_inputs["enable_thinking"] if self.enable_mm else None), think_end_id=(self.model_config.think_end_id if self.enable_mm else -1), - need_think_end=(self.share_inputs["need_think_end"] if self.enable_mm else None), - reasoning_index=(self.share_inputs["reasoning_index"] if self.enable_mm else None), + need_think_end=(self.share_inputs["need_think_end"][:num_running_requests] if self.enable_mm else None), + reasoning_index=(self.share_inputs["reasoning_index"][:num_running_requests] if self.enable_mm else None), stop_token_ids=self.share_inputs["stop_seqs"], stop_seqs_len=self.share_inputs["stop_seqs_len"], ) @@ -1397,6 +1408,10 @@ class at the server level, which is too granular for ModelRunner. self._update_chunked_prefill(model_forward_batch) self._add_cache(model_forward_batch) + + self.seq_lens_this_time_buffer[:num_running_requests].copy_( + self.share_inputs["seq_lens_this_time"][:num_running_requests], False + ) return None def _add_cache(self, model_forward_batch) -> None: diff --git a/fastdeploy/worker/gpu_worker.py b/fastdeploy/worker/gpu_worker.py index 084b4f0f2d8..ad780e21ad1 100644 --- a/fastdeploy/worker/gpu_worker.py +++ b/fastdeploy/worker/gpu_worker.py @@ -175,20 +175,21 @@ def initialize_cache(self, num_gpu_blocks: int) -> None: def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + num_running_request: int = None, ) -> Optional[ModelRunnerOutput]: """ """ - output = self.model_runner.execute_model(model_forward_batch) + output = self.model_runner.execute_model(model_forward_batch, num_running_request) return output - def preprocess_new_task(self, req_dicts: List[Request]) -> None: + def preprocess_new_task(self, req_dicts: List[Request], num_running_requests: int) -> None: """Process new requests and then start the decode loop TODO(gongshaotian):The scheduler should schedule the handling of prefill, and workers and modelrunners should not perceive it. """ if envs.ENABLE_V1_KVCACHE_SCHEDULER: - self.model_runner.insert_tasks_v1(req_dicts=req_dicts) + self.model_runner.insert_tasks_v1(req_dicts=req_dicts, num_running_requests=num_running_requests) else: - self.model_runner.insert_prefill_inputs(req_dicts=req_dicts) + self.model_runner.insert_prefill_inputs(req_dicts=req_dicts, num_running_requests=num_running_requests) def graph_optimize_and_warm_up_model(self) -> None: """ diff --git a/fastdeploy/worker/iluvatar_model_runner.py b/fastdeploy/worker/iluvatar_model_runner.py index a84ab7118ae..526f3361e0e 100644 --- a/fastdeploy/worker/iluvatar_model_runner.py +++ b/fastdeploy/worker/iluvatar_model_runner.py @@ -142,9 +142,10 @@ def _init_logits_processor(self, request): schemata_key, ) - def insert_prefill_inputs(self, req_dicts: List[Request]): + def insert_prefill_inputs(self, req_dicts: List[Request], num_running_requests: int = None): """ Process inputs for prefill tasks and insert it to share_inputs buffer + num_running_requests: batch_size TODO(gongshaotian): Refactor this func """ @@ -176,7 +177,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): self.share_inputs["input_ids"][idx : idx + 1, 0] = request.prompt_token_ids[0] self.share_inputs["seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["seq_lens_decoder"][idx : idx + 1] = length - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = 1 + self.seq_lens_this_time_buffer[idx : idx + 1] = 1 self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = length self.share_inputs["prompt_lens"][idx : idx + 1] = length @@ -188,7 +189,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): request.draft_token_ids[0:num_prefill_send_token], dtype="int64", ) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = num_prefill_send_token + self.seq_lens_this_time_buffer[idx : idx + 1] = num_prefill_send_token else: self.share_inputs["pre_ids"][idx : idx + 1] = -1 self.share_inputs["step_idx"][idx : idx + 1] = 0 @@ -199,7 +200,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): request.set("chunk_idx", 1) logger.info(f"prefill_chunk_info: {request.prefill_chunk_info}") token_chunk_size = request.prefill_chunk_info[0] - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = token_chunk_size + self.seq_lens_this_time_buffer[idx : idx + 1] = token_chunk_size self.share_inputs["input_ids"][idx, :token_chunk_size] = np.array( request.prompt_token_ids[:token_chunk_size] ) @@ -211,7 +212,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): else: self.share_inputs["seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = request.get("seq_lens_decoder", 0) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["prompt_lens"][idx : idx + 1] = length @@ -262,6 +263,7 @@ def insert_prefill_inputs(self, req_dicts: List[Request]): self.sampler.apply_logits_processor(idx, request.get("logits_processor"), prefill_tokens) self.share_inputs["not_need_stop"][0] = True + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decode_len: int): """Set dummy prefill inputs to share_inputs""" @@ -281,7 +283,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decod self.share_inputs["input_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["prompt_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["eos_token_id"][:] = np.array([2], dtype="int64").reshape(-1, 1) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = input_length + self.seq_lens_this_time_buffer[idx : idx + 1] = input_length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 @@ -297,6 +299,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int, expected_decod self.share_inputs["block_tables"][idx : idx + 1, :block_num] = np.arange( idx * block_num, (idx + 1) * block_num, 1 ) + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer def _init_share_inputs(self, max_num_seqs: int): """Initialize all share buffers for model inputs. @@ -342,7 +345,7 @@ def _init_share_inputs(self, max_num_seqs: int): self.share_inputs["max_dec_len"] = paddle.full([max_num_seqs, 1], self.model_config.max_length, dtype="int64") self.share_inputs["min_length"] = paddle.full([max_num_seqs, 1], self.model_config.min_length, dtype="int64") self.share_inputs["max_length"] = paddle.full([max_num_seqs, 1], self.model_config.max_length, dtype="int64") - self.share_inputs["seq_lens_this_time"] = paddle.full(max_num_seqs, 0, dtype="int32") + self.seq_lens_this_time_buffer = paddle.full(max_num_seqs, 0, dtype="int32") self.share_inputs["seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["seq_lens_decoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["step_seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") @@ -859,6 +862,7 @@ def _get_skip_idx(self, model_forward_batch): def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + num_running_requests: int = None, ) -> Optional[ModelRunnerOutput]: """ The Entrance of model execute. @@ -866,6 +870,7 @@ def execute_model( model_forward_batch: 'Request' contains information related to prompt and is an abstract class at the server level, which is too granular for ModelRunner. We plan to replace it with 'ModelForwardBatch'. + num_running_requests: batch_size intermediate_tensors: """ # Note(@wufeisheng): If `not_need_stop`` is False, it means the current worker is in an idle state. @@ -986,6 +991,9 @@ class at the server level, which is too granular for ModelRunner. self._update_chunked_prefill(model_forward_batch) self._add_cache(model_forward_batch) + self.seq_lens_this_time_buffer[:num_running_requests].copy_( + self.share_inputs["seq_lens_this_time"][:num_running_requests], False + ) return None def _add_cache(self, model_forward_batch) -> None: diff --git a/fastdeploy/worker/iluvatar_worker.py b/fastdeploy/worker/iluvatar_worker.py index 6c390584f88..cd899619bbe 100644 --- a/fastdeploy/worker/iluvatar_worker.py +++ b/fastdeploy/worker/iluvatar_worker.py @@ -106,17 +106,18 @@ def initialize_cache(self, num_gpu_blocks: int) -> None: def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + num_running_requests: int = None, ) -> Optional[ModelRunnerOutput]: """ """ - output = self.model_runner.execute_model(model_forward_batch) + output = self.model_runner.execute_model(model_forward_batch, num_running_requests) return output - def preprocess_new_task(self, req_dicts: List[Request]) -> None: + def preprocess_new_task(self, req_dicts: List[Request], num_running_requests: int) -> None: """Process new requests and then start the decode loop TODO(gongshaotian):The scheduler should schedule the handling of prefill, and workers and modelrunners should not perceive it. """ - self.model_runner.insert_prefill_inputs(req_dicts=req_dicts) + self.model_runner.insert_prefill_inputs(req_dicts=req_dicts, num_running_requests=num_running_requests) def graph_optimize_and_warm_up_model(self) -> None: """ diff --git a/fastdeploy/worker/worker_process.py b/fastdeploy/worker/worker_process.py index 54f7019c871..ebe08669d65 100644 --- a/fastdeploy/worker/worker_process.py +++ b/fastdeploy/worker/worker_process.py @@ -257,11 +257,11 @@ def event_loop_ep(self) -> None: f"num_insert_requests: {len(req_dicts)}" ) # Process prefill inputs - self.worker.preprocess_new_task(req_dicts) + self.worker.preprocess_new_task(req_dicts, num_running_requests) # Execute model to generate token. The generated token will be written to the buffer. # These generated tokens can be obtained through get_output op. - self.worker.execute_model() + self.worker.execute_model(num_running_requests) def event_loop_normal(self) -> None: """Main event loop for Paddle Distrubuted Workers. @@ -338,7 +338,7 @@ def event_loop_normal(self) -> None: ) # Process prefill inputs - self.worker.preprocess_new_task(req_dicts) + self.worker.preprocess_new_task(req_dicts, num_running_requests) if not self.worker.model_runner.not_need_stop(): if self.ranks > 1: @@ -349,7 +349,7 @@ def event_loop_normal(self) -> None: # Execute model to generate token. The generated token will be written to the buffer. # These generated tokens can be obtained through get_output op. - self.worker.execute_model(req_dicts) + self.worker.execute_model(req_dicts, num_running_requests) self.exist_prefill_task_signal.value[0] = self.worker.exist_prefill() def initialize_kv_cache(self) -> None: diff --git a/fastdeploy/worker/xpu_model_runner.py b/fastdeploy/worker/xpu_model_runner.py index 53ca380207d..b61e84b9f0b 100644 --- a/fastdeploy/worker/xpu_model_runner.py +++ b/fastdeploy/worker/xpu_model_runner.py @@ -373,7 +373,7 @@ def __init__(self, fd_config: FDConfig, device: str, rank: int, local_rank: int) # Forward meta store the global meta information of the forward self.forward_meta: ForwardMeta = None - def insert_tasks_v1(self, req_dicts: List[Request]): + def insert_tasks_v1(self, req_dicts: List[Request], num_running_requests: int = None): """ Process scheduler output tasks, used when ENABLE_V1_KVCACHE_SCHEDULER=1 """ @@ -403,7 +403,7 @@ def insert_tasks_v1(self, req_dicts: List[Request]): ) self.share_inputs["stop_flags"][idx : idx + 1] = False self.share_inputs["seq_lens_decoder"][idx : idx + 1] = prefill_start_index - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["step_seq_lens_decoder"][idx : idx + 1] = 0 self.share_inputs["prompt_lens"][idx : idx + 1] = len(input_ids) @@ -425,7 +425,7 @@ def insert_tasks_v1(self, req_dicts: List[Request]): logger.debug(f"Handle preempted request {request} at idx {idx}") self.share_inputs["block_tables"][idx : idx + 1, :] = -1 self.share_inputs["stop_flags"][idx : idx + 1] = True - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = 0 + self.seq_lens_this_time_buffer[idx : idx + 1] = 0 self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 self.share_inputs["seq_lens_encoder"][idx : idx + 1] = 0 self.share_inputs["is_block_step"][idx : idx + 1] = False @@ -462,8 +462,9 @@ def insert_tasks_v1(self, req_dicts: List[Request]): ) if has_prefill_task: self.share_inputs["not_need_stop"][0] = True + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] - def process_prefill_inputs(self, req_dicts: List[Request]): + def process_prefill_inputs(self, req_dicts: List[Request], num_running_requests: int = None): """Process inputs for prefill tasks and update share_inputs buffer""" req_len = len(req_dicts) for i in range(req_len): @@ -482,7 +483,7 @@ def process_prefill_inputs(self, req_dicts: List[Request]): self.share_inputs["penalty_score"][idx : idx + 1] = request.get("repetition_penalty", 1.0) self.share_inputs["frequency_score"][idx : idx + 1] = request.get("frequency_penalty", 0.0) self.share_inputs["presence_score"][idx : idx + 1] = request.get("presence_penalty", 0.0) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = length + self.seq_lens_this_time_buffer[idx : idx + 1] = length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = length self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 @@ -524,6 +525,7 @@ def process_prefill_inputs(self, req_dicts: List[Request]): ) self.share_inputs["not_need_stop"][0] = True + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer[:num_running_requests] def _init_share_inputs(self, max_num_seqs: int): """Initialize all share buffers for model inputs. @@ -569,7 +571,7 @@ def _init_share_inputs(self, max_num_seqs: int): self.share_inputs["max_length"] = paddle.full( [max_num_seqs, 1], self.model_config.max_model_len, dtype="int64" ) - self.share_inputs["seq_lens_this_time"] = paddle.full(max_num_seqs, 0, dtype="int32") + self.seq_lens_this_time_buffer = paddle.full(max_num_seqs, 0, dtype="int32") self.share_inputs["seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["seq_lens_decoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") self.share_inputs["step_seq_lens_encoder"] = paddle.full([max_num_seqs, 1], 0, dtype="int32") @@ -811,7 +813,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int): idx = i self.share_inputs["input_ids"][idx : idx + 1, :input_length] = np.array([5] * input_length) self.share_inputs["eos_token_id"][:] = np.array([2], dtype="int64").reshape(-1, 1) - self.share_inputs["seq_lens_this_time"][idx : idx + 1] = input_length + self.seq_lens_this_time_buffer[idx : idx + 1] = input_length self.share_inputs["step_seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_encoder"][idx : idx + 1] = input_length self.share_inputs["seq_lens_decoder"][idx : idx + 1] = 0 @@ -827,6 +829,7 @@ def _dummy_prefill_inputs(self, num_tokens: int, batch_size: int): self.share_inputs["block_tables"][idx : idx + 1, :block_num] = np.arange( idx * block_num, (idx + 1) * block_num, 1 ) + self.share_inputs["seq_lens_this_time"] = self.seq_lens_this_time_buffer def _dummy_run( self, @@ -851,6 +854,7 @@ def execute_model( self, model_forward_batch: Optional[List[Request]] = None, is_dummy_run: bool = False, + num_running_requests: int = None, ) -> Optional[ModelRunnerOutput]: """ The Entrance of model execute. @@ -858,6 +862,7 @@ def execute_model( model_forward_batch: 'Request' contains information related to prompt and is an abstract class at the server level, which is too granular for ModelRunner. We plan to replace it with 'ModelForwardBatch'. + num_running_requests: batch_size intermediate_tensors: """ # 1. Prepare inputs of model and decoder. @@ -918,6 +923,10 @@ class at the server level, which is too granular for ModelRunner. self.cache_config.block_size, self.cache_config.enc_dec_block_num, ) + if num_running_requests is not None: + self.seq_lens_this_time_buffer[:num_running_requests].copy_( + self.share_inputs["seq_lens_this_time"][:num_running_requests], False + ) return None diff --git a/fastdeploy/worker/xpu_worker.py b/fastdeploy/worker/xpu_worker.py index 16b51d2e567..a5993abb458 100644 --- a/fastdeploy/worker/xpu_worker.py +++ b/fastdeploy/worker/xpu_worker.py @@ -145,9 +145,14 @@ def initialize_cache(self, num_gpu_blocks: int) -> None: def execute_model( self, model_forward_batch: Optional[List[Request]] = None, + is_dummy_run: bool = False, + num_running_requests: Optional[int] = None, ) -> Optional[ModelRunnerOutput]: """ """ - output = self.model_runner.execute_model(model_forward_batch) + if is_dummy_run: + output = self.model_runner.execute_model(model_forward_batch) + else: + output = self.model_runner.execute_model(model_forward_batch, num_running_requests) return output def exist_prefill(self): @@ -156,15 +161,15 @@ def exist_prefill(self): """ return self.model_runner.exist_prefill() - def preprocess_new_task(self, req_dicts: List[Request]) -> None: + def preprocess_new_task(self, req_dicts: List[Request], num_running_requests: int) -> None: """Process new requests and then start the decode loop TODO(gongshaotian):The scheduler should schedule the handling of prefill, and workers and modelrunners should not perceive it. """ if envs.ENABLE_V1_KVCACHE_SCHEDULER: - self.model_runner.insert_tasks_v1(req_dicts=req_dicts) + self.model_runner.insert_tasks_v1(req_dicts=req_dicts, num_running_requests=num_running_requests) else: - self.model_runner.process_prefill_inputs(req_dicts=req_dicts) + self.model_runner.process_prefill_inputs(req_dicts=req_dicts, num_running_requests=num_running_requests) def check_health(self) -> bool: """ """ From f672a34f95c6d4fab17f85c49e4ef3346fb41636 Mon Sep 17 00:00:00 2001 From: Sunny-bot1 <68891411+Sunny-bot1@users.noreply.github.com> Date: Wed, 6 Aug 2025 15:47:27 +0800 Subject: [PATCH 12/17] [FIX 2.1]fix bad_words when sending requests consecutively (#3199) * fix bad_words * fix log * fix log --- fastdeploy/engine/sampling_params.py | 24 ++++++++++++---------- fastdeploy/worker/gcu_model_runner.py | 14 +++++++------ fastdeploy/worker/gpu_model_runner.py | 14 +++++++------ fastdeploy/worker/iluvatar_model_runner.py | 14 +++++++------ fastdeploy/worker/xpu_model_runner.py | 14 +++++++------ 5 files changed, 45 insertions(+), 35 deletions(-) diff --git a/fastdeploy/engine/sampling_params.py b/fastdeploy/engine/sampling_params.py index 46d9fd8acfb..1cd77d2b160 100644 --- a/fastdeploy/engine/sampling_params.py +++ b/fastdeploy/engine/sampling_params.py @@ -218,20 +218,22 @@ def update_from_tokenizer(self, tokenizer): prompt_token_ids = tokenizer.encode(text=prompt, add_special_tokens=False)["input_ids"] if len(prompt_token_ids) != 1: - logger.warning( - f"Skip bad_words: {prompt}." - f"Bad words should be a single token." - f"Got tokens: {prompt_token_ids}." - ) + if not add_prefix_space: + logger.warning( + f"Skip bad_words: <{prompt}>." + f"Bad words should be a single token." + f"Got tokens: {prompt_token_ids}." + ) continue if prompt_token_ids[0] > tokenizer.vocab_size: - logger.warning( - f"Skip bad_words: {prompt}." - f"All token id values should be satisfying:" - f" 0 <= token_id < {tokenizer.vocab_size}." - f"Got token: {prompt_token_ids}." - ) + if not add_prefix_space: + logger.warning( + f"Skip bad_words: <{prompt}>." + f"All token id values should be satisfying:" + f" 0 <= token_id < {tokenizer.vocab_size}." + f"Got token: {prompt_token_ids}." + ) continue if prompt_token_ids not in self._bad_words_token_ids: diff --git a/fastdeploy/worker/gcu_model_runner.py b/fastdeploy/worker/gcu_model_runner.py index e35dc4d23ad..e0086b50378 100644 --- a/fastdeploy/worker/gcu_model_runner.py +++ b/fastdeploy/worker/gcu_model_runner.py @@ -272,13 +272,15 @@ def get_attr_from_request(request, attr, default_value=None): request.block_tables, dtype="int32" ) - if request.get("bad_words_token_ids") is not None: + if request.get("bad_words_token_ids") is not None and len(request.get("bad_words_token_ids")) > 0: bad_words_len = len(request.get("bad_words_token_ids")) - if bad_words_len > 0: - self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len - self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( - request.get("bad_words_token_ids"), dtype="int64" - ) + self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len + self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( + request.get("bad_words_token_ids"), dtype="int64" + ) + else: + self.share_inputs["bad_tokens_len"][idx : idx + 1] = 1 + self.share_inputs["bad_tokens"][idx : idx + 1, :] = np.array([-1], dtype="int64") if request.get("stop_token_ids") is not None and request.get("stop_seqs_len") is not None: stop_seqs_num = len(request.get("stop_seqs_len")) diff --git a/fastdeploy/worker/gpu_model_runner.py b/fastdeploy/worker/gpu_model_runner.py index c551364ef2a..447f249d491 100644 --- a/fastdeploy/worker/gpu_model_runner.py +++ b/fastdeploy/worker/gpu_model_runner.py @@ -495,13 +495,15 @@ def get_attr_from_request(request, attr, default_value=None): request.block_tables, dtype="int32" ) - if request.get("bad_words_token_ids") is not None: + if request.get("bad_words_token_ids") is not None and len(request.get("bad_words_token_ids")) > 0: bad_words_len = len(request.get("bad_words_token_ids")) - if bad_words_len > 0: - self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len - self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( - request.get("bad_words_token_ids"), dtype="int64" - ) + self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len + self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( + request.get("bad_words_token_ids"), dtype="int64" + ) + else: + self.share_inputs["bad_tokens_len"][idx : idx + 1] = 1 + self.share_inputs["bad_tokens"][idx : idx + 1, :] = np.array([-1], dtype="int64") if request.get("stop_token_ids") is not None and request.get("stop_seqs_len") is not None: stop_seqs_num = len(request.get("stop_seqs_len")) diff --git a/fastdeploy/worker/iluvatar_model_runner.py b/fastdeploy/worker/iluvatar_model_runner.py index 526f3361e0e..4a7aaaf8d0d 100644 --- a/fastdeploy/worker/iluvatar_model_runner.py +++ b/fastdeploy/worker/iluvatar_model_runner.py @@ -243,13 +243,15 @@ def insert_prefill_inputs(self, req_dicts: List[Request], num_running_requests: request.block_tables, dtype="int32" ) - if request.get("bad_words_token_ids") is not None: + if request.get("bad_words_token_ids") is not None and len(request.get("bad_words_token_ids")) > 0: bad_words_len = len(request.get("bad_words_token_ids")) - if bad_words_len > 0: - self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len - self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( - request.get("bad_words_token_ids"), dtype="int64" - ) + self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len + self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( + request.get("bad_words_token_ids"), dtype="int64" + ) + else: + self.share_inputs["bad_tokens_len"][idx : idx + 1] = 1 + self.share_inputs["bad_tokens"][idx : idx + 1, :] = np.array([-1], dtype="int64") if request.get("stop_token_ids") is not None and request.get("stop_seqs_len") is not None: stop_seqs_num = len(request.get("stop_seqs_len")) diff --git a/fastdeploy/worker/xpu_model_runner.py b/fastdeploy/worker/xpu_model_runner.py index b61e84b9f0b..a153e555659 100644 --- a/fastdeploy/worker/xpu_model_runner.py +++ b/fastdeploy/worker/xpu_model_runner.py @@ -507,13 +507,15 @@ def process_prefill_inputs(self, req_dicts: List[Request], num_running_requests: request.block_tables, dtype="int32" ) - if request.get("bad_words_token_ids") is not None: + if request.get("bad_words_token_ids") is not None and len(request.get("bad_words_token_ids")) > 0: bad_words_len = len(request.get("bad_words_token_ids")) - if bad_words_len > 0: - self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len - self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( - request.get("bad_words_token_ids"), dtype="int64" - ) + self.share_inputs["bad_tokens_len"][idx : idx + 1] = bad_words_len + self.share_inputs["bad_tokens"][idx : idx + 1, :bad_words_len] = np.array( + request.get("bad_words_token_ids"), dtype="int64" + ) + else: + self.share_inputs["bad_tokens_len"][idx : idx + 1] = 1 + self.share_inputs["bad_tokens"][idx : idx + 1, :] = np.array([-1], dtype="int64") if request.get("stop_token_ids") is not None and request.get("stop_seqs_len") is not None: stop_seqs_num = len(request.get("stop_seqs_len")) From 5d3bf308f6e5546a55ef3945a4ffc5a57eb6ee1a Mon Sep 17 00:00:00 2001 From: sg263 <1021937542@qq.com> Date: Thu, 7 Aug 2025 11:10:55 +0800 Subject: [PATCH 13/17] merge develop trace FD_START (#3253) Co-authored-by: shige --- fastdeploy/entrypoints/openai/api_server.py | 3 ++- fastdeploy/envs.py | 2 ++ fastdeploy/metrics/trace_util.py | 21 +++++++++++++++++++-- 3 files changed, 23 insertions(+), 3 deletions(-) diff --git a/fastdeploy/entrypoints/openai/api_server.py b/fastdeploy/entrypoints/openai/api_server.py index 5161d2d2586..714425e2fe8 100644 --- a/fastdeploy/entrypoints/openai/api_server.py +++ b/fastdeploy/entrypoints/openai/api_server.py @@ -45,7 +45,7 @@ get_filtered_metrics, main_process_metrics, ) -from fastdeploy.metrics.trace_util import inject_to_metadata, instrument +from fastdeploy.metrics.trace_util import fd_start_span, inject_to_metadata, instrument from fastdeploy.utils import ( FlexibleArgumentParser, api_server_logger, @@ -270,6 +270,7 @@ def launch_api_server() -> None: api_server_logger.info(f"launch Fastdeploy api server... port: {args.port}") api_server_logger.info(f"args: {args.__dict__}") + fd_start_span("FD_START") try: uvicorn.run( diff --git a/fastdeploy/envs.py b/fastdeploy/envs.py index 901ef3f5a0a..9c857d3d6d0 100644 --- a/fastdeploy/envs.py +++ b/fastdeploy/envs.py @@ -80,6 +80,8 @@ "EXPORTER_OTLP_HEADERS": lambda: os.getenv("EXPORTER_OTLP_HEADERS"), # enable kv cache block scheduler v1 (no need for kv_cache_ratio) "ENABLE_V1_KVCACHE_SCHEDULER": lambda: int(os.getenv("ENABLE_V1_KVCACHE_SCHEDULER", "0")), + # set trace attribute job_id. + "FD_JOB_ID": lambda: os.getenv("FD_JOB_ID"), } diff --git a/fastdeploy/metrics/trace_util.py b/fastdeploy/metrics/trace_util.py index e51446e7754..8b391dd6656 100644 --- a/fastdeploy/metrics/trace_util.py +++ b/fastdeploy/metrics/trace_util.py @@ -1,4 +1,5 @@ import json +import os from fastapi import FastAPI from opentelemetry import trace @@ -176,7 +177,22 @@ def start_span(span_name, request, kind=trace.SpanKind.CLIENT): return # extract Trace context from request.metadata.trace_carrier ctx = extract_from_metadata(request) - with tracer.start_as_current_span(span_name, context=ctx, kind=kind): + with tracer.start_as_current_span(span_name, context=ctx, kind=kind) as span: + span.set_attribute("job_id", os.getenv("FD_JOB_ID", default="null")) + pass + except: + pass + + +def fd_start_span(span_name, kind=trace.SpanKind.CLIENT): + """ + when fd start, start a new span show start success + """ + try: + if not traces_enable: + return + with tracer.start_as_current_span(span_name, kind=kind) as span: + span.set_attribute("job_id", os.getenv("FD_JOB_ID", default="null")) pass except: pass @@ -191,7 +207,8 @@ def start_span_request(span_name, request, kind=trace.SpanKind.CLIENT): return # extract Trace context from request.metadata.trace_carrier ctx = extract_from_request(request) - with tracer.start_as_current_span(span_name, context=ctx, kind=kind): + with tracer.start_as_current_span(span_name, context=ctx, kind=kind) as span: + span.set_attribute("job_id", os.getenv("FD_JOB_ID", default="null")) pass except: pass From 1b6f482c151871634517f60e3e6de55ca49fc6d1 Mon Sep 17 00:00:00 2001 From: JYChen Date: Thu, 7 Aug 2025 19:11:37 +0800 Subject: [PATCH 14/17] [Cherry-pick] fix stop seq (#3263) * fix out-bound value for stop sequence * catch error if there are out-of-bounds value * check in offline mode --- fastdeploy/engine/engine.py | 20 ++++++++++++++++++++ fastdeploy/entrypoints/engine_client.py | 21 +++++++++++++++++++++ 2 files changed, 41 insertions(+) diff --git a/fastdeploy/engine/engine.py b/fastdeploy/engine/engine.py index c76a6432dd1..88067ed068a 100644 --- a/fastdeploy/engine/engine.py +++ b/fastdeploy/engine/engine.py @@ -530,6 +530,26 @@ def add_requests(self, task, sampling_params=None, **kwargs): llm_logger.error(error_msg) raise EngineError(error_msg, error_code=400) + if request.get("stop_seqs_len") is not None: + stop_seqs_len = request.get("stop_seqs_len") + max_stop_seqs_num = int(envs.FD_MAX_STOP_SEQS_NUM) + if len(stop_seqs_len) > max_stop_seqs_num: + error_msg = ( + f"Length of stop ({stop_seqs_len}) exceeds the limit max_model_len({max_stop_seqs_num})." + "Please reduce the number of stop or set a lager max_stop_seqs_num by `FD_MAX_STOP_SEQS_NUM`" + ) + llm_logger.error(error_msg) + raise EngineError(error_msg, error_code=400) + stop_seqs_max_len = int(envs.FD_STOP_SEQS_MAX_LEN) + for single_stop_seq_len in stop_seqs_len: + if single_stop_seq_len > stop_seqs_max_len: + error_msg = ( + f"Length of stop_seqs({single_stop_seq_len}) exceeds the limit max_model_len({stop_seqs_max_len})." + "Please reduce the length of stop sequences or set a larger stop_seqs_max_len by `FD_STOP_SEQS_MAX_LEN`" + ) + llm_logger.error(error_msg) + raise EngineError(error_msg, error_code=400) + if self.guided_decoding_checker is not None: request, err_msg = self.guided_decoding_checker.schema_format(request) if err_msg is not None: diff --git a/fastdeploy/entrypoints/engine_client.py b/fastdeploy/entrypoints/engine_client.py index 9be9eccb4a9..230cc319640 100644 --- a/fastdeploy/entrypoints/engine_client.py +++ b/fastdeploy/entrypoints/engine_client.py @@ -19,6 +19,7 @@ import numpy as np +from fastdeploy import envs from fastdeploy.input.preprocess import InputPreprocessor from fastdeploy.inter_communicator import IPCSignal, ZmqClient from fastdeploy.metrics.work_metrics import work_process_metrics @@ -144,6 +145,26 @@ def add_requests(self, task): api_server_logger.error(error_msg) raise EngineError(error_msg, error_code=400) + if "stop_seqs_len" in task: + stop_seqs_len = task["stop_seqs_len"] + max_stop_seqs_num = int(envs.FD_MAX_STOP_SEQS_NUM) + if len(stop_seqs_len) > max_stop_seqs_num: + error_msg = ( + f"Length of stop ({stop_seqs_len}) exceeds the limit max_model_len({max_stop_seqs_num})." + "Please reduce the number of stop or set a lager max_stop_seqs_num by `FD_MAX_STOP_SEQS_NUM`" + ) + api_server_logger.error(error_msg) + raise EngineError(error_msg, error_code=400) + stop_seqs_max_len = int(envs.FD_STOP_SEQS_MAX_LEN) + for single_stop_seq_len in stop_seqs_len: + if single_stop_seq_len > stop_seqs_max_len: + error_msg = ( + f"Length of stop_seqs({single_stop_seq_len}) exceeds the limit max_model_len({stop_seqs_max_len})." + "Please reduce the length of stop sequences or set a larger stop_seqs_max_len by `FD_STOP_SEQS_MAX_LEN`" + ) + api_server_logger.error(error_msg) + raise EngineError(error_msg, error_code=400) + task["preprocess_end_time"] = time.time() preprocess_cost_time = task["preprocess_end_time"] - task["preprocess_start_time"] api_server_logger.info( From 6706ccb37e9d654767bdf74df852de30c6abdb90 Mon Sep 17 00:00:00 2001 From: ltd0924 <32387785+ltd0924@users.noreply.github.com> Date: Fri, 8 Aug 2025 20:11:32 +0800 Subject: [PATCH 15/17] [BugFix] fix too many open files problem (#3275) --- fastdeploy/entrypoints/engine_client.py | 4 +- fastdeploy/entrypoints/openai/api_server.py | 100 ++++++++++++++---- fastdeploy/entrypoints/openai/serving_chat.py | 16 ++- .../entrypoints/openai/serving_completion.py | 13 ++- fastdeploy/envs.py | 2 + fastdeploy/utils.py | 66 ++++++++++++ 6 files changed, 178 insertions(+), 23 deletions(-) diff --git a/fastdeploy/entrypoints/engine_client.py b/fastdeploy/entrypoints/engine_client.py index 230cc319640..7c92390bc37 100644 --- a/fastdeploy/entrypoints/engine_client.py +++ b/fastdeploy/entrypoints/engine_client.py @@ -24,7 +24,7 @@ from fastdeploy.inter_communicator import IPCSignal, ZmqClient from fastdeploy.metrics.work_metrics import work_process_metrics from fastdeploy.platforms import current_platform -from fastdeploy.utils import EngineError, api_server_logger +from fastdeploy.utils import EngineError, StatefulSemaphore, api_server_logger class EngineClient: @@ -44,6 +44,7 @@ def __init__( reasoning_parser=None, data_parallel_size=1, enable_logprob=False, + workers=1, ): input_processor = InputPreprocessor( tokenizer, @@ -76,6 +77,7 @@ def __init__( suffix=pid, create=False, ) + self.semaphore = StatefulSemaphore((envs.FD_SUPPORT_MAX_CONNECTIONS + workers - 1) // workers) def create_zmq_client(self, model, mode): """ diff --git a/fastdeploy/entrypoints/openai/api_server.py b/fastdeploy/entrypoints/openai/api_server.py index 714425e2fe8..6562bfac3c7 100644 --- a/fastdeploy/entrypoints/openai/api_server.py +++ b/fastdeploy/entrypoints/openai/api_server.py @@ -14,15 +14,17 @@ # limitations under the License. """ +import asyncio import os import threading import time +from collections.abc import AsyncGenerator from contextlib import asynccontextmanager from multiprocessing import current_process import uvicorn import zmq -from fastapi import FastAPI, Request +from fastapi import FastAPI, HTTPException, Request from fastapi.responses import JSONResponse, Response, StreamingResponse from prometheus_client import CONTENT_TYPE_LATEST @@ -48,6 +50,7 @@ from fastdeploy.metrics.trace_util import fd_start_span, inject_to_metadata, instrument from fastdeploy.utils import ( FlexibleArgumentParser, + StatefulSemaphore, api_server_logger, console_logger, is_port_available, @@ -60,6 +63,13 @@ parser.add_argument("--workers", default=1, type=int, help="number of workers") parser.add_argument("--metrics-port", default=8001, type=int, help="port for metrics server") parser.add_argument("--controller-port", default=-1, type=int, help="port for controller server") +parser.add_argument( + "--max-waiting-time", + default=-1, + type=int, + help="max waiting time for connection, if set value -1 means no waiting time limit", +) +parser.add_argument("--max-concurrency", default=512, type=int, help="max concurrency") parser = EngineArgs.add_cli_args(parser) args = parser.parse_args() args.model = retrive_model_from_server(args.model, args.revision) @@ -115,10 +125,11 @@ async def lifespan(app: FastAPI): args.reasoning_parser, args.data_parallel_size, args.enable_logprob, + args.workers, ) app.state.dynamic_load_weight = args.dynamic_load_weight - chat_handler = OpenAIServingChat(engine_client, pid, args.ips) - completion_handler = OpenAIServingCompletion(engine_client, pid, args.ips) + chat_handler = OpenAIServingChat(engine_client, pid, args.ips, args.max_waiting_time) + completion_handler = OpenAIServingCompletion(engine_client, pid, args.ips, args.max_waiting_time) engine_client.create_zmq_client(model=pid, mode=zmq.PUSH) engine_client.pid = pid app.state.engine_client = engine_client @@ -140,6 +151,41 @@ async def lifespan(app: FastAPI): instrument(app) +MAX_CONCURRENT_CONNECTIONS = (args.max_concurrency + args.workers - 1) // args.workers +connection_semaphore = StatefulSemaphore(MAX_CONCURRENT_CONNECTIONS) + + +@asynccontextmanager +async def connection_manager(): + """ + async context manager for connection manager + """ + try: + await asyncio.wait_for(connection_semaphore.acquire(), timeout=0.001) + yield + except asyncio.TimeoutError: + api_server_logger.info(f"Reach max request release: {connection_semaphore.status()}") + if connection_semaphore.locked(): + connection_semaphore.release() + raise HTTPException(status_code=429, detail="Too many requests") + + +def wrap_streaming_generator(original_generator: AsyncGenerator): + """ + Wrap an async generator to release the connection semaphore when the generator is finished. + """ + + async def wrapped_generator(): + try: + async for chunk in original_generator: + yield chunk + finally: + api_server_logger.debug(f"release: {connection_semaphore.status()}") + connection_semaphore.release() + + return wrapped_generator + + # TODO 传递真实引擎值 通过pid 获取状态 @app.get("/health") def health(request: Request) -> Response: @@ -202,16 +248,23 @@ async def create_chat_completion(request: ChatCompletionRequest): status, msg = app.state.engine_client.is_workers_alive() if not status: return JSONResponse(content={"error": "Worker Service Not Healthy"}, status_code=304) - inject_to_metadata(request) - generator = await app.state.chat_handler.create_chat_completion(request) - - if isinstance(generator, ErrorResponse): - return JSONResponse(content=generator.model_dump(), status_code=generator.code) - - elif isinstance(generator, ChatCompletionResponse): - return JSONResponse(content=generator.model_dump()) - - return StreamingResponse(content=generator, media_type="text/event-stream") + try: + async with connection_manager(): + inject_to_metadata(request) + generator = await app.state.chat_handler.create_chat_completion(request) + if isinstance(generator, ErrorResponse): + connection_semaphore.release() + return JSONResponse(content={"detail": generator.model_dump()}, status_code=generator.code) + elif isinstance(generator, ChatCompletionResponse): + connection_semaphore.release() + return JSONResponse(content=generator.model_dump()) + else: + wrapped_generator = wrap_streaming_generator(generator) + return StreamingResponse(content=wrapped_generator(), media_type="text/event-stream") + + except HTTPException as e: + api_server_logger.error(f"Error in chat completion: {str(e)}") + return JSONResponse(status_code=e.status_code, content={"detail": e.detail}) @app.post("/v1/completions") @@ -224,13 +277,20 @@ async def create_completion(request: CompletionRequest): if not status: return JSONResponse(content={"error": "Worker Service Not Healthy"}, status_code=304) - generator = await app.state.completion_handler.create_completion(request) - if isinstance(generator, ErrorResponse): - return JSONResponse(content=generator.model_dump(), status_code=generator.code) - elif isinstance(generator, CompletionResponse): - return JSONResponse(content=generator.model_dump()) - - return StreamingResponse(content=generator, media_type="text/event-stream") + try: + async with connection_manager(): + generator = await app.state.completion_handler.create_completion(request) + if isinstance(generator, ErrorResponse): + connection_semaphore.release() + return JSONResponse(content=generator.model_dump(), status_code=generator.code) + elif isinstance(generator, CompletionResponse): + connection_semaphore.release() + return JSONResponse(content=generator.model_dump()) + else: + wrapped_generator = wrap_streaming_generator(generator) + return StreamingResponse(content=wrapped_generator(), media_type="text/event-stream") + except HTTPException as e: + return JSONResponse(status_code=e.status_code, content={"detail": e.detail}) @app.get("/update_model_weight") diff --git a/fastdeploy/entrypoints/openai/serving_chat.py b/fastdeploy/entrypoints/openai/serving_chat.py index 3e74c89df37..5b4abbab2ab 100644 --- a/fastdeploy/entrypoints/openai/serving_chat.py +++ b/fastdeploy/entrypoints/openai/serving_chat.py @@ -49,10 +49,11 @@ class OpenAIServingChat: OpenAI-style chat completions serving """ - def __init__(self, engine_client, pid, ips): + def __init__(self, engine_client, pid, ips, max_waiting_time): self.engine_client = engine_client self.pid = pid self.master_ip = ips + self.max_waiting_time = max_waiting_time self.host_ip = get_host_ip() if self.master_ip is not None: if isinstance(self.master_ip, list): @@ -94,6 +95,15 @@ async def create_chat_completion(self, request: ChatCompletionRequest): del current_req_dict + try: + api_server_logger.debug(f"{self.engine_client.semaphore.status()}") + if self.max_waiting_time < 0: + await self.engine_client.semaphore.acquire() + else: + await asyncio.wait_for(self.engine_client.semaphore.acquire(), timeout=self.max_waiting_time) + except Exception: + return ErrorResponse(code=408, message=f"Request queued time exceed {self.max_waiting_time}") + if request.stream: return self.chat_completion_stream_generator(request, request_id, request.model, prompt_token_ids) else: @@ -310,6 +320,8 @@ async def chat_completion_stream_generator( yield f"data: {error_data}\n\n" finally: dealer.close() + self.engine_client.semaphore.release() + api_server_logger.info(f"release {self.engine_client.semaphore.status()}") yield "data: [DONE]\n\n" async def chat_completion_full_generator( @@ -384,6 +396,8 @@ async def chat_completion_full_generator( break finally: dealer.close() + self.engine_client.semaphore.release() + api_server_logger.info(f"release {self.engine_client.semaphore.status()}") choices = [] output = final_res["outputs"] diff --git a/fastdeploy/entrypoints/openai/serving_completion.py b/fastdeploy/entrypoints/openai/serving_completion.py index 87b6444df94..e5f50773a83 100644 --- a/fastdeploy/entrypoints/openai/serving_completion.py +++ b/fastdeploy/entrypoints/openai/serving_completion.py @@ -40,11 +40,12 @@ class OpenAIServingCompletion: - def __init__(self, engine_client, pid, ips): + def __init__(self, engine_client, pid, ips, max_waiting_time): self.engine_client = engine_client self.pid = pid self.master_ip = ips self.host_ip = get_host_ip() + self.max_waiting_time = max_waiting_time if self.master_ip is not None: if isinstance(self.master_ip, list): self.master_ip = self.master_ip[0] @@ -114,6 +115,14 @@ async def create_completion(self, request: CompletionRequest): del current_req_dict + try: + if self.max_waiting_time < 0: + await self.engine_client.semaphore.acquire() + else: + await asyncio.wait_for(self.engine_client.semaphore.acquire(), timeout=self.max_waiting_time) + except Exception: + return ErrorResponse(code=408, message=f"Request queued time exceed {self.max_waiting_time}") + if request.stream: return self.completion_stream_generator( request=request, @@ -223,6 +232,7 @@ async def completion_full_generator( finally: if dealer is not None: dealer.close() + self.engine_client.semaphore.release() async def completion_stream_generator( self, @@ -372,6 +382,7 @@ async def completion_stream_generator( del request if dealer is not None: dealer.close() + self.engine_client.semaphore.release() yield "data: [DONE]\n\n" def request_output_to_completion_response( diff --git a/fastdeploy/envs.py b/fastdeploy/envs.py index 9c857d3d6d0..f5aa5dc7e95 100644 --- a/fastdeploy/envs.py +++ b/fastdeploy/envs.py @@ -82,6 +82,8 @@ "ENABLE_V1_KVCACHE_SCHEDULER": lambda: int(os.getenv("ENABLE_V1_KVCACHE_SCHEDULER", "0")), # set trace attribute job_id. "FD_JOB_ID": lambda: os.getenv("FD_JOB_ID"), + # support max connections + "FD_SUPPORT_MAX_CONNECTIONS": lambda: 768, } diff --git a/fastdeploy/utils.py b/fastdeploy/utils.py index 9ea25000c73..5d68c7681e2 100644 --- a/fastdeploy/utils.py +++ b/fastdeploy/utils.py @@ -15,6 +15,7 @@ """ import argparse +import asyncio import codecs import importlib import logging @@ -291,6 +292,16 @@ def extract_tar(tar_path, output_dir): raise RuntimeError(f"Extraction failed: {e!s}") +def get_limited_max_value(max_value): + def validator(value): + value = float(value) + if value > max_value: + raise argparse.ArgumentTypeError(f"The value cannot exceed {max_value}") + return value + + return validator + + def download_model(url, output_dir, temp_tar): """ 下载模型,并将其解压到指定目录。 @@ -596,6 +607,61 @@ def version(): return content +class StatefulSemaphore: + __slots__ = ("_semaphore", "_max_value", "_acquired_count", "_last_reset") + + """ + StatefulSemaphore is a class that wraps an asyncio.Semaphore and provides additional stateful information. + """ + + def __init__(self, value: int): + """ + StatefulSemaphore constructor + """ + if value < 0: + raise ValueError("Value must be non-negative.") + self._semaphore = asyncio.Semaphore(value) + self._max_value = value + self._acquired_count = 0 + self._last_reset = time.monotonic() + + async def acquire(self): + await self._semaphore.acquire() + self._acquired_count += 1 + + def release(self): + self._semaphore.release() + + self._acquired_count = max(0, self._acquired_count - 1) + + def locked(self) -> bool: + return self._semaphore.locked() + + @property + def available(self) -> int: + return self._max_value - self._acquired_count + + @property + def acquired(self) -> int: + return self._acquired_count + + @property + def max_value(self) -> int: + return self._max_value + + @property + def uptime(self) -> float: + return time.monotonic() - self._last_reset + + def status(self) -> dict: + return { + "available": self.available, + "acquired": self.acquired, + "max_value": self.max_value, + "uptime": round(self.uptime, 2), + } + + llm_logger = get_logger("fastdeploy", "fastdeploy.log") data_processor_logger = get_logger("data_processor", "data_processor.log") scheduler_logger = get_logger("scheduler", "scheduler.log") From ac04c5bf4f847d706919b9554b234294834cf927 Mon Sep 17 00:00:00 2001 From: lizexu <2694294196@qq.com> Date: Sat, 9 Aug 2025 11:37:34 +0800 Subject: [PATCH 16/17] fix ep --- fastdeploy/worker/gpu_model_runner.py | 5 +++-- fastdeploy/worker/worker_process.py | 3 ++- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/fastdeploy/worker/gpu_model_runner.py b/fastdeploy/worker/gpu_model_runner.py index 447f249d491..bdd164f17ca 100644 --- a/fastdeploy/worker/gpu_model_runner.py +++ b/fastdeploy/worker/gpu_model_runner.py @@ -661,6 +661,7 @@ def _init_share_inputs(self, max_num_seqs: int): # Initialize rotary position embedding tmp_position_ids = paddle.arange(self.parallel_config.max_model_len).reshape((1, -1)) + # self.share_inputs["seq_lens_this_time"]=0 # TODO(gongshaotian): move to models if not self.enable_mm: @@ -777,7 +778,7 @@ def _prepare_inputs(self) -> None: output_padding_offset, ) = pre_process( self.share_inputs["input_ids"], - self.share_inputs["seq_lens_this_time"], + getattr(self.share_inputs, "seq_lens_this_time", self.seq_lens_this_time_buffer), self.speculative_decoding, (self.share_inputs["draft_tokens"] if self.speculative_decoding else None), self.share_inputs["seq_lens_encoder"], @@ -861,7 +862,7 @@ def initialize_forward_meta(self): max_len_tensor_cpu=self.share_inputs["max_len_tensor_cpu"], seq_lens_encoder=self.share_inputs["seq_lens_encoder"], seq_lens_decoder=self.share_inputs["seq_lens_decoder"], - seq_lens_this_time=self.share_inputs["seq_lens_this_time"], + seq_lens_this_time=getattr(self.share_inputs, "seq_lens_this_time", self.seq_lens_this_time_buffer), batch_id_per_token=self.share_inputs["batch_id_per_token"], cu_seqlens_q=self.share_inputs["cu_seqlens_q"], cu_seqlens_k=self.share_inputs["cu_seqlens_k"], diff --git a/fastdeploy/worker/worker_process.py b/fastdeploy/worker/worker_process.py index ebe08669d65..bd650923cbf 100644 --- a/fastdeploy/worker/worker_process.py +++ b/fastdeploy/worker/worker_process.py @@ -244,7 +244,7 @@ def event_loop_ep(self) -> None: """ while True: self.worker_healthy_live_signal.value[self.local_rank % self.max_chips_per_node] = int(time.time()) - + num_running_requests = 0 if self.fd_config.parallel_config.tensor_parallel_rank == 0 and self.task_queue.num_tasks() > 0: tasks, read_finish = self.task_queue.get_tasks() @@ -271,6 +271,7 @@ def event_loop_normal(self) -> None: self.nnode = int((self.parallel_config.tensor_parallel_size + 7) // 8) mp_num_per_node = self.parallel_config.tensor_parallel_size // self.nnode req_ids = [] + num_running_requests = 0 while True: if self.local_rank == 0: if self.model_weights_status.value[0] != 0: From 2530d45cbbf09fd1ce6475d860ecfa80f019fee1 Mon Sep 17 00:00:00 2001 From: lizexu <2694294196@qq.com> Date: Sat, 9 Aug 2025 15:00:18 +0800 Subject: [PATCH 17/17] fix --- fastdeploy/worker/gpu_model_runner.py | 1 - 1 file changed, 1 deletion(-) diff --git a/fastdeploy/worker/gpu_model_runner.py b/fastdeploy/worker/gpu_model_runner.py index bdd164f17ca..13dc0bf0fdb 100644 --- a/fastdeploy/worker/gpu_model_runner.py +++ b/fastdeploy/worker/gpu_model_runner.py @@ -661,7 +661,6 @@ def _init_share_inputs(self, max_num_seqs: int): # Initialize rotary position embedding tmp_position_ids = paddle.arange(self.parallel_config.max_model_len).reshape((1, -1)) - # self.share_inputs["seq_lens_this_time"]=0 # TODO(gongshaotian): move to models if not self.enable_mm: