diff --git a/apps/ai-game-creator-shell/src-tauri/Cargo.lock b/apps/ai-game-creator-shell/src-tauri/Cargo.lock
index 21f1d6233..7a8d82d5f 100644
--- a/apps/ai-game-creator-shell/src-tauri/Cargo.lock
+++ b/apps/ai-game-creator-shell/src-tauri/Cargo.lock
@@ -3931,6 +3931,7 @@ dependencies = [
"serde",
"serde_json",
"tokio",
+ "tracing",
]
[[package]]
diff --git a/docs/project-memory/shared-memory/team-conventions.md b/docs/project-memory/shared-memory/team-conventions.md
index 84df79771..3e3c71e64 100644
--- a/docs/project-memory/shared-memory/team-conventions.md
+++ b/docs/project-memory/shared-memory/team-conventions.md
@@ -32,6 +32,8 @@
- 修改 `/api/external/v1` 时,同批更新 `docs/openapi/genarrative-external-v1.openapi.json` 与契约测试。
- 修改 SpacetimeDB schema 时遵守字段追加/default 约束,同步 migration、表目录、生成绑定,并运行 schema 检查;删除、改名、重排或改类型前先确认迁移计划。
- 日志不递归输出完整配置、应用状态或 provider client;新增字段默认不进入安全摘要。
+- HTTP 横切能力集中在 Axum/Tower 中间件:正常与降级路由复用追踪层;指标与 trace 使用 `MatchedPath` 模板及固定兜底,不把请求 ID、实际资源 ID 或 query 放入指标标签。在途请求通过 RAII guard 覆盖 Future 取消与 panic unwind;请求执行和响应体存活分别计量,不能把 handler 耗时当作 SSE 全生命周期。
+- 业务依赖在组合根显式装配,Axum `FromRef` 只抽取可浅拷贝的窄能力。项目元数据与 External API 鉴权不持有完整 `AppState`,测试经相同接口注入替代依赖。集中鉴权仍保留方法级 fallback、公开入口、MCP 和 body limit 顺序;Provider span 跳过完整参数,不隐藏计费、重试、幂等或事务规则。
- 中文文案、注释和文档保持 UTF-8,优先局部补丁,不擅自翻译成英文。
## 文档生命周期
diff --git a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md
index dfab9a03e..a65dd79ef 100644
--- a/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md
+++ b/docs/【后端架构】server-rs与SpacetimeDB数据契约-2026-05-15.md
@@ -47,6 +47,15 @@ SpacetimeDB 版本口径:当前 Rust crate `spacetimedb`、`spacetimedb-sdk`
npm run check:server-rs-ddd
```
+## 依赖装配与横切能力
+
+- 启动入口负责构造共享依赖,业务入口通过 Axum `State` / `FromRef` 获取所需能力。项目元数据链路与 External API 鉴权使用窄状态;窄状态不得持有完整 `AppState`、通过 `Deref` 暴露完整配置,或提供 `root_state()` 逃逸。仅在需要替换或隔离测试的外部能力处抽取接口,生产实现仍复用现有 `spacetime-client` facade,不新建数据库访问通道。
+- 后台、External API 和编辑器路由按同一种身份策略聚合鉴权,公开登录、公开文档、公开精选读取及 MCP 保持独立边界。路由、HTTP 方法、404/405、请求大小限制、授权和错误响应合同保持不变;后台 member 的实时权限校验和资源 owner 校验继续由现有权威逻辑执行。
+- HTTP 观测使用统一 Tower 追踪层与 `MatchedPath` 模板;请求计数通过 RAII 覆盖取消与 unwind,响应体存活单独统计。具体指标合同见开发运维文档。
+- Provider 函数级追踪采用现有 tracing 能力,显式列出 provider、operation、model 等非秘密字段,跳过完整参数、配置、凭据和消息正文;异步 span 覆盖实际执行与等待,不能只记录 Future 创建。装饰与追踪保持原始成功值、错误、重试次数、流式回调和取消语义。
+- 扣费、退款、幂等、资源登记、重试判定与事务边界保持显式业务流程;不引入自动扫描 IOC 容器、通用 AOP 框架或跨外部副作用的隐式事务。
+- 验收覆盖:所有受保护路由的匿名拒绝及公开入口可访问;原有 404/405 和 body limit;窄依赖可独立构造和替换;Provider 成功/失败、异步追踪归属与脱敏;已有 API、计费和幂等回归。真实本地服务 smoke 必须记录所用数据库及依赖可用性,不能用单元测试代替运行时证据。
+
## `spacetime-client` mapper 组织
`spacetime-client` 的 Cargo `lib.path` 指向 `src/active.rs`,现役 mapper 聚合入口是 `src/active/mapper.rs`;原旧 facade 和 mapper 已删除。
diff --git a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md
index 92e4d4fc6..2b884aad3 100644
--- a/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md
+++ b/docs/【开发运维】本地开发验证与生产运维-2026-05-15.md
@@ -808,13 +808,15 @@ OpenTelemetry 现阶段默认开启 OTLP traces / metrics / logs,但本地日
- 应用日志按进程查看:父 API 使用 `journalctl -u genarrative-api.service`,独立 BgFilter worker 使用 `journalctl -u genarrative-bgfilter-worker.service`;Nginx 日志仍写文件。日志等级继续用 `GENARRATIVE_API_LOG` / `RUST_LOG` 控制,例如 `info,tower_http=info,spacetime_client=info`。
- debug exporter / Rider 转发都会同时接收 traces、metrics 和 logs。
- api-server 会随 metrics 发送进程级指标:`process.memory.usage`、`process.memory.virtual`、`process.cpu.time`、`genarrative.process.cpu.usage_percent`、`process.thread.count`、`genarrative.process.memory.private`;Windows 额外发送 `process.windows.handle.count`,Linux 额外发送 `process.unix.file_descriptor.count`。这些指标只描述当前进程,不携带请求、用户或作品 label。
-- HTTP 运行态补充发送 `genarrative.http.server.response_bodies.in_flight` 与 `genarrative.http.server.request_permits.available`,后者带低基数 `pool=default|gallery|detail|admin` label,用于区分业务 handler / 背压 permit 是否仍被占用;拼图广场热点缓存补充发送 `genarrative.puzzle_gallery.cache.*` 指标,记录 fresh hit、stale hit、未命中、后台刷新开始 / 失败、重建耗时和预序列化 data JSON 字节数。
+- HTTP 请求在途指标 `http.server.active_requests` 从进入观测中间件计数,到生成 `Response` 时结束;使用 RAII guard 保证正常返回、Future 被取消释放和 panic unwind 时都按原始 method / route 标签且仅递减一次。`http.server.request.duration` 采用同一执行区间,仅记录已生成响应的请求,不表示响应已发送到客户端,也不表示 SSE 已结束。响应体存活另由 `genarrative.http.server.response_bodies.in_flight` 统计;背压 permit 由 `genarrative.http.server.request_permits.available{pool=default|admin}` 统计。
+- 正常服务与 SpacetimeDB 不可用时的降级路由复用同一 HTTP `TraceLayer` 构造函数。请求上下文在追踪层外侧,错误归一化和响应头回写在追踪层内侧,保证提前拒绝和降级响应仍记录同一 request ID 与最终状态。
+- `platform-llm` 的普通与流式调用统一生成 `llm.request` 子 span,字段白名单为 `provider`、`operation`、`api_kind`、`model`,继承调用方当前 span;范围覆盖响应读取、流式回调和异步等待,取消释放 Future 时结束。完整 client/config、API Key、请求/响应正文不进入该 span;既有 Provider 错误和重试策略保持原样,不能用遥测包装吞掉失败或自动重放。
- 外部 API 失败统一发送 OTLP 并落库。当前 VectorEngine 图片生成 / 编辑失败由 `platform-image` provider 输出结构化日志字段,字段包括 provider、endpoint、failure_stage、status、source、source_chain、source_chain_depth、timeout、retryable、latency_ms、prompt_chars、reference_image_count、实际 provider `image_model`、request_params 和 raw_excerpt;发生模型回退时另带 `fallback_from_model` / `fallback_to_model`。图片编辑请求参数日志还会带 reference_image_bytes_total,并在 request_params.referenceImages 中记录每个 multipart `image` part 的 fileName、mimeType 和 bytes,不记录 API key 或原始图片 bytes;`api-server` 再记录指标 `genarrative.external_api.failures{provider,failure_stage,status_class,retryable}`,并写入 `tracking_event`,`event_key = external_api_call_failure`、`module_key = external-api`、`scope_kind = module`、`scope_id = provider`。调用方能拿到身份上下文时,失败事件还会在行级 `user_id` / `owner_user_id` / `profile_id` 和 `metadata_json.userId` / `metadata_json.profileId` / `metadata_json.requestId` / `metadata_json.errorSource` 中记录触发者、草稿 / 作品作用域、请求标识和传输错误链。排障时先按 provider / failureStage / imageModel 聚合,再下钻 userId / profileId,最后结合 request 日志、errorSource 和上游响应 excerpt 判断是模型不可用、限流、超时、解析失败还是未返回图片。
- OSS 平台适配器也输出结构化日志,覆盖 `sign_post_object`、`sign_get_object_url`、`head_object` 和 `put_object`。排查资产签名、上传或确认失败时,先按 `provider=aliyun-oss` 与 `operation` 过滤,再看 `object_key` / `key_prefix`、`status`、`status_class`、`error_kind`、`content_length`、`content_type` 和 `elapsed_ms`;角色动画逐帧额外按 `frame_index`、`operation=source_put|final_put|final_head`、`attempt/max_attempts`、`will_retry`、`oss_code` 和 `oss_request_id` 对齐同一对象的请求尝试。`请求 OSS 失败` 时,`timeout/connect/transport=true` 表示传输类失败,OSS PutObject 的 `status=400, oss_code=RequestTimeout, timeout=true`、`status=429` 或 `500–599` 表示暂时性失败,PUT 的 `status=400`、`oss_code` 为空且 `timeout=true` 或 `transport=true`(message 含「错误响应体读取失败」,即 400 错误体读取超时/断流)也会重试;除这两类例外外,其他 400、401/403/404、配置、URL 和签名错误是确定性失败,不会重试。最终帧 HEAD 失败只会重试 HEAD,不会重复 PUT。日志不得包含 AccessKey、policy、signature、Authorization header、完整 signed URL 或 OSS 错误响应体;`oss_request_id` 只用于关联 OSS 服务端排障。排查 generated 图片重复下载时,先确认前端输入是否为 `/generated-*` legacy path 或可归一化的 `https://*.oss-*.aliyuncs.com/generated-*`;正确链路应先调 `/api/assets/read-url`,再由浏览器请求 signed URL,且同一路径、同一 `refreshKey` 版本和未临近过期的 signed URL 应复用。新上传 generated 私有对象应带 `Cache-Control: public, max-age=31536000, immutable`;旧对象若只有 `ETag` / `Last-Modified`,浏览器会走 304 协商缓存而不是长期强缓存,可通过刷新 OSS 元数据或 CDN 配置补齐。
- SpacetimeDB 观测分为两类:procedure / reducer 调用继续用 `genarrative.spacetime.procedure.*`,订阅本地 cache 读使用 `genarrative.spacetime.read.*`。`read=list_puzzle_gallery` 表示拼图广场当前从 `puzzle_gallery_card_view` 本地 cache 读取,不再每个 HTTP 请求调用 `list_puzzle_gallery` procedure。
- 本地 Windows 直连压测的内存高水位要结合 K6 VU / 连接数解释。250 RPS 下过高 `PREALLOCATED_VUS` 可能让 300 个本地 Established 连接把 `api-server` private memory 瞬时推到 GB 级,且 `/healthz` 小响应也能复现;若压测结束后回落、`response_bodies.in_flight` 和背压 permit 未显示业务积压,应优先按连接 / 发送链路高水位处理,而不是判断为 SpacetimeDB 或 JSON 缓存泄漏。
- Rider 的 Logs 面板只展示 log event 自身字段,不会自动展开父 span 的全部 attributes;请求完成日志会直接带 `request_id`、`http.request.method`、`http.route`、`url.scheme`、`url.path`、`http.response.status_code`、`status_class`、`latency_ms` 和 `slow_request`,完整链路继续到 Traces 面板按 trace/span 查看。
-- 指标 label 只允许低基数字段:HTTP 使用 `method`、`route`、`status_class`,SpacetimeDB 调用使用 `procedure`、`status_class`;`request_id` 只进入 trace/log attribute,不进入 metric label。
+- 指标 label 只允许低基数字段:HTTP 使用 `http.request.method`、`http.route`、`status_class`,SpacetimeDB 调用使用 `procedure`、`status_class`;`request_id` 只进入 trace/log attribute,不进入 metric label。HTTP trace、完成日志与指标共用 Axum `MatchedPath` 路由模板,例如 `/api/editor/projects/{project_id}`;未匹配路由或无路由模板的降级请求只使用 `/api/*`、`/admin/api/*`、`other` 三种固定兜底,不把实际 ID 或 query 写入 route 标签。原来按 `/api/*` 聚合的监控查询应改为按路由模板汇总。
常见外部服务变量:
diff --git a/server-rs/Cargo.lock b/server-rs/Cargo.lock
index 5de6bdfbe..614e6da6b 100644
--- a/server-rs/Cargo.lock
+++ b/server-rs/Cargo.lock
@@ -4101,6 +4101,8 @@ dependencies = [
"serde",
"serde_json",
"tokio",
+ "tracing",
+ "tracing-subscriber",
]
[[package]]
diff --git a/server-rs/crates/api-server/src/app.rs b/server-rs/crates/api-server/src/app.rs
index 4c34bda86..ded1dafe0 100644
--- a/server-rs/crates/api-server/src/app.rs
+++ b/server-rs/crates/api-server/src/app.rs
@@ -1,3 +1,5 @@
+use std::time::Duration;
+
use axum::{
Router,
body::Body,
@@ -9,8 +11,8 @@ use axum::{
};
use serde_json::json;
use tower_http::{
- classify::ServerErrorsFailureClass,
- trace::{DefaultOnRequest, TraceLayer},
+ classify::{ServerErrorsAsFailures, ServerErrorsFailureClass, SharedClassifier},
+ trace::{DefaultOnBodyChunk, DefaultOnEos, DefaultOnRequest, TraceLayer},
};
use tracing::{Level, Span, error, info_span};
@@ -85,66 +87,7 @@ pub fn build_router(state: AppState) -> Router {
state.clone(),
record_http_observability,
))
- // 当前阶段先统一挂接 HTTP tracing,后续 request_id、响应头与错误中间件继续在这里扩展。
- .layer(
- TraceLayer::new_for_http()
- .make_span_with(|request: &Request
| {
- let request_id =
- resolve_request_id(request).unwrap_or_else(|| "unknown".to_string());
- let route = crate::telemetry::observability_route(request.uri().path());
- let scheme = crate::telemetry::resolve_request_scheme(request.headers());
- let span_name = format!("{} {}", request.method(), route);
-
- info_span!(
- "http.request",
- otel.kind = "server",
- otel.name = %span_name,
- otel.status_code = tracing::field::Empty,
- http.response.status_code = tracing::field::Empty,
- method = %request.method(),
- http.request.method = %request.method(),
- http.route = %route,
- url.scheme = %scheme,
- url.path = %request.uri().path(),
- request_id = %request_id,
- status = tracing::field::Empty,
- latency_ms = tracing::field::Empty,
- )
- })
- .on_request(DefaultOnRequest::new().level(Level::INFO))
- .on_response(
- |response: &axum::response::Response,
- latency: std::time::Duration,
- span: &Span| {
- let latency_ms = latency.as_millis().min(u64::MAX as u128) as u64;
- let status = response.status().as_u16();
- span.record("status", status);
- span.record("http.response.status_code", status);
- span.record(
- "otel.status_code",
- if response.status().is_server_error() {
- "ERROR"
- } else {
- "OK"
- },
- );
- span.record("latency_ms", latency_ms);
- },
- )
- .on_failure(
- |failure: ServerErrorsFailureClass,
- latency: std::time::Duration,
- span: &Span| {
- let latency_ms = latency.as_millis().min(u64::MAX as u128) as u64;
- error!(
- parent: span,
- latency_ms,
- failure = %failure,
- "http request failed"
- );
- },
- ),
- )
+ .layer(http_trace_layer())
// request_id 中间件先进入请求链,确保后续 tracing、错误处理和响应头层都能复用同一份请求标识。
.layer(middleware::from_fn(attach_request_context))
.with_state(state)
@@ -159,68 +102,72 @@ pub fn build_spacetime_unavailable_router(message: String) -> Router {
// 依赖不可用模式不挂业务 state,统一返回 503,并继续保留 request_id / API 版本 / 耗时响应头。
.layer(middleware::from_fn(normalize_error_response))
.layer(middleware::from_fn(propagate_request_id_header))
- .layer(
- TraceLayer::new_for_http()
- .make_span_with(|request: &Request| {
- let request_id =
- resolve_request_id(request).unwrap_or_else(|| "unknown".to_string());
- let route = crate::telemetry::observability_route(request.uri().path());
- let scheme = crate::telemetry::resolve_request_scheme(request.headers());
- let span_name = format!("{} {}", request.method(), route);
-
- info_span!(
- "http.request",
- otel.kind = "server",
- otel.name = %span_name,
- otel.status_code = tracing::field::Empty,
- http.response.status_code = tracing::field::Empty,
- method = %request.method(),
- http.request.method = %request.method(),
- http.route = %route,
- url.scheme = %scheme,
- url.path = %request.uri().path(),
- request_id = %request_id,
- status = tracing::field::Empty,
- latency_ms = tracing::field::Empty,
- )
- })
- .on_request(DefaultOnRequest::new().level(Level::INFO))
- .on_response(
- |response: &axum::response::Response,
- latency: std::time::Duration,
- span: &Span| {
- let latency_ms = latency.as_millis().min(u64::MAX as u128) as u64;
- let status = response.status().as_u16();
- span.record("status", status);
- span.record("http.response.status_code", status);
- span.record(
- "otel.status_code",
- if response.status().is_server_error() {
- "ERROR"
- } else {
- "OK"
- },
- );
- span.record("latency_ms", latency_ms);
- },
- )
- .on_failure(
- |failure: ServerErrorsFailureClass,
- latency: std::time::Duration,
- span: &Span| {
- let latency_ms = latency.as_millis().min(u64::MAX as u128) as u64;
- error!(
- parent: span,
- latency_ms,
- failure = %failure,
- "http request failed"
- );
- },
- ),
- )
+ .layer(http_trace_layer())
.layer(middleware::from_fn(attach_request_context))
}
+type HttpTraceLayer = TraceLayer<
+ SharedClassifier,
+ fn(&Request) -> Span,
+ DefaultOnRequest,
+ fn(&Response, Duration, &Span),
+ DefaultOnBodyChunk,
+ DefaultOnEos,
+ fn(ServerErrorsFailureClass, Duration, &Span),
+>;
+
+fn http_trace_layer() -> HttpTraceLayer {
+ TraceLayer::new_for_http()
+ .make_span_with(make_http_request_span as fn(&Request) -> Span)
+ .on_request(DefaultOnRequest::new().level(Level::INFO))
+ .on_response(record_http_response_span as fn(&Response, Duration, &Span))
+ .on_failure(record_http_failure as fn(ServerErrorsFailureClass, Duration, &Span))
+}
+
+fn make_http_request_span(request: &Request) -> Span {
+ let request_id = resolve_request_id(request).unwrap_or_else(|| "unknown".to_string());
+ let route = crate::telemetry::observability_route(request);
+ let scheme = crate::telemetry::resolve_request_scheme(request.headers());
+ let span_name = format!("{} {}", request.method(), route);
+
+ info_span!(
+ "http.request",
+ otel.kind = "server",
+ otel.name = %span_name,
+ otel.status_code = tracing::field::Empty,
+ http.response.status_code = tracing::field::Empty,
+ method = %request.method(),
+ http.request.method = %request.method(),
+ http.route = %route,
+ url.scheme = %scheme,
+ url.path = %request.uri().path(),
+ request_id = %request_id,
+ status = tracing::field::Empty,
+ latency_ms = tracing::field::Empty,
+ )
+}
+
+fn record_http_response_span(response: &Response, latency: Duration, span: &Span) {
+ let latency_ms = latency.as_millis().min(u64::MAX as u128) as u64;
+ let status = response.status().as_u16();
+ span.record("status", status);
+ span.record("http.response.status_code", status);
+ span.record(
+ "otel.status_code",
+ if response.status().is_server_error() {
+ "ERROR"
+ } else {
+ "OK"
+ },
+ );
+ span.record("latency_ms", latency_ms);
+}
+
+fn record_http_failure(failure: ServerErrorsFailureClass, latency: Duration, span: &Span) {
+ let latency_ms = latency.as_millis().min(u64::MAX as u128) as u64;
+ error!(parent: span, latency_ms, failure = %failure, "http request failed");
+}
+
#[derive(Clone, Debug)]
struct SpacetimeUnavailableState {
message: std::sync::Arc,
@@ -308,6 +255,201 @@ mod tests {
const TEST_PASSWORD: &str = "secret123";
const INTERNAL_TEST_SECRET: &str = "test-internal-secret";
+ mod http_tracing {
+ use std::{
+ collections::BTreeMap,
+ sync::{Arc, Mutex},
+ };
+
+ use axum::extract::FromRef;
+ use tracing::{
+ Metadata, Subscriber,
+ field::{Field, Visit},
+ instrument::WithSubscriber,
+ span::{Attributes, Id, Record},
+ };
+
+ use super::*;
+ use crate::state::{BackpressureState, HttpRequestPermitPoolKind};
+
+ #[derive(Clone, Default)]
+ struct HttpSpanCapture(Arc>>>);
+
+ struct SpanFields<'a>(&'a mut BTreeMap);
+
+ impl Visit for SpanFields<'_> {
+ fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
+ self.0
+ .insert(field.name().to_string(), format!("{value:?}"));
+ }
+
+ fn record_str(&mut self, field: &Field, value: &str) {
+ self.0.insert(field.name().to_string(), value.to_string());
+ }
+ }
+
+ impl Subscriber for HttpSpanCapture {
+ fn enabled(&self, metadata: &Metadata<'_>) -> bool {
+ metadata.is_span() && metadata.name() == "http.request"
+ }
+
+ fn new_span(&self, attributes: &Attributes<'_>) -> Id {
+ let mut spans = self.0.lock().expect("span capture should lock");
+ let mut fields = BTreeMap::new();
+ attributes.record(&mut SpanFields(&mut fields));
+ spans.push(fields);
+ Id::from_u64(spans.len() as u64)
+ }
+
+ fn record(&self, span: &Id, values: &Record<'_>) {
+ let mut spans = self.0.lock().expect("span capture should lock");
+ values.record(&mut SpanFields(&mut spans[span.into_u64() as usize - 1]));
+ }
+
+ fn record_follows_from(&self, _: &Id, _: &Id) {}
+ fn event(&self, _: &tracing::Event<'_>) {}
+ fn enter(&self, _: &Id) {}
+ fn exit(&self, _: &Id) {}
+ }
+
+ async fn assert_rejection_observed(
+ app: Router,
+ request: Request,
+ expected_status: StatusCode,
+ expected_route: &str,
+ expected_code: &str,
+ ) -> (Value, String) {
+ let expected_request_id = request.headers().get("x-request-id").cloned();
+ let path = request.uri().path().to_string();
+ let capture = HttpSpanCapture::default();
+ let response = app
+ .oneshot(request)
+ .with_subscriber(capture.clone())
+ .await
+ .expect("rejected request should complete");
+
+ assert_eq!(response.status(), expected_status);
+ let request_id = response.headers()["x-request-id"]
+ .to_str()
+ .expect("request id should be a header string")
+ .to_string();
+ assert!(!request_id.is_empty());
+ assert_ne!(request_id, "unknown");
+ if let Some(expected_request_id) = expected_request_id {
+ assert_eq!(response.headers()["x-request-id"], expected_request_id);
+ }
+ for name in ["x-api-version", "x-route-version"] {
+ assert_eq!(response.headers()[name], shared_contracts::api::API_VERSION);
+ }
+ assert!(
+ response.headers()["x-response-time-ms"]
+ .to_str()
+ .unwrap()
+ .parse::()
+ .is_ok()
+ );
+ if expected_status == StatusCode::TOO_MANY_REQUESTS {
+ assert_eq!(response.headers()["retry-after"], "1");
+ }
+ let payload = read_json_response(response).await;
+ assert_eq!(payload["error"]["code"], expected_code);
+
+ let spans = capture.0.lock().expect("span capture should lock");
+ assert_eq!(
+ spans.len(),
+ 1,
+ "each rejected request should have one HTTP span"
+ );
+ let fields = &spans[0];
+ assert_eq!(fields["request_id"], request_id);
+ assert_eq!(fields["http.route"], expected_route);
+ assert_eq!(fields["otel.name"], format!("GET {expected_route}"));
+ assert_eq!(fields["url.path"], path);
+ assert_eq!(fields["http.request.method"], "GET");
+ assert_eq!(fields["status"], expected_status.as_u16().to_string());
+ assert_eq!(
+ fields["http.response.status_code"],
+ expected_status.as_u16().to_string()
+ );
+ assert_eq!(
+ fields["otel.status_code"],
+ if expected_status.is_server_error() {
+ "ERROR"
+ } else {
+ "OK"
+ }
+ );
+ assert!(fields["latency_ms"].parse::().is_ok());
+ (payload, request_id)
+ }
+
+ #[tokio::test]
+ async fn auth_rejections_share_the_matched_route_template() {
+ let app =
+ build_router(AppState::new(AppConfig::default()).expect("state should build"));
+ for project_id in ["project-one", "project-two"] {
+ assert_rejection_observed(
+ app.clone(),
+ Request::builder()
+ .uri(format!(
+ "/api/editor/projects/{project_id}/agent-conversations"
+ ))
+ .header("x-request-id", format!("req-trace-{project_id}"))
+ .body(Body::empty())
+ .expect("request should build"),
+ StatusCode::UNAUTHORIZED,
+ "/api/editor/projects/{project_id}/agent-conversations",
+ "UNAUTHORIZED",
+ )
+ .await;
+ }
+ }
+
+ #[tokio::test]
+ async fn backpressure_rejection_keeps_generated_context_and_headers() {
+ let config = AppConfig {
+ max_concurrent_requests: Some(1),
+ ..AppConfig::default()
+ };
+ let state = AppState::new(config).expect("state should build");
+ let (_, pool) = BackpressureState::from_ref(&state)
+ .request_permit_pool(HttpRequestPermitPoolKind::Default)
+ .expect("default request pool should exist");
+ let _held_permit = pool
+ .try_acquire_owned()
+ .expect("pool should have one permit");
+
+ let (payload, request_id) = assert_rejection_observed(
+ build_router(state),
+ Request::builder()
+ .uri("/api/editor/projects/project-one/agent-conversations")
+ .body(Body::empty())
+ .expect("request should build"),
+ StatusCode::TOO_MANY_REQUESTS,
+ "/api/editor/projects/{project_id}/agent-conversations",
+ "TOO_MANY_REQUESTS",
+ )
+ .await;
+ assert_eq!(payload["meta"]["requestId"], request_id);
+ }
+
+ #[tokio::test]
+ async fn unavailable_router_rejection_keeps_generated_context_and_headers() {
+ let (payload, request_id) = assert_rejection_observed(
+ build_spacetime_unavailable_router("test unavailable".to_string()),
+ Request::builder()
+ .uri("/api/auth/login-options")
+ .body(Body::empty())
+ .expect("request should build"),
+ StatusCode::SERVICE_UNAVAILABLE,
+ "/api/*",
+ "SERVICE_UNAVAILABLE",
+ )
+ .await;
+ assert_eq!(payload["meta"]["requestId"], request_id);
+ }
+ }
+
async fn seed_phone_user_with_password(
state: &AppState,
phone_number: &str,
diff --git a/server-rs/crates/api-server/src/editor_project.rs b/server-rs/crates/api-server/src/editor_project.rs
index b880b06e6..100d89b41 100644
--- a/server-rs/crates/api-server/src/editor_project.rs
+++ b/server-rs/crates/api-server/src/editor_project.rs
@@ -12,7 +12,7 @@ use std::{
use axum::{
Json,
- extract::{Extension, Path, Query, State, rejection::JsonRejection},
+ extract::{Extension, FromRef, Path, Query, State, rejection::JsonRejection},
http::{HeaderMap, HeaderValue, StatusCode},
};
use module_assets::{
@@ -122,7 +122,7 @@ use crate::{
platform_errors::map_oss_error,
prompt::editor_scene::build_editor_scene_prompt,
request_context::RequestContext,
- state::AppState,
+ state::{AppState, EditorMediaStorageState, EditorProjectState},
work_author::{ORPHAN_WORK_AUTHOR_PUBLIC_USER_CODE, ORPHAN_WORK_OWNER_USER_ID},
};
@@ -1893,19 +1893,19 @@ pub struct EditorAssetPayload {
}
pub async fn load_recent_editor_project(
- State(state): State,
+ State(state): State,
Extension(request_context): Extension,
Extension(authenticated): Extension,
) -> Result, AppError> {
let owner_user_id = authenticated.claims().user_id().to_string();
let project = state
- .spacetime_client()
+ .projects()
.get_recent_editor_project(owner_user_id)
.await
.map_err(map_editor_project_error)?;
let project = match project {
Some(project) => Some(editor_project_payload_from_record(
- repair_editor_project_record_inline_media(&state, project).await,
+ state.repair_project_media(project).await,
)),
None => None,
};
@@ -1931,20 +1931,20 @@ pub async fn get_editor_generation_pricing(
}
pub async fn list_editor_projects(
- State(state): State,
+ State(state): State,
Extension(request_context): Extension,
Extension(authenticated): Extension,
) -> Result, AppError> {
let owner_user_id = authenticated.claims().user_id().to_string();
let project_records = state
- .spacetime_client()
+ .projects()
.list_editor_projects(owner_user_id)
.await
.map_err(map_editor_project_error)?;
let mut projects = Vec::with_capacity(project_records.len());
for project in project_records {
projects.push(editor_project_payload_from_record(
- repair_editor_project_record_inline_media(&state, project).await,
+ state.repair_project_media(project).await,
));
}
@@ -2049,7 +2049,7 @@ fn parse_editor_timestamp_micros(value: &str) -> Option {
}
pub async fn create_editor_project(
- State(state): State,
+ State(state): State,
headers: HeaderMap,
Extension(request_context): Extension,
Extension(authenticated): Extension,
@@ -2073,7 +2073,7 @@ pub async fn create_editor_project(
})
.unwrap_or_else(|| build_prefixed_uuid_id(EDITOR_PROJECT_ID_PREFIX));
let project = state
- .spacetime_client()
+ .projects()
.create_editor_project(EditorProjectCreateRecordInput {
project_id,
owner_user_id,
@@ -2092,20 +2092,20 @@ pub async fn create_editor_project(
}
pub async fn get_editor_project(
- State(state): State,
+ State(state): State,
Path(project_id): Path,
Extension(request_context): Extension,
Extension(authenticated): Extension,
) -> Result, AppError> {
let project = state
- .spacetime_client()
+ .projects()
.get_editor_project(EditorProjectGetRecordInput {
project_id,
owner_user_id: authenticated.claims().user_id().to_string(),
})
.await
.map_err(map_editor_project_error)?;
- let project = repair_editor_project_record_inline_media(&state, project).await;
+ let project = state.repair_project_media(project).await;
Ok(json_success_body(
Some(&request_context),
@@ -2116,7 +2116,7 @@ pub async fn get_editor_project(
}
pub async fn save_editor_project_layout(
- State(state): State,
+ State(state): State,
Path(project_id): Path,
Extension(request_context): Extension,
Extension(authenticated): Extension,
@@ -2128,7 +2128,7 @@ pub async fn save_editor_project_layout(
let owner_user_id = authenticated.claims().user_id().to_string();
let updated_at_micros = current_utc_micros();
let ack = state
- .spacetime_client()
+ .projects()
.save_editor_project_layout_v2_ack(EditorProjectLayoutSaveV2RecordInput {
project_id,
owner_user_id,
@@ -2152,14 +2152,14 @@ pub async fn save_editor_project_layout(
}
pub async fn rename_editor_project(
- State(state): State,
+ State(state): State,
Path(project_id): Path,
Extension(request_context): Extension,
Extension(authenticated): Extension,
Json(payload): Json,
) -> Result, AppError> {
let project = state
- .spacetime_client()
+ .projects()
.rename_editor_project(EditorProjectRenameRecordInput {
project_id,
owner_user_id: authenticated.claims().user_id().to_string(),
@@ -2178,13 +2178,13 @@ pub async fn rename_editor_project(
}
pub async fn delete_editor_project(
- State(state): State,
+ State(state): State,
Path(project_id): Path,
Extension(request_context): Extension,
Extension(authenticated): Extension,
) -> Result, AppError> {
let deleted_project_id = state
- .spacetime_client()
+ .projects()
.delete_editor_project(EditorProjectDeleteRecordInput {
project_id,
owner_user_id: authenticated.claims().user_id().to_string(),
@@ -7754,7 +7754,7 @@ pub async fn snap_editor_image_to_pixel_art(
}
result = async move {
let mut uploaded = upload_editor_generated_image_object_prepared(
- &state,
+ &EditorMediaStorageState::from_ref(&state),
owner_user_id.as_str(),
persistence_identity.task_id.as_str(),
prepared_upload,
@@ -8946,8 +8946,8 @@ pub(crate) async fn extract_editor_ui_design_assets_for_owner(
))
}
-async fn repair_editor_project_record_inline_media(
- state: &AppState,
+pub(crate) async fn repair_editor_project_record_inline_media(
+ state: &EditorMediaStorageState,
mut record: EditorProjectRecord,
) -> EditorProjectRecord {
let owner_user_id = record.owner_user_id.clone();
@@ -8976,7 +8976,7 @@ async fn repair_editor_asset_library_inline_media(
}
async fn repair_editor_project_resource_inline_media(
- state: &AppState,
+ state: &EditorMediaStorageState,
owner_user_id: &str,
resource: EditorProjectResourceRecord,
) -> EditorProjectResourceRecord {
@@ -9058,7 +9058,7 @@ async fn repair_editor_asset_inline_media(
}
let repair = persist_editor_legacy_inline_image(
- state,
+ &EditorMediaStorageState::from_ref(state),
owner_user_id,
asset.image_src.as_str(),
asset.task_id.as_deref().unwrap_or(asset.asset_id.as_str()),
@@ -9112,7 +9112,7 @@ async fn repair_editor_asset_inline_media(
}
async fn persist_editor_legacy_inline_image(
- state: &AppState,
+ state: &EditorMediaStorageState,
owner_user_id: &str,
image_src: &str,
task_id: &str,
@@ -9133,11 +9133,14 @@ async fn persist_editor_legacy_inline_image(
mime_type: decoded.format.mime_type,
extension: decoded.format.extension,
};
- persist_editor_generated_image(
+ persist_editor_generated_image_data(
state,
owner_user_id,
task_id,
- &image,
+ GeneratedImageAssetDataUrl {
+ format: normalize_generated_image_asset_mime(image.mime_type.as_str()),
+ bytes: image.bytes,
+ },
prompt,
actual_prompt,
asset_kind.unwrap_or(EDITOR_LEGACY_INLINE_IMAGE_ASSET_KIND),
@@ -10407,7 +10410,11 @@ pub(crate) async fn complete_editor_canvas_generation(
})
.await
.map_err(map_editor_project_error)?;
- let project = repair_editor_project_record_inline_media(state, project).await;
+ let project = repair_editor_project_record_inline_media(
+ &EditorMediaStorageState::from_ref(state),
+ project,
+ )
+ .await;
let expected_revision = project.canvas.revision;
let viewport = project.viewport.clone();
let project_payload = editor_project_payload_from_record(project);
@@ -10683,7 +10690,11 @@ pub(crate) async fn prepare_editor_canvas_generation_layout(
})
.await
.map_err(map_editor_project_error)?;
- let project = repair_editor_project_record_inline_media(state, project).await;
+ let project = repair_editor_project_record_inline_media(
+ &EditorMediaStorageState::from_ref(state),
+ project,
+ )
+ .await;
let expected_revision = project.canvas.revision;
let viewport = project.viewport.clone();
let project_payload = editor_project_payload_from_record(project);
@@ -10776,7 +10787,11 @@ async fn prepare_editor_canvas_background_removal_layout(
})
.await
.map_err(map_editor_project_error)?;
- let project = repair_editor_project_record_inline_media(state, project).await;
+ let project = repair_editor_project_record_inline_media(
+ &EditorMediaStorageState::from_ref(state),
+ project,
+ )
+ .await;
let expected_revision = project.canvas.revision;
let viewport = project.viewport.clone();
let project_payload = editor_project_payload_from_record(project);
@@ -12123,7 +12138,7 @@ pub(crate) async fn persist_editor_generated_image(
provider: &str,
) -> Result {
persist_editor_generated_image_data(
- state,
+ &EditorMediaStorageState::from_ref(state),
owner_user_id,
task_id,
GeneratedImageAssetDataUrl {
@@ -12162,7 +12177,7 @@ pub(crate) async fn prepare_editor_generated_image(
}))
})?;
let uploaded = upload_editor_generated_image_object_data(
- state,
+ &EditorMediaStorageState::from_ref(state),
caller.owner_user_id.as_str(),
task_id,
GeneratedImageAssetDataUrl {
@@ -12207,7 +12222,7 @@ async fn persist_editor_generated_image_owned(
bytes: image.bytes,
};
persist_editor_generated_image_data(
- state,
+ &EditorMediaStorageState::from_ref(state),
owner_user_id,
task_id,
image_data,
@@ -12223,7 +12238,7 @@ async fn persist_editor_generated_image_owned(
}
async fn persist_editor_generated_image_data(
- state: &AppState,
+ state: &EditorMediaStorageState,
owner_user_id: &str,
task_id: &str,
image: GeneratedImageAssetDataUrl,
@@ -12270,7 +12285,7 @@ async fn persist_editor_generated_image_data(
#[allow(clippy::too_many_arguments)]
async fn upload_editor_generated_image_object_data(
- state: &AppState,
+ state: &EditorMediaStorageState,
owner_user_id: &str,
task_id: &str,
image: GeneratedImageAssetDataUrl,
@@ -12306,7 +12321,7 @@ async fn upload_editor_generated_image_object_data(
}
async fn upload_editor_generated_image_object_prepared(
- state: &AppState,
+ state: &EditorMediaStorageState,
owner_user_id: &str,
task_id: &str,
prepared: GeneratedImageAssetPreparedPut,
@@ -13679,6 +13694,9 @@ pub(crate) fn current_utc_micros() -> i64 {
i64::try_from(duration.as_micros()).expect("current unix micros should fit in i64")
}
+#[cfg(test)]
+mod metadata_tests;
+
#[cfg(test)]
mod tests {
use super::*;
diff --git a/server-rs/crates/api-server/src/editor_project/metadata_tests.rs b/server-rs/crates/api-server/src/editor_project/metadata_tests.rs
new file mode 100644
index 000000000..cf3b2b04f
--- /dev/null
+++ b/server-rs/crates/api-server/src/editor_project/metadata_tests.rs
@@ -0,0 +1,364 @@
+use super::*;
+use crate::{
+ request_context::attach_request_context,
+ state::project_metadata::{EditorProjectMediaRepair, EditorProjectRepository},
+};
+use axum::{
+ Router,
+ body::{Body, to_bytes},
+ http::Request,
+ middleware,
+ routing::{get, patch},
+};
+use futures_util::future::BoxFuture;
+use platform_auth::{AccessTokenClaims, AuthProvider, BindingStatus};
+use spacetime_client::EditorProjectLayoutSaveV2AckRecord;
+use std::sync::Mutex;
+use tower::ServiceExt;
+
+#[derive(Default)]
+struct RecordingProjects {
+ calls: Mutex>,
+ error: Option<&'static str>,
+}
+
+impl RecordingProjects {
+ fn finish(
+ &self,
+ operation: &'static str,
+ input: Value,
+ result: T,
+ ) -> BoxFuture<'_, Result> {
+ self.calls.lock().unwrap().push((operation, input));
+ Box::pin(async move {
+ match self.error {
+ Some(message) => Err(SpacetimeClientError::Procedure(message.to_string())),
+ None => Ok(result),
+ }
+ })
+ }
+}
+
+fn project_record(project_id: &str, owner: &str, title: &str) -> EditorProjectRecord {
+ let viewport = EditorCanvasViewportRecord {
+ x: 1.0,
+ y: 2.0,
+ scale: 1.0,
+ };
+ EditorProjectRecord {
+ project_id: project_id.to_string(),
+ owner_user_id: owner.to_string(),
+ title: title.to_string(),
+ canvas: EditorCanvasRecord {
+ canvas_id: "canvas-fixture".to_string(),
+ project_id: project_id.to_string(),
+ title: title.to_string(),
+ viewport: viewport.clone(),
+ layers: json!([]),
+ revision: 7,
+ layout_storage_version: 2,
+ background_color: None,
+ created_at: "0.000000Z".to_string(),
+ updated_at: "0.000000Z".to_string(),
+ },
+ viewport,
+ layers: json!([]),
+ resources: vec![],
+ created_at: "0.000000Z".to_string(),
+ updated_at: "0.000000Z".to_string(),
+ }
+}
+
+impl EditorProjectRepository for RecordingProjects {
+ fn get_recent_editor_project(
+ &self,
+ owner_user_id: String,
+ ) -> BoxFuture<'_, Result