test[tcp]: 简化异步和同步测试用例

This commit is contained in:
jaysunxiao committed 2021-07-01 14:05:55 +08:00
1 parent c7d3fe1d69
commit b9d6f825f0
19 files changed
+110 -233

No files matched your search

@@ -304,9 +304,9 @@ public class PacketDispatcher implements IPacketDispatcher {
// 调用PacketReceiver
PacketBus.submit(session, packet, packetAttachment);
} catch (Exception e) {
logger.error(StringUtils.format("e[{}][{}]未知exception异常[e:{}]", session.getAttribute(AttributeType.UID), session.getSid(), e.getMessage()), e);
logger.error(StringUtils.format("e[uid:{}][sid:{}]未知exception异常", session.getAttribute(AttributeType.UID), session.getSid(), e.getMessage()), e);
} catch (Throwable t) {
logger.error(StringUtils.format("e[{}][{}]未知error错误[t:{}]", session.getAttribute(AttributeType.UID), session.getSid(), t.getMessage()), t);
logger.error(StringUtils.format("e[uid:{}][sid:{}]未知error错误", session.getAttribute(AttributeType.UID), session.getSid(), t.getMessage()), t);
} finally {
// 如果有服务器在处理同步或者异步消息的时候由于错误没有返回给客户端消息,则可能会残留serverAttachment,所以先移除
if (packetAttachment != null) {
@@ -29,7 +29,7 @@ import org.springframework.context.support.ClassPathXmlApplicationContext;
@Ignore
public class ServerTest {
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("server_config.xml");
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("config.xml");
@Test
public void startServer() {
@@ -15,16 +15,10 @@ package com.zfoo.net.core.tcp.client;
import com.zfoo.net.NetContext;
import com.zfoo.net.core.tcp.TcpClient;
import com.zfoo.net.packet.CM_AsyncMess0;
import com.zfoo.net.packet.CM_SyncMess;
import com.zfoo.net.packet.SM_AsyncMess0;
import com.zfoo.net.packet.SM_SyncMess;
import com.zfoo.net.packet.tcp.TcpHelloRequest;
import com.zfoo.net.packet.tcp.*;
import com.zfoo.net.session.SessionUtils;
import com.zfoo.protocol.exception.ExceptionUtils;
import com.zfoo.protocol.util.FileUtils;
import com.zfoo.protocol.util.JsonUtils;
import com.zfoo.protocol.util.StringUtils;
import com.zfoo.util.ThreadUtils;
import com.zfoo.util.net.HostAndPort;
import org.junit.Ignore;
@@ -33,8 +27,8 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;
/**
* @author jaysunxiao
@@ -42,12 +36,12 @@ import java.util.concurrent.Executors;
*/
@Ignore
public class TcpClientTest {
private static final Logger logger = LoggerFactory.getLogger(TcpClientTest.class);
private static final ExecutorService executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors() * 2);
@Test
public void startClientTest() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("client_config.xml");
public void startClient0() {
var context = new ClassPathXmlApplicationContext("config.xml");
SessionUtils.printSessionInfo();
var client = new TcpClient(HostAndPort.valueOf(NetContext.getConfigManager().getLocalConfig().getHostConfig().getAddressMap().get("server0")));
var session = client.start();
@@ -67,24 +61,24 @@ public class TcpClientTest {
@Test
public void syncClientTest() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("client_config.xml");
var context = new ClassPathXmlApplicationContext("config.xml");
SessionUtils.printSessionInfo();
var client = new TcpClient(HostAndPort.valueOf(NetContext.getConfigManager().getLocalConfig().getHostConfig().getAddressMap().get("server1")));
var session = client.start();
for (int i = 0; i < 10000; i++) {
int it = i;
var executorSize = Runtime.getRuntime().availableProcessors() * 2;
var executor = Executors.newFixedThreadPool(executorSize);
var atomicInteger = new AtomicInteger(0);
for (int i = 0; i < executorSize; i++) {
var thread = new Thread(() -> {
try {
CM_SyncMess cm = new CM_SyncMess();
cm.setA("Hello, this is client!");
cm.setId(it);
SM_SyncMess sm = NetContext.getDispatcher().syncAsk(session, cm, SM_SyncMess.class, null).packet();
var info = StringUtils.MULTIPLE_HYPHENS + FileUtils.LS
+ JsonUtils.object2String(cm) + FileUtils.LS
+ JsonUtils.object2String(sm) + FileUtils.LS;
System.out.println(info);
for (int j = 0; j < 10000; j++) {
var ask = new SyncMessAsk();
ask.setMessage("Hello, this is sync client!");
var answer = NetContext.getDispatcher().syncAsk(session, ask, SyncMessAnswer.class, null).packet();
logger.info("同步请求[{}]收到结果[{}]", atomicInteger.incrementAndGet(), JsonUtils.object2String(answer));
}
} catch (Exception e) {
logger.error(ExceptionUtils.getMessage(e));
}
@@ -97,30 +91,29 @@ public class TcpClientTest {
@Test
public void asyncClientTest() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("client_config.xml");
var client = new TcpClient(HostAndPort.valueOf(NetContext.getConfigManager().getLocalConfig().getHostConfig().getAddressMap().get("server0")));
var session = client.start();
var context = new ClassPathXmlApplicationContext("config.xml");
var client1 = new TcpClient(HostAndPort.valueOf(NetContext.getConfigManager().getLocalConfig().getHostConfig().getAddressMap().get("server1")));
var session1 = client1.start();
for (int i = 0; i < 10000; i++) {
Thread thread = new Thread(() -> {
try {
CM_AsyncMess0 cm = new CM_AsyncMess0();
cm.setA("Hello, client0 -> server0!");
var executorSize = Runtime.getRuntime().availableProcessors() * 2;
var executor = Executors.newFixedThreadPool(executorSize);
var asyncResponse = NetContext.getDispatcher().asyncAsk(session, cm, SM_AsyncMess0.class, null);
asyncResponse.whenComplete(sm -> {
var info = StringUtils.MULTIPLE_HYPHENS + FileUtils.LS
+ JsonUtils.object2String(cm) + FileUtils.LS
+ JsonUtils.object2String(sm) + FileUtils.LS;
System.out.println(info);
for (int i = 0; i < executorSize; i++) {
var thread = new Thread(() -> {
for (int j = 0; j < 1000; j++) {
var ask = new AsyncMess0Ask();
ask.setMessage("Hello, client0 -> server0!");
var answer = NetContext.getDispatcher().asyncAsk(session1, ask, AsyncMess0Answer.class, null);
answer.whenComplete(sm -> {
logger.info("异步请求收到结果[{}]", JsonUtils.object2String(answer));
}
);
} catch (Exception e) {
e.printStackTrace();
}
});
executor.submit(thread);
executor.execute(thread);
}
SessionUtils.printSessionInfo();
ThreadUtils.sleep(Long.MAX_VALUE);
@@ -14,13 +14,9 @@ package com.zfoo.net.core.tcp.server;
import com.zfoo.net.NetContext;
import com.zfoo.net.dispatcher.model.anno.PacketReceiver;
import com.zfoo.net.packet.*;
import com.zfoo.net.packet.tcp.TcpHelloRequest;
import com.zfoo.net.packet.tcp.TcpHelloResponse;
import com.zfoo.net.packet.tcp.*;
import com.zfoo.net.session.model.Session;
import com.zfoo.protocol.util.FileUtils;
import com.zfoo.protocol.util.JsonUtils;
import com.zfoo.protocol.util.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
@@ -39,20 +35,21 @@ public class TcpServerPacketController {
logger.info("receive [packet:{}] from client", JsonUtils.object2String(request));
var response = new TcpHelloResponse();
response.setMessage("Hello, this is the udp server!");
response.setMessage("Hello, this is the tcp server!");
NetContext.getDispatcher().send(session, response);
}
@PacketReceiver
public void atCM_SyncMess(Session session, CM_SyncMess cm) {
public void atSyncMessAsk(Session session, SyncMessAsk ask) {
logger.info("receive [packet:{}] from client", JsonUtils.object2String(ask));
// 测试超时
// ThreadUtils.sleep(Integer.MAX_VALUE);
// 测试正常返回
SM_SyncMess sm = new SM_SyncMess();
sm.setA("Hello, this is server!");
sm.setId(cm.getId());
var answer = new SyncMessAnswer();
answer.setMessage("Hello, this is sync server answer!");
// 测试返回不是预期的消息
// SM_Int sm = new SM_Int();
@@ -60,48 +57,35 @@ public class TcpServerPacketController {
// 测试错误返回
// var sm = ErrorResponse.valueOf(1, 1, "this is error response");
NetContext.getDispatcher().send(session, sm);
var info = StringUtils.MULTIPLE_HYPHENS + FileUtils.LS
+ JsonUtils.object2String(cm) + FileUtils.LS
+ JsonUtils.object2String(sm) + FileUtils.LS;
System.out.println(info);
NetContext.getDispatcher().send(session, answer);
}
// client0->server0->server1->server0->client0
@PacketReceiver
public void atCM_AsyncMess0(Session session, CM_AsyncMess0 cm0) {
CM_AsyncMess1 cm1 = new CM_AsyncMess1();
cm1.setA("Hello, server0 -> server1");
public void atAsyncMess0Ask(Session session, AsyncMess0Ask ask0) {
var ask1 = new AsyncMess1Ask();
ask1.setMessage("Hello, server0 -> server1");
var server1 = NetContext.getSessionManager().getClientSession(0L);
NetContext.getDispatcher().asyncAsk(server1, cm1, SM_AsyncMess1.class, null)
NetContext.getDispatcher().asyncAsk(server1, ask1, AsyncMess1Answer.class, null)
.whenComplete(sm_asyncMess0 -> {
SM_AsyncMess0 sm = new SM_AsyncMess0();
sm.setA("Hello, server0 -> client0!");
NetContext.getDispatcher().send(session, sm);
var info = StringUtils.MULTIPLE_HYPHENS + FileUtils.LS
+ JsonUtils.object2String(cm0) + FileUtils.LS
+ JsonUtils.object2String(sm) + FileUtils.LS;
System.out.println(info);
var answer = new AsyncMess0Answer();
answer.setMessage("Hello, server0 -> client0!");
NetContext.getDispatcher().send(session, answer);
});
}
@PacketReceiver
public void atCM_AsyncMess1(Session session, CM_AsyncMess1 cm) {
public void atAsyncMess1Ask(Session session, AsyncMess1Ask ask) {
// 测试超时
// ThreadUtils.sleep(Integer.MAX_VALUE);
// 测试正常返回
SM_AsyncMess1 sm = new SM_AsyncMess1();
sm.setA("Hello, server1 -> server0!");
var answer = new AsyncMess1Answer();
answer.setMessage("Hello, server1 -> server0!");
// 测试返回不是预期的消息
// SM_Int sm = new SM_Int();
@@ -109,13 +93,7 @@ public class TcpServerPacketController {
// 测试错误返回
// var sm = ErrorResponse.valueOf(1, 1, "this is error response");
NetContext.getDispatcher().send(session, sm);
var info = StringUtils.MULTIPLE_HYPHENS + FileUtils.LS
+ JsonUtils.object2String(cm) + FileUtils.LS
+ JsonUtils.object2String(sm) + FileUtils.LS;
System.out.println(info);
NetContext.getDispatcher().send(session, answer);
}
}
@@ -33,7 +33,7 @@ import java.util.concurrent.Executors;
@Ignore
public class TcpServerTest {
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("server_config.xml");
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("config.xml");
private static final ExecutorService executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());
/**
@@ -31,7 +31,7 @@ public class UdpClientTest {
@Test
public void startClientTest() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("client_config.xml");
var context = new ClassPathXmlApplicationContext("config.xml");
var hostAndPort = HostAndPort.valueOf(NetContext.getConfigManager().getLocalConfig().getHostConfig().getAddressMap().get("server0"));
var client = new UdpClient(hostAndPort);
@@ -29,7 +29,7 @@ public class UdpServerTest {
@Test
public void startServerTest() {
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("server_config.xml");
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("config.xml");
var server = new UdpServer(HostAndPort.valueOf(NetContext.getConfigManager().getLocalConfig().getHostConfig().getAddressMap().get("server0")));
server.start();
ThreadUtils.sleep(Long.MAX_VALUE);
@@ -31,7 +31,7 @@ import java.util.concurrent.Executors;
@Ignore
public class WebsocketServerTest {
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("server_config.xml");
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("config.xml");
private static final ExecutorService executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());
@@ -1,6 +1,5 @@
/*
* Copyright (C) 2020 The zfoo Authors
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
* in compliance with the License. You may obtain a copy of the License at
*
@@ -11,7 +10,7 @@
* See the License for the specific language governing permissions and limitations under the License.
*/
package com.zfoo.net.packet;
package com.zfoo.net.packet.tcp;
import com.zfoo.protocol.IPacket;
@@ -19,13 +18,12 @@ import com.zfoo.protocol.IPacket;
* @author jaysunxiao
* @version 3.0
*/
public class SM_AsyncMess0 implements IPacket {
public class AsyncMess0Answer implements IPacket {
public static final transient short PROTOCOL_ID = 1153;
private String a;
private String message;
private int id;
@Override
public short protocolId() {
@@ -33,20 +31,12 @@ public class SM_AsyncMess0 implements IPacket {
}
public String getA() {
return a;
public String getMessage() {
return message;
}
public void setA(String a) {
this.a = a;
public void setMessage(String message) {
this.message = message;
}
public int getId() {
return id;
}
public void setId(int id) {
this.id = id;
}
}
@@ -1,6 +1,5 @@
/*
* Copyright (C) 2020 The zfoo Authors
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
* in compliance with the License. You may obtain a copy of the License at
*
@@ -11,7 +10,7 @@
* See the License for the specific language governing permissions and limitations under the License.
*/
package com.zfoo.net.packet;
package com.zfoo.net.packet.tcp;
import com.zfoo.protocol.IPacket;
@@ -19,32 +18,23 @@ import com.zfoo.protocol.IPacket;
* @author jaysunxiao
* @version 3.0
*/
public class CM_AsyncMess0 implements IPacket {
public class AsyncMess0Ask implements IPacket {
public static final transient short PROTOCOL_ID = 1152;
private String a;
private int id;
private String message;
@Override
public short protocolId() {
return PROTOCOL_ID;
}
public String getA() {
return a;
public String getMessage() {
return message;
}
public void setA(String a) {
this.a = a;
public void setMessage(String message) {
this.message = message;
}
public int getId() {
return id;
}
public void setId(int id) {
this.id = id;
}
}
@@ -1,6 +1,5 @@
/*
* Copyright (C) 2020 The zfoo Authors
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
* in compliance with the License. You may obtain a copy of the License at
*
@@ -11,7 +10,7 @@
* See the License for the specific language governing permissions and limitations under the License.
*/
package com.zfoo.net.packet;
package com.zfoo.net.packet.tcp;
import com.zfoo.protocol.IPacket;
@@ -19,13 +18,11 @@ import com.zfoo.protocol.IPacket;
* @author jaysunxiao
* @version 3.0
*/
public class SM_AsyncMess1 implements IPacket {
public class AsyncMess1Answer implements IPacket {
public static final transient short PROTOCOL_ID = 1155;
private String a;
private int id;
private String message;
@Override
public short protocolId() {
@@ -33,20 +30,12 @@ public class SM_AsyncMess1 implements IPacket {
}
public String getA() {
return a;
public String getMessage() {
return message;
}
public void setA(String a) {
this.a = a;
public void setMessage(String message) {
this.message = message;
}
public int getId() {
return id;
}
public void setId(int id) {
this.id = id;
}
}
@@ -1,6 +1,5 @@
/*
* Copyright (C) 2020 The zfoo Authors
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
* in compliance with the License. You may obtain a copy of the License at
*
@@ -11,7 +10,7 @@
* See the License for the specific language governing permissions and limitations under the License.
*/
package com.zfoo.net.packet;
package com.zfoo.net.packet.tcp;
import com.zfoo.protocol.IPacket;
@@ -19,32 +18,22 @@ import com.zfoo.protocol.IPacket;
* @author jaysunxiao
* @version 3.0
*/
public class CM_AsyncMess1 implements IPacket {
public class AsyncMess1Ask implements IPacket {
public static final transient short PROTOCOL_ID = 1154;
private String a;
private int id;
private String message;
@Override
public short protocolId() {
return PROTOCOL_ID;
}
public String getA() {
return a;
public String getMessage() {
return message;
}
public void setA(String a) {
this.a = a;
}
public int getId() {
return id;
}
public void setId(int id) {
this.id = id;
public void setMessage(String message) {
this.message = message;
}
}
@@ -1,6 +1,5 @@
/*
* Copyright (C) 2020 The zfoo Authors
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
* in compliance with the License. You may obtain a copy of the License at
*
@@ -11,7 +10,7 @@
* See the License for the specific language governing permissions and limitations under the License.
*/
package com.zfoo.net.packet;
package com.zfoo.net.packet.tcp;
import com.zfoo.protocol.IPacket;
@@ -19,13 +18,11 @@ import com.zfoo.protocol.IPacket;
* @author jaysunxiao
* @version 3.0
*/
public class SM_SyncMess implements IPacket {
public class SyncMessAnswer implements IPacket {
public static final transient short PROTOCOL_ID = 1151;
private String a;
private int id;
private String message;
@Override
public short protocolId() {
@@ -33,20 +30,12 @@ public class SM_SyncMess implements IPacket {
}
public String getA() {
return a;
public String getMessage() {
return message;
}
public void setA(String a) {
this.a = a;
public void setMessage(String message) {
this.message = message;
}
public int getId() {
return id;
}
public void setId(int id) {
this.id = id;
}
}
@@ -1,6 +1,5 @@
/*
* Copyright (C) 2020 The zfoo Authors
*
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except
* in compliance with the License. You may obtain a copy of the License at
*
@@ -11,7 +10,7 @@
* See the License for the specific language governing permissions and limitations under the License.
*/
package com.zfoo.net.packet;
package com.zfoo.net.packet.tcp;
import com.zfoo.protocol.IPacket;
@@ -19,32 +18,23 @@ import com.zfoo.protocol.IPacket;
* @author jaysunxiao
* @version 3.0
*/
public class CM_SyncMess implements IPacket {
public class SyncMessAsk implements IPacket {
public static final transient short PROTOCOL_ID = 1150;
private String a;
private int id;
private String message;
@Override
public short protocolId() {
return PROTOCOL_ID;
}
public String getA() {
return a;
public String getMessage() {
return message;
}
public void setA(String a) {
this.a = a;
public void setMessage(String message) {
this.message = message;
}
public int getId() {
return id;
}
public void setId(int id) {
this.id = id;
}
}
@@ -36,7 +36,7 @@ import java.util.Set;
*/
public class ProtocolTest {
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("server_config.xml");
private static final ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("config.xml");
private static final IPacketService packetService = NetContext.getPacketService();
private static SignalPacketAttachment attachment = new SignalPacketAttachment();
-31
View File
@@ -1,31 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:context="http://www.springframework.org/schema/context"
xmlns:net="http://www.zfoo.com/schema/net"
xsi:schemaLocation="
http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans-4.0.xsd
http://www.springframework.org/schema/context
http://www.springframework.org/schema/context/spring-context-4.0.xsd
http://www.zfoo.com/schema/net
http://www.zfoo.com/schema/net-1.0.xsd">
<context:component-scan base-package="com.zfoo"/>
<net:config id="config" protocol-location="protocol.xml">
<net:host center="direct connect" user="jaysunxiao" password="123456">
<net:address name="server0" url="127.0.0.1:9000"/>
<net:address name="server1" url="127.0.0.1:9001"/>
<net:address name="client0" url="127.0.0.1:9000"/>
<net:address name="client1" url="127.0.0.1:9001"/>
</net:host>
</net:config>
</beans>
File renamed without changes.
+6 -6
View File
@@ -48,12 +48,12 @@
<protocol id="1121" location="com.zfoo.net.packet.CM_Set" enhance="false"/>
<protocol id="1150" location="com.zfoo.net.packet.CM_SyncMess" enhance="false"/>
<protocol id="1151" location="com.zfoo.net.packet.SM_SyncMess" enhance="false"/>
<protocol id="1152" location="com.zfoo.net.packet.CM_AsyncMess0" enhance="false"/>
<protocol id="1153" location="com.zfoo.net.packet.SM_AsyncMess0" enhance="false"/>
<protocol id="1154" location="com.zfoo.net.packet.CM_AsyncMess1" enhance="false"/>
<protocol id="1155" location="com.zfoo.net.packet.SM_AsyncMess1" enhance="false"/>
<protocol id="1150" location="com.zfoo.net.packet.tcp.SyncMessAsk" enhance="false"/>
<protocol id="1151" location="com.zfoo.net.packet.tcp.SyncMessAnswer" enhance="false"/>
<protocol id="1152" location="com.zfoo.net.packet.tcp.AsyncMess0Ask" enhance="false"/>
<protocol id="1153" location="com.zfoo.net.packet.tcp.AsyncMess0Answer" enhance="false"/>
<protocol id="1154" location="com.zfoo.net.packet.tcp.AsyncMess1Ask" enhance="false"/>
<protocol id="1155" location="com.zfoo.net.packet.tcp.AsyncMess1Answer" enhance="false"/>
<protocol id="1165" location="com.zfoo.net.packet.csharp.CM_CSharpRequest" enhance="false"/>
<protocol id="1166" location="com.zfoo.net.packet.csharp.CSharpObjectA" enhance="false"/>
@@ -228,9 +228,9 @@ public class EntityCaches<PK extends Comparable<PK>, E extends IEntity<PK>> impl
updateList.clear();
} catch (Exception e) {
logger.error("数据库持久化器[{}]的持久化过程中exception异常退出[e:{}]", entityDef.getClazz().getSimpleName(), e);
logger.error("数据库持久化器[{}]的持久化过程中exception异常退出", entityDef.getClazz().getSimpleName(), e);
} catch (Throwable t) {
logger.error("数据库持久化器[{}]的持久化过程中throwable异常退出[t:{}]", entityDef.getClazz().getSimpleName(), t);
logger.error("数据库持久化器[{}]的持久化过程中throwable异常退出", entityDef.getClazz().getSimpleName(), t);
} finally {
}
}