summaryrefslogtreecommitdiff
path: root/plugins/warp/src/container
diff options
context:
space:
mode:
authorMason Reed <mason@vector35.com>2025-07-13 14:39:14 -0400
committerMason Reed <mason@vector35.com>2025-07-15 12:34:43 -0400
commit4a49ba509bdc0b4ffa650fc8461738c2085161c7 (patch)
tree0c55d84afde06cfe1416f8f199df985b41a0f1ff /plugins/warp/src/container
parent9c80b724bc28eca0b1b4fad1b955d4fd2ba48ab5 (diff)
[WARP] Add network container
This is going to be disabled by default on this upcoming stable, however users may enable it once we deploy the public server. The data from the server is done through `Container::fetch_functions` independent of the nonblocking function lookup functions. The sidebar has been updated to drive fetching so that when users navigate to a new function the fetcher will kick off. This fetcher operates on a separate thread, in the event of a user navigating to many functions before the current fetch has completed they all will be batched together in a single fetch. Networked container currently is limited to just function prototypes, other type information separate from the function object will be omitted.
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
+ }
+ }
+ }
+}