@@ -0,0 +1,200 @@
package com.njcn.message.mq ;
import com.alibaba.fastjson.JSON ;
import com.njcn.message.message.RecallMessage ;
import com.njcn.message.messagedto.MessageDataDTO ;
import com.njcn.middle.stream.autoconfig.RedisStreamAutoConfiguration ;
import com.njcn.mq.autoconfig.MqCoreAutoConfiguration ;
import com.njcn.mq.driver.redis.RedisStreamMqDriver ;
import com.njcn.mq.driver.redis.RedisStreamMqDriverAutoConfiguration ;
import com.njcn.mq.driver.rocketmq.RocketMqDriver ;
import com.njcn.mq.driver.rocketmq.RocketMqDriverAutoConfiguration ;
import com.njcn.mq.spi.MqDriver ;
import com.njcn.stat.api.MessAnalysisFeignClient ;
import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties ;
import org.apache.rocketmq.spring.core.RocketMQTemplate ;
import org.junit.jupiter.api.Assumptions ;
import org.junit.jupiter.api.Test ;
import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration ;
import org.springframework.boot.test.context.runner.ApplicationContextRunner ;
import org.springframework.data.redis.connection.stream.MapRecord ;
import org.springframework.data.redis.connection.stream.ReadOffset ;
import org.springframework.data.redis.connection.stream.StreamOffset ;
import org.springframework.data.redis.connection.stream.StreamReadOptions ;
import org.springframework.data.redis.connection.stream.StreamRecords ;
import org.springframework.data.redis.core.StringRedisTemplate ;
import java.net.InetSocketAddress ;
import java.net.Socket ;
import java.util.HashMap ;
import java.util.List ;
import java.util.Map ;
import static org.assertj.core.api.Assertions.assertThat ;
import static org.mockito.ArgumentMatchers.anyList ;
import static org.mockito.Mockito.mock ;
import static org.mockito.Mockito.timeout ;
import static org.mockito.Mockito.verify ;
/**
* MQ starter 切片接入验证测试( ApplicationContextRunner, 不启动全应用) 。
* 真实 Redis 要求: localhost:6379( docker redis-stream-it) 。
* 无 Redis 时,依赖 Redis 的用例通过 Assumptions 跳过。
*
* <p>注意: withConfiguration(AutoConfigurations.of(...)) 会对 autoconfig 类名做字母序排序,
* 导致 MqCoreAutoConfiguration 在 driver autoconfig 之前处理,@ConditionalOnBean(MqDriver.class)
* 因此失败。改用 withUserConfiguration(...) 按依赖顺序显式列出,避免排序问题。
*/
public class MqSliceWiringTest {
private static boolean redisAvailable ( ) {
try {
Socket s = new Socket ( ) ;
s . connect ( new InetSocketAddress ( " localhost " , 6379 ) , 500 ) ;
s . close ( ) ;
return true ;
} catch ( Exception e ) {
return false ;
}
}
/**
* redis-stream 模式的基础 runner。
* 按依赖顺序列出: Redis → Phase1 stream → redis driver → rocket driver(inactive) → mq core → 业务 bean
* 这样 @ConditionalOnBean(MqDriver.class) 在 MqCoreAutoConfiguration 被解析时已能找到 MqDriver 定义。
*/
private ApplicationContextRunner redisStreamRunner ( ) {
return new ApplicationContextRunner ( )
. withUserConfiguration (
RedisAutoConfiguration . class , // 1. StringRedisTemplate
RedisStreamAutoConfiguration . class , // 2. StreamKeyBuilder, RedisStreamProperties
RedisStreamMqDriverAutoConfiguration . class , // 3. MqDriver(RedisStreamMqDriver)
RocketMqDriverAutoConfiguration . class , // 4. inactive for redis-stream
MqCoreAutoConfiguration . class , // 5. MqTemplate, MqListenerRegistry
FrontDataMqListener . class , // 6. 上行 listener
RecallMqSender . class // 7. 下行 sender
)
. withBean ( MessAnalysisFeignClient . class ,
( ) - > mock ( MessAnalysisFeignClient . class ) )
. withPropertyValues (
" mq.type=redis-stream " ,
" spring.redis.host=localhost " ,
" spring.redis.port=6379 "
) ;
}
// ------------------------------------------------------------------
// 1. Driver 二选一
// ------------------------------------------------------------------
@Test
void testRedisStreamDriverSelected ( ) {
Assumptions . assumeTrue ( redisAvailable ( ) , " Redis not available at localhost:6379, skipping " ) ;
redisStreamRunner ( ) . run ( ctx - > {
assertThat ( ctx ) . hasSingleBean ( MqDriver . class ) ;
assertThat ( ctx . getBean ( MqDriver . class ) ) . isInstanceOf ( RedisStreamMqDriver . class ) ;
} ) ;
}
@Test
void testRocketMqDriverSelected ( ) {
// rocket 测试不依赖真实 Redis, 仅引入 rocket 相关 + core
new ApplicationContextRunner ( )
. withUserConfiguration (
RocketMqDriverAutoConfiguration . class , // 1. MqDriver(RocketMqDriver)
RedisStreamMqDriverAutoConfiguration . class , // 2. inactive for rocketmq
MqCoreAutoConfiguration . class , // 3. MqTemplate, MqListenerRegistry
FrontDataMqListener . class , // 4. inactive for rocketmq
RecallMqSender . class // 5. sender
)
. withBean ( MessAnalysisFeignClient . class ,
( ) - > mock ( MessAnalysisFeignClient . class ) )
. withBean ( RocketMQTemplate . class ,
( ) - > mock ( RocketMQTemplate . class ) )
. withBean ( RocketMQProperties . class , RocketMQProperties : : new )
. withPropertyValues ( " mq.type=rocketmq " )
. run ( ctx - > {
assertThat ( ctx ) . hasSingleBean ( MqDriver . class ) ;
assertThat ( ctx . getBean ( MqDriver . class ) ) . isInstanceOf ( RocketMqDriver . class ) ;
} ) ;
}
// ------------------------------------------------------------------
// 2. 上行(需要真实 Redis)
// 向 LN_Topic stream XADD 一条消息,验证 FrontDataMqListener 触发了
// messAnalysisFeignClient.analysis(...)
// ------------------------------------------------------------------
@Test
void testUpstreamFrontData ( ) {
Assumptions . assumeTrue ( redisAvailable ( ) , " Redis not available at localhost:6379, skipping " ) ;
redisStreamRunner ( ) . run ( ctx - > {
StringRedisTemplate redis = ctx . getBean ( StringRedisTemplate . class ) ;
MessAnalysisFeignClient mockClient = ctx . getBean ( MessAnalysisFeignClient . class ) ;
// 构造 RedisStreamMqDriver 消费时期望的字段格式
// enc=plain → body 即为原始 JSON( 参见 MessageCodec.encodeBody(json, false))
// stream key = StreamKeyBuilder.buildStream("LN_Topic") = "LN_Topic"
// (默认 envIsolation=false, 无 env 前缀)
MessageDataDTO dto = new MessageDataDTO ( ) ;
dto . setDataType ( 1 ) ;
String json = JSON . toJSONString ( dto ) ;
Map < String , String > fields = new HashMap < String , String > ( ) ;
fields . put ( " key " , " test-key-upstream-1 " ) ;
fields . put ( " tag " , " " ) ;
fields . put ( " enc " , " plain " ) ;
fields . put ( " body " , json ) ;
fields . put ( " ts " , String . valueOf ( System . currentTimeMillis ( ) ) ) ;
// XADD 到 LN_Topic( MqListenerRegistry.start() 已在 context 启动时订阅该 stream)
redis . opsForStream ( ) . add (
StreamRecords . newRecord ( ) . in ( " LN_Topic " ) . ofMap ( fields )
) ;
// 等待消费线程处理( blockMs 默认 2000ms, timeout=5000ms 足够)
verify ( mockClient , timeout ( 5000 ) ) . analysis ( anyList ( ) ) ;
} ) ;
}
// ------------------------------------------------------------------
// 3. 下行(需要真实 Redis)
// 通过 RecallMqSender 发送 RecallMessage, 验证 stream 内容正确
// ------------------------------------------------------------------
@Test
void testDownstreamRecall ( ) {
Assumptions . assumeTrue ( redisAvailable ( ) , " Redis not available at localhost:6379, skipping " ) ;
redisStreamRunner ( ) . run ( ctx - > {
StringRedisTemplate redis = ctx . getBean ( StringRedisTemplate . class ) ;
RecallMqSender sender = ctx . getBean ( RecallMqSender . class ) ;
// 清理,确保只有本次发送的消息
redis . delete ( " n1_recall_Topic " ) ;
RecallMessage msg = new RecallMessage ( ) ;
msg . setGuid ( " g1 " ) ;
// MqTemplate.send("n1_recall_Topic", msg) → RedisStreamMqDriver XADD
// stream key = StreamKeyBuilder.buildStream("n1_recall_Topic") = "n1_recall_Topic"
sender . send ( msg , " n1 " ) ;
Long size = redis . opsForStream ( ) . size ( " n1_recall_Topic " ) ;
assertThat ( size ) . isGreaterThanOrEqualTo ( 1L ) ;
@SuppressWarnings ( " unchecked " )
List < MapRecord < String , Object , Object > > records =
redis . opsForStream ( ) . read (
StreamReadOptions . empty ( ) . count ( 1L ) ,
StreamOffset . create ( " n1_recall_Topic " , ReadOffset . from ( " 0-0 " ) )
) ;
assertThat ( records ) . isNotEmpty ( ) ;
Map < Object , Object > firstFields = records . get ( 0 ) . getValue ( ) ;
// DefaultMqTemplate.send(topic, payload) 使用 encodeBody(json, false) → enc="plain"
// 因此 body 即为原始 JSON, 无需 gzip 解压
String body = ( String ) firstFields . get ( " body " ) ;
RecallMessage result = JSON . parseObject ( body , RecallMessage . class ) ;
assertThat ( result . getGuid ( ) ) . isEqualTo ( " g1 " ) ;
} ) ;
}
}