Conversation
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request adds comprehensive support for DeepSeek V4.1 model serving on Ascend hardware. It implements a full suite of framework and model-specific integrations, ranging from preprocessing and speculative decoding to advanced KV-cache and MoE scheduling. The changes also include significant kernel-level optimizations and new benchmarking tools to ensure high-performance execution. Highlights
New Features🧠 You can now enable Memory (public preview) to help Gemini Code Assist learn from your team's feedback. This makes future code reviews more consistent and personalized to your project's style. Click here to enable Memory in your admin console. Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize the Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counterproductive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for GitHub and other Google products, sign up here. Footnotes
|
bca8654 to
a633d9d
Compare
|
👋 Hi! Thank you for contributing to the vLLM Ascend project. The following points will speed up your PR merge:
If CI fails, you can run linting and testing checks locally according Contributing and Testing. Tip 💡 Consider Linking a Related Issue or RFCYour PR title contains the [Feature] tag, indicating a bug fix or new feature. Linking a related issue or RFC in the PR description is strongly encouraged — it gives reviewers helpful context and speeds up the review. You can use any of these keywords:
🙏 Thanks for helping us keep the project well-organized! |
There was a problem hiding this comment.
Code Review
Suggested PR Title:
[Attention][Feature] Introduce Two-Level TopK Candidate Selection for QuantLightningIndexerV2 on Arch22Suggested PR Summary:
### What this PR does / why we need it?
This pull request introduces the two-level TopK (candidate block selection) feature for the `QuantLightningIndexerV2` operator, specifically targeting the `arch22` (910b) platform. It adds the necessary host tiling, kernel, and service vector implementations to support candidate block generation (source mode) and candidate block filtering (consumer mode). Additionally, it introduces several helper headers (such as `init_output.h`, `attn_buffer.h`, and `attn_buffer_manager.h`), benchmarks, comprehensive documentation, and a robust pytest-based testing suite.
Feedback on the implementation:
- **Memory Alignment**: In `init_output.h`, `singleCoreMaxElementNum` must be aligned to 32-byte boundaries to prevent runtime alignment faults during `DataCopy`/`DataCopyPad` on cores > 0.
- **Unpacking Error**: In `quant_lightning_indexer.py`, the 3-tuple returned by `quant_lightning_indexer_candidate` is incorrectly unpacked into a 2-tuple in `quant_lightning_indexer_candidate_consumer`, which will cause a runtime `ValueError`.
- **Optional Outputs**: In `quant_lightning_indexer.cpp`, unused optional outputs should be passed as undefined `at::Tensor()` instead of empty `{0}` tensors to avoid ACLNN validation failures.
- **GE InferShape**: In `quant_lightning_indexer_v2_infershape.cpp`, checking for null on optional output shapes breaks the optional contract and can fail graph compilation when the output is unconnected.
- **CMake Cache**: In `CMakeLists.txt`, `FORCE` should be used when setting `quant_lightning_indexer_v2_depends` to ensure cached values are overwritten during incremental builds.
- **Device Mismatch**: In `quant_lightning_indexer.py`, ensure `seqused_k` and `cmp_residual_k` are explicitly moved to the correct device to avoid runtime device mismatches.
- **Strict Aliasing**: In `vf_flash_decode_arch35.h`, avoid reference casting `(__ubuf__ float *&)` to prevent strict aliasing violations.
### Does this PR introduce _any_ user-facing change?
Yes, it introduces the `quant_lightning_indexer_candidate` PyTorch operator and its associated Python bindings (`quant_lightning_indexer_candidate_source` and `quant_lightning_indexer_candidate_consumer`) for two-level TopK candidate selection.
### How was this patch tested?
The patch was tested using the newly added pytest suite under `csrc/attention/quant_lightning_indexer_v2/tests/pytest/`, covering various scenarios including BSND/TND layouts, different quantization modes, and candidate modes.| uint64_t singleCoreMaxElementNum = (totalElementNum + vecCoreNum - 1U) / vecCoreNum; | ||
| uint32_t tmpBlockIdx = AscendC::GetBlockIdx(); | ||
| uint64_t gmOffset = tmpBlockIdx * singleCoreMaxElementNum; |
There was a problem hiding this comment.
In Ascend C, global memory (GM) accesses via DataCopy and DataCopyPad require the start address to be 32-byte aligned. Since singleCoreMaxElementNum is not guaranteed to be a multiple of the 32-byte alignment size (e.g., 8 elements for float, 16 elements for half/float16), gmOffset for cores > 0 can easily become unaligned. This will lead to runtime alignment faults, memory corruption, or undefined behavior during DataCopy/DataCopyPad.\n\nTo resolve this, the partitioning logic must align singleCoreMaxElementNum to the 32-byte boundary (i.e., round it up to a multiple of 32 / sizeof(T) elements), and adjust the actual count for the last core accordingly.
uint64_t alignElements = 32 / sizeof(T);\n uint64_t singleCoreMaxElementNum = ((totalElementNum + vecCoreNum - 1U) / vecCoreNum + alignElements - 1U) / alignElements * alignElements;\n uint32_t tmpBlockIdx = AscendC::GetBlockIdx();\n uint64_t gmOffset = tmpBlockIdx * singleCoreMaxElementNum;| idx, vals, _ = quant_lightning_indexer_candidate( | ||
| q_in, k, w_in, q_descale_in, k_descale, topk, quant_mode, | ||
| candidate_topk_index=cand_in, | ||
| cu_seqlens_q=cu_seqlens_q, cu_seqlens_k=None, seqused_q=None, seqused_k=seqused_q, | ||
| cmp_residual_k=cmp_residual_k, block_table=block_table, output_idx_offset=output_idx_offset, | ||
| metadata=metadata, max_seqlen_q=max_seqlen_q, layout_q=layout_q, layout_k="PA_BBND", | ||
| mask_mode=mask_mode, cmp_ratio=cmp_ratio, | ||
| candidate_mode=2, candidate_topk_blocks=topk_blocks, candidate_block_size=int(candidate_block_size)) |
There was a problem hiding this comment.
The quant_lightning_indexer_candidate operator is defined in the schema to return a 3-tuple of Tensors: (Tensor, Tensor, Tensor). However, in quant_lightning_indexer_candidate_consumer, the return value is unpacked into a 2-tuple idx, vals. This will unconditionally raise a ValueError: too many values to unpack (expected 2) at runtime whenever this function is called.\n\nTo fix this, unpack the 3-tuple correctly by capturing the third unused tensor (e.g., using _).
idx, vals, _ = quant_lightning_indexer_candidate(\n q_in, k, w_in, q_descale_in, k_descale, topk, quant_mode,\n candidate_topk_index=cand_in,\n cu_seqlens_q=cu_seqlens_q, cu_seqlens_k=None, seqused_q=None, seqused_k=seqused_q,\n cmp_residual_k=cmp_residual_k, block_table=block_table, output_idx_offset=output_idx_offset,\n metadata=metadata, max_seqlen_q=max_seqlen_q, layout_q=layout_q, layout_k="PA_BBND",\n mask_mode=mask_mode, cmp_ratio=cmp_ratio,\n candidate_mode=2, candidate_topk_blocks=topk_blocks, candidate_block_size=int(candidate_block_size))| at::Tensor sparseIndicesOut = std::get<0>(quantLightningIndexerOutput); | ||
| at::Tensor sparseValuesOut = std::get<1>(quantLightningIndexerOutput); | ||
|
|
||
| int64_t keyHeadNum = (keyLayoutStr == "TND") ? key.size(DIM_1) : key.size(DIM_2); | ||
| at::Tensor candidateTopkIndexOut; | ||
| if (candidateMode == 1) { | ||
| at::SmallVector<int64_t, SIZE> candSize; | ||
| if (queryLayoutStr == "BSND") { | ||
| candSize = {query.size(DIM_0), query.size(DIM_1), keyHeadNum, candidateTopkBlocks}; | ||
| } else { | ||
| candSize = {query.size(DIM_0), keyHeadNum, candidateTopkBlocks}; | ||
| } |
There was a problem hiding this comment.
In ACLNN, optional outputs that are not used/requested by the caller should be passed as a null pointer (nullptr / undefined at::Tensor()). Passing an empty tensor with shape {0} (via at::empty({0}, ...)) when candidateMode != 1 can cause ACLNN parameter validation failures or runtime crashes because the shape {0} does not match the expected shape of the optional output.\n\nTo fix this, initialize candidateTopkIndexOut as an undefined/null tensor at::Tensor() when candidateMode != 1.
at::Tensor candidateTopkIndexOut;\n if (candidateMode == 1) {\n at::SmallVector<int64_t, SIZE> candSize;\n if (queryLayoutStr == "BSND") {\n candSize = {query.size(DIM_0), query.size(DIM_1), keyHeadNum, candidateTopkBlocks};\n } else {\n candSize = {query.size(DIM_0), keyHeadNum, candidateTopkBlocks};\n }\n candidateTopkIndexOut = at::empty(candSize, query.options().dtype(at::kInt));\n } else {\n candidateTopkIndexOut = at::Tensor();\n }| gert::Shape *candidateTopkIndexShape = context->GetOutputShape(CANDIDATE_TOPK_INDEX_OUTPUT_INDEX); | ||
| OP_CHECK_NULL_WITH_CONTEXT(context, candidateTopkIndexShape); | ||
| const int32_t *candidate_mode = attrs->GetAttrPointer<int32_t>(ATTR_CANDIDATE_MODE_INDEX); | ||
| uint32_t candidateMode = (candidate_mode != nullptr) ? static_cast<uint32_t>(*candidate_mode) : 3U; | ||
| if (candidateMode == CANDIDATE_MODE_SOURCE) { | ||
| OP_CHECK_IF(inputLayoutQueryPtrStr != "BSND", | ||
| OP_LOGE("QuantLightningIndexerV2", | ||
| "candidate_mode=1 only supports layout_q=BSND, but got %s.", | ||
| inputLayoutQueryPtrStr.c_str()), | ||
| return GRAPH_FAILED); | ||
| const int64_t *candidate_topk_blocks = attrs->GetAttrPointer<int64_t>(ATTR_CANDIDATE_TOPK_BLOCKS_INDEX); | ||
| int64_t candBlocks = (candidate_topk_blocks != nullptr) ? *candidate_topk_blocks : CANDIDATE_TOPK_BLOCKS_FIX; | ||
| *candidateTopkIndexShape = *sparseIndicesShape; | ||
| candidateTopkIndexShape->SetDim(sparseIndicesShape->GetDimNum() - 1, candBlocks); | ||
| } else { | ||
| candidateTopkIndexShape->SetDimNum(1); | ||
| candidateTopkIndexShape->SetDim(0, 0); | ||
| } |
There was a problem hiding this comment.
In GE (Graph Engine), optional outputs that are not connected in the graph will have their corresponding output shape pointers as nullptr in InferShape. Unconditionally checking for null and returning GRAPH_FAILED via OP_CHECK_NULL_WITH_CONTEXT breaks the optional contract of the output, causing graph compilation to fail if candidate_topk_index_out is not connected.\n\nTo fix this, check if candidateTopkIndexShape is nullptr and gracefully skip setting its shape, rather than failing the entire compilation.
gert::Shape *candidateTopkIndexShape = context->GetOutputShape(CANDIDATE_TOPK_INDEX_OUTPUT_INDEX);\n if (candidateTopkIndexShape != nullptr) {\n const int32_t *candidate_mode = attrs->GetAttrPointer<int32_t>(ATTR_CANDIDATE_MODE_INDEX);\n uint32_t candidateMode = (candidate_mode != nullptr) ? static_cast<uint32_t>(*candidate_mode) : 3U;\n if (candidateMode == CANDIDATE_MODE_SOURCE) {\n OP_CHECK_IF(inputLayoutQueryPtrStr != "BSND",\n OP_LOGE("QuantLightningIndexerV2",\n "candidate_mode=1 only supports layout_q=BSND, but got %s.",\n inputLayoutQueryPtrStr.c_str()),\n return GRAPH_FAILED);\n const int64_t *candidate_topk_blocks = attrs->GetAttrPointer<int64_t>(ATTR_CANDIDATE_TOPK_BLOCKS_INDEX);\n int64_t candBlocks = (candidate_topk_blocks != nullptr) ? *candidate_topk_blocks : CANDIDATE_TOPK_BLOCKS_FIX;\n *candidateTopkIndexShape = *sparseIndicesShape;\n candidateTopkIndexShape->SetDim(sparseIndicesShape->GetDimNum() - 1, candBlocks);\n } else {\n candidateTopkIndexShape->SetDimNum(1);\n candidateTopkIndexShape->SetDim(0, 0);\n }\n }| if (BUILD_OPEN_PROJECT) | ||
| # Only the Ascend 950 implementation requires LightningIndexerV2 headers. | ||
| if (ASCEND_COMPUTE_UNIT STREQUAL "ascend950") | ||
| set(quant_lightning_indexer_v2_depends attention/lightning_indexer_v2 CACHE INTERNAL "Dependencies for quant_lightning_indexer_v2") |
There was a problem hiding this comment.
In CMake, when changing a cache variable's type or value, FORCE is required to overwrite any existing cached entry in CMakeCache.txt. Without FORCE, incremental builds or existing build directories will retain the old cached value (which was set unconditionally for all compute units), ignoring the new if (ASCEND_COMPUTE_UNIT STREQUAL "ascend950") condition.\n\nTo fix this, add FORCE to the set command.
set(quant_lightning_indexer_v2_depends attention/lightning_indexer_v2 CACHE INTERNAL "Dependencies for quant_lightning_indexer_v2" FORCE)
| s2_t = seqused_k.to(torch.int64) # (B,) 每 batch key 有效长度 | ||
| res_t = cmp_residual_k.to(torch.int64) if cmp_residual_k is not None else torch.zeros_like(s2_t) | ||
| act_k_t = s2_t * int(cmp_ratio) + res_t |
There was a problem hiding this comment.
If seqused_k or cmp_residual_k are on CPU, s2_t and res_t will also be on CPU. Later in the function, they are indexed by row_b (which is on NPU), leading to a PyTorch device mismatch runtime error.\n\nTo prevent this, explicitly specify device=device when converting these tensors.
| s2_t = seqused_k.to(torch.int64) # (B,) 每 batch key 有效长度 | |
| res_t = cmp_residual_k.to(torch.int64) if cmp_residual_k is not None else torch.zeros_like(s2_t) | |
| act_k_t = s2_t * int(cmp_ratio) + res_t | |
| s2_t = seqused_k.to(device=device, dtype=torch.int64) # (B,) 每 batch key 有效长度\n res_t = cmp_residual_k.to(device=device, dtype=torch.int64) if cmp_residual_k is not None else torch.zeros_like(s2_t)\n act_k_t = s2_t * int(cmp_ratio) + res_t |
| Reg::LoadAlign<T, Reg::LoadDist::DIST_BLK>(vregLse, | ||
| (__ubuf__ float *&)lseUb + splitKVIndex * dealRowCount * 8 + k * 8); |
There was a problem hiding this comment.
Using reference casting (__ubuf__ float *&) to perform pointer arithmetic is unsafe and violates strict aliasing rules, which can lead to undefined behavior or compiler misoptimizations. Instead, cast the pointer value directly using reinterpret_cast<__ubuf__ float *> or a standard value cast ( __ubuf__ float * ).
Reg::LoadAlign<T, Reg::LoadDist::DIST_BLK>(vregLse,\n reinterpret_cast<__ubuf__ float *>(lseUb) + splitKVIndex * dealRowCount * 8 + k * 8);9eb4370 to
d70396b
Compare
d70396b to
0e89b7b
Compare
Add Quant Lightning Indexer v2 enhancements, Sparse Flash MLA and metadata operators, Triton compressor helpers, and focused operator tests required by DeepSeek V4.1. Signed-off-by: GDzhu01 <116337067+GDzhu01@users.noreply.github.com>
Signed-off-by: GDzhu01 <116337067+GDzhu01@users.noreply.github.com>
Signed-off-by: GDzhu01 <116337067+GDzhu01@users.noreply.github.com>
Signed-off-by: GDzhu01 <116337067+GDzhu01@users.noreply.github.com>
Signed-off-by: GDzhu01 <116337067+GDzhu01@users.noreply.github.com>
Signed-off-by: GDzhu01 <116337067+GDzhu01@users.noreply.github.com>
| # projection in the checkpoint. | ||
| del self.hc_head_fn, self.hc_head_base, self.hc_head_scale, self.hc_norm | ||
| topology = build_layer_plan(self.config) | ||
| max_tokens = vllm_config.scheduler_config.max_num_batched_tokens |
There was a problem hiding this comment.
【可靠性】candidate buffer 仅按 scheduler max_num_batched_tokens 分配,与后面 Engram 使用 max scheduler/cudagraph capture 的策略不一致。较大的捕获 shape 会使 shared.candidates 切片不足并在 indexer copy_ 处失败。建议统一静态 token 容量的计算函数,避免各模块各自取值。
| and "content_blocks" in merged[-1] | ||
| and merged[-1].get("task") is None | ||
| ): | ||
| merged[-1]["content_blocks"].extend(content_blocks) |
There was a problem hiding this comment.
【安全、可靠性】显式提供 content_blocks 时未验证其为 list;若传入字符串或字典,extend 会按字符/键拆入前一条消息,随后渲染阶段再以 block.get 崩溃。建议在合并前规范化并严格校验每个 block 为受支持对象。
| raise ValueError("A candidate source must use the unfiltered position TopK") | ||
| if uses_candidate_filter and candidates is None: | ||
| raise RuntimeError("V4.1 candidate-filtering indexer ran before its source") | ||
| if self.width != 128 or self.n_heads not in (32, 64): |
There was a problem hiding this comment.
【可靠性】算子约束校验了 width/head 数,却没有验证 0 <= rope_width <= width 且旋转宽度满足算子对齐要求。rope_width>width 会产生负 partial_slice 起点,错误配置可能旋转错误区间。建议把 RoPE 几何契约纳入这里的集中校验。
| else: | ||
| raise TypeError(f"Unsupported V4.1 cache spec: {type(spec).__name__}") | ||
|
|
||
| num_reqs = int(getattr(common, "num_reqs", common.seq_lens.shape[0])) |
There was a problem hiding this comment.
【可靠性】num_reqs/num_input_tokens 未与构造时的 max_reqs/max_tokens 容量核对,超限时多个预分配 slice 会静默变短,错误直到 copy_ 或自定义算子处才暴露。建议保存容量并在 build 开头统一 fail-fast,错误中给出实际值与上限。
| if candidate_block_size != 8: | ||
| raise ValueError("The current A3 candidate kernel requires candidate_block_size=8") | ||
| candidate_shape = (query.shape[0], 1, candidate_topk_blocks) | ||
| if uses_candidate_filter and (candidates.shape != candidate_shape or candidates.dtype != torch.int32): |
There was a problem hiding this comment.
【安全、可靠性】candidate consumer 只检查 shape/dtype,没有检查 device 与 query/算子一致。CPU 或其他 NPU 上形状正确的 candidates 会直接进入自定义算子,可能报底层错误甚至访问非法地址。建议同时校验 device、contiguous 和取值范围(-1 或合法 block ID)。
| raise ValueError("Every DeepSeek V4.1 KV source must also be an index source") | ||
| if candidate_source not in kv_sources: | ||
| raise ValueError("DeepSeek V4.1 candidate_source_layer must be a KV source") | ||
| if candidate_topk_blocks <= 0 or candidate_block_size <= 0 or index_topk <= 0: |
There was a problem hiding this comment.
【可靠性、可维护性】这里只要求候选参数为正,但下游 A3 内核实际还要求 candidate_topk_blocks<=2048 且为 64 倍数、candidate_block_size==8、index_topk<=2048。当前模型能加载成功后在首次请求才失败。建议在 layer plan 阶段校验完整硬件契约,实现启动即失败。
| mask.new_zeros(capacity), | ||
| ) | ||
| buffers, mask_buffer = self._engram_input_buffers | ||
| padded_mask = mask_buffer[:output_tokens] |
There was a problem hiding this comment.
【可靠性】output_tokens 可由 padded_tokens 提升到 _engram_max_tokens 以上,但底层 buffers 仍只按固定 capacity 分配;切片会静默截断,返回张量长度小于请求值。建议在切片前拒绝 output_tokens>capacity,或按已知最大 capture size 一次性正确分配。
| build_compressor_metadata=True, | ||
| ): | ||
| super().__init__(kv_cache_spec, layer_names, vllm_config, device) | ||
| max_tokens = getattr(vllm_config.scheduler_config, "max_num_batched_tokens", 4096) |
There was a problem hiding this comment.
【可靠性】metadata builder 的全部静态缓冲只按 max_num_batched_tokens 分配,而图捕获最大 token 数可能由 max_cudagraph_capture_size 决定并更大。build 中随后按 num_input_tokens 切片/写入,会得到短缓冲并导致 copy_ 失败或图地址不稳定。建议容量取两者最大值,并在 build 入口做显式上界检查。
| """Expose images nested in reference tool_result/content_blocks to vLLM.""" | ||
| result = [] | ||
| for block in blocks: | ||
| if block.get("type") == "tool_result": |
There was a problem hiding this comment.
【安全、可靠性】blocks 来自外部消息,但循环直接调用 block.get;列表中混入字符串、None 或其他非对象值会抛 AttributeError 并形成 500。建议在递归入口验证 blocks/list 元素类型,并返回带消息位置的协议错误。
| for idx, tc in enumerate(msg["tool_calls"]): | ||
| tc_id = tc.get("id") or tc.get("function", {}).get("id", "") | ||
| if tc_id: | ||
| last_tool_call_order[tc_id] = idx |
There was a problem hiding this comment.
【可靠性】重复的 tool_call id 会在字典中后写覆盖前写,排序结果与调用顺序不再一一对应,却没有任何错误提示。建议检测重复/空 ID,并验证每个 tool_result 恰好对应一个先前调用,避免结果串线。
| # Deliberately do not call the V4 dSPark constructor: V4.1 has delayed | ||
| # mHC state between blocks and no terminal hc_head parameters. | ||
| torch.nn.Module.__init__(self) | ||
| assert vllm_config.speculative_config is not None |
There was a problem hiding this comment.
【可靠性】speculative_config 属于运行配置,这里及 CausalLM 构造器使用 assert 检查;优化模式会移除检查并在下一行产生 NoneType 异常。建议改为显式 ValueError,并对 draft_model_config、target_layer_ids 等必需字段一起做启动期校验。
| response_format = response_format.model_dump(by_alias=True) | ||
| schema = None | ||
| if response_format and response_format.get("type") == "json_schema": | ||
| schema = response_format["json_schema"]["schema"] |
There was a problem hiding this comment.
【安全、可靠性】json_schema 分支直接索引 response_format[json_schema][schema],对缺字段或错误类型会泄漏 KeyError/TypeError。建议按 OpenAI 结构显式验证 type、json_schema 对象和 schema 对象,并统一转换为 ValueError,便于 API 层返回 4xx。
| # history reads. | ||
| compressed_list = compressed.tolist() | ||
| position_list = positions.tolist() | ||
| page_indices = block_table[request_ids, positions // block_size].tolist() |
There was a problem hiding this comment.
【可靠性】update 未验证 input_ids、positions、request_ids 三者等长,也未验证 block_size>0 与 request/position 范围;随后 zip 会静默截断写入,而 history 仍按 len(input_ids) 分配,产生部分未更新的哈希。建议入口做完整 shape/range 校验,并使用 strict=True。
| result = torch.empty((ids.numel(), self.width), dtype=torch.bfloat16, device=device) | ||
| else: | ||
| result = output | ||
| if result.shape != (ids.numel(), self.width) or result.device != device: |
There was a problem hiding this comment.
【可靠性】复用 output 时只校验 shape/device,不校验 dtype、布局和可写性;非 BF16 buffer 会发生静默转换,非连续 view 还可能让 collective/赋值行为与固定输出契约不一致。建议要求 BF16、contiguous 且无重叠,或明确支持的布局。
| page_indices = block_table[request_ids, positions // block_size].tolist() | ||
| if len(input_ids) < _PAGE_WRITE_NUMPY_MIN_TOKENS: | ||
| for token, position, page in zip(compressed_list, position_list, page_indices): | ||
| if page not in self.pages: |
There was a problem hiding this comment.
【安全、可靠性】block_table 中的 -1/越界 page ID 会被当作普通字典键创建页面,多个无效请求可能共享 pages[-1] 并相互污染历史。建议在任何写入前要求 page>=0,并验证 request/block 索引落在表范围内。
| event.record(torch.npu.current_stream(device)) | ||
| self._offload_events[source_ptr] = event | ||
|
|
||
| def load_checkpoint(self, model_path, key, chunk_rows=65536): |
There was a problem hiding this comment.
【可靠性】chunk_rows 可由调用方传入,但未要求为正;0 会让 range 直接抛错,负值则可能完全跳过加载并让未初始化权重进入推理。建议显式校验 chunk_rows>0,并在结束后确认本 rank 行区间已完整加载。
| class DeepseekV41SWASpec(AscendSlidingWindowMLASpec): | ||
| def is_uniform_with_collection(self, specs): | ||
| return all( | ||
| isinstance(s, DeepseekV41SWASpec) and s.sliding_window == self.sliding_window for s in specs.values() |
There was a problem hiding this comment.
【可靠性】SWA uniform 判定只比较 sliding_window,未比较 block_size、dtype、head_size、num_kv_heads 等实际页面几何。不同布局可能被合并进同一 UniformType group,后续按一个 spec 分配/寻址。建议比较所有影响 page layout 的字段。
| torch.empty((0, self.primes.shape[0], columns), dtype=torch.int64, device="cpu"), | ||
| torch.empty(0, dtype=torch.bool, device="cpu"), | ||
| ) | ||
| compressed = self.token_map[input_ids] |
There was a problem hiding this comment.
【安全、可靠性】input_ids 直接用于 token_map 高级索引。图填充或异常输入中的负 token ID 会从末尾反向取值而不是报错,过大 ID 才报 IndexError,可能污染 n-gram 历史。建议先验证范围,或明确把占位/填充值映射到 pad/barrier。
| if vocab_size != config.engram_compressed_vocab_size: | ||
| raise ValueError(f"Engram compressed vocabulary mismatch: {vocab_size}") | ||
| self.token_map = torch.tensor(token_map, dtype=torch.int64) | ||
| self.pad_id = token_map[config.engram_pad_id] |
There was a problem hiding this comment.
【可靠性】engram_pad_id 直接索引 token_map,没有先验证 0<=id<len(token_map)。错误 checkpoint 会抛无上下文 IndexError;负 ID 还会被 Python 当作反向索引而静默选错 token。建议显式校验 pad/image token ID 契约。
| def _get_thread_local_tokenizer(tokenizer): | ||
| cached = getattr(_TOKENIZER_THREAD_LOCAL, "tokenizer", None) | ||
| source_id = getattr(_TOKENIZER_THREAD_LOCAL, "source_id", None) | ||
| if cached is None or source_id != id(tokenizer): |
There was a problem hiding this comment.
【可维护性、可靠性】thread-local 缓存只保存源 tokenizer 的 id,不持有源对象;源释放后 Python 可能复用该 id,新 tokenizer 会错误命中旧 deepcopy。建议同时保存源对象/弱引用并用 is 判断身份,或使用带对象生命周期的缓存。
|
|
||
| def set_rows(self, start, rows): | ||
| """Load BF16 rows into local storage without allocating a BF16 table copy.""" | ||
| end = start + rows.shape[0] |
There was a problem hiding this comment.
【安全、可靠性】set_rows 没有检查 start>=0、end<=local_rows、输入 width 与 dtype。负 start 会按 Python 反向切片,越界 end 可能出现不一致 copy,导致权重被写到错误位置。建议把它作为明确的本地 shard 写入契约校验。
| mm_prompt_updates, | ||
| ) | ||
| placeholders: dict[str, list[PlaceholderFeaturesInfo]] = {modality: [] for modality in base_placeholders} | ||
| ordered = sorted( |
There was a problem hiding this comment.
【可靠性】sorted 的 tuple 第三项是 PlaceholderFeaturesInfo;若两个同 modality 的占位符 start_idx 相同,前两项相等后 Python 会尝试比较不可排序对象并抛 TypeError。建议只对 (start_idx, stable_sequence_no) 使用 key 排序,并显式处理重叠占位符。
| sizes.append(current) | ||
| per_ngram.append(tuple(sizes)) | ||
| primes.append(tuple(per_ngram)) | ||
| return cls( |
There was a problem hiding this comment.
【安全、可靠性】构造出的各 n-gram/head 素数桶通过 offset 拼接,但从未核对每层桶总容量是否 <= 对应 num_embeddings。配置较小时 hash 会生成超出 Engram 表行数的 ID,直到分布式路由才报错。建议逐层校验 sum(primes[layer]) 与表行数,并在加载期失败。
| self._offload_buffer_index = {} | ||
| self._offload_events = {} | ||
| # Ceil partition leaves at most size-1 unused rows, never a replica. | ||
| self.shard_rows = (rows + query_group.size - 1) // query_group.size |
There was a problem hiding this comment.
【可靠性】rows、width 与 query_group.size/rank 未做基础合法性校验;size=0 会除零,负 rows/width 会进入非法分配,越界 rank 会计算出错误 shard。建议在计算 shard_rows 前集中校验 rows/width/size>0 且 0<=rank<size。
| max_h_float = max_w_float * r | ||
| cell = patch_size * downsample_ratio | ||
| if max_w_float < 1.0: | ||
| return (max_n_token - 2) // 2 * cell, cell |
There was a problem hiding this comment.
【可靠性、多模态】即便 max_n_token>2,极端宽高比与很小预算仍可能让 (max_n_token-2)//2 或 (max_n_token-3) 为 0,返回零高度/宽度,随后 PIL resize/reshape 失败。建议保证求解结果至少一个 LLM cell,并对无法容纳 start/image/newline/end 的预算明确报错。
| "key_and_gate_basis": "original", | ||
| "runtime_delta_rotation": False, | ||
| } | ||
| if any(rotation.get(name) != value for name, value in supported_rotation.items()): |
There was a problem hiding this comment.
【安全、可维护性】这里仅核对已知四个 rotation 字段的值,未知字段会被保留并静默接受,与注释中的 fail closed 不一致。未来 checkpoint 增加改变数学语义的字段时,旧 runtime 可能误判为兼容。建议同时要求 key 集合完全匹配,或显式拒绝未知键。
| assert self.image_start is not None | ||
| assert self.image_end is not None | ||
| assert self.image_newline is not None | ||
| span[types == IMAGE_START] = self.image_start.to(dtype) |
There was a problem hiding this comment.
【高效】image_start/end/newline 参数固定以 FP32 创建,但每次构造图片 span 都调用 .to(dtype),在 BF16 模型上会反复生成三个临时张量,也不利于图捕获地址稳定。建议加载后缓存目标 dtype 视图,或让参数按 model_config.dtype 注册并在权重加载时一次转换。
| # V4.1 uses one image token for every span role. The adjacent reserved | ||
| # token is used only for vLLM's ratio-2 compressor-alignment row. | ||
| self.image_sentinel_base_id = self.image_token_id | ||
| self.image_pad_token_id = self.image_token_id + 1 |
There was a problem hiding this comment.
【安全、可靠性】image_pad_token_id 直接按 image_token_id+1 推导,未验证仍在 vocab 范围内、不是 eos/pad 等已有特殊 token,也未尊重 checkpoint 显式声明。建议优先读取配置值并校验唯一性/词表边界,避免静默改写普通 token 的 embedding 语义。
| return 1 | ||
|
|
||
| def is_uniform_with_collection(self, specs): | ||
| return all(type(s) is type(self) and s.block_size == self.block_size for s in specs.values()) |
There was a problem hiding this comment.
【可靠性】CircularBuffer uniform 判定只比较具体类型与 block_size,未比较 dtype、num_kv_heads、head_size,因此不同 page_size 的状态缓存可能被判为 uniform。建议比较所有决定 real_page_size_bytes 的字段,或直接比较规范化页面布局。
|
|
||
| def is_uniform_with_collection(self, specs): | ||
| return all( | ||
| isinstance(s, DeepseekV41DraftSWASpec) |
There was a problem hiding this comment.
【可靠性】Draft SWA 的 uniform 判定只比较 block_size/window,虽然 post_init 固定了 dtype 与 head 数,却未固定 head_size 等页面几何。不同 draft spec 仍可能被合并并按同一布局分配。建议复用完整 layout compatibility 比较,而不是手写部分字段。
What this PR does / why we need it?
Adds the framework and model integration for DeepSeek V4.1 on Ascend:
Depends on #16422. The operator implementation and operator-level tests are kept in that PR; this PR contains the framework/features and their tests. Until #16422 merges, GitHub's main-based aggregate diff includes the dependency commit; reviewers can inspect commits
0ef4335faandd70396b96independently.Does this PR introduce any user-facing change?
Yes. It registers and enables DeepSeek V4.1 serving on Ascend, including eager and graph execution paths, DSA context parallelism, DSpark speculative decoding, and optional Engram support.
How was this patch tested?
FULL_DECODE_ONLY: capture completed (2/2), health/model endpoint passed17 * 23returned391; two-sentence Rayleigh-scattering answer was identical and correctRemote evidence:
eager job:
job-20260912T124029Z-47cb97dbgraph job:
job-20260912T124917Z-0914d504graph service endpoint inside the remote container:
http://127.0.0.1:30165vLLM main: vllm-project/vllm@a97dacb