[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