Commit 61493bc7 authored by Paras Garg's avatar Paras Garg

Updated

parent d16dbc01
Set up disni
export LD_LIBRARY_PATH=/usr/local/lib
How to setup RocksDB jar
1. Compiling RocksDb java from source code
make -j8 rocksjava DEBUG_LEVEL=0
......
app.HOST=192.168.200.20
app.MASTER=true
app.cpu_affinity=4
#this is how to pair threads to cpu cores
app.cpu_affinity=2
//No of thread for each connectin
app.rdma_cluster_size=1
app.rdma_receive_queue=16
app.rdma_send_queue=16
app.rdma_polling=false
......
app.HOST=192.168.200.20
app.MASTER=false
app.cpu_affinity=4
app.rdma_cluster_size=3
app.rdma_receive_queue=16
app.rdma_send_queue=16
app.rdma_polling=false
......@@ -8,6 +9,7 @@ app.rdma_max_inline=0
app.rdma_server_port=1921
app.NETWORK_HANDLER_THREADS=10
app.REPLICATION_HANDLER_THREADS=10
app.db_path=/tmp/rocks
#The below properties are for master node only.
app.follower_registration_port=9876
......
apply plugin: "application"
apply plugin: "java-library"
application
{
mainClassName = "hpdos_rdma_offloaded.MetadataServer"
}
mainClassName = "hpdos_rdma_offloaded.MetadataServer"
version '1.0-SNAPSHOT'
repositories
......
package hpdos_rdma_offloaded;
import java.io.File;
import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.util.Properties;
......@@ -18,13 +20,15 @@ import hpdos_rdma_offloaded.handler.RegistrationHandler;
import hpdos_rdma_offloaded.handler.StorageHandler;
import hpdos_rdma_offloaded.protocol.Request;
import hpdos_rdma_offloaded.protocol.Response;
import hpdos_rdma_offloaded.service.MetadataServerService;
import hpdos_rdma_offloaded.service.FollowerMetadataService;
import hpdos_rdma_offloaded.service.MasterMetadataServerService;
public class MetadataServer implements Runnable{
private String host;
private int poolsize;
private int recvQueue;
private int sendQueue;
private int clusterSize;
private int wqSize = recvQueue;
private boolean polling;
private int maxinline = 0;
......@@ -34,49 +38,68 @@ public class MetadataServer implements Runnable{
boolean isMaster;
StorageHandler storageHandler;
NetworkHandler networkHandler;
MetadataServerService metadataServerService;
MasterMetadataServerService metadataServerService = null;
FollowerMetadataService followerMetadataServerService = null;
RegistrationHandler registrationHandler = null;
public final static Properties properties= new Properties();;
public final static Properties properties= new Properties();
String propertyPath;
public MetadataServer(String propertyPath) throws IOException,RocksDBException {
public MetadataServer(String propertyPath) throws Exception,RocksDBException {
this.propertyPath = propertyPath;
InputStream inputStream = new FileInputStream(propertyPath);
properties.load(inputStream);
boolean isMaster = Boolean.valueOf((String)properties.get("app.MASTER"));
if(isMaster)
{
System.out.println("STARTING MASTER SERVICES.");
}
else
{
System.out.println("STARTING FOLLOWER SERVICES.");
}
isMaster = Boolean.valueOf((String)properties.get("app.MASTER"));
storageHandler = new StorageHandler(properties.getProperty("app.db_path"));
networkHandler = new NetworkHandler(this.storageHandler, properties);
metadataServerService = new MetadataServerService(this.networkHandler);
registrationHandler = new RegistrationHandler(properties);
this.isMaster = Boolean.valueOf((String)properties.get("app.MASTER"));
host = (String) properties.get("app.HOST");
poolsize = Integer.valueOf( (String)properties.get("app.cpu_affinity"));
recvQueue = Integer.valueOf( (String)properties.get("app.rdma_receive_queue"));
sendQueue = Integer.valueOf( (String) properties.get("app.rdma_send_queue"));
host = properties.getProperty("app.HOST");
poolsize = Integer.valueOf( properties.getProperty("app.cpu_affinity"));
recvQueue = Integer.valueOf(properties.getProperty("app.rdma_receive_queue"));
sendQueue = Integer.valueOf( properties.getProperty("app.rdma_send_queue"));
clusterSize = Integer.valueOf(properties.getProperty("app.rdma_cluster_size"));
wqSize = recvQueue;
maxinline = Integer.valueOf((String) properties.get("app.rdma_max_inline"));
rdma_port = Integer.valueOf((String) properties.get("app.rdma_server_port"));
polling = Boolean.valueOf((String) properties.get("app.rdma_polling"));
maxinline = Integer.valueOf(properties.getProperty("app.rdma_max_inline"));
rdma_port = Integer.valueOf(properties.getProperty("app.rdma_server_port"));
polling = Boolean.valueOf(properties.getProperty("app.rdma_polling"));
int registrationPort = Integer.valueOf(properties.getProperty("app.follower_registration_port"));
clusterAffinities = new long[poolsize];
for (int i = 0; i < poolsize; i++)
{
long cpu = 1L << i;
clusterAffinities[i] = cpu;
System.out.println(cpu);
}
isMaster = Boolean.valueOf(properties.getProperty("app.MASTER"));
storageHandler = new StorageHandler(properties.getProperty("app.db_path"));
networkHandler = new NetworkHandler(isMaster);
// metadataServerService = new MasterMetadataServerService(this.networkHandler,storageHandler);
registrationHandler = new RegistrationHandler(isMaster,host,registrationPort, networkHandler);
if(isMaster)
{
System.out.println("STARTING MASTER SERVICES.");
metadataServerService = new MasterMetadataServerService(this.networkHandler,storageHandler);
}
else
{
System.out.println("STARTING FOLLOWER SERVICES.");
followerMetadataServerService = new FollowerMetadataService(storageHandler);
}
}
public void run(){
// MetadataServerService rpcService = new MetadataServerService();
DaRPCServerGroup<Request, Response> group = null;
try
{
group = DaRPCServerGroup.createServerGroup(this.metadataServerService, this.clusterAffinities, -1, maxinline, polling, recvQueue, sendQueue, wqSize, 32);
if (isMaster) {
group = DaRPCServerGroup.createServerGroup(this.metadataServerService, this.clusterAffinities, 1000, maxinline, polling, recvQueue, sendQueue, wqSize, clusterSize);
} else {
group = DaRPCServerGroup.createServerGroup(this.followerMetadataServerService, this.clusterAffinities, 1000, maxinline, polling, recvQueue, sendQueue, wqSize, clusterSize);
}
RdmaServerEndpoint<DaRPCServerEndpoint<Request, Response>> serverEp = null;
serverEp = group.createServerEndpoint();
InetSocketAddress address = new InetSocketAddress(host, rdma_port);
......@@ -102,7 +125,10 @@ public class MetadataServer implements Runnable{
public static void main(String[] args) throws Exception
{
RocksDB.loadLibrary();
new MetadataServer(args[0]).run();
System.out.println("Rocksdb Exception occur");
//var process = Runtime.getRuntime().exec("/bin/sh export LD_LIBRARY_PATH=/usr/local/lib");
//process.wait();
// System.out.println("Exit value "+process.exitValue());
Runnable r = new MetadataServer(args[0]);
new Thread(r).start();
}
}
package hpdos_rdma_offloaded.handler;
import java.net.SocketAddress;
import java.net.InetSocketAddress;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Properties;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import com.ibm.darpc.DaRPCClientEndpoint;
import com.ibm.darpc.DaRPCClientGroup;
import com.ibm.darpc.DaRPCFuture;
import com.ibm.darpc.DaRPCStream;
import org.rocksdb.RocksDBException;
import hpdos_rdma_offloaded.lib.Follower;
import hpdos_rdma_offloaded.lib.Packet;
import hpdos_rdma_offloaded.protocol.InvalidationRequest;
import hpdos_rdma_offloaded.protocol.InvalidationResponse;
import hpdos_rdma_offloaded.protocol.InvalidationRpcProtocol;
import hpdos_rdma_offloaded.protocol.Request;
import hpdos_rdma_offloaded.protocol.RequestType;
import hpdos_rdma_offloaded.protocol.Response;
import hpdos_rdma_offloaded.protocol.RpcProtocol;
public class NetworkHandler {
StorageHandler storageHandler;
private ExecutorService executorService;
public HashMap<String, DaRPCStream<Request, Response>> streams;
// public static List<Follower> followers = new ArrayList<>();
// public HashMap<String, DaRPCStream<Request, Response>> streams;
public static List <DaRPCStream<InvalidationRequest, InvalidationResponse>> invalidationStreams = new ArrayList<>();
public static List<Follower> followers = new ArrayList<>();
boolean isMaster;
ConcurrentHashMap<SocketAddress,DaRPCStream<InvalidationRequest, InvalidationResponse>> invalidationClientEpStreams;
public NetworkHandler(boolean isMaster){
this.executorService = Executors.newFixedThreadPool(10);
this.isMaster = isMaster;
}
public NetworkHandler(){}
public NetworkHandler(StorageHandler storageHandler, Properties properties){
this.storageHandler = storageHandler;
this.executorService = Executors.newFixedThreadPool(10);
this.isMaster = Boolean.valueOf((String)properties.get("app.MASTER"));
public void setClientStreams(ConcurrentHashMap<SocketAddress,DaRPCStream<InvalidationRequest, InvalidationResponse>> invalidationClientEpStreams)
{
this.invalidationClientEpStreams = invalidationClientEpStreams;
}
public void connectToSals() throws Exception{
// RpcProtocol rpcProtocol = new RpcProtocol();
// DaRPCClientGroup<Request, Response> group = DaRPCClientGroup.createClientGroup(rpcProtocol, 100, 0, 16, 16);
// InetSocketAddress address = new InetSocketAddress("192.168.200.50", 1919);
// DaRPCClientEndpoint<Request, Response> clientEp = group.createEndpoint();
// clientEp.connect(address, 1000);
// DaRPCStream<Request, Response> stream = clientEp.createStream();
// if(!stream.isEmpty()){
// this.streams.put("192.168.200.50", stream);
// }
public void connectToSal(String ip) throws Exception{
Thread.sleep(3000);
System.out.println("Making rdma connection to client for invalidation request.");
// Follower follower = new Follower();
// follower.setFollowerId(uid);
// follower.setIpAddress(ip);
InvalidationRpcProtocol rpcProtocol = new InvalidationRpcProtocol();
DaRPCClientGroup<InvalidationRequest, InvalidationResponse> group = DaRPCClientGroup.createClientGroup(rpcProtocol, 100, 0, 16, 16);
InetSocketAddress address = new InetSocketAddress(ip, 1922);
DaRPCClientEndpoint<InvalidationRequest, InvalidationResponse> clientEp = group.createEndpoint();
clientEp.connect(address, 1000);
DaRPCStream<InvalidationRequest, InvalidationResponse> stream = clientEp.createStream();
NetworkHandler.invalidationStreams.add(stream);
}
public void connectToFollower(String uid , String ip) throws Exception {
......@@ -66,15 +80,13 @@ public class NetworkHandler {
clientEp.connect(address, 1000);
DaRPCStream<Request, Response> stream = clientEp.createStream();
follower.setStream(stream);
ReplicationHandler.followers.add(follower);
NetworkHandler.followers.add(follower);
}
// Change the Futre to CompletableFuture to achieve correct asynchronousness
public void create(Packet packet) throws InterruptedException, ExecutionException,RocksDBException{
System.out.println("Received create request for key/value: " + new String(packet.getKey()) + "/" + new String(packet.getValue()));
storageHandler.create(packet.getKey(), packet.getValue());
if (this.isMaster)
{
System.out.println("Starting replication");
......@@ -85,32 +97,27 @@ public class NetworkHandler {
return replicationResponse;
});
Response response = futureReplication.get();
System.out.println("Replicating complete Ack" + response.getAck());
System.out.println("Replication complete");
}
}
public byte[] read(Packet packet) throws RocksDBException,InterruptedException,ExecutionException{
public void read(Packet packet) throws RocksDBException,InterruptedException,ExecutionException{
System.out.println("Received read request for key/value: " + new String(packet.getKey()));
byte[] value = storageHandler.read(packet.getKey());
if(value != null)
System.out.println("Got value "+ new String(value));
if (this.isMaster)
{
System.out.println("Starting replication");
Future<Response> futureReplication = this.executorService.submit(()->{
Response replicationResponse = new Response();
// Write code to replicate the data to other nics
ReplicationHandler.replicateMetadata(packet);
// ReplicationHandler.replicateMetadata(packet);
return replicationResponse;
});
Response response = futureReplication.get();
System.out.println("Replicating complete Ack "+ response.getAck());
}
return value;
}
public void update(Packet packet) throws InterruptedException, ExecutionException,RocksDBException{
this.storageHandler.update(packet.getKey(), packet.getValue());
if (isMaster) {
Response response = new Response();
Future<Response> futureReplication = this.executorService.submit(()->{
......@@ -122,7 +129,7 @@ public class NetworkHandler {
});
Future<Boolean> futureInvalidation = this.executorService.submit(()->{
sendInvalidationRequest(packet);
// sendInvalidationRequest(packet);
return false;
});
response = futureReplication.get();
......@@ -130,13 +137,11 @@ public class NetworkHandler {
System.out.println("Replicating complete Ack "+ response.getAck());
}
// Write code to parse the responses here.
}
// To implement delete
public void delete(Packet packet) throws RocksDBException,InterruptedException,ExecutionException
{
this.storageHandler.delete(packet.getKey());
if (this.isMaster) {
System.out.println("Starting replication");
Future<Response> futureReplication = this.executorService.submit(()->{
......@@ -149,27 +154,82 @@ public class NetworkHandler {
System.out.println("Replicating complete "+ response.getAck());
}
}
public void sendInvalidateRequest(Request request) throws InterruptedException, ExecutionException{
/* Paras method
public DaRPCFuture<InvalidationRequest, InvalidationResponse> sendInvalidateRequest(Request request){
System.out.println("Sending Invalidation Request");
InvalidationRequest invalidationRequest = new InvalidationRequest();
invalidationRequest.setKey(request.getKey());
ArrayList<DaRPCFuture<InvalidationRequest,InvalidationResponse>> requestFutures;
requestFutures = new ArrayList<>();
for(var ep : invalidationClientEpStreams.keySet())
{
try
{
requestFutures.add(invalidationClientEpStreams.get(ep).
request(invalidationRequest, new InvalidationResponse(), false));
}catch(Exception e)
{
e.printStackTrace();
}
public void sendInvalidationRequest(Packet packet) throws InterruptedException, ExecutionException{
Set<Callable<Response>> callables = new HashSet<>();
for (DaRPCStream<Request, Response> stream: streams.values()) {
callables.add(() -> {
Request request = new Request();
Response response = new Response();
request.setRequestType(RequestType.INVALIDATE);
request.setKey(null);
request.setValue(null);
response = stream.request(request, response, false).get();
return response;
}
return requestFutures;
*/
// Pramod method
ExecutorService executorService = Executors.newFixedThreadPool(10);
System.out.println("Sending Invalidation Request");
InvalidationRequest invalidationRequest = new InvalidationRequest();
invalidationRequest.setKey(request.getKey());
Set<Callable<DaRPCFuture<InvalidationRequest, InvalidationResponse>>> callables = new HashSet<>();
ArrayList<DaRPCFuture<InvalidationRequest,InvalidationResponse>> requestFutures = new ArrayList<>();;
for (DaRPCStream<InvalidationRequest, InvalidationResponse> invalidationStream : NetworkHandler.invalidationStreams) {
callables.add(()->{
InvalidationResponse response = new InvalidationResponse();
DaRPCFuture<InvalidationRequest, InvalidationResponse> future = invalidationStream.request(invalidationRequest, response, false);
return future;
});
}
List<Future<DaRPCFuture<InvalidationRequest, InvalidationResponse>>> futures = executorService.invokeAll(callables);
List<Future<Response>> futures = this.executorService.invokeAll(callables);
for (Future<DaRPCFuture<InvalidationRequest, InvalidationResponse>> future: futures) {
InvalidationResponse response = future.get().get();
}
}
public boolean consumeReplicationRequestFutures(ArrayList<DaRPCFuture<InvalidationRequest,InvalidationResponse>> requestFutures )
throws InterruptedException,ExecutionException{
for (Future<Response> future: futures) {
Response response;
response = future.get();
//poll once
var it = requestFutures.iterator();
while(it.hasNext())
{
var future = it.next();
if(future.isDone())
{
it.remove();
}
}
if(requestFutures.size() == 0)
return true;
//if not done now wait for 100 milliseconds before sending response
requestFutures.get(0).get(100,TimeUnit.MILLISECONDS);
it = requestFutures.iterator();
while(it.hasNext())
{
var future = it.next();
if(future.isDone())
{
it.remove();
}
}
if(requestFutures.size() == 0)
return true;
//cancel each future
requestFutures.forEach(future -> future.cancel(true));
return false;
}
}
......@@ -7,31 +7,34 @@ import java.lang.ClassNotFoundException;
import java.net.InetAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.Properties;
import java.util.UUID;
public class RegistrationHandler {
public ServerSocket server;
public int port = 9876;
boolean isMaster;
String registrationServerIp;
int follower_registration_port;
int sal_registration_port;
public RegistrationHandler(Properties properties) throws IOException {
this.isMaster = Boolean.valueOf((String)properties.get("app.MASTER"));
NetworkHandler networkHandler;
if (this.isMaster) {
this.follower_registration_port = Integer.valueOf((String)properties.get("app.follower_registration_port"));
this.sal_registration_port = Integer.valueOf((String)properties.get("app.follower_registration_port"));
public RegistrationHandler(boolean isMaster, String registrationServerIP, int registrationPort, NetworkHandler networkHandler) throws IOException {
this.registrationServerIp = registrationServerIP;
this.follower_registration_port = 9876;
this.sal_registration_port = 9875;
this.networkHandler = networkHandler;
if (isMaster)
{
startFollowerRegistrationHandlerServer();
startSalRegistrationHandlerServer();
} else {
}
else {
// Follower follower = new Follower();
// follower.setFollowerId( UUID.randomUUID().toString());
// follower.setIpAddress("192.168.200.20");
String followerUUID = UUID.randomUUID().toString();
String followerIp = "192.168.200.20";
String followerIp = registrationServerIP;
String message = followerUUID + ";" + followerIp;
sendRegistrationRequest(message);
}
......@@ -39,10 +42,10 @@ public class RegistrationHandler {
private void sendRegistrationRequest(String message) throws IOException {
System.out.println("Sending a registration request");
InetAddress host = InetAddress.getByName("192.168.200.20");
InetAddress host = InetAddress.getByName(this.registrationServerIp);
Socket socket = null;
ObjectOutputStream oos = null;
socket = new Socket(host.getHostName(), 9876);
socket = new Socket(this.registrationServerIp, 9876);
oos = new ObjectOutputStream(socket.getOutputStream());
oos.writeObject(message);
}
......@@ -50,10 +53,12 @@ public class RegistrationHandler {
public void startFollowerRegistrationHandlerServer(){
Runnable serverTask = new Runnable() {
@Override
public void run() {
try {
InetAddress addr = InetAddress.getByName("192.168.200.20");
ServerSocket serverSocket = new ServerSocket(9876, 50, addr);
public void run()
{
try
{
InetAddress addr = InetAddress.getByName(registrationServerIp);
ServerSocket serverSocket = new ServerSocket(follower_registration_port, 50, addr);
System.out.println("Started the follower registration service...");
while (true) {
Socket socket = serverSocket.accept();
......@@ -65,8 +70,8 @@ public class RegistrationHandler {
String followerIp = arr[1];
System.out.println("Got follower ip: " + followerIp + " , with UUID of: " + followerUUID);
NetworkHandler networkHandler = new NetworkHandler();
networkHandler.connectToFollower(followerUUID, followerIp);
NetworkHandler networkHandler = new NetworkHandler();
networkHandler.connectToFollower(followerUUID, followerIp);
System.out.println("RDMA connection establised");
// Add the follower details to the respective class/component
// clientProcessingPool.submit(new ClientTask(clientSocket));
......@@ -86,16 +91,18 @@ public class RegistrationHandler {
@Override
public void run() {
try {
ServerSocket serverSocket = new ServerSocket(9875);
InetAddress addr = InetAddress.getByName(registrationServerIp);
ServerSocket serverSocket = new ServerSocket(sal_registration_port, 50, addr);
System.out.println("Started the SAL registration service...");
while (true) {
Socket socket = serverSocket.accept();
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
String message = (String) ois.readObject();
// Send a response back with the list of followers to the SAL
String ip = (String) ois.readObject();
NetworkHandler networkHandler = new NetworkHandler();
networkHandler.connectToSal(ip);
// clientProcessingPool.submit(new ClientTask(clientSocket));
}
} catch (IOException | ClassNotFoundException e) {
} catch (Exception e) {
System.err.println("Unable to process client request");
e.printStackTrace();
}
......
......@@ -4,6 +4,9 @@ import org.rocksdb.RocksDB;
import org.rocksdb.Options;
import org.rocksdb.RocksDBException;
import hpdos_rdma_offloaded.protocol.Request;
public class StorageHandler implements AutoCloseable{
RocksDB rocksDB;
......@@ -27,16 +30,32 @@ public class StorageHandler implements AutoCloseable{
public void create(byte[] key,byte[] value) throws RocksDBException{
rocksDB.put(key,value);
}
public void create(Request request) throws RocksDBException {
System.out.println("Writing to local");
rocksDB.put(request.getKey(),request.getValue());
}
public byte[] read(byte[] key)throws RocksDBException
{
return rocksDB.get(key);
}
public byte[] read(Request request)throws RocksDBException
{
return rocksDB.get(request.getKey());
}
public void update(byte[] key,byte[] value) throws RocksDBException
{
rocksDB.put(key,value);
}
public void update(Request request) throws RocksDBException
{
rocksDB.put(request.getKey(),request.getValue());
}
public void delete(byte[] key) throws RocksDBException
{
rocksDB.delete(key);
}
public void delete(Request request) throws RocksDBException
{
rocksDB.delete(request.getKey());
}
}
\ No newline at end of file
......@@ -6,6 +6,8 @@ public interface AckType
public static int NOTFOUND = 1;
public static int NOTALLOWED = 2;
public static int DBFAILED = 3;
public static int SUCCESS_WITH_VALUE = 4;
public static int FAILED = 4;
public static int SUCCESS_WITH_VALUE = 5;
public static int INVALID_REQUEST = 6;
}
\ No newline at end of file
package hpdos_rdma_offloaded.rdma_invalidation_protocol;
package hpdos_rdma_offloaded.protocol;
import java.io.IOException;
import java.nio.ByteBuffer;
......@@ -6,34 +6,31 @@ import java.nio.ByteBuffer;
import com.ibm.darpc.DaRPCMessage;
public class InvalidationRequest implements DaRPCMessage {
public String key;
private static char[] dst = new char[100];
private int SERIALIZED_SIZE = 256;
public byte[] key;
private static int SERIALIZED_SIZE = 128 ;
@Override
public int write(ByteBuffer buffer) throws IOException {
buffer.asCharBuffer().put(key);
buffer.put(key);
return SERIALIZED_SIZE;
}
@Override
public void update(ByteBuffer buffer) throws IOException {
buffer.asCharBuffer().get(dst, 0, 99);
String s = new String(dst);
s = s.trim();
key = s;
if(key == null)
key = new byte[128];
buffer.get(key);
}
@Override
public int size() {
return SERIALIZED_SIZE;
}
public String getKey() {
public byte[] getKey() {
return key;
}
public void setKey(String key) {
public void setKey(byte[] key) {
this.key = key;
}
}
package hpdos_rdma_offloaded.rdma_invalidation_protocol;
package hpdos_rdma_offloaded.protocol;
import java.io.IOException;
import java.nio.ByteBuffer;
......@@ -6,23 +6,32 @@ import java.nio.ByteBuffer;