Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 5 additions & 4 deletions crates/bench/benches/special.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ use spacetimedb_lib::sats::{self, bsatn};
use spacetimedb_lib::{bsatn::ToBsatn as _, ProductValue};
use spacetimedb_schema::schema::TableSchema;
use spacetimedb_schema::table_name::TableName;
use spacetimedb_table::page_pool::PagePool;
use spacetimedb_table::tiered::{PageEvictionPolicy, PageManager};
use spacetimedb_testing::modules::{Csharp, ModuleLanguage, Rust};
use std::sync::Arc;
use std::sync::OnceLock;
Expand Down Expand Up @@ -141,24 +141,25 @@ fn serialize_benchmarks<
let mut table = spacetimedb_table::table::Table::new(
Arc::new(table_schema),
spacetimedb_table::indexes::SquashedOffset::COMMITTED_STATE,
PageManager::new_for_test().into(),
PageEvictionPolicy::NeverEvict,
);
let pool = PagePool::new_for_test();
let mut blob_store = spacetimedb_table::blob_store::HashMapBlobStore::default();

let ptrs = data_pv
.elements
.iter()
.map(|row| {
table
.insert(&pool, &mut blob_store, row.as_product().unwrap())
.insert(&mut blob_store, row.as_product().unwrap())
.unwrap()
.1
.pointer()
})
.collect::<Vec<_>>();
let refs = ptrs
.into_iter()
.map(|ptr| table.get_row_ref(&blob_store, ptr).unwrap())
.map(|ptr| table.get_row_ref(&blob_store, ptr).unwrap().unwrap())
.collect::<Vec<_>>();
group.bench_function(format!("bflatn_to_bsatn_slow_path/count={count}"), |b| {
b.iter(|| {
Expand Down
15 changes: 9 additions & 6 deletions crates/bench/src/spacetime_raw.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,10 @@ use crate::{
schemas::{table_name, BenchTable, IndexStrategy},
ResultBench,
};
use spacetimedb::db::relational_db::{tests_utils::TestDB, RelationalDB};
use spacetimedb::{
db::relational_db::{tests_utils::TestDB, RelationalDB},
error::DatastoreError,
};
use spacetimedb_datastore::execution_context::Workload;
use spacetimedb_primitives::{ColId, IndexId, TableId};
use spacetimedb_sats::{bsatn, AlgebraicValue};
Expand Down Expand Up @@ -120,8 +123,8 @@ impl BenchDatabase for SpacetimeRaw {
.db
.iter_mut(tx, *table_id)?
.take(row_count as usize)
.map(|row| row.to_product_value())
.collect::<Vec<_>>();
.map(|row| row.map(|r| r.to_product_value()))
.collect::<Result<Vec<_>, DatastoreError>>()?;

assert_eq!(rows.len(), row_count as usize, "not enough rows found for update_bulk!");
let mut scratch = Vec::new();
Expand All @@ -133,7 +136,7 @@ impl BenchDatabase for SpacetimeRaw {
.db
.iter_by_col_eq_mut(tx, *table_id, 0, &row.elements[0])?
.next()
.expect("failed to find row during update!")
.expect("failed to find row during update!")?
.pointer();

assert_eq!(
Expand All @@ -160,7 +163,7 @@ impl BenchDatabase for SpacetimeRaw {
fn iterate(&mut self, table_id: &Self::TableId) -> ResultBench<()> {
self.db.with_auto_commit(Workload::Internal, |tx| {
for row in self.db.iter_mut(tx, *table_id)? {
black_box(row);
black_box(row?);
}
Ok(())
})
Expand All @@ -174,7 +177,7 @@ impl BenchDatabase for SpacetimeRaw {
) -> ResultBench<()> {
self.db.with_auto_commit(Workload::Internal, |tx| {
for row in self.db.iter_by_col_eq_mut(tx, *table_id, col_id, &value)? {
black_box(row);
black_box(row?);
}
Ok(())
})
Expand Down
10 changes: 7 additions & 3 deletions crates/core/src/db/environment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,15 +32,18 @@ pub fn get(state: &impl StateView, key: &str) -> Result<Option<String>, Environm
state
.iter_by_col_eq(ST_ENV_ID, StEnvFields::Key, &AlgebraicValue::String(key.into()))?
.next()
.map(|row| Ok(StEnvRow::try_from(row)?.value))
.map(|row| {
let row = row.and_then(StEnvRow::try_from)?;
Ok(row.value)
})
.transpose()
}

pub fn snapshot(state: &impl StateView) -> Result<BTreeMap<String, String>, EnvironmentError> {
state
.iter(ST_ENV_ID)?
.map(|row| {
let row = StEnvRow::try_from(row)?;
let row = StEnvRow::try_from(row?)?;
Ok((row.key, row.value))
})
.collect()
Expand Down Expand Up @@ -81,7 +84,8 @@ fn delete(db: &RelationalDB, tx: &mut MutTx, key: &str) -> Result<bool, Environm
let pointer = tx
.iter_by_col_eq(ST_ENV_ID, StEnvFields::Key, &AlgebraicValue::String(key.into()))?
.next()
.map(|row| row.pointer());
.map(|row| row.map(|row| row.pointer()))
.transpose()?;
if let Some(pointer) = pointer {
db.delete(tx, ST_ENV_ID, [pointer]);
return Ok(true);
Expand Down
50 changes: 33 additions & 17 deletions crates/core/src/host/instance_env.rs
Original file line number Diff line number Diff line change
Expand Up @@ -197,10 +197,10 @@ impl ChunkedWriter {
self.chunks
}

pub fn collect_iter(
pub fn collect_iter<T: ToBsatn>(
pool: &mut ChunkPool,
iter: impl Iterator<Item = impl ToBsatn>,
) -> (Vec<Vec<u8>>, usize, usize) {
iter: impl Iterator<Item = Result<T, DatastoreError>>,
) -> Result<(Vec<Vec<u8>>, usize, usize), DatastoreError> {
// Track the number of rows and the number of bytes scanned by the iterator.
let mut rows_scanned = 0;
let mut bytes_scanned = 0;
Expand All @@ -209,6 +209,7 @@ impl ChunkedWriter {
// Consume the iterator, serializing each `item`,
// while allowing a chunk to be created at boundaries.
for item in iter {
let item = item?;
// Write the item directly to the BSATN `chunked_writer` buffer.
item.to_bsatn_extend(&mut chunked_writer.curr).unwrap();
// Flush at item boundaries.
Expand All @@ -222,7 +223,7 @@ impl ChunkedWriter {
// Update (BSATN) bytes scanned
bytes_scanned += chunks.iter().map(|chunk| chunk.len()).sum::<usize>();

(chunks, rows_scanned, bytes_scanned)
Ok((chunks, rows_scanned, bytes_scanned))
}
}

Expand Down Expand Up @@ -347,8 +348,9 @@ impl InstanceEnv {
.iter(ST_MODULE_ID)
.map_err(DBError::from)?
.next()
.ok_or_else(|| fail("database program is not initialized"))?;
let current_hash = spacetimedb_datastore::system_tables::read_hash_from_col(row, StModuleFields::ProgramHash)
.ok_or_else(|| fail("database program is not initialized"))?
.map_err(DBError::from)?;
let current_hash = spacetimedb_datastore::system_tables::read_hash_from_col(&row, StModuleFields::ProgramHash)
.map_err(DBError::from)?;
if current_hash != *hash {
return Err(fail("module was replaced while this function was running"));
Expand Down Expand Up @@ -421,7 +423,7 @@ impl InstanceEnv {
/// and return the full length of the BSATN.
///
/// Assumes that the full encoding of `cols` will fit in `buffer`.
fn project_cols_bsatn(buffer: &mut [u8], cols: ColList, row_ref: RowRef<'_>) -> usize {
fn project_cols_bsatn(buffer: &mut [u8], cols: ColList, row_ref: &RowRef<'_>) -> usize {
// We get back a col-list with the columns with generated values.
// Write those back to `buffer` and then the encoded length to `row_len`.
let (_, count) = CountWriter::run(buffer, |writer| {
Expand Down Expand Up @@ -453,7 +455,7 @@ impl InstanceEnv {
let (row_len, row_ptr, insert_flags) = stdb
.insert(tx, table_id, buffer)
.map(|(gen_cols, row_ref, insert_flags)| {
let row_len = Self::project_cols_bsatn(buffer, gen_cols, row_ref);
let row_len = Self::project_cols_bsatn(buffer, gen_cols, &row_ref);
(row_len, row_ref.pointer(), insert_flags)
})
.inspect_err(
Expand Down Expand Up @@ -535,7 +537,7 @@ impl InstanceEnv {
let (row_len, row_ptr, update_flags) = stdb
.update(tx, table_id, index_id, buffer)
.map(|(gen_cols, row_ref, update_flags)| {
let row_len = Self::project_cols_bsatn(buffer, gen_cols, row_ref);
let row_len = Self::project_cols_bsatn(buffer, gen_cols, &row_ref);
(row_len, row_ref.pointer(), update_flags)
})
.inspect_err(
Expand Down Expand Up @@ -576,7 +578,10 @@ impl InstanceEnv {
let (table_id, _, iter) = stdb.index_scan_point(tx, index_id, point)?;
Self::require_module_table(table_id)?;
// Re. `SmallVec`, `delete_by_field` only cares about 1 element, so optimize for that.
let rows_to_delete = iter.map(|row_ref| row_ref.pointer()).collect::<SmallVec<[_; 1]>>();
let rows_to_delete = iter
.map(|row_ref| row_ref.map(|rr| rr.pointer()))
.collect::<Result<SmallVec<[_; 1]>, DatastoreError>>()
.map_err(DBError::from)?;

Ok(Self::datastore_delete_by_index_scan(stdb, tx, table_id, rows_to_delete))
}
Expand All @@ -598,8 +603,14 @@ impl InstanceEnv {
Self::require_module_table(table_id)?;
// Re. `SmallVec`, `delete_by_field` only cares about 1 element, so optimize for that.
let rows_to_delete = match iter {
IndexScanPointOrRange::Point(_, iter) => iter.map(|row_ref| row_ref.pointer()).collect(),
IndexScanPointOrRange::Range(iter) => iter.map(|row_ref| row_ref.pointer()).collect(),
IndexScanPointOrRange::Point(_, iter) => iter
.map(|row_ref| row_ref.map(|rr| rr.pointer()))
.collect::<Result<_, DatastoreError>>()
.map_err(DBError::from)?,
IndexScanPointOrRange::Range(iter) => iter
.map(|row_ref| row_ref.map(|rr| rr.pointer()))
.collect::<Result<_, DatastoreError>>()
.map_err(DBError::from)?,
};

Ok(Self::datastore_delete_by_index_scan(stdb, tx, table_id, rows_to_delete))
Expand Down Expand Up @@ -733,7 +744,7 @@ impl InstanceEnv {
let iter = self.relational_db().iter_mut(tx, table_id)?;

// Scan the index and serialize rows to BSATN.
let (chunks, rows_scanned, bytes_scanned) = ChunkedWriter::collect_iter(pool, iter);
let (chunks, rows_scanned, bytes_scanned) = ChunkedWriter::collect_iter(pool, iter).map_err(DBError::from)?;

// Record the number of rows and the number of bytes scanned by the iterator.
tx.metrics.bytes_scanned += bytes_scanned;
Expand All @@ -758,7 +769,7 @@ impl InstanceEnv {
Self::require_module_table(table_id)?;

// Scan the index and serialize rows to BSATN.
let (chunks, rows_scanned, bytes_scanned) = ChunkedWriter::collect_iter(pool, iter);
let (chunks, rows_scanned, bytes_scanned) = ChunkedWriter::collect_iter(pool, iter).map_err(DBError::from)?;

// Record the number of rows and the number of bytes scanned by the iterator.
tx.metrics.index_seeks += 1;
Expand Down Expand Up @@ -790,8 +801,13 @@ impl InstanceEnv {

// Scan the index and serialize rows to BSATN.
let (point, (chunks, rows_scanned, bytes_scanned)) = match iter {
IndexScanPointOrRange::Point(point, iter) => (Some(point), ChunkedWriter::collect_iter(pool, iter)),
IndexScanPointOrRange::Range(iter) => (None, ChunkedWriter::collect_iter(pool, iter)),
IndexScanPointOrRange::Point(point, iter) => (
Some(point),
ChunkedWriter::collect_iter(pool, iter).map_err(DBError::from)?,
),
IndexScanPointOrRange::Range(iter) => {
(None, ChunkedWriter::collect_iter(pool, iter).map_err(DBError::from)?)
}
};

// Record the number of rows and the number of bytes scanned by the iterator.
Expand Down Expand Up @@ -1618,7 +1634,7 @@ mod test {
&schema,
&BTreeMap::from([("A".into(), "uncommitted".into()), ("MISSING".into(), "".into())]),
)?;
assert!(tx.views_for_refresh().any(|dependency| dependency == &view));
assert!(tx.views_for_refresh()?.any(|dependency| dependency == &view));
let (_, metrics, reducer) = db.rollback_mut_tx(tx);
db.report_mut_tx_metrics(reducer, metrics, None);
env.start_funcall(
Expand Down
9 changes: 8 additions & 1 deletion crates/core/src/host/module_host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3100,7 +3100,14 @@ impl ModuleHost {
let mut abi_duration = Duration::ZERO;
let mut trapped = false;
let mut num_views_evaluated = 0;
for view_call in tx.views_for_refresh().cloned().collect::<Vec<_>>() {
let views_for_refresh = match tx.views_for_refresh() {
Ok(views) => views.cloned().collect::<Vec<_>>(),
Err(error) => {
outcome = ViewOutcome::Failed(format!("failed to find views for refresh: {error}"));
Vec::new()
}
};
for view_call in views_for_refresh {
let resolved = match resolve_view_for_refresh(&tx, module_def, &view_call) {
Ok(resolved) => resolved,
Err(err) => {
Expand Down
5 changes: 4 additions & 1 deletion crates/core/src/host/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ impl SchedulerStarter {

// Find all Scheduled tables
for st_scheduled_row in self.db.iter(&tx, ST_SCHEDULED_ID)? {
let st_scheduled_row = st_scheduled_row?;
let table_id = st_scheduled_row.read_col(StScheduledFields::TableId)?;
let function_name =
Arc::<str>::from(st_scheduled_row.read_col::<Box<str>>(StScheduledFields::ReducerName)?);
Expand All @@ -117,6 +118,7 @@ impl SchedulerStarter {

// Insert each entry (row) in the scheduled table into `queue`.
for scheduled_row in self.db.iter(&tx, table_id)? {
let scheduled_row = scheduled_row?;
let (schedule_id, schedule_at) = get_schedule_from_row(&scheduled_row, id_column, at_column)?;
let row_hash = scheduled_row_hash(&scheduled_row)?;
// calculate duration left to call the scheduled reducer
Expand Down Expand Up @@ -1032,7 +1034,8 @@ fn get_schedule_row_mut<'a>(
) -> anyhow::Result<Option<RowRef<'a>>> {
Ok(db
.iter_by_col_eq_mut(tx, id.table_id, id.id_column, &id.schedule_id.into())?
.next())
.next()
.transpose()?)
}

/// Helper to get `schedule_id` and `schedule_at`
Expand Down
8 changes: 7 additions & 1 deletion crates/core/src/host/v8/syscall/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ use crate::host::wasm_common::{RowIterIdx, TimingSpan, TimingSpanIdx};
use anyhow::Context;
use bytes::Bytes;
use spacetimedb_datastore::locking_tx_datastore::{FuncCallType, MutTxId, ViewCallInfo};
use spacetimedb_engine::error::DBError;
use spacetimedb_lib::{ConnectionId, Identity, RawModuleDef, Timestamp};
use spacetimedb_primitives::{ColId, IndexId, ProcedureId, TableId, ViewFnPtr};
use spacetimedb_sats::bsatn;
Expand Down Expand Up @@ -744,7 +745,12 @@ fn refresh_views(
hooks: &HookFunctions<'_>,
module_def: &ModuleDef,
) -> SysCallResult<MutTxId> {
let views_for_refresh = tx.views_for_refresh().cloned().collect::<Vec<_>>();
let views_for_refresh = tx
.views_for_refresh()
.map_err(DBError::from)
.map_err(NodesError::from)?
.cloned()
.collect::<Vec<_>>();
let stdb = get_env(scope)?.instance_env.relational_db().clone();
let database_identity = *get_env(scope)?.instance_env.database_identity();
let mut tx_slot = get_env(scope)?.instance_env.tx.clone();
Expand Down
24 changes: 13 additions & 11 deletions crates/core/src/host/wasm_common/module_host_actor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -651,13 +651,13 @@ struct UpdateEffects {
}

impl UpdateEffects {
fn after_migration(result: crate::db::update::UpdateResult, tx: &MutTxId) -> Self {
fn after_migration(result: crate::db::update::UpdateResult, tx: &MutTxId) -> Result<Self, DatastoreError> {
use crate::db::update::UpdateResult;
Self {
Ok(Self {
refresh_views: matches!(result, UpdateResult::EvaluateSubscribedViews)
|| tx.views_for_refresh().next().is_some(),
|| tx.views_for_refresh()?.next().is_some(),
disconnect_clients: matches!(result, UpdateResult::RequiresClientDisconnect),
}
})
}

fn committed(
Expand Down Expand Up @@ -741,9 +741,10 @@ impl InstanceCommon {
let row = tx
.iter(ST_MODULE_ID)?
.next()
.context("database program is not initialized")?;
.context("database program is not initialized")?
.context("failed to load database program")?;
let current_hash =
spacetimedb_datastore::system_tables::read_hash_from_col(row, StModuleFields::ProgramHash)?;
spacetimedb_datastore::system_tables::read_hash_from_col(&row, StModuleFields::ProgramHash)?;
anyhow::ensure!(
current_hash == old_module_info.module_hash,
"database program changed before publication"
Expand Down Expand Up @@ -794,7 +795,7 @@ impl InstanceCommon {
};
let durable_offset = stdb.durable_tx_offset();

let effects = UpdateEffects::after_migration(res, &tx);
let effects = UpdateEffects::after_migration(res, &tx)?;
let res = if effects.refresh_views {
// Resolve surviving materializations through the new module,
// even when this migration also requires client disconnection.
Expand Down Expand Up @@ -841,9 +842,10 @@ impl InstanceCommon {
let row = tx
.iter(ST_MODULE_ID)?
.next()
.context("database program is not initialized")?;
.context("database program is not initialized")?
.context("failed to load database program")?;
anyhow::ensure!(
read_hash_from_col(row, StModuleFields::ProgramHash)? == self.info.module_hash,
read_hash_from_col(&row, StModuleFields::ProgramHash)? == self.info.module_hash,
"database program changed before publication"
);
crate::db::environment::replace(db, tx, self.info.module_def.environment(), &environment)?;
Expand Down Expand Up @@ -2232,8 +2234,8 @@ mod tests {
);
let result = update::update_database(&db, &mut tx, AuthCtx::for_testing(), plan, &TestLogger)?;
assert!(matches!(result, update::UpdateResult::RequiresClientDisconnect));
assert!(tx.views_for_refresh().any(|dirty| *dirty == call));
let effects = UpdateEffects::after_migration(result, &tx);
assert!(tx.views_for_refresh()?.any(|dirty| *dirty == call));
let effects = UpdateEffects::after_migration(result, &tx)?;
assert!(effects.refresh_views);
assert!(effects.disconnect_clients);
let calls = collect_subscribed_view_calls(&tx, &new, Identity::ZERO)?;
Expand Down
Loading
Loading