详细讲解用Go(编排并发)+Rust(性能安全)在本地运行AI Assistant,完全避免云依赖和数据外泄,适合金融/医疗等敏感场景。
原文首发于 tamiz.pro。
人工智能领域的主流叙事一直由基于云、API 驱动的模型主导。虽然这种方法提供了可扩展性,但它带来了关键的延迟、对外部服务的依赖,以及有关数据泄露的重大隐私隐患。对于任务关键应用、金融分析或医疗系统而言,无法保证数据驻留和离线运行是不可接受的。
解决方案在于"本地优先"架构,其中 AI 智能体完全在本地或设备上运行。然而,构建此类系统不仅仅需要下载一个 LLM 权重文件;它要求一个强大的基础设施层,能够管理状态、内存安全和实时并发。本文探讨如何使用两种强大的编程语言来构建这种基础设施:Go 因其出色的并发原语和编排中的开发速度,以及 Rust 因其内存安全、零成本抽象和性能关键的推理执行而被选择。
我们将分析一个安全、本地优先 AI 智能体的架构,从概念模型到实现细节,重点关注编排层(Go)和执行层(Rust)之间的边界。
构建本地优先 AI 助手不仅仅是一个软件工程挑战;它是一个系统架构问题。核心矛盾在于灵活性(交换模型、调整提示和处理复杂工作流的能力)和性能/安全性(最小化延迟并防止内存损坏或数据泄露)之间。
为了解决这个问题,我们采用微内核架构:
编排器(Go):处理用户界面、API 网关、会话管理、工具调用和高级逻辑。Go 的 goroutine 允许它以最小的内存开销管理数千个并发智能体会话。
引擎(Rust):处理繁重的工作:模型加载、分词、推理和内存管理。Rust 确保关键路径(数据被处理且可能敏感的地方)不存在竞态条件、缓冲区溢出和未定义行为。
桥接(FFI/gRPC):语言之间的薄、严格类型化边界,确保跨越边界的数据被验证和高效序列化。
这种分离允许团队在 Go 中快速迭代智能体的逻辑(利用其庞大的 HTTP 服务器、数据库驱动程序和 UI 框架生态系统),同时在 Rust 中维护一个加固的、高性能的核心。
本地优先 AI 是计算密集型的。与云推理不同(你可以无限地水平扩展),本地推理受主机硬件限制(CPU/RAM/GPU)。Rust 被选为引擎层有三个主要原因:
AI 模型通常处理非结构化文本,可能具有对抗性。格式错误的输入可能导致 C/C++ 库中的缓冲区溢出。Rust 的所有权模型在编译时保证内存安全。对于本地优先智能体,这是一个安全特性,而不仅仅是性能特性。如果推理引擎由于内存错误而崩溃,它将关闭整个本地服务。Rust 防止了这一点。
WebAssembly(Wasm)和实时系统需要确定性行为。Rust 没有垃圾回收暂停,确保令牌生成延迟保持一致,这对于期望实时流式响应的聊天界面至关重要。
现代 Rust ML 生态系统,包括 Burn、Candle(由 Hugging Face 开发)和 tch-rs(PyTorch 绑定),提供了对最先进模型的访问。这些库针对 SIMD 指令和多线程矩阵乘法进行了优化,从消费级硬件中榨取最后一滴性能。
虽然 Rust 在原始计算方面更优越,但 Go 在并发模式和生态系统集成中表现出色。AI 智能体很少只是一个模型;它是一个系统,该系统:
Go 的 goroutine 模型非常适合这一点。每个用户会话可以被分配一个专用 goroutine,允许系统以小内存占用处理数千个并发用户。此外,Go 的标准库为构建 HTTP/2 服务器、处理 WebSocket 流以用于实时令牌传递以及与 SQL/NoSQL 数据库交互提供了出色的工具。
让我们可视化一个典型请求中的数据流:
如果数据被过度复制,进程间通信(IPC)或外部函数接口(FFI)调用的成本可能很高。为了缓解这一问题,我们尽可能使用零复制策略。
对于 gRPC,我们定义了一个最小化有效负载大小的架构。对于 FFI(Go 调用 Rust),我们谨慎使用 unsafe 块来传递预分配缓冲区的指针,避免 JSON 的序列化/反序列化开销。
我们将使用 Candle(由 Hugging Face 开发)来实现本地推理的简洁性和性能。Candle 设计成可嵌入的,在 CPU 和 GPU 上都表现良好。
rust-engine/
├── Cargo.toml
├── src/
│ ├── main.rs # 入口点,gRPC 服务器
│ ├── inference.rs # 模型加载和执行逻辑
│ ├── tokenizer.rs # 分词实用程序
│ └── lib.rs # FFI/gRPC 的导出
└── proto/
└── agent.proto # gRPC 定义
首先,我们定义 Go 和 Rust 之间的契约。这确保了类型安全和明确的期望。
// proto/agent.proto
syntax = "proto3";
package agent;
service InferenceService {
// 将令牌流式传输到编排器
rpc GenerateStream (Request) returns (stream Response);
// 健康检查和模型状态
rpc GetModelStatus (Empty) returns (Status);
}
message Request {
string prompt = 1;
int32 max_tokens = 2;
float temperature = 3;
repeated string history = 4; // 前面的对话轮次
}
message Response {
string token = 1;
bool is_end = 2;
}
message Status {
string model_name = 1;
bool is_ready = 2;
}
message Empty {}
我们使用 tonic 处理 gRPC,candle 处理推理。这个设置允许 Rust 二进制文件充当 Go 客户端连接到的独立服务。
// src/inference.rs
use candle::{Device, Tensor, DType};
use candle_transformers::generation::LogitsProcessor;
use candle_transformers::models::llama as llama_model;
use std::sync::Arc;
pub struct LlamaEngine {
model: Arc<llama_model::Llama>,
device: Device,
tokenizer: tokenizers::Tokenizer,
}
impl LlamaEngine {
pub fn new(model_path: &str, tokenizer_path: &str, device: Device) -> Result<Self, Box<dyn std::error::Error>> {
// 加载分词器
let tokenizer = tokenizers::Tokenizer::from_file(tokenizer_path)?;
// 加载模型权重
let model = llama_model::Llama::from_preferred(model_path, &device)?;
Ok(LlamaEngine {
model: Arc::new(model),
device,
tokenizer,
})
}
pub fn generate(&self, prompt: &str, max_tokens: usize, temperature: f64) -> Result<String, Box<dyn std::error::Error>> {
// 分词输入
let tokens = self.tokenizer.encode(prompt, true)?;
let tokens = tokens.get_ids();
// 准备输入张量
let input = Tensor::new(tokens, &self.device)?;
// 初始化 logits 处理器
let mut logits_processor = LogitsProcessor::new(temperature, Some(42), None);
let mut output_tokens = Vec::new();
let mut current_tokens = input.clone();
// 推理循环
for _ in 0..max_tokens {
let logits = self.model.forward(¤t_tokens, 0)?;
let logits = logits.squeeze(0)?;
let next_token = logits_processor.sample(&logits)?;
output_tokens.push(next_token);
// 若到达序列末尾则中断
if next_token == self.tokenizer.tokenizer.eos_token_id() {
break;
}
// 准备下一个输入
current_tokens = Tensor::new(&[next_token], &self.device)?;
}
// 解码输出 token
let decoded = self.tokenizer.decode(&output_tokens, true)?;
Ok(decoded)
}
}
// src/main.rs
use tonic::{transport::Server, Request, Response, Status};
use tokio_stream::{Stream, StreamExt};
mod inference;
mod agent {
include!(concat!(env!("OUT_DIR"), "/agent.rs"));
}
struct InferenceServiceImpl {
engine: inference::LlamaEngine,
}
#[tonic::async_trait]
impl agent::inference_service_server::InferenceService for InferenceServiceImpl {
async fn generate_stream(
&self,
request: Request<agent::Request>,
) -> Result<Response<Self::GenerateStreamStream>, Status> {
let req = request.into_inner();
let prompt = req.prompt;
let max_tokens = req.max_tokens as usize;
let temperature = req.temperature;
// 同步生成 token 以简化演示,
// 在生产环境中,你会流式传输中间 token
let result = match self.engine.generate(&prompt, max_tokens, temperature) {
Ok(text) => text,
Err(e) => return Err(Status::internal(e.to_string())),
};
// 对于流式传输,我们会分割结果并逐块输出
// 这里出于演示目的返回单个响应
Ok(Response::new(agent::Response {
token: result,
is_end: true,
}))
}
async fn get_model_status(
&self,
_request: Request<agent::Empty>,
) -> Result<Response<agent::Status>, Status> {
Ok(Response::new(agent::Status {
model_name: "Llama-2-7b".to_string(),
is_ready: true,
}))
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let engine = inference::LlamaEngine::new(
"./models/llama-2-7b",
"./models/tokenizer.json",
Device::Cpu,
)?;
let svc = agent::inference_service_server::InferenceServiceServer::new(
InferenceServiceImpl { engine }
);
Server::builder()
.add_service(svc)
.serve("127.0.0.1:50051".parse()?)
.await?;
Ok(())
}
Go 应用充当大脑的角色,管理对话的状态并与 Rust 引擎协调。
我们需要一种方式来管理多个用户。对于演示来说,简单的内存映射就足够了,但在生产环境中,你会使用 Redis 或 PostgreSQL。
// go-orchestrator/internal/session/store.go
package session
import (
"sync"
)
type Session struct {
ID string
History []string
}
type Store struct {
mu sync.RWMutex
sessions map[string]*Session
}
func NewStore() *Store {
return &Store{
sessions: make(map[string]*Session),
}
}
func (s *Store) GetOrCreate(id string) *Session {
s.mu.RLock()
sess, ok := s.sessions[id]
s.mu.RUnlock()
if !ok {
s.mu.Lock()
// 双重检查锁定
if sess, ok = s.sessions[id]; !ok {
sess = &Session{ID: id, History: make([]string, 0)}
s.sessions[id] = sess
}
s.mu.Unlock()
}
return sess
}
Go 应用需要连接到 Rust gRPC 服务器。我们使用生成的 protobuf 客户端。
// go-orchestrator/internal/ai/engine.go
package ai
import (
"context"
"log"
"github.com/yourname/agent/proto/agent"
"google.golang.org/grpc"
)
type RustEngineClient struct {
client agent.InferenceServiceClient
}
func NewRustEngineClient(addr string) (*RustEngineClient, error) {
conn, err := grpc.Dial(addr, grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
return nil, err
}
return &RustEngineClient{
client: agent.NewInferenceServiceClient(conn),
}, nil
}
func (c *RustEngineClient) Generate(ctx context.Context, prompt string, history []string, maxTokens int, temperature float32) (string, error) {
resp, err := c.client.GenerateStream(ctx, &agent.Request{
Prompt: prompt,
MaxTokens: int32(maxTokens),
Temperature: temperature,
History: history,
})
if err != nil {
return "", err
}
return resp.Token, nil
}
Go 的强项是处理许多并发连接。我们使用 WebSocket 服务器将响应流式传输回客户端。
// go-orchestrator/internal/handler/ai_handler.go
package handler
import (
"context"
"log"
"net/http"
"time"
"github.com/gorilla/websocket"
"github.com/yourname/agent/go-orchestrator/internal/ai"
"github.com/yourname/agent/go-orchestrator/internal/session"
)
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool {
return true
},
}
func AIHandler(engine *ai.RustEngineClient, store *session.Store) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
log.Println("Upgrade error:", err)
return
}
defer conn.Close()
// 为这个连接生成唯一的会话 ID
sessionID := "user-" + fmt.Sprintf("%d", time.Now().UnixNano())
sess := store.GetOrCreate(sessionID)
ctx := context.Background()
for {
_, msg, err := conn.ReadMessage()
if err != nil {
break
}
// 处理消息
response, err := engine.Generate(ctx, string(msg), sess.History, 500, 0.7)
if err != nil {
conn.WriteMessage(websocket.TextMessage, []byte("Error: "+err.Error()))
continue
}
// 发送响应
conn.WriteMessage(websocket.TextMessage, []byte(response))
// 更新历史记录
sess.History = append(sess.History, string(msg), response)
// 限制历史大小以管理内存
if len(sess.History) > 20 {
sess.History = sess.History[len(sess.History)-20:]
}
}
}
}
构建本地优先 AI 智能体引入了独特的安全挑战。与云 API 不同,在云 API 中你可以实施速率限制和 WAF,你的本地应用是防御的第一道防线。
即使在本地,用户也可以尝试提示注入。Go 编排器应该在将输入发送到 Rust 引擎之前对其进行清理。实现一个预处理层,该层过滤恶意模式或限制上下文窗口以防止 Rust 层中的缓冲区溢出尝试。
如果 Rust 引擎被编译为共享库(.so 或 .dylib)并链接到 Go 二进制文件中,Rust 中的漏洞可能会危及整个进程。为了缓解这一点,考虑将 Rust 引擎作为单独的进程运行(如 gRPC 示例所示)并通过 IPC 进行通信。这样,Rust 引擎中的崩溃不会导致 Go 编排器崩溃,你可以为 Rust 进程实施更严格的操作系统级沙箱化。
本地优先意味着数据保留在设备上。确保会话历史和 AI 处理的任何敏感数据在静止时加密。使用 Go 的 crypto 包在写入磁盘之前加密会话存储。
本地推理受硬件限制。为了使智能体响应迅速,你必须优化模型。
大型语言模型通常以 FP16 或 BF16 形式存储。量化到 INT8 甚至 INT4 可以将内存使用减少 2-4 倍,质量损失最小。Candle 库开箱即支持量化模型。加载量化模型以在消费级硬件上适配更大的上下文窗口。
如果你有多个用户,Rust 引擎可以对请求进行批处理。但是,对于本地优先的单用户智能体,批处理不太相关。对于多用户场景,Rust 引擎可以维护批处理队列,使用多线程矩阵乘法并行处理多个提示。
构建安全的本地优先 AI 智能体是一个复杂但有价值的工程挑战。通过利用 Go 的并发和生态系统,以及 Rust 的性能和内存安全性,你可以创建一个既灵活又强大的系统。
这种架构提供了:
隐私:数据永远不会离开设备。
性能:通过优化的 Rust 代码实现低延迟推理。
安全:内存安全的推理和隔离的进程。
可扩展性:Go 处理许多并发会话的能力。
随着 AI 从以云为中心的模式转向以分布式、边缘为中心的模式,这种混合方法将成为高保障应用的标准。可以先构建 Rust 推理引擎的原型,然后围绕它开发 Go 编排器,并逐步添加安全和优化层。
问:我可以使用 Python 作为编排层吗?答:可以,Python 经常用于 AI 开发。不过,对于高并发的本地服务,Go 或 Rust 更具优势。Python 的 GIL 限制了真正的并行处理能力,因此与 Go 的 goroutine 相比,它不太适合处理大量并发的 WebSocket 连接。
问:在本地优先应用中,如何处理模型更新?答:在你的 gRPC 协议中实现版本控制系统。当新模型发布时,Go 编排器可以触发下载新的权重(经过量化),并重新加载 Rust 引擎进程,而无须重启整个应用程序。
问:这种方法适用于移动设备吗?答:适用。相同的架构可以通过 Go(借助 gomobile)或 Rust(借助 Swift/Kotlin 绑定)适配 iOS/Android。关键是让推理引擎(Rust)保持原生运行,并将逻辑层(Go)设计为轻量层,或者同样采用原生实现。
如需采取进一步措施,你可以考虑屏蔽此人和/或举报滥用行为。