Compare commits
	
		
			No commits in common. "0367d70ef34e132caa4dea6a865dc5c38693546f" and "5ccd432566c6b95d7c764a2ef67b8a793f8775a2" have entirely different histories.
		
	
	
		
			0367d70ef3
			...
			5ccd432566
		
	
		
							
								
								
									
										153
									
								
								src/db/vm.rs
									
									
									
									
									
								
							
							
								
								
								
								
								
									
									
								
							
						
						
									
										153
									
								
								src/db/vm.rs
									
									
									
									
									
								
							| @ -191,120 +191,63 @@ impl NewVmReq { | ||||
|     } | ||||
| 
 | ||||
|     pub async fn submit_error(db: &Surreal<Client>, id: &str, error: String) -> Result<(), Error> { | ||||
|         let tx_query = String::from( | ||||
|             " | ||||
|             BEGIN TRANSACTION; | ||||
| 
 | ||||
|                 LET $new_vm_req = $new_vm_req_input; | ||||
|                 LET $error = $error_input; | ||||
|                 LET $record = (select * from $new_vm_req)[0]; | ||||
|                 LET $admin = $record.in; | ||||
| 
 | ||||
|                 if $record == None {{ | ||||
|                     THROW 'vm req not exist ' + <string>$new_vm_req | ||||
|                 }}; | ||||
| 
 | ||||
|                 UPDATE $new_vm_req SET error = $error; | ||||
| 
 | ||||
|                 UPDATE $admin SET tmp_locked -= $record.locked_nano; | ||||
|                 UPDATE $admin SET balance += $record.locked_nano; | ||||
| 
 | ||||
|             COMMIT TRANSACTION;",
 | ||||
|         ); | ||||
| 
 | ||||
|         log::trace!("submit_new_vm_err query: {tx_query}"); | ||||
| 
 | ||||
|         let mut query_resp = db | ||||
|             .query(tx_query) | ||||
|             .bind(("new_vm_req_input", RecordId::from((NEW_VM_REQ, id)))) | ||||
|             .bind(("error_input", error)) | ||||
|             .await?; | ||||
| 
 | ||||
|         let query_err = query_resp.take_errors(); | ||||
|         if !query_err.is_empty() { | ||||
|             log::trace!("errors in submit_new_vm_err: {query_err:?}"); | ||||
|             let tx_fail_err_str = | ||||
|                 String::from("The query was not executed due to a failed transaction"); | ||||
| 
 | ||||
|             if query_err.contains_key(&4) && query_err[&4].to_string() != tx_fail_err_str { | ||||
|                 log::error!("vm req not exist: {}", query_err[&4]); | ||||
|                 return Err(Error::ContractNotFound); | ||||
|             } else { | ||||
|                 log::error!("Unknown error in submit_new_vm_err: {query_err:?}"); | ||||
|                 return Err(Error::Unknown("submit_new_vm_req".to_string())); | ||||
|             } | ||||
|         #[derive(Serialize)] | ||||
|         struct NewVmError { | ||||
|             error: String, | ||||
|         } | ||||
| 
 | ||||
|         let _: Option<Self> = db.update((NEW_VM_REQ, id)).merge(NewVmError { error }).await?; | ||||
|         // TODO: IMP refund tmp_locked
 | ||||
|         Ok(()) | ||||
|     } | ||||
| 
 | ||||
|     pub async fn submit(self, db: &Surreal<Client>) -> Result<(), Error> { | ||||
|         let locked_nano = self.locked_nano; | ||||
|         let tx_query = format!(" | ||||
|         let account = self.admin.key().to_string(); | ||||
|         let vm_id = self.id.key().to_string(); | ||||
|         let vm_node = self.vm_node.key().to_string(); | ||||
|         // TODO: check for possible injection and maybe use .bind()
 | ||||
|         let query = format!( | ||||
|             " | ||||
|             BEGIN TRANSACTION; | ||||
| 
 | ||||
|                 LET $account = $account_input; | ||||
|                 LET $new_vm_req = $new_vm_req_input; | ||||
|                 LET $vm_node = $vm_node_input; | ||||
|                 LET $hostname = $hostname_input; | ||||
|                 LET $dtrfs_url = $dtrfs_url_input; | ||||
|                 LET $dtrfs_sha = $dtrfs_sha_input; | ||||
|                 LET $kernel_url = $kernel_url_input; | ||||
|                 LET $kernel_sha = $kernel_sha_input; | ||||
| 
 | ||||
|                 UPDATE $account SET balance -= {locked_nano}; | ||||
|                 IF $account.balance < 0 {{ | ||||
|                     THROW 'Insufficient funds.' | ||||
|                 }}; | ||||
|                 UPDATE $account SET tmp_locked += {locked_nano}; | ||||
|                 RELATE | ||||
|                     $account | ||||
|                     ->$new_vm_req | ||||
|                     ->$vm_node | ||||
|                 CONTENT {{ | ||||
|                     created_at: time::now(), hostname: $hostname, vcpus: {}, memory_mb: {}, disk_size_gb: {}, | ||||
|                     extra_ports: {:?}, public_ipv4: {}, public_ipv6: {}, | ||||
|                     dtrfs_url: $dtrfs_url, dtrfs_sha: $dtrfs_sha, kernel_url: $kernel_url, kernel_sha: $kernel_sha, | ||||
|                     price_per_unit: {}, locked_nano: {locked_nano}, error: '' | ||||
|                 }}; | ||||
|             
 | ||||
|             COMMIT TRANSACTION;",
 | ||||
|             UPDATE account:{account} SET balance -= {locked_nano}; | ||||
|             IF account:{account}.balance < 0 {{ | ||||
|                 THROW 'Insufficient funds.' | ||||
|             }}; | ||||
|             UPDATE account:{account} SET tmp_locked += {locked_nano}; | ||||
|             RELATE | ||||
|                 account:{account} | ||||
|                 ->new_vm_req:{vm_id} | ||||
|                 ->vm_node:{vm_node} | ||||
|             CONTENT {{ | ||||
|                 created_at: time::now(), hostname: '{}', vcpus: {}, memory_mb: {}, disk_size_gb: {}, | ||||
|                 extra_ports: {:?}, public_ipv4: {:?}, public_ipv6: {:?}, | ||||
|                 dtrfs_url: '{}', dtrfs_sha: '{}', kernel_url: '{}', kernel_sha: '{}', | ||||
|                 price_per_unit: {}, locked_nano: {locked_nano}, error: '' | ||||
|             }}; | ||||
|             COMMIT TRANSACTION; | ||||
|         ",
 | ||||
|             self.hostname, | ||||
|             self.vcpus, | ||||
|             self.memory_mb, | ||||
|             self.disk_size_gb, | ||||
|             self.extra_ports, | ||||
|             self.public_ipv4, | ||||
|             self.public_ipv6, | ||||
|             self.dtrfs_url, | ||||
|             self.dtrfs_sha, | ||||
|             self.kernel_url, | ||||
|             self.kernel_sha, | ||||
|             self.price_per_unit | ||||
|         ); | ||||
|         //let _: Vec<Self> = db.insert(NEW_VM_REQ).relation(self).await?;
 | ||||
|         let mut query_resp = db.query(query).await?; | ||||
|         let resp_err = query_resp.take_errors(); | ||||
| 
 | ||||
|         log::trace!("submit_new_vm_req query: {tx_query}"); | ||||
| 
 | ||||
|         let mut query_resp = db | ||||
|             .query(tx_query) | ||||
|             .bind(("account_input", self.admin)) | ||||
|             .bind(("new_vm_req_input", self.id)) | ||||
|             .bind(("vm_node_input", self.vm_node)) | ||||
|             .bind(("hostname_input", self.hostname)) | ||||
|             .bind(("dtrfs_url_input", self.dtrfs_url)) | ||||
|             .bind(("dtrfs_sha_input", self.dtrfs_sha)) | ||||
|             .bind(("kernel_url_input", self.kernel_url)) | ||||
|             .bind(("kernel_sha_input", self.kernel_sha)) | ||||
|             .await?; | ||||
| 
 | ||||
|         let query_err = query_resp.take_errors(); | ||||
|         if !query_err.is_empty() { | ||||
|             log::trace!("errors in submit_new_vm_req: {query_err:?}"); | ||||
|             let tx_fail_err_str = | ||||
|                 String::from("The query was not executed due to a failed transaction"); | ||||
| 
 | ||||
|             if query_err.contains_key(&9) && query_err[&9].to_string() != tx_fail_err_str { | ||||
|                 log::error!("Transaction error: {}", query_err[&9]); | ||||
|                 return Err(Error::InsufficientFunds); | ||||
|             } else { | ||||
|                 log::error!("Unknown error in submit_new_vm_req: {query_err:?}"); | ||||
|                 return Err(Error::Unknown("submit_new_vm_req".to_string())); | ||||
|             } | ||||
|         if let Some(surrealdb::Error::Api(surrealdb::error::Api::Query(tx_query_error))) = | ||||
|             resp_err.get(&1) | ||||
|         { | ||||
|             log::error!("Transaction error: {tx_query_error}"); | ||||
|             return Err(Error::InsufficientFunds); | ||||
|         } | ||||
| 
 | ||||
|         Ok(()) | ||||
| @ -329,18 +272,14 @@ impl WrappedMeasurement { | ||||
|             UPDATE_VM_REQ => UPDATE_VM_REQ, | ||||
|             _ => NEW_VM_REQ, | ||||
|         }; | ||||
| 
 | ||||
|         let mut args_stream = db | ||||
|         let mut resp = db | ||||
|             .query(format!("live select error from {table} where id = {NEW_VM_REQ}:{vm_id};")) | ||||
|             .query(format!( | ||||
|                 "live select * from measurement_args where id = measurement_args:{vm_id};" | ||||
|             )) | ||||
|             .await? | ||||
|             .stream::<Notification<vm_proto::MeasurementArgs>>(0)?; | ||||
| 
 | ||||
|         let mut error_stream = db | ||||
|             .query(format!("live select error from {table} where id = {NEW_VM_REQ}:{vm_id};")) | ||||
|             .await? | ||||
|             .stream::<Notification<ErrorFromTable>>(0)?; | ||||
|             .await?; | ||||
|         let mut error_stream = resp.stream::<Notification<ErrorFromTable>>(0)?; | ||||
|         let mut args_stream = resp.stream::<Notification<vm_proto::MeasurementArgs>>(1)?; | ||||
| 
 | ||||
|         let args: Option<vm_proto::MeasurementArgs> = | ||||
|             db.delete(("measurement_args", vm_id)).await?; | ||||
|  | ||||
| @ -46,7 +46,7 @@ pub async fn create_new_vm( | ||||
|     assert!(new_vm_resp.uuid.len() == 40); | ||||
| 
 | ||||
|     // wait for update db
 | ||||
|     tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; | ||||
|     tokio::time::sleep(tokio::time::Duration::from_millis(700)).await; | ||||
| 
 | ||||
|     let vm_req_db: Option<db::NewVmReq> = db.select((NEW_VM_REQ, new_vm_resp.uuid.clone())).await?; | ||||
| 
 | ||||
|  | ||||
| @ -8,10 +8,7 @@ use tokio::sync::mpsc; | ||||
| use tokio_stream::wrappers::ReceiverStream; | ||||
| use tonic::transport::Channel; | ||||
| 
 | ||||
| pub async fn mock_vm_daemon( | ||||
|     brain_channel: &Channel, | ||||
|     daemon_error: Option<String>, | ||||
| ) -> Result<String> { | ||||
| pub async fn mock_vm_daemon(brain_channel: &Channel) -> Result<String> { | ||||
|     let mut daemon_client = BrainVmDaemonClient::new(brain_channel.clone()); | ||||
|     let daemon_key = Key::new(); | ||||
| 
 | ||||
| @ -29,7 +26,7 @@ pub async fn mock_vm_daemon( | ||||
|         rx, | ||||
|     )); | ||||
| 
 | ||||
|     tokio::spawn(daemon_engine(daemon_msg_tx.clone(), brain_msg_rx, daemon_error)); | ||||
|     tokio::spawn(daemon_engine(daemon_msg_tx.clone(), brain_msg_rx)); | ||||
| 
 | ||||
|     tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; | ||||
| 
 | ||||
| @ -105,7 +102,6 @@ pub async fn daemon_msg_sender( | ||||
| pub async fn daemon_engine( | ||||
|     tx: mpsc::Sender<vm_proto::VmDaemonMessage>, | ||||
|     mut rx: mpsc::Receiver<vm_proto::BrainVmMessage>, | ||||
|     new_vm_err: Option<String>, | ||||
| ) -> Result<()> { | ||||
|     log::info!("daemon engine vm_daemon"); | ||||
|     while let Some(brain_msg) = rx.recv().await { | ||||
| @ -124,7 +120,7 @@ pub async fn daemon_engine( | ||||
|                 let new_vm_resp = vm_proto::NewVmResp { | ||||
|                     uuid: new_vm_req.uuid.clone(), | ||||
|                     args, | ||||
|                     error: new_vm_err.clone().unwrap_or_default(), | ||||
|                     error: String::new(), | ||||
|                 }; | ||||
| 
 | ||||
|                 let res_data = vm_proto::VmDaemonMessage { | ||||
|  | ||||
| @ -116,7 +116,7 @@ async fn test_report_node() { | ||||
|     let db = prepare_test_db().await.unwrap(); | ||||
| 
 | ||||
|     let brain_channel = run_service_for_stream().await.unwrap(); | ||||
|     let daemon_key = mock_vm_daemon(&brain_channel, None).await.unwrap(); | ||||
|     let daemon_key = mock_vm_daemon(&brain_channel).await.unwrap(); | ||||
| 
 | ||||
|     let key = Key::new(); | ||||
| 
 | ||||
|  | ||||
| @ -2,30 +2,24 @@ use common::prepare_test_env::{prepare_test_db, run_service_for_stream}; | ||||
| use common::test_utils::Key; | ||||
| use common::vm_cli_utils::{airdrop, create_new_vm, user_list_vm_contracts}; | ||||
| use common::vm_daemon_utils::{mock_vm_daemon, register_vm_node}; | ||||
| use detee_shared::vm_proto; | ||||
| use detee_shared::vm_proto::brain_vm_cli_client::BrainVmCliClient; | ||||
| use detee_shared::vm_proto::brain_vm_daemon_client::BrainVmDaemonClient; | ||||
| use detee_shared::vm_proto::{ExtendVmReq, ListVmContractsReq, NewVmReq}; | ||||
| use futures::StreamExt; | ||||
| use std::vec; | ||||
| use surreal_brain::constants::{ACCOUNT, ACTIVE_VM, NEW_VM_REQ, TOKEN_DECIMAL}; | ||||
| use surreal_brain::constants::{ACCOUNT, ACTIVE_VM}; | ||||
| use surreal_brain::db::prelude as db; | ||||
| use surreal_brain::db::vm::ActiveVm; | ||||
| 
 | ||||
| mod common; | ||||
| 
 | ||||
| #[tokio::test] | ||||
| async fn test_vm_creation() { | ||||
|     let db = prepare_test_db().await.unwrap(); | ||||
|     /* | ||||
|     env_logger::builder() | ||||
|         .filter_level(log::LevelFilter::Debug) | ||||
|         .filter_module("tungstenite", log::LevelFilter::Debug) | ||||
|         .filter_module("tokio_tungstenite", log::LevelFilter::Debug) | ||||
|         .init(); | ||||
|     */ | ||||
|     // env_logger::builder().filter_level(log::LevelFilter::Error).init();
 | ||||
| 
 | ||||
|     let brain_channel = run_service_for_stream().await.unwrap(); | ||||
|     let daemon_key = mock_vm_daemon(&brain_channel, None).await.unwrap(); | ||||
|     let daemon_key = mock_vm_daemon(&brain_channel).await.unwrap(); | ||||
| 
 | ||||
|     let key = Key::new(); | ||||
| 
 | ||||
| @ -35,60 +29,11 @@ async fn test_vm_creation() { | ||||
|     let grpc_error_message = new_vm_resp.err().unwrap().to_string(); | ||||
|     assert!(grpc_error_message.contains("Insufficient funds")); | ||||
| 
 | ||||
|     let airdrop_amount = 10; | ||||
|     airdrop(&brain_channel, &key.pubkey, airdrop_amount).await.unwrap(); | ||||
|     airdrop(&brain_channel, &key.pubkey, 10).await.unwrap(); | ||||
| 
 | ||||
|     let new_vm_id = create_new_vm(&db, &key, &daemon_key, &brain_channel).await.unwrap(); | ||||
|     let active_vm: Option<db::ActiveVm> = db.select((ACTIVE_VM, new_vm_id)).await.unwrap(); | ||||
|     let active_vm: Option<ActiveVm> = db.select((ACTIVE_VM, new_vm_id)).await.unwrap(); | ||||
|     assert!(active_vm.is_some()); | ||||
| 
 | ||||
|     let daemon_key_02 = | ||||
|         mock_vm_daemon(&brain_channel, Some("something went wrong 04".to_string())).await.unwrap(); | ||||
| 
 | ||||
|     let locking_nano = 100; | ||||
|     let mut new_vm_req = vm_proto::NewVmReq { | ||||
|         admin_pubkey: key.pubkey.clone(), | ||||
|         node_pubkey: daemon_key_02.clone(), | ||||
|         price_per_unit: 1200, | ||||
|         extra_ports: vec![8080, 8081], | ||||
|         locked_nano: locking_nano, | ||||
|         ..Default::default() | ||||
|     }; | ||||
| 
 | ||||
|     let mut client_vm_cli = BrainVmCliClient::new(brain_channel.clone()); | ||||
|     let new_vm_resp = client_vm_cli | ||||
|         .new_vm(key.sign_request(new_vm_req.clone()).unwrap()) | ||||
|         .await | ||||
|         .unwrap() | ||||
|         .into_inner(); | ||||
| 
 | ||||
|     assert!(!new_vm_resp.error.is_empty()); | ||||
| 
 | ||||
|     let vm_req_db: Option<db::NewVmReq> = | ||||
|         db.select((NEW_VM_REQ, new_vm_resp.uuid.clone())).await.unwrap(); | ||||
| 
 | ||||
|     if let Some(new_vm_req) = vm_req_db { | ||||
|         assert!(!new_vm_req.error.is_empty()); | ||||
|         assert_eq!(new_vm_resp.error, new_vm_req.error); | ||||
|     } | ||||
| 
 | ||||
|     let acc_db: db::Account = db.select((ACCOUNT, key.pubkey.clone())).await.unwrap().unwrap(); | ||||
|     assert_eq!(acc_db.balance, airdrop_amount * TOKEN_DECIMAL); | ||||
|     assert_eq!(acc_db.tmp_locked, 0); | ||||
| 
 | ||||
|     new_vm_req.node_pubkey = daemon_key; | ||||
| 
 | ||||
|     let new_vm_resp = | ||||
|         client_vm_cli.new_vm(key.sign_request(new_vm_req).unwrap()).await.unwrap().into_inner(); | ||||
|     assert!(new_vm_resp.error.is_empty()); | ||||
| 
 | ||||
|     tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; | ||||
|     let active_vm: Option<db::ActiveVm> = db.select((ACTIVE_VM, new_vm_resp.uuid)).await.unwrap(); | ||||
|     assert!(active_vm.is_some()); | ||||
| 
 | ||||
|     let acc_db: db::Account = db.select((ACCOUNT, key.pubkey.clone())).await.unwrap().unwrap(); | ||||
|     assert_eq!(acc_db.balance, airdrop_amount * TOKEN_DECIMAL - locking_nano); | ||||
|     assert_eq!(acc_db.tmp_locked, 0); | ||||
| } | ||||
| 
 | ||||
| #[tokio::test] | ||||
|  | ||||
| @ -29,7 +29,7 @@ async fn test_brain_message() { | ||||
|     prepare_test_db().await.unwrap(); | ||||
| 
 | ||||
|     let brain_channel = run_service_for_stream().await.unwrap(); | ||||
|     let daemon_key = mock_vm_daemon(&brain_channel, None).await.unwrap(); | ||||
|     let daemon_key = mock_vm_daemon(&brain_channel).await.unwrap(); | ||||
|     let mut cli_client = BrainVmCliClient::new(brain_channel.clone()); | ||||
| 
 | ||||
|     let cli_key = Key::new(); | ||||
|  | ||||
		Loading…
	
		Reference in New Issue
	
	Block a user