1e186369c9
Project CI / AI game creator shell Rust crates (push) Successful in 1m24s
Project CI / AI game creator shell Rust smoke (push) Successful in 1m56s
Project CI / AI game creator shell Rust lane 1/2 (push) Has been cancelled
Project CI / Frontend tests (push) Has been cancelled
Project CI / Backend tests (push) Has been cancelled
Project CI / Repository checks (push) Has been cancelled
Project CI / AI game creator shell web tests (push) Has been cancelled
Project CI / AI game creator shell Rust lane 2/2 (push) Has been cancelled
Project CI / Native shell tests (push) Has been cancelled
Reviewed-on: https://git.genarrative.world/git/GenarrativeAI/Genarrative/pulls/446 Co-authored-by: Linghong <ink29535@proton.me> Co-committed-by: Linghong <ink29535@proton.me>
229 lines
7.4 KiB
Rust
229 lines
7.4 KiB
Rust
//! 独立后台上传轮次;网络等待不占用 writer 或业务锁。
|
|
use super::{contract::Route, store};
|
|
use crate::platform_session::{current_platform_session, PlatformSessionSnapshot};
|
|
use std::fs::{File, OpenOptions};
|
|
use std::path::{Path, PathBuf};
|
|
use std::sync::atomic::{AtomicU64, Ordering};
|
|
use std::time::{Duration, Instant};
|
|
|
|
const INTERVAL: Duration = Duration::from_secs(15 * 60);
|
|
const REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
|
|
const CYCLE_BUDGET: Duration = Duration::from_secs(120);
|
|
const MAX_BATCHES: usize = 20;
|
|
static FAILURES: AtomicU64 = AtomicU64::new(0);
|
|
|
|
fn failure() {
|
|
let _ = FAILURES.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| {
|
|
Some(n.saturating_add(1))
|
|
});
|
|
}
|
|
|
|
pub(super) fn start(config_dir: PathBuf) {
|
|
tauri::async_runtime::spawn(async move {
|
|
let Ok(client) = reqwest::Client::builder()
|
|
.redirect(reqwest::redirect::Policy::none())
|
|
.timeout(REQUEST_TIMEOUT)
|
|
.build()
|
|
else {
|
|
return;
|
|
};
|
|
let root = config_dir.join("analytics");
|
|
let mut next = tokio::time::Instant::now() + INTERVAL;
|
|
loop {
|
|
tokio::time::sleep_until(next).await;
|
|
// 休眠恢复只跑一次,下一轮从当前时间重新计时。
|
|
if advance_schedule(tokio::time::Instant::now(), &mut next) {
|
|
cycle(&root, &client).await;
|
|
}
|
|
}
|
|
});
|
|
}
|
|
|
|
fn advance_schedule(now: tokio::time::Instant, next: &mut tokio::time::Instant) -> bool {
|
|
if now < *next {
|
|
return false;
|
|
}
|
|
*next = now + INTERVAL;
|
|
true
|
|
}
|
|
|
|
fn claim(root: &Path) -> std::io::Result<File> {
|
|
if !store::safe_metadata(root)?.is_dir() {
|
|
return Err(store::invalid_data());
|
|
}
|
|
let path = root.join("upload.lock");
|
|
if path.exists() {
|
|
store::safe_metadata(&path)?;
|
|
}
|
|
let lock = OpenOptions::new()
|
|
.read(true)
|
|
.write(true)
|
|
.create(true)
|
|
.truncate(false)
|
|
.open(path)?;
|
|
lock.try_lock().map_err(|_| store::invalid_data())?;
|
|
Ok(lock)
|
|
}
|
|
|
|
fn matches_identity(route: &Route, session: &PlatformSessionSnapshot) -> bool {
|
|
route.user_id.is_some()
|
|
&& route.destination_origin.is_some()
|
|
&& *route
|
|
== Route::from_identity(Some(session.user_id.clone()), Some(&session.api_base_url))
|
|
}
|
|
|
|
fn acknowledged(value: &serde_json::Value, batch: &store::UploadBatch) -> bool {
|
|
value.get("ok") == Some(&serde_json::Value::Bool(true))
|
|
&& value.get("error") == Some(&serde_json::Value::Null)
|
|
&& value.pointer("/data/acknowledged_batch_ids")
|
|
== Some(&serde_json::json!([batch.batch_id]))
|
|
&& value
|
|
.pointer("/data/event_count")
|
|
.and_then(serde_json::Value::as_u64)
|
|
== Some(batch.event_count as u64)
|
|
}
|
|
|
|
fn stop_after(status: u16) -> bool {
|
|
matches!(status, 401 | 403 | 429 | 500..=599)
|
|
}
|
|
|
|
pub(super) async fn cycle(root: &Path, client: &reqwest::Client) {
|
|
let Ok(_lock) = claim(root) else { return };
|
|
let started = Instant::now();
|
|
let Ok(candidates) = store::upload_candidates(root) else {
|
|
failure();
|
|
return;
|
|
};
|
|
let mut attempted = 0;
|
|
for path in candidates {
|
|
if attempted >= MAX_BATCHES || started.elapsed() >= CYCLE_BUDGET {
|
|
break;
|
|
}
|
|
let Ok(batch) = store::load_upload_batch(path) else {
|
|
continue;
|
|
};
|
|
let Some(session) = current_platform_session() else {
|
|
break;
|
|
};
|
|
if !matches_identity(&batch.route, &session) || !batch.exists() {
|
|
continue;
|
|
}
|
|
let Some(origin) = &batch.route.destination_origin else {
|
|
continue;
|
|
};
|
|
attempted += 1;
|
|
// 身份只在发起前读取;在途请求保持原凭据,切换账号后仍可清理原批次。
|
|
let response = client
|
|
.post(format!("{origin}/api/agc/analytics/batches"))
|
|
.bearer_auth(&session.access_token)
|
|
.header("x-genarrative-response-envelope", "1")
|
|
.timeout(REQUEST_TIMEOUT.min(CYCLE_BUDGET.saturating_sub(started.elapsed())))
|
|
.json(&batch.request)
|
|
.send()
|
|
.await;
|
|
let Ok(mut response) = response else {
|
|
failure();
|
|
break;
|
|
};
|
|
let status = response.status().as_u16();
|
|
if status != 200 {
|
|
failure();
|
|
if stop_after(status) {
|
|
break;
|
|
}
|
|
continue;
|
|
}
|
|
let mut bytes = Vec::new();
|
|
let mut valid = true;
|
|
let mut unavailable = false;
|
|
loop {
|
|
match response.chunk().await {
|
|
Ok(Some(chunk)) if bytes.len() + chunk.len() <= 16 * 1024 => {
|
|
bytes.extend_from_slice(&chunk)
|
|
}
|
|
Ok(None) => break,
|
|
Err(_) => {
|
|
valid = false;
|
|
unavailable = true;
|
|
break;
|
|
}
|
|
_ => {
|
|
valid = false;
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
if valid && serde_json::from_slice(&bytes).is_ok_and(|value| acknowledged(&value, &batch)) {
|
|
if batch.acknowledge().is_err() {
|
|
failure();
|
|
}
|
|
} else {
|
|
failure();
|
|
}
|
|
if unavailable {
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn first_cycle_waits_fifteen_minutes_and_resume_does_not_catch_up() {
|
|
let start = tokio::time::Instant::now();
|
|
let mut next = start + INTERVAL;
|
|
assert!(!advance_schedule(start, &mut next));
|
|
assert!(!advance_schedule(
|
|
start + INTERVAL - Duration::from_secs(1),
|
|
&mut next
|
|
));
|
|
assert!(advance_schedule(start + INTERVAL, &mut next));
|
|
let resumed = start + INTERVAL * 8;
|
|
assert!(advance_schedule(resumed, &mut next));
|
|
assert!(!advance_schedule(resumed, &mut next));
|
|
assert_eq!(next, resumed + INTERVAL);
|
|
}
|
|
|
|
#[test]
|
|
fn identity_does_not_claim_anonymous_or_other_account_or_platform() {
|
|
let session = PlatformSessionSnapshot {
|
|
user_id: "A".into(),
|
|
access_token: "test".into(),
|
|
api_base_url: "https://example.com/api".into(),
|
|
identity_generation: 1,
|
|
revision: 1,
|
|
};
|
|
let route = Route::from_identity(Some("A".into()), Some("https://example.com"));
|
|
assert!(matches_identity(&route, &session));
|
|
for route in [
|
|
Route::from_identity(None, Some("https://example.com")),
|
|
Route::from_identity(Some("B".into()), Some("https://example.com")),
|
|
Route::from_identity(Some("A".into()), Some("https://other.example")),
|
|
] {
|
|
assert!(!matches_identity(&route, &session));
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn upload_lock_is_nonblocking_and_single_owner() {
|
|
let dir = tempfile::tempdir().unwrap();
|
|
let lock = claim(dir.path()).unwrap();
|
|
assert!(claim(dir.path()).is_err());
|
|
drop(lock);
|
|
assert!(claim(dir.path()).is_ok());
|
|
}
|
|
|
|
#[test]
|
|
fn toxic_batch_can_be_skipped_but_auth_and_service_failures_end_cycle() {
|
|
for status in [400, 409, 413] {
|
|
assert!(!stop_after(status));
|
|
}
|
|
for status in [401, 403, 429, 500, 503] {
|
|
assert!(stop_after(status));
|
|
}
|
|
assert_eq!(INTERVAL, Duration::from_secs(900));
|
|
}
|
|
}
|