Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in
Toggle navigation
H
hpdos
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Analytics
Analytics
CI / CD
Repository
Value Stream
Wiki
Wiki
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
SYNERG
hpdos
Commits
2b6c1a94
Commit
2b6c1a94
authored
Nov 29, 2021
by
Pramod S
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
update api
parent
58600282
Changes
14
Hide whitespace changes
Inline
Side-by-side
Showing
14 changed files
with
193 additions
and
103 deletions
+193
-103
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/MetadataServer.java
...ed/src/main/java/hpdos_rdma_offloaded/MetadataServer.java
+1
-1
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/NetworkHandler.java
...ain/java/hpdos_rdma_offloaded/handler/NetworkHandler.java
+13
-16
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/NetworkHandlerM.java
...in/java/hpdos_rdma_offloaded/handler/NetworkHandlerM.java
+136
-35
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/ReplicationHandler.java
...java/hpdos_rdma_offloaded/handler/ReplicationHandler.java
+1
-1
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/StorageHandler.java
...ain/java/hpdos_rdma_offloaded/handler/StorageHandler.java
+3
-4
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/protocol/Request.java
.../src/main/java/hpdos_rdma_offloaded/protocol/Request.java
+0
-2
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/protocol/Response.java
...src/main/java/hpdos_rdma_offloaded/protocol/Response.java
+0
-2
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/rdma/DaRPCClientEndpointM.java
.../java/hpdos_rdma_offloaded/rdma/DaRPCClientEndpointM.java
+4
-4
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/rdma/DaRPCClientGroupM.java
...ain/java/hpdos_rdma_offloaded/rdma/DaRPCClientGroupM.java
+1
-1
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/rdma/DaRPCFutureM.java
...src/main/java/hpdos_rdma_offloaded/rdma/DaRPCFutureM.java
+3
-2
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/rdma/DaRPCStreamM.java
...src/main/java/hpdos_rdma_offloaded/rdma/DaRPCStreamM.java
+4
-4
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/service/FollowerMetadataService.java
...hpdos_rdma_offloaded/service/FollowerMetadataService.java
+2
-5
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/service/MasterMetadataServerService.java
...s_rdma_offloaded/service/MasterMetadataServerService.java
+11
-11
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/service/MetadataServiceMaster.java
...a/hpdos_rdma_offloaded/service/MetadataServiceMaster.java
+14
-15
No files found.
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/MetadataServer.java
View file @
2b6c1a94
...
...
@@ -19,9 +19,9 @@ import hpdos_rdma_offloaded.handler.StorageHandler;
import
hpdos_rdma_offloaded.protocol.Request
;
import
hpdos_rdma_offloaded.protocol.Response
;
import
hpdos_rdma_offloaded.protocol.RpcProtocol
;
import
hpdos_rdma_offloaded.rdma.DaRPCClientGroupM
;
import
hpdos_rdma_offloaded.service.FollowerMetadataService
;
import
hpdos_rdma_offloaded.service.MetadataServiceMaster
;
import
rdma.DaRPCClientGroupM
;
public
class
MetadataServer
{
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/NetworkHandler.java
View file @
2b6c1a94
...
...
@@ -89,12 +89,12 @@ public class NetworkHandler {
// Change the Futre to CompletableFuture to achieve correct asynchronousness
public
void
create
(
Packet
packet
)
throws
InterruptedException
,
ExecutionException
,
RocksDBException
{
ReplicationHandler
replicationHandler
=
new
ReplicationHandler
();
logger
.
info
(
"Received create request for key/value: "
+
new
String
(
packet
.
getKey
())
+
"/"
+
new
String
(
packet
.
getValue
()));
logger
.
debug
(
"Received create request for key/value: "
+
new
String
(
packet
.
getKey
())
+
"/"
+
new
String
(
packet
.
getValue
()));
Stopwatch
stopWatch
=
Stopwatch
.
createUnstarted
();
stopWatch
.
start
();
if
(
this
.
isMaster
)
{
// logger.
info
("Starting replication");
// logger.
debug
("Starting replication");
// Future<Response> futureReplication = this.executorService.submit(()->{
// Response replicationResponse = new Response();
// // Write code to replicate the data to other nics
...
...
@@ -122,21 +122,20 @@ public class NetworkHandler {
stopWatch
.
stop
();
}
logger
.
info
(
"Future time: "
+
stopWatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Future time: "
+
stopWatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopWatch
.
reset
();
// stopWatch.start();
// Response response = futureReplication.get();
// stopWatch.stop();
// logger.
info
("Replication time: " + stopWatch.elapsed(TimeUnit.MICROSECONDS));
logger
.
info
(
"Replication complete"
);
// logger.
debug
("Replication time: " + stopWatch.elapsed(TimeUnit.MICROSECONDS));
logger
.
debug
(
"Replication complete"
);
}
public
void
read
(
Packet
packet
)
throws
RocksDBException
,
InterruptedException
,
ExecutionException
{
logger
.
info
(
"Received read request for key/value: "
+
new
String
(
packet
.
getKey
()));
logger
.
debug
(
"Received read request for key/value: "
+
new
String
(
packet
.
getKey
()));
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
...
...
@@ -161,7 +160,7 @@ public class NetworkHandler {
return
replicationResponse
;
});
stopWatch
.
stop
();
logger
.
info
(
"Future time: "
+
stopWatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Future time: "
+
stopWatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopWatch
.
reset
();
// Future<Boolean> futureInvalidation = this.executorService.submit(()->{
// sendInvalidationRequest(packet);
...
...
@@ -171,13 +170,13 @@ public class NetworkHandler {
stopWatch
.
start
();
response
=
futureReplication
.
get
();
stopWatch
.
stop
();
logger
.
info
(
"Replication time: "
+
stopWatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Replication time: "
+
stopWatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopWatch
.
reset
();
// stopWatch.start();
// futureInvalidation.get();
// stopWatch.stop();
// logger.
info
("Invalidation time: " + stopWatch.elapsed(TimeUnit.MICROSECONDS));
logger
.
info
(
"Replicating complete Ack "
+
response
.
getAck
());
// logger.
debug
("Invalidation time: " + stopWatch.elapsed(TimeUnit.MICROSECONDS));
logger
.
debug
(
"Replicating complete Ack "
+
response
.
getAck
());
}
}
...
...
@@ -187,7 +186,7 @@ public class NetworkHandler {
{
ReplicationHandler
replicationHandler
=
new
ReplicationHandler
();
if
(
this
.
isMaster
)
{
logger
.
info
(
"Starting replication"
);
logger
.
debug
(
"Starting replication"
);
Future
<
Response
>
futureReplication
=
this
.
executorService
.
submit
(()->{
Response
replicationResponse
=
new
Response
();
// Write code to replicate the data to other nics
...
...
@@ -195,14 +194,12 @@ public class NetworkHandler {
return
replicationResponse
;
});
Response
response
=
futureReplication
.
get
();
logger
.
info
(
"Replicating complete "
+
response
.
getAck
());
logger
.
debug
(
"Replicating complete "
+
response
.
getAck
());
}
}
//Paras method
/* public ArrayList<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;
...
...
@@ -225,7 +222,7 @@ public class NetworkHandler {
/*
public void sendInvalidateRequest(Request request) throws InterruptedException, ExecutionException{
ExecutorService executorService = Executors.newWorkStealingPool();
logger.
info
("Sending Invalidation Request");
logger.
debug
("Sending Invalidation Request");
InvalidationRequest invalidationRequest = new InvalidationRequest();
invalidationRequest.setKey(request.getKey());
Set<Callable<DaRPCFuture<InvalidationRequest, InvalidationResponse>>> callables = new HashSet<>();
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/NetworkHandlerM.java
View file @
2b6c1a94
package
hpdos_rdma_offloaded.handler
;
import
java.io.IOException
;
import
org.apache.log4j.Logger
;
import
org.checkerframework.checker.units.qual.degrees
;
import
java.lang.System.Logger.Level
;
import
java.util.concurrent.CopyOnWriteArrayList
;
import
com.ibm.disni.util.NativeAffinity
;
...
...
@@ -8,15 +12,15 @@ import com.ibm.disni.util.NativeAffinity;
import
hpdos_rdma_offloaded.protocol.InvalidationRequest
;
import
hpdos_rdma_offloaded.protocol.InvalidationResponse
;
import
hpdos_rdma_offloaded.protocol.Request
;
import
hpdos_rdma_offloaded.protocol.RequestType
;
import
hpdos_rdma_offloaded.protocol.Response
;
import
rdma.DaRPCFutureM
;
import
rdma.DaRPCStreamM
;
import
hpdos_rdma_offloaded.
rdma.DaRPCFutureM
;
import
hpdos_rdma_offloaded.
rdma.DaRPCStreamM
;
public
class
NetworkHandlerM
implements
Runnable
{
private
CopyOnWriteArrayList
<
DaRPCStreamM
<
InvalidationRequest
,
InvalidationResponse
>>
invalidationStreams
;
private
CopyOnWriteArrayList
<
DaRPCStreamM
<
Request
,
Response
>>
replicationStreams
;
final
static
Logger
logger
=
Logger
.
getLogger
(
NetworkHandlerM
.
class
);
public
NetworkHandlerM
()
{
...
...
@@ -30,50 +34,135 @@ public class NetworkHandlerM implements Runnable{
{
this
.
replicationStreams
.
add
(
stream
);
}
public
void
run
()
void
processCloseInvalidationStream
(
DaRPCStreamM
<
InvalidationRequest
,
InvalidationResponse
>
stream
)
{
NativeAffinity
.
setAffinity
(
1
<<
1
);
while
(
true
)
{
for
(
var
stream
:
invalidationStreams
)
stream
.
completedList
.
forEach
(
future
->
{
try
{
DaRPCFutureM
<
InvalidationRequest
,
InvalidationResponse
>
future
=
stream
.
processStream
();
if
(
future
!=
null
&&
future
.
metadataResponse
.
invalidationTotal
==
future
.
metadataResponse
.
invalidationDone
.
incrementAndGet
())
if
(
future
.
metadataResponse
.
invalidationTotal
==
future
.
metadataResponse
.
invalidationDone
.
incrementAndGet
())
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
{
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
try
{
future
.
metadataResponse
.
event
.
triggerResponse
();
}
catch
(
IOException
ex
)
{
logger
.
error
(
ex
);
}
}
}
catch
(
Exception
exx
)
});
stream
.
endpoint
.
pendingFutures
.
forEach
((
ticket
,
future
)->{
if
(
future
.
metadataResponse
.
invalidationTotal
==
future
.
metadataResponse
.
invalidationDone
.
incrementAndGet
())
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
{
try
{
future
.
metadataResponse
.
event
.
triggerResponse
();
}
catch
(
IOException
ex
)
{
logger
.
error
(
ex
);
}
}
});
}
void
processInvalidationStream
(
DaRPCStreamM
<
InvalidationRequest
,
InvalidationResponse
>
stream
)
{
try
{
DaRPCFutureM
<
InvalidationRequest
,
InvalidationResponse
>
future
=
stream
.
processStream
();
if
(
future
!=
null
)
{
exx
.
printStackTrace
();
//logger.debug("future"+future.metadataResponse.invalidationTotal+" "+future.metadataResponse.invalidationDone);
//logger.debug(future.metadataResponse.state);
if
(
future
.
metadataResponse
.
invalidationTotal
==
future
.
metadataResponse
.
invalidationDone
.
incrementAndGet
())
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
future
.
metadataResponse
.
event
.
triggerResponse
();
}
}
catch
(
Exception
exx
)
{
exx
.
printStackTrace
();
}
for
(
var
stream
:
replicationStreams
)
}
void
processCloseReplicationStream
(
DaRPCStreamM
<
Request
,
Response
>
stream
)
{
stream
.
completedList
.
forEach
(
future
->
{
try
{
DaRPCFutureM
<
Request
,
Response
>
future
=
stream
.
processStream
();
if
(
future
!=
null
&&
future
.
metadataResponse
.
replicatonTotal
==
future
.
metadataResponse
.
replicationDone
.
incrementAndGet
())
{
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
if
(
future
.
metadataResponse
.
replicationDone
.
incrementAndGet
()
==
future
.
metadataResponse
.
replicatonTotal
)
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
{
try
{
future
.
metadataResponse
.
event
.
triggerResponse
();
}
catch
(
IOException
ex
)
{
logger
.
error
(
ex
);
}
}
}
catch
(
Exception
exx
)
});
stream
.
endpoint
.
pendingFutures
.
forEach
((
ticket
,
future
)->{
if
(
future
.
metadataResponse
.
replicationDone
.
incrementAndGet
()
==
future
.
metadataResponse
.
replicatonTotal
)
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
{
try
{
future
.
metadataResponse
.
event
.
triggerResponse
();
}
catch
(
IOException
ex
)
{
logger
.
error
(
ex
);
}
}
});
}
void
processReplicationStream
(
DaRPCStreamM
<
Request
,
Response
>
stream
)
{
try
{
DaRPCFutureM
<
Request
,
Response
>
future
=
stream
.
processStream
();
if
(
future
!=
null
)
{
exx
.
printStackTrace
();
if
(
future
.
metadataResponse
.
replicationDone
.
incrementAndGet
()
==
future
.
metadataResponse
.
replicatonTotal
)
if
(
Response
.
DONE
==
future
.
metadataResponse
.
state
.
incrementAndGet
())
future
.
metadataResponse
.
event
.
triggerResponse
();
}
}
catch
(
Exception
exx
)
{
exx
.
printStackTrace
();
}
}
public
void
run
()
{
NativeAffinity
.
setAffinity
(
1
<<
1
);
while
(
true
)
{
invalidationStreams
.
forEach
(
stream
->
{
if
(
stream
.
endpoint
.
isClosed
())
{
processCloseInvalidationStream
(
stream
);
invalidationStreams
.
remove
(
stream
);
}
else
{
processInvalidationStream
(
stream
);
}
});
replicationStreams
.
forEach
(
stream
->
{
if
(
!
stream
.
endpoint
.
isClosed
())
{
processCloseReplicationStream
(
stream
);
replicationStreams
.
remove
(
stream
);
}
else
{
processReplicationStream
(
stream
);
}
});
}
}
public
void
sendInvalidations
(
Request
request
,
Response
response
)
{
InvalidationRequest
inRequest
=
new
InvalidationRequest
();
//logger.debug("invdlidation send "+invalidationStreams.size());
inRequest
.
setKey
(
request
.
getKey
());
response
.
invalidationTotal
=
invalidationStreams
.
size
();
...
...
@@ -81,19 +170,25 @@ public class NetworkHandlerM implements Runnable{
{
try
{
stream
.
request
(
inRequest
,
new
InvalidationResponse
(),
response
,
fals
e
);
stream
.
request
(
inRequest
,
new
InvalidationResponse
(),
response
,
tru
e
);
}
catch
(
IOException
ex
)
{
System
.
out
.
println
(
ex
);
logger
.
error
(
ex
);
}
}
if
(
invalidationStreams
.
size
()
==
0
)
{
response
.
state
.
incrementAndGet
();
}
}
public
void
sendReplications
(
Request
request
,
Response
response
)
{
Request
repRequest
=
new
Request
();
logger
.
debug
(
"replication send "
+
replicationStreams
.
size
());
/*Request repRequest = new Request();
repRequest.setRequestType(request.getRequestType());
repRequest.setKey(request.getKey());
...
...
@@ -102,18 +197,24 @@ public class NetworkHandlerM implements Runnable{
repRequest.setValue(request.getValue());
}
*/
response
.
replicatonTotal
=
replicationStreams
.
size
();
for
(
var
stream
:
replicationStreams
)
{
try
{
stream
.
request
(
re
pRequest
,
new
Response
(),
response
,
fals
e
);
stream
.
request
(
re
quest
,
new
Response
(),
response
,
tru
e
);
}
catch
(
IOException
ex
)
{
System
.
out
.
println
(
ex
);
logger
.
error
(
ex
);
}
}
}
if
(
replicationStreams
.
size
()
==
0
)
{
response
.
state
.
incrementAndGet
();
}
}
}
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/ReplicationHandler.java
View file @
2b6c1a94
...
...
@@ -53,7 +53,7 @@ public class ReplicationHandler {
Response response = future.get().get();
}
stopWatch.stop();
logger.
info
("Replication future time: " + stopWatch.elapsed(TimeUnit.MICROSECONDS));
logger.
debug
("Replication future time: " + stopWatch.elapsed(TimeUnit.MICROSECONDS));
stopWatch.reset();*/
}
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/handler/StorageHandler.java
View file @
2b6c1a94
...
...
@@ -15,24 +15,23 @@ public class StorageHandler implements AutoCloseable{
final
static
Logger
logger
=
Logger
.
getLogger
(
StorageHandler
.
class
);
public
StorageHandler
(
String
dbpath
)
throws
RocksDBException
{
logger
.
debug
(
"Creating RocksDB"
);
logger
.
warn
(
"Creating RocksDB"
);
this
.
rockDbOptions
=
new
Options
();
rockDbOptions
.
setCreateIfMissing
(
true
);
this
.
rocksDB
=
RocksDB
.
open
(
rockDbOptions
,
dbpath
);
logger
.
debug
(
"Created RocksDB"
);
logger
.
warn
(
"Created RocksDB"
);
}
public
void
close
()
{
rocksDB
.
close
();
rockDbOptions
.
close
();
logger
.
debug
(
"Closing RocksDB instance"
);
logger
.
warn
(
"Closing RocksDB instance"
);
}
public
void
create
(
byte
[]
key
,
byte
[]
value
)
throws
RocksDBException
{
rocksDB
.
put
(
key
,
value
);
}
public
void
create
(
Request
request
)
throws
RocksDBException
{
logger
.
debug
(
"Writing to local"
);
rocksDB
.
put
(
request
.
getKey
(),
request
.
getValue
());
}
public
byte
[]
read
(
byte
[]
key
)
throws
RocksDBException
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/protocol/Request.java
View file @
2b6c1a94
...
...
@@ -16,7 +16,6 @@ public class Request implements DaRPCMessage
@Override
public
int
write
(
ByteBuffer
buffer
)
throws
IOException
{
// System.out.println("Request Write Method");
buffer
.
putInt
(
requestType
);
if
(
key
==
null
)
buffer
.
put
(
key
);
...
...
@@ -34,7 +33,6 @@ public class Request implements DaRPCMessage
@Override
public
void
update
(
ByteBuffer
buffer
)
throws
IOException
{
// System.out.println("Request update method");
requestType
=
buffer
.
getInt
();
if
(
key
==
null
||
key
.
length
!=
128
)
this
.
key
=
new
byte
[
128
];
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/protocol/Response.java
View file @
2b6c1a94
...
...
@@ -49,7 +49,6 @@ public class Response implements DaRPCMessage{
@Override
public
int
write
(
ByteBuffer
buffer
)
throws
IOException
{
// System.out.println("Response write Method");
buffer
.
putInt
(
ack
);
if
(
ack
==
AckType
.
SUCCESS_WITH_VALUE
)
{
...
...
@@ -62,7 +61,6 @@ public class Response implements DaRPCMessage{
@Override
public
void
update
(
ByteBuffer
buffer
)
throws
IOException
{
// System.out.println("esponse update method");
ack
=
buffer
.
getInt
();
if
(
ack
==
AckType
.
SUCCESS_WITH_VALUE
)
{
...
...
code/hpdos_rdma_offloaded/src/main/java/rdma/DaRPCClientEndpointM.java
→
code/hpdos_rdma_offloaded/src/main/java/
hpdos_rdma_offloaded/
rdma/DaRPCClientEndpointM.java
View file @
2b6c1a94
package
rdma
;
package
hpdos_rdma_offloaded.
rdma
;
/*
* DaRPC: Data Center Remote Procedure Call
*
...
...
@@ -39,7 +39,7 @@ import com.ibm.darpc.DaRPCEndpointGroup;
public
class
DaRPCClientEndpointM
<
R
extends
DaRPCMessage
,
T
extends
DaRPCMessage
>
extends
DaRPCEndpoint
<
R
,
T
>
{
private
static
final
Logger
logger
=
LoggerFactory
.
getLogger
(
"com.ibm.darpc"
);
p
rivate
ConcurrentHashMap
<
Integer
,
DaRPCFutureM
<
R
,
T
>>
pendingFutures
;
p
ublic
ConcurrentHashMap
<
Integer
,
DaRPCFutureM
<
R
,
T
>>
pendingFutures
;
private
AtomicInteger
ticketCount
;
private
int
streamCount
;
private
IbvWC
[]
wcList
;
...
...
@@ -86,7 +86,7 @@ public class DaRPCClientEndpointM<R extends DaRPCMessage, T extends DaRPCMessage
public
void
dispatchReceive
(
ByteBuffer
recvBuffer
,
int
ticket
,
int
recvIndex
)
throws
IOException
{
DaRPCFutureM
<
R
,
T
>
future
=
pendingFutures
.
get
(
ticket
);
if
(
future
==
null
){
logger
.
info
(
"no pending future (receive) for ticket "
+
ticket
);
logger
.
debug
(
"no pending future (receive) for ticket "
+
ticket
);
throw
new
IOException
(
"no pending future (receive) for ticket "
+
ticket
);
}
future
.
getReceiveMessage
().
update
(
recvBuffer
);
...
...
@@ -102,7 +102,7 @@ public class DaRPCClientEndpointM<R extends DaRPCMessage, T extends DaRPCMessage
public
void
dispatchSend
(
int
ticket
)
throws
IOException
{
DaRPCFutureM
<
R
,
T
>
future
=
pendingFutures
.
get
(
ticket
);
if
(
future
==
null
){
logger
.
info
(
"no pending future (send) for ticket "
+
ticket
);
logger
.
debug
(
"no pending future (send) for ticket "
+
ticket
);
throw
new
IOException
(
"no pending future (send) for ticket "
+
ticket
);
}
if
(
future
.
touch
()){
...
...
code/hpdos_rdma_offloaded/src/main/java/rdma/DaRPCClientGroupM.java
→
code/hpdos_rdma_offloaded/src/main/java/
hpdos_rdma_offloaded/
rdma/DaRPCClientGroupM.java
View file @
2b6c1a94
package
rdma
;
package
hpdos_rdma_offloaded.
rdma
;
/*
* DaRPC: Data Center Remote Procedure Call
...
...
code/hpdos_rdma_offloaded/src/main/java/rdma/DaRPCFutureM.java
→
code/hpdos_rdma_offloaded/src/main/java/
hpdos_rdma_offloaded/
rdma/DaRPCFutureM.java
View file @
2b6c1a94
package
rdma
;
package
hpdos_rdma_offloaded.
rdma
;
/*
* DaRPC: Data Center Remote Procedure Call
...
...
@@ -59,6 +59,7 @@ public class DaRPCFutureM<R extends DaRPCMessage, T extends DaRPCMessage> implem
this
.
response
=
response
;
this
.
metadataResponse
=
metadataResponse
;
this
.
stream
=
stream
;
this
.
endpoint
=
endpoint
;
this
.
status
=
new
AtomicInteger
(
RPC_PENDING
);
...
...
@@ -134,7 +135,7 @@ public class DaRPCFutureM<R extends DaRPCMessage, T extends DaRPCMessage> implem
endpoint
.
pollOnce
();
}
catch
(
Exception
e
){
status
.
set
(
RPC_ERROR
);
logger
.
info
(
e
.
getMessage
());
logger
.
debug
(
e
.
getMessage
());
}
}
return
status
.
get
()
>
0
;
...
...
code/hpdos_rdma_offloaded/src/main/java/rdma/DaRPCStreamM.java
→
code/hpdos_rdma_offloaded/src/main/java/
hpdos_rdma_offloaded/
rdma/DaRPCStreamM.java
View file @
2b6c1a94
package
rdma
;
package
hpdos_rdma_offloaded.
rdma
;
/*
* DaRPC: Data Center Remote Procedure Call
...
...
@@ -34,11 +34,11 @@ import com.ibm.darpc.DaRPCMessage;
public
class
DaRPCStreamM
<
R
extends
DaRPCMessage
,
T
extends
DaRPCMessage
>
{
private
static
final
Logger
logger
=
LoggerFactory
.
getLogger
(
"com.ibm.darpc"
);
p
rivate
DaRPCClientEndpointM
<
R
,
T
>
endpoint
;
p
rivate
LinkedBlockingDeque
<
DaRPCFutureM
<
R
,
T
>>
completedList
;
p
ublic
DaRPCClientEndpointM
<
R
,
T
>
endpoint
;
p
ublic
LinkedBlockingDeque
<
DaRPCFutureM
<
R
,
T
>>
completedList
;
DaRPCStreamM
(
DaRPCClientEndpointM
<
R
,
T
>
endpoint
,
int
streamId
)
throws
IOException
{
logger
.
info
(
"new direct rpc stream"
);
logger
.
debug
(
"new direct rpc stream"
);
this
.
endpoint
=
endpoint
;
this
.
completedList
=
new
LinkedBlockingDeque
<
DaRPCFutureM
<
R
,
T
>>();
}
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/service/FollowerMetadataService.java
View file @
2b6c1a94
...
...
@@ -33,17 +33,14 @@ public class FollowerMetadataService extends RpcProtocol implements DaRPCService
Request
request
=
event
.
getReceiveMessage
();
Response
response
=
event
.
getSendMessage
();
logger
.
info
(
"Received "
+
request
.
getRequestType
()+
" Request for key: "
+
new
String
(
request
.
getKey
()));
logger
.
debug
(
"Received "
+
request
.
getRequestType
()+
" Request for key: "
+
new
String
(
request
.
getKey
()));
if
(
request
.
getValue
()
!=
null
)
logger
.
info
(
" value : "
+
new
String
(
request
.
getValue
())
);
// System.out.println();
logger
.
debug
(
" value : "
+
new
String
(
request
.
getValue
())
);
try
{
if
(
request
.
getRequestType
()
==
RequestType
.
PUT
)
{
this
.
storageHandler
.
create
(
request
);
response
.
setAck
(
AckType
.
SUCCESS
);
response
.
setKey
(
null
);
response
.
setValue
(
null
);
}
else
if
(
request
.
getRequestType
()
==
RequestType
.
GET
)
{
byte
[]
value
=
this
.
storageHandler
.
read
(
request
);
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/service/MasterMetadataServerService.java
View file @
2b6c1a94
...
...
@@ -54,11 +54,11 @@ public class MasterMetadataServerService extends RpcProtocol implements DaRPCSer
Request
request
=
event
.
getReceiveMessage
();
Response
response
=
event
.
getSendMessage
();
Packet
packet
=
new
Packet
(
request
.
getRequestType
(),
request
.
getKey
(),
request
.
getValue
());
logger
.
info
(
"Received "
+
request
.
getRequestType
()+
" Request for key: "
+
new
String
(
request
.
getKey
()));
logger
.
debug
(
"Received "
+
request
.
getRequestType
()+
" Request for key: "
+
new
String
(
request
.
getKey
()));
Stopwatch
stopwatch
=
Stopwatch
.
createUnstarted
();
long
startTime
,
endTime
,
timeTaken
;
if
(
request
.
getValue
()
!=
null
)
logger
.
info
(
" value : "
+
new
String
(
request
.
getValue
()));
logger
.
debug
(
" value : "
+
new
String
(
request
.
getValue
()));
try
{
if
(
request
.
getRequestType
()
==
RequestType
.
PUT
)
{
startTime
=
System
.
nanoTime
();
...
...
@@ -66,27 +66,27 @@ public class MasterMetadataServerService extends RpcProtocol implements DaRPCSer
endTime
=
System
.
nanoTime
();
timeTaken
=
(
endTime
-
startTime
)
/
1000
;
ExperimentStatistics
.
collectStatistics
(
timeTaken
,
"localWriteTime"
);
logger
.
info
(
"Time to write to local cache: "
+
timeTaken
);
logger
.
debug
(
"Time to write to local cache: "
+
timeTaken
);
// stopwatch.reset();
// stopwatch.start();
startTime
=
System
.
nanoTime
();
this
.
networkHandler
.
create
(
packet
);
endTime
=
System
.
nanoTime
();
timeTaken
=
(
endTime
-
startTime
)
/
1000
;
logger
.
info
(
"Time to replicate: "
+
timeTaken
);
logger
.
debug
(
"Time to replicate: "
+
timeTaken
);
ExperimentStatistics
.
collectStatistics
(
timeTaken
,
"replicationTime"
);
// stopwatch.stop();
// logger.
info
("Time to replicate: " + stopwatch.elapsed(TimeUnit.MICROSECONDS));
// logger.
debug
("Time to replicate: " + stopwatch.elapsed(TimeUnit.MICROSECONDS));
// stopwatch.reset();
startTime
=
System
.
nanoTime
();
// this.networkHandler.sendInvalidateRequest(request);
endTime
=
System
.
nanoTime
();
timeTaken
=
(
endTime
-
startTime
)
/
1000
;
logger
.
info
(
"Time to invalidate: "
+
timeTaken
);
logger
.
debug
(
"Time to invalidate: "
+
timeTaken
);
ExperimentStatistics
.
collectStatistics
(
timeTaken
,
"invalidationTime"
);
ExperimentStatistics
.
collectStatistics
(
0
,
"write"
);
// stopwatch.reset();
// logger.
info
("Time to invalidate: " + stopwatch.elapsed(TimeUnit.MICROSECONDS));
// logger.
debug
("Time to invalidate: " + stopwatch.elapsed(TimeUnit.MICROSECONDS));
// stopwatch.reset();
//Yet to consume isDone
// var isDone = this.networkHandler.consumeReplicationRequestFutures(replicationRequestFutures);
...
...
@@ -99,7 +99,7 @@ public class MasterMetadataServerService extends RpcProtocol implements DaRPCSer
stopwatch
.
start
();
byte
[]
value
=
this
.
storageHandler
.
read
(
request
);
stopwatch
.
stop
();
logger
.
info
(
"Time to read: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Time to read: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopwatch
.
reset
();
if
(
value
==
null
)
{
...
...
@@ -115,17 +115,17 @@ public class MasterMetadataServerService extends RpcProtocol implements DaRPCSer
stopwatch
.
start
();
this
.
storageHandler
.
delete
(
request
);
stopwatch
.
stop
();
logger
.
info
(
"Time to delete to local cache: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Time to delete to local cache: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopwatch
.
reset
();
stopwatch
.
start
();
this
.
networkHandler
.
delete
(
packet
);
stopwatch
.
stop
();
logger
.
info
(
"Time to delete: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Time to delete: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopwatch
.
reset
();
stopwatch
.
start
();
// this.networkHandler.sendInvalidateRequest(request);
stopwatch
.
stop
();
logger
.
info
(
"Time to invalidate delete: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Time to invalidate delete: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopwatch
.
reset
();
// var replicationRequestFutures = this.networkHandler.sendReplicateRequest(request);
//Yet to consume isDone
...
...
code/hpdos_rdma_offloaded/src/main/java/hpdos_rdma_offloaded/service/MetadataServiceMaster.java
View file @
2b6c1a94
...
...
@@ -25,8 +25,7 @@ import hpdos_rdma_offloaded.protocol.Request;
import
hpdos_rdma_offloaded.protocol.RequestType
;
import
hpdos_rdma_offloaded.protocol.Response
;
import
hpdos_rdma_offloaded.protocol.RpcProtocol
;
import
rdma.DaRPCClientGroupM
;
import
hpdos_rdma_offloaded.rdma.DaRPCClientGroupM
;
public
class
MetadataServiceMaster
extends
RpcProtocol
implements
DaRPCService
<
Request
,
Response
>{
NetworkHandlerM
networkHandler
;
...
...
@@ -48,11 +47,11 @@ public class MetadataServiceMaster extends RpcProtocol implements DaRPCService<R
public
void
processServerEvent
(
DaRPCServerEvent
<
Request
,
Response
>
event
)
throws
IOException
{
Request
request
=
event
.
getReceiveMessage
();
Response
response
=
event
.
getSendMessage
();
logger
.
info
(
"Received "
+
request
.
getRequestType
()+
" Request for key: "
+
new
String
(
request
.
getKey
()));
logger
.
debug
(
"Received "
+
request
.
getRequestType
()+
" Request for key: "
+
new
String
(
request
.
getKey
()));
Stopwatch
stopwatch
=
Stopwatch
.
createUnstarted
();
long
startTime
,
endTime
,
timeTaken
;
if
(
request
.
getValue
()
!=
null
)
logger
.
info
(
" value : "
+
new
String
(
request
.
getValue
()));
logger
.
debug
(
" value : "
+
new
String
(
request
.
getValue
()));
try
{
if
(
request
.
getRequestType
()
==
RequestType
.
PUT
)
{
response
.
event
=
event
;
...
...
@@ -63,8 +62,8 @@ public class MetadataServiceMaster extends RpcProtocol implements DaRPCService<R
endTime
=
System
.
nanoTime
();
timeTaken
=
(
endTime
-
startTime
)
/
1000
;
response
.
state
.
getAndIncrement
();
//
ExperimentStatistics.collectStatistics(timeTaken, "localWriteTime");
logger
.
info
(
"Time to write to local cache: "
+
timeTaken
);
//
ExperimentStatistics.collectStatistics(timeTaken, "localWriteTime");
logger
.
debug
(
"Time to write to local cache: "
+
timeTaken
);
response
.
setAck
(
AckType
.
SUCCESS
);
response
.
setKey
(
null
);
...
...
@@ -74,7 +73,7 @@ public class MetadataServiceMaster extends RpcProtocol implements DaRPCService<R
stopwatch
.
start
();
byte
[]
value
=
this
.
storageHandler
.
read
(
request
);
stopwatch
.
stop
();
logger
.
info
(
"Time to read: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Time to read: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopwatch
.
reset
();
if
(
value
==
null
)
{
...
...
@@ -94,7 +93,7 @@ public class MetadataServiceMaster extends RpcProtocol implements DaRPCService<R
this
.
storageHandler
.
delete
(
request
);
stopwatch
.
stop
();
response
.
state
.
getAndIncrement
();
logger
.
info
(
"Time to delete to local cache: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
logger
.
debug
(
"Time to delete to local cache: "
+
stopwatch
.
elapsed
(
TimeUnit
.
MICROSECONDS
));
stopwatch
.
reset
();
response
.
event
=
event
;
...
...
@@ -119,18 +118,19 @@ public class MetadataServiceMaster extends RpcProtocol implements DaRPCService<R
ForkJoinPool
.
commonPool
().
submit
(()->{
try
{
System
.
out
.
println
(
rpcClientEndpoint
.
getDstAddr
()+
" "
+
rpcClientEndpoint
.
getSrcAddr
());
logger
.
debug
(
rpcClientEndpoint
.
getDstAddr
()+
" "
+
rpcClientEndpoint
.
getSrcAddr
());
String
clientIP
=
((
InetSocketAddress
)
rpcClientEndpoint
.
getDstAddr
()).
getHostName
();
InetSocketAddress
clientSocketAddress
=
new
InetSocketAddress
(
clientIP
,
1921
);
System
.
out
.
println
(
"Creating Endpoint"
);
logger
.
debug
(
"Creating Endpoint"
);
var
invalidationClientEp
=
invalidationClientGroup
.
createEndpoint
();
invalidationClientEp
.
connect
(
clientSocketAddress
,
100
);
System
.
out
.
println
(
"accepted"
);
invalidationClientEp
.
connect
(
clientSocketAddress
,
100
0
);
if
(
invalidationClientEp
.
isConnected
()){
logger
.
debug
(
"accepted"
);
this
.
networkHandler
.
addInvalidationStream
(
invalidationClientEp
.
createStream
());
}
else
{
System
.
out
.
println
(
"Not able to Connect to client for invalidation"
);
logger
.
error
(
"Not able to Connect to client for invalidation"
);
}
}
catch
(
Exception
e
){
...
...
@@ -143,8 +143,7 @@ public class MetadataServiceMaster extends RpcProtocol implements DaRPCService<R
public
void
close
(
DaRPCServerEndpoint
<
Request
,
Response
>
rpcClientEndpoint
)
{
try
{
System
.
out
.
println
(
"closed"
+((
InetSocketAddress
)
rpcClientEndpoint
.
getDstAddr
()).
getHostName
());
logger
.
debug
(
"closed"
+((
InetSocketAddress
)
rpcClientEndpoint
.
getDstAddr
()).
getHostName
());
}
catch
(
IOException
e
)
{
e
.
printStackTrace
();
}
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment