[CVS spice] Make sure that the netevent package works when multithreaded system

pdonald-yCVjj/[email protected] 17 May 2004 06:15:22 -0000
Newsgroups gmane.comp.java.spice.cvs
Message-ID <[email protected]>
Commit in spice/sandbox/netevent/src/test/org/codehaus/spice/netevent on MAIN

NetRuntime.java +2 -2 1.1 -> 1.2

TestEventHandler.java +44 -21 1.3 -> 1.4

TestServer.java +59 -102 1.12 -> 1.13

+105 -125

3 modified files

Make sure that the netevent package works when multithreaded system

----------

spice/sandbox/netevent/src/test/org/codehaus/spice/netevent

NetRuntime.java 1.1 -> 1.2

diff -u -r1.1 -r1.2
--- NetRuntime.java 17 May 2004 05:46:01 -0000 1.1
+++ NetRuntime.java 17 May 2004 06:15:21 -0000 1.2
@@ -16,7 +16,7 @@

/**
* @author Peter Donald

- * @version $Revision: 1.1 $ $Date: 2004/05/17 05:46:01 $

+ * @version $Revision: 1.2 $ $Date: 2004/05/17 06:15:21 $

*/
class NetRuntime
{

@@ -84,7 +84,7 @@

final DefaultEventQueue queue2 = new DefaultEventQueue( new UnboundedFifoBuffer( 15 ) );

final SelectableChannelEventSource source1 = new SelectableChannelEventSource( queue1 );

- source1.setSelectTimeout( 0 );

+ source1.setSelectTimeout( 200 );

final DefaultBufferManager bufferManager = new DefaultBufferManager();
final ChannelEventHandler ceh = new ChannelEventHandler( source1,

----------

spice/sandbox/netevent/src/test/org/codehaus/spice/netevent

TestEventHandler.java 1.3 -> 1.4

diff -u -r1.3 -r1.4
--- TestEventHandler.java 23 Jan 2004 04:26:21 -0000 1.3
+++ TestEventHandler.java 17 May 2004 06:15:21 -0000 1.4
@@ -15,7 +15,7 @@

/**
* @author Peter Donald

- * @version $Revision: 1.3 $ $Date: 2004/01/23 04:26:21 $

+ * @version $Revision: 1.4 $ $Date: 2004/05/17 06:15:21 $

*/
class TestEventHandler
extends AbstractEventHandler

@@ -28,11 +28,13 @@

private final long _transmitCount;
private final long _receiveCount;
private final boolean _closeOnReceive;

+ private int m_connectCount;
+ private int m_disconnectCount;

- public TestEventHandler( final String header,
- final long transmitCount,
- final long receiveCount,
- final boolean closeOnReceive )

+ TestEventHandler( final String header,
+ final long transmitCount,
+ final long receiveCount,
+ final boolean closeOnReceive )

{
_header = header;
_transmitCount = transmitCount;

@@ -40,6 +42,16 @@

_closeOnReceive = closeOnReceive;
}

+ int getConnectCount()
+ {
+ return m_connectCount;
+ }
+
+ int getDisconnectCount()
+ {
+ return m_disconnectCount;
+ }
+

/**
* @see EventHandler#handleEvent(Object)
*/

@@ -47,7 +59,7 @@

{
if( event instanceof InputDataPresentEvent )
{

- final InputDataPresentEvent e = (InputDataPresentEvent)event;

+ final InputDataPresentEvent e = (InputDataPresentEvent) event;

final ChannelTransport transport = e.getTransport();
final int available = transport.getInputStream().available();
if( available == _receiveCount && _closeOnReceive )

@@ -57,18 +69,32 @@

}
else if( event instanceof ChannelClosedEvent )
{

- final ChannelClosedEvent ce = (ChannelClosedEvent)event;

+ m_disconnectCount++;
+ final ChannelClosedEvent ce = (ChannelClosedEvent) event;

final ChannelTransport transport = ce.getTransport();

+ status( transport, "Channel Diconnected." );

receiveData( transport );
}
else if( event instanceof ConnectEvent )
{

- final ConnectEvent ce = (ConnectEvent)event;

+ m_connectCount++;
+ final ConnectEvent ce = (ConnectEvent) event;

final ChannelTransport transport = ce.getTransport();

+ setupUserData( transport );
+
+ status( transport, "Channel Connected." );

transmitData( transport );
}
}

+ private void status( final ChannelTransport transport, final String msg )
+ {
+ final int connected = m_connectCount - m_disconnectCount;
+ final String message = _header + " (" + transport.getUserData() + "): " +
+ ( msg + " Current=" + connected + " Connected= " + m_connectCount + " Disconnected=" + m_disconnectCount );
+ System.out.println( message );
+ }
+

private void receiveData( final ChannelTransport transport )
{
final MultiBufferInputStream in = transport.getInputStream();

@@ -80,7 +106,7 @@

{
for( int i = 0; i < count; i++ )
{

- sb.append( (char)in.read() );

+ sb.append( (char) in.read() );

}
}
catch( IOException e )

@@ -88,16 +114,11 @@

e.printStackTrace();
}

- output( transport, "Received " + available + " Sample: " + sb );

+ status( transport, "Received " + available + " Sample: " + sb );

}

private void transmitData( final ChannelTransport transport )
{

- final SocketChannel channel = (SocketChannel)transport.getChannel();
- final Socket socket = channel.socket();
- final String conn =
- socket.getLocalPort() + "<->" + socket.getPort();
- transport.setUserData( conn );

final OutputStream outputStream = transport.getOutputStream();
try
{

@@ -108,11 +129,11 @@

}
for( int i = 0; i < transmitCount; i++ )
{

- final byte ch = DATA[ i % DATA.length ];

+ final byte ch = DATA[i % DATA.length];

outputStream.write( ch );
}
outputStream.flush();

- output( transport, "Transmitting " + transmitCount );

+ status( transport, "Transmitting " + transmitCount );

}
catch( final IOException ioe )
{

@@ -120,10 +141,12 @@

}
}

- private void output( final ChannelTransport transport, final String text )

+ private void setupUserData( final ChannelTransport transport )

{

- final String message =
- _header + " (" + transport.getUserData() + "): " + text;
- System.out.println( message );

+ final SocketChannel channel = (SocketChannel) transport.getChannel();
+ final Socket socket = channel.socket();
+ final String conn = socket.getLocalPort() + "<->" + socket.getPort();
+ transport.setUserData( conn );

}

+

}

----------

spice/sandbox/netevent/src/test/org/codehaus/spice/netevent

TestServer.java 1.12 -> 1.13

diff -u -r1.12 -r1.13
--- TestServer.java 10 Feb 2004 00:54:50 -0000 1.12
+++ TestServer.java 17 May 2004 06:15:21 -0000 1.13
@@ -1,144 +1,101 @@

package org.codehaus.spice.netevent;

-import java.io.IOException;

import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.nio.channels.SelectionKey;
import java.nio.channels.ServerSocketChannel;
import java.nio.channels.SocketChannel;

-import java.util.ArrayList;
-import java.util.Arrays;
-import org.codehaus.spice.event.impl.DefaultEventQueue;

import org.codehaus.spice.event.impl.EventPump;

-import org.codehaus.spice.event.impl.collections.UnboundedFifoBuffer;
-import org.codehaus.spice.netevent.buffers.DefaultBufferManager;
-import org.codehaus.spice.netevent.handlers.ChannelEventHandler;
-import org.codehaus.spice.netevent.source.SelectableChannelEventSource;

/**
* @author Peter Donald

- * @version $Revision: 1.12 $ $Date: 2004/02/10 00:54:50 $

+ * @version $Revision: 1.13 $ $Date: 2004/05/17 06:15:21 $

*/
public class TestServer
{

+ private static final int MAX_CONNECTIONS = 500;
+ private static final int SERVER_PORT = 1980;
+

private static boolean c_done;

- private static SelectableChannelEventSource c_clientSocketSouce;

public static void main( final String[] args )
throws Exception
{

- final EventPump[] serverSidePumps = createServerSidePumps();
- final EventPump[] clientSidePumps = createClientSidePumps();
- final ArrayList pumpList = new ArrayList();
- pumpList.addAll( Arrays.asList( serverSidePumps ) );
- pumpList.addAll( Arrays.asList( clientSidePumps ) );
- final EventPump[] pumps =
- (EventPump[])pumpList.toArray( new EventPump[ pumpList.size() ] );

+ final NetRuntime server = NetRuntime.createRuntime( "SV", 5, -1, false );
+ final NetRuntime client = NetRuntime.createRuntime( "CL", -1, 5, true );

- final Runnable runnable = new Runnable()
- {
- public void run()
- {
- doPump( pumps );
- }
- };
- final Thread thread = new Thread( runnable );
- thread.start();

+ startPumps( server.getPumps() );
+ startPumps( client.getPumps() );
+
+ final TestEventHandler ch = client.getTestEventHandler();
+
+ final ServerSocketChannel ss = ServerSocketChannel.open();
+ server.getSource().registerChannel( ss, SelectionKey.OP_ACCEPT, null );
+ ss.socket().bind( new InetSocketAddress( SERVER_PORT ) );

+ int count = 0;

while( !c_done )
{
Thread.sleep( 50 );

- final SocketChannel channel = SocketChannel.open();
- c_clientSocketSouce.registerChannel( channel,
- SelectionKey.OP_CONNECT,
- null );
- final InetSocketAddress address =
- new InetSocketAddress( InetAddress.getLocalHost(), 1980 );
- channel.connect( address );

+ if( MAX_CONNECTIONS > count )
+ {
+ count++;
+ final SocketChannel channel = SocketChannel.open();
+ System.out.println( "Creating Client Connection: " + count );
+ client.getSource().registerChannel( channel, SelectionKey.OP_CONNECT, null );
+ final InetSocketAddress address = new InetSocketAddress(
+ InetAddress.getLocalHost(), SERVER_PORT );
+ channel.connect( address );
+ }
+
+ if( MAX_CONNECTIONS <= ch.getConnectCount() &&
+ ch.getConnectCount() == ch.getDisconnectCount() )
+ {
+ c_done = true;
+ }

}
System.exit( 1 );
}

- private static EventPump[] createServerSidePumps()
- throws IOException

+ private static void startPumps( final EventPump[] pumps )

{

- final DefaultEventQueue queue1 =
- new DefaultEventQueue( new UnboundedFifoBuffer( 15 ) );
- final DefaultEventQueue queue2 =
- new DefaultEventQueue( new UnboundedFifoBuffer( 15 ) );
-
- final SelectableChannelEventSource source1 =
- new SelectableChannelEventSource( queue1 );
-
- final ServerSocketChannel channel = ServerSocketChannel.open();
- source1.registerChannel( channel,
- SelectionKey.OP_ACCEPT,
- null );
- channel.socket().bind( new InetSocketAddress( 1980 ) );
-
- final DefaultBufferManager bufferManager =
- new DefaultBufferManager();
-
- final ChannelEventHandler handler1 =
- new ChannelEventHandler( source1, queue1, queue2, bufferManager );
-
- final TestEventHandler handler2 =
- new TestEventHandler( "SV", 5, -1, false );
-
- final EventPump pump1 = new EventPump( source1, handler1 );
- pump1.setBatchSize( 10 );
-
- final EventPump pump2 = new EventPump( queue2, handler2 );
- pump1.setBatchSize( 10 );
-
- return new EventPump[]{pump1, pump2};

+ for( int i = 0; i < pumps.length; i++ )
+ {
+ final EventPump pump = pumps[i];
+ startThread( pump );
+ }

}

- private static EventPump[] createClientSidePumps()
- throws IOException

+ private static void startThread( final EventPump pump )

{

- final DefaultEventQueue queue1 =
- new DefaultEventQueue( new UnboundedFifoBuffer( 15 ) );
- final DefaultEventQueue queue2 =
- new DefaultEventQueue( new UnboundedFifoBuffer( 15 ) );
-
- c_clientSocketSouce = new SelectableChannelEventSource( queue1 );
-
- final DefaultBufferManager bufferManager = new DefaultBufferManager();
-
- final ChannelEventHandler handler1 =
- new ChannelEventHandler( c_clientSocketSouce, queue1, queue2,
- bufferManager );
-
- final TestEventHandler handler2 =
- new TestEventHandler( "CL", -1, 5, true );
-
- final EventPump pump1 = new EventPump( c_clientSocketSouce, handler1 );
- pump1.setBatchSize( 10 );
-
- final EventPump pump2 = new EventPump( queue2, handler2 );
- pump1.setBatchSize( 10 );

+ final Runnable runnable = new Runnable()
+ {
+ public void run()
+ {
+ doPump( pump );
+ }
+ };
+ final Thread thread = new Thread( runnable );
+ thread.setName( pump.getName() );
+ thread.start();

- return new EventPump[]{pump1, pump2};

+ thread.setPriority( Thread.NORM_PRIORITY - 1 );

}

- private static void doPump( final EventPump[] pumps )

+ private static void doPump( final EventPump pump )

{

- for( int i = 0; i < 1000; i++ )

+ try

{

- for( int j = 0; j < pumps.length; j++ )

+ System.out.println( "Entering Thread " + Thread.currentThread().getName() );
+ while( !c_done )

{

- pumps[ j ].refresh();
- try
- {
- Thread.sleep( 2 );
- }
- catch( InterruptedException e )
- {
- e.printStackTrace();
- }

+ pump.refresh();

}

+ System.out.println( "Exiting Thread " + Thread.currentThread().getName() );
+ }
+ catch( final Throwable e )
+ {
+ e.printStackTrace();

}

- c_done = true;

}
}

CVSspam 0.2.8