ZY-TAKEOVER-SNAPSHOT 20260815: 铸渊接管前现场冻结(Codex 未提交改动全部入册)· 冰朔面谕接管+可回退
This commit is contained in:
parent
ebe2ecf262
commit
ffed065841
22 changed files with 4293 additions and 567 deletions
|
|
@ -0,0 +1,853 @@
|
|||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
use ring::digest::{digest, SHA256};
|
||||
use rusqlite::{params, Connection, OptionalExtension, Transaction, TransactionBehavior};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::fs;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||
use tauri::{AppHandle, Manager};
|
||||
use uuid::Uuid;
|
||||
|
||||
const KERNEL_SCHEMA: &str = "hololake.personal-channel-kernel/v1";
|
||||
const DATABASE_SCHEMA_VERSION: i64 = 1;
|
||||
const ZERO_HASH: &str = "0000000000000000000000000000000000000000000000000000000000000000";
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct InitializePersonalChannelInput {
|
||||
pub display_name: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct CreatePersonalChannelTaskInput {
|
||||
pub title: String,
|
||||
pub purpose: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct TransitionPersonalChannelTaskInput {
|
||||
pub task_id: String,
|
||||
pub expected_status: String,
|
||||
pub next_status: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct PersonalChannelIdentity {
|
||||
pub human_subject_id: String,
|
||||
pub display_name: String,
|
||||
pub channel_id: String,
|
||||
pub created_at_unix_ms: i64,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct PersonalChannelTask {
|
||||
pub task_id: String,
|
||||
pub title: String,
|
||||
pub purpose: String,
|
||||
pub status: String,
|
||||
pub created_at_unix_ms: i64,
|
||||
pub updated_at_unix_ms: i64,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct PersonalChannelEventProjection {
|
||||
pub sequence: i64,
|
||||
pub event_id: String,
|
||||
pub kind: String,
|
||||
pub task_id: Option<String>,
|
||||
pub summary: String,
|
||||
pub occurred_at_unix_ms: i64,
|
||||
pub event_hash: String,
|
||||
pub receipt_id: String,
|
||||
pub receipt_hash: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct PersonalChannelIntegrity {
|
||||
pub state: &'static str,
|
||||
pub schema_version: i64,
|
||||
pub event_count: i64,
|
||||
pub receipt_count: i64,
|
||||
pub last_event_hash: String,
|
||||
pub last_receipt_hash: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct PersonalChannelSnapshot {
|
||||
pub schema: &'static str,
|
||||
pub state: &'static str,
|
||||
pub identity: Option<PersonalChannelIdentity>,
|
||||
pub current_task: Option<PersonalChannelTask>,
|
||||
pub recent_events: Vec<PersonalChannelEventProjection>,
|
||||
pub integrity: PersonalChannelIntegrity,
|
||||
pub storage: &'static str,
|
||||
pub authority: &'static str,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct EventHashPayload<'a> {
|
||||
schema: &'static str,
|
||||
sequence: i64,
|
||||
event_id: &'a str,
|
||||
human_subject_id: &'a str,
|
||||
kind: &'a str,
|
||||
task_id: Option<&'a str>,
|
||||
summary: &'a str,
|
||||
occurred_at_unix_ms: i64,
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn get_personal_channel_snapshot(
|
||||
app: AppHandle,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
let database = personal_channel_database(&app)?;
|
||||
tauri::async_runtime::spawn_blocking(move || snapshot_at(&database))
|
||||
.await
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_JOIN_FAILED: {error}"))?
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn initialize_personal_channel(
|
||||
app: AppHandle,
|
||||
input: InitializePersonalChannelInput,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
let database = personal_channel_database(&app)?;
|
||||
tauri::async_runtime::spawn_blocking(move || initialize_at(&database, input))
|
||||
.await
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_JOIN_FAILED: {error}"))?
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn create_personal_channel_task(
|
||||
app: AppHandle,
|
||||
input: CreatePersonalChannelTaskInput,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
let database = personal_channel_database(&app)?;
|
||||
tauri::async_runtime::spawn_blocking(move || create_task_at(&database, input))
|
||||
.await
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_JOIN_FAILED: {error}"))?
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub async fn transition_personal_channel_task(
|
||||
app: AppHandle,
|
||||
input: TransitionPersonalChannelTaskInput,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
let database = personal_channel_database(&app)?;
|
||||
tauri::async_runtime::spawn_blocking(move || transition_task_at(&database, input))
|
||||
.await
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_JOIN_FAILED: {error}"))?
|
||||
}
|
||||
|
||||
fn personal_channel_database(app: &AppHandle) -> Result<PathBuf, String> {
|
||||
let app_data = app
|
||||
.path()
|
||||
.app_data_dir()
|
||||
.map_err(|error| format!("HOLOLAKE_APP_DATA_UNAVAILABLE: {error}"))?;
|
||||
let root = app_data.join("personal-channel-v1");
|
||||
create_private_directory(&root)?;
|
||||
Ok(root.join("personal-channel.sqlite3"))
|
||||
}
|
||||
|
||||
fn create_private_directory(path: &Path) -> Result<(), String> {
|
||||
fs::create_dir_all(path)
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_STORAGE_UNAVAILABLE: {error}"))?;
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
fs::set_permissions(path, fs::Permissions::from_mode(0o700)).map_err(|error| {
|
||||
format!("HOLOLAKE_PERSONAL_CHANNEL_STORAGE_PERMISSION_FAILED: {error}")
|
||||
})?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn open_database(path: &Path) -> Result<Connection, String> {
|
||||
if let Some(parent) = path.parent() {
|
||||
create_private_directory(parent)?;
|
||||
}
|
||||
let connection = Connection::open(path)
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_DATABASE_UNAVAILABLE: {error}"))?;
|
||||
connection
|
||||
.busy_timeout(Duration::from_secs(5))
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_DATABASE_UNAVAILABLE: {error}"))?;
|
||||
connection
|
||||
.execute_batch(
|
||||
"PRAGMA foreign_keys = ON;
|
||||
PRAGMA journal_mode = DELETE;
|
||||
PRAGMA synchronous = FULL;
|
||||
PRAGMA trusted_schema = OFF;
|
||||
CREATE TABLE IF NOT EXISTS kernel_meta (
|
||||
key TEXT PRIMARY KEY NOT NULL,
|
||||
value TEXT NOT NULL
|
||||
);
|
||||
INSERT OR IGNORE INTO kernel_meta(key, value) VALUES ('schema_version', '1');
|
||||
CREATE TABLE IF NOT EXISTS identities (
|
||||
singleton INTEGER PRIMARY KEY CHECK(singleton = 1),
|
||||
human_subject_id TEXT NOT NULL UNIQUE,
|
||||
display_name TEXT NOT NULL,
|
||||
channel_id TEXT NOT NULL UNIQUE,
|
||||
created_at_unix_ms INTEGER NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS tasks (
|
||||
task_id TEXT PRIMARY KEY NOT NULL,
|
||||
human_subject_id TEXT NOT NULL,
|
||||
title TEXT NOT NULL,
|
||||
purpose TEXT NOT NULL,
|
||||
status TEXT NOT NULL CHECK(status IN ('ACTIVE', 'COMPLETED')),
|
||||
created_at_unix_ms INTEGER NOT NULL,
|
||||
updated_at_unix_ms INTEGER NOT NULL
|
||||
);
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS one_active_personal_task
|
||||
ON tasks(status) WHERE status = 'ACTIVE';
|
||||
CREATE TABLE IF NOT EXISTS events (
|
||||
sequence INTEGER PRIMARY KEY NOT NULL,
|
||||
event_id TEXT NOT NULL UNIQUE,
|
||||
human_subject_id TEXT NOT NULL,
|
||||
kind TEXT NOT NULL,
|
||||
task_id TEXT,
|
||||
summary TEXT NOT NULL,
|
||||
occurred_at_unix_ms INTEGER NOT NULL,
|
||||
previous_event_hash TEXT NOT NULL,
|
||||
payload_sha256 TEXT NOT NULL,
|
||||
event_hash TEXT NOT NULL UNIQUE
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS receipts (
|
||||
sequence INTEGER PRIMARY KEY NOT NULL,
|
||||
receipt_id TEXT NOT NULL UNIQUE,
|
||||
event_sequence INTEGER NOT NULL UNIQUE REFERENCES events(sequence),
|
||||
event_hash TEXT NOT NULL,
|
||||
payload_sha256 TEXT NOT NULL,
|
||||
previous_receipt_hash TEXT NOT NULL,
|
||||
receipt_hash TEXT NOT NULL UNIQUE,
|
||||
issued_at_unix_ms INTEGER NOT NULL
|
||||
);",
|
||||
)
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_SCHEMA_INVALID: {error}"))?;
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
fs::set_permissions(path, fs::Permissions::from_mode(0o600)).map_err(|error| {
|
||||
format!("HOLOLAKE_PERSONAL_CHANNEL_STORAGE_PERMISSION_FAILED: {error}")
|
||||
})?;
|
||||
}
|
||||
let version: String = connection
|
||||
.query_row(
|
||||
"SELECT value FROM kernel_meta WHERE key = 'schema_version'",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_SCHEMA_INVALID: {error}"))?;
|
||||
if version != DATABASE_SCHEMA_VERSION.to_string() {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_SCHEMA_UNSUPPORTED".into());
|
||||
}
|
||||
Ok(connection)
|
||||
}
|
||||
|
||||
fn initialize_at(
|
||||
database: &Path,
|
||||
input: InitializePersonalChannelInput,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
let display_name = validated_text(&input.display_name, 80, "DISPLAY_NAME")?;
|
||||
let mut connection = open_database(database)?;
|
||||
verify_integrity(&connection)?;
|
||||
let transaction = connection
|
||||
.transaction_with_behavior(TransactionBehavior::Immediate)
|
||||
.map_err(database_write_error)?;
|
||||
let already_initialized: bool = transaction
|
||||
.query_row(
|
||||
"SELECT EXISTS(SELECT 1 FROM identities WHERE singleton = 1)",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)
|
||||
.map_err(database_read_error)?;
|
||||
if already_initialized {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_ALREADY_INITIALIZED".into());
|
||||
}
|
||||
let human_subject_id = format!("human-local-{}", Uuid::new_v4());
|
||||
let channel_id = format!("channel-local-{}", Uuid::new_v4());
|
||||
let created_at = now_unix_ms()?;
|
||||
transaction
|
||||
.execute(
|
||||
"INSERT INTO identities(singleton, human_subject_id, display_name, channel_id, created_at_unix_ms)
|
||||
VALUES(1, ?1, ?2, ?3, ?4)",
|
||||
params![human_subject_id, display_name, channel_id, created_at],
|
||||
)
|
||||
.map_err(database_write_error)?;
|
||||
append_event(
|
||||
&transaction,
|
||||
&human_subject_id,
|
||||
"CHANNEL_INITIALIZED",
|
||||
None,
|
||||
&format!("{display_name} 建立了个人频道"),
|
||||
created_at,
|
||||
)?;
|
||||
transaction.commit().map_err(database_write_error)?;
|
||||
snapshot_at(database)
|
||||
}
|
||||
|
||||
fn create_task_at(
|
||||
database: &Path,
|
||||
input: CreatePersonalChannelTaskInput,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
let title = validated_text(&input.title, 160, "TASK_TITLE")?;
|
||||
let purpose = validated_text(&input.purpose, 1_000, "TASK_PURPOSE")?;
|
||||
let mut connection = open_database(database)?;
|
||||
verify_integrity(&connection)?;
|
||||
let transaction = connection
|
||||
.transaction_with_behavior(TransactionBehavior::Immediate)
|
||||
.map_err(database_write_error)?;
|
||||
let human_subject_id = require_identity_id(&transaction)?;
|
||||
let active_exists: bool = transaction
|
||||
.query_row(
|
||||
"SELECT EXISTS(SELECT 1 FROM tasks WHERE status = 'ACTIVE')",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)
|
||||
.map_err(database_read_error)?;
|
||||
if active_exists {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_ACTIVE_TASK_EXISTS".into());
|
||||
}
|
||||
let task_id = format!("task-local-{}", Uuid::new_v4());
|
||||
let observed_at = now_unix_ms()?;
|
||||
transaction
|
||||
.execute(
|
||||
"INSERT INTO tasks(task_id, human_subject_id, title, purpose, status, created_at_unix_ms, updated_at_unix_ms)
|
||||
VALUES(?1, ?2, ?3, ?4, 'ACTIVE', ?5, ?5)",
|
||||
params![task_id, human_subject_id, title, purpose, observed_at],
|
||||
)
|
||||
.map_err(database_write_error)?;
|
||||
append_event(
|
||||
&transaction,
|
||||
&human_subject_id,
|
||||
"TASK_STARTED",
|
||||
Some(&task_id),
|
||||
&format!("开始:{title}"),
|
||||
observed_at,
|
||||
)?;
|
||||
transaction.commit().map_err(database_write_error)?;
|
||||
snapshot_at(database)
|
||||
}
|
||||
|
||||
fn transition_task_at(
|
||||
database: &Path,
|
||||
input: TransitionPersonalChannelTaskInput,
|
||||
) -> Result<PersonalChannelSnapshot, String> {
|
||||
validate_identifier(&input.task_id, "TASK")?;
|
||||
if input.expected_status != "ACTIVE" || input.next_status != "COMPLETED" {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_TASK_TRANSITION_INVALID".into());
|
||||
}
|
||||
let mut connection = open_database(database)?;
|
||||
verify_integrity(&connection)?;
|
||||
let transaction = connection
|
||||
.transaction_with_behavior(TransactionBehavior::Immediate)
|
||||
.map_err(database_write_error)?;
|
||||
let human_subject_id = require_identity_id(&transaction)?;
|
||||
let task = transaction
|
||||
.query_row(
|
||||
"SELECT title, status FROM tasks WHERE task_id = ?1",
|
||||
params![input.task_id],
|
||||
|row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
|
||||
)
|
||||
.optional()
|
||||
.map_err(database_read_error)?
|
||||
.ok_or("HOLOLAKE_PERSONAL_CHANNEL_TASK_NOT_FOUND")?;
|
||||
if task.1 != input.expected_status {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_TASK_STATE_CONFLICT".into());
|
||||
}
|
||||
let observed_at = now_unix_ms()?;
|
||||
let changed = transaction
|
||||
.execute(
|
||||
"UPDATE tasks SET status = 'COMPLETED', updated_at_unix_ms = ?1
|
||||
WHERE task_id = ?2 AND status = 'ACTIVE'",
|
||||
params![observed_at, input.task_id],
|
||||
)
|
||||
.map_err(database_write_error)?;
|
||||
if changed != 1 {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_TASK_STATE_CONFLICT".into());
|
||||
}
|
||||
append_event(
|
||||
&transaction,
|
||||
&human_subject_id,
|
||||
"TASK_COMPLETED",
|
||||
Some(&input.task_id),
|
||||
&format!("完成:{}", task.0),
|
||||
observed_at,
|
||||
)?;
|
||||
transaction.commit().map_err(database_write_error)?;
|
||||
snapshot_at(database)
|
||||
}
|
||||
|
||||
fn require_identity_id(transaction: &Transaction<'_>) -> Result<String, String> {
|
||||
transaction
|
||||
.query_row(
|
||||
"SELECT human_subject_id FROM identities WHERE singleton = 1",
|
||||
[],
|
||||
|row| row.get(0),
|
||||
)
|
||||
.optional()
|
||||
.map_err(database_read_error)?
|
||||
.ok_or_else(|| "HOLOLAKE_PERSONAL_CHANNEL_NOT_INITIALIZED".into())
|
||||
}
|
||||
|
||||
fn append_event(
|
||||
transaction: &Transaction<'_>,
|
||||
human_subject_id: &str,
|
||||
kind: &str,
|
||||
task_id: Option<&str>,
|
||||
summary: &str,
|
||||
occurred_at_unix_ms: i64,
|
||||
) -> Result<(), String> {
|
||||
let (last_sequence, previous_event_hash, previous_receipt_hash) = transaction
|
||||
.query_row(
|
||||
"SELECT e.sequence, e.event_hash, r.receipt_hash
|
||||
FROM events e JOIN receipts r ON r.event_sequence = e.sequence
|
||||
ORDER BY e.sequence DESC LIMIT 1",
|
||||
[],
|
||||
|row| {
|
||||
Ok((
|
||||
row.get::<_, i64>(0)?,
|
||||
row.get::<_, String>(1)?,
|
||||
row.get::<_, String>(2)?,
|
||||
))
|
||||
},
|
||||
)
|
||||
.optional()
|
||||
.map_err(database_read_error)?
|
||||
.unwrap_or((0, ZERO_HASH.into(), ZERO_HASH.into()));
|
||||
let sequence = last_sequence + 1;
|
||||
let event_id = format!("event-local-{}", Uuid::new_v4());
|
||||
let payload = EventHashPayload {
|
||||
schema: KERNEL_SCHEMA,
|
||||
sequence,
|
||||
event_id: &event_id,
|
||||
human_subject_id,
|
||||
kind,
|
||||
task_id,
|
||||
summary,
|
||||
occurred_at_unix_ms,
|
||||
};
|
||||
let payload_bytes = serde_json::to_vec(&payload)
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_EVENT_INVALID: {error}"))?;
|
||||
let payload_sha256 = sha256_hex(&payload_bytes);
|
||||
let event_hash =
|
||||
sha256_hex(format!("event-chain/v1\n{previous_event_hash}\n{payload_sha256}").as_bytes());
|
||||
let receipt_seed = sha256_hex(format!("receipt-id/v1\n{event_hash}").as_bytes());
|
||||
let receipt_id = format!("HLR-{}", &receipt_seed[..24]);
|
||||
let receipt_hash = sha256_hex(
|
||||
format!(
|
||||
"receipt-chain/v1\n{previous_receipt_hash}\n{receipt_id}\n{event_hash}\n{payload_sha256}"
|
||||
)
|
||||
.as_bytes(),
|
||||
);
|
||||
transaction
|
||||
.execute(
|
||||
"INSERT INTO events(sequence, event_id, human_subject_id, kind, task_id, summary,
|
||||
occurred_at_unix_ms, previous_event_hash, payload_sha256, event_hash)
|
||||
VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
|
||||
params![
|
||||
sequence,
|
||||
event_id,
|
||||
human_subject_id,
|
||||
kind,
|
||||
task_id,
|
||||
summary,
|
||||
occurred_at_unix_ms,
|
||||
previous_event_hash,
|
||||
payload_sha256,
|
||||
event_hash
|
||||
],
|
||||
)
|
||||
.map_err(database_write_error)?;
|
||||
transaction
|
||||
.execute(
|
||||
"INSERT INTO receipts(sequence, receipt_id, event_sequence, event_hash, payload_sha256,
|
||||
previous_receipt_hash, receipt_hash, issued_at_unix_ms)
|
||||
VALUES(?1, ?2, ?1, ?3, ?4, ?5, ?6, ?7)",
|
||||
params![
|
||||
sequence,
|
||||
receipt_id,
|
||||
event_hash,
|
||||
payload_sha256,
|
||||
previous_receipt_hash,
|
||||
receipt_hash,
|
||||
occurred_at_unix_ms
|
||||
],
|
||||
)
|
||||
.map_err(database_write_error)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn snapshot_at(database: &Path) -> Result<PersonalChannelSnapshot, String> {
|
||||
let connection = open_database(database)?;
|
||||
let integrity = verify_integrity(&connection)?;
|
||||
let identity = connection
|
||||
.query_row(
|
||||
"SELECT human_subject_id, display_name, channel_id, created_at_unix_ms
|
||||
FROM identities WHERE singleton = 1",
|
||||
[],
|
||||
|row| {
|
||||
Ok(PersonalChannelIdentity {
|
||||
human_subject_id: row.get(0)?,
|
||||
display_name: row.get(1)?,
|
||||
channel_id: row.get(2)?,
|
||||
created_at_unix_ms: row.get(3)?,
|
||||
})
|
||||
},
|
||||
)
|
||||
.optional()
|
||||
.map_err(database_read_error)?;
|
||||
let current_task = connection
|
||||
.query_row(
|
||||
"SELECT task_id, title, purpose, status, created_at_unix_ms, updated_at_unix_ms
|
||||
FROM tasks WHERE status = 'ACTIVE' LIMIT 1",
|
||||
[],
|
||||
|row| {
|
||||
Ok(PersonalChannelTask {
|
||||
task_id: row.get(0)?,
|
||||
title: row.get(1)?,
|
||||
purpose: row.get(2)?,
|
||||
status: row.get(3)?,
|
||||
created_at_unix_ms: row.get(4)?,
|
||||
updated_at_unix_ms: row.get(5)?,
|
||||
})
|
||||
},
|
||||
)
|
||||
.optional()
|
||||
.map_err(database_read_error)?;
|
||||
let mut statement = connection
|
||||
.prepare(
|
||||
"SELECT e.sequence, e.event_id, e.kind, e.task_id, e.summary, e.occurred_at_unix_ms,
|
||||
e.event_hash, r.receipt_id, r.receipt_hash
|
||||
FROM events e JOIN receipts r ON r.event_sequence = e.sequence
|
||||
ORDER BY e.sequence DESC LIMIT 12",
|
||||
)
|
||||
.map_err(database_read_error)?;
|
||||
let recent_events = statement
|
||||
.query_map([], |row| {
|
||||
Ok(PersonalChannelEventProjection {
|
||||
sequence: row.get(0)?,
|
||||
event_id: row.get(1)?,
|
||||
kind: row.get(2)?,
|
||||
task_id: row.get(3)?,
|
||||
summary: row.get(4)?,
|
||||
occurred_at_unix_ms: row.get(5)?,
|
||||
event_hash: row.get(6)?,
|
||||
receipt_id: row.get(7)?,
|
||||
receipt_hash: row.get(8)?,
|
||||
})
|
||||
})
|
||||
.map_err(database_read_error)?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(database_read_error)?;
|
||||
Ok(PersonalChannelSnapshot {
|
||||
schema: KERNEL_SCHEMA,
|
||||
state: if identity.is_some() {
|
||||
"READY"
|
||||
} else {
|
||||
"UNINITIALIZED"
|
||||
},
|
||||
identity,
|
||||
current_task,
|
||||
recent_events,
|
||||
integrity,
|
||||
storage: "LOCAL_PRIVATE_SQLITE_SINGLE_HOLOLAKE_OWNER",
|
||||
authority: "LOCAL_HUMAN_CONFIRMED_IDENTITY_NOT_SERVER_AUTHORITY",
|
||||
})
|
||||
}
|
||||
|
||||
fn verify_integrity(connection: &Connection) -> Result<PersonalChannelIntegrity, String> {
|
||||
let schema_version: i64 = connection
|
||||
.query_row(
|
||||
"SELECT value FROM kernel_meta WHERE key = 'schema_version'",
|
||||
[],
|
||||
|row| row.get::<_, String>(0),
|
||||
)
|
||||
.map_err(database_read_error)?
|
||||
.parse()
|
||||
.map_err(|_| "HOLOLAKE_PERSONAL_CHANNEL_SCHEMA_INVALID".to_string())?;
|
||||
if schema_version != DATABASE_SCHEMA_VERSION {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_SCHEMA_UNSUPPORTED".into());
|
||||
}
|
||||
let receipt_count: i64 = connection
|
||||
.query_row("SELECT COUNT(*) FROM receipts", [], |row| row.get(0))
|
||||
.map_err(database_read_error)?;
|
||||
let mut statement = connection
|
||||
.prepare(
|
||||
"SELECT e.sequence, e.event_id, e.human_subject_id, e.kind, e.task_id, e.summary,
|
||||
e.occurred_at_unix_ms, e.previous_event_hash, e.payload_sha256, e.event_hash,
|
||||
r.receipt_id, r.event_hash, r.payload_sha256, r.previous_receipt_hash, r.receipt_hash
|
||||
FROM events e LEFT JOIN receipts r ON r.event_sequence = e.sequence
|
||||
ORDER BY e.sequence ASC",
|
||||
)
|
||||
.map_err(database_read_error)?;
|
||||
let mut rows = statement.query([]).map_err(database_read_error)?;
|
||||
let mut expected_sequence = 1_i64;
|
||||
let mut previous_event_hash = ZERO_HASH.to_string();
|
||||
let mut previous_receipt_hash = ZERO_HASH.to_string();
|
||||
let mut event_count = 0_i64;
|
||||
while let Some(row) = rows.next().map_err(database_read_error)? {
|
||||
let sequence: i64 = row.get(0).map_err(database_read_error)?;
|
||||
let event_id: String = row.get(1).map_err(database_read_error)?;
|
||||
let human_subject_id: String = row.get(2).map_err(database_read_error)?;
|
||||
let kind: String = row.get(3).map_err(database_read_error)?;
|
||||
let task_id: Option<String> = row.get(4).map_err(database_read_error)?;
|
||||
let summary: String = row.get(5).map_err(database_read_error)?;
|
||||
let occurred_at: i64 = row.get(6).map_err(database_read_error)?;
|
||||
let stored_previous_event: String = row.get(7).map_err(database_read_error)?;
|
||||
let stored_payload: String = row.get(8).map_err(database_read_error)?;
|
||||
let stored_event_hash: String = row.get(9).map_err(database_read_error)?;
|
||||
let receipt_id: Option<String> = row.get(10).map_err(database_read_error)?;
|
||||
let receipt_event_hash: Option<String> = row.get(11).map_err(database_read_error)?;
|
||||
let receipt_payload: Option<String> = row.get(12).map_err(database_read_error)?;
|
||||
let stored_previous_receipt: Option<String> = row.get(13).map_err(database_read_error)?;
|
||||
let stored_receipt_hash: Option<String> = row.get(14).map_err(database_read_error)?;
|
||||
if sequence != expected_sequence || stored_previous_event != previous_event_hash {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_INTEGRITY_FAILED".into());
|
||||
}
|
||||
let payload = EventHashPayload {
|
||||
schema: KERNEL_SCHEMA,
|
||||
sequence,
|
||||
event_id: &event_id,
|
||||
human_subject_id: &human_subject_id,
|
||||
kind: &kind,
|
||||
task_id: task_id.as_deref(),
|
||||
summary: &summary,
|
||||
occurred_at_unix_ms: occurred_at,
|
||||
};
|
||||
let payload_sha256 = sha256_hex(
|
||||
&serde_json::to_vec(&payload)
|
||||
.map_err(|error| format!("HOLOLAKE_PERSONAL_CHANNEL_EVENT_INVALID: {error}"))?,
|
||||
);
|
||||
let event_hash = sha256_hex(
|
||||
format!("event-chain/v1\n{previous_event_hash}\n{payload_sha256}").as_bytes(),
|
||||
);
|
||||
let receipt_id = receipt_id.ok_or("HOLOLAKE_PERSONAL_CHANNEL_RECEIPT_MISSING")?;
|
||||
let expected_receipt_id = format!(
|
||||
"HLR-{}",
|
||||
&sha256_hex(format!("receipt-id/v1\n{event_hash}").as_bytes())[..24]
|
||||
);
|
||||
let receipt_hash = sha256_hex(
|
||||
format!(
|
||||
"receipt-chain/v1\n{previous_receipt_hash}\n{receipt_id}\n{event_hash}\n{payload_sha256}"
|
||||
)
|
||||
.as_bytes(),
|
||||
);
|
||||
if stored_payload != payload_sha256
|
||||
|| stored_event_hash != event_hash
|
||||
|| receipt_id != expected_receipt_id
|
||||
|| receipt_event_hash.as_deref() != Some(event_hash.as_str())
|
||||
|| receipt_payload.as_deref() != Some(payload_sha256.as_str())
|
||||
|| stored_previous_receipt.as_deref() != Some(previous_receipt_hash.as_str())
|
||||
|| stored_receipt_hash.as_deref() != Some(receipt_hash.as_str())
|
||||
{
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_INTEGRITY_FAILED".into());
|
||||
}
|
||||
previous_event_hash = event_hash;
|
||||
previous_receipt_hash = receipt_hash;
|
||||
expected_sequence += 1;
|
||||
event_count += 1;
|
||||
}
|
||||
if receipt_count != event_count {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_INTEGRITY_FAILED".into());
|
||||
}
|
||||
let identity_count: i64 = connection
|
||||
.query_row("SELECT COUNT(*) FROM identities", [], |row| row.get(0))
|
||||
.map_err(database_read_error)?;
|
||||
if identity_count > 1 || (identity_count == 1 && event_count == 0) {
|
||||
return Err("HOLOLAKE_PERSONAL_CHANNEL_INTEGRITY_FAILED".into());
|
||||
}
|
||||
Ok(PersonalChannelIntegrity {
|
||||
state: "PASS_100",
|
||||
schema_version,
|
||||
event_count,
|
||||
receipt_count,
|
||||
last_event_hash: previous_event_hash,
|
||||
last_receipt_hash: previous_receipt_hash,
|
||||
})
|
||||
}
|
||||
|
||||
fn validated_text(value: &str, maximum_chars: usize, kind: &str) -> Result<String, String> {
|
||||
let trimmed = value.trim();
|
||||
let count = trimmed.chars().count();
|
||||
if count == 0
|
||||
|| count > maximum_chars
|
||||
|| trimmed.chars().any(|character| character.is_control())
|
||||
{
|
||||
return Err(format!("HOLOLAKE_PERSONAL_CHANNEL_{kind}_INVALID"));
|
||||
}
|
||||
Ok(trimmed.to_string())
|
||||
}
|
||||
|
||||
fn validate_identifier(value: &str, kind: &str) -> Result<(), String> {
|
||||
if value.is_empty()
|
||||
|| value.len() > 128
|
||||
|| !value
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b':'))
|
||||
{
|
||||
return Err(format!("HOLOLAKE_PERSONAL_CHANNEL_{kind}_ID_INVALID"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn now_unix_ms() -> Result<i64, String> {
|
||||
let millis = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.map_err(|error| format!("HOLOLAKE_SYSTEM_CLOCK_INVALID: {error}"))?
|
||||
.as_millis();
|
||||
i64::try_from(millis).map_err(|_| "HOLOLAKE_SYSTEM_CLOCK_INVALID".into())
|
||||
}
|
||||
|
||||
fn sha256_hex(value: &[u8]) -> String {
|
||||
digest(&SHA256, value)
|
||||
.as_ref()
|
||||
.iter()
|
||||
.map(|byte| format!("{byte:02x}"))
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn database_read_error(error: rusqlite::Error) -> String {
|
||||
format!("HOLOLAKE_PERSONAL_CHANNEL_DATABASE_UNREADABLE: {error}")
|
||||
}
|
||||
|
||||
fn database_write_error(error: rusqlite::Error) -> String {
|
||||
format!("HOLOLAKE_PERSONAL_CHANNEL_DATABASE_WRITE_FAILED: {error}")
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use tempfile::TempDir;
|
||||
|
||||
fn database(temp: &TempDir) -> PathBuf {
|
||||
temp.path().join("personal-channel.sqlite3")
|
||||
}
|
||||
|
||||
fn initialize(database: &Path) -> PersonalChannelSnapshot {
|
||||
initialize_at(
|
||||
database,
|
||||
InitializePersonalChannelInput {
|
||||
display_name: "冰朔".into(),
|
||||
},
|
||||
)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn initialization_creates_identity_event_and_receipt_atomically() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let snapshot = initialize(&database(&temp));
|
||||
assert_eq!(snapshot.state, "READY");
|
||||
assert_eq!(snapshot.identity.unwrap().display_name, "冰朔");
|
||||
assert_eq!(snapshot.integrity.event_count, 1);
|
||||
assert_eq!(snapshot.integrity.receipt_count, 1);
|
||||
assert_eq!(snapshot.recent_events[0].kind, "CHANNEL_INITIALIZED");
|
||||
assert!(snapshot.recent_events[0].receipt_id.starts_with("HLR-"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn restart_reads_the_same_identity_task_event_and_receipt_chain() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let database = database(&temp);
|
||||
initialize(&database);
|
||||
let created = create_task_at(
|
||||
&database,
|
||||
CreatePersonalChannelTaskInput {
|
||||
title: "完成第一阶段闭环".into(),
|
||||
purpose: "让身份、任务、事件与回执在重启后仍然可见".into(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
let before_hash = created.integrity.last_receipt_hash.clone();
|
||||
drop(created);
|
||||
let after_restart = snapshot_at(&database).unwrap();
|
||||
assert_eq!(
|
||||
after_restart.current_task.unwrap().title,
|
||||
"完成第一阶段闭环"
|
||||
);
|
||||
assert_eq!(after_restart.integrity.event_count, 2);
|
||||
assert_eq!(after_restart.integrity.last_receipt_hash, before_hash);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn one_active_task_is_enforced_and_completion_is_receipted() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let database = database(&temp);
|
||||
initialize(&database);
|
||||
let created = create_task_at(
|
||||
&database,
|
||||
CreatePersonalChannelTaskInput {
|
||||
title: "当前任务".into(),
|
||||
purpose: "验证单一当前焦点".into(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
create_task_at(
|
||||
&database,
|
||||
CreatePersonalChannelTaskInput {
|
||||
title: "冲突任务".into(),
|
||||
purpose: "不应被创建".into(),
|
||||
},
|
||||
)
|
||||
.unwrap_err(),
|
||||
"HOLOLAKE_PERSONAL_CHANNEL_ACTIVE_TASK_EXISTS"
|
||||
);
|
||||
let completed = transition_task_at(
|
||||
&database,
|
||||
TransitionPersonalChannelTaskInput {
|
||||
task_id: created.current_task.unwrap().task_id,
|
||||
expected_status: "ACTIVE".into(),
|
||||
next_status: "COMPLETED".into(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
assert!(completed.current_task.is_none());
|
||||
assert_eq!(completed.recent_events[0].kind, "TASK_COMPLETED");
|
||||
assert_eq!(completed.integrity.event_count, 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn changed_event_bytes_fail_closed_on_readback() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let database = database(&temp);
|
||||
initialize(&database);
|
||||
let connection = open_database(&database).unwrap();
|
||||
connection
|
||||
.execute(
|
||||
"UPDATE events SET summary = '被篡改' WHERE sequence = 1",
|
||||
[],
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
snapshot_at(&database).unwrap_err(),
|
||||
"HOLOLAKE_PERSONAL_CHANNEL_INTEGRITY_FAILED"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn local_identity_is_singleton_and_not_recreated() {
|
||||
let temp = TempDir::new().unwrap();
|
||||
let database = database(&temp);
|
||||
initialize(&database);
|
||||
assert_eq!(
|
||||
initialize_at(
|
||||
&database,
|
||||
InitializePersonalChannelInput {
|
||||
display_name: "另一个人".into(),
|
||||
},
|
||||
)
|
||||
.unwrap_err(),
|
||||
"HOLOLAKE_PERSONAL_CHANNEL_ALREADY_INITIALIZED"
|
||||
);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue