diff --git a/vm2/framework/starcoin-framework/tests/in_place_delayed_field_tests.move b/vm2/framework/starcoin-framework/tests/in_place_delayed_field_tests.move new file mode 100644 index 0000000000..d1cd1fb932 --- /dev/null +++ b/vm2/framework/starcoin-framework/tests/in_place_delayed_field_tests.move @@ -0,0 +1,48 @@ +#[test_only] +module starcoin_framework::in_place_delayed_field_tests { + use std::signer; + use starcoin_framework::aggregator_v2; + + struct DelayedFieldHolder has key { + marker: u64, + amount: aggregator_v2::Aggregator, + } + + fun publish_holder(account: &signer, marker: u64, start: u64) { + move_to( + account, + DelayedFieldHolder { + marker, + amount: aggregator_v2::create_unbounded_aggregator_with_value(start), + }, + ); + } + + fun add_amount_only(addr: address, delta: u64) acquires DelayedFieldHolder { + let holder = borrow_global_mut(addr); + aggregator_v2::add(&mut holder.amount, delta); + } + + fun read_amount(addr: address): u64 acquires DelayedFieldHolder { + let holder = borrow_global(addr); + aggregator_v2::read(&holder.amount) + } + + fun read_marker(addr: address): u64 acquires DelayedFieldHolder { + let holder = borrow_global(addr); + holder.marker + } + + #[test(account = @starcoin_framework)] + fun test_write_delayed_field_in_place(account: signer) acquires DelayedFieldHolder { + let addr = signer::address_of(&account); + publish_holder(&account, 7, 10); + + // Only mutate delayed field (aggregator) in Move layer. + add_amount_only(addr, 2); + add_amount_only(addr, 5); + + assert!(read_marker(addr) == 7, 0); + assert!(read_amount(addr) == 17, 1); + } +} diff --git a/vm2/vm-runtime/src/move_vm_ext/session.rs b/vm2/vm-runtime/src/move_vm_ext/session.rs index 4ae15c13dc..8fffc970d3 100644 --- a/vm2/vm-runtime/src/move_vm_ext/session.rs +++ b/vm2/vm-runtime/src/move_vm_ext/session.rs @@ -847,3 +847,123 @@ impl SessionExt<'_, '_> { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::data_cache::AsMoveResolver; + use bytes::Bytes; + use move_core_types::{ + identifier::Identifier, + language_storage::StructTag, + value::{MoveStructLayout, MoveTypeLayout}, + }; + use starcoin_vm_runtime_types::abstract_write_op::AbstractResourceWriteOp; + use starcoin_vm_types::state_store::{ + in_memory_state_view::InMemoryStateView, + state_value::{StateValue, StateValueMetadata}, + }; + use std::{collections::HashMap, sync::Arc}; + + fn delayed_layout() -> Arc { + Arc::new(MoveTypeLayout::Struct(MoveStructLayout::Runtime(vec![ + MoveTypeLayout::U64, + MoveTypeLayout::Native( + move_core_types::value::IdentifierMappingKind::Aggregator, + Box::new(MoveTypeLayout::U64), + ), + ]))) + } + + #[test] + fn convert_change_set_produces_in_place_write_for_delayed_read_only_key() { + let key = StateKey::raw(b"delayed-read-only-key"); + let state_view = InMemoryStateView::new(HashMap::new()); + let resolver = state_view.as_move_resolver(); + let woc = WriteOpConverter::new(&resolver, false); + + let aggregator_change_set = AggregatorChangeSet { + aggregator_v1_changes: BTreeMap::new(), + delayed_field_changes: BTreeMap::new(), + reads_needing_exchange: BTreeMap::from([( + key.clone(), + (StateValueMetadata::none(), 16, delayed_layout()), + )]), + group_reads_needing_exchange: BTreeMap::new(), + }; + + let (change_set, _) = SessionExt::convert_change_set( + &woc, + ChangeSet::new(), + ResourceGroupChangeSet::V1(BTreeMap::new()), + vec![], + TableChangeSet::default(), + aggregator_change_set, + false, + ) + .expect("convert_change_set should succeed"); + + assert!(matches!( + change_set.resource_write_set().get(&key), + Some(AbstractResourceWriteOp::InPlaceDelayedFieldChange(_)) + )); + } + + #[test] + fn convert_change_set_filters_in_place_when_same_key_has_resource_write() { + let addr = AccountAddress::ONE; + let tag = StructTag { + address: addr, + module: Identifier::new("DelayedOnly").unwrap(), + name: Identifier::new("SameTxn").unwrap(), + type_args: vec![], + }; + let key = resource_state_key(&addr, &tag).unwrap(); + + let mut state_data = HashMap::new(); + state_data.insert(key.clone(), StateValue::new_legacy(Bytes::from(vec![9u8]))); + let state_view = InMemoryStateView::new(state_data); + let resolver = state_view.as_move_resolver(); + let woc = WriteOpConverter::new(&resolver, false); + + let mut change_set = ChangeSet::new(); + change_set + .add_account_changeset( + addr, + AccountChangeSet::from_modules_resources( + BTreeMap::new(), + BTreeMap::from([( + tag, + MoveStorageOp::Modify((Bytes::from(vec![1u8, 2u8]), None)), + )]), + ), + ) + .expect("add account changeset should succeed"); + + let aggregator_change_set = AggregatorChangeSet { + aggregator_v1_changes: BTreeMap::new(), + delayed_field_changes: BTreeMap::new(), + reads_needing_exchange: BTreeMap::from([( + key.clone(), + (StateValueMetadata::none(), 16, delayed_layout()), + )]), + group_reads_needing_exchange: BTreeMap::new(), + }; + + let (converted, _) = SessionExt::convert_change_set( + &woc, + change_set, + ResourceGroupChangeSet::V1(BTreeMap::new()), + vec![], + TableChangeSet::default(), + aggregator_change_set, + false, + ) + .expect("convert_change_set should succeed"); + + assert!(!matches!( + converted.resource_write_set().get(&key), + Some(AbstractResourceWriteOp::InPlaceDelayedFieldChange(_)) + )); + } +} diff --git a/vm2/vm-runtime/src/parallel_executor/mod.rs b/vm2/vm-runtime/src/parallel_executor/mod.rs index fcc8cfeca7..b9566c6bf0 100644 --- a/vm2/vm-runtime/src/parallel_executor/mod.rs +++ b/vm2/vm-runtime/src/parallel_executor/mod.rs @@ -45,7 +45,7 @@ use starcoin_vm_types::{ state_store::state_key::StateKey, state_store::StateView, transaction::{Transaction, TransactionOutput, TransactionStatus}, - write_set::WriteOp, + write_set::{TransactionWrite, WriteOp}, }; use std::collections::{BTreeMap, HashMap, HashSet}; use std::sync::Arc; @@ -425,6 +425,7 @@ fn materialize_parallel_outputs( &vm_output, &mapping, delayed_field_cache, + state_view, max_value_nest_depth, )? }; @@ -457,11 +458,21 @@ fn materialize_parallel_outputs( let mut has_agg_v1 = false; let mut group_touches = 0u64; let mut group_touch_counts: HashMap = HashMap::new(); + let mut delayed_resource_touch_counts: HashMap = HashMap::new(); for (_, output) in outputs.iter() { if !output.output.aggregator_v1_delta_set().is_empty() { has_agg_v1 = true; } for (key, op) in output.output.resource_write_set() { + if matches!( + op, + AbstractResourceWriteOp::WriteWithDelayedFields(_) + | AbstractResourceWriteOp::InPlaceDelayedFieldChange(_) + ) { + *delayed_resource_touch_counts + .entry(key.clone()) + .or_insert(0) += 1; + } if is_group_write_op(op) { group_touches += 1; *group_touch_counts.entry(key.clone()).or_insert(0) += 1; @@ -469,7 +480,10 @@ fn materialize_parallel_outputs( } } let has_group_dup = group_touch_counts.values().any(|count| *count > 1); - let needs_sequential = has_agg_v1 || has_group_dup; + let has_delayed_resource_dup = delayed_resource_touch_counts + .values() + .any(|count| *count > 1); + let needs_sequential = has_agg_v1 || has_group_dup || has_delayed_resource_dup; outputs.sort_by_key(|(idx, _)| *idx); if !needs_sequential { @@ -493,9 +507,10 @@ fn materialize_parallel_outputs( info!( target: "vm-bench", - "materialize sequential: agg_v1={} group_dup={} group_touches={}", + "materialize sequential: agg_v1={} group_dup={} delayed_resource_dup={} group_touches={}", has_agg_v1, has_group_dup, + has_delayed_resource_dup, group_touches ); let mut state_cache = StateViewCache::new(state_view); @@ -632,6 +647,7 @@ fn materialize_resource_write_set_no_groups( output: &VMOutput, mapping: &impl ValueToIdentifierMapping, delayed_field_cache: &DelayedFieldCache, + state_view: &impl StateView, max_value_nest_depth: Option, ) -> Result, VMStatus> { let mut patched = Vec::new(); @@ -662,6 +678,7 @@ fn materialize_resource_write_set_no_groups( metadata, mapping, delayed_field_cache, + state_view, max_value_nest_depth, )?, AbstractResourceWriteOp::WriteResourceGroup(_) @@ -736,6 +753,7 @@ pub(crate) fn materialize_resource_write_set( metadata, mapping, delayed_field_cache, + state_view, max_value_nest_depth, )? } @@ -830,17 +848,37 @@ fn materialize_in_place_change( metadata: &starcoin_vm_types::state_store::state_value::StateValueMetadata, mapping: &impl ValueToIdentifierMapping, delayed_field_cache: &DelayedFieldCache, + state_view: &impl StateView, max_value_nest_depth: Option, ) -> Result { - let base = delayed_field_cache.get_base_value(key).ok_or_else(|| { - VMStatus::error( - StatusCode::DELAYED_MATERIALIZATION_CODE_INVARIANT_ERROR, - Some(format!( - "Missing cached base value for delayed field exchange: {:?}", - key - )), - ) - })?; + let base = if let Some(cached_base) = delayed_field_cache.get_base_value(key) { + cached_base + } else { + let state_value = state_view + .get_state_value(key) + .map_err(|err| { + VMStatus::error( + StatusCode::UNKNOWN_INVARIANT_VIOLATION_ERROR, + Some(format!( + "Failed to read base state value for delayed field exchange {:?}: {:?}", + key, err + )), + ) + })? + .ok_or_else(|| { + VMStatus::error( + StatusCode::DELAYED_MATERIALIZATION_CODE_INVARIANT_ERROR, + Some(format!( + "Missing base state value for delayed field exchange: {:?}", + key + )), + ) + })?; + + let base_from_state = WriteOp::from_state_value(Some(state_value)); + delayed_field_cache.insert_base_value(key.clone(), base_from_state.clone(), true); + base_from_state + }; let bytes = base .bytes() .ok_or_else(|| { @@ -1134,6 +1172,7 @@ mod tests { use starcoin_aggregator::bounded_math::SignedU128; use starcoin_aggregator::delta_change_set::DeltaOp; use starcoin_aggregator::delta_math::DeltaHistory; + use starcoin_aggregator::types::DelayedFieldValue; use starcoin_vm_runtime_types::change_set::VMChangeSet; use starcoin_vm_runtime_types::module_write_set::ModuleWriteSet; use starcoin_vm_runtime_types::resolver::ResourceGroupSize; @@ -1145,6 +1184,8 @@ mod tests { use starcoin_vm_types::write_set::WriteOp; use std::time::{Duration, Instant}; + mod tests_delayed_from_state; + struct TestMapping { value: u64, } @@ -1224,6 +1265,7 @@ mod tests { &StateValueMetadata::none(), &mapping, &delayed_field_cache, + &InMemoryStateView::new(HashMap::new()), Some(DEFAULT_MAX_VALUE_NEST_DEPTH), ) .unwrap(); @@ -1342,6 +1384,116 @@ mod tests { (outputs, InMemoryStateView::new(state_data), agg_key) } + fn build_delayed_only_dependency_case( + consumer_count: usize, + ) -> ( + Vec<(usize, StarcoinTransactionOutput)>, + InMemoryStateView, + VersionedDelayedFields, + StateKey, + ) { + let address = AccountAddress::from_hex_literal("0x1").unwrap(); + let struct_tag = StructTag { + address, + module: Identifier::new("DelayedOnly").unwrap(), + name: Identifier::new("Resource").unwrap(), + type_args: vec![], + }; + let state_key = StateKey::resource(&address, &struct_tag).unwrap(); + let layout = Arc::new(MoveTypeLayout::Struct(MoveStructLayout::Runtime(vec![ + MoveTypeLayout::U64, + MoveTypeLayout::Native( + move_core_types::value::IdentifierMappingKind::Aggregator, + Box::new(MoveTypeLayout::U64), + ), + ]))); + let delayed_id = DelayedFieldID::new_with_width(101, 8); + let delayed_value = Value::struct_(Struct::pack(vec![ + Value::u64(1), + Value::delayed_value(delayed_id), + ])); + let delayed_bytes = ValueSerDeContext::new(Some(DEFAULT_MAX_VALUE_NEST_DEPTH)) + .with_delayed_fields_serde() + .serialize(&delayed_value, layout.as_ref()) + .unwrap() + .unwrap(); + let delayed_bytes = Bytes::from(delayed_bytes); + + let mut outputs = Vec::with_capacity(consumer_count + 1); + let mut first_write_set = BTreeMap::new(); + first_write_set.insert( + state_key.clone(), + AbstractResourceWriteOp::WriteWithDelayedFields(WriteWithDelayedFieldsOp { + write_op: WriteOp::Modification { + data: delayed_bytes.clone(), + metadata: StateValueMetadata::none(), + }, + layout: layout.clone(), + materialized_size: Some(delayed_bytes.len() as u64), + }), + ); + outputs.push(( + 0, + StarcoinTransactionOutput::new( + VMOutput::new( + VMChangeSet::new( + first_write_set, + vec![], + BTreeMap::new(), + BTreeMap::new(), + BTreeMap::new(), + ), + ModuleWriteSet::empty(), + FeeStatement::zero(), + TransactionStatus::Keep(KeptVMStatus::Executed), + TransactionAuxiliaryData::None, + ), + HashMap::new(), + ), + )); + + for txn_idx in 1..=consumer_count { + let mut write_set = BTreeMap::new(); + write_set.insert( + state_key.clone(), + AbstractResourceWriteOp::InPlaceDelayedFieldChange(InPlaceDelayedFieldChangeOp { + layout: layout.clone(), + materialized_size: delayed_bytes.len() as u64, + metadata: StateValueMetadata::none(), + }), + ); + outputs.push(( + txn_idx, + StarcoinTransactionOutput::new( + VMOutput::new( + VMChangeSet::new( + write_set, + vec![], + BTreeMap::new(), + BTreeMap::new(), + BTreeMap::new(), + ), + ModuleWriteSet::empty(), + FeeStatement::zero(), + TransactionStatus::Keep(KeptVMStatus::Executed), + TransactionAuxiliaryData::None, + ), + HashMap::new(), + ), + )); + } + + let delayed_fields = VersionedDelayedFields::empty(); + delayed_fields.set_base_value(delayed_id, DelayedFieldValue::Aggregator(7)); + + ( + outputs, + InMemoryStateView::new(HashMap::new()), + delayed_fields, + state_key, + ) + } + fn materialize_parallel_outputs_legacy_all_seq( mut outputs: Vec<(usize, StarcoinTransactionOutput)>, delayed_fields: VersionedDelayedFields, @@ -1454,6 +1606,34 @@ mod tests { assert_eq!(mixed, legacy); } + #[test] + fn delayed_only_dependency_materialization_matches_legacy_all_seq() { + let (mixed_inputs, mixed_state_view, mixed_delayed_fields, mixed_key) = + build_delayed_only_dependency_case(64); + let (legacy_inputs, legacy_state_view, legacy_delayed_fields, legacy_key) = + build_delayed_only_dependency_case(64); + assert_eq!(mixed_key, legacy_key); + + let mixed = materialize_parallel_outputs( + mixed_inputs, + mixed_delayed_fields, + Arc::new(DelayedFieldCache::default()), + &mixed_state_view, + Some(DEFAULT_MAX_VALUE_NEST_DEPTH), + ) + .expect("delayed-only dependency case should materialize successfully"); + let legacy = materialize_parallel_outputs_legacy_all_seq( + legacy_inputs, + legacy_delayed_fields, + Arc::new(DelayedFieldCache::default()), + &legacy_state_view, + Some(DEFAULT_MAX_VALUE_NEST_DEPTH), + ) + .expect("legacy all-seq materialization should succeed"); + + assert_eq!(mixed, legacy); + } + #[test] #[ignore = "manual benchmark; run with -- --ignored --nocapture"] fn bench_mixed_materialization_sparse_conflicts() { diff --git a/vm2/vm-runtime/src/parallel_executor/tests/tests_delayed_from_state.rs b/vm2/vm-runtime/src/parallel_executor/tests/tests_delayed_from_state.rs new file mode 100644 index 0000000000..fac75f0a50 --- /dev/null +++ b/vm2/vm-runtime/src/parallel_executor/tests/tests_delayed_from_state.rs @@ -0,0 +1,106 @@ +use super::*; +use bytes::Bytes; +use move_core_types::account_address::AccountAddress; +use move_core_types::identifier::Identifier; +use move_core_types::language_storage::StructTag; +use move_core_types::value::{MoveStructLayout, MoveTypeLayout}; +use move_core_types::vm_status::KeptVMStatus; +use move_vm_types::delayed_values::delayed_field_id::DelayedFieldID; +use move_vm_types::value_serde::ValueSerDeContext; +use move_vm_types::values::{Struct, Value}; +use starcoin_aggregator::types::DelayedFieldValue; +use starcoin_vm_runtime_types::change_set::VMChangeSet; +use starcoin_vm_runtime_types::module_write_set::ModuleWriteSet; +use starcoin_vm_types::fee_statement::FeeStatement; +use starcoin_vm_types::state_store::in_memory_state_view::InMemoryStateView; +use starcoin_vm_types::state_store::state_value::StateValue; +use starcoin_vm_types::state_store::state_value::StateValueMetadata; +use starcoin_vm_types::transaction::TransactionAuxiliaryData; + +fn build_delayed_only_from_state_case( + txn_count: usize, +) -> ( + Vec<(usize, StarcoinTransactionOutput)>, + InMemoryStateView, + VersionedDelayedFields, +) { + let address = AccountAddress::from_hex_literal("0x1").unwrap(); + let struct_tag = StructTag { + address, + module: Identifier::new("DelayedOnly").unwrap(), + name: Identifier::new("FromState").unwrap(), + type_args: vec![], + }; + let state_key = StateKey::resource(&address, &struct_tag).unwrap(); + let layout = Arc::new(MoveTypeLayout::Struct(MoveStructLayout::Runtime(vec![ + MoveTypeLayout::U64, + MoveTypeLayout::Native( + move_core_types::value::IdentifierMappingKind::Aggregator, + Box::new(MoveTypeLayout::U64), + ), + ]))); + let delayed_id = DelayedFieldID::new_with_width(202, 8); + let state_value = Value::struct_(Struct::pack(vec![ + Value::u64(5), + Value::delayed_value(delayed_id), + ])); + let state_bytes = ValueSerDeContext::new(Some(DEFAULT_MAX_VALUE_NEST_DEPTH)) + .with_delayed_fields_serde() + .serialize(&state_value, layout.as_ref()) + .unwrap() + .unwrap(); + let state_bytes = Bytes::from(state_bytes); + + let mut outputs = Vec::with_capacity(txn_count); + for txn_idx in 0..txn_count { + let mut write_set = BTreeMap::new(); + write_set.insert( + state_key.clone(), + AbstractResourceWriteOp::InPlaceDelayedFieldChange(InPlaceDelayedFieldChangeOp { + layout: layout.clone(), + materialized_size: state_bytes.len() as u64, + metadata: StateValueMetadata::none(), + }), + ); + outputs.push(( + txn_idx, + StarcoinTransactionOutput::new( + VMOutput::new( + VMChangeSet::new( + write_set, + vec![], + BTreeMap::new(), + BTreeMap::new(), + BTreeMap::new(), + ), + ModuleWriteSet::empty(), + FeeStatement::zero(), + TransactionStatus::Keep(KeptVMStatus::Executed), + TransactionAuxiliaryData::None, + ), + HashMap::new(), + ), + )); + } + + let delayed_fields = VersionedDelayedFields::empty(); + delayed_fields.set_base_value(delayed_id, DelayedFieldValue::Aggregator(11)); + + let mut state_data = HashMap::new(); + state_data.insert(state_key, StateValue::new_legacy(state_bytes)); + + (outputs, InMemoryStateView::new(state_data), delayed_fields) +} + +#[test] +fn delayed_only_from_state_without_cached_base_should_materialize() { + let (outputs, state_view, delayed_fields) = build_delayed_only_from_state_case(32); + let _ = materialize_parallel_outputs( + outputs, + delayed_fields, + Arc::new(DelayedFieldCache::default()), + &state_view, + Some(DEFAULT_MAX_VALUE_NEST_DEPTH), + ) + .expect("delayed-only in-place changes should materialize from state base"); +}