Files
calculet-npu-research-archive/reports/Calculet-NPU-异步执行器与生产稳定性实施规格-20260802.md

19 KiB
Raw Permalink Blame History

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=1max_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 强制不变量

  1. DeviceExecutor 独占 VirtualDeviceCalbinKvManager
  2. 只有 submitter 线程调用 infer();只有 completion 线程对已提交 job 调用 Wait()
  3. 一个 BufferSlot 同一时刻最多属于一个 JobContext
  4. MapBuf() 后 host vector 在完成前不得 resize、move、swap 或释放。
  5. job 的 generation_id 必须与 pool/model generation 一致。
  6. KV reserve 成功不等于 KV commit;只有设备完成且输出通过校验才 commit。
  7. cancel 不释放底层资源,只阻止结果交付;仍需完成 Wait() 或执行已验证的 reset。
  8. 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_idseq_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() 检查顺序

  1. token_count > 0
  2. past_len + token_count <= max_seq,使用 64 位加法防溢出。
  3. 当前 calbin batch=1 时,seq_count == 1
  4. token 与 position 指针非空,长度一致;embed 路径单独检查元素数。
  5. tensor exact name 存在,dtype、rank、shape 和 ByteSize 与 manifest 一致。
  6. input_bytes <= tensor.ByteSize()slice offset/size 均在范围内。
  7. 两个 SetCsrByName() 都返回 CalrtSuccess
  8. 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 标记 Unknownbuffer 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/400KV/队列满 RESOURCE_EXHAUSTED/429deadline 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/completionin-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 等待到 deadlineCPU 不忙等
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 个 promptNPU 改造前后 greedy token 序列完全一致;若参考实现允许浮点误差,另报 top-1 一致率、top-5 overlap 和 logits max_abs/max_rel/cosine
  • 输入长度覆盖 1、2、15、16、17、127、128、1024、4095、4096、4095940960 加 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 是否部分更新。
  • ResetCCUResetConfigurationReset 对 DMA、KV、另一芯粒和 host buffer 的影响。

厂商交付必须包含最小 runner、时序图、错误注入结果和恢复后的 golden;口头说明不足以把这些能力从 C1/C3 提升为 C0。