flushtest.java
来自「JGRoups源码」· Java 代码 · 共 572 行 · 第 1/2 页
JAVA
572 行
System.out.println("=== fetching the state ====");
c2.getState(null, 10000);
// Util.sleep(2000);
if(sendMessages){
c1.send(new Message());
c1.send(new Message());
c2.send(new Message());
c2.send(new Message());
}
checkNonStateTransferMemberSequence(receiver3);
checkBlockStateUnBlockSequence(receiver);
checkBlockStateUnBlockSequence(receiver2);
c2.close();
c2=null;
// Util.sleep(2000);
c2=createChannel();
receiver2=new MyReceiver(c2,"c2");
c2.setReceiver(receiver2);
c2.connect("bla");
// Util.sleep(2000);
if(sendMessages){
c1.send(new Message());
}
checkEventSequence(receiver);
checkEventSequence(receiver2);
checkEventSequence(receiver3);
System.out.println("=== fetching the state ====");
c3.getState(null, 10000);
if(sendMessages){
c1.send(new Message());
c2.send(new Message());
}
// Util.sleep(2000);
checkNonStateTransferMemberSequence(receiver2);
checkBlockStateUnBlockSequence(receiver);
checkBlockStateUnBlockSequence(receiver3);
}
private void checkBlockStateUnBlockSequence(MyReceiver receiver) {
List events = receiver.getEvents();
String name = receiver.getName();
assertNotNull(events);
assertEquals("Should have three events [block,get|setstate,unblock] but " + name + " has "
+ events, 3, events.size());
Object obj=events.remove(0);
assertTrue(name, obj instanceof BlockEvent);
obj=events.remove(0);
assertTrue(name, obj instanceof GetStateEvent || obj instanceof SetStateEvent);
obj=events.remove(0);
assertTrue(name, obj instanceof UnblockEvent);
receiver.clear();
}
private void checkEventSequence(MyReceiver receiver) {
List events = receiver.getEvents();
String eventString = "[" + receiver.getName() + ",events:" + events;
assertNotNull(events);
int size = events.size();
for (int i = 0; i < size; i++)
{
Object event = events.get(i);
if(event instanceof BlockEvent)
{
if(i+1<size)
{
assertTrue("After Block should be View " + eventString,events.get(i+1) instanceof View);
}
if(i!=0)
{
assertTrue("Before Block should be Unblock " + eventString,events.get(i-1) instanceof UnblockEvent);
}
}
if(event instanceof View)
{
if(i+1<size)
{
assertTrue("After View should be Unblock " + eventString,events.get(i+1) instanceof UnblockEvent);
}
assertTrue("Before View should be Block " +eventString,events.get(i-1) instanceof BlockEvent);
}
if(event instanceof UnblockEvent)
{
if(i+1<size)
{
assertTrue("After UnBlock should be Block " + eventString,events.get(i+1) instanceof BlockEvent);
}
assertTrue("Before UnBlock should be View "+eventString,events.get(i-1) instanceof View);
}
}
receiver.clear();
}
private void checkNonStateTransferMemberSequence(MyReceiver receiver) {
List events = receiver.getEvents();
assertNotNull(events);
assertEquals("Should have two events [block,unblock] but " + receiver.getName() + " has "
+ events, 2, events.size());
Object obj = events.remove(0);
assertTrue(obj instanceof BlockEvent);
obj = events.remove(0);
assertTrue(obj instanceof UnblockEvent);
receiver.clear();
}
private Channel createChannel() throws ChannelException {
Channel ret=new JChannel(CONFIG);
ret.setOpt(Channel.BLOCK, Boolean.TRUE);
Protocol flush=((JChannel)ret).getProtocolStack().findProtocol("FLUSH");
if(flush != null) {
Properties p=new Properties();
p.setProperty("timeout", "0");
flush.setProperties(p);
// send timeout up and down the stack, so other protocols can use the same value too
Map map = new HashMap();
map.put("flush_timeout", new Long(0));
flush.passUp(new Event(Event.CONFIG, map));
flush.passDown(new Event(Event.CONFIG, map));
}
return ret;
}
public static Test suite() {
return new TestSuite(FlushTest.class);
}
public static void main(String[] args) {
junit.textui.TestRunner.run(FlushTest.suite());
}
private static class MyReceiver extends ExtendedReceiverAdapter {
List events;
String name;
boolean verbose = true;
Channel channel = null;
public MyReceiver(Channel ch,String name) {
this.name=name;
channel = ch;
events=Collections.synchronizedList(new LinkedList());
}
public MyReceiver(String name) {
this.name=name;
events=Collections.synchronizedList(new LinkedList());
}
public String getLocalAddress()
{
String address = "";
if (channel != null)
{
address = channel.getLocalAddress().toString();
}
return address;
}
public String getName()
{
return name;
}
public void clear() {
events.clear();
}
public List getEvents() {return new LinkedList(events);}
public void block() {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() +"]: BLOCK");
events.add(new BlockEvent());
}
public void unblock() {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() + "]: UNBLOCK");
events.add(new UnblockEvent());
}
public void viewAccepted(View new_view) {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() + "]: " + new_view);
events.add(new_view);
}
public byte[] getState() {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() + "]: GetStateEvent");
events.add(new GetStateEvent(null, null));
return new byte[]{'b', 'e', 'l', 'a'};
}
public void setState(byte[] state) {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() + "]: SetStateEvent");
events.add(new SetStateEvent(null, null));
}
public void getState(OutputStream ostream) {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() + "]: GetStateEvent streamed");
events.add(new GetStateEvent(null, null));
byte [] payload = new byte[]{'b', 'e', 'l', 'a'};
try{
ostream.write(payload);
}
catch (IOException e){
e.printStackTrace();
}
finally{
Util.close(ostream);
}
}
public void setState(InputStream istream) {
if(verbose)
System.out.println("[" + name + ":" +getLocalAddress() + "]: SetStateEvent streamed");
events.add(new SetStateEvent(null, null));
byte [] payload = new byte[4];
try{
istream.read(payload);
}
catch (IOException e){
e.printStackTrace();
}
finally{
Util.close(istream);
}
}
}
private static class MySimpleReplier extends ExtendedReceiverAdapter {
Channel channel;
boolean handle_requests=false;
public MySimpleReplier(Channel channel, boolean handle_requests) {
this.channel=channel;
this.handle_requests=handle_requests;
}
public void receive(Message msg) {
Message reply=new Message(msg.getSrc());
try {
System.out.print("-- MySimpleReplier[" + channel.getLocalAddress() + "]: received message from " + msg.getSrc());
if(handle_requests) {
System.out.println(", sending reply");
channel.send(reply);
}
else
System.out.println("\n");
}
catch(Exception e) {
e.printStackTrace();
}
}
public void viewAccepted(View new_view) {
System.out.println("-- MySimpleReplier[" + channel.getLocalAddress() + "]: viewAccepted(" + new_view + ")");
}
public void block() {
System.out.println("-- MySimpleReplier[" + channel.getLocalAddress() + "]: block()");
}
public void unblock() {
System.out.println("-- MySimpleReplier[" + channel.getLocalAddress() + "]: unblock()");
}
}
}
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?