19 KiB
Calculet NPU 异步执行器与生产稳定性实施规格
日期:2026-08-02
适用基线:CalRT 0.7.6、llama.cpp fd9bd632、当前 Qwen3 batch=1 双芯粒 calbin
1. 文档目的和边界
本文不是架构草图,而是改造当前 llama-calrt.cpp 的实施规格。第一阶段只做 C2 工程改造,不假设 Runtime 支持多物理设备,也不把未验证的 parallel/ping-pong 作为生产能力。
完成后应得到:
infer()、CSR、Wait()的所有错误都能向上返回,不再abort()。- KV 不足有界等待,不再 CPU 无限忙等。
- 每个在途任务独占 InputBuf、OutputBuf 和所有 host backing storage。
- 请求线程、设备提交线程、完成线程解耦;Phase A 仍保持
max_inflight=1。 - 取消、超时、设备异常、模型切换都有可验证的状态机。
- Phase B 才在隔离 runner 中验证 ping/pong;验证不通过可无代码回退到 1 in-flight。
2. 当前代码级问题清单
| ID | 当前实现 | 直接风险 | 第一落点 |
|---|---|---|---|
| S01 | KvManager::Apply() 失败后 do...while(!ok) |
无 sleep、无 deadline、无法取消,可能占满 CPU 永久挂起 | admission queue |
| S02 | infer 后二次 Apply() 失败 abort() |
单请求资源问题杀死整个服务 | KV transaction + error return |
| S03 | SetCsrByName() 返回码未检查 |
CSR 名称或值错误仍提交图 | prepare_job() |
| S04 | infer() 返回码未检查 |
submit 失败仍进入 Wait() |
submit_job() |
| S05 | Wait() 返回码未检查 |
timeout/CCU exception 被当作成功 | complete_job() |
| S06 | context 只有一组 prefill/decode vectors | 第二任务可覆盖第一任务 DMA 内存 | BufferLease 独占 backing |
| S07 | 每次创建 Runtime buffer | 热路径分配和初始化开销 | generation-scoped pool |
| S08 | substring 选择子模型且最后匹配覆盖前者 | 名称碰撞时选错图 | exact manifest |
| S09 | prefill output D0=16 |
新布局下 slice/untile 错误 | layout descriptor |
| S10 | 整个 BF16 vocab 先转临时 FP32 vector | 多一次分配、遍历和拷贝 | 直接转换到目标区 |
| S11 | 计时存在 std::map 且按 seq 共享 |
并发读写数据竞争、历史值残留 | per-job timings + metrics |
| S12 | dump 路径硬编码用户目录 | 隐私、磁盘和可移植性风险 | 默认禁用的受控 dump sink |
3. 代码改动边界
建议保持 llama.cpp 上层调度接口稳定,新建四个文件并收缩原有文件职责:
src/llama-calrt.cpp 只保留 llama batch 适配和同步兼容入口
src/llama-calrt.h 对上接口、manifest 和结果类型
src/llama-calrt-executor.h/.cpp DeviceExecutor、JobContext、队列和 completion
src/llama-calrt-buffer-pool.h/.cpp BufferPool、BufferLease、layout 校验
src/llama-calrt-kv-txn.h/.cpp KV reservation/commit/rollback
src/llama-calrt-error.h/.cpp CalRT 错误分类、日志和服务错误映射
llama-kv-cache-calrt-adapter.cpp 只声明已经真实支持的语义。空实现不能继续被当成成功;无法实现的 seq_cp/keep/add/div/state_* 应显式返回“不支持”或在调用前禁用对应上层功能。
4. 精确数据类型
4.1 模型清单
enum class CalTaskKind : uint8_t { PrefillIds, PrefillEmbeds, Decode, Vision };
struct TensorContract {
std::string name;
std::vector<uint64_t> shape;
calrt::PrimitiveType dtype;
uint64_t byte_size;
bool sliceable;
};
struct SubmodelContract {
CalTaskKind kind;
std::string exact_name;
TensorContract token_or_embed;
TensorContract position;
TensorContract logits;
std::string csr_cur_seq;
std::string csr_total_seq;
uint32_t max_batch;
uint32_t max_seq;
uint32_t logits_tile_d0;
};
struct ModelGeneration {
uint64_t generation_id;
std::string calbin_sha256;
std::string calrt_version;
std::unordered_map<CalTaskKind, SubmodelContract> submodels;
};
当前产物 manifest 的固定值是:prefill/decode max_batch=1、max_seq=40960、token/position 为 S32、logits 为 BF16 [1,1,151936]。logits_tile_d0=16 只能由当前图 metadata 解析后写入 manifest,不能作为代码默认值。
4.2 任务与 buffer
enum class JobState : uint8_t {
Created, Admitted, Preparing, Queued, Submitted,
Completing, Succeeded, Failed, CancelRequested, Quarantined
};
enum class EngineSlot : int8_t { Auto = -1, Ping = 0, Pong = 1 };
struct JobTimings {
int64_t accepted_us = 0;
int64_t admitted_us = 0;
int64_t queued_us = 0;
int64_t submitted_us = 0;
int64_t completed_us = 0;
double h2d_ms = 0;
double device_wait_ms = 0;
double d2h_ms = 0;
double postprocess_ms = 0;
};
struct BufferStorage {
std::vector<uint8_t> input0;
std::vector<uint8_t> position;
std::vector<uint8_t> output0;
};
struct BufferSlot {
uint64_t generation_id;
CalTaskKind task;
EngineSlot engine;
std::unique_ptr<calrt::CalrtInputBuf> input;
std::unique_ptr<calrt::CalrtOutputBuf> output;
BufferStorage host;
bool sliced = false;
bool quarantined = false;
};
struct JobContext {
uint64_t job_id;
uint64_t generation_id;
uint64_t request_id;
calrt::KvSeqId seq_id;
CalTaskKind task;
EngineSlot engine;
uint32_t past_len;
uint32_t token_count;
std::chrono::steady_clock::time_point deadline;
std::atomic<JobState> state{JobState::Created};
std::atomic<bool> cancel_requested{false};
std::optional<BufferLease> buffers;
std::optional<KvReservation> kv;
JobTimings timings;
Status result;
std::promise<InferenceResult> promise;
};
BufferLease 必须是 move-only。它的析构只能在 Succeeded/Failed 且确认 Runtime 不再 DMA 时把 slot 归还池;Submitted 状态析构应触发断言并转入 quarantine,而不是静默复用。
5. 所有权与锁不变量
5.1 强制不变量
DeviceExecutor独占VirtualDevice、Calbin和KvManager。- 只有 submitter 线程调用
infer();只有 completion 线程对已提交 job 调用Wait()。 - 一个
BufferSlot同一时刻最多属于一个JobContext。 MapBuf()后 host vector 在完成前不得 resize、move、swap 或释放。- job 的
generation_id必须与 pool/model generation 一致。 - KV reserve 成功不等于 KV commit;只有设备完成且输出通过校验才 commit。
- cancel 不释放底层资源,只阻止结果交付;仍需完成
Wait()或执行已验证的 reset。 - reset 后所有旧 generation buffer 一律 quarantine 并重建。
5.2 锁和顺序
| 锁 | 保护对象 | 可持锁调用 Runtime 吗 |
|---|---|---|
lifecycle_mu |
model generation、Ready/Draining/Recovering | configure/reset 可在独占阶段调用;不能同时持 queue 锁 |
admission_mu |
KV 等待队列、token budget | 否 |
pool_mu |
free/leased/quarantine 集合 | 否 |
submit_mu |
submit queue 条件变量 | 否 |
completion_mu |
submitted job 队列 | Wait() 前必须释放 |
若一次操作需要多个锁,固定顺序为:lifecycle_mu -> admission_mu -> pool_mu -> submit_mu。completion 线程不可持任何上述锁执行阻塞 Wait()。
6. 端到端调用流
6.1 同步兼容入口
StatusOr<InferenceResult> calrt_context::decode_sync(const llama_batch &batch) {
auto request = build_request(batch);
CAL_ASSIGN_OR_RETURN(auto handle, executor_->enqueue(std::move(request)));
return handle.wait_until(request.deadline);
}
原 int calrt_decode(...) 可暂时把 Status 映射成 llama 返回码,但日志必须保留 job_id、seq_id、generation、CalRT 原始错误码和阶段。
6.2 enqueue()
validate request
-> exact submodel lookup
-> deadline/cancel check
-> admission reserve KV + token budget
-> acquire BufferLease
-> copy inputs + set slices + set CSR
-> push submit queue
-> return JobHandle
任一步失败都必须按逆序释放已取得资源。尚未调用 infer() 时 buffer 可正常归还;已经提交后只能由 completion/recovery 路径处理。
6.3 prepare_job() 检查顺序
token_count > 0。past_len + token_count <= max_seq,使用 64 位加法防溢出。- 当前 calbin batch=1 时,
seq_count == 1。 - token 与 position 指针非空,长度一致;embed 路径单独检查元素数。
- tensor exact name 存在,dtype、rank、shape 和 ByteSize 与 manifest 一致。
input_bytes <= tensor.ByteSize();slice offset/size 均在范围内。- 两个
SetCsrByName()都返回CalrtSuccess。 MapBuf()和SliceTensor()的返回/状态均检查。
6.4 submitter
void DeviceExecutor::submit_one(std::shared_ptr<JobContext> job) {
if (job->cancel_requested.load()) {
fail_before_submit(job, Status::Cancelled());
return;
}
const auto rc = use_fixed_engine_
? calrt::infer_with_fixed_task_type(vdev_.get(), model(job),
job->buffers->input(), job->buffers->output(), to_task_type(job->engine))
: calrt::infer(vdev_.get(), model(job),
job->buffers->input(), job->buffers->output());
if (rc != CalrtSuccess) {
rollback_before_device_execution(job, from_calrt(rc, "submit"));
return;
}
transition(job, JobState::Queued, JobState::Submitted);
completion_queue_.push(job);
}
Phase A 的 semaphore 容量固定为 1。Phase B 只有 runner 验证后才改成 2,并给每个 slot 固定 Ping/Pong;不能把并发 HTTP 请求数直接当作 in-flight 上限。
6.5 completion
pop submitted job (不持锁)
-> OutputBuf::Wait()
-> 读取 OutputBuf::GetStatus()
-> rc/status 双重判定
-> 检查 tensor shape/size 和 finite logits
-> 直接 BF16 -> 目标 FP32 区
-> KV commit
-> reset slice/tensors
-> promise success
-> release BufferLease
若 Wait() 返回失败、状态为 DONE_CCU_EXCEPTION 或输出校验失败:不 commit KV;相关 sequence 标记 Unknown;buffer quarantine;触发故障控制器判定是否 drain/reset。
7. KV 事务实现
enum class KvTxnState : uint8_t { Reserved, Submitted, Committed, RolledBack, Unknown };
struct KvReservation {
calrt::KvSeqId seq_id;
uint32_t old_len;
uint32_t target_len;
uint64_t generation_id;
KvTxnState state;
};
class KvCoordinator {
public:
StatusOr<KvReservation> reserve(KvSeqId id, uint32_t target,
std::chrono::steady_clock::time_point deadline);
Status mark_submitted(KvReservation &txn);
Status commit(KvReservation &txn);
Status rollback(KvReservation &txn);
void mark_unknown(KvReservation &txn, std::string_view reason);
};
0.7.6 的 KvManager::Apply() 返回 bool 而无错误原因。短期规则:
- reserve 前先用
canAllocate(1);false 时进入条件变量等待,而非轮询。 - 每次唤醒只重试一次
Apply(),重新检查 deadline 和 shutdown。 - submit 前 Apply 的目标长度应由现有语义 runner 固化,不能同时改变长度口径与异步架构。
- submit 失败后若 Runtime 没执行图,可将软件 transaction rollback;硬件 KV 是否已经变化必须通过 vendor 契约或 golden 证明。
- submit 成功后任何错误默认
Unknown,不能声称请求级 rollback 已完成。
8. 错误分类和服务动作
| CalRT 错误 | 类别 | 当前 job | 新 admission | 自动动作 |
|---|---|---|---|---|
| 1 内存分配、14 越界、21 dtype、22 shape | Request/Artifact | 失败 | 同模型可继续或隔离产物 | 记录 manifest 差异;不 reset |
| 2 已注册、7 已配置、16 缺配置 | Lifecycle | 失败 | 暂停 | drain,核对 generation/config |
| 4 非法值、5 非法配置、19 非法 calbin、10/11/12 文件 | Artifact | 失败 | 拒绝该 generation | 回滚上个已知好模型 |
| 8 driver 不兼容、9 RT 过旧、23 依赖错误 | Compatibility | 失败 | 停止 | 人工修复,不自动重试 |
| 13 timeout、17 busy | Transient/Unknown | 失败或重试前排队 | 降低 admission | 有界退避;达到阈值 drain |
| 18 unavailable、24 map broken、25 soft reset 失败、28 crash | Device | 失败 | 立即停止 | quarantine -> reset/restart/escalate |
| 27 不支持、29 方向错误 | Programming | 失败 | 相关功能关闭 | 不 reset,修代码/配置 |
| 15/20 未预期/未知 | Unknown | 失败 | 暂停 | 保守按 device generation 污染处理 |
推荐 HTTP/gRPC 映射:输入越界 INVALID_ARGUMENT/400;KV/队列满 RESOURCE_EXHAUSTED/429;deadline DEADLINE_EXCEEDED/504;设备恢复中 UNAVAILABLE/503;内部 shape/calbin 错误 INTERNAL/500 并摘除 generation。
9. 超时、取消与恢复
9.1 三类 deadline
| deadline | 开始 | 结束 | 超时动作 |
|---|---|---|---|
| admission | 请求进入服务 | KV+buffer 取得 | 从等待队列移除,不碰设备 |
| queue | job prepared | infer() 成功 |
未提交可直接 rollback |
| execution | infer() 成功 |
Wait() 返回 |
不释放 DMA buffer,进入 Recovering |
0.7.6 没有已确认的 job cancel API。execution timeout 后不能 detach completion 再复用资源;默认步骤为停止 admission、保留所有 submitted buffer、等待短 grace period,然后走厂商确认过的 reset。没有 reset 契约前,最安全的最终恢复手段是进程退出,由 supervisor 重启。
9.2 设备恢复状态机
Ready -> Draining -> Recovering -> Configuring -> Warming -> Ready
| | | |
+--------> Failed <------+-------------+
恢复需要 generation+1,重新创建 Calbin/KvManager/buffer pool,跑固定 prefill+decode golden,成功后才接流。旧 completion 只能写旧 job promise,不得更新新 generation 的 KV 或指标 map。
10. 直接优化项
10.1 logits 转换
把:
std::vector<float> src(element_num);
for (...) src[i] = ggml_bf16_to_fp32(src_bf16[i]);
memcpy(dst, src.data() + base, dict_len * sizeof(float));
改为仅转换需要的 vocab 区:
for (size_t i = 0; i < dict_len; ++i) {
dst[i] = ggml_bf16_to_fp32(src_bf16[base + i]);
}
边界检查使用 base <= tensor_size && dict_len <= tensor_size - base,避免 base + dict_len 溢出。后续再基于 profiling 选择 SIMD 或 NPU top-k;后者是 C3。
10.2 exact submodel lookup
启动时枚举一次 GetAllModels(),验证 manifest 中每个 exact name 恰好出现一次。禁止运行时 substring fallback。缺失、多余或同 kind 重复都阻止 model generation 进入 Ready。
10.3 metrics guard
server.cpp 所有 cal_ctx->get_* 必须完全置于 #ifdef USE_CALRT 内。计时从共享 std::map<seq_id,double> 改到 InferenceResult.timings,请求结束即聚合到直方图,避免 seq id 复用污染。
11. 分阶段提交计划
| PR | 代码范围 | 行为变化 | 必须测试 |
|---|---|---|---|
| PR1 | error wrapper、所有返回码、metrics guard | 错误可见,不改调度 | 错误映射单测、现有 golden |
| PR2 | manifest exact lookup、tensor/layout 校验 | 启动 fail-fast | 缺模型/重名/shape/dtype 注入 |
| PR3 | KV 有界 admission、移除 abort | 资源不足变排队/错误 | deadline/cancel/容量用尽 |
| PR4 | BufferPool + move-only lease | 减少分配、明确所有权 | 地址稳定、归还重置、quarantine |
| PR5 | JobContext + submit/completion,in-flight=1 | 线程解耦、行为等价 | TSAN、取消、shutdown、1000 次循环 |
| PR6 | logits 直接转换 | CPU/内存优化 | bitwise/误差、边界、perf |
| PR7 | ping/pong runner,默认关闭 | 只形成实验能力 | 双 sequence 串扰、状态、fault |
| PR8 | 产品灰度 max_inflight=2 |
可能提高 overlap | 24h soak、P99、自动降级 |
每个 PR 都可独立回滚。PR7/PR8 不应与 batch calbin 接入混在同一变更中。
12. 故障注入表
| FI | 注入点 | 注入方式 | 预期结果 |
|---|---|---|---|
| FI01 | exact lookup | manifest 改错一个字符 | 启动失败,未 configure |
| FI02 | CSR | 使用不存在名称 | job prepare 失败,未 submit |
| FI03 | tensor slice | size=ByteSize+1 | 本地越界拒绝,未调用 Runtime |
| FI04 | KV capacity | mock canAllocate=false |
等待到 deadline,CPU 不忙等 |
| FI05 | submit | wrapper 返回 DeviceBusy |
job rollback,无 Wait |
| FI06 | submit | wrapper 返回 DeviceUnavailable |
admission 停止,进入恢复 |
| FI07 | Wait | 返回 Timeout |
buffer/KV quarantine,不交付 logits |
| FI08 | status | DONE_CCU_EXCEPTION |
即使 Wait=success 也失败并恢复 |
| FI09 | output | tensor size 小于 vocab | 检测失败,无越界写 |
| FI10 | output | 注入 NaN/Inf | generation 隔离,golden 复测 |
| FI11 | cancel before submit | 设置 cancel flag | 无设备提交,资源释放 |
| FI12 | cancel after submit | 完成前取消 | 仍 Wait,结果丢弃,资源最终释放 |
| FI13 | shutdown | 有 1 个 queued + 1 submitted | queued 取消,submitted drain |
| FI14 | stale completion | reset 后返回旧 job | 不写新 generation |
| FI15 | pool reset | 故意不 UndoSlice | 归还检查失败,slot quarantine |
13. 验收门槛
13.1 正确性
- 固定 100 个 prompt,NPU 改造前后 greedy token 序列完全一致;若参考实现允许浮点误差,另报 top-1 一致率、top-5 overlap 和 logits
max_abs/max_rel/cosine。 - 输入长度覆盖 1、2、15、16、17、127、128、1024、4095、4096、40959;40960 加 decode 必须被边界检查拒绝。
- 两个 seq id 交替 1000 token 无输出或 KV 串扰。
13.2 稳定性
- Phase A 24h,零进程 abort、零无限等待、零未解释 CCU exception。
- RSS、Runtime device memory 和 pool leased 数在 warmup 后无单调增长;终止后 leased=0。
- 取消/超时/错误注入 1000 次后仍可通过 golden,或进入明确 Failed 状态由 supervisor 恢复。
13.3 性能
- PR1-PR5 行为重构阶段,单请求 TTFT/TPOT P50 退化不超过 3%,P99 不超过 5%。
- logits 直接转换单独报告 CPU time 和分配次数;不得只看端到端平均值。
- Phase B 只有在 aggregate tokens/s 提升至少 10%、单请求 TPOT P99 退化不超过 5%、正确性和 24h 稳定性均通过时才默认启用。
14. 厂商阻断项
以下项不能由当前源码推断:
EnableParallelMode(true)的线程安全、最大 in-flight 和关闭时 drain 语义。- PING/PONG 是否真的重叠 H2D/compute/D2H,以及 engine slot 的复用时点。
- execution timeout 后安全 cancel 方法。
- submit 已成功但 Wait 失败时 KV 是否部分更新。
ResetCCU、ResetConfiguration、Reset对 DMA、KV、另一芯粒和 host buffer 的影响。
厂商交付必须包含最小 runner、时序图、错误注入结果和恢复后的 golden;口头说明不足以把这些能力从 C1/C3 提升为 C0。