合并远程 master 分支更新
Project CI / AI game creator shell Rust shard 2/4 (pull_request) Successful in 6m55s
Project CI / AI game creator shell Rust shard 3/4 (pull_request) Successful in 6m41s
Project CI / AI game creator shell Rust shard 1/4 (pull_request) Successful in 7m10s
Project CI / AI game creator shell Rust shard 4/4 (pull_request) Successful in 4m37s
Project CI / AI game creator shell Rust smoke (pull_request) Successful in 1m59s
Project CI / AI game creator shell Rust crates (pull_request) Successful in 3m14s
Project CI / Frontend tests (pull_request) Successful in 6m53s
Project CI / Repository checks (pull_request) Successful in 7m7s
Project CI / Backend tests (pull_request) Successful in 11m20s
Project CI / Native shell tests (pull_request) Successful in 13m36s
Project CI / AI game creator shell web tests (pull_request) Failing after 7m11s
Project CI / AI game creator shell Rust shard 2/4 (pull_request) Successful in 6m55s
Project CI / AI game creator shell Rust shard 3/4 (pull_request) Successful in 6m41s
Project CI / AI game creator shell Rust shard 1/4 (pull_request) Successful in 7m10s
Project CI / AI game creator shell Rust shard 4/4 (pull_request) Successful in 4m37s
Project CI / AI game creator shell Rust smoke (pull_request) Successful in 1m59s
Project CI / AI game creator shell Rust crates (pull_request) Successful in 3m14s
Project CI / Frontend tests (pull_request) Successful in 6m53s
Project CI / Repository checks (pull_request) Successful in 7m7s
Project CI / Backend tests (pull_request) Successful in 11m20s
Project CI / Native shell tests (pull_request) Successful in 13m36s
Project CI / AI game creator shell web tests (pull_request) Failing after 7m11s
将 origin/master 最新变更合入策划提示词修复分支
This commit is contained in:
@@ -3002,6 +3002,29 @@ impl CodexAppServerConnection {
|
||||
self.inner.workspace_mode,
|
||||
direct_client_turn_id,
|
||||
);
|
||||
// 在发送前冻结本轮模型和归属;目录标识不能冒充上游返回的实际型号。
|
||||
// 记录失败仅留安全诊断,不阻断回合或重试付费请求。
|
||||
let _model_usage_guard =
|
||||
if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject {
|
||||
if backfill_project_model_usage_at(history_root, llm).is_err() {
|
||||
app_log!("project.model_usage.backfill_failed");
|
||||
}
|
||||
let context = ProjectModelUsageContext {
|
||||
root: history_root.to_path_buf(),
|
||||
client_turn_id: direct_tool_call_turn_id.clone(),
|
||||
thread_id: Some(thread_id.clone()),
|
||||
requested_model: model.to_string(),
|
||||
};
|
||||
if record_project_model_request_at(&context, llm.custom_enabled).is_err() {
|
||||
app_log!("project.model_usage.request_write_failed");
|
||||
}
|
||||
self.inner
|
||||
._provider_proxy
|
||||
.as_ref()
|
||||
.map(|proxy| proxy.begin_model_usage(context))
|
||||
} else {
|
||||
None
|
||||
};
|
||||
apply_game_creator_codex_app_server_reasoning_effort(&mut params, &request);
|
||||
if let Some(schema) = game_creator_codex_cli_tool_output_schema(&request) {
|
||||
params["outputSchema"] = schema;
|
||||
|
||||
@@ -4,7 +4,13 @@ use axum::http::{HeaderMap, HeaderName, Request, Response, StatusCode};
|
||||
use axum::routing::any;
|
||||
use axum::Router;
|
||||
use futures::StreamExt;
|
||||
use std::sync::Arc;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
mod model_usage;
|
||||
|
||||
use model_usage::ModelResponseObserver;
|
||||
|
||||
type ActiveModelUsage = Arc<Mutex<Option<Arc<crate::project::ProjectModelUsageContext>>>>;
|
||||
|
||||
pub(crate) const CODEX_PROVIDER_PROXY_PROTOCOL: &str = "genarrative-codex-provider-proxy.v1";
|
||||
|
||||
@@ -19,12 +25,34 @@ struct CodexProviderProxyState {
|
||||
downstream_bearer_token: String,
|
||||
main_site_upstream: bool,
|
||||
client: reqwest::Client,
|
||||
model_usage: ActiveModelUsage,
|
||||
}
|
||||
|
||||
pub(crate) struct CodexProviderProxy {
|
||||
base_url: String,
|
||||
downstream_bearer_token: String,
|
||||
task: tokio::task::JoinHandle<()>,
|
||||
model_usage: ActiveModelUsage,
|
||||
}
|
||||
|
||||
pub(crate) struct CodexProviderModelUsageGuard {
|
||||
active: ActiveModelUsage,
|
||||
registration: Arc<crate::project::ProjectModelUsageContext>,
|
||||
}
|
||||
|
||||
impl Drop for CodexProviderModelUsageGuard {
|
||||
fn drop(&mut self) {
|
||||
let mut active = self
|
||||
.active
|
||||
.lock()
|
||||
.unwrap_or_else(|error| error.into_inner());
|
||||
if active
|
||||
.as_ref()
|
||||
.is_some_and(|context| Arc::ptr_eq(context, &self.registration))
|
||||
{
|
||||
*active = None;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl CodexProviderProxy {
|
||||
@@ -35,6 +63,21 @@ impl CodexProviderProxy {
|
||||
pub(crate) fn downstream_bearer_token(&self) -> &str {
|
||||
&self.downstream_bearer_token
|
||||
}
|
||||
|
||||
pub(crate) fn begin_model_usage(
|
||||
&self,
|
||||
context: crate::project::ProjectModelUsageContext,
|
||||
) -> CodexProviderModelUsageGuard {
|
||||
let registration = Arc::new(context);
|
||||
*self
|
||||
.model_usage
|
||||
.lock()
|
||||
.unwrap_or_else(|error| error.into_inner()) = Some(Arc::clone(®istration));
|
||||
CodexProviderModelUsageGuard {
|
||||
active: Arc::clone(&self.model_usage),
|
||||
registration,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for CodexProviderProxy {
|
||||
@@ -128,6 +171,12 @@ async fn proxy_codex_provider_request(
|
||||
if request.method() != axum::http::Method::POST || request.uri().path() != "/responses" {
|
||||
return proxy_error(StatusCode::NOT_FOUND, "provider proxy route not found");
|
||||
}
|
||||
// 在读请求体或等待上游之前冻结归属,迟到响应不能使用下一回合的项目上下文。
|
||||
let model_usage = state
|
||||
.model_usage
|
||||
.lock()
|
||||
.unwrap_or_else(|error| error.into_inner())
|
||||
.clone();
|
||||
let upstream_url = format!("{}{path_and_query}", state.upstream_base_url);
|
||||
let (parts, body) = request.into_parts();
|
||||
let body = match to_bytes(body, CODEX_PROVIDER_PROXY_MAX_REQUEST_BYTES).await {
|
||||
@@ -171,9 +220,28 @@ async fn proxy_codex_provider_request(
|
||||
};
|
||||
let status = upstream.status();
|
||||
let upstream_headers = upstream.headers().clone();
|
||||
let stream = upstream
|
||||
.bytes_stream()
|
||||
.map(|chunk| chunk.map_err(|_| std::io::Error::other("provider response stream failed")));
|
||||
let observer = ModelResponseObserver::new(model_usage, status, &upstream_headers);
|
||||
let stream = futures::stream::unfold(
|
||||
(upstream.bytes_stream().boxed(), observer),
|
||||
|(mut upstream, mut observer)| async move {
|
||||
match upstream.next().await {
|
||||
Some(chunk) => {
|
||||
match &chunk {
|
||||
Ok(bytes) => observer.observe(bytes),
|
||||
Err(_) => observer.failed(),
|
||||
}
|
||||
Some((
|
||||
chunk.map_err(|_| std::io::Error::other("provider response stream failed")),
|
||||
(upstream, observer),
|
||||
))
|
||||
}
|
||||
None => {
|
||||
observer.finish();
|
||||
None
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
let mut response = Response::builder().status(status);
|
||||
if let Some(headers) = response.headers_mut() {
|
||||
let mut stripped_limit_headers = 0_usize;
|
||||
@@ -227,12 +295,14 @@ pub(crate) async fn start_codex_provider_proxy(
|
||||
let address = listener
|
||||
.local_addr()
|
||||
.map_err(|error| format!("读取 Codex Provider 代理地址失败:{error}"))?;
|
||||
let model_usage = Arc::new(Mutex::new(None));
|
||||
let state = Arc::new(CodexProviderProxyState {
|
||||
upstream_base_url,
|
||||
upstream_bearer_token: upstream_bearer_token.to_string(),
|
||||
downstream_bearer_token: downstream_bearer_token.clone(),
|
||||
main_site_upstream,
|
||||
client,
|
||||
model_usage: Arc::clone(&model_usage),
|
||||
});
|
||||
let app = Router::new()
|
||||
.fallback(any(proxy_codex_provider_request))
|
||||
@@ -244,6 +314,7 @@ pub(crate) async fn start_codex_provider_proxy(
|
||||
base_url: format!("http://127.0.0.1:{}", address.port()),
|
||||
downstream_bearer_token,
|
||||
task,
|
||||
model_usage,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -253,6 +324,318 @@ mod tests {
|
||||
use axum::routing::post;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
#[derive(Clone)]
|
||||
struct ModelFixture {
|
||||
status: StatusCode,
|
||||
content_type: &'static str,
|
||||
bytes: Vec<u8>,
|
||||
entered: Option<Arc<tokio::sync::Notify>>,
|
||||
release: Option<Arc<tokio::sync::Notify>>,
|
||||
}
|
||||
|
||||
async fn model_fixture_response(State(fixture): State<ModelFixture>) -> Response<Body> {
|
||||
if let Some(entered) = fixture.entered {
|
||||
entered.notify_one();
|
||||
}
|
||||
if let Some(release) = fixture.release {
|
||||
release.notified().await;
|
||||
}
|
||||
// 刻意拆开 UTF-8 与 CRLF;代理仍必须逐字节保留完整响应。
|
||||
let chunks = fixture
|
||||
.bytes
|
||||
.into_iter()
|
||||
.map(|byte| Ok::<_, std::io::Error>(axum::body::Bytes::from(vec![byte])));
|
||||
Response::builder()
|
||||
.status(fixture.status)
|
||||
.header("content-type", fixture.content_type)
|
||||
.body(Body::from_stream(futures::stream::iter(chunks)))
|
||||
.expect("model fixture response")
|
||||
}
|
||||
|
||||
async fn start_model_fixture(
|
||||
fixture: ModelFixture,
|
||||
) -> (CodexProviderProxy, tokio::task::JoinHandle<()>) {
|
||||
let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
|
||||
.await
|
||||
.expect("bind model upstream");
|
||||
let address = listener.local_addr().expect("model upstream address");
|
||||
let app = Router::new()
|
||||
.route("/responses", post(model_fixture_response))
|
||||
.with_state(fixture);
|
||||
let task = tokio::spawn(async move {
|
||||
let _ = axum::serve(listener, app).await;
|
||||
});
|
||||
let proxy = start_codex_provider_proxy(
|
||||
&format!("http://127.0.0.1:{}", address.port()),
|
||||
"fixture-provider-key",
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.expect("start model provider proxy");
|
||||
(proxy, task)
|
||||
}
|
||||
|
||||
fn model_context(
|
||||
root: &std::path::Path,
|
||||
turn: &str,
|
||||
) -> crate::project::ProjectModelUsageContext {
|
||||
crate::project::ProjectModelUsageContext {
|
||||
root: root.to_path_buf(),
|
||||
client_turn_id: Some(turn.to_string()),
|
||||
thread_id: Some("thread_fixture".to_string()),
|
||||
requested_model: "requested-alias".to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
fn model_records(root: &std::path::Path) -> Vec<serde_json::Value> {
|
||||
let path = root.join(".agent/model-usage.jsonl");
|
||||
if !path.exists() {
|
||||
return Vec::new();
|
||||
}
|
||||
std::fs::read_to_string(path)
|
||||
.expect("model usage file")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str(line).expect("model usage json"))
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn wait_for_model_records(
|
||||
root: &std::path::Path,
|
||||
count: usize,
|
||||
) -> Vec<serde_json::Value> {
|
||||
tokio::time::timeout(std::time::Duration::from_secs(5), async {
|
||||
loop {
|
||||
// 写入线程可能正写一行,只有完整 JSONL 才是测试的落盘证据。
|
||||
if let Ok(text) = std::fs::read_to_string(root.join(".agent/model-usage.jsonl")) {
|
||||
let records: Result<Vec<serde_json::Value>, _> =
|
||||
text.lines().map(serde_json::from_str).collect();
|
||||
if let Ok(records) = records {
|
||||
if records.len() >= count {
|
||||
return records;
|
||||
}
|
||||
}
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("model records persisted")
|
||||
}
|
||||
|
||||
async fn fetch_model_fixture(proxy: &CodexProviderProxy) -> Vec<u8> {
|
||||
reqwest::Client::new()
|
||||
.post(format!("{}/responses", proxy.base_url()))
|
||||
.bearer_auth(proxy.downstream_bearer_token())
|
||||
.body("{\"input\":\"private prompt fixture\"}")
|
||||
.send()
|
||||
.await
|
||||
.expect("model fixture response")
|
||||
.bytes()
|
||||
.await
|
||||
.expect("model fixture body")
|
||||
.to_vec()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_usage_proxy_records_json_model_and_passes_original_bytes() {
|
||||
let root = tempfile::tempdir().expect("model project");
|
||||
let bytes = br#"{"id":"resp_json","model":"actual-main-v1","output":[{"text":"private output fixture","model":"not-main"}]}"#.to_vec();
|
||||
let (proxy, task) = start_model_fixture(ModelFixture {
|
||||
status: StatusCode::OK,
|
||||
content_type: "application/json; charset=utf-8",
|
||||
bytes: bytes.clone(),
|
||||
entered: None,
|
||||
release: None,
|
||||
})
|
||||
.await;
|
||||
let _guard = proxy.begin_model_usage(model_context(root.path(), "turn_json"));
|
||||
assert_eq!(fetch_model_fixture(&proxy).await, bytes);
|
||||
let records = wait_for_model_records(root.path(), 1).await;
|
||||
assert_eq!(records.len(), 1);
|
||||
assert_eq!(records[0]["requestedModel"], "requested-alias");
|
||||
assert_eq!(records[0]["modelName"], "actual-main-v1");
|
||||
assert_eq!(records[0]["clientTurnId"], "turn_json");
|
||||
assert_eq!(records[0]["responseId"], "resp_json");
|
||||
assert_eq!(records[0]["source"], "provider-response");
|
||||
assert_eq!(records[0]["historicalModelConfirmed"], true);
|
||||
let persisted = serde_json::to_string(&records).expect("serialize records");
|
||||
for forbidden in [
|
||||
"private prompt",
|
||||
"private output",
|
||||
"fixture-provider-key",
|
||||
"not-main",
|
||||
] {
|
||||
assert!(!persisted.contains(forbidden));
|
||||
}
|
||||
task.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_usage_proxy_records_sse_models_once_and_passes_original_bytes() {
|
||||
let root = tempfile::tempdir().expect("model project");
|
||||
let bytes = concat!(
|
||||
"event: response.created\r\n",
|
||||
"data: {\"type\":\"response.created\",\r\n",
|
||||
"data: \"response\":{\"id\":\"resp_sse\",\"model\":\"actual-main-v2\",\"text\":\"隐私正文\"}}\r\n\r\n",
|
||||
"data: {\"type\":\"response.in_progress\",\"response\":{\"id\":\"resp_sse\",\"model\":\"actual-main-v2\"}}\n\n",
|
||||
"data: {\"type\":\"response.output_text.delta\",\"response\":{\"model\":\"not-main\"}}\n\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_sse\",\"model\":\"actual-main-v3\"}}\n\n",
|
||||
"data: [DONE]\n\n"
|
||||
).as_bytes().to_vec();
|
||||
let (proxy, task) = start_model_fixture(ModelFixture {
|
||||
status: StatusCode::OK,
|
||||
content_type: "text/event-stream",
|
||||
bytes: bytes.clone(),
|
||||
entered: None,
|
||||
release: None,
|
||||
})
|
||||
.await;
|
||||
let _guard = proxy.begin_model_usage(model_context(root.path(), "turn_sse"));
|
||||
assert_eq!(fetch_model_fixture(&proxy).await, bytes);
|
||||
let records = wait_for_model_records(root.path(), 2).await;
|
||||
assert_eq!(records.len(), 2);
|
||||
assert_eq!(records[0]["modelName"], "actual-main-v2");
|
||||
assert_eq!(records[1]["modelName"], "actual-main-v3");
|
||||
assert!(records
|
||||
.iter()
|
||||
.all(|record| record["responseId"] == "resp_sse"));
|
||||
assert!(!serde_json::to_string(&records)
|
||||
.unwrap()
|
||||
.contains("隐私正文"));
|
||||
task.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_usage_proxy_freezes_request_owner_and_old_guard_preserves_new_owner() {
|
||||
let old_root = tempfile::tempdir().expect("old model project");
|
||||
let new_root = tempfile::tempdir().expect("new model project");
|
||||
let entered = Arc::new(tokio::sync::Notify::new());
|
||||
let release = Arc::new(tokio::sync::Notify::new());
|
||||
let bytes = br#"{"id":"resp_late","model":"actual-late"}"#.to_vec();
|
||||
let (proxy, task) = start_model_fixture(ModelFixture {
|
||||
status: StatusCode::OK,
|
||||
content_type: "application/json",
|
||||
bytes: bytes.clone(),
|
||||
entered: Some(Arc::clone(&entered)),
|
||||
release: Some(Arc::clone(&release)),
|
||||
})
|
||||
.await;
|
||||
let old_guard = proxy.begin_model_usage(model_context(old_root.path(), "turn_old"));
|
||||
let proxy = Arc::new(proxy);
|
||||
let requester = Arc::clone(&proxy);
|
||||
let request = tokio::spawn(async move { fetch_model_fixture(&requester).await });
|
||||
tokio::time::timeout(std::time::Duration::from_secs(5), entered.notified())
|
||||
.await
|
||||
.expect("upstream entered");
|
||||
let new_guard = proxy.begin_model_usage(model_context(new_root.path(), "turn_new"));
|
||||
drop(old_guard);
|
||||
release.notify_one();
|
||||
assert_eq!(request.await.expect("old response"), bytes);
|
||||
assert_eq!(
|
||||
wait_for_model_records(old_root.path(), 1).await[0]["clientTurnId"],
|
||||
"turn_old"
|
||||
);
|
||||
assert!(model_records(new_root.path()).is_empty());
|
||||
|
||||
release.notify_one();
|
||||
assert_eq!(fetch_model_fixture(&proxy).await, bytes);
|
||||
assert_eq!(
|
||||
wait_for_model_records(new_root.path(), 1).await[0]["clientTurnId"],
|
||||
"turn_new"
|
||||
);
|
||||
drop(new_guard);
|
||||
assert!(proxy.model_usage.lock().expect("active model").is_none());
|
||||
task.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_usage_proxy_passes_response_while_record_file_is_locked() {
|
||||
let root = tempfile::tempdir().expect("locked model project");
|
||||
std::fs::create_dir(root.path().join(".agent")).expect("create agent directory");
|
||||
let path =
|
||||
crate::project::resolve_local_project_path(root.path(), ".agent/model-usage.jsonl")
|
||||
.expect("model usage path");
|
||||
let lock = crate::project::project_append_lock_for(&path).expect("model usage lock");
|
||||
let held = lock
|
||||
.lock("model usage fixture")
|
||||
.expect("hold model usage lock");
|
||||
let bytes = br#"{"id":"resp_locked","model":"actual-main"}"#.to_vec();
|
||||
let (proxy, task) = start_model_fixture(ModelFixture {
|
||||
status: StatusCode::OK,
|
||||
content_type: "application/json",
|
||||
bytes: bytes.clone(),
|
||||
entered: None,
|
||||
release: None,
|
||||
})
|
||||
.await;
|
||||
let _guard = proxy.begin_model_usage(model_context(root.path(), "turn_locked"));
|
||||
let response = tokio::time::timeout(
|
||||
std::time::Duration::from_secs(2),
|
||||
fetch_model_fixture(&proxy),
|
||||
)
|
||||
.await
|
||||
.expect("record lock must not block the response");
|
||||
assert_eq!(response, bytes);
|
||||
assert!(model_records(root.path()).is_empty());
|
||||
drop(held);
|
||||
assert_eq!(
|
||||
wait_for_model_records(root.path(), 1).await[0]["modelName"],
|
||||
"actual-main"
|
||||
);
|
||||
task.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_usage_observer_keeps_confirmed_sse_on_stream_error_but_rejects_partial_json() {
|
||||
let root = tempfile::tempdir().expect("stream error project");
|
||||
let context = Arc::new(model_context(root.path(), "turn_failed"));
|
||||
let mut headers = HeaderMap::new();
|
||||
headers.insert("content-type", "text/event-stream".parse().unwrap());
|
||||
let mut observer =
|
||||
ModelResponseObserver::new(Some(Arc::clone(&context)), StatusCode::OK, &headers);
|
||||
observer.observe(b"data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp_early\",\"model\":\"confirmed-before-error\"}}\n\n");
|
||||
observer.failed();
|
||||
observer.observe(
|
||||
b"data: {\"type\":\"response.completed\",\"response\":{\"model\":\"after-error\"}}\n\n",
|
||||
);
|
||||
observer.finish();
|
||||
let records = wait_for_model_records(root.path(), 1).await;
|
||||
assert_eq!(records.len(), 1);
|
||||
assert_eq!(records[0]["modelName"], "confirmed-before-error");
|
||||
|
||||
headers.insert("content-type", "application/json".parse().unwrap());
|
||||
let mut observer =
|
||||
ModelResponseObserver::new(Some(Arc::clone(&context)), StatusCode::OK, &headers);
|
||||
observer.observe(br#"{"id":"resp_partial","model":"partial-json"}"#);
|
||||
observer.failed();
|
||||
observer.finish();
|
||||
assert_eq!(model_records(root.path()).len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn model_usage_proxy_does_not_invent_model_for_error_or_missing_field() {
|
||||
for (status, content_type, bytes) in [
|
||||
(StatusCode::BAD_REQUEST, "application/json", br#"{"model":"error-model"}"#.to_vec()),
|
||||
(StatusCode::OK, "application/json", br#"{"output":[{"model":"nested-model"}]}"#.to_vec()),
|
||||
(StatusCode::OK, "application/json", br#"{"model":"incomplete""#.to_vec()),
|
||||
(StatusCode::OK, "text/event-stream", b"data: {\"type\":\"response.failed\",\"response\":{\"model\":\"failed-model\"}}\n\n".to_vec()),
|
||||
] {
|
||||
let root = tempfile::tempdir().expect("model project");
|
||||
let (proxy, task) = start_model_fixture(ModelFixture {
|
||||
status,
|
||||
content_type,
|
||||
bytes: bytes.clone(),
|
||||
entered: None,
|
||||
release: None,
|
||||
})
|
||||
.await;
|
||||
let _guard = proxy.begin_model_usage(model_context(root.path(), "turn_empty"));
|
||||
assert_eq!(fetch_model_fixture(&proxy).await, bytes);
|
||||
assert!(model_records(root.path()).is_empty());
|
||||
task.abort();
|
||||
}
|
||||
}
|
||||
|
||||
async fn fake_upstream(
|
||||
State(calls): State<Arc<AtomicUsize>>,
|
||||
headers: HeaderMap,
|
||||
|
||||
@@ -0,0 +1,320 @@
|
||||
use axum::http::{HeaderMap, StatusCode};
|
||||
use serde_json::Value;
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
|
||||
const MAX_OBSERVATION_BYTES: usize = 1024 * 1024;
|
||||
const MAX_MODELS_PER_RESPONSE: usize = 64;
|
||||
const MAX_IDENTIFIER_BYTES: usize = 256;
|
||||
|
||||
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
|
||||
struct ModelObservation {
|
||||
model: String,
|
||||
response_id: Option<String>,
|
||||
}
|
||||
|
||||
enum ResponseParser {
|
||||
Json(Vec<u8>),
|
||||
Sse(SseParser),
|
||||
Disabled,
|
||||
}
|
||||
|
||||
pub(super) struct ModelResponseObserver {
|
||||
context: Option<Arc<crate::project::ProjectModelUsageContext>>,
|
||||
parser: ResponseParser,
|
||||
seen: HashSet<ModelObservation>,
|
||||
writer: Option<tokio::sync::mpsc::Sender<ModelObservation>>,
|
||||
}
|
||||
|
||||
impl ModelResponseObserver {
|
||||
pub(super) fn new(
|
||||
context: Option<Arc<crate::project::ProjectModelUsageContext>>,
|
||||
status: StatusCode,
|
||||
headers: &HeaderMap,
|
||||
) -> Self {
|
||||
let content_type = headers
|
||||
.get("content-type")
|
||||
.and_then(|value| value.to_str().ok())
|
||||
.unwrap_or_default()
|
||||
.split(';')
|
||||
.next()
|
||||
.unwrap_or_default()
|
||||
.trim();
|
||||
let parser = if !status.is_success() || context.is_none() {
|
||||
ResponseParser::Disabled
|
||||
} else if content_type.eq_ignore_ascii_case("text/event-stream") {
|
||||
ResponseParser::Sse(SseParser::default())
|
||||
} else if content_type.eq_ignore_ascii_case("application/json") {
|
||||
ResponseParser::Json(Vec::new())
|
||||
} else {
|
||||
ResponseParser::Disabled
|
||||
};
|
||||
Self {
|
||||
context,
|
||||
parser,
|
||||
seen: HashSet::new(),
|
||||
writer: None,
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn observe(&mut self, bytes: &[u8]) {
|
||||
match &mut self.parser {
|
||||
ResponseParser::Json(buffer) => {
|
||||
if bytes.len() > MAX_OBSERVATION_BYTES.saturating_sub(buffer.len()) {
|
||||
self.parser = ResponseParser::Disabled;
|
||||
} else {
|
||||
buffer.extend_from_slice(bytes);
|
||||
}
|
||||
}
|
||||
ResponseParser::Sse(parser) => {
|
||||
// 逐字节识别换行,UTF-8 只在完整事件中解析;不改写转发的原始块。
|
||||
for byte in bytes {
|
||||
if let Some(observation) = parser.push(*byte) {
|
||||
Self::record(&self.context, &mut self.seen, &mut self.writer, observation);
|
||||
if self.seen.len() >= MAX_MODELS_PER_RESPONSE {
|
||||
self.parser = ResponseParser::Disabled;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
ResponseParser::Disabled => {}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn failed(&mut self) {
|
||||
// SSE 已确认的初始事件已排队写入;不完整 JSON 不构成型号证据。
|
||||
self.parser = ResponseParser::Disabled;
|
||||
}
|
||||
|
||||
pub(super) fn finish(&mut self) {
|
||||
if let ResponseParser::Json(bytes) =
|
||||
std::mem::replace(&mut self.parser, ResponseParser::Disabled)
|
||||
{
|
||||
if let Ok(value) = serde_json::from_slice::<Value>(&bytes) {
|
||||
if let Some(observation) = response_observation(&value) {
|
||||
Self::record(&self.context, &mut self.seen, &mut self.writer, observation);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn record(
|
||||
context: &Option<Arc<crate::project::ProjectModelUsageContext>>,
|
||||
seen: &mut HashSet<ModelObservation>,
|
||||
writer: &mut Option<tokio::sync::mpsc::Sender<ModelObservation>>,
|
||||
observation: ModelObservation,
|
||||
) {
|
||||
if seen.len() >= MAX_MODELS_PER_RESPONSE || !seen.insert(observation.clone()) {
|
||||
return;
|
||||
}
|
||||
let Some(context) = context else {
|
||||
return;
|
||||
};
|
||||
let writer = writer.get_or_insert_with(|| {
|
||||
let (sender, mut receiver) =
|
||||
tokio::sync::mpsc::channel::<ModelObservation>(MAX_MODELS_PER_RESPONSE);
|
||||
let context = Arc::clone(context);
|
||||
// 单响应顺序写入;项目追加锁和磁盘 I/O 不得阻塞上游响应转发。
|
||||
tokio::spawn(async move {
|
||||
while let Some(observation) = receiver.recv().await {
|
||||
let context = Arc::clone(&context);
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
crate::project::record_project_model_response_at(
|
||||
&context,
|
||||
&observation.model,
|
||||
observation.response_id.as_deref(),
|
||||
)
|
||||
})
|
||||
.await;
|
||||
if !matches!(result, Ok(Ok(()))) {
|
||||
app_log!("agent.direct_codex.model_usage.response_record_failed");
|
||||
}
|
||||
}
|
||||
});
|
||||
sender
|
||||
});
|
||||
if writer.try_send(observation).is_err() {
|
||||
app_log!("agent.direct_codex.model_usage.response_record_queue_unavailable");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn identifier(value: Option<&Value>) -> Option<String> {
|
||||
let value = value?.as_str()?.trim();
|
||||
(!value.is_empty()
|
||||
&& value.len() <= MAX_IDENTIFIER_BYTES
|
||||
&& !value.chars().any(char::is_control))
|
||||
.then(|| value.to_string())
|
||||
}
|
||||
|
||||
fn response_observation(value: &Value) -> Option<ModelObservation> {
|
||||
Some(ModelObservation {
|
||||
model: identifier(value.get("model"))?,
|
||||
response_id: identifier(value.get("id")),
|
||||
})
|
||||
}
|
||||
|
||||
fn known_event(value: &str) -> bool {
|
||||
matches!(
|
||||
value,
|
||||
"response.created" | "response.in_progress" | "response.completed"
|
||||
)
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct SseParser {
|
||||
line: Vec<u8>,
|
||||
line_bytes: usize,
|
||||
data: Vec<u8>,
|
||||
event: Option<String>,
|
||||
event_bytes: usize,
|
||||
discarded: bool,
|
||||
skip_lf: bool,
|
||||
}
|
||||
|
||||
impl SseParser {
|
||||
fn push(&mut self, byte: u8) -> Option<ModelObservation> {
|
||||
if self.skip_lf {
|
||||
self.skip_lf = false;
|
||||
if byte == b'\n' {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
if byte == b'\r' || byte == b'\n' {
|
||||
self.skip_lf = byte == b'\r';
|
||||
return self.finish_line();
|
||||
}
|
||||
self.line_bytes = self.line_bytes.saturating_add(1);
|
||||
self.event_bytes = self.event_bytes.saturating_add(1);
|
||||
if self.event_bytes > MAX_OBSERVATION_BYTES {
|
||||
self.discarded = true;
|
||||
self.line.clear();
|
||||
self.data.clear();
|
||||
self.event = None;
|
||||
} else if !self.discarded {
|
||||
self.line.push(byte);
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn finish_line(&mut self) -> Option<ModelObservation> {
|
||||
if self.line_bytes == 0 {
|
||||
let observation = (!self.discarded).then(|| self.parse_event()).flatten();
|
||||
self.line.clear();
|
||||
self.data.clear();
|
||||
self.event = None;
|
||||
self.event_bytes = 0;
|
||||
self.discarded = false;
|
||||
return observation;
|
||||
}
|
||||
self.line_bytes = 0;
|
||||
if !self.discarded {
|
||||
// 换行同样计入事件预算,避免无限 data 空行绕过有界缓冲。
|
||||
self.event_bytes = self.event_bytes.saturating_add(1);
|
||||
if self.event_bytes > MAX_OBSERVATION_BYTES {
|
||||
self.discarded = true;
|
||||
self.data.clear();
|
||||
self.event = None;
|
||||
} else if let Some((name, value)) = self
|
||||
.line
|
||||
.iter()
|
||||
.position(|byte| *byte == b':')
|
||||
.map(|colon| (&self.line[..colon], &self.line[colon + 1..]))
|
||||
{
|
||||
let value = value.strip_prefix(b" ").unwrap_or(value);
|
||||
match name {
|
||||
b"data" => {
|
||||
self.data.extend_from_slice(value);
|
||||
self.data.push(b'\n');
|
||||
}
|
||||
b"event" => {
|
||||
self.event = Some(
|
||||
std::str::from_utf8(value)
|
||||
.ok()
|
||||
.filter(|value| known_event(value))
|
||||
.unwrap_or("")
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
self.line.clear();
|
||||
None
|
||||
}
|
||||
|
||||
fn parse_event(&self) -> Option<ModelObservation> {
|
||||
let value: Value = serde_json::from_slice(&self.data).ok()?;
|
||||
let payload_type = value.get("type").and_then(Value::as_str);
|
||||
let event = payload_type.or(self.event.as_deref())?;
|
||||
if !known_event(event)
|
||||
|| self
|
||||
.event
|
||||
.as_deref()
|
||||
.is_some_and(|header_event| header_event != event)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
response_observation(value.get("response")?)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn parse_sse(bytes: &[u8]) -> Vec<ModelObservation> {
|
||||
let mut parser = SseParser::default();
|
||||
bytes.iter().filter_map(|byte| parser.push(*byte)).collect()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn model_usage_parser_accepts_split_utf8_crlf_multiline_and_cr_events() {
|
||||
let event = concat!(
|
||||
": keepalive\r\n",
|
||||
"event: response.created\r\n",
|
||||
"data: {\"type\":\"response.created\",\r\n",
|
||||
"data: \"response\":{\"id\":\"resp_1\",\"model\":\"模型-v1\"}}\r\n\r\n",
|
||||
"event: response.completed\r",
|
||||
"data: {\"response\":{\"id\":\"resp_1\",\"model\":\"模型-v1\"}}\r\r"
|
||||
);
|
||||
let parsed = parse_sse(event.as_bytes());
|
||||
assert_eq!(parsed.len(), 2);
|
||||
assert_eq!(parsed[0].model, "模型-v1");
|
||||
assert_eq!(parsed[0].response_id.as_deref(), Some("resp_1"));
|
||||
assert_eq!(parsed[0], parsed[1]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn model_usage_parser_rejects_untrusted_locations_and_invalid_events() {
|
||||
for event in [
|
||||
"data: {\"type\":\"response.output_text.delta\",\"response\":{\"model\":\"fake\"}}\n\n",
|
||||
"data: {\"type\":\"response.failed\",\"response\":{\"model\":\"fake\"}}\n\n",
|
||||
"data: {\"type\":\"response.created\",\"model\":\"fake\"}\n\n",
|
||||
"data: {\"type\":\"response.created\",\"response\":{\"output\":[{\"model\":\"fake\"}]}}\n\n",
|
||||
"event: response.failed\ndata: {\"type\":\"response.created\",\"response\":{\"model\":\"fake\"}}\n\n",
|
||||
"data: [DONE]\n\n",
|
||||
"data: invalid json\n\n",
|
||||
"data: {\"type\":\"response.created\",\"response\":{\"model\":\"fake\"}}",
|
||||
] {
|
||||
assert!(parse_sse(event.as_bytes()).is_empty());
|
||||
}
|
||||
assert!(
|
||||
response_observation(&serde_json::json!({"output": [{"model": "fake"}]})).is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn model_usage_parser_recovers_after_oversize_and_invalid_utf8_events() {
|
||||
let valid = b"data: {\"type\":\"response.created\",\"response\":{\"model\":\"real\"}}\n\n";
|
||||
let mut bytes = b"data: ".to_vec();
|
||||
bytes.extend(std::iter::repeat_n(b'x', MAX_OBSERVATION_BYTES + 1));
|
||||
bytes.extend_from_slice(b"\n\ndata: \xff\n\n");
|
||||
bytes.extend_from_slice(valid);
|
||||
let parsed = parse_sse(&bytes);
|
||||
assert_eq!(parsed.len(), 1);
|
||||
assert_eq!(parsed[0].model, "real");
|
||||
}
|
||||
}
|
||||
@@ -981,7 +981,18 @@ pub(crate) fn get_local_game_manifest_sync(
|
||||
return Err(format!("不支持通过 manifest 执行命令:{command_id}"));
|
||||
}
|
||||
enforce_project_permission_policy(root, command_id)?;
|
||||
read_manifest_for_project_with_godot_root_calibration(root)
|
||||
let manifest = read_manifest_for_project_with_godot_root_calibration(root)?;
|
||||
if command_id == "project.status" {
|
||||
match load_game_creator_app_config() {
|
||||
Ok(config) => {
|
||||
if backfill_project_model_usage_at(root, &config.llm).is_err() {
|
||||
app_log!("project.model_usage.backfill_failed");
|
||||
}
|
||||
}
|
||||
Err(_) => app_log!("project.model_usage.backfill_config_unavailable"),
|
||||
}
|
||||
}
|
||||
Ok(manifest)
|
||||
}
|
||||
|
||||
/// 读取资源画布的持久化布局(`.agent/workbench/resource-layouts/*.json`)。
|
||||
|
||||
@@ -14,6 +14,7 @@ mod external_editor_bindings;
|
||||
mod filesystem;
|
||||
mod manifest;
|
||||
mod memory;
|
||||
mod model_usage;
|
||||
mod resource_dependency_graph;
|
||||
mod resource_editor;
|
||||
mod resource_layout;
|
||||
@@ -32,6 +33,7 @@ pub(crate) use external_editor_bindings::*;
|
||||
pub(crate) use filesystem::*;
|
||||
pub(crate) use manifest::*;
|
||||
pub(crate) use memory::*;
|
||||
pub(crate) use model_usage::*;
|
||||
pub(crate) use resource_dependency_graph::*;
|
||||
pub(crate) use resource_editor::*;
|
||||
pub(crate) use resource_layout::*;
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -16,6 +16,8 @@
|
||||
|
||||
## 开发中
|
||||
|
||||
- AGC 主模型追溯保存在项目 `.agent/model-usage.jsonl`,请求目录标识与响应确认的型号分别记录;旧项目当前配置补录必须标注来源,不冒充历史事实。仅保存有界模型与回合身份字段,不保存配置、凭据或对话正文,不增加 UI 展示。详见 AGC 实施计划“项目主模型使用记录”。
|
||||
|
||||
- Agent 提示词正文与工具说明放在所属组件的 `prompts/`;AGC 通过现有 Prompt Bundle 编译加载,服务端独立 crate 编译包含自己的提示词文件。代码负责变量填充、结构化 schema 与执行校验。
|
||||
|
||||
- 策划 Agent 的顾问态由用户指示驱动,不自主推进项目、主动安排下一步或提交阶段审批;完成单次请求不结束顾问态。五个策划阶段的审批用于检阅已完成产物,关键选择先问询;过程文档按需记录且不重复正式正文。顶层设计按需保留易混淆方向及排除理由,提示词精简应保留这些行为与设计边界。详见策划 Agent 生产迁移与工作区浏览方案。
|
||||
|
||||
@@ -1,5 +1,23 @@
|
||||
# AI 游戏创作智能体 App 实施计划
|
||||
|
||||
## 项目主模型使用记录
|
||||
|
||||
- 项目在 `.agent/model-usage.jsonl` 保存主模型记录,供复制或压缩完整项目后查询,不新增界面展示。每条只保存版本、记录时间、来源、请求模型标识、可确认的模型名称,以及可用的客户端回合和线程身份;不保存提示词、响应正文、连接地址或凭据。
|
||||
- DirectProject 每次提交回合前保存本次请求的模型快照,来源为 `turn-request`。官方目录标识与真实型号分开:官方标识只写 `requestedModel`,不能把 `platform-default`、目录 ID、`Direct Codex` 或 `codex-app-server` 当成模型名;自定义路由保存明确提交的型号。响应返回的模型名以 `provider-response` 追加,使用请求发出时冻结的项目/回合归属,后续切模型不会覆盖历史。
|
||||
- 模型名从成功 Responses JSON 的顶层 `model` 或 SSE `response.created` / `response.in_progress` / `response.completed` 中的 `response.model` 提取;只读取白名单字段,有界处理分块和异常数据,响应内容仍原样转发。上游没有返回型号时保留请求证据,不推测实际模型。
|
||||
- 旧项目通过 `project.status` 读取 manifest 成功后补录一次:在有界范围内优先恢复项目内可明确识别的主模型运行记录,超出历史扫描预算时跳过对应候选;没有找到可信历史型号时写入当前配置快照并标注 `current-config-backfill`,不声称其为历史事实。官方配置只有目录标识时保留标识,模型名为空;以后真实回合仍继续追加。补录幂等,不覆盖或改写原项目历史。
|
||||
- 记录复用项目受控路径与追加锁;无写权限、损坏或超限时保留已有文件并写安全诊断,不因此中断项目打开、模型响应或触发额外付费重试。响应记录在阻塞工作线程中追加,文件锁等待不阻塞响应原始块转发。补录仅发生在项目 `project.status` 命令和本次回合入口,不扫描机器上的其它项目或私人会话。
|
||||
- 验收覆盖:请求与响应的模型差异、跨回合/项目隔离、切模型保留历史、旧项目补录及幂等、分块 SSE/JSON、不含模型/失败响应不伪造型号,以及凭据/正文不落盘。使用本地 HTTP fixture 验证代理透传和落盘,真实供应商调用另行报告。
|
||||
|
||||
| 合同 | 自动化证据入口 |
|
||||
| --- | --- |
|
||||
| 来源区分、切模型保留历史、补录幂等、损坏/超限保护 | `project::model_usage::tests` |
|
||||
| JSON/SSE 字节透传、模型落盘与跨回合归属 | `agent::codex_provider_proxy::tests::model_usage_proxy_*` |
|
||||
| 分块 UTF-8、换行、多行事件和有界解析 | `agent::codex_provider_proxy::model_usage::tests` |
|
||||
| 文件锁争用不阻塞响应、流中断仍保留已确认型号 | Provider 代理的锁争用与流错误 fixture |
|
||||
|
||||
记录格式版本为 `schemaVersion: 1`。`recordedAtMs` 是记录时间,`historicalModelConfirmed` 仅在响应观测或可信历史恢复时为 `true`;它在请求快照和当前配置补录时为 `false`。上述本地 fixture 不替代真实供应商或安装包验收。
|
||||
|
||||
## 2026-09-17 GameCreationApp 资源 kind:唯一词汇表、严格解析与 `app_log!` 留痕
|
||||
|
||||
本节覆盖 2026-09-15 节里关于「canonical 字符串列表 / legacy 别名表 / `tracing` 留痕 / ts-rs 生成路径」的表述;枚举成员集合、「不迁移、不静默转换」的总体口径不变。
|
||||
|
||||
Reference in New Issue
Block a user