hadean-std 0.1.6

Hadean stdlib. Requires Hadean Rust.
Documentation
//! A scalable distributed dataframe datastructure.

use std::{cmp,env,mem,ops,iter,marker,collections,hash,slice,ptr};
use hadean::{Sender,Receiver,Connection,Process,ChannelType,Channel,pid,spawn,ProcessSendable};
use list::{Leaf,LeafIndex};
use hashmap;

/// A scalable distributed dataframe datastructure.
/// `Vec<T>` is a dataframe datastructure that holds a collection of elements of type `T`.
///
/// It enables computation to be performed in parallel across many elements at the same time.
///
/// Here's an example showing addition of data to a `Vec`, and `map`, `filter`, and `reduce` being called upon it:
///
/// ```
/// use hadean_std::dataframe::Vec;
///
/// let mut data_frame: Vec<String> = Vec::new();
/// data_frame.push(String::from("http://google.com"));
/// data_frame.push(String::from("http://bbc.co.uk"));
/// data_frame.push(String::from("https://hadean.com"));
/// let mut data_frame: Vec<String> = data_frame.map(|mut val|{val.push_str("/xyz");val});
/// data_frame.filter(|val|val.len() > 20);
/// data_frame.reduce(0u64, |acc, val| acc + val.len() as u64);
/// ```
///
pub struct Vec<'a,T> where T: ProcessSendable {
	list: OwnedOrBorrowed<'a,Leaf<u8>>,
	start: LeafIndex,
	end: LeafIndex,
	phantom: marker::PhantomData<T>
}
enum OwnedOrBorrowed<'a,T> where T: 'a {
	Owned(Box<T>),
	Borrowed(&'a mut T)
}
impl<'a,T> OwnedOrBorrowed<'a,T> {
	fn list(&self) -> &T {
		match self {
			&OwnedOrBorrowed::Owned(ref list) => &*list,
			&OwnedOrBorrowed::Borrowed(ref list) => *list
		}
	}
	fn list_mut(&mut self) -> &mut T {
		match self {
			&mut OwnedOrBorrowed::Owned(ref mut list) => &mut *list,
			&mut OwnedOrBorrowed::Borrowed(ref mut list) => *list
		}
	}
}
impl<'a,T> Vec<'a,T> where T: ProcessSendable {
	/// Constructs a new, empty `Vec<T>`.
	pub fn new() -> Vec<'a,T> {
		let mut leaf: Box<Leaf<u8>> = box unsafe{mem::uninitialized()};
		let (start, end) = Leaf::init(&mut leaf);
		Vec{list:OwnedOrBorrowed::Owned(leaf),start:start,end:end,phantom:marker::PhantomData}
	}
	/// Constructs a new, empty `Vec<T>` in the given `List<u8>`. `list.len(&start, &end)` must equal `0`.
	pub fn with_list(list: &'a mut Leaf<u8>, start: LeafIndex, end: LeafIndex) -> Vec<'a,T> {
		assert!(start == end && list.len(&start, &end) == 0);
		Vec{list:OwnedOrBorrowed::Borrowed(list),start:start,end:end,phantom:marker::PhantomData}
	}
	/// Append `val` to the front of the dataframe.
	pub fn push_front(&mut self, val: T) {
		val.processsendable_write(&mut |buf,len| {
			let mut start = self.start.clone_right();
			let mut start2 = start.clone_right();
			self.list.list_mut().replace(&mut start, &mut start2, unsafe{slice::from_raw_parts(buf, len)}.iter().map(|&x|x));
		});
	}
	/// Append `val` to the dataframe.
	pub fn push(&mut self, val: T) {
		val.processsendable_write(&mut |buf,len| {
			let mut end = self.end.clone_left();
			let mut end2 = end.clone_right();
			self.list.list_mut().replace(&mut end, &mut end2, unsafe{slice::from_raw_parts(buf, len)}.iter().map(|&x|x));
		});
	}
	/// Map each element of the dataframe using `f`.
	pub fn map<F,T1>(mut self, f: F) -> Vec<'a,T1> where F: Fn(T) -> T1 + ProcessSendable, T1: ProcessSendable {
		let mut iter = self.start.clone_right();
		while iter != self.end {
			// println!("mapping");
			let mut start = iter.clone_left();
			let mut res: T = unsafe{mem::uninitialized()};
			res.processsendable_read(&mut |buf,len| {
				let start2 = iter.clone_left();
				for _ in 0..len {
					self.list.list_mut().increment(&mut iter);
				}
				let res = self.list.list_mut().read(&start2, &iter);
				unsafe{ptr::copy_nonoverlapping(res.as_ptr(), buf, len)};
			});
			// println!("a: {}", res);
			let res = f(res);
			// println!("b: {}", res);
			self.list.list_mut().replace(&mut start, &mut iter, iter::empty());
			res.processsendable_write(&mut |buf,len| {
				let mut start = iter.clone_left();
				self.list.list_mut().replace(&mut start, &mut iter, unsafe{slice::from_raw_parts(buf, len)}.iter().map(|&x|x));
			});
			// println!("written");
		}
		Vec{list:self.list,start:self.start,end:self.end,phantom:marker::PhantomData}
	}
	/// Filter elements of the dataframe, keeping only those that satisfy `f`.
	pub fn filter<F>(&mut self, f: F) where F: Fn(&T) -> bool + ProcessSendable {
		let mut iter = self.start.clone_right();
		while iter != self.end {
			// println!("mapping");
			let mut start = iter.clone_left();
			let mut res: T = unsafe{mem::uninitialized()};
			res.processsendable_read(&mut |buf,len| {
				let start2 = iter.clone_left();
				for _ in 0..len {
					self.list.list_mut().increment(&mut iter);
				}
				let res = self.list.list_mut().read(&start2, &iter);
				unsafe{ptr::copy_nonoverlapping(res.as_ptr(), buf, len)};
			});
			// println!("a: {}", res);
			if !f(&res) {
				self.list.list_mut().replace(&mut start, &mut iter, iter::empty());
			}
			// println!("written");
		}
	}
	// /// Filter elements of the dataframe, keeping only those that satisfy `f`.
	// pub fn group_by<F,K,V>(mut self, f: F) -> hashmap::HashMap<K,Vec<V>> where F: Fn(T) -> (K,V) + ProcessSendable, K: hash::Hash + cmp::Eq + ProcessSendable, Vec<V>: ProcessSendable {
	// 	let mut ret = hashmap::HashMap::new();
	// 	let mut iter = self.start.clone_right();
	// 	while &iter != self.end {
	// 		// println!("mapping");
	// 		let mut start = iter.clone_left();
	// 		let mut res: T = unsafe{mem::uninitialized()};
	// 		res.processsendable_read(&mut |buf,len| {
	// 			let start2 = iter.clone_left();
	// 			for _ in 0..len {
	// 				self.list.list_mut().increment(&mut iter);
	// 			}
	// 			let res = self.list.list_mut().read(&start2, &iter);
	// 			unsafe{ptr::copy_nonoverlapping(res.as_ptr(), buf, len)};
	// 		});
	// 		// println!("a: {}", res);
	// 		let (key, val) = f(res);
	// 		ret.insert(key, vec![val]);
	// 		// println!("written");
	// 	}
	// 	ret
	// }
	/// Reduce the dataframe into one value.
	///
	/// # Examples
	///
	/// ```
	/// let mut data_frame: Vec<String> = Vec::new();
	/// data_frame.push(String::from("http://google.com"));
	/// data_frame.push(String::from("http://bbc.co.uk"));
	/// data_frame.push(String::from("https://hadean.com"));
	/// let total_length = data_frame.reduce(0, |acc, val| acc + val.len() as u64);
	/// ```
	///
	pub fn reduce<F,T1>(&mut self, initial: T1, f: F) -> T1 where F: Fn(T1, &T) -> T1 + ProcessSendable, T1: ProcessSendable {
		let mut acc = initial;
		let mut iter = self.start.clone_right();
		while iter != self.end {
			// println!("mapping");
			let mut res: T = unsafe{mem::uninitialized()};
			res.processsendable_read(&mut |buf,len| {
				let start2 = iter.clone_left();
				for _ in 0..len {
					self.list.list_mut().increment(&mut iter);
				}
				let res = self.list.list_mut().read(&start2, &iter);
				unsafe{ptr::copy_nonoverlapping(res.as_ptr(), buf, len)};
			});
			// println!("a: {}", res);
			acc = f(acc, &res);
			// println!("written");
		}
		acc
	}
	/// Returns an iterator over the `Vec`.
	///
	/// # Examples
	///
	/// ```
	/// let mut data_frame: Vec<String> = Vec::new();
	/// data_frame.push(String::from("http://google.com"));
	/// data_frame.push(String::from("http://bbc.co.uk"));
	/// data_frame.push(String::from("https://hadean.com"));
	/// let mut iterator = x.iter();
	///
	/// assert_eq!(iterator.next(), Some(String::from("http://google.com")));
	/// assert_eq!(iterator.next(), Some(String::from("http://bbc.co.uk")));
	/// assert_eq!(iterator.next(), Some(String::from("https://hadean.com")));
	/// assert_eq!(iterator.next(), None);
	/// ```
	pub fn iter<'b>(&'b mut self) -> VecIter<'a,'b,T> {
		let iter = self.start.clone_right();
		VecIter{list:self,iter:iter}
	}
}

pub struct VecIter<'a,'b,T> where T: 'a + ProcessSendable, 'a:'b {
	list: &'b mut Vec<'a,T>,
	iter: LeafIndex
}
impl<'a,'b,T> iter::Iterator for VecIter<'a,'b,T> where T: ProcessSendable {
	type Item = T; // .iter() should be over &T... but that's not possible here
	fn next(&mut self) -> Option<Self::Item> {
		if self.iter == self.list.end {
			None
		} else {
			let mut res: T = unsafe{mem::uninitialized()};
			res.processsendable_read(&mut |buf,len| {
				let iter = self.iter.clone_left();
				for _ in 0..len {
					self.list.list.list_mut().increment(&mut self.iter);
				}
				let res = self.list.list.list_mut().read(&iter, &self.iter);
				unsafe{ptr::copy_nonoverlapping(res.as_ptr(), buf, len)};
			});
			Some(res)
		}
	}
}

#[cfg(test)]
mod tests {
	use super::*;
	#[test]
	// #[ignore]
	fn dataframe() {
		let mut data_frame: Vec<String> = Vec::new();
		data_frame.push(String::from("http://google.com"));
		data_frame.push(String::from("http://bbc.co.uk"));
		data_frame.push(String::from("http://hadean.com"));
		let mut data_frame: Vec<String> = data_frame.map(|mut val|{val.push_str("/xyz");val});
		data_frame.filter(|val|val.len() > 20);
		for string in data_frame.iter() {
			println!("{}", string);
		}
		let total_length = data_frame.reduce(0u64, |acc, val| acc + val.len() as u64);
		println!("{}", total_length);
	}
}