Преглед на файлове

solana: Add POC for watching pubkeys for transactions.

parazyd преди 4 години
родител
ревизия
5e3a5c9870
променени са 5 файла, в които са добавени 699 реда и са изтрити 55 реда
  1. 554 50
      Cargo.lock
  2. 18 4
      Cargo.toml
  3. 7 0
      src/bin/drk.rs
  4. 117 0
      src/bin/solana-poc.rs
  5. 3 1
      src/rpc/jsonrpc.rs

Файловите разлики са ограничени, защото са твърде много
+ 554 - 50
Cargo.lock


+ 18 - 4
Cargo.toml

@@ -69,6 +69,7 @@ async-h1 = "2.3.0"
 async-native-tls = "0.3.3"
 async-native-tls = "0.3.3"
 toml = "0.5.8"
 toml = "0.5.8"
 
 
+
 # serde
 # serde
 serde = { version = "1.0.126", features = ["derive"]}
 serde = { version = "1.0.126", features = ["derive"]}
 
 
@@ -77,17 +78,30 @@ zeromq = { git="https://github.com/zeromq/zmq.rs", default-features = false, fea
 bytes = "1.0.1"
 bytes = "1.0.1"
 
 
 # wallet deps
 # wallet deps
-rocksdb = "0.16.0"
+rocksdb = {version = "0.16.0", default-features = false, features = ["lz4"]}
 dirs = "3.0.2"
 dirs = "3.0.2"
 
 
-bitcoin = "0.27.0"
-secp256k1 = "0.20.3"
-electrum-client = "0.8.0"
+## Cashier Solana Dependencies
+solana-sdk = {version = "1.7.11", optional = true}
+solana-client = {version = "1.7.11", optional = true}
+tokio-tungstenite = {version = "0.15.0", optional = true}
+tokio = {version = "1.11.0", features = ["full"], optional = true}
+
+## Cashier Bitcoin Dependencies
+# TODO: Make optional, cashier should have "features" to enable different chains
+# See https://doc.rust-lang.org/cargo/reference/features.html
+bitcoin = {version = "0.27.0", optional = false }
+secp256k1 = {version = "0.20.3", optional = false }
+electrum-client = {version = "0.8.0", optional = false }
 
 
 [dependencies.rusqlite]
 [dependencies.rusqlite]
 version = "0.25.1"
 version = "0.25.1"
 features = ["bundled", "sqlcipher"]
 features = ["bundled", "sqlcipher"]
 
 
+[features]
+sol = ["solana-sdk", "solana-client", "tokio-tungstenite", "tokio"]
+#btc = ["bitcoin", "secp256k1", "electrum-client"]
+
 [[bin]]
 [[bin]]
 name = "lisp"
 name = "lisp"
 path = "lisp/lisp.rs"
 path = "lisp/lisp.rs"

+ 7 - 0
src/bin/drk.rs

@@ -40,6 +40,13 @@ impl Drk {
                 debug!(target: "DRK", "<-- {:?}", e);
                 debug!(target: "DRK", "<-- {:?}", e);
                 return Err(Error::JsonRpcError(e.error.message.to_string()));
                 return Err(Error::JsonRpcError(e.error.message.to_string()));
             }
             }
+
+            JsonResult::Notif(n) => {
+                debug!(target: "DRK", "<-- {:?}", n);
+                return Err(Error::JsonRpcError(
+                    "Unexpected reply from server".to_string(),
+                ));
+            }
         };
         };
     }
     }
 
 

+ 117 - 0
src/bin/solana-poc.rs

@@ -0,0 +1,117 @@
+// $ cargo run --features sol --bin solana-poc
+// $ solana transfer 10 $pubkey
+use futures::{SinkExt, StreamExt};
+use rand::rngs::OsRng;
+use serde::Serialize;
+use serde_json::{json, Value};
+use solana_client::{blockhash_query::BlockhashQuery, rpc_client::RpcClient};
+use solana_sdk::{
+    native_token::lamports_to_sol, pubkey::Pubkey, signature::Signer, signer::keypair::Keypair,
+    system_instruction, transaction::Transaction,
+};
+use std::sync::{Arc, Mutex};
+use tokio_tungstenite::{connect_async, tungstenite::protocol::Message};
+
+use drk::rpc::{jsonrpc, jsonrpc::JsonResult};
+
+//const RPC_SERVER: &'static str = "https://api.mainnet-beta.solana.com";
+//const WSS_SERVER: &'static str = "wss://api.mainnet-beta.solana.com";
+//const RPC_SERVER: &'static str = "https://api.devnet.solana.com";
+//const WSS_SERVER: &'static str = "wss://api.devnet.solana.com";
+const RPC_SERVER: &'static str = "http://localhost:8899";
+const WSS_SERVER: &'static str = "ws://localhost:8900";
+
+// https://docs.solana.com/developing/clients/jsonrpc-api#accountsubscribe
+#[derive(Serialize)]
+struct SubscribeParams {
+    encoding: Value,
+    commitment: Value,
+}
+
+// Example function to show how to transfer `amount` lamports
+fn transfer_lamports(from: &Keypair, to: &Pubkey, amount: u64) {
+    let rpc = RpcClient::new(RPC_SERVER.to_string());
+    let instruction = system_instruction::transfer(&from.pubkey(), to, amount);
+
+    let mut tx = Transaction::new_with_payer(&[instruction], Some(&from.pubkey()));
+    let bhq = BlockhashQuery::default();
+    match bhq.get_blockhash_and_fee_calculator(&rpc, rpc.commitment()) {
+        Err(_) => panic!("Couldn't connect to RPC"),
+        Ok(v) => tx.sign(&[from], v.0),
+    }
+
+    let _signature = rpc.send_and_confirm_transaction(&tx);
+}
+
+#[tokio::main]
+async fn main() -> Result<(), &'static str> {
+    let keypair = Keypair::generate(&mut OsRng);
+    println!("Pubkey: {:?}", keypair.pubkey());
+
+    let rpc = RpcClient::new(RPC_SERVER.to_string());
+    let balance = rpc.get_balance(&keypair.pubkey()).unwrap();
+    let account_bal = Arc::new(Mutex::new(balance));
+
+    // Parameters for subscription to events related to `pubkey`.
+    let sub_params = SubscribeParams {
+        encoding: json!("jsonParsed"),
+        // XXX: Use "finalized" for 100% certainty.
+        commitment: json!("confirmed"),
+    };
+
+    let sub_msg = jsonrpc::request(
+        json!("accountSubscribe"),
+        json!([json!(keypair.pubkey().to_string()), json!(sub_params)]),
+    );
+
+    // WebSocket handshake/connect
+    let (ws_stream, _) = connect_async(WSS_SERVER)
+        .await
+        .expect("Failed to connect to WebSocket server");
+
+    let (mut write, read) = ws_stream.split();
+
+    // Send the subscription request
+    write
+        .send(Message::Text(serde_json::to_string(&sub_msg).unwrap()))
+        .await
+        .unwrap();
+    println!("Subscribed to events for {:?}", keypair.pubkey());
+
+    // Subscription ID so we can map our notifications to our pubkey
+    // when we do multiple subscriptions and also do `accountUnsubscribe`.
+    let sub_id = Arc::new(Mutex::new(0));
+
+    let read_future = read.for_each(|message| async {
+        let data = message.unwrap().into_text().unwrap();
+        let v: JsonResult = serde_json::from_str(&data).unwrap();
+        match v {
+            JsonResult::Resp(r) => {
+                println!(
+                    "Successfully subscribed with ID: {:?}",
+                    r.result.as_i64().unwrap()
+                );
+                *sub_id.lock().unwrap() = r.result.as_i64().unwrap();
+            }
+
+            JsonResult::Err(e) => {
+                println!("Error on subscription: {:?}", e.error.message.to_string());
+            }
+
+            JsonResult::Notif(n) => {
+                println!("Got WebSocket notification: {:?}", n);
+                println!(
+                    "Old balance: {:?} SOL",
+                    lamports_to_sol(*account_bal.lock().unwrap())
+                );
+                let new_bal = n.params["result"]["value"]["lamports"].as_u64().unwrap();
+                *account_bal.lock().unwrap() = new_bal;
+                println!("New balance: {:?} SOL", lamports_to_sol(new_bal));
+            }
+        }
+    });
+
+    read_future.await;
+
+    Ok(())
+}

+ 3 - 1
src/rpc/jsonrpc.rs

@@ -1,13 +1,15 @@
+use std::str;
+
 use rand::Rng;
 use rand::Rng;
 use serde::{Deserialize, Serialize};
 use serde::{Deserialize, Serialize};
 use serde_json::{json, Value};
 use serde_json::{json, Value};
-use std::str;
 
 
 #[derive(Serialize, Deserialize, Debug)]
 #[derive(Serialize, Deserialize, Debug)]
 #[serde(untagged)]
 #[serde(untagged)]
 pub enum JsonResult {
 pub enum JsonResult {
     Resp(JsonResponse),
     Resp(JsonResponse),
     Err(JsonError),
     Err(JsonError),
+    Notif(JsonNotification),
 }
 }
 
 
 #[derive(Serialize, Deserialize, Debug)]
 #[derive(Serialize, Deserialize, Debug)]

Някои файлове не бяха показани, защото твърде много файлове са промени