mirror of
https://github.com/sudosylabs/vnidrop.git
synced 2026-08-12 05:29:57 +02:00
refactor(core): open AppDataStores instead of exporting SqlitePool
Introduce persistence::open_all so domain stores (invitation, targeted, blocked) are constructed once. Extract TargetedTransferStore and own it on CoreInner; soft-close raw pool access for new callers. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
3
.gitignore
vendored
3
.gitignore
vendored
@@ -25,6 +25,9 @@ bin/
|
||||
# Local design export scratch
|
||||
output/
|
||||
.scratch/
|
||||
|
||||
# Local ADRs (not tracked — agent/session decisions)
|
||||
docs/adr/
|
||||
.screenshots
|
||||
apple/RELEASE-MACOS.md
|
||||
apple/Generated/*.xcconfig
|
||||
|
||||
39
CONTEXT.md
Normal file
39
CONTEXT.md
Normal file
@@ -0,0 +1,39 @@
|
||||
# VniDrop
|
||||
|
||||
Local peer-to-peer file transfer. This glossary is the product/core ubiquitous language — not an implementation guide.
|
||||
|
||||
## Transfers
|
||||
|
||||
**Invitation transfer**:
|
||||
A share anyone with the ticket can request, subject to approval and access policy. Ordinary multi-recipient send/receive.
|
||||
_Avoid_: contact send, held offer, reusable share offer
|
||||
|
||||
**Targeted transfer**:
|
||||
A transfer bound to one saved-device relationship: immutable sender, receiver, manifest, and content identity; requires explicit approval before content.
|
||||
_Avoid_: contact transfer, private share
|
||||
|
||||
**Saved device**:
|
||||
A remote app identity this installation has mutually consented to remember, with directional grants at a relationship generation.
|
||||
_Avoid_: contact, person, account
|
||||
|
||||
**Device relationship**:
|
||||
The durable pairing state between this installation and a remote endpoint (pending, saved, forgotten/blocked lifecycle).
|
||||
_Avoid_: contact record, friendship
|
||||
|
||||
## Persistence (core)
|
||||
|
||||
**Domain store**:
|
||||
The module that owns schema and queries for one domain (invitation history, targeted transfers, blocked devices, relationship rows, secret metadata). Callers use store methods — never a raw SQL pool.
|
||||
_Avoid_: repository-for-everything, DAO, database layer
|
||||
|
||||
**Invitation repository**:
|
||||
The domain store for invitation-transfer history, artifacts, receiver requests, and related events. Today’s type name may still be `Repository`.
|
||||
_Avoid_: “the database”, AppDataStores
|
||||
|
||||
**AppDataStores**:
|
||||
The bag of concrete domain stores opened together for one app-data profile (one SQLite pool, every schema applied once).
|
||||
_Avoid_: Repository (for the bag), Persistence (as a type name), DbContext
|
||||
|
||||
**Persistence open**:
|
||||
Creating the profile’s SQLite pool, applying all domain schemas, and returning `AppDataStores`. The only place that may touch pool creation for app data.
|
||||
_Avoid_: Repository::open as the global DB entry (once migrated), sqlite_pool export
|
||||
@@ -54,8 +54,15 @@ src/
|
||||
receive.rs # receive, download, export, OutputSinkFile
|
||||
lifecycle.rs # cancel share, delete, status, access mode, shutdown
|
||||
provider.rs # provider messages, per-peer transfer progress
|
||||
saved_devices.rs # experimental saved-device pairing, forget, block
|
||||
targeted.rs # saved-device targeted transfers
|
||||
persistence.rs # AppDataStores / persistence open (domain stores)
|
||||
repository.rs # invitation-transfer domain store (not raw pool export)
|
||||
device_relationship/ # mutual consent + grants
|
||||
targeted_transfer/ # targeted protocol + store adapter
|
||||
blocked_devices.rs
|
||||
secure_secret/ # custody + platform credential adapters
|
||||
filesystem.rs # collect sources, atomic publish, path rules
|
||||
repository.rs # SQLite
|
||||
approval.rs / handshake.rs / ticket.rs / access_policy.rs / event_hub.rs
|
||||
api.rs # UniFFI records/enums
|
||||
tests/ # crate-private unit tests
|
||||
@@ -63,7 +70,8 @@ tests/ # public-API integration tests + support/
|
||||
```
|
||||
|
||||
**Do not** reassemble a single huge `runtime.rs`. Prefer new focused modules if a
|
||||
file approaches ~800 LoC of non-test code.
|
||||
file approaches ~800 LoC of non-test code. Do not add new `SqlitePool` call sites —
|
||||
open domain stores via `persistence::open_all`.
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ mod grant;
|
||||
mod handshake;
|
||||
mod logging;
|
||||
mod pairing_eligibility;
|
||||
mod persistence;
|
||||
mod repository;
|
||||
mod runtime;
|
||||
mod secret;
|
||||
|
||||
61
crates/vnidrop/src/persistence.rs
Normal file
61
crates/vnidrop/src/persistence.rs
Normal file
@@ -0,0 +1,61 @@
|
||||
//! Persistence open: one SQLite pool, every domain schema, [`AppDataStores`].
|
||||
//!
|
||||
//! Runtime talks to domain stores — not a raw pool. Unmigrated modules may still
|
||||
//! take [`AppDataStores::pool_for_unmigrated`] until their own stores deepen.
|
||||
|
||||
use std::{path::Path, str::FromStr};
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use sqlx::{
|
||||
sqlite::{SqliteConnectOptions, SqlitePoolOptions},
|
||||
SqlitePool,
|
||||
};
|
||||
|
||||
use crate::{
|
||||
blocked_devices::BlockStore, repository::Repository, targeted_transfer::TargetedTransferStore,
|
||||
};
|
||||
|
||||
/// Concrete domain stores for one app-data profile.
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct AppDataStores {
|
||||
/// Invitation-transfer history and related invitation tables.
|
||||
pub(crate) invitation: Repository,
|
||||
/// Targeted-transfer durable rows.
|
||||
pub(crate) targeted: TargetedTransferStore,
|
||||
/// Identity-wide deny list.
|
||||
pub(crate) blocked: BlockStore,
|
||||
pool: SqlitePool,
|
||||
}
|
||||
|
||||
impl AppDataStores {
|
||||
/// Temporary: remaining domain modules still construct on a shared pool.
|
||||
/// Do not add new callers — migrate them to domain stores instead.
|
||||
pub(crate) fn pool_for_unmigrated(&self) -> SqlitePool {
|
||||
self.pool.clone()
|
||||
}
|
||||
}
|
||||
|
||||
/// Create the profile pool, apply all domain schemas, return [`AppDataStores`].
|
||||
pub(crate) async fn open_all(app_data_dir: &Path) -> Result<AppDataStores> {
|
||||
let db_path = app_data_dir.join("vnidrop.sqlite3");
|
||||
let options = SqliteConnectOptions::from_str("sqlite://")?
|
||||
.filename(db_path)
|
||||
.create_if_missing(true);
|
||||
let pool = SqlitePoolOptions::new()
|
||||
.max_connections(4)
|
||||
.connect_with(options)
|
||||
.await
|
||||
.context("failed to open app data sqlite")?;
|
||||
|
||||
// Invitation ensure_schema still orchestrates cross-domain schemas until
|
||||
// each domain store owns its ensure_schema call from this path alone.
|
||||
let invitation = Repository::from_pool(pool.clone());
|
||||
invitation.ensure_schema().await?;
|
||||
|
||||
Ok(AppDataStores {
|
||||
targeted: TargetedTransferStore::new(pool.clone()),
|
||||
blocked: BlockStore::new(pool.clone()),
|
||||
invitation,
|
||||
pool,
|
||||
})
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
use std::{path::Path, str::FromStr};
|
||||
#[cfg(test)]
|
||||
use std::path::Path;
|
||||
|
||||
#[cfg(test)]
|
||||
use std::sync::{
|
||||
@@ -7,10 +8,7 @@ use std::sync::{
|
||||
};
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use sqlx::{
|
||||
sqlite::{SqliteConnectOptions, SqlitePoolOptions},
|
||||
Row, SqlitePool,
|
||||
};
|
||||
use sqlx::{Row, SqlitePool};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::{
|
||||
@@ -105,27 +103,23 @@ pub(crate) struct PendingDeliveryReceiptInsert<'a> {
|
||||
}
|
||||
|
||||
impl Repository {
|
||||
pub(crate) async fn open(app_data_dir: &Path) -> Result<Self> {
|
||||
let db_path = app_data_dir.join("vnidrop.sqlite3");
|
||||
let options = SqliteConnectOptions::from_str("sqlite://")?
|
||||
.filename(db_path)
|
||||
.create_if_missing(true);
|
||||
let pool = SqlitePoolOptions::new()
|
||||
.max_connections(4)
|
||||
.connect_with(options)
|
||||
.await?;
|
||||
let repository = Self {
|
||||
pub(crate) fn from_pool(pool: SqlitePool) -> Self {
|
||||
Self {
|
||||
pool,
|
||||
#[cfg(test)]
|
||||
fail_next_write: Arc::new(AtomicBool::new(false)),
|
||||
#[cfg(test)]
|
||||
fail_receive_history_after_dependants: Arc::new(AtomicBool::new(false)),
|
||||
};
|
||||
repository.ensure_schema().await?;
|
||||
Ok(repository)
|
||||
}
|
||||
}
|
||||
|
||||
async fn ensure_schema(&self) -> Result<()> {
|
||||
/// Test/helper entry: opens [`AppDataStores`] and returns the invitation store.
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn open(app_data_dir: &Path) -> Result<Self> {
|
||||
Ok(crate::persistence::open_all(app_data_dir).await?.invitation)
|
||||
}
|
||||
|
||||
pub(crate) async fn ensure_schema(&self) -> Result<()> {
|
||||
// The app owns this SQLite file. Keep migrations explicit so future
|
||||
// desktop/mobile releases can move user history forward in place.
|
||||
sqlx::query(
|
||||
@@ -374,10 +368,6 @@ impl Repository {
|
||||
BlockStore::new(self.pool.clone())
|
||||
}
|
||||
|
||||
pub(crate) fn sqlite_pool(&self) -> SqlitePool {
|
||||
self.pool.clone()
|
||||
}
|
||||
|
||||
#[allow(
|
||||
dead_code,
|
||||
reason = "the private custody seam is activated by platform credential adapters"
|
||||
|
||||
@@ -96,6 +96,8 @@ pub(super) struct CoreInner {
|
||||
pub(super) router: Router,
|
||||
pub(super) store: FsStore,
|
||||
pub(super) repository: Repository,
|
||||
pub(super) targeted_transfers: crate::targeted_transfer::TargetedTransferStore,
|
||||
pub(super) blocked_devices: crate::blocked_devices::BlockStore,
|
||||
_profile_lock: Option<ProfileLock>,
|
||||
pub(super) secret_custody: Option<Arc<crate::secure_secret::SecretCustody>>,
|
||||
pub(super) event_hub: Arc<EventHub>,
|
||||
@@ -150,7 +152,10 @@ impl CoreInner {
|
||||
) -> Result<Arc<Self>> {
|
||||
tokio::fs::create_dir_all(&app_data_dir).await?;
|
||||
init_logging(&app_data_dir)?;
|
||||
let repository = Repository::open(&app_data_dir).await?;
|
||||
let stores = crate::persistence::open_all(&app_data_dir).await?;
|
||||
let repository = stores.invitation.clone();
|
||||
let targeted_transfers = stores.targeted.clone();
|
||||
let blocked_devices = stores.blocked.clone();
|
||||
let (secret_key, secret_custody, profile_lock) = match identity_mode {
|
||||
IdentityMode::Legacy => (load_or_create_secret(&app_data_dir).await?, None, None),
|
||||
IdentityMode::Protected {
|
||||
@@ -378,7 +383,7 @@ impl CoreInner {
|
||||
limits.offer_timeout_ms,
|
||||
);
|
||||
let device_relationships = Arc::new(DeviceRelationshipService::new(
|
||||
repository.sqlite_pool(),
|
||||
stores.pool_for_unmigrated(),
|
||||
secret_custody.clone(),
|
||||
pairing_eligibility.clone(),
|
||||
event_hub.clone(),
|
||||
@@ -401,7 +406,7 @@ impl CoreInner {
|
||||
TargetedTransferProtocol::new(
|
||||
device_relationships.clone(),
|
||||
targeted_offers.clone(),
|
||||
repository.sqlite_pool(),
|
||||
targeted_transfers.clone(),
|
||||
limits.clone(),
|
||||
endpoint.id().to_string(),
|
||||
relay_mode,
|
||||
@@ -416,6 +421,8 @@ impl CoreInner {
|
||||
router,
|
||||
store,
|
||||
repository,
|
||||
targeted_transfers,
|
||||
blocked_devices,
|
||||
_profile_lock: profile_lock,
|
||||
secret_custody: secret_custody.clone(),
|
||||
event_hub,
|
||||
|
||||
@@ -49,8 +49,7 @@ impl CoreInner {
|
||||
accepted: bool,
|
||||
) -> Result<bool, crate::error::VnidropError> {
|
||||
if self
|
||||
.repository
|
||||
.blocked_devices()
|
||||
.blocked_devices
|
||||
.is_blocked(&peer_endpoint_id)
|
||||
.await
|
||||
.unwrap_or(true)
|
||||
@@ -97,8 +96,7 @@ impl CoreInner {
|
||||
peer_endpoint_id: String,
|
||||
) -> Result<(), crate::error::VnidropError> {
|
||||
let now = now_ms();
|
||||
self.repository
|
||||
.blocked_devices()
|
||||
self.blocked_devices
|
||||
.block_endpoint(&peer_endpoint_id, now)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
@@ -120,8 +118,7 @@ impl CoreInner {
|
||||
&self,
|
||||
peer_endpoint_id: String,
|
||||
) -> Result<(), crate::error::VnidropError> {
|
||||
self.repository
|
||||
.blocked_devices()
|
||||
self.blocked_devices
|
||||
.unblock_endpoint(&peer_endpoint_id)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
@@ -132,8 +129,7 @@ impl CoreInner {
|
||||
pub(super) async fn list_blocked_devices(
|
||||
&self,
|
||||
) -> Result<Vec<String>, crate::error::VnidropError> {
|
||||
self.repository
|
||||
.blocked_devices()
|
||||
self.blocked_devices
|
||||
.list_blocked()
|
||||
.await
|
||||
.map_err(VnidropError::repository)
|
||||
|
||||
@@ -21,7 +21,7 @@ use crate::{
|
||||
SubmitTargetedOffer, TargetedOfferResponse, TargetedTransferProtocol,
|
||||
},
|
||||
reconstruct_authorization, TargetedAuthorization, TargetedAuthorizationDraft,
|
||||
TargetedTransferRole, TargetedTransferRow, TargetedTransferStore,
|
||||
TargetedTransferRole, TargetedTransferRow,
|
||||
},
|
||||
ticket::VnidropTicket,
|
||||
util::{non_empty, now_ms},
|
||||
@@ -36,8 +36,8 @@ impl CoreInner {
|
||||
std::time::Duration::from_millis(self.limits.offer_timeout_ms)
|
||||
}
|
||||
|
||||
pub(super) fn targeted_store(&self) -> TargetedTransferStore {
|
||||
TargetedTransferStore::new(self.repository.sqlite_pool())
|
||||
pub(super) fn targeted_store(&self) -> crate::targeted_transfer::TargetedTransferStore {
|
||||
self.targeted_transfers.clone()
|
||||
}
|
||||
|
||||
pub(super) async fn list_pending_targeted_offers(&self) -> Vec<PendingTargetedOffer> {
|
||||
|
||||
@@ -7,6 +7,7 @@ mod auth;
|
||||
pub(crate) mod inbox;
|
||||
pub(crate) mod protocol;
|
||||
mod state;
|
||||
mod store;
|
||||
|
||||
pub(crate) use auth::{
|
||||
auth_secret_material, reconstruct_authorization, TargetedAuthorization,
|
||||
@@ -14,485 +15,6 @@ pub(crate) use auth::{
|
||||
};
|
||||
pub(crate) use inbox::{RespondError, TargetedOfferInbox};
|
||||
pub(crate) use protocol::TargetedTransferProtocol;
|
||||
|
||||
use sqlx::{Row, SqlitePool};
|
||||
|
||||
use crate::{
|
||||
api::{TargetedTransfer, TargetedTransferState},
|
||||
error::VnidropError,
|
||||
util::now_ms,
|
||||
pub(crate) use store::{
|
||||
ensure_schema, state_as_str, TargetedTransferRole, TargetedTransferRow, TargetedTransferStore,
|
||||
};
|
||||
|
||||
pub(crate) async fn ensure_schema(pool: &SqlitePool) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
CREATE TABLE IF NOT EXISTS targeted_transfers (
|
||||
id TEXT PRIMARY KEY,
|
||||
protocol_transfer_id INTEGER NOT NULL UNIQUE,
|
||||
sender_endpoint_id TEXT NOT NULL,
|
||||
receiver_endpoint_id TEXT NOT NULL,
|
||||
manifest_id TEXT NOT NULL,
|
||||
content_hash TEXT NOT NULL,
|
||||
transfer_name TEXT NOT NULL,
|
||||
file_count INTEGER NOT NULL,
|
||||
total_size INTEGER NOT NULL,
|
||||
verified_bytes INTEGER NOT NULL DEFAULT 0,
|
||||
blob_ticket TEXT,
|
||||
authorization_secret_handle TEXT,
|
||||
role TEXT NOT NULL DEFAULT 'sender',
|
||||
state TEXT NOT NULL,
|
||||
created_at INTEGER NOT NULL,
|
||||
updated_at INTEGER NOT NULL
|
||||
);
|
||||
"#,
|
||||
)
|
||||
.execute(pool)
|
||||
.await?;
|
||||
let columns = sqlx::query("PRAGMA table_info(targeted_transfers)")
|
||||
.fetch_all(pool)
|
||||
.await?;
|
||||
let has = |name: &str| columns.iter().any(|row| row.get::<String, _>(1) == name);
|
||||
if !has("verified_bytes") {
|
||||
sqlx::query(
|
||||
"ALTER TABLE targeted_transfers ADD COLUMN verified_bytes INTEGER NOT NULL DEFAULT 0",
|
||||
)
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
if !has("blob_ticket") {
|
||||
sqlx::query("ALTER TABLE targeted_transfers ADD COLUMN blob_ticket TEXT")
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
if !has("authorization_secret_handle") {
|
||||
sqlx::query("ALTER TABLE targeted_transfers ADD COLUMN authorization_secret_handle TEXT")
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
if !has("role") {
|
||||
sqlx::query(
|
||||
"ALTER TABLE targeted_transfers ADD COLUMN role TEXT NOT NULL DEFAULT 'sender'",
|
||||
)
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct TargetedTransferStore {
|
||||
pool: SqlitePool,
|
||||
}
|
||||
|
||||
impl TargetedTransferStore {
|
||||
pub(crate) fn new(pool: SqlitePool) -> Self {
|
||||
Self { pool }
|
||||
}
|
||||
|
||||
pub(crate) async fn insert(&self, transfer: &TargetedTransferRow) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO targeted_transfers (
|
||||
id, protocol_transfer_id, sender_endpoint_id, receiver_endpoint_id,
|
||||
manifest_id, content_hash, transfer_name, file_count, total_size,
|
||||
verified_bytes, blob_ticket, authorization_secret_handle, role,
|
||||
state, created_at, updated_at
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)
|
||||
"#,
|
||||
)
|
||||
.bind(&transfer.id)
|
||||
.bind(transfer.protocol_transfer_id as i64)
|
||||
.bind(&transfer.sender_endpoint_id)
|
||||
.bind(&transfer.receiver_endpoint_id)
|
||||
.bind(&transfer.manifest_id)
|
||||
.bind(&transfer.content_hash)
|
||||
.bind(&transfer.transfer_name)
|
||||
.bind(transfer.file_count as i64)
|
||||
.bind(transfer.total_size as i64)
|
||||
.bind(transfer.verified_bytes as i64)
|
||||
.bind(&transfer.blob_ticket)
|
||||
.bind(&transfer.authorization_secret_handle)
|
||||
.bind(role_as_str(transfer.role))
|
||||
.bind(state_as_str(transfer.state))
|
||||
.bind(transfer.created_at)
|
||||
.bind(transfer.updated_at)
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn set_state(
|
||||
&self,
|
||||
id: &str,
|
||||
from: TargetedTransferState,
|
||||
to: TargetedTransferState,
|
||||
) -> Result<(), VnidropError> {
|
||||
from.validate_transition_to(to)?;
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = ?2, updated_at = ?3
|
||||
WHERE id = ?1 AND state = ?4
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(state_as_str(to))
|
||||
.bind(now_ms())
|
||||
.bind(state_as_str(from))
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
if result.rows_affected() == 0 {
|
||||
return Err(VnidropError::InvalidTransition {
|
||||
reason: format!("{} -> {}", state_as_str(from), state_as_str(to)),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Transition from any non-terminal state; used by cancel/delete.
|
||||
pub(crate) async fn set_state_from_any(
|
||||
&self,
|
||||
id: &str,
|
||||
to: TargetedTransferState,
|
||||
) -> Result<(), VnidropError> {
|
||||
let Some(row) = self.get_row(id).await? else {
|
||||
return Err(VnidropError::invalid_input(anyhow::anyhow!(
|
||||
"unknown targeted transfer"
|
||||
)));
|
||||
};
|
||||
if row.state == to {
|
||||
return Ok(());
|
||||
}
|
||||
row.state.validate_transition_to(to)?;
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = ?2, updated_at = ?3
|
||||
WHERE id = ?1 AND state = ?4
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(state_as_str(to))
|
||||
.bind(now_ms())
|
||||
.bind(state_as_str(row.state))
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
if result.rows_affected() == 0 {
|
||||
return Err(VnidropError::InvalidTransition {
|
||||
reason: format!("{} -> {}", state_as_str(row.state), state_as_str(to)),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn set_verified_bytes(
|
||||
&self,
|
||||
id: &str,
|
||||
verified_bytes: u64,
|
||||
) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET verified_bytes = ?2, updated_at = ?3
|
||||
WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(verified_bytes as i64)
|
||||
.bind(now_ms())
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn store_authorization(
|
||||
&self,
|
||||
id: &str,
|
||||
blob_ticket: &str,
|
||||
authorization_secret_handle: &str,
|
||||
) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET blob_ticket = ?2,
|
||||
authorization_secret_handle = ?3,
|
||||
updated_at = ?4
|
||||
WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(blob_ticket)
|
||||
.bind(authorization_secret_handle)
|
||||
.bind(now_ms())
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn clear_authorization(&self, id: &str) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET blob_ticket = NULL,
|
||||
authorization_secret_handle = NULL,
|
||||
verified_bytes = 0,
|
||||
updated_at = ?2
|
||||
WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(now_ms())
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn get(&self, id: &str) -> Result<Option<TargetedTransfer>, VnidropError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT id, sender_endpoint_id, receiver_endpoint_id, manifest_id,
|
||||
file_count, total_size, verified_bytes, state, created_at, updated_at
|
||||
FROM targeted_transfers WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
row.map(row_to_transfer).transpose()
|
||||
}
|
||||
|
||||
pub(crate) async fn get_row(
|
||||
&self,
|
||||
id: &str,
|
||||
) -> Result<Option<TargetedTransferRow>, VnidropError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT id, protocol_transfer_id, sender_endpoint_id, receiver_endpoint_id,
|
||||
manifest_id, content_hash, transfer_name, file_count, total_size,
|
||||
verified_bytes, blob_ticket, authorization_secret_handle, role,
|
||||
state, created_at, updated_at
|
||||
FROM targeted_transfers WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
row.map(row_to_full).transpose()
|
||||
}
|
||||
|
||||
pub(crate) async fn list(&self) -> Result<Vec<TargetedTransfer>, VnidropError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT id, sender_endpoint_id, receiver_endpoint_id, manifest_id,
|
||||
file_count, total_size, verified_bytes, state, created_at, updated_at
|
||||
FROM targeted_transfers
|
||||
ORDER BY updated_at DESC
|
||||
"#,
|
||||
)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
rows.into_iter().map(row_to_transfer).collect()
|
||||
}
|
||||
|
||||
pub(crate) async fn list_resumable_sender_rows(
|
||||
&self,
|
||||
) -> Result<Vec<TargetedTransferRow>, VnidropError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT id, protocol_transfer_id, sender_endpoint_id, receiver_endpoint_id,
|
||||
manifest_id, content_hash, transfer_name, file_count, total_size,
|
||||
verified_bytes, blob_ticket, authorization_secret_handle, role,
|
||||
state, created_at, updated_at
|
||||
FROM targeted_transfers
|
||||
WHERE role = 'sender'
|
||||
AND state IN ('approved', 'connecting', 'transferring', 'interrupted')
|
||||
"#,
|
||||
)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
rows.into_iter().map(row_to_full).collect()
|
||||
}
|
||||
|
||||
pub(crate) async fn cancel_by_peer(&self, peer_endpoint_id: &str) -> Result<u64, VnidropError> {
|
||||
let now = now_ms();
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = 'cancelled', updated_at = ?2
|
||||
WHERE (sender_endpoint_id = ?1 OR receiver_endpoint_id = ?1)
|
||||
AND state NOT IN ('completed', 'declined', 'cancelled', 'failed', 'deleted')
|
||||
"#,
|
||||
)
|
||||
.bind(peer_endpoint_id)
|
||||
.bind(now)
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(result.rows_affected())
|
||||
}
|
||||
|
||||
pub(crate) async fn protocol_ids_for_peer(
|
||||
&self,
|
||||
peer_endpoint_id: &str,
|
||||
) -> Result<Vec<u64>, VnidropError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT protocol_transfer_id FROM targeted_transfers
|
||||
WHERE sender_endpoint_id = ?1 OR receiver_endpoint_id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(peer_endpoint_id)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(rows
|
||||
.into_iter()
|
||||
.map(|row| row.get::<i64, _>(0) as u64)
|
||||
.collect())
|
||||
}
|
||||
|
||||
pub(crate) async fn mark_interrupted_in_flight(&self) -> Result<u64, VnidropError> {
|
||||
let now = now_ms();
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = 'interrupted', updated_at = ?1
|
||||
WHERE state IN ('connecting', 'transferring')
|
||||
"#,
|
||||
)
|
||||
.bind(now)
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(result.rows_affected())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub(crate) enum TargetedTransferRole {
|
||||
Sender,
|
||||
Receiver,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub(crate) struct TargetedTransferRow {
|
||||
pub(crate) id: String,
|
||||
pub(crate) protocol_transfer_id: u64,
|
||||
pub(crate) sender_endpoint_id: String,
|
||||
pub(crate) receiver_endpoint_id: String,
|
||||
pub(crate) manifest_id: String,
|
||||
pub(crate) content_hash: String,
|
||||
pub(crate) transfer_name: String,
|
||||
pub(crate) file_count: u64,
|
||||
pub(crate) total_size: u64,
|
||||
pub(crate) verified_bytes: u64,
|
||||
pub(crate) blob_ticket: Option<String>,
|
||||
pub(crate) authorization_secret_handle: Option<String>,
|
||||
pub(crate) role: TargetedTransferRole,
|
||||
pub(crate) state: TargetedTransferState,
|
||||
pub(crate) created_at: i64,
|
||||
pub(crate) updated_at: i64,
|
||||
}
|
||||
|
||||
fn row_to_transfer(row: sqlx::sqlite::SqliteRow) -> Result<TargetedTransfer, VnidropError> {
|
||||
Ok(TargetedTransfer {
|
||||
id: row.get("id"),
|
||||
sender_endpoint_id: row.get("sender_endpoint_id"),
|
||||
receiver_endpoint_id: row.get("receiver_endpoint_id"),
|
||||
manifest_id: row.get("manifest_id"),
|
||||
file_count: row.get::<i64, _>("file_count") as u64,
|
||||
total_size: row.get::<i64, _>("total_size") as u64,
|
||||
verified_bytes: row.try_get::<i64, _>("verified_bytes").unwrap_or(0) as u64,
|
||||
state: parse_state(&row.get::<String, _>("state"))?,
|
||||
created_at: row.get("created_at"),
|
||||
updated_at: row.get("updated_at"),
|
||||
})
|
||||
}
|
||||
|
||||
fn row_to_full(row: sqlx::sqlite::SqliteRow) -> Result<TargetedTransferRow, VnidropError> {
|
||||
Ok(TargetedTransferRow {
|
||||
id: row.get("id"),
|
||||
protocol_transfer_id: row.get::<i64, _>("protocol_transfer_id") as u64,
|
||||
sender_endpoint_id: row.get("sender_endpoint_id"),
|
||||
receiver_endpoint_id: row.get("receiver_endpoint_id"),
|
||||
manifest_id: row.get("manifest_id"),
|
||||
content_hash: row.get("content_hash"),
|
||||
transfer_name: row.get("transfer_name"),
|
||||
file_count: row.get::<i64, _>("file_count") as u64,
|
||||
total_size: row.get::<i64, _>("total_size") as u64,
|
||||
verified_bytes: row.try_get::<i64, _>("verified_bytes").unwrap_or(0) as u64,
|
||||
blob_ticket: row.try_get("blob_ticket").ok().flatten(),
|
||||
authorization_secret_handle: row.try_get("authorization_secret_handle").ok().flatten(),
|
||||
role: parse_role(
|
||||
&row.try_get::<String, _>("role")
|
||||
.unwrap_or_else(|_| "sender".to_string()),
|
||||
)?,
|
||||
state: parse_state(&row.get::<String, _>("state"))?,
|
||||
created_at: row.get("created_at"),
|
||||
updated_at: row.get("updated_at"),
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn state_as_str(state: TargetedTransferState) -> &'static str {
|
||||
match state {
|
||||
TargetedTransferState::Preparing => "preparing",
|
||||
TargetedTransferState::Offering => "offering",
|
||||
TargetedTransferState::AwaitingApproval => "awaiting_approval",
|
||||
TargetedTransferState::Approved => "approved",
|
||||
TargetedTransferState::Connecting => "connecting",
|
||||
TargetedTransferState::Transferring => "transferring",
|
||||
TargetedTransferState::Interrupted => "interrupted",
|
||||
TargetedTransferState::Completed => "completed",
|
||||
TargetedTransferState::Declined => "declined",
|
||||
TargetedTransferState::Cancelled => "cancelled",
|
||||
TargetedTransferState::Failed => "failed",
|
||||
TargetedTransferState::Deleted => "deleted",
|
||||
}
|
||||
}
|
||||
|
||||
fn role_as_str(role: TargetedTransferRole) -> &'static str {
|
||||
match role {
|
||||
TargetedTransferRole::Sender => "sender",
|
||||
TargetedTransferRole::Receiver => "receiver",
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_role(value: &str) -> Result<TargetedTransferRole, VnidropError> {
|
||||
match value {
|
||||
"sender" => Ok(TargetedTransferRole::Sender),
|
||||
"receiver" => Ok(TargetedTransferRole::Receiver),
|
||||
other => Err(VnidropError::repository(anyhow::anyhow!(
|
||||
"unknown targeted transfer role: {other}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_state(value: &str) -> Result<TargetedTransferState, VnidropError> {
|
||||
match value {
|
||||
"preparing" => Ok(TargetedTransferState::Preparing),
|
||||
"offering" => Ok(TargetedTransferState::Offering),
|
||||
"awaiting_approval" => Ok(TargetedTransferState::AwaitingApproval),
|
||||
"approved" => Ok(TargetedTransferState::Approved),
|
||||
"connecting" => Ok(TargetedTransferState::Connecting),
|
||||
"transferring" => Ok(TargetedTransferState::Transferring),
|
||||
"interrupted" => Ok(TargetedTransferState::Interrupted),
|
||||
"completed" => Ok(TargetedTransferState::Completed),
|
||||
"declined" => Ok(TargetedTransferState::Declined),
|
||||
"cancelled" => Ok(TargetedTransferState::Cancelled),
|
||||
"failed" => Ok(TargetedTransferState::Failed),
|
||||
"deleted" => Ok(TargetedTransferState::Deleted),
|
||||
other => Err(VnidropError::repository(anyhow::anyhow!(
|
||||
"unknown targeted transfer state: {other}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -55,7 +55,7 @@ impl TargetedTransferProtocol {
|
||||
pub(crate) fn new(
|
||||
relationships: std::sync::Arc<DeviceRelationshipService>,
|
||||
inbox: TargetedOfferInbox,
|
||||
pool: sqlx::SqlitePool,
|
||||
store: TargetedTransferStore,
|
||||
limits: crate::api::CoreLimits,
|
||||
local_endpoint_id: String,
|
||||
relay_mode: CoreRelayMode,
|
||||
@@ -64,7 +64,7 @@ impl TargetedTransferProtocol {
|
||||
Self {
|
||||
relationships,
|
||||
inbox,
|
||||
store: TargetedTransferStore::new(pool),
|
||||
store,
|
||||
limits,
|
||||
local_endpoint_id,
|
||||
relay_mode,
|
||||
|
||||
486
crates/vnidrop/src/targeted_transfer/store.rs
Normal file
486
crates/vnidrop/src/targeted_transfer/store.rs
Normal file
@@ -0,0 +1,486 @@
|
||||
//! Durable targeted-transfer rows (schema + queries).
|
||||
//!
|
||||
//! This is the domain store adapter for targeted transfers. Callers use store
|
||||
//! methods — not a raw SQL pool.
|
||||
|
||||
use sqlx::{Row, SqlitePool};
|
||||
|
||||
use crate::{
|
||||
api::{TargetedTransfer, TargetedTransferState},
|
||||
error::VnidropError,
|
||||
util::now_ms,
|
||||
};
|
||||
|
||||
pub(crate) async fn ensure_schema(pool: &SqlitePool) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
CREATE TABLE IF NOT EXISTS targeted_transfers (
|
||||
id TEXT PRIMARY KEY,
|
||||
protocol_transfer_id INTEGER NOT NULL UNIQUE,
|
||||
sender_endpoint_id TEXT NOT NULL,
|
||||
receiver_endpoint_id TEXT NOT NULL,
|
||||
manifest_id TEXT NOT NULL,
|
||||
content_hash TEXT NOT NULL,
|
||||
transfer_name TEXT NOT NULL,
|
||||
file_count INTEGER NOT NULL,
|
||||
total_size INTEGER NOT NULL,
|
||||
verified_bytes INTEGER NOT NULL DEFAULT 0,
|
||||
blob_ticket TEXT,
|
||||
authorization_secret_handle TEXT,
|
||||
role TEXT NOT NULL DEFAULT 'sender',
|
||||
state TEXT NOT NULL,
|
||||
created_at INTEGER NOT NULL,
|
||||
updated_at INTEGER NOT NULL
|
||||
);
|
||||
"#,
|
||||
)
|
||||
.execute(pool)
|
||||
.await?;
|
||||
let columns = sqlx::query("PRAGMA table_info(targeted_transfers)")
|
||||
.fetch_all(pool)
|
||||
.await?;
|
||||
let has = |name: &str| columns.iter().any(|row| row.get::<String, _>(1) == name);
|
||||
if !has("verified_bytes") {
|
||||
sqlx::query(
|
||||
"ALTER TABLE targeted_transfers ADD COLUMN verified_bytes INTEGER NOT NULL DEFAULT 0",
|
||||
)
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
if !has("blob_ticket") {
|
||||
sqlx::query("ALTER TABLE targeted_transfers ADD COLUMN blob_ticket TEXT")
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
if !has("authorization_secret_handle") {
|
||||
sqlx::query("ALTER TABLE targeted_transfers ADD COLUMN authorization_secret_handle TEXT")
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
if !has("role") {
|
||||
sqlx::query(
|
||||
"ALTER TABLE targeted_transfers ADD COLUMN role TEXT NOT NULL DEFAULT 'sender'",
|
||||
)
|
||||
.execute(pool)
|
||||
.await?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct TargetedTransferStore {
|
||||
pool: SqlitePool,
|
||||
}
|
||||
|
||||
impl TargetedTransferStore {
|
||||
pub(crate) fn new(pool: SqlitePool) -> Self {
|
||||
Self { pool }
|
||||
}
|
||||
|
||||
pub(crate) async fn insert(&self, transfer: &TargetedTransferRow) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO targeted_transfers (
|
||||
id, protocol_transfer_id, sender_endpoint_id, receiver_endpoint_id,
|
||||
manifest_id, content_hash, transfer_name, file_count, total_size,
|
||||
verified_bytes, blob_ticket, authorization_secret_handle, role,
|
||||
state, created_at, updated_at
|
||||
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)
|
||||
"#,
|
||||
)
|
||||
.bind(&transfer.id)
|
||||
.bind(transfer.protocol_transfer_id as i64)
|
||||
.bind(&transfer.sender_endpoint_id)
|
||||
.bind(&transfer.receiver_endpoint_id)
|
||||
.bind(&transfer.manifest_id)
|
||||
.bind(&transfer.content_hash)
|
||||
.bind(&transfer.transfer_name)
|
||||
.bind(transfer.file_count as i64)
|
||||
.bind(transfer.total_size as i64)
|
||||
.bind(transfer.verified_bytes as i64)
|
||||
.bind(&transfer.blob_ticket)
|
||||
.bind(&transfer.authorization_secret_handle)
|
||||
.bind(role_as_str(transfer.role))
|
||||
.bind(state_as_str(transfer.state))
|
||||
.bind(transfer.created_at)
|
||||
.bind(transfer.updated_at)
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn set_state(
|
||||
&self,
|
||||
id: &str,
|
||||
from: TargetedTransferState,
|
||||
to: TargetedTransferState,
|
||||
) -> Result<(), VnidropError> {
|
||||
from.validate_transition_to(to)?;
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = ?2, updated_at = ?3
|
||||
WHERE id = ?1 AND state = ?4
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(state_as_str(to))
|
||||
.bind(now_ms())
|
||||
.bind(state_as_str(from))
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
if result.rows_affected() == 0 {
|
||||
return Err(VnidropError::InvalidTransition {
|
||||
reason: format!("{} -> {}", state_as_str(from), state_as_str(to)),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Transition from any non-terminal state; used by cancel/delete.
|
||||
pub(crate) async fn set_state_from_any(
|
||||
&self,
|
||||
id: &str,
|
||||
to: TargetedTransferState,
|
||||
) -> Result<(), VnidropError> {
|
||||
let Some(row) = self.get_row(id).await? else {
|
||||
return Err(VnidropError::invalid_input(anyhow::anyhow!(
|
||||
"unknown targeted transfer"
|
||||
)));
|
||||
};
|
||||
if row.state == to {
|
||||
return Ok(());
|
||||
}
|
||||
row.state.validate_transition_to(to)?;
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = ?2, updated_at = ?3
|
||||
WHERE id = ?1 AND state = ?4
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(state_as_str(to))
|
||||
.bind(now_ms())
|
||||
.bind(state_as_str(row.state))
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
if result.rows_affected() == 0 {
|
||||
return Err(VnidropError::InvalidTransition {
|
||||
reason: format!("{} -> {}", state_as_str(row.state), state_as_str(to)),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn set_verified_bytes(
|
||||
&self,
|
||||
id: &str,
|
||||
verified_bytes: u64,
|
||||
) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET verified_bytes = ?2, updated_at = ?3
|
||||
WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(verified_bytes as i64)
|
||||
.bind(now_ms())
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn store_authorization(
|
||||
&self,
|
||||
id: &str,
|
||||
blob_ticket: &str,
|
||||
authorization_secret_handle: &str,
|
||||
) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET blob_ticket = ?2,
|
||||
authorization_secret_handle = ?3,
|
||||
updated_at = ?4
|
||||
WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(blob_ticket)
|
||||
.bind(authorization_secret_handle)
|
||||
.bind(now_ms())
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn clear_authorization(&self, id: &str) -> Result<(), VnidropError> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET blob_ticket = NULL,
|
||||
authorization_secret_handle = NULL,
|
||||
verified_bytes = 0,
|
||||
updated_at = ?2
|
||||
WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.bind(now_ms())
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn get(&self, id: &str) -> Result<Option<TargetedTransfer>, VnidropError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT id, sender_endpoint_id, receiver_endpoint_id, manifest_id,
|
||||
file_count, total_size, verified_bytes, state, created_at, updated_at
|
||||
FROM targeted_transfers WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
row.map(row_to_transfer).transpose()
|
||||
}
|
||||
|
||||
pub(crate) async fn get_row(
|
||||
&self,
|
||||
id: &str,
|
||||
) -> Result<Option<TargetedTransferRow>, VnidropError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT id, protocol_transfer_id, sender_endpoint_id, receiver_endpoint_id,
|
||||
manifest_id, content_hash, transfer_name, file_count, total_size,
|
||||
verified_bytes, blob_ticket, authorization_secret_handle, role,
|
||||
state, created_at, updated_at
|
||||
FROM targeted_transfers WHERE id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
row.map(row_to_full).transpose()
|
||||
}
|
||||
|
||||
pub(crate) async fn list(&self) -> Result<Vec<TargetedTransfer>, VnidropError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT id, sender_endpoint_id, receiver_endpoint_id, manifest_id,
|
||||
file_count, total_size, verified_bytes, state, created_at, updated_at
|
||||
FROM targeted_transfers
|
||||
ORDER BY updated_at DESC
|
||||
"#,
|
||||
)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
rows.into_iter().map(row_to_transfer).collect()
|
||||
}
|
||||
|
||||
pub(crate) async fn list_resumable_sender_rows(
|
||||
&self,
|
||||
) -> Result<Vec<TargetedTransferRow>, VnidropError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT id, protocol_transfer_id, sender_endpoint_id, receiver_endpoint_id,
|
||||
manifest_id, content_hash, transfer_name, file_count, total_size,
|
||||
verified_bytes, blob_ticket, authorization_secret_handle, role,
|
||||
state, created_at, updated_at
|
||||
FROM targeted_transfers
|
||||
WHERE role = 'sender'
|
||||
AND state IN ('approved', 'connecting', 'transferring', 'interrupted')
|
||||
"#,
|
||||
)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
rows.into_iter().map(row_to_full).collect()
|
||||
}
|
||||
|
||||
pub(crate) async fn cancel_by_peer(&self, peer_endpoint_id: &str) -> Result<u64, VnidropError> {
|
||||
let now = now_ms();
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = 'cancelled', updated_at = ?2
|
||||
WHERE (sender_endpoint_id = ?1 OR receiver_endpoint_id = ?1)
|
||||
AND state NOT IN ('completed', 'declined', 'cancelled', 'failed', 'deleted')
|
||||
"#,
|
||||
)
|
||||
.bind(peer_endpoint_id)
|
||||
.bind(now)
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(result.rows_affected())
|
||||
}
|
||||
|
||||
pub(crate) async fn protocol_ids_for_peer(
|
||||
&self,
|
||||
peer_endpoint_id: &str,
|
||||
) -> Result<Vec<u64>, VnidropError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT protocol_transfer_id FROM targeted_transfers
|
||||
WHERE sender_endpoint_id = ?1 OR receiver_endpoint_id = ?1
|
||||
"#,
|
||||
)
|
||||
.bind(peer_endpoint_id)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(rows
|
||||
.into_iter()
|
||||
.map(|row| row.get::<i64, _>(0) as u64)
|
||||
.collect())
|
||||
}
|
||||
|
||||
pub(crate) async fn mark_interrupted_in_flight(&self) -> Result<u64, VnidropError> {
|
||||
let now = now_ms();
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE targeted_transfers
|
||||
SET state = 'interrupted', updated_at = ?1
|
||||
WHERE state IN ('connecting', 'transferring')
|
||||
"#,
|
||||
)
|
||||
.bind(now)
|
||||
.execute(&self.pool)
|
||||
.await
|
||||
.map_err(VnidropError::repository)?;
|
||||
Ok(result.rows_affected())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub(crate) enum TargetedTransferRole {
|
||||
Sender,
|
||||
Receiver,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub(crate) struct TargetedTransferRow {
|
||||
pub(crate) id: String,
|
||||
pub(crate) protocol_transfer_id: u64,
|
||||
pub(crate) sender_endpoint_id: String,
|
||||
pub(crate) receiver_endpoint_id: String,
|
||||
pub(crate) manifest_id: String,
|
||||
pub(crate) content_hash: String,
|
||||
pub(crate) transfer_name: String,
|
||||
pub(crate) file_count: u64,
|
||||
pub(crate) total_size: u64,
|
||||
pub(crate) verified_bytes: u64,
|
||||
pub(crate) blob_ticket: Option<String>,
|
||||
pub(crate) authorization_secret_handle: Option<String>,
|
||||
pub(crate) role: TargetedTransferRole,
|
||||
pub(crate) state: TargetedTransferState,
|
||||
pub(crate) created_at: i64,
|
||||
pub(crate) updated_at: i64,
|
||||
}
|
||||
|
||||
fn row_to_transfer(row: sqlx::sqlite::SqliteRow) -> Result<TargetedTransfer, VnidropError> {
|
||||
Ok(TargetedTransfer {
|
||||
id: row.get("id"),
|
||||
sender_endpoint_id: row.get("sender_endpoint_id"),
|
||||
receiver_endpoint_id: row.get("receiver_endpoint_id"),
|
||||
manifest_id: row.get("manifest_id"),
|
||||
file_count: row.get::<i64, _>("file_count") as u64,
|
||||
total_size: row.get::<i64, _>("total_size") as u64,
|
||||
verified_bytes: row.try_get::<i64, _>("verified_bytes").unwrap_or(0) as u64,
|
||||
state: parse_state(&row.get::<String, _>("state"))?,
|
||||
created_at: row.get("created_at"),
|
||||
updated_at: row.get("updated_at"),
|
||||
})
|
||||
}
|
||||
|
||||
fn row_to_full(row: sqlx::sqlite::SqliteRow) -> Result<TargetedTransferRow, VnidropError> {
|
||||
Ok(TargetedTransferRow {
|
||||
id: row.get("id"),
|
||||
protocol_transfer_id: row.get::<i64, _>("protocol_transfer_id") as u64,
|
||||
sender_endpoint_id: row.get("sender_endpoint_id"),
|
||||
receiver_endpoint_id: row.get("receiver_endpoint_id"),
|
||||
manifest_id: row.get("manifest_id"),
|
||||
content_hash: row.get("content_hash"),
|
||||
transfer_name: row.get("transfer_name"),
|
||||
file_count: row.get::<i64, _>("file_count") as u64,
|
||||
total_size: row.get::<i64, _>("total_size") as u64,
|
||||
verified_bytes: row.try_get::<i64, _>("verified_bytes").unwrap_or(0) as u64,
|
||||
blob_ticket: row.try_get("blob_ticket").ok().flatten(),
|
||||
authorization_secret_handle: row.try_get("authorization_secret_handle").ok().flatten(),
|
||||
role: parse_role(
|
||||
&row.try_get::<String, _>("role")
|
||||
.unwrap_or_else(|_| "sender".to_string()),
|
||||
)?,
|
||||
state: parse_state(&row.get::<String, _>("state"))?,
|
||||
created_at: row.get("created_at"),
|
||||
updated_at: row.get("updated_at"),
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn state_as_str(state: TargetedTransferState) -> &'static str {
|
||||
match state {
|
||||
TargetedTransferState::Preparing => "preparing",
|
||||
TargetedTransferState::Offering => "offering",
|
||||
TargetedTransferState::AwaitingApproval => "awaiting_approval",
|
||||
TargetedTransferState::Approved => "approved",
|
||||
TargetedTransferState::Connecting => "connecting",
|
||||
TargetedTransferState::Transferring => "transferring",
|
||||
TargetedTransferState::Interrupted => "interrupted",
|
||||
TargetedTransferState::Completed => "completed",
|
||||
TargetedTransferState::Declined => "declined",
|
||||
TargetedTransferState::Cancelled => "cancelled",
|
||||
TargetedTransferState::Failed => "failed",
|
||||
TargetedTransferState::Deleted => "deleted",
|
||||
}
|
||||
}
|
||||
|
||||
fn role_as_str(role: TargetedTransferRole) -> &'static str {
|
||||
match role {
|
||||
TargetedTransferRole::Sender => "sender",
|
||||
TargetedTransferRole::Receiver => "receiver",
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_role(value: &str) -> Result<TargetedTransferRole, VnidropError> {
|
||||
match value {
|
||||
"sender" => Ok(TargetedTransferRole::Sender),
|
||||
"receiver" => Ok(TargetedTransferRole::Receiver),
|
||||
other => Err(VnidropError::repository(anyhow::anyhow!(
|
||||
"unknown targeted transfer role: {other}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_state(value: &str) -> Result<TargetedTransferState, VnidropError> {
|
||||
match value {
|
||||
"preparing" => Ok(TargetedTransferState::Preparing),
|
||||
"offering" => Ok(TargetedTransferState::Offering),
|
||||
"awaiting_approval" => Ok(TargetedTransferState::AwaitingApproval),
|
||||
"approved" => Ok(TargetedTransferState::Approved),
|
||||
"connecting" => Ok(TargetedTransferState::Connecting),
|
||||
"transferring" => Ok(TargetedTransferState::Transferring),
|
||||
"interrupted" => Ok(TargetedTransferState::Interrupted),
|
||||
"completed" => Ok(TargetedTransferState::Completed),
|
||||
"declined" => Ok(TargetedTransferState::Declined),
|
||||
"cancelled" => Ok(TargetedTransferState::Cancelled),
|
||||
"failed" => Ok(TargetedTransferState::Failed),
|
||||
"deleted" => Ok(TargetedTransferState::Deleted),
|
||||
other => Err(VnidropError::repository(anyhow::anyhow!(
|
||||
"unknown targeted transfer state: {other}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
@@ -22,6 +22,8 @@ mod limits_tests;
|
||||
mod network_config_tests;
|
||||
#[path = "tests/pairing_eligibility.rs"]
|
||||
mod pairing_eligibility_tests;
|
||||
#[path = "tests/persistence.rs"]
|
||||
mod persistence_tests;
|
||||
#[path = "tests/platform_contract_android.rs"]
|
||||
mod platform_contract_android_tests;
|
||||
#[cfg(any(target_os = "macos", target_os = "ios"))]
|
||||
|
||||
@@ -1,8 +1,7 @@
|
||||
use crate::{blocked_devices::BlockStore, repository::Repository};
|
||||
use crate::{blocked_devices::BlockStore, persistence, repository::Repository};
|
||||
|
||||
async fn store(temp: &tempfile::TempDir) -> BlockStore {
|
||||
let repository = Repository::open(temp.path()).await.unwrap();
|
||||
repository.blocked_devices()
|
||||
persistence::open_all(temp.path()).await.unwrap().blocked
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
@@ -24,7 +23,7 @@ async fn block_list_persists_and_unblocks() {
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn opening_repository_drops_unreleased_prototype_tables() {
|
||||
async fn opening_app_data_drops_unreleased_prototype_tables() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let db = temp.path().join("vnidrop.sqlite3");
|
||||
{
|
||||
@@ -45,8 +44,8 @@ async fn opening_repository_drops_unreleased_prototype_tables() {
|
||||
}
|
||||
}
|
||||
|
||||
let repository = Repository::open(temp.path()).await.unwrap();
|
||||
let pool = repository.sqlite_pool();
|
||||
let stores = persistence::open_all(temp.path()).await.unwrap();
|
||||
let pool = stores.pool_for_unmigrated();
|
||||
for table in ["contacts", "grants_issued", "grants_held", "held_offers"] {
|
||||
let row = sqlx::query(&format!(
|
||||
"SELECT COUNT(*) AS n FROM sqlite_master WHERE type = 'table' AND name = '{table}'"
|
||||
@@ -58,9 +57,7 @@ async fn opening_repository_drops_unreleased_prototype_tables() {
|
||||
assert_eq!(n, 0, "{table} must be dropped without migration");
|
||||
}
|
||||
|
||||
assert!(repository
|
||||
.blocked_devices()
|
||||
.is_blocked("keep-me")
|
||||
.await
|
||||
.unwrap());
|
||||
assert!(stores.blocked.is_blocked("keep-me").await.unwrap());
|
||||
// Invitation repository remains reachable from the bag.
|
||||
let _ = Repository::open(temp.path()).await.unwrap();
|
||||
}
|
||||
|
||||
13
crates/vnidrop/src/tests/persistence.rs
Normal file
13
crates/vnidrop/src/tests/persistence.rs
Normal file
@@ -0,0 +1,13 @@
|
||||
//! Persistence open returns domain stores without exporting a raw pool to callers.
|
||||
|
||||
use crate::persistence;
|
||||
|
||||
#[tokio::test]
|
||||
async fn open_all_returns_invitation_targeted_and_blocked_stores() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let stores = persistence::open_all(temp.path()).await.unwrap();
|
||||
|
||||
assert!(stores.blocked.list_blocked().await.unwrap().is_empty());
|
||||
assert!(stores.targeted.list().await.unwrap().is_empty());
|
||||
assert!(stores.invitation.list_transfers().await.unwrap().is_empty());
|
||||
}
|
||||
Reference in New Issue
Block a user