summaryrefslogtreecommitdiff
path: root/plugins/warp/src/container
diff options
context:
space:
mode:
Diffstat (limited to 'plugins/warp/src/container')
-rw-r--r--plugins/warp/src/container/network.rs332
-rw-r--r--plugins/warp/src/container/network/client.rs228
2 files changed, 549 insertions, 11 deletions
diff --git a/plugins/warp/src/container/network.rs b/plugins/warp/src/container/network.rs
index ffbe6108..2a2d7c65 100644
--- a/plugins/warp/src/container/network.rs
+++ b/plugins/warp/src/container/network.rs
@@ -1,13 +1,323 @@
-pub struct NetworkContainer {}
+use crate::container::disk::DiskContainer;
+use crate::container::{Container, ContainerError, ContainerResult, SourceId, SourcePath};
+use std::collections::HashMap;
+use std::fmt::{Debug, Display, Formatter};
+use warp::chunk::{Chunk, ChunkKind, CompressionType};
+use warp::r#type::guid::TypeGUID;
+use warp::r#type::{ComputedType, Type};
+use warp::signature::chunk::SignatureChunk;
+use warp::signature::function::{Function, FunctionGUID};
+use warp::target::Target;
+use warp::{WarpFile, WarpFileHeader};
-// TODO: The container is populated as the user is navigating a binary
-// TODO: We need to have a few helper functions here to post and pull
-// TODO: Then in the interface we operate off the network cache
-// TODO: The network cache could just be a disk container? Or disk container sources?
-// TODO: We should also store the cache on the filesystem for a certain time, will need to timestamp
-// TODO: When we commit we need to actually POST i believe.
-// TODO: There needs to be a setting that adjusts the sweep size of functions at the cursor.
-// TODO: Probably need a callback or something to tell the network containers to refresh from the network.
-// TODO: The network container should never instantiate itself, unless its gurenteed to not have any data in it?
+pub mod client;
-// TODO: Need to PUSH chunks and PULL chunks
+pub use client::NetworkClient;
+
+/// This is the id on the server for the [`Target`], we can get it via [`NetworkClient::query_target_id`].
+pub type NetworkTargetId = i32;
+
+pub struct NetworkContainer {
+ client: NetworkClient,
+ /// This is the store that the interface will write to; then we have special functions for pulling
+ /// and pushing to the network source.
+ cache: DiskContainer,
+ /// Populated when targets are queried.
+ known_targets: HashMap<Target, Option<NetworkTargetId>>,
+ /// Populated with function sources are queried.
+ known_function_sources: HashMap<FunctionGUID, Vec<SourceId>>,
+ /// Populated when user adds function, this is used for writing back to the server.
+ added_chunks: HashMap<SourceId, Vec<Chunk<'static>>>,
+}
+
+impl NetworkContainer {
+ pub fn new(client: NetworkClient) -> Self {
+ Self {
+ cache: DiskContainer::new("Network Container".to_string(), HashMap::new()),
+ client,
+ known_targets: HashMap::new(),
+ known_function_sources: HashMap::new(),
+ added_chunks: HashMap::new(),
+ }
+ }
+
+ /// Gets the network id for the `target`, this will be used in later function queries.
+ ///
+ /// **This is blocking**
+ ///
+ /// # Caching policy
+ ///
+ /// The [`NetworkTargetId`] is unique and immutable, so they will be persisted indefinitely.
+ pub fn get_target_id(&mut self, target: &Target) -> Option<NetworkTargetId> {
+ // It's highly probable we have previously queried the target, check that first.
+ if let Some(target_id) = self.known_targets.get(target) {
+ return target_id.clone();
+ }
+
+ let target_id = self.client.query_target_id(target);
+ // Keep the target id so the next lookup is free.
+ self.known_targets.insert(target.clone(), target_id);
+ target_id
+ }
+
+ /// Pulls sources for the set of unseen function guids.
+ ///
+ /// **This is blocking**
+ ///
+ /// # Caching policy
+ ///
+ /// When we get the source, we store the results indefinitely in the container; this is fine
+ /// for now as the requests for functions come at the request of some user interaction. Any guid
+ /// with no sources will still be cached.
+ pub fn get_unseen_functions_source(
+ &mut self,
+ target: Option<&Target>,
+ guids: &[FunctionGUID],
+ ) -> HashMap<SourceId, Vec<FunctionGUID>> {
+ let Some(target_id) = target.and_then(|t| self.get_target_id(t)) else {
+ log::debug!("Cannot query functions source without a target, skipping...");
+ return HashMap::new();
+ };
+
+ // Split guids into known and unknown
+ let (_known, unknown): (Vec<_>, Vec<_>) = guids
+ .into_iter()
+ .cloned()
+ .partition(|guid| self.known_function_sources.contains_key(guid));
+
+ let mut result: HashMap<SourceId, Vec<FunctionGUID>> = HashMap::new();
+ // Only query server for unknown guids if we have any.
+ if !unknown.is_empty() {
+ if let Some(queried_results) = self
+ .client
+ .query_functions_source(Some(target_id), &unknown)
+ {
+ // Cache the new results, this means we will not try and contact the server for that guids source.
+ // NOTE: Here we do not just simply list the queried results because we also
+ // want to cache function guids which have no source, this is important so that we never
+ // attempt to contact the server for that guid.
+ for guid in &unknown {
+ let sources = queried_results
+ .keys()
+ .filter(|source_id| queried_results[source_id].contains(guid))
+ .cloned()
+ .collect();
+ self.known_function_sources.insert(*guid, sources);
+ }
+
+ for (source_id, guids) in queried_results {
+ result.entry(source_id).or_default().extend(guids);
+ }
+ }
+ }
+
+ result
+ }
+
+ /// Pulls function metadata from the server and adds it into the container cache.
+ ///
+ /// **This is blocking**
+ ///
+ /// # Caching policy
+ ///
+ /// Every request we store the returned objects on disk, this means that users will first
+ /// query against the disk objects, then the server. This also means we need to cache functions f
+ /// or which we have not received any functions for, as otherwise we would keep trying to query it.
+ pub fn pull_functions(
+ &mut self,
+ target: &Target,
+ source: &SourceId,
+ functions: &[FunctionGUID],
+ ) {
+ let target_id = self.get_target_id(target);
+ if let Some(file) = self
+ .client
+ .query_functions(target_id, Some(*source), functions)
+ {
+ log::debug!("Got {} chunks from server", file.chunks.len());
+ for chunk in &file.chunks {
+ match &chunk.kind {
+ ChunkKind::Signature(sc) => {
+ let functions: Vec<_> = sc.functions().collect();
+ match self.cache.add_functions(target, source, &functions) {
+ Ok(_) => log::debug!(
+ "Added {} functions into cached source '{}'",
+ functions.len(),
+ source
+ ),
+ Err(err) => log::error!(
+ "Failed to add {} function into cached source '{}': {}",
+ functions.len(),
+ source,
+ err
+ ),
+ }
+ }
+ // TODO; Probably want to pull type in with this.
+ ChunkKind::Type(_) => {}
+ }
+ }
+ }
+ }
+
+ /// Push a file to the network source.
+ ///
+ /// **This is blocking**
+ pub fn push_file(&mut self, source_id: SourceId, file: &WarpFile) {
+ self.client.push_file(source_id, file);
+ }
+}
+
+impl Container for NetworkContainer {
+ fn sources(&self) -> ContainerResult<Vec<SourceId>> {
+ self.cache.sources()
+ }
+
+ fn add_source(&mut self, path: SourcePath) -> ContainerResult<SourceId> {
+ // TODO: How do we want to let users create new sources?
+ log::error!("NetworkContainer::add_source not allowed");
+ Err(ContainerError::CannotCreateSource(path))
+ }
+
+ fn commit_source(&mut self, source: &SourceId) -> ContainerResult<bool> {
+ let chunks = self
+ .added_chunks
+ .remove(source)
+ .ok_or(ContainerError::SourceNotFound(source.clone()))?;
+ let file = WarpFile::new(WarpFileHeader::new(), chunks);
+ self.push_file(*source, &file);
+ Ok(true)
+ }
+
+ fn is_source_writable(&self, source: &SourceId) -> ContainerResult<bool> {
+ // TODO: This is retrievable from /users/me/sources we will grab it when connecting.
+ log::error!("NetworkContainer::is_source_writable not allowed");
+ Err(ContainerError::SourceNotWritable(source.clone()))
+ }
+
+ fn is_source_uncommitted(&self, source: &SourceId) -> ContainerResult<bool> {
+ Ok(self.added_chunks.contains_key(source))
+ }
+
+ fn source_path(&self, source: &SourceId) -> ContainerResult<SourcePath> {
+ self.cache.source_path(source)
+ }
+
+ fn add_computed_types(
+ &mut self,
+ source: &SourceId,
+ types: &[ComputedType],
+ ) -> ContainerResult<()> {
+ self.cache.add_computed_types(source, types)
+ }
+
+ fn remove_types(&mut self, source: &SourceId, guids: &[TypeGUID]) -> ContainerResult<()> {
+ self.cache.remove_types(source, guids)
+ }
+
+ fn add_functions(
+ &mut self,
+ target: &Target,
+ source: &SourceId,
+ functions: &[Function],
+ ) -> ContainerResult<()> {
+ let signature_chunk = SignatureChunk::new(functions).ok_or(
+ ContainerError::CorruptedData("signature chunk failed to validate"),
+ )?;
+ let chunk = Chunk::new_with_target(
+ ChunkKind::Signature(signature_chunk),
+ CompressionType::None,
+ target.clone(),
+ );
+ self.added_chunks.entry(*source).or_default().push(chunk);
+ Ok(())
+ }
+
+ fn remove_functions(
+ &mut self,
+ target: &Target,
+ source: &SourceId,
+ functions: &[Function],
+ ) -> ContainerResult<()> {
+ // TODO: Wont persist, need to add remote removal.
+ self.cache.remove_functions(target, source, functions)
+ }
+
+ fn fetch_functions(
+ &mut self,
+ target: &Target,
+ functions: &[FunctionGUID],
+ ) -> ContainerResult<()> {
+ // NOTE: Blocking request to get the mapped function sources.
+ let mapped_unseen_functions = self.get_unseen_functions_source(Some(&target), functions);
+
+ // Actually get the function data for the unseen guids, we really only want to do this once per
+ // session, anymore, and this is annoying!
+ for (source, unseen_guids) in mapped_unseen_functions {
+ // NOTE: Blocking request to get the function data in the container cache.
+ self.pull_functions(&target, &source, &unseen_guids);
+ }
+
+ Ok(())
+ }
+
+ fn sources_with_type_guid(&self, guid: &TypeGUID) -> ContainerResult<Vec<SourceId>> {
+ self.cache.sources_with_type_guid(guid)
+ }
+
+ fn sources_with_type_guids(
+ &self,
+ guids: &[TypeGUID],
+ ) -> ContainerResult<HashMap<TypeGUID, Vec<SourceId>>> {
+ self.cache.sources_with_type_guids(guids)
+ }
+
+ fn type_guids_with_name(
+ &self,
+ source: &SourceId,
+ name: &str,
+ ) -> ContainerResult<Vec<TypeGUID>> {
+ self.cache.type_guids_with_name(source, name)
+ }
+
+ fn type_with_guid(&self, source: &SourceId, guid: &TypeGUID) -> ContainerResult<Option<Type>> {
+ self.cache.type_with_guid(source, guid)
+ }
+
+ fn sources_with_function_guid(
+ &self,
+ target: &Target,
+ guid: &FunctionGUID,
+ ) -> ContainerResult<Vec<SourceId>> {
+ self.cache.sources_with_function_guid(target, guid)
+ }
+
+ fn sources_with_function_guids(
+ &self,
+ target: &Target,
+ guids: &[FunctionGUID],
+ ) -> ContainerResult<HashMap<FunctionGUID, Vec<SourceId>>> {
+ self.cache.sources_with_function_guids(target, guids)
+ }
+
+ fn functions_with_guid(
+ &self,
+ target: &Target,
+ source: &SourceId,
+ guid: &FunctionGUID,
+ ) -> ContainerResult<Vec<Function>> {
+ self.cache.functions_with_guid(target, source, guid)
+ }
+}
+
+impl Debug for NetworkContainer {
+ fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
+ f.debug_struct("NetworkContainer").finish()
+ }
+}
+
+impl Display for NetworkContainer {
+ fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
+ f.debug_struct("NetworkContainer").finish()
+ }
+}
diff --git a/plugins/warp/src/container/network/client.rs b/plugins/warp/src/container/network/client.rs
new file mode 100644
index 00000000..f77f1118
--- /dev/null
+++ b/plugins/warp/src/container/network/client.rs
@@ -0,0 +1,228 @@
+use crate::container::network::NetworkTargetId;
+use crate::container::SourceId;
+use reqwest::blocking::Client;
+use reqwest::header::{HeaderMap, HeaderValue, AUTHORIZATION};
+use reqwest::StatusCode;
+use serde_json::json;
+use std::collections::HashMap;
+use std::str::FromStr;
+use warp::signature::function::FunctionGUID;
+use warp::target::Target;
+use warp::WarpFile;
+
+/// Responsible for sending and receiving data from the server.
+///
+/// NOTE: **All requests are blocking**.
+#[derive(Clone, Debug)]
+pub struct NetworkClient {
+ client: Client,
+ server_url: String,
+}
+
+impl NetworkClient {
+ pub fn new(
+ server_url: String,
+ server_token: Option<String>,
+ https_proxy: Option<String>,
+ ) -> reqwest::Result<Self> {
+ let version_info = binaryninja::version_info();
+ // TODO: IIRC we had a user agent format already for some other thing.
+ let client_agent = format!(
+ "Binary Ninja/{}.{}.{}",
+ version_info.major, version_info.minor, version_info.build
+ );
+ // TODO: This might want to be kept for the request header?
+ let mut headers = HeaderMap::new();
+ if let Some(token) = &server_token {
+ headers.insert(
+ AUTHORIZATION,
+ HeaderValue::from_str(&format!("Bearer {}", token)).unwrap(),
+ );
+ }
+ // TODO: Configurable timeout?
+ let mut client_builder = Client::builder()
+ .connect_timeout(std::time::Duration::from_secs(10))
+ .default_headers(headers)
+ .user_agent(client_agent);
+ if let Some(https_proxy) = https_proxy {
+ client_builder = client_builder.proxy(reqwest::Proxy::all(&https_proxy)?);
+ }
+ Ok(Self {
+ client: client_builder.build()?,
+ server_url,
+ })
+ }
+
+ /// Check to see the status of the server.
+ ///
+ /// This is useful if you want to fail early and prevent constructing a network container to a
+ /// server that is unresponsive.
+ ///
+ /// Route: `api/v1/status`
+ pub fn status(&self) -> reqwest::Result<StatusCode> {
+ let status_url = format!("{}/api/v1/status", self.server_url);
+ let resp = self.client.get(&status_url).send()?;
+ Ok(resp.status())
+ }
+
+ /// Query the [`NetworkTargetId`] for the given [`Target`].
+ ///
+ /// NOTE: **THIS IS BLOCKING**
+ ///
+ /// Route: `api/v1/targets/query` (TODO: Comment about the query)
+ pub fn query_target_id(&self, target: &Target) -> Option<NetworkTargetId> {
+ let query_target_url = format!("{}/api/v1/targets/query", self.server_url);
+
+ let mut query = HashMap::new();
+ if let Some(platform) = &target.platform {
+ query.insert("platform", platform);
+ }
+ if let Some(architecture) = &target.architecture {
+ query.insert("architecture", architecture);
+ }
+
+ // NOTE: This is blocking.
+ let target_id: NetworkTargetId = self
+ .client
+ .get(query_target_url)
+ .query(&query)
+ .send()
+ .ok()?
+ .json::<NetworkTargetId>()
+ .ok()?;
+
+ Some(target_id)
+ }
+
+ fn query_functions_body(
+ target: Option<NetworkTargetId>,
+ source: Option<SourceId>,
+ guids: &[FunctionGUID],
+ ) -> serde_json::Value {
+ let guids_str: Vec<String> = guids.iter().map(|g| g.to_string()).collect();
+ // TODO: The limit here needs to be somewhat flexible. But 1000 will do for now.
+ let mut body = json!({
+ "format": "flatbuffer",
+ "guids": guids_str,
+ "limit": 1000
+ });
+ if let Some(target_id) = target {
+ body["target_id"] = json!(target_id);
+ }
+ if let Some(source_id) = source {
+ body["source_id"] = json!(source_id.to_string());
+ }
+ body
+ }
+
+ /// Query the functions, returning the warp file response containing the entries.
+ ///
+ /// NOTE: **THIS IS BLOCKING**
+ ///
+ /// Route: `api/v1/functions/query` (TODO: Comment about the query)
+ pub fn query_functions(
+ &self,
+ target: Option<NetworkTargetId>,
+ source: Option<SourceId>,
+ guids: &[FunctionGUID],
+ ) -> Option<WarpFile<'static>> {
+ let query_functions_url = format!("{}/api/v1/functions/query", self.server_url);
+ let payload = Self::query_functions_body(target, source, guids);
+
+ // Make the POST request
+ let response = self
+ .client
+ .post(&query_functions_url)
+ .json(&payload)
+ .send()
+ .ok()?;
+ if !response.status().is_success() {
+ log::error!("Failed to query functions: {}", response.status());
+ return None;
+ }
+
+ // Get response bytes and convert to WarpFile
+ let bytes = response.bytes().ok()?;
+ WarpFile::from_owned_bytes(bytes.to_vec())
+ }
+
+ /// Query the functions, returning the sources and the corresponding function guids.
+ ///
+ /// NOTE: **THIS IS BLOCKING**
+ ///
+ /// Route: `api/v1/functions/query/source` (TODO: Comment about the query)
+ pub fn query_functions_source(
+ &self,
+ target: Option<NetworkTargetId>,
+ guids: &[FunctionGUID],
+ ) -> Option<HashMap<SourceId, Vec<FunctionGUID>>> {
+ let query_functions_source_url =
+ format!("{}/api/v1/functions/query/source", self.server_url);
+ let payload = Self::query_functions_body(target, None, guids);
+
+ // Make the POST request
+ let response = self
+ .client
+ .post(&query_functions_source_url)
+ .json(&payload)
+ .send()
+ .ok()?;
+ if !response.status().is_success() {
+ log::error!("Failed to query functions source: {}", response.status());
+ return None;
+ }
+
+ // Mapping of source id to function guids
+ let json_response: HashMap<String, Vec<String>> = response.json().ok()?;
+ let mapped_function_guids = json_response
+ .into_iter()
+ .filter_map(|(source_str, guid_strs)| {
+ let source_id = SourceId::from_str(&source_str).ok()?;
+ let guids = guid_strs
+ .into_iter()
+ .filter_map(|guid_str| FunctionGUID::from_str(&guid_str).ok())
+ .collect();
+ Some((source_id, guids))
+ })
+ .collect();
+
+ Some(mapped_function_guids)
+ }
+
+ /// Pushes the file to the remote source.
+ ///
+ /// NOTE: **THIS IS BLOCKING**
+ ///
+ /// Route: `api/v1/files/{source}`
+ pub fn push_file(&self, source_id: SourceId, file: &WarpFile) -> bool {
+ let push_file_url = format!("{}/api/v1/files/{}", self.server_url, source_id.to_string());
+
+ // Convert WarpFile to bytes
+ let file_bytes = file.to_bytes();
+
+ // Create the form part with the file
+ let form = reqwest::blocking::multipart::Form::new().part(
+ "file",
+ reqwest::blocking::multipart::Part::bytes(file_bytes)
+ .file_name("data.warp")
+ .mime_str("application/octet-stream")
+ .unwrap(),
+ );
+
+ // Send the request
+ match self.client.post(&push_file_url).multipart(form).send() {
+ Ok(response) => {
+ if response.status().is_success() {
+ true
+ } else {
+ log::error!("Failed to push file: {}", response.status());
+ false
+ }
+ }
+ Err(e) => {
+ log::error!("Failed to send push request: {}", e);
+ false
+ }
+ }
+ }
+}