5bf036bb81
Project CI / AI game creator shell Rust shard 4/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 3/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 2/4 (push) Has been cancelled
Project CI / AI game creator shell Rust smoke (push) Has been cancelled
Project CI / AI game creator shell Rust crates (push) Has been cancelled
Project CI / Backend tests (push) Has been cancelled
Project CI / Native shell tests (push) Has been cancelled
Project CI / Frontend tests (push) Has been cancelled
Project CI / Repository checks (push) Has been cancelled
Project CI / AI game creator shell web tests (push) Has been cancelled
Project CI / AI game creator shell Rust shard 1/4 (push) Has been cancelled
HTTP 请求取消后,在途计数原先无法释放;项目元数据和 External API 鉴权依赖完整 AppState,相同鉴权和追踪配置也散落在多个入口。本次集中装配这些依赖和横切能力,保持现有公开 API、权限、计费、幂等和事务规则。 ## 修改 - 91 个受保护路由集中应用鉴权;保留方法级 404/405/HEAD/Allow、公开入口、MCP、精选缓存头及 2/4 MiB 请求限制。 - 七个项目元数据入口改用缓存的 EditorProjectState,External/MCP 鉴权改用 ExternalApiAuthState;生产实现复用 SpacetimeClient,媒体修复维持原有上传和登记顺序。 - RAII 覆盖请求 Future 取消和 panic unwind 的计数清理;正常与降级服务复用 TraceLayer,指标采用 MatchedPath 模板及固定兜底。 - LLM 普通与流式调用增加跳过参数的异步 span,保持父上下文、流式回调、错误与重试行为;补充替代依赖测试并同步锁文件和文档。 ## 验证 | 验证面 | 结果 | | --- | --- | | platform-llm 完整本地回归 | 161 个单元测试、3 个集成测试通过;1 个真实 Provider 用例按原配置忽略 | | api-server 完整回归 | 执行时 1075 通过、12 失败、6 忽略;其中 1 个新增公开读取 fixture 断言已修正,14 个路由契约回归随后全部通过;剩余 11 个是下述既有 Windows 失败 | | 窄依赖、取消与追踪 | 元数据 owner/幂等/revision、鉴权及 MCP 错误传播、取消/panic/流式响应、追踪父子关系与敏感参数省略均通过 | | 实际本地服务 | 独立 SpacetimeDB 上 102/102 检查通过,两个动态项目 ID 的路由模板及请求 ID 日志核验 3/3 通过 | | 编译与边界 | api-server cargo check、AGC 锁文件下 platform-llm cargo check、rustfmt、编码、文档索引、DDD 与 diff 检查通过 | | 合入最新 master 后 | 后端源码及锁文件保持已测内容;再次通过 14 个路由契约测试、3 个 Provider 追踪测试及编码/文档/DDD/diff 检查 | 实际服务检查覆盖 health/ready、两账号登录、项目 CRUD、幂等重复、跨 owner 拒绝、revision 冲突、External/MCP 读取、Key 撤销及 404/405。使用既有 test 环境的本地 Router 拒绝 fixture,未调用真实付费 Provider;自建服务已关闭,原开发实例保留。 ## 已知测试限制 API 全量测试尚未全绿:11 个 wallet_refund_outbox 用例在 Windows 的目录同步处失败。其生产文件与变更前内容一致;标准库隔离复现确认 File::open(目录) 返回 OS 5,而普通文件写入、同步及 hard_link 正常。这个已有的目录持久化问题未混入本次重构,也未通过跳过或弱化相关断言掩盖。 --------- Co-authored-by: kdletters <61648117+kdletters@users.noreply.github.com> Reviewed-on: http://192.168.35.82/git/GenarrativeAI/Genarrative/pulls/425
333 lines
12 KiB
Rust
333 lines
12 KiB
Rust
use std::{
|
|
collections::BTreeMap,
|
|
io::{Read, Write},
|
|
net::TcpListener,
|
|
sync::{Arc, Mutex, mpsc},
|
|
thread,
|
|
time::Duration,
|
|
};
|
|
|
|
use tokio::sync::oneshot;
|
|
use tracing::{
|
|
Instrument, Subscriber,
|
|
field::{Field, Visit},
|
|
instrument::WithSubscriber,
|
|
span::{Attributes, Id, Record},
|
|
};
|
|
use tracing_subscriber::{Layer, layer::Context, prelude::*, registry::LookupSpan};
|
|
|
|
use super::{LlmClient, LlmConfig, LlmError, LlmProvider, LlmRunRequest};
|
|
|
|
const PRIVATE_INPUT: &str = "PRIVATE_MESSAGE_MUST_NOT_ENTER_SPAN";
|
|
const PRIVATE_KEY: &str = "PRIVATE_API_KEY_MUST_NOT_ENTER_SPAN";
|
|
|
|
#[derive(Clone, Debug, Default)]
|
|
struct CapturedSpan {
|
|
name: String,
|
|
parent: Option<u64>,
|
|
fields: BTreeMap<String, String>,
|
|
closed: bool,
|
|
}
|
|
|
|
#[derive(Clone, Default)]
|
|
struct Capture {
|
|
spans: Arc<Mutex<BTreeMap<u64, CapturedSpan>>>,
|
|
delta_parents: Arc<Mutex<Vec<u64>>>,
|
|
}
|
|
|
|
struct Fields<'a>(&'a mut BTreeMap<String, String>);
|
|
|
|
impl Visit for Fields<'_> {
|
|
fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
|
|
self.0.insert(field.name().into(), format!("{value:?}"));
|
|
}
|
|
|
|
fn record_str(&mut self, field: &Field, value: &str) {
|
|
self.0.insert(field.name().into(), value.into());
|
|
}
|
|
}
|
|
|
|
impl<S: Subscriber + for<'a> LookupSpan<'a>> Layer<S> for Capture {
|
|
fn on_new_span(&self, attributes: &Attributes<'_>, id: &Id, context: Context<'_, S>) {
|
|
let span = context.span(id).expect("registered span");
|
|
let mut captured = CapturedSpan {
|
|
name: attributes.metadata().name().into(),
|
|
parent: span.parent().map(|parent| parent.id().into_u64()),
|
|
..CapturedSpan::default()
|
|
};
|
|
attributes.record(&mut Fields(&mut captured.fields));
|
|
self.spans.lock().unwrap().insert(id.into_u64(), captured);
|
|
}
|
|
|
|
fn on_record(&self, id: &Id, values: &Record<'_>, _: Context<'_, S>) {
|
|
let mut spans = self.spans.lock().unwrap();
|
|
values.record(&mut Fields(
|
|
&mut spans.get_mut(&id.into_u64()).unwrap().fields,
|
|
));
|
|
}
|
|
|
|
fn on_close(&self, id: Id, _: Context<'_, S>) {
|
|
self.spans
|
|
.lock()
|
|
.unwrap()
|
|
.get_mut(&id.into_u64())
|
|
.unwrap()
|
|
.closed = true;
|
|
}
|
|
|
|
fn on_event(&self, event: &tracing::Event<'_>, context: Context<'_, S>) {
|
|
if event.metadata().target() == "llm_observability_test_delta" {
|
|
let parent = context
|
|
.event_span(event)
|
|
.expect("delta should have a parent");
|
|
self.delta_parents
|
|
.lock()
|
|
.unwrap()
|
|
.push(parent.id().into_u64());
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Capture {
|
|
fn provider_span(&self) -> (u64, CapturedSpan) {
|
|
let spans = self.spans.lock().unwrap();
|
|
let provider_spans = spans
|
|
.iter()
|
|
.filter(|(_, span)| span.name == "llm.request")
|
|
.collect::<Vec<_>>();
|
|
assert_eq!(provider_spans.len(), 1);
|
|
let (id, span) = provider_spans[0];
|
|
(*id, span.clone())
|
|
}
|
|
|
|
fn assert_completed(&self, operation: &str) {
|
|
let (_, span) = self.provider_span();
|
|
assert!(
|
|
span.closed,
|
|
"provider span must close after completion or cancellation"
|
|
);
|
|
assert_eq!(
|
|
span.fields.get("operation").map(String::as_str),
|
|
Some(operation)
|
|
);
|
|
assert_eq!(
|
|
span.fields.get("provider").map(String::as_str),
|
|
Some("openai_compatible")
|
|
);
|
|
assert_eq!(
|
|
span.fields.get("api_kind").map(String::as_str),
|
|
Some("openai_chat")
|
|
);
|
|
assert_eq!(
|
|
span.fields.get("model").map(String::as_str),
|
|
Some("requested-model")
|
|
);
|
|
let spans = self.spans.lock().unwrap();
|
|
assert_eq!(
|
|
spans[&span.parent.expect("HTTP/application parent")].name,
|
|
"test.request"
|
|
);
|
|
let captured = format!("{spans:?}");
|
|
assert!(!captured.contains(PRIVATE_INPUT));
|
|
assert!(!captured.contains(PRIVATE_KEY));
|
|
assert!(!span.fields.contains_key("self"));
|
|
assert!(!span.fields.contains_key("request"));
|
|
}
|
|
}
|
|
|
|
struct ProviderFixture {
|
|
client: LlmClient,
|
|
entered: oneshot::Receiver<()>,
|
|
release: mpsc::Sender<()>,
|
|
server: thread::JoinHandle<()>,
|
|
}
|
|
|
|
// 使用本地替代上游和通道控制响应时点,不依赖实际 Provider 或计时猜测。
|
|
fn provider_fixture(status: &str, content_type: &str, body: &str) -> ProviderFixture {
|
|
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
|
|
let address = listener.local_addr().unwrap();
|
|
let response = format!(
|
|
"HTTP/1.1 {status}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
|
|
body.len()
|
|
);
|
|
let (entered, entered_rx) = oneshot::channel();
|
|
let (release, release_rx) = mpsc::channel();
|
|
let server = thread::spawn(move || {
|
|
listener.set_nonblocking(true).unwrap();
|
|
let deadline = std::time::Instant::now() + Duration::from_secs(10);
|
|
let mut stream = loop {
|
|
match listener.accept() {
|
|
Ok((stream, _)) => break stream,
|
|
Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
|
|
assert!(
|
|
std::time::Instant::now() < deadline,
|
|
"request did not arrive"
|
|
);
|
|
thread::sleep(Duration::from_millis(5));
|
|
}
|
|
Err(error) => panic!("fixture accept: {error}"),
|
|
}
|
|
};
|
|
stream.set_nonblocking(false).unwrap();
|
|
stream
|
|
.set_read_timeout(Some(Duration::from_secs(5)))
|
|
.unwrap();
|
|
let mut bytes = Vec::new();
|
|
let mut chunk = [0; 4096];
|
|
loop {
|
|
let count = stream.read(&mut chunk).unwrap();
|
|
assert_ne!(count, 0);
|
|
bytes.extend_from_slice(&chunk[..count]);
|
|
if let Some(headers_end) = bytes.windows(4).position(|part| part == b"\r\n\r\n") {
|
|
let headers = String::from_utf8_lossy(&bytes[..headers_end]);
|
|
let length = headers
|
|
.lines()
|
|
.find_map(|line| {
|
|
let (name, value) = line.split_once(':')?;
|
|
name.eq_ignore_ascii_case("content-length")
|
|
.then(|| value.trim().parse::<usize>().unwrap())
|
|
})
|
|
.unwrap_or(0);
|
|
if bytes.len() >= headers_end + 4 + length {
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
entered.send(()).unwrap();
|
|
release_rx.recv_timeout(Duration::from_secs(10)).unwrap();
|
|
// 取消用例已关闭客户端连接,允许写回失败。
|
|
let _ = stream.write_all(response.as_bytes());
|
|
});
|
|
let config = LlmConfig::new(
|
|
LlmProvider::OpenAiCompatible,
|
|
format!("http://{address}"),
|
|
PRIVATE_KEY.into(),
|
|
"default-model".into(),
|
|
5_000,
|
|
0,
|
|
1,
|
|
)
|
|
.unwrap();
|
|
ProviderFixture {
|
|
client: LlmClient::new(config).unwrap(),
|
|
entered: entered_rx,
|
|
release,
|
|
server,
|
|
}
|
|
}
|
|
|
|
fn request() -> LlmRunRequest {
|
|
LlmRunRequest::single_turn("system", PRIVATE_INPUT)
|
|
.with_openai_chat()
|
|
.with_model("requested-model")
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn provider_span_covers_awaited_execution_and_keeps_parent_without_arguments() {
|
|
let fixture = provider_fixture(
|
|
"200 OK",
|
|
"application/json",
|
|
r#"{"choices":[{"message":{"content":"completed"},"finish_reason":"stop"}]}"#,
|
|
);
|
|
let capture = Capture::default();
|
|
let response =
|
|
async {
|
|
async {
|
|
let mut future = Box::pin(fixture.client.run(request()));
|
|
assert!(capture.spans.lock().unwrap().values().all(|span| span.name != "llm.request"));
|
|
tokio::select! {
|
|
response = &mut future => panic!("request returned before released: {response:?}"),
|
|
entered = fixture.entered => entered.unwrap(),
|
|
}
|
|
assert!(!capture.provider_span().1.closed);
|
|
// 挂起时不得把 Provider span 留在当前异步执行上下文。
|
|
assert_eq!(tracing::Span::current().metadata().unwrap().name(), "test.request");
|
|
fixture.release.send(()).unwrap();
|
|
future.await.unwrap()
|
|
}.instrument(tracing::info_span!("test.request")).await
|
|
}
|
|
.with_subscriber(tracing_subscriber::registry().with(capture.clone()))
|
|
.await;
|
|
fixture.server.join().unwrap();
|
|
assert_eq!(response.text, "completed");
|
|
capture.assert_completed("run");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn stream_callbacks_inherit_provider_span_and_keep_result() {
|
|
let fixture = provider_fixture(
|
|
"200 OK",
|
|
"text/event-stream",
|
|
"data: {\"choices\":[{\"delta\":{\"content\":\"hello\"},\"finish_reason\":null}]}\n\ndata: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\ndata: [DONE]\n\n",
|
|
);
|
|
let capture = Capture::default();
|
|
let mut deltas = Vec::new();
|
|
let response =
|
|
async {
|
|
async {
|
|
let mut future = Box::pin(fixture.client.stream_run(request(), |delta| {
|
|
tracing::info!(target: "llm_observability_test_delta", "delta received");
|
|
deltas.push(delta.delta_text.clone());
|
|
}));
|
|
tokio::select! {
|
|
response = &mut future => panic!("stream returned before released: {response:?}"),
|
|
entered = fixture.entered => entered.unwrap(),
|
|
}
|
|
fixture.release.send(()).unwrap();
|
|
future.await.unwrap()
|
|
}.instrument(tracing::info_span!("test.request")).await
|
|
}
|
|
.with_subscriber(tracing_subscriber::registry().with(capture.clone()))
|
|
.await;
|
|
fixture.server.join().unwrap();
|
|
assert_eq!(response.text, "hello");
|
|
assert_eq!(deltas.concat(), "hello");
|
|
capture.assert_completed("stream_run");
|
|
let id = capture.provider_span().0;
|
|
let parents = capture.delta_parents.lock().unwrap();
|
|
assert!(!parents.is_empty());
|
|
assert!(parents.iter().all(|parent| *parent == id));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn provider_span_preserves_upstream_errors_and_closes_on_cancellation() {
|
|
let fixture = provider_fixture(
|
|
"401 Unauthorized",
|
|
"application/json",
|
|
r#"{"error":{"message":"upstream-rejected"}}"#,
|
|
);
|
|
let capture = Capture::default();
|
|
let error = async {
|
|
async {
|
|
fixture.release.send(()).unwrap();
|
|
fixture.client.run(request()).await.unwrap_err()
|
|
}
|
|
.instrument(tracing::info_span!("test.request"))
|
|
.await
|
|
}
|
|
.with_subscriber(tracing_subscriber::registry().with(capture.clone()))
|
|
.await;
|
|
fixture.server.join().unwrap();
|
|
assert!(
|
|
matches!(error, LlmError::Upstream { status_code: 401, message } if message.contains("upstream-rejected"))
|
|
);
|
|
capture.assert_completed("run");
|
|
|
|
let fixture = provider_fixture("200 OK", "application/json", "{}");
|
|
let capture = Capture::default();
|
|
async {
|
|
async {
|
|
let mut future = Box::pin(fixture.client.run(request()));
|
|
tokio::select! {
|
|
response = &mut future => panic!("request returned before cancellation: {response:?}"),
|
|
entered = fixture.entered => entered.unwrap(),
|
|
}
|
|
assert!(!capture.provider_span().1.closed);
|
|
drop(future);
|
|
fixture.release.send(()).unwrap();
|
|
}.instrument(tracing::info_span!("test.request")).await
|
|
}.with_subscriber(tracing_subscriber::registry().with(capture.clone())).await;
|
|
fixture.server.join().unwrap();
|
|
capture.assert_completed("run");
|
|
}
|