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>
249 lines
7.9 KiB
Rust
249 lines
7.9 KiB
Rust
//! 实例所有权与本地会话状态;仅后台 writer 访问磁盘。
|
|
use super::contract::Route;
|
|
use super::store::{invalid_data, read_bounded, safe_metadata};
|
|
use serde::{Deserialize, Serialize};
|
|
use std::fs::{self, File, OpenOptions};
|
|
use std::io::{self, Write};
|
|
use std::path::Path;
|
|
use std::time::{Duration, SystemTime};
|
|
|
|
#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)]
|
|
#[serde(rename_all = "snake_case")]
|
|
pub(crate) enum LifecycleState {
|
|
Active,
|
|
Closed,
|
|
Incomplete,
|
|
}
|
|
|
|
#[derive(Clone, Debug, Serialize, Deserialize)]
|
|
#[serde(deny_unknown_fields)]
|
|
pub(crate) struct SessionRecord {
|
|
pub schema_version: u32,
|
|
pub editor_session_id: String,
|
|
pub route: Route,
|
|
pub lifecycle_state: LifecycleState,
|
|
pub focus_interval_id: Option<String>,
|
|
pub updated_at: String,
|
|
pub incomplete_detected_at: Option<String>,
|
|
}
|
|
|
|
impl SessionRecord {
|
|
pub(crate) fn validate(&self) -> bool {
|
|
self.updated_at.len() <= 64
|
|
&& self
|
|
.incomplete_detected_at
|
|
.as_ref()
|
|
.is_none_or(|s| s.len() <= 64)
|
|
&& self.route.user_id.as_ref().is_none_or(|s| s.len() <= 4096)
|
|
&& self
|
|
.route
|
|
.destination_origin
|
|
.as_ref()
|
|
.is_none_or(|s| s.len() <= 4096)
|
|
&& self.schema_version == 1
|
|
&& uuid::Uuid::parse_str(&self.editor_session_id).is_ok()
|
|
&& self.route.validate()
|
|
&& self
|
|
.focus_interval_id
|
|
.as_ref()
|
|
.is_none_or(|id| uuid::Uuid::parse_str(id).is_ok())
|
|
&& chrono::DateTime::parse_from_rfc3339(&self.updated_at).is_ok()
|
|
&& self
|
|
.incomplete_detected_at
|
|
.as_ref()
|
|
.is_none_or(|s| chrono::DateTime::parse_from_rfc3339(s).is_ok())
|
|
&& (self.lifecycle_state == LifecycleState::Incomplete)
|
|
== self.incomplete_detected_at.is_some()
|
|
}
|
|
}
|
|
|
|
pub(super) fn claim(instance: &Path) -> io::Result<File> {
|
|
let path = instance.join("owner.lock");
|
|
if path.exists() {
|
|
safe_metadata(&path)?;
|
|
}
|
|
let file = OpenOptions::new()
|
|
.read(true)
|
|
.write(true)
|
|
.create(true)
|
|
.truncate(false)
|
|
.open(path)?;
|
|
file.try_lock().map_err(|_| invalid_data())?;
|
|
Ok(file)
|
|
}
|
|
|
|
fn existing_owner(instance: &Path) -> io::Result<File> {
|
|
let path = instance.join("owner.lock");
|
|
if !safe_metadata(&path)?.is_file() {
|
|
return Err(invalid_data());
|
|
}
|
|
let file = OpenOptions::new().read(true).write(true).open(path)?;
|
|
file.try_lock().map_err(|_| invalid_data())?;
|
|
Ok(file)
|
|
}
|
|
|
|
pub(super) fn write(instance: &Path, record: &SessionRecord, session: &str) -> io::Result<()> {
|
|
if !record.validate() || record.editor_session_id != session {
|
|
return Err(invalid_data());
|
|
}
|
|
let bytes = serde_json::to_vec(record)?;
|
|
if bytes.len() > 16 * 1024 {
|
|
return Err(invalid_data());
|
|
}
|
|
let target = instance.join("session.json");
|
|
if target.exists() {
|
|
safe_metadata(&target)?;
|
|
}
|
|
let temporary = instance.join(format!(".session-{}.tmp", uuid::Uuid::new_v4()));
|
|
let result = (|| {
|
|
let mut file = OpenOptions::new()
|
|
.write(true)
|
|
.create_new(true)
|
|
.open(&temporary)?;
|
|
file.write_all(&bytes)?;
|
|
file.sync_all()?;
|
|
drop(file);
|
|
fs::rename(&temporary, target)
|
|
})();
|
|
if result.is_err() {
|
|
let _ = fs::remove_file(temporary);
|
|
}
|
|
result
|
|
}
|
|
|
|
fn read(instance: &Path) -> io::Result<SessionRecord> {
|
|
let record: SessionRecord =
|
|
serde_json::from_slice(&read_bounded(&instance.join("session.json"), 16 * 1024)?)?;
|
|
if !record.validate()
|
|
|| instance.file_name().and_then(|s| s.to_str()) != Some(record.editor_session_id.as_str())
|
|
{
|
|
return Err(invalid_data());
|
|
}
|
|
Ok(record)
|
|
}
|
|
|
|
pub(super) fn recover(instances: &Path, current: &str) {
|
|
let Ok(entries) = fs::read_dir(instances) else {
|
|
return;
|
|
};
|
|
for entry in entries.flatten() {
|
|
if entry.file_name() == current || !safe_metadata(&entry.path()).is_ok_and(|m| m.is_dir()) {
|
|
continue;
|
|
}
|
|
let Ok(_owner) = existing_owner(&entry.path()) else {
|
|
continue;
|
|
};
|
|
let Ok(mut record) = read(&entry.path()) else {
|
|
continue;
|
|
};
|
|
if record.lifecycle_state != LifecycleState::Active {
|
|
continue;
|
|
}
|
|
record.lifecycle_state = LifecycleState::Incomplete;
|
|
record.incomplete_detected_at = Some(super::contract::timestamp_now());
|
|
let _ = write(&entry.path(), &record, &record.editor_session_id);
|
|
}
|
|
}
|
|
|
|
pub(super) fn prune(
|
|
instances: &Path,
|
|
current: &str,
|
|
reserve: u64,
|
|
limit: u64,
|
|
retention: Duration,
|
|
size: impl Fn() -> u64,
|
|
) {
|
|
let Ok(entries) = fs::read_dir(instances) else {
|
|
return;
|
|
};
|
|
let mut candidates = Vec::new();
|
|
let mut inactive = Vec::new();
|
|
for entry in entries.flatten() {
|
|
if entry.file_name() == current
|
|
|| uuid::Uuid::parse_str(&entry.file_name().to_string_lossy()).is_err()
|
|
|| !safe_metadata(&entry.path()).is_ok_and(|m| m.is_dir())
|
|
{
|
|
continue;
|
|
}
|
|
// 文件锁是存活证明;损坏或未写完的 JSON 不能永久阻止队列清理。
|
|
let Ok(_owner) = existing_owner(&entry.path()) else {
|
|
continue;
|
|
};
|
|
let Ok(files) = fs::read_dir(entry.path()) else {
|
|
continue;
|
|
};
|
|
for file in files.flatten() {
|
|
let name = file.file_name();
|
|
let name = name.to_string_lossy();
|
|
let temporary = name
|
|
.strip_prefix(".session-")
|
|
.and_then(|name| name.strip_suffix(".tmp"))
|
|
.is_some_and(|id| uuid::Uuid::parse_str(id).is_ok());
|
|
if name != "session.json" && !temporary {
|
|
continue;
|
|
}
|
|
let Ok(meta) = safe_metadata(&file.path()) else {
|
|
continue;
|
|
};
|
|
if !meta.is_file() {
|
|
continue;
|
|
}
|
|
let Ok(modified) = meta.modified() else {
|
|
continue;
|
|
};
|
|
candidates.push((modified, entry.path(), file.path()));
|
|
}
|
|
inactive.push(entry.path());
|
|
}
|
|
candidates.sort_by_key(|(modified, _, _)| *modified);
|
|
for (modified, instance, path) in candidates {
|
|
let expired = SystemTime::now()
|
|
.duration_since(modified)
|
|
.is_ok_and(|age| age >= retention);
|
|
if !expired && size().saturating_add(reserve) <= limit {
|
|
continue;
|
|
}
|
|
let Ok(_owner) = existing_owner(&instance) else {
|
|
continue;
|
|
};
|
|
if safe_metadata(&path).is_ok_and(|meta| meta.is_file()) {
|
|
let _ = fs::remove_file(path);
|
|
}
|
|
}
|
|
for instance in inactive {
|
|
let Ok(_owner) = existing_owner(&instance) else {
|
|
continue;
|
|
};
|
|
remove_empty_instance(&instance);
|
|
}
|
|
}
|
|
|
|
fn remove_empty_instance(instance: &Path) {
|
|
// 只删除空目录及普通 owner.lock,不递归删除未知文件或跟随链接。
|
|
let batches = instance.join("batches");
|
|
if safe_metadata(&batches).is_ok_and(|meta| meta.is_dir()) {
|
|
let _ = fs::remove_dir(&batches);
|
|
}
|
|
let Ok(entries) = fs::read_dir(instance) else {
|
|
return;
|
|
};
|
|
let Ok(entries) = entries.collect::<Result<Vec<_>, _>>() else {
|
|
return;
|
|
};
|
|
if entries.len() != 1 || entries[0].file_name() != "owner.lock" {
|
|
return;
|
|
}
|
|
let owner = instance.join("owner.lock");
|
|
if !safe_metadata(&owner).is_ok_and(|meta| meta.is_file()) {
|
|
return;
|
|
}
|
|
// 调用者仍持有锁;实例 UUID 不复用,其他清理者不能同时取得所有权。
|
|
if fs::remove_file(owner).is_ok() {
|
|
let _ = fs::remove_dir(instance);
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
#[path = "session_tests.rs"]
|
|
mod tests;
|