Files
Genarrative/server-rs/crates/platform-oss/src/template_library.rs
T
kdletters 4951b71d71
Project CI / AI game creator shell Rust crates (push) Successful in 2m46s
Project CI / AI game creator shell Rust smoke (push) Successful in 3m37s
Project CI / AI game creator shell Rust lane 2/2 (push) Failing after 6m5s
Project CI / AI game creator shell Rust lane 1/2 (push) Failing after 6m40s
Project CI / Repository checks (push) Successful in 5m37s
Project CI / Frontend tests (push) Successful in 6m21s
Project CI / Backend tests (push) Successful in 9m56s
Project CI / Native shell tests (push) Successful in 10m57s
Project CI / AI game creator shell web tests (push) Successful in 5m29s
项目快照按渠道分区,后台项目工程支持渠道筛选与游标分页
项目快照对象键升级为 agc/project-snapshots/v2/{channel}/{user}/{project}/,渠道取部署配置 GENARRATIVE_AGC_PROJECT_SNAPSHOT_CHANNEL(缺省沿用客户端下载渠道),非法渠道在写入处失败关闭
后台新增渠道列表接口,项目工程列表与下载接受 channel,游标绑定渠道并拒绝跨渠道复用
后台项目工程用户列改为素材查询口径:昵称 + 陶泥号 + 用户详情入口,由 api-server 解析作者信息
后台项目工程改为游标分页:每页 20/50/100、上一页/下一页与当前页提示,翻页失败保留当前页
部署环境示例补充快照渠道配置,并同步运维、技术方案、里程碑与决策记录文档
2026-09-21 16:57:16 +08:00

1246 lines
44 KiB
Rust

//! 模板清单与 CLI 共享发布锁;清单写入结果不明时保留锁,不猜测或重试提交。
use std::{collections::BTreeMap, fmt, time::Duration};
use reqwest::{Client, Method, Response, StatusCode, Url, redirect::Policy};
use serde_json::{Value, json};
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
use crate::{OssClient, sha256_hex, signed_request_builder};
const BUCKET: &str = "agc-dev";
const ENDPOINT: &str = "oss-rg-china-mainland.aliyuncs.com";
const PUBLIC_BASE: &str = "https://agc-dev.oss-rg-china-mainland.aliyuncs.com/";
const INDEX_KEY: &str = "templates/index.json";
const LOCK_KEY: &str = "templates/.publish-lock.json";
const MAX_INDEX_BYTES: usize = 4 * 1024 * 1024;
const MAX_OBJECT_BYTES: usize = 5 * 1024 * 1024;
/// 模板包(`application/zip`)是后台导入路径唯一的超大对象,单独给出上限。
const MAX_TEMPLATE_ZIP_BYTES: usize = 64 * 1024 * 1024;
const MAX_CONTROL_BYTES: usize = 64 * 1024;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum TemplateStoreError {
NotFound,
Busy,
Unavailable,
Invalid,
Transport,
Uncertain,
}
impl fmt::Display for TemplateStoreError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(match self {
Self::NotFound => "模板库对象不存在",
Self::Busy => "模板库正在修改或发布锁归属已变化,请稍后刷新",
Self::Unavailable => "模板库存储暂不可用或不满足安全发布条件",
Self::Invalid => "模板库存储请求或对象内容无效",
Self::Transport => "模板库存储请求失败,请稍后重试",
Self::Uncertain => "模板保存或发布锁状态需要核对,已停止自动操作,请联系运维",
})
}
}
impl std::error::Error for TemplateStoreError {}
#[derive(Clone)]
pub struct TemplateLibraryStore {
oss: OssClient,
http: Client,
#[cfg(test)]
test_base: Option<Url>,
}
// 不实现 Clone/Drop:只能显式释放一次,任务取消或进程退出时留锁。
pub struct TemplatePublishSession {
store: TemplateLibraryStore,
owner: String,
release_allowed: bool,
objects_verified: bool,
}
impl TemplateLibraryStore {
pub fn new(oss: OssClient) -> Result<Self, TemplateStoreError> {
if oss.config.bucket() != BUCKET || oss.config.endpoint() != ENDPOINT {
return Err(TemplateStoreError::Invalid);
}
Ok(Self {
oss,
http: build_http_client()?,
#[cfg(test)]
test_base: None,
})
}
pub async fn read_index(&self) -> Result<Vec<u8>, TemplateStoreError> {
self.read_key(INDEX_KEY, MAX_INDEX_BYTES).await
}
pub async fn read_public_index() -> Result<Vec<u8>, TemplateStoreError> {
let http = build_http_client()?;
let url = Url::parse(&format!("{PUBLIC_BASE}{INDEX_KEY}"))
.map_err(|_| TemplateStoreError::Invalid)?;
read_public_index_at(&http, url).await
}
pub async fn begin_publish(
&self,
owner: String,
) -> Result<TemplatePublishSession, TemplateStoreError> {
if owner.is_empty()
|| owner.len() > 128
|| !owner
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'-')
{
return Err(TemplateStoreError::Invalid);
}
let response = self
.request(
Method::GET,
None,
Some("versioning"),
None,
None,
"no-cache",
false,
)
.await?;
if response.status() != StatusCode::OK {
return Err(TemplateStoreError::Unavailable);
}
let versioning = read_bytes(response, MAX_CONTROL_BYTES).await?;
if !is_unversioned_configuration(&versioning) {
return Err(TemplateStoreError::Unavailable);
}
let created_at = OffsetDateTime::now_utc()
.format(&Rfc3339)
.map_err(|_| TemplateStoreError::Invalid)?;
let body = serde_json::to_vec(&json!({ "owner": owner, "createdAt": created_at }))
.map_err(|_| TemplateStoreError::Invalid)?;
let acquired = self
.request(
Method::PUT,
Some(LOCK_KEY),
None,
Some(body),
Some("application/json"),
"no-store",
true,
)
.await
.map_err(|_| TemplateStoreError::Uncertain)?;
match acquired.status() {
StatusCode::OK => Ok(TemplatePublishSession {
store: self.clone(),
owner,
release_allowed: true,
objects_verified: true,
}),
StatusCode::CONFLICT => Err(TemplateStoreError::Busy),
status if is_certain_rejection(status) => Err(TemplateStoreError::Unavailable),
_ => Err(TemplateStoreError::Uncertain),
}
}
async fn read_key(&self, key: &str, max_bytes: usize) -> Result<Vec<u8>, TemplateStoreError> {
let response = self
.request(Method::GET, Some(key), None, None, None, "no-cache", false)
.await?;
read_success_bytes(response, max_bytes).await
}
#[allow(clippy::too_many_arguments)]
async fn request(
&self,
method: Method,
key: Option<&str>,
query: Option<&str>,
body: Option<Vec<u8>>,
content_type: Option<&str>,
cache_control: &str,
create_only: bool,
) -> Result<Response, TemplateStoreError> {
let mut url = self.target_url(key)?;
url.set_query(query);
let mut headers =
BTreeMap::from([("cache-control".to_string(), cache_control.to_string())]);
if create_only {
headers.insert("x-oss-forbid-overwrite".to_string(), "true".to_string());
}
let mut request = signed_request_builder(
&self.http,
&self.oss.config,
method,
key,
url,
content_type,
&headers,
)
.map_err(|_| TemplateStoreError::Invalid)?;
if let Some(body) = body {
request = request
.header(reqwest::header::CONTENT_LENGTH, body.len())
.body(body);
}
request
.send()
.await
.map_err(|_| TemplateStoreError::Transport)
}
fn target_url(&self, key: Option<&str>) -> Result<Url, TemplateStoreError> {
#[cfg(test)]
let base = self
.test_base
.as_ref()
.map(Url::as_str)
.unwrap_or(PUBLIC_BASE);
#[cfg(not(test))]
let base = PUBLIC_BASE;
let mut url = Url::parse(base).map_err(|_| TemplateStoreError::Invalid)?;
if let Some(key) = key {
url.path_segments_mut()
.map_err(|_| TemplateStoreError::Invalid)?
.clear()
.extend(key.split('/'));
}
Ok(url)
}
}
impl TemplatePublishSession {
pub async fn read_index(&self) -> Result<Vec<u8>, TemplateStoreError> {
self.store.read_index().await
}
pub async fn read_object(
&self,
key: &str,
max_bytes: usize,
) -> Result<Vec<u8>, TemplateStoreError> {
validate_content_key(key)?;
if max_bytes == 0 || max_bytes > MAX_OBJECT_BYTES {
return Err(TemplateStoreError::Invalid);
}
self.store.read_key(key, max_bytes).await
}
pub async fn put_immutable(
&mut self,
key: &str,
bytes: Vec<u8>,
content_type: &str,
) -> Result<(), TemplateStoreError> {
if !self.release_allowed {
return Err(TemplateStoreError::Uncertain);
}
if !self.objects_verified {
return Err(TemplateStoreError::Invalid);
}
validate_content_key(key)?;
let parts: Vec<_> = key.split('/').collect();
let max_bytes = match content_type {
"application/json" => MAX_INDEX_BYTES,
// 模板包按上传字节原样发布,见决策「后台模板上传」;仍受内容寻址与回读校验约束。
"application/zip" => MAX_TEMPLATE_ZIP_BYTES,
"image/png" | "image/jpeg" | "image/webp" => MAX_OBJECT_BYTES,
_ => return Err(TemplateStoreError::Invalid),
};
if bytes.is_empty()
|| bytes.len() > max_bytes
|| parts.len() != 6
|| parts[3] != "sha256"
|| parts[4] != sha256_hex(&bytes)
{
return Err(TemplateStoreError::Invalid);
}
self.objects_verified = false;
let response = self
.store
.request(
Method::PUT,
Some(key),
None,
Some(bytes.clone()),
Some(content_type),
"public, max-age=31536000, immutable",
true,
)
.await?;
if !response.status().is_success() && response.status() != StatusCode::CONFLICT {
return Err(TemplateStoreError::Unavailable);
}
let existing = self.store.read_key(key, max_bytes).await?;
if existing != bytes {
return Err(TemplateStoreError::Invalid);
}
self.objects_verified = true;
Ok(())
}
pub async fn commit_index(&mut self, bytes: Vec<u8>) -> Result<(), TemplateStoreError> {
if !self.release_allowed {
return Err(TemplateStoreError::Uncertain);
}
if !self.objects_verified || bytes.is_empty() || bytes.len() > MAX_INDEX_BYTES {
return Err(TemplateStoreError::Invalid);
}
self.release_allowed = false;
let response = self
.store
.request(
Method::PUT,
Some(INDEX_KEY),
None,
Some(bytes.clone()),
Some("application/json"),
"no-store",
false,
)
.await
.map_err(|_| TemplateStoreError::Uncertain)?;
if response.status().is_success() {
self.release_allowed = true;
} else if is_certain_rejection(response.status()) {
self.release_allowed = true;
return Err(TemplateStoreError::Unavailable);
} else {
return Err(TemplateStoreError::Uncertain);
}
if self.store.read_index().await? != bytes {
return Err(TemplateStoreError::Invalid);
}
Ok(())
}
pub async fn finish(self) -> Result<(), TemplateStoreError> {
if !self.release_allowed {
return Err(TemplateStoreError::Uncertain);
}
let bytes = self.store.read_key(LOCK_KEY, MAX_CONTROL_BYTES).await?;
let lock: Value =
serde_json::from_slice(&bytes).map_err(|_| TemplateStoreError::Invalid)?;
if lock.get("owner").and_then(Value::as_str) != Some(self.owner.as_str()) {
return Err(TemplateStoreError::Busy);
}
let response = self
.store
.request(
Method::DELETE,
Some(LOCK_KEY),
None,
None,
None,
"no-store",
false,
)
.await
.map_err(|_| TemplateStoreError::Uncertain)?;
if response.status().is_success() {
Ok(())
} else if is_certain_rejection(response.status()) {
Err(TemplateStoreError::Unavailable)
} else {
Err(TemplateStoreError::Uncertain)
}
}
}
fn build_http_client() -> Result<Client, TemplateStoreError> {
Client::builder()
.redirect(Policy::none())
.retry(reqwest::retry::never())
.timeout(Duration::from_secs(30))
.build()
.map_err(|_| TemplateStoreError::Unavailable)
}
async fn read_public_index_at(http: &Client, url: Url) -> Result<Vec<u8>, TemplateStoreError> {
let response = http
.get(url)
.header(reqwest::header::CACHE_CONTROL, "no-cache")
.send()
.await
.map_err(|_| TemplateStoreError::Transport)?;
read_success_bytes(response, MAX_INDEX_BYTES).await
}
async fn read_success_bytes(
response: Response,
max_bytes: usize,
) -> Result<Vec<u8>, TemplateStoreError> {
match response.status() {
StatusCode::OK => read_bytes(response, max_bytes).await,
StatusCode::NOT_FOUND => Err(TemplateStoreError::NotFound),
_ => Err(TemplateStoreError::Unavailable),
}
}
async fn read_bytes(
mut response: Response,
max_bytes: usize,
) -> Result<Vec<u8>, TemplateStoreError> {
if response
.content_length()
.is_some_and(|length| length > max_bytes as u64)
{
return Err(TemplateStoreError::Invalid);
}
let mut bytes = Vec::new();
while let Some(chunk) = response
.chunk()
.await
.map_err(|_| TemplateStoreError::Transport)?
{
if chunk.len() > max_bytes - bytes.len() {
return Err(TemplateStoreError::Invalid);
}
bytes.extend_from_slice(&chunk);
}
Ok(bytes)
}
fn is_certain_rejection(status: StatusCode) -> bool {
status.is_client_error() && status != StatusCode::REQUEST_TIMEOUT
}
fn validate_content_key(key: &str) -> Result<(), TemplateStoreError> {
if !key.starts_with("templates/v1/")
|| key.len() > 1024
|| key.chars().any(|ch| {
ch.is_control() || ch.is_whitespace() || matches!(ch, '\\' | ':' | '%' | '?' | '#')
})
|| key
.split('/')
.any(|part| part.is_empty() || part == "." || part == "..")
|| key.split('/').count() < 4
{
return Err(TemplateStoreError::Invalid);
}
Ok(())
}
// 只接受官方空配置的 XML 子集;未知元素、状态、命名空间或畸形 XML 均失败关闭。
fn is_unversioned_configuration(bytes: &[u8]) -> bool {
let Ok(text) = std::str::from_utf8(bytes) else {
return false;
};
let mut text = text.trim_start_matches('\u{feff}').trim();
if let Some(declaration) = text.strip_prefix("<?xml") {
let Some((attributes, remaining)) = declaration.split_once("?>") else {
return false;
};
let Some(attributes) = xml_attributes(attributes) else {
return false;
};
if attributes.get("version").map(String::as_str) != Some("1.0")
|| attributes.iter().any(|(key, value)| match key.as_str() {
"version" => false,
"encoding" => !value.eq_ignore_ascii_case("UTF-8"),
"standalone" => !matches!(value.as_str(), "yes" | "no"),
_ => true,
})
{
return false;
}
text = remaining.trim();
}
let Some(root) = text.strip_prefix("<VersioningConfiguration") else {
return false;
};
let Some((attributes, remaining)) = root.split_once('>') else {
return false;
};
let (attributes, empty) = attributes
.strip_suffix('/')
.map_or((attributes, false), |value| (value, true));
let Some(attributes) = xml_attributes(attributes) else {
return false;
};
if attributes
.iter()
.any(|(key, value)| key != "xmlns" || value != "http://doc.oss-cn-hangzhou.aliyuncs.com")
{
return false;
}
if empty {
remaining.trim().is_empty()
} else {
remaining.trim() == "</VersioningConfiguration>"
}
}
fn xml_attributes(mut text: &str) -> Option<BTreeMap<String, String>> {
let mut result = BTreeMap::new();
while !text.is_empty() {
if !text.starts_with(char::is_whitespace) {
return None;
}
text = text.trim_start();
if text.is_empty() {
break;
}
let name_end = text.find(|ch: char| !ch.is_ascii_alphanumeric() && ch != '-')?;
let (name, remaining) = text.split_at(name_end);
if name.is_empty() {
return None;
}
let remaining = remaining.trim_start().strip_prefix('=')?.trim_start();
let quote = remaining.chars().next()?;
if !matches!(quote, '\'' | '"') {
return None;
}
let (value, rest) = remaining[1..].split_once(quote)?;
if value.contains(['<', '>', '&'])
|| result.insert(name.to_string(), value.to_string()).is_some()
{
return None;
}
text = rest;
}
Some(result)
}
#[cfg(test)]
mod tests {
use super::*;
use hmac::{Hmac, Mac};
use sha2::{Digest, Sha256};
use std::sync::{Arc, Mutex};
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
net::{TcpListener, TcpStream},
task::JoinHandle,
};
#[derive(Clone)]
struct RecordedRequest {
method: String,
target: String,
headers: BTreeMap<String, String>,
body: Vec<u8>,
}
struct MockState {
objects: BTreeMap<String, Vec<u8>>,
requests: Vec<RecordedRequest>,
versioning_status: u16,
versioning_body: Vec<u8>,
index_status: Option<u16>,
disconnect_index: bool,
disconnect_acquire: bool,
delete_status: Option<u16>,
chunked_reads: bool,
}
impl Default for MockState {
fn default() -> Self {
Self {
objects: BTreeMap::from([(INDEX_KEY.to_string(), b"{\"templates\":[]}".to_vec())]),
requests: Vec::new(),
versioning_status: 200,
versioning_body:
br#"<VersioningConfiguration xmlns="http://doc.oss-cn-hangzhou.aliyuncs.com"/>"#
.to_vec(),
index_status: None,
disconnect_index: false,
disconnect_acquire: false,
delete_status: None,
chunked_reads: false,
}
}
}
struct MockServer {
store: TemplateLibraryStore,
state: Arc<Mutex<MockState>>,
task: JoinHandle<()>,
}
impl Drop for MockServer {
fn drop(&mut self) {
self.task.abort();
}
}
fn test_oss(bucket: &str) -> OssClient {
OssClient::new(
crate::OssConfig::new(
bucket.to_string(),
ENDPOINT.to_string(),
"test-id".to_string(),
"test-secret".to_string(),
600,
600,
20 * 1024 * 1024,
200,
)
.unwrap(),
)
}
impl MockServer {
async fn start() -> Self {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let base = Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap();
let mut store = TemplateLibraryStore::new(test_oss(BUCKET)).unwrap();
store.test_base = Some(base);
store.http = Client::builder()
.no_proxy()
.redirect(Policy::none())
.retry(reqwest::retry::never())
.timeout(Duration::from_secs(2))
.build()
.unwrap();
let state = Arc::new(Mutex::new(MockState::default()));
let server_state = state.clone();
let task = tokio::spawn(async move {
loop {
let Ok((socket, _)) = listener.accept().await else {
break;
};
let state = server_state.clone();
tokio::spawn(serve_connection(socket, state));
}
});
Self { store, state, task }
}
fn count(&self, method: &str, key: &str) -> usize {
self.state
.lock()
.unwrap()
.requests
.iter()
.filter(|request| request.method == method && request.target == format!("/{key}"))
.count()
}
}
async fn read_request(socket: &mut TcpStream) -> Option<RecordedRequest> {
let mut bytes = Vec::new();
let header_end = loop {
if let Some(position) = bytes.windows(4).position(|window| window == b"\r\n\r\n") {
break position;
}
let mut buffer = [0; 8192];
let read = socket.read(&mut buffer).await.ok()?;
if read == 0 {
return None;
}
bytes.extend_from_slice(&buffer[..read]);
};
let head = std::str::from_utf8(&bytes[..header_end]).ok()?;
let mut lines = head.split("\r\n");
let mut request_line = lines.next()?.split_whitespace();
let method = request_line.next()?.to_string();
let target = request_line.next()?.to_string();
let headers: BTreeMap<_, _> = lines
.filter_map(|line| {
let (name, value) = line.split_once(':')?;
Some((name.to_ascii_lowercase(), value.trim().to_string()))
})
.collect();
let length = headers
.get("content-length")
.and_then(|value| value.parse::<usize>().ok())
.unwrap_or(0);
// 模板包上限高于图片 / 元数据,替身必须按同一口径放行。
if length > MAX_TEMPLATE_ZIP_BYTES + 1024 {
return None;
}
let start = header_end + 4;
while bytes.len() - start < length {
let mut buffer = [0; 8192];
let read = socket.read(&mut buffer).await.ok()?;
if read == 0 {
return None;
}
bytes.extend_from_slice(&buffer[..read]);
}
Some(RecordedRequest {
method,
target,
headers,
body: bytes[start..start + length].to_vec(),
})
}
fn mock_response(
state: &mut MockState,
request: RecordedRequest,
) -> Option<(u16, Vec<u8>, bool)> {
state.requests.push(request.clone());
if request.target == "/?versioning" {
return Some((
state.versioning_status,
state.versioning_body.clone(),
false,
));
}
let key = request.target.trim_start_matches('/');
match request.method.as_str() {
"GET" => Some(match state.objects.get(key) {
Some(bytes) => (200, bytes.clone(), state.chunked_reads),
None => (404, b"private upstream details".to_vec(), false),
}),
"PUT" => {
if request
.headers
.get("x-oss-forbid-overwrite")
.map(String::as_str)
== Some("true")
&& state.objects.contains_key(key)
{
return Some((409, b"FileAlreadyExists".to_vec(), false));
}
if key == INDEX_KEY
&& let Some(status) = state.index_status
{
return Some((status, b"private upstream details".to_vec(), false));
}
state.objects.insert(key.to_string(), request.body);
if (key == INDEX_KEY && state.disconnect_index)
|| (key == LOCK_KEY && state.disconnect_acquire)
{
return None;
}
Some((200, Vec::new(), false))
}
"DELETE" => {
if let Some(status) = state.delete_status {
return Some((status, b"private upstream details".to_vec(), false));
}
state.objects.remove(key);
Some((204, Vec::new(), false))
}
_ => Some((405, Vec::new(), false)),
}
}
async fn serve_connection(mut socket: TcpStream, state: Arc<Mutex<MockState>>) {
let Some(request) = read_request(&mut socket).await else {
return;
};
let response = { mock_response(&mut state.lock().unwrap(), request) };
let Some((status, body, chunked)) = response else {
return;
};
let mut bytes = if chunked {
format!("HTTP/1.1 {status} Test\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n{:x}\r\n", body.len()).into_bytes()
} else {
format!("HTTP/1.1 {status} Test\r\nContent-Length: {}\r\nLocation: /redirected\r\nConnection: close\r\n\r\n", body.len()).into_bytes()
};
bytes.extend_from_slice(&body);
if chunked {
bytes.extend_from_slice(b"\r\n0\r\n\r\n");
}
let _ = socket.write_all(&bytes).await;
}
fn content_key(bytes: &[u8]) -> String {
format!(
"templates/v1/demo/sha256/{}/template.json",
sha256_hex(bytes)
)
}
#[test]
fn fixed_target_and_strict_unversioned_xml_fail_closed() {
assert!(TemplateLibraryStore::new(test_oss("resource-bucket")).is_err());
for valid in [
"<VersioningConfiguration/>",
"<VersioningConfiguration></VersioningConfiguration>",
"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\n<VersioningConfiguration xmlns=\"http://doc.oss-cn-hangzhou.aliyuncs.com\"/>",
] {
assert!(is_unversioned_configuration(valid.as_bytes()), "{valid}");
}
for invalid in [
"",
"<Error/>",
"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>",
"<VersioningConfiguration><Status>Suspended</Status></VersioningConfiguration>",
"<VersioningConfiguration><Status/></VersioningConfiguration>",
"<VersioningConfiguration xmlns=\"unknown\"/>",
"<VersioningConfiguration unknown=\"value\"/>",
"<VersioningConfiguration>",
"<VersioningConfiguration/><Other/>",
"<VersioningConfigurationInvalid/>",
"<?xml version=\"2.0\"?><VersioningConfiguration/>",
"<!DOCTYPE root><VersioningConfiguration/>",
] {
assert!(
!is_unversioned_configuration(invalid.as_bytes()),
"{invalid}"
);
}
}
#[test]
fn v4_signs_versioning_query_and_create_only_header() {
let client = build_http_client().unwrap();
let oss = test_oss(BUCKET);
for (method, key, query, content_type, create_only) in [
(Method::GET, None, Some("versioning"), None, false),
(
Method::PUT,
Some(LOCK_KEY),
None,
Some("application/json"),
true,
),
] {
let mut url = Url::parse(PUBLIC_BASE).unwrap();
url.set_path(key.unwrap_or(""));
url.set_query(query);
let mut headers =
BTreeMap::from([("cache-control".to_string(), "no-store".to_string())]);
if create_only {
headers.insert("x-oss-forbid-overwrite".to_string(), "true".to_string());
}
let request = signed_request_builder(
&client,
&oss.config,
method.clone(),
key,
url,
content_type,
&headers,
)
.unwrap()
.build()
.unwrap();
let signed_at = request.headers()["x-oss-date"].to_str().unwrap();
let authorization = request.headers()["authorization"].to_str().unwrap();
let scope = authorization
.split("Credential=test-id/")
.nth(1)
.unwrap()
.split(',')
.next()
.unwrap();
let uri = key.map_or("/agc-dev/".to_string(), |key| format!("/agc-dev/{key}"));
let query = if query.is_some() { "versioning" } else { "" };
let content_header =
content_type.map_or(String::new(), |value| format!("content-type:{value}\n"));
let condition_header = if create_only {
"x-oss-forbid-overwrite:true\n"
} else {
""
};
let canonical = format!(
"{method}\n{uri}\n{query}\ncache-control:no-store\n{content_header}host:agc-dev.oss-rg-china-mainland.aliyuncs.com\nx-oss-content-sha256:UNSIGNED-PAYLOAD\nx-oss-date:{signed_at}\n{condition_header}\ncache-control;host\nUNSIGNED-PAYLOAD"
);
// 独立按官方 Python SDK oss2/auth.py 的 V4 算法计算,不调用生产签名 helper。
let string_to_sign = format!(
"OSS4-HMAC-SHA256\n{signed_at}\n{scope}\n{:x}",
Sha256::digest(canonical.as_bytes())
);
let hmac = |key: &[u8], message: &[u8]| {
let mut mac = Hmac::<Sha256>::new_from_slice(key).unwrap();
mac.update(message);
mac.finalize().into_bytes().to_vec()
};
let mut signing_key = b"aliyun_v4test-secret".to_vec();
for part in scope.split('/') {
signing_key = hmac(&signing_key, part.as_bytes());
}
let signature = hmac(&signing_key, string_to_sign.as_bytes())
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
assert_eq!(
authorization,
format!(
"OSS4-HMAC-SHA256 Credential=test-id/{scope},AdditionalHeaders=cache-control;host,Signature={signature}"
)
);
if create_only {
assert_eq!(request.headers()["x-oss-forbid-overwrite"], "true");
}
}
}
#[tokio::test]
async fn publish_roundtrip_uses_shared_owner_and_verifies_before_committing() {
let server = MockServer::start().await;
let legacy = "templates/v1/demo/template.json";
server
.state
.lock()
.unwrap()
.objects
.insert(legacy.to_string(), b"legacy metadata".to_vec());
let mut session = server
.store
.begin_publish("rust-owner".to_string())
.await
.unwrap();
assert_eq!(
session.read_object(legacy, MAX_INDEX_BYTES).await.unwrap(),
b"legacy metadata"
);
let metadata = b"{\"title\":\"new title\"}".to_vec();
let key = content_key(&metadata);
session
.put_immutable(&key, metadata.clone(), "application/json")
.await
.unwrap();
let index = b"{\"templates\":[],\"inactiveTemplates\":[]}".to_vec();
session.commit_index(index.clone()).await.unwrap();
session.finish().await.unwrap();
let state = server.state.lock().unwrap();
assert_eq!(state.objects[INDEX_KEY], index);
assert_eq!(state.objects[&key], metadata);
assert!(!state.objects.contains_key(LOCK_KEY));
let acquired = state
.requests
.iter()
.find(|r| r.method == "PUT" && r.target == format!("/{LOCK_KEY}"))
.unwrap();
let owner: Value = serde_json::from_slice(&acquired.body).unwrap();
assert_eq!(owner["owner"], "rust-owner");
assert!(owner["createdAt"].as_str().unwrap().ends_with('Z'));
assert_eq!(acquired.headers["x-oss-forbid-overwrite"], "true");
assert_eq!(acquired.headers["cache-control"], "no-store");
let upload = state
.requests
.iter()
.find(|r| r.method == "PUT" && r.target == format!("/{key}"))
.unwrap();
assert_eq!(
upload.headers["cache-control"],
"public, max-age=31536000, immutable"
);
let commit = state
.requests
.iter()
.position(|r| r.method == "PUT" && r.target == format!("/{INDEX_KEY}"))
.unwrap();
let verified = state
.requests
.iter()
.position(|r| r.method == "GET" && r.target == format!("/{key}"))
.unwrap();
assert!(verified < commit);
assert_eq!(state.requests[commit].headers["cache-control"], "no-store");
assert_eq!(
state
.requests
.iter()
.filter(|r| r.method == "DELETE")
.count(),
1
);
}
#[tokio::test]
async fn unsafe_or_unknown_versioning_never_sends_a_put() {
for (status, xml) in [
(
200,
"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>",
),
(
200,
"<VersioningConfiguration><Status>Suspended</Status></VersioningConfiguration>",
),
(200, "<Error>unknown</Error>"),
(200, ""),
(403, "<Error>private access details</Error>"),
(302, "<VersioningConfiguration/>"),
(500, "<VersioningConfiguration/>"),
] {
let server = MockServer::start().await;
{
let mut state = server.state.lock().unwrap();
state.versioning_status = status;
state.versioning_body = xml.as_bytes().to_vec();
}
assert!(matches!(
server.store.begin_publish("owner".to_string()).await,
Err(TemplateStoreError::Unavailable)
));
assert!(
server
.state
.lock()
.unwrap()
.requests
.iter()
.all(|r| r.method == "GET")
);
assert_eq!(server.state.lock().unwrap().requests.len(), 1);
}
}
#[tokio::test]
async fn javascript_owner_lock_is_preserved_and_publishers_are_mutually_exclusive() {
let server = MockServer::start().await;
let js_lock =
br#"{"owner":"js-cli-owner","createdAt":"2026-09-19T00:00:00.000Z"}"#.to_vec();
server
.state
.lock()
.unwrap()
.objects
.insert(LOCK_KEY.to_string(), js_lock.clone());
assert!(matches!(
server.store.begin_publish("rust-owner".to_string()).await,
Err(TemplateStoreError::Busy)
));
assert_eq!(server.state.lock().unwrap().objects[LOCK_KEY], js_lock);
assert_eq!(server.count("DELETE", LOCK_KEY), 0);
server.state.lock().unwrap().objects.remove(LOCK_KEY);
let (first, second) = tokio::join!(
server.store.begin_publish("publisher-one".to_string()),
server.store.begin_publish("publisher-two".to_string()),
);
let mut winner = match (first, second) {
(Ok(session), Err(TemplateStoreError::Busy))
| (Err(TemplateStoreError::Busy), Ok(session)) => session,
_ => panic!("exactly one publisher must acquire the shared lock"),
};
winner
.commit_index(b"{\"newIndex\":true}".to_vec())
.await
.unwrap();
winner.finish().await.unwrap();
let later = server
.store
.begin_publish("later-owner".to_string())
.await
.unwrap();
assert_eq!(later.read_index().await.unwrap(), b"{\"newIndex\":true}");
later.finish().await.unwrap();
}
#[tokio::test]
async fn acquisition_disconnect_leaves_an_unclaimed_lock_without_delete_or_retry() {
let server = MockServer::start().await;
server.state.lock().unwrap().disconnect_acquire = true;
assert!(matches!(
server
.store
.begin_publish("uncertain-owner".to_string())
.await,
Err(TemplateStoreError::Uncertain)
));
assert!(server.state.lock().unwrap().objects.contains_key(LOCK_KEY));
assert_eq!(server.count("PUT", LOCK_KEY), 1);
assert_eq!(server.count("DELETE", LOCK_KEY), 0);
}
#[tokio::test]
async fn immutable_conflict_requires_identical_readback_and_blocks_commit_on_mismatch() {
for same in [true, false] {
let server = MockServer::start().await;
let body = b"metadata".to_vec();
let key = content_key(&body);
server.state.lock().unwrap().objects.insert(
key.clone(),
if same {
body.clone()
} else {
b"other content".to_vec()
},
);
let mut session = server
.store
.begin_publish("owner".to_string())
.await
.unwrap();
let result = session.put_immutable(&key, body, "application/json").await;
if same {
assert_eq!(result, Ok(()));
} else {
assert_eq!(result, Err(TemplateStoreError::Invalid));
assert_eq!(
session.commit_index(b"{}".to_vec()).await,
Err(TemplateStoreError::Invalid)
);
assert_eq!(server.count("PUT", INDEX_KEY), 0);
}
session.finish().await.unwrap();
assert_eq!(server.count("GET", &key), 1);
}
}
#[tokio::test]
async fn template_zip_objects_accept_larger_bodies_than_images() {
let server = MockServer::start().await;
let mut session = server
.store
.begin_publish("zip-owner".to_string())
.await
.unwrap();
// 6 MiB 超过图片 / 元数据的 5 MiB 上限:模板包必须放行,并仍按自身字节摘要定位与回读。
let zip = vec![7_u8; 6 * 1024 * 1024];
let key = format!("templates/v1/demo/sha256/{}/template.zip", sha256_hex(&zip));
assert_eq!(
session
.put_immutable(&key, zip.clone(), "application/zip")
.await,
Ok(())
);
assert_eq!(server.count("PUT", &key), 1);
// 未知内容类型仍然失败关闭。
let body = b"tiny".to_vec();
let unknown_key = format!("templates/v1/demo/sha256/{}/cover.png", sha256_hex(&body));
assert_eq!(
session
.put_immutable(&unknown_key, body, "application/octet-stream")
.await,
Err(TemplateStoreError::Invalid)
);
session.finish().await.unwrap();
}
#[tokio::test]
async fn uncertain_commit_never_retries_or_releases_even_if_the_write_reached_storage() {
for status in [Some(500), Some(408), Some(302), None] {
let server = MockServer::start().await;
{
let mut state = server.state.lock().unwrap();
state.index_status = status;
state.disconnect_index = status.is_none();
}
let mut session = server
.store
.begin_publish("owner".to_string())
.await
.unwrap();
let next = b"{\"committed\":true}".to_vec();
assert_eq!(
session.commit_index(next.clone()).await,
Err(TemplateStoreError::Uncertain)
);
assert_eq!(
session.commit_index(next.clone()).await,
Err(TemplateStoreError::Uncertain)
);
assert_eq!(session.finish().await, Err(TemplateStoreError::Uncertain));
assert_eq!(server.count("PUT", INDEX_KEY), 1);
assert_eq!(server.count("GET", INDEX_KEY), 0);
assert_eq!(server.count("DELETE", LOCK_KEY), 0);
let state = server.state.lock().unwrap();
assert!(state.objects.contains_key(LOCK_KEY));
if status.is_none() {
assert_eq!(state.objects[INDEX_KEY], next);
}
}
}
#[tokio::test]
async fn certain_rejection_releases_but_foreign_owner_and_release_failure_do_not() {
let server = MockServer::start().await;
server.state.lock().unwrap().index_status = Some(403);
let mut session = server
.store
.begin_publish("owner".to_string())
.await
.unwrap();
assert_eq!(
session.commit_index(b"{}".to_vec()).await,
Err(TemplateStoreError::Unavailable)
);
session.finish().await.unwrap();
assert!(!server.state.lock().unwrap().objects.contains_key(LOCK_KEY));
let session = server
.store
.begin_publish("owner".to_string())
.await
.unwrap();
server
.state
.lock()
.unwrap()
.objects
.insert(LOCK_KEY.to_string(), br#"{"owner":"other-owner"}"#.to_vec());
assert_eq!(session.finish().await, Err(TemplateStoreError::Busy));
assert_eq!(server.count("DELETE", LOCK_KEY), 1);
server.state.lock().unwrap().objects.remove(LOCK_KEY);
let session = server
.store
.begin_publish("owner".to_string())
.await
.unwrap();
server.state.lock().unwrap().delete_status = Some(503);
assert_eq!(session.finish().await, Err(TemplateStoreError::Uncertain));
assert_eq!(server.count("DELETE", LOCK_KEY), 2);
assert!(server.state.lock().unwrap().objects.contains_key(LOCK_KEY));
}
#[tokio::test]
async fn index_reads_are_bounded_and_public_reads_never_send_credentials() {
let server = MockServer::start().await;
let url = server.store.target_url(Some(INDEX_KEY)).unwrap();
assert_eq!(
read_public_index_at(&server.store.http, url).await.unwrap(),
b"{\"templates\":[]}"
);
assert!(
!server.state.lock().unwrap().requests[0]
.headers
.contains_key("authorization")
);
for chunked in [false, true] {
{
let mut state = server.state.lock().unwrap();
state.chunked_reads = chunked;
state
.objects
.insert(INDEX_KEY.to_string(), vec![b' '; MAX_INDEX_BYTES + 1]);
}
assert_eq!(
server.store.read_index().await,
Err(TemplateStoreError::Invalid)
);
}
server.state.lock().unwrap().objects.remove(INDEX_KEY);
assert_eq!(
server.store.read_index().await,
Err(TemplateStoreError::NotFound)
);
}
#[tokio::test]
async fn object_access_cannot_reach_the_lock_index_or_other_prefixes() {
let server = MockServer::start().await;
let mut session = server
.store
.begin_publish("owner".to_string())
.await
.unwrap();
let before = server.state.lock().unwrap().requests.len();
for key in [
LOCK_KEY,
INDEX_KEY,
"agc/project-snapshots/v1/a/b",
"agc/project-snapshots/v2/dev/a/b",
"templates/v1/a/../index.json",
"templates/v1/a/%2e%2e/index.json",
"templates/v1/a/file?x",
"/templates/v1/a/file",
] {
assert_eq!(
session.read_object(key, MAX_INDEX_BYTES).await,
Err(TemplateStoreError::Invalid)
);
assert_eq!(
session
.put_immutable(key, b"{}".to_vec(), "application/json")
.await,
Err(TemplateStoreError::Invalid)
);
}
assert_eq!(server.state.lock().unwrap().requests.len(), before);
session.finish().await.unwrap();
}
}