Commit aaec213e authored by wuyuhang's avatar wuyuhang

Initial commit

parents
Pipeline #4369 failed with stages
target/
!.mvn/wrapper/maven-wrapper.jar
!**/src/main/**/target/
!**/src/test/**/target/
### IntelliJ IDEA ###
.idea/modules.xml
.idea/jarRepositories.xml
.idea/compiler.xml
.idea/libraries/
*.iws
*.iml
*.ipr
### Eclipse ###
.apt_generated
.classpath
.factorypath
.project
.settings
.springBeans
.sts4-cache
### NetBeans ###
/nbproject/private/
/nbbuild/
/dist/
/nbdist/
/.nb-gradle/
build/
!**/src/main/**/build/
!**/src/test/**/build/
### VS Code ###
.vscode/
### Mac OS ###
.DS_Store
\ No newline at end of file
# 默认忽略的文件
/shelf/
/workspace.xml
# 基于编辑器的 HTTP 客户端请求
/httpRequests/
# 依赖于环境的 Maven 主目录路径
/mavenHomeManager.xml
# Datasource local storage ignored files
/dataSources/
/dataSources.local.xml
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="Encoding">
<file url="file://$PROJECT_DIR$/src/main/java" charset="UTF-8" />
<file url="file://$PROJECT_DIR$/src/main/resources" charset="UTF-8" />
</component>
</project>
\ No newline at end of file
<component name="InspectionProjectProfileManager">
<profile version="1.0">
<option name="myName" value="Project Default" />
<inspection_tool class="AliAccessStaticViaInstance" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliArrayNamingShouldHaveBracket" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliControlFlowStatementWithoutBraces" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliDeprecation" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliEqualsAvoidNull" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliLongLiteralsEndingWithLowercaseL" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliMissingOverrideAnnotation" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="AliWrapperTypeEquality" enabled="true" level="WARNING" enabled_by_default="true" />
<inspection_tool class="MapOrSetKeyShouldOverrideHashCodeEquals" enabled="true" level="WARNING" enabled_by_default="true" />
</profile>
</component>
\ No newline at end of file
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="ExternalStorageConfigurationManager" enabled="true" />
<component name="MavenProjectsManager">
<option name="originalFiles">
<list>
<option value="$PROJECT_DIR$/pom.xml" />
</list>
</option>
</component>
<component name="ProjectRootManager" version="2" languageLevel="JDK_1_8" default="true" project-jdk-name="1.8" project-jdk-type="JavaSDK">
<output url="file://$PROJECT_DIR$/out" />
</component>
</project>
\ No newline at end of file
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="VcsDirectoryMappings">
<mapping directory="$PROJECT_DIR$/.." vcs="Git" />
<mapping directory="$PROJECT_DIR$" vcs="Git" />
</component>
</project>
\ No newline at end of file
import http.server
import threading
import json
import random
import math
PORT = 8765
HTML = r"""<!DOCTYPE html>
<html><head><meta charset="utf-8"><title>Particle Life</title>
<style>
* { margin:0; padding:0; box-sizing:border-box; }
body { background:#0a0a0f; overflow:hidden; font-family:'Courier New',monospace; }
canvas { display:block; }
#ui { position:fixed; top:16px; left:16px; color:#aaa; font-size:13px; z-index:10; }
#ui h1 { font-size:16px; color:#e0e0e0; margin-bottom:8px; letter-spacing:1px; }
#ui p { margin:2px 0; opacity:0.7; }
#stats { position:fixed; bottom:16px; left:16px; color:#555; font-size:11px; }
#matrix { position:fixed; top:16px; right:16px; }
#matrix canvas { border:1px solid #222; border-radius:4px; }
button { background:#1a1a2e; color:#aaa; border:1px solid #333; padding:4px 12px;
cursor:pointer; border-radius:3px; font-family:inherit; font-size:12px; margin-top:8px; }
button:hover { background:#2a2a4e; color:#ddd; }
</style></head><body>
<div id="ui">
<h1>PARTICLE LIFE</h1>
<p>Particles: <span id="pcount">600</span></p>
<p>Species: <span id="scount">6</span></p>
<button onclick="randomize()">Randomize Rules</button>
<button onclick="reset()">Reset</button>
</div>
<div id="matrix"><canvas id="mcanvas" width="120" height="120"></canvas></div>
<div id="stats"><span id="fps">0</span> fps</div>
<canvas id="canvas"></canvas>
<script>
const canvas = document.getElementById('canvas');
const ctx = canvas.getContext('2d');
const mcanvas = document.getElementById('mcanvas');
const mctx = mcanvas.getContext('2d');
let W, H;
function resize() { W = canvas.width = innerWidth; H = canvas.height = innerHeight; }
resize();
addEventListener('resize', resize);
const NUM_SPECIES = 6;
const NUM_PARTICLES = 600;
const RMAX = 80;
const BETA = 0.3;
const FRICTION = 0.15;
const DT = 0.02;
const COLORS = [
[255, 80, 80], [80, 200, 255], [80, 255, 120],
[255, 200, 60], [200, 80, 255], [255, 120, 200]
];
let attraction = [];
let particles = [];
function randomize() {
attraction = [];
for (let i = 0; i < NUM_SPECIES; i++) {
attraction[i] = [];
for (let j = 0; j < NUM_SPECIES; j++) {
attraction[i][j] = (Math.random() * 2 - 1);
}
}
drawMatrix();
}
function initParticles() {
particles = [];
for (let i = 0; i < NUM_PARTICLES; i++) {
particles.push({
x: Math.random() * W, y: Math.random() * H,
vx: 0, vy: 0,
species: Math.floor(Math.random() * NUM_SPECIES)
});
}
}
function reset() { randomize(); initParticles(); }
function force(r, a) {
if (r < BETA) return r / BETA - 1;
if (r < 1) return a * (1 - Math.abs(2 * r - 1 - BETA) / (1 - BETA));
return 0;
}
// Spatial hash for performance
const CELL = RMAX;
let grid = {};
function hashKey(x, y) { return (Math.floor(x/CELL)) + ',' + (Math.floor(y/CELL)); }
function buildGrid() {
grid = {};
for (let i = 0; i < particles.length; i++) {
const k = hashKey(particles[i].x, particles[i].y);
if (!grid[k]) grid[k] = [];
grid[k].push(i);
}
}
function neighbors(x, y) {
const cx = Math.floor(x/CELL), cy = Math.floor(y/CELL);
const result = [];
for (let dx = -1; dx <= 1; dx++)
for (let dy = -1; dy <= 1; dy++) {
const k = (cx+dx) + ',' + (cy+dy);
if (grid[k]) for (const idx of grid[k]) result.push(idx);
}
return result;
}
function update() {
buildGrid();
for (let i = 0; i < particles.length; i++) {
const p = particles[i];
let fx = 0, fy = 0;
const nbrs = neighbors(p.x, p.y);
for (const j of nbrs) {
if (i === j) continue;
const q = particles[j];
let dx = q.x - p.x, dy = q.y - p.y;
// Wrap around
if (dx > W/2) dx -= W; if (dx < -W/2) dx += W;
if (dy > H/2) dy -= H; if (dy < -H/2) dy += H;
const dist = Math.sqrt(dx*dx + dy*dy);
if (dist > 0 && dist < RMAX) {
const r = dist / RMAX;
const a = attraction[p.species][q.species];
const f = force(r, a);
fx += (dx / dist) * f;
fy += (dy / dist) * f;
}
}
p.vx = (p.vx + fx * DT * RMAX) * (1 - FRICTION);
p.vy = (p.vy + fy * DT * RMAX) * (1 - FRICTION);
p.x = ((p.x + p.vx) % W + W) % W;
p.y = ((p.y + p.vy) % H + H) % H;
}
}
function drawMatrix() {
const s = 120 / NUM_SPECIES;
mctx.clearRect(0, 0, 120, 120);
for (let i = 0; i < NUM_SPECIES; i++)
for (let j = 0; j < NUM_SPECIES; j++) {
const a = attraction[i][j];
const c = a > 0 ? `rgba(80,255,120,${a})` : `rgba(255,80,80,${-a})`;
mctx.fillStyle = c;
mctx.fillRect(j * s, i * s, s, s);
}
}
let lastTime = performance.now(), frameCount = 0;
function draw() {
ctx.fillStyle = 'rgba(10,10,15,0.25)';
ctx.fillRect(0, 0, W, H);
for (const p of particles) {
const c = COLORS[p.species];
const speed = Math.sqrt(p.vx*p.vx + p.vy*p.vy);
const alpha = Math.min(1, 0.4 + speed * 0.3);
ctx.fillStyle = `rgba(${c[0]},${c[1]},${c[2]},${alpha})`;
const r = 2 + speed * 0.5;
ctx.beginPath();
ctx.arc(p.x, p.y, r, 0, Math.PI * 2);
ctx.fill();
}
frameCount++;
const now = performance.now();
if (now - lastTime > 1000) {
document.getElementById('fps').textContent =
Math.round(frameCount * 1000 / (now - lastTime));
frameCount = 0; lastTime = now;
}
update();
requestAnimationFrame(draw);
}
randomize();
initParticles();
draw();
</script></body></html>"""
class Handler(http.server.BaseHTTPRequestHandler):
def do_GET(self):
self.send_response(200)
self.send_header('Content-Type', 'text/html; charset=utf-8')
self.end_headers()
self.wfile.write(HTML.encode())
def log_message(self, *a):
pass
print(f"Serving Particle Life on http://localhost:{PORT}")
http.server.HTTPServer(('127.0.0.1', PORT), Handler).serve_forever()
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.monitor</groupId>
<artifactId>d2d-monitor</artifactId>
<version>1.0.0</version>
<packaging>jar</packaging>
<name>D2D Monitor Simulator</name>
<description>数据传输管理软件监控接口测试程序</description>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.7.18</version>
<relativePath/>
</parent>
<properties>
<java.version>1.8</java.version>
<netty.version>4.1.115.Final</netty.version>
</properties>
<dependencies>
<!-- Spring Boot Web -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- Netty - UDP Multicast -->
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-all</artifactId>
<version>${netty.version}</version>
</dependency>
<!-- Lombok -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!-- Jackson JSON -->
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
</plugin>
</plugins>
</build>
</project>
\ No newline at end of file
package com.jk;
//TIP 要<b>运行</b>代码,请按 <shortcut actionId="Run"/> 或
// 点击装订区域中的 <icon src="AllIcons.Actions.Execute"/> 图标。
public class Main {
public static void main(String[] args) {
//TIP 当文本光标位于高亮显示的文本处时按 <shortcut actionId="ShowIntentionActions"/>
// 查看 IntelliJ IDEA 建议如何修正。
System.out.printf("Hello and welcome!");
for (int i = 1; i <= 5; i++) {
//TIP 按 <shortcut actionId="Debug"/> 开始调试代码。我们已经设置了一个 <icon src="AllIcons.Debugger.Db_set_breakpoint"/> 断点
// 但您始终可以通过按 <shortcut actionId="ToggleLineBreakpoint"/> 添加更多断点。
System.out.println("i = " + i);
}
}
}
\ No newline at end of file
package com.monitor.d2d;
import com.monitor.d2d.netty.NettyTcpPushServer;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableScheduling;
@SpringBootApplication
@EnableScheduling
public class D2dMonitorApplication {
public static void main(String[] args) {
/*Thread thread = new Thread(() -> {
NettyTcpPushServer server = new NettyTcpPushServer();
try {
server.main(null);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
});
thread.start();*/
SpringApplication.run(D2dMonitorApplication.class, args);
}
}
\ No newline at end of file
package com.monitor.d2d.config;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
@Data
@Component
@ConfigurationProperties(prefix = "d2d")
public class AppProperties {
/** 监控软件信源地址 SID(十六进制,如 11651101) */
// private String sid = "11651101";
//
// /** 数传软件信宿地址 DID(十六进制,如 26000004) */
// private String did = "26000004";
//
// /** 命令发送目标组播地址 */
// private String sendHost = "238.1.0.103";
//
// /** 命令发送目标端口 */
// private int sendPort = 8889;
//
// /** 上报接收监听组播组 */
// private String recvGroup = "238.1.0.3";
//
// /** 上报接收监听端口 */
// private int recvPort = 8888;
// 244
private String sid = "26125010";
private String did = "26045002";
private String sendHost = "232.18.185.21";
private int sendPort = 8888;
private String recvGroup = "232.18.185.11";
private int recvPort = 8888;
/** 网络接口名称,留空则自动选择第一个支持组播的接口 ens192*/
private String networkInterface = "wlan1";
/** 状态查询间隔(秒),0=不自动查询 */
private int queryInterval = 10;
/** 解析SID为long */
public long parseSid() {
return Long.parseUnsignedLong(sid, 16);
}
/** 解析DID为long */
public long parseDid() {
return Long.parseUnsignedLong(did, 16);
}
}
\ No newline at end of file
package com.monitor.d2d.netty;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.handler.logging.LogLevel;
import io.netty.handler.logging.LoggingHandler;
import io.netty.handler.timeout.IdleState;
import io.netty.handler.timeout.IdleStateEvent;
import io.netty.handler.timeout.IdleStateHandler;
import io.netty.handler.traffic.ChannelTrafficShapingHandler;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
/**
* Netty高性能TCP推送服务端
* 客户端连接后全速循环下发固定十六进制报文,无连接时不发包、不消耗CPU
*/
public class NettyTcpPushServer {
// ===================== 配置区 =====================
private static final int LISTEN_PORT = 27586;
// Netty 发送缓冲区高低水位,控制推送上限
private static final int HIGH_WATER_MARK = 1024 * 1024 * 8;
private static final int LOW_WATER_MARK = 1024 * 1024 * 4;
// 固定下发报文十六进制
static final String PUSH_HEX_MSG = "8001001003046e0004066610021100b50c000000000000000000000000008a0301870500000000000000308f050000000000ffffffffffffffff328f050000000000338f050000000000348f050000000000358f050000000000368f050000000000378f050000000000388f050000000000398f0500000000003a8f0500000000003b8f0500000000003c8f0500000000003d8f0500000000003e8f0500000000003f8f050000000000408f050000000000418f050000000000428f050000000000438f050000000000448f050000000000458f050000000000468f050000000000478f050000000000488f050000000000498f0500000000004a8f0500000000004b8f0500000000004c8f0500000000004d8f0500000000004e8f0500000000004f8f050000000000508f050000000000518f050000000000528f050000000000538f050000000000548f050000000000558f050000000000568f050000000000578f050000000000588f050000000000598f0500000000005a8f0500000000005b8f0500000000005c8f0500000000005d8f0500000000005e8f0500000000005f8f050000000000608f050000000000618f050000000000628f050000000000638f050000000000648f050000000000658f050000000000668f050000000000678f050000000000688f050000000000698f0500000000006a8f0500000000006b8f0500000000006c8f0500000000006d8f0500000000006e8f0500000000006f8f050000000000708f050000000000718f050000000000728f050000000000738f050000000000748f050000000000758f050000000000768f050000000000778f050000000000788f050000000000798f0500000000007a8f0500000000007b8f0500000000007c8f0500000000007d8f0500000000007e8f0500000000007f8f050000000000808f050000000000818f050000000000828f050000000000838f050000000000848f050000000000858f050000000000868f050000000000878f050000000000888f050000000000898f0500000000008a8f0500000000008b8f0500000000008c8f0500000000008d8f0500000000008e8f0500000000008f8f050000000000908f050000000000918f050000000000928f050000000000938f050000000000948f050000000000958f050000000000968f050000000000978f050000000000988f050000000000998f0500000000009a8f0500000000009b8f0500000000009c8f0500000000009d8f0500000000009e8f0500000000009f8f050000000000";
static final String END_HEX_MSG = "8001001003046e0004066607040200b50c000000000000000000000000008a03018705000000000000003005050000000000ffffffffffffffff328f050000000000338f050000000000348f050000000000358f050000000000368f050000000000378f050000000000388f050000000000398f0500000000003a8f0500000000003b8f0500000000003c8f0500000000003d8f0500000000003e8f0500000000003f8f050000000000408f050000000000418f050000000000428f050000000000438f050000000000448f050000000000458f050000000000468f050000000000478f050000000000488f050000000000498f0500000000004a8f0500000000004b8f0500000000004c8f0500000000004d8f0500000000004e8f0500000000004f8f050000000000508f050000000000518f050000000000528f050000000000538f050000000000548f050000000000558f050000000000568f050000000000578f050000000000588f050000000000598f0500000000005a8f0500000000005b8f0500000000005c8f0500000000005d8f0500000000005e8f0500000000005f8f050000000000608f050000000000618f050000000000628f050000000000638f050000000000648f050000000000658f050000000000668f050000000000678f050000000000688f050000000000698f0500000000006a8f0500000000006b8f0500000000006c8f0500000000006d8f0500000000006e8f0500000000006f8f050000000000708f050000000000718f050000000000728f050000000000738f050000000000748f050000000000758f050000000000768f050000000000778f050000000000788f050000000000798f0500000000007a8f0500000000007b8f0500000000007c8f0500000000007d8f0500000000007e8f0500000000007f8f050000000000808f050000000000818f050000000000828f050000000000838f050000000000848f050000000000858f050000000000868f050000000000878f050000000000888f050000000000898f0500000000008a8f0500000000008b8f0500000000008c8f0500000000008d8f0500000000008e8f0500000000008f8f050000000000908f050000000000918f050000000000928f050000000000938f050000000000948f050000000000958f050000000000968f050000000000978f050000000000988f050000000000998f0500000000009a8f0500000000009b8f0500000000009c8f0500000000009d8f0500000000009e8f0500000000009f8f050000000000";
// ===================== 全局静态资源 =====================
// 全局预编译报文Buf,所有通道复用
private static ByteBuf GLOBAL_PUSH_BUF;
private static ByteBuf GLOBAL_END_BUF;
// 在线客户端通道集合
public static final Set<Channel> ONLINE_CHANNELS = ConcurrentHashMap.newKeySet();
// 全局统计指标
public static final AtomicLong TOTAL_PACK_SENT = new AtomicLong(0);
public static final AtomicLong TOTAL_BYTE_SENT = new AtomicLong(0);
public static final AtomicLong END_BYTE_SENT = new AtomicLong(0);
// 单条报文字节长度
public static int MSG_BYTE_LEN;
public static int END_BYTE_LEN;
static {
// 启动一次性转换十六进制报文
byte[] msgBytes = hexToBytes(PUSH_HEX_MSG);
MSG_BYTE_LEN = msgBytes.length;
GLOBAL_PUSH_BUF = PooledByteBufAllocator.DEFAULT.buffer(MSG_BYTE_LEN);
GLOBAL_PUSH_BUF.writeBytes(msgBytes);
// 统计初始字节基数
TOTAL_BYTE_SENT.set(0);
// 启动一次性转换十六进制报文
byte[] endMsgBytes = hexToBytes(END_HEX_MSG);
END_BYTE_LEN = endMsgBytes.length;
GLOBAL_END_BUF = PooledByteBufAllocator.DEFAULT.buffer(END_BYTE_LEN);
GLOBAL_END_BUF.writeBytes(endMsgBytes);
// 统计初始字节基数
END_BYTE_SENT.set(0);
}
public static void main(String[] args) throws InterruptedException {
NioEventLoopGroup bossGroup = new NioEventLoopGroup(1);
NioEventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap bootstrap = new ServerBootstrap();
bootstrap.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.handler(new LoggingHandler(LogLevel.INFO))
.option(ChannelOption.SO_BACKLOG, 1024)
.childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
.childOption(ChannelOption.TCP_NODELAY, true) // 关闭Nagle,极致推送
.childOption(ChannelOption.SO_SNDBUF, 4 * 1024 * 1024)
.childOption(ChannelOption.SO_RCVBUF, 4 * 1024 * 1024)
.childOption(ChannelOption.WRITE_BUFFER_HIGH_WATER_MARK, HIGH_WATER_MARK)
.childOption(ChannelOption.WRITE_BUFFER_LOW_WATER_MARK, LOW_WATER_MARK)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();
// pipeline.addLast(new LoggingHandler(LogLevel.DEBUG)); // 调试打开,压测关闭
pipeline.addLast(new IdleStateHandler(0, 50, 0));
pipeline.addLast(new PushHandler());
}
});
// 绑定端口启动服务
ChannelFuture bindFuture = bootstrap.bind(LISTEN_PORT).sync();
System.out.println("TCP推送服务启动成功,监听端口:" + LISTEN_PORT);
// 定时打印推送统计
long lastPackCount = 0;
long lastByteCount = 0;
while (true) {
Thread.sleep(1000);
long currPack = TOTAL_PACK_SENT.get();
long currByte = TOTAL_BYTE_SENT.get();
long qps = currPack - lastPackCount;
long bytePerSec = currByte - lastByteCount;
lastPackCount = currPack;
lastByteCount = currByte;
System.out.printf("[统计]在线连接:%d | QPS:%d | 每秒字节:%d | 总发包:%d%n",
ONLINE_CHANNELS.size(), qps, bytePerSec, currPack);
}
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
/** 推送处理器,单个客户端通道独立循环推送 */
static class PushHandler extends ChannelInboundHandlerAdapter {
private int i = 0;
private Channel clientChannel;
@Override
public void channelActive(ChannelHandlerContext ctx) throws InterruptedException {
clientChannel = ctx.channel();
ONLINE_CHANNELS.add(clientChannel);
System.out.println("客户端接入:" + clientChannel.remoteAddress() + ",开始推送报文");
// 通道就绪立即启动推送循环
startPushLoop();
}
/** 全速推送循环:无休眠,写缓冲满自动暂停 */
private void startPushLoop() throws InterruptedException {
// 通道存活 且 当前缓冲区未达高水位(可写)持续发送
while (clientChannel.isActive() && clientChannel.isWritable()) {
// 复用全局Buf,duplicate不拷贝底层数组,retain防止释放
ByteBuf sendBuf = GLOBAL_PUSH_BUF.duplicate().retain();
clientChannel.write(sendBuf);
TOTAL_PACK_SENT.incrementAndGet();
TOTAL_BYTE_SENT.addAndGet(MSG_BYTE_LEN);
if(TOTAL_PACK_SENT.get() >= 10000000) {
for (int j = 0; j < 5; j++) {
ByteBuf endBuf = GLOBAL_END_BUF.duplicate().retain();
clientChannel.write(endBuf);
END_BYTE_SENT.addAndGet(END_BYTE_LEN);
}
clientChannel.flush();
return;
}
}
// 批量刷写,减少系统调用
clientChannel.flush();
}
/** 缓冲区水位变化回调,可写后恢复推送 */
@Override
public void channelWritabilityChanged(ChannelHandlerContext ctx) throws InterruptedException {
if(clientChannel.isWritable()) {
startPushLoop();
}
}
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
if (evt instanceof IdleStateEvent) {
IdleStateEvent e = (IdleStateEvent) evt;
if (e.state() == IdleState.WRITER_IDLE) {
System.out.println("写超时,可能对端消费异常,主动断开:" + ctx.channel().remoteAddress());
ctx.close();
}
}
}
/** 客户端断开,停止推送,移除在线集合 */
@Override
public void channelInactive(ChannelHandlerContext ctx) {
Channel ch = ctx.channel();
ONLINE_CHANNELS.remove(ch);
System.out.println("客户端断开:" + ch.remoteAddress() + ",停止该通道推送");
}
/** 异常直接关闭通道 */
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
cause.printStackTrace();
ctx.close();
}
}
public static byte[] hexToBytes(String hexStr) {
if (hexStr == null || hexStr.trim().isEmpty()) {
return new byte[0];
}
int length = hexStr.length();
byte[] bytes = new byte[length / 2];
for (int i = 0; i < length; i += 2) {
int high = Character.digit(hexStr.charAt(i), 16);
int low = Character.digit(hexStr.charAt(i + 1), 16);
bytes[i / 2] = (byte) ((high << 4) | low);
}
return bytes;
}
}
\ No newline at end of file
package com.monitor.d2d.netty;
import lombok.Data;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Map;
/**
* 收到的上报帧记录
*/
@Data
public class ReportEntry {
private static final DateTimeFormatter FMT = DateTimeFormatter.ofPattern("HH:mm:ss.SSS");
private LocalDateTime recvTime;
private long sid;
private long did;
private long bid;
private String frameSummary;
private Map<String, Object> parsed;
public String getRecvTimeStr() {
return recvTime != null ? recvTime.format(FMT) : "";
}
public String getBidName() {
switch ((int)(bid & 0xFFFFFFFFL)) {
case 0x0000F001: return "过程控制命令";
case 0x0000F002: return "控制响应令";
case 0x0000F003: return "状态查询令";
case 0x0000F004: return "状态上报令";
default: return String.format("未知BID(%08X)", bid);
}
}
}
\ No newline at end of file
package com.monitor.d2d.netty;
import com.monitor.d2d.protocol.*;
import com.monitor.d2d.service.ReportService;
import com.monitor.d2d.netty.ReportEntry;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.channel.socket.DatagramPacket;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* 接收上报帧,解析后存入ReportService
*/
@Slf4j
@Component
@ChannelHandler.Sharable
public class ReportInboundHandler extends SimpleChannelInboundHandler<DatagramPacket> {
@Autowired
private ReportService reportService;
@Override
protected void channelRead0(ChannelHandlerContext ctx, DatagramPacket msg) {
ByteBuf buf = msg.content();
if (buf.readableBytes() < D2dFrame.HEADER_LENGTH) {
log.warn("收到过短的UDP包,忽略");
return;
}
D2dFrame frame = D2dFrame.decode(buf);
if (frame == null) return;
if(61444 == frame.getBid()){
return;
}
log.info("← 收到帧: {}", frame.summary());
ReportEntry entry = new ReportEntry();
entry.setRecvTime(LocalDateTime.now());
entry.setSid(frame.getSid());
entry.setDid(frame.getDid());
entry.setBid(frame.getBid());
entry.setFrameSummary(frame.summary());
entry.setParsed(parseData(frame));
reportService.add(entry);
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
log.error("接收异常: {}", cause.getMessage(), cause);
}
// ── 数据解析 ──────────────────────────────────────────────────────
private Map<String, Object> parseData(D2dFrame frame) {
Map<String, Object> result = new LinkedHashMap<>();
try {
if (frame.getBid() == BidType.STATUS_REPORT) {
parseStatusReport(frame.getData(), result);
} else if (frame.getBid() == BidType.CONTROL_RESPONSE) {
parseControlResponse(frame.getData(), result);
} else {
result.put("rawHex", bytesToHex(frame.getData(), 64));
}
} catch (Exception e) {
result.put("parseError", e.getMessage());
result.put("rawHex", bytesToHex(frame.getData(), 64));
}
return result;
}
/** 解析状态上报令 BID=0x0000F004 */
private void parseStatusReport(byte[] data, Map<String, Object> out) {
if (data == null || data.length < 2) return;
ProtoReader r = new ProtoReader(data);
out.put("type", "状态上报令");
out.put("No", r.readU16());
// Unit1 start marker
if (r.remaining() < 2) return;
int unitStart = r.readU16(); // 0x01FB
out.put("unitStartMarker", String.format("0x%04X", unitStart));
// 系统状态信息
if (r.remaining() < 1) return;
int monitorFlag = r.readU8();
out.put("monitorFlag", monitorFlag == 1 ? "分控" : monitorFlag == 2 ? "本控" : "未知(" + monitorFlag + ")");
if (r.remaining() < 1) return;
int devStatus = r.readU8();
out.put("deviceStatus", devStatus == 1 ? "故障" : devStatus == 2 ? "正常" : devStatus == 3 ? "维护" : "未知(" + devStatus + ")");
if (r.remaining() < 4) return; out.put("cpuPercent", r.readFloat());
if (r.remaining() < 4) return; out.put("memoryPercent", r.readFloat());
if (r.remaining() < 4) return; out.put("storagePercent", r.readFloat());
if (r.remaining() < 4) return; out.put("networkRateKbps", r.readU32());
if (r.remaining() < 1) return;
int dbConn = r.readU8();
out.put("dbConnection", dbConn == 1 ? "断开" : dbConn == 2 ? "连接" : "未知(" + dbConn + ")");
if (r.remaining() < 1) return; out.put("storageTimeLimitDays", r.readU8());
if (r.remaining() < 1) return; out.put("storagePercent2", r.readU8());
if (r.remaining() < 1) return; out.put("transmitProtocol1", r.readU8());
if (r.remaining() < 1) return; out.put("transmitProtocol2", r.readU8());
if (r.remaining() < 1) return; out.put("recvTaskAmount", r.readU8());
if (r.remaining() < 4) return; out.put("recvSpeedKbps", r.readU32());
if (r.remaining() < 1) return; out.put("sendTaskAmount", r.readU8());
if (r.remaining() < 4) return; out.put("sendSpeedKbps", r.readU32());
if (r.remaining() < 128) return; out.put("storagePath", r.readFixedString(128));
if (r.remaining() < 1) return;
int swNum = r.readU8();
out.put("softwareUnitCount", swNum);
StringBuilder swInfo = new StringBuilder();
for (int i = 0; i < swNum && r.remaining() >= 2; i++) {
int uid = r.readU8();
int ust = r.readU8();
swInfo.append(String.format("[%02X=%s]", uid, ust == 0 ? "运行" : ust == 1 ? "断开" : ust == 2 ? "维护" : "异常"));
}
out.put("softwareUnits", swInfo.toString());
if (r.remaining() < 1) return;
int devNum = r.readU8();
out.put("deviceCount", devNum);
StringBuilder devInfo = new StringBuilder();
for (int i = 0; i < devNum && r.remaining() >= 24; i++) {
String name = r.readFixedString(20);
long uac = r.readU32();
devInfo.append(String.format("[%s UAC=%08X]", name, uac));
}
out.put("devices", devInfo.toString());
}
/** 解析控制响应令 BID=0x0000F002 */
private void parseControlResponse(byte[] data, Map<String, Object> out) {
if (data == null || data.length < 2) return;
ProtoReader r = new ProtoReader(data);
int cmdId = r.readU16();
out.put("type", "控制响应令");
out.put("cmdId", String.format("0x%04X", cmdId));
out.put("cmdName", CmdId.nameOf(cmdId));
switch (cmdId) {
case 0x5404: parse5404Response(r, out); break;
case 0x5406: parse5406Response(r, out); break;
case 0x5422: parse5422Response(r, out); break;
case 0x5419: parse5419Response(r, out); break;
case 0x5420: parse5420Response(r, out); break;
case 0x5421: parse5421Response(r, out); break;
default: parseSimpleResponse(r, out); break;
}
}
/** 简单控制结果(大多数命令上报格式:CMD_ID + 控制结果U8) */
private void parseSimpleResponse(ProtoReader r, Map<String, Object> out) {
if (r.remaining() < 1) return;
int result = r.readU8();
out.put("result", decodeSimpleResult(result));
}
/** 数据接收控制上报 5404 */
private void parse5404Response(ProtoReader r, Map<String, Object> out) {
if (r.remaining() < 9) return;
out.put("planNo", r.readFixedString(9));
out.put("srcAddr", r.readFixedString(20));
out.put("spacecraft", r.readFixedString(6));
out.put("orbitNo", r.readFixedString(10));
out.put("channelNo", String.format("%08X", r.readU32()));
out.put("transferNo", r.readU32());
out.put("protocol", r.readU8());
int procCount = r.remaining() >= 1 ? r.readU8() : 0;
for (int i = 0; i < procCount && r.remaining() >= 2; i++) {
int procId = r.readU8();
int procResult = r.readU8();
String key = "proc" + procId;
out.put(key + "Id", procId);
out.put(key + "Result", decodeProcessResult(procId, procResult));
if (procId == 2 && r.remaining() >= 8) {
out.put("completedFiles", r.readU32());
int curCount = (int) r.readU32();
out.put("currentFileCount", curCount);
StringBuilder files = new StringBuilder();
for (int j = 0; j < curCount && r.remaining() >= 256; j++) {
files.append(r.readFixedString(256)).append(";");
}
out.put("fileNames", files.toString());
if (r.remaining() >= 12) {
out.put("recvRateKbps", r.readU32());
out.put("recvSizeBytes", r.readU64());
}
} else if (procId == 3 && r.remaining() >= 100) {
r.skip(100); // 预留
}
}
}
/** 数据发送控制上报 5406 */
private void parse5406Response(ProtoReader r, Map<String, Object> out) {
if (r.remaining() < 4) return;
out.put("transferNo", r.readU32());
out.put("planNo", r.readFixedString(9));
out.put("srcAddr", r.readFixedString(20));
out.put("spacecraft", r.readFixedString(6));
out.put("orbitNo", r.readFixedString(10));
out.put("channelNo", String.format("%08X", r.readU32()));
out.put("protocol", r.readU8());
out.put("dataCenter", r.readFixedString(20));
out.put("maxRateKbps", r.readU32());
int procCount = r.remaining() >= 1 ? r.readU8() : 0;
for (int i = 0; i < procCount && r.remaining() >= 2; i++) {
int procId = r.readU8();
int procResult = r.readU8();
out.put("proc" + procId + "Result", decodeProcessResult(procId, procResult));
if (procId == 2 && r.remaining() >= 8) {
out.put("sentFiles", r.readU32());
int curCnt = (int) r.readU32();
out.put("currentSendCount", curCnt);
StringBuilder files = new StringBuilder();
for (int j = 0; j < curCnt && r.remaining() >= 256; j++) {
files.append(r.readFixedString(256)).append(";");
}
out.put("sendingFiles", files.toString());
if (r.remaining() >= 8) out.put("sentBytes", r.readU64());
if (r.remaining() >= 4) out.put("sendRateKbps", r.readU32());
if (r.remaining() >= 4) out.put("forwardDelayMs", r.readU32());
if (r.remaining() >= 4) out.put("rttMs", r.readU32());
if (r.remaining() >= 4) out.put("sendWindowBit", r.readU32());
if (r.remaining() >= 10) r.skip(10); // 预留
}
}
}
/** 返向数据任务控制上报 5422 */
private void parse5422Response(ProtoReader r, Map<String, Object> out) {
if (r.remaining() < 1) return;
int primary = r.readU8();
out.put("primaryFlag", primary == 0x0F ? "主机" : primary == (byte)0xF0 ? "备机" : "未知");
out.put("planNo", r.readFixedString(9));
out.put("terminalStation", r.readFixedString(20));
out.put("relaySatTaskCode", String.format("%04X", r.readU16()));
out.put("relayAntenna", String.format("%04X", r.readU16()));
out.put("spacecraft", r.readFixedString(6));
out.put("orbitNo", r.readFixedString(10));
out.put("transferNo", r.readU32());
out.put("mode", r.readU8() == 1 ? "实战" : "联试");
// 简化:其余过程信息
out.put("remaining", r.remaining() + "字节待解析");
}
/** 数据接收快传上报 5419 */
private void parse5419Response(ProtoReader r, Map<String, Object> out) {
out.put("fileName", r.readFixedString(256));
if (r.remaining() >= 1) {
int s = r.readU8();
out.put("status", s == 1 ? "开始接收" : s == 2 ? "正在发送" : s == 3 ? "执行成功" : "执行失败");
}
}
/** 数据发送快传上报 5420 */
private void parse5420Response(ProtoReader r, Map<String, Object> out) {
out.put("fileName", r.readFixedString(256));
if (r.remaining() >= 1) {
int s = r.readU8();
out.put("status", s == 1 ? "开始发送" : s == 2 ? "正在发送" : s == 3 ? "执行成功" : "执行失败");
}
}
/** 任务查询 5421 - 数传→监控,要求监控重发任务 */
private void parse5421Response(ProtoReader r, Map<String, Object> out) {
if (r.remaining() >= 4) {
long flag = r.readU32();
out.put("moduleStatus", flag == 0 ? "模块已启动" : "模块未启动");
}
}
// ── 枚举值解码 ────────────────────────────────────────────────────
private String decodeSimpleResult(int v) {
switch (v) {
case 1: return "正常完成";
case 2: return "异常结束";
case 3: return "分控不执行";
default: return "未知(" + v + ")";
}
}
private String decodeProcessResult(int procId, int v) {
if (procId == 1) {
switch (v) {
case 1: return "接收成功";
case 2: return "接收失败";
case 3: return "分控不执行";
case 4: return "卫星代号不存在";
case 5: return "数据接收设备不存在";
case 6: return "任务已存在";
case 7: return "任务不存在";
case 8: return "任务标识不存在";
case 9: return "任务不存在(停止失败)";
default: return "未知(" + v + ")";
}
} else if (procId == 2) {
switch (v) {
case 1: return "未执行";
case 2: return "正在执行";
case 3: return "执行成功";
case 4: return "执行失败";
case 5: return "暂停";
default: return "未知(" + v + ")";
}
} else if (procId == 3) {
switch (v) {
case 1: return "未执行";
case 3: return "执行成功";
case 4: return "执行失败";
case 5: return "超时结束";
default: return "未知(" + v + ")";
}
}
return "(" + v + ")";
}
private String bytesToHex(byte[] bytes, int maxLen) {
if (bytes == null) return "(null)";
int len = Math.min(bytes.length, maxLen);
StringBuilder sb = new StringBuilder();
for (int i = 0; i < len; i++) {
sb.append(String.format("%02X ", bytes[i]));
}
if (bytes.length > maxLen) sb.append("...(").append(bytes.length).append("字节)");
return sb.toString().trim();
}
}
\ No newline at end of file
package com.monitor.d2d.netty;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.channel.Channel;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.channel.ChannelOption;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.DatagramPacket;
import io.netty.channel.socket.nio.NioDatagramChannel;
import lombok.extern.slf4j.Slf4j;
import java.net.InetSocketAddress;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.LockSupport;
/**
* UDP PDXP 模拟数据推送器。
*
* cmd5422(返向数据任务控制)命令中,若检测到终端站端传输层协议为 UDP
* (0x50=UDP单播 / 0x60=UDP指定源组播 / 0x70=UDP任意源组播),
* 则视为"触发UDP参数":延迟5秒后启动一个独立的 Netty UDP 发送线程,
* 本地绑定 terminalLinks.terminalIp,向 terminalLinks.dataTransIp:dataTransPort
* 以约 900Mbps 的速率持续推送模拟 PDXP 报文。
*
* 报文内容与 {@link NettyTcpPushServer#PUSH_HEX_MSG} 完全一致,
* 仅是承载协议由 TCP 换成 UDP,用于联调/压测场景下模拟真实数据流。
*
* 注意:bindIp(terminalIp)需为本机已实际配置的IP地址,否则Netty绑定会失败。
*/
@Slf4j
public class UdpPdxpPusher {
/** 目标速率:1000 Mbps(与工程实际吞吐目标一致,单位 bit/s) */
private static final long TARGET_BITS_PER_SEC = 1000_000_000L;
private static final long TARGET_BYTES_PER_SEC = TARGET_BITS_PER_SEC / 8;
/** 触发后延迟启动时间 */
private static final long START_DELAY_MS = 5000L;
/** 限速滑动窗口 */
private static final long WINDOW_NANOS = 20_000_000L; // 20ms
/** 预编译模拟报文字节,内容与 NettyTcpPushServer.PUSH_HEX_MSG 一致 */
private static final byte[] MSG_BYTES = NettyTcpPushServer.hexToBytes(NettyTcpPushServer.PUSH_HEX_MSG);
private static final int MSG_LEN = MSG_BYTES.length;
/** 每个任务(按key区分)对应一个推送实例,避免重复启动 */
private static final ConcurrentHashMap<String, UdpPdxpPusher> RUNNING = new ConcurrentHashMap<>();
/** 延迟调度 + 统计打印共用的调度线程池 */
private static final ScheduledExecutorService SCHEDULER =
Executors.newScheduledThreadPool(2, r -> {
Thread t = new Thread(r, "udp-pdxp-scheduler");
t.setDaemon(true);
return t;
});
private final String key;
private final String bindIp;
private final String targetIp;
private final int targetPort;
private final AtomicBoolean running = new AtomicBoolean(false);
private final AtomicLong totalBytesSent = new AtomicLong(0);
private final AtomicLong totalPackSent = new AtomicLong(0);
private NioEventLoopGroup group;
private Channel channel;
private ScheduledFuture<?> statsFuture;
private UdpPdxpPusher(String key, String bindIp, String targetIp, int targetPort) {
this.key = key;
this.bindIp = bindIp;
this.targetIp = targetIp;
this.targetPort = targetPort;
}
/**
* cmd5422 检测到UDP传输参数后调用:延迟5秒启动模拟推送。
* 若同一 key 已有任务在跑,会先停止旧任务再重新计时启动(应对参数变更/重复下发)。
*
* @param key 任务唯一标识(建议:计划号 + 终端站标识 + terminalIp + dataTransIp + dataTransPort)
* @param bindIp 本地绑定地址(对应 terminalLinks.terminalIp,需为本机已配置的IP)
* @param targetIp 目标地址(对应 terminalLinks.dataTransIp)
* @param targetPort 目标端口(对应 terminalLinks.dataTransPort)
*/
public static void triggerDelayedStart(String key, String bindIp, String targetIp, int targetPort) {
stop(key);
log.info("→ [UDP-PDXP] 检测到UDP传输参数,任务[{}]将在{}秒后启动模拟推送 {} → {}:{}",
key, START_DELAY_MS / 1000, bindIp, targetIp, targetPort);
SCHEDULER.schedule(() -> doStart(key, bindIp, targetIp, targetPort), START_DELAY_MS, TimeUnit.MILLISECONDS);
}
private static void doStart(String key, String bindIp, String targetIp, int targetPort) {
UdpPdxpPusher pusher = new UdpPdxpPusher(key, bindIp, targetIp, targetPort);
RUNNING.put(key, pusher);
try {
pusher.start();
} catch (Exception e) {
log.error("[UDP-PDXP] 任务[{}]启动失败(绑定地址 {} 是否为本机已配置的IP?)", key, bindIp, e);
RUNNING.remove(key);
}
}
/** 停止指定任务的推送(例如收到 cmd5422 停止指令时调用) */
public static void stop(String key) {
UdpPdxpPusher pusher = RUNNING.remove(key);
if (pusher != null) {
pusher.shutdown();
}
}
private void start() throws InterruptedException {
group = new NioEventLoopGroup(1);
Bootstrap b = new Bootstrap();
b.group(group)
.channel(NioDatagramChannel.class)
.option(ChannelOption.SO_SNDBUF, 8 * 1024 * 1024)
.handler(new ChannelInboundHandlerAdapter());
// 绑定本地地址为 terminalLinks.terminalIp(端口任意,系统分配)
channel = b.bind(new InetSocketAddress(bindIp, 0)).sync().channel();
running.set(true);
InetSocketAddress target = new InetSocketAddress(targetIp, targetPort);
log.info("→ [UDP-PDXP] 任务[{}]开始以约{}Mbps速率推送PDXP模拟数据 {} → {}",
key, TARGET_BITS_PER_SEC / 1_000_000, bindIp, target);
// 独立发送线程做限速发送循环,不占用EventLoop
Thread senderThread = new Thread(() -> sendLoop(target), "udp-pdxp-sender-" + key);
senderThread.setDaemon(true);
senderThread.start();
statsFuture = SCHEDULER.scheduleAtFixedRate(this::logStats, 1, 1, TimeUnit.SECONDS);
}
/**
* 限速发送循环:按固定时间窗口累计已发字节数,
* 一旦达到该窗口应发字节数即挂起剩余时间,从而把整体速率钳制在约 900Mbps。
*/
private void sendLoop(InetSocketAddress target) {
long bytesPerWindow = TARGET_BYTES_PER_SEC * WINDOW_NANOS / 1_000_000_000L;
long windowStart = System.nanoTime();
long sentInWindow = 0;
int i = 0;
while (running.get() && channel.isActive() && i++ < 1) {
ByteBuf buf = channel.alloc().buffer(MSG_LEN);
buf.writeBytes(MSG_BYTES);
channel.writeAndFlush(new DatagramPacket(buf, target));
totalPackSent.incrementAndGet();
totalBytesSent.addAndGet(MSG_LEN);
sentInWindow += MSG_LEN;
if (sentInWindow >= bytesPerWindow) {
long elapsed = System.nanoTime() - windowStart;
long remain = WINDOW_NANOS - elapsed;
if (remain > 0) {
LockSupport.parkNanos(remain);
}
windowStart = System.nanoTime();
sentInWindow = 0;
}
}
shutdown();
}
private void logStats() {
if (!running.get()) return;
log.info("[UDP-PDXP][{}] 累计发包={} 累计字节={} (约{}Mbps)",
key, totalPackSent.get(), totalBytesSent.get(),
String.format("%.1f", totalBytesSent.get() * 8.0 / 1_000_000));
}
private void shutdown() {
running.set(false);
if (statsFuture != null) statsFuture.cancel(false);
if (channel != null) channel.close();
if (group != null) group.shutdownGracefully();
log.info("[UDP-PDXP] 任务[{}]已停止,累计发送{}包/{}字节", key, totalPackSent.get(), totalBytesSent.get());
}
}
package com.monitor.d2d.netty;
import com.monitor.d2d.config.AppProperties;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelOption;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.InternetProtocolFamily;
import io.netty.channel.socket.nio.NioDatagramChannel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.net.*;
import java.util.Enumeration;
/**
* UDP 组播接收端
* 加入组播组 238.1.0.3,监听端口 8888,接收数传软件上报
*/
@Slf4j
@Component
public class UdpReceiver {
@Autowired
private AppProperties props;
@Autowired
private ReportInboundHandler reportHandler;
private NioEventLoopGroup group;
private Channel channel;
@PostConstruct
public void start() throws Exception {
NetworkInterface ni = resolveNetworkInterface();
InetAddress groupAddr = InetAddress.getByName(props.getRecvGroup());
group = new NioEventLoopGroup(1);
Bootstrap b = new Bootstrap();
b.group(group)
.channelFactory(() -> new NioDatagramChannel(InternetProtocolFamily.IPv4))
.option(ChannelOption.SO_REUSEADDR, true)
.option(ChannelOption.IP_MULTICAST_IF, ni)
.option(ChannelOption.IP_MULTICAST_TTL, 32)
.handler(reportHandler);
channel = b.bind(new InetSocketAddress(props.getRecvPort())).sync().channel();
// 加入组播组
((NioDatagramChannel) channel).joinGroup(groupAddr, ni, null, channel.newPromise()).sync();
log.info("UDP接收端已启动,已加入组播组 {}:{} 网络接口: {}",
props.getRecvGroup(), props.getRecvPort(), ni.getName());
}
/**
* 自动或按配置选择网络接口
*/
private NetworkInterface resolveNetworkInterface() throws SocketException {
String ifName = props.getNetworkInterface();
if (ifName != null && !ifName.isEmpty()) {
NetworkInterface ni = NetworkInterface.getByName(ifName);
if (ni != null) return ni;
log.warn("配置的网络接口 [{}] 不存在,尝试自动选择", ifName);
}
// 自动选择:第一个支持组播且已启动的非回环接口
Enumeration<NetworkInterface> nis = NetworkInterface.getNetworkInterfaces();
while (nis.hasMoreElements()) {
NetworkInterface ni = nis.nextElement();
if (ni.isUp() && ni.supportsMulticast() && !ni.isLoopback()) {
log.info("自动选择网络接口: {}", ni.getName());
return ni;
}
}
// 回退到回环(仅用于本地测试)
NetworkInterface lo = NetworkInterface.getByName("lo");
if (lo != null) {
log.warn("未找到可用组播接口,回退到回环接口(仅限本地测试)");
return lo;
}
throw new IllegalStateException("无法找到可用的网络接口");
}
@PreDestroy
public void stop() {
if (channel != null) channel.close();
if (group != null) group.shutdownGracefully();
log.info("UDP接收端已关闭");
}
}
\ No newline at end of file
package com.monitor.d2d.netty;
import com.monitor.d2d.config.AppProperties;
import com.monitor.d2d.protocol.D2dFrame;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.DatagramPacket;
import io.netty.channel.socket.nio.NioDatagramChannel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.net.InetSocketAddress;
import java.net.NetworkInterface;
/**
* UDP 组播发送端
* 发送命令帧到 238.1.0.103:8889
*/
@Slf4j
@Component
public class UdpSender {
@Autowired
private AppProperties props;
private NioEventLoopGroup group;
private Channel channel;
private InetSocketAddress targetAddress;
@PostConstruct
public void start() throws Exception {
targetAddress = new InetSocketAddress(props.getSendHost(), props.getSendPort());
group = new NioEventLoopGroup(1);
Bootstrap b = new Bootstrap();
b.group(group)
.channel(NioDatagramChannel.class)
.option(ChannelOption.IP_MULTICAST_IF, NetworkInterface.getByName(props.getNetworkInterface()))
.option(ChannelOption.SO_REUSEADDR, true)
.handler(new SimpleChannelInboundHandler<Object>() {
@Override
protected void channelRead0(ChannelHandlerContext ctx, Object msg) {
// 发送端不处理接收
}
});
channel = b.bind(0).sync().channel();
log.info("UDP发送端已就绪,目标地址:{}:{}", props.getSendHost(), props.getSendPort());
}
/**
* 发送完整帧
*/
public void send(D2dFrame frame) {
if (channel == null || !channel.isActive()) {
log.error("UDP发送通道未就绪");
return;
}
ByteBuf buf = frame.encode();
DatagramPacket packet = new DatagramPacket(buf, targetAddress);
channel.writeAndFlush(packet).addListener(future -> {
if (!future.isSuccess()) {
log.error("UDP发送失败: {}", future.cause().getMessage());
} else {
log.debug("已发送帧: {}", frame.summary());
}
});
}
@PreDestroy
public void stop() {
if (channel != null) channel.close();
if (group != null) group.shutdownGracefully();
log.info("UDP发送端已关闭");
}
}
\ No newline at end of file
package com.monitor.d2d.protocol;
/**
* 信息类别 BID 常量(见接口文档第5章)
*/
public final class BidType {
/** 过程控制命令 - 监控→数传 */
public static final long PROCESS_CONTROL = 0x0000F001L;
/** 控制响应令 - 数传→监控 */
public static final long CONTROL_RESPONSE = 0x0000F002L;
/** 状态查询令 - 监控→数传 */
public static final long STATUS_QUERY = 0x0000F003L;
/** 状态上报令 - 数传→监控 */
public static final long STATUS_REPORT = 0x0000F004L;
private BidType() {}
}
\ No newline at end of file
package com.monitor.d2d.protocol;
/**
* 过程控制命令编码 CMD_ID(见接口文档6.2节)
*
* 高字节固定为 'T'=0x54,低字节为顺序编号01~1B
*/
public final class CmdId {
public static final short CONTROL_MODE = 0x5401; // 6.2.1 控制方式控制
public static final short TRANSPORT_PROTOCOL = 0x5402; // 6.2.2 传输层协议控制
public static final short STORAGE_PATH = 0x5403; // 6.2.3 存储路径控制
public static final short DATA_RECEIVE = 0x5404; // 6.2.4 数据接收控制
public static final short DATA_RECEIVE_END = 0x5405; // 6.2.5 数据接收结束告知
public static final short DATA_SEND = 0x5406; // 6.2.6 数据发送控制
public static final short RATE_ADJUST = 0x5407; // 6.2.7 速率调整控制
public static final short SPACECRAFT_ADJUST = 0x5408; // 6.2.8 航天器调整控制
public static final short SPACECRAFT_DELETE = 0x5409; // 6.2.9 航天器删除控制
public static final short DEVICE_ADJUST = 0x5410; // 6.2.10 设备调整控制
public static final short DEVICE_DELETE = 0x5411; // 6.2.11 设备删除控制
public static final short CENTER_ADJUST = 0x5412; // 6.2.12 中心调整控制
public static final short CENTER_DELETE = 0x5413; // 6.2.13 中心删除控制
public static final short SOFTWARE_CONTROL = 0x5414; // 6.2.14 软件控制
public static final short MODULE_CONTROL = 0x5415; // 6.2.15 模块控制
public static final short TRANSMISSION_TEST = 0x5416; // 6.2.16 传输测试控制
public static final short STORAGE_PERIOD = 0x5417; // 6.2.17 存储周期控制
public static final short DATA_DELETE = 0x5418; // 6.2.18 数据删除控制
public static final short FAST_RECV_REPORT = 0x5419; // 6.2.19 数据接收快传上报(仅接收)
public static final short FAST_SEND_REPORT = 0x5420; // 6.2.20 数据发送快传上报(仅接收)
public static final short TASK_QUERY = 0x5421; // 6.2.21 任务查询(数传→监控,仅接收)
public static final short RETURN_DATA_TASK = 0x5422; // 6.2.22 返向数据任务控制
public static final short RELAY_SAT_PARAM = 0x5423; // 6.2.23 中继卫星参数控制
public static final short FAST_TRANSFER_PERIOD = 0x5424; // 6.2.24 快传任务周期调整控制
public static final short FAST_TRANSFER_TASK = 0x5425; // 6.2.25 快传任务调整控制
public static final short SPACE_BER = 0x5426; // 6.2.26 天基误码率控制
public static final short TASK_ADJUST = 0x5427; // 6.2.27 任务调整控制(天基任务速率)
private CmdId() {}
public static String nameOf(int cmdId) {
switch (cmdId & 0xFFFF) {
case 0x5401: return "控制方式控制";
case 0x5402: return "传输层协议控制";
case 0x5403: return "存储路径控制";
case 0x5404: return "数据接收控制";
case 0x5405: return "数据接收结束告知";
case 0x5406: return "数据发送控制";
case 0x5407: return "速率调整控制";
case 0x5408: return "航天器调整控制";
case 0x5409: return "航天器删除控制";
case 0x5410: return "设备调整控制";
case 0x5411: return "设备删除控制";
case 0x5412: return "中心调整控制";
case 0x5413: return "中心删除控制";
case 0x5414: return "软件控制";
case 0x5415: return "模块控制";
case 0x5416: return "传输测试控制";
case 0x5417: return "存储周期控制";
case 0x5418: return "数据删除控制";
case 0x5419: return "数据接收快传上报";
case 0x5420: return "数据发送快传上报";
case 0x5421: return "任务查询";
case 0x5422: return "返向数据任务控制";
case 0x5423: return "中继卫星参数控制";
case 0x5424: return "快传任务周期调整";
case 0x5425: return "快传任务调整控制";
case 0x5426: return "天基误码率控制";
case 0x5427: return "天基任务速率调整";
default: return "未知命令[" + Integer.toHexString(cmdId) + "]";
}
}
}
\ No newline at end of file
package com.monitor.d2d.protocol;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import lombok.Data;
import java.time.LocalDate;
import java.time.LocalTime;
/**
* 应用层数据帧格式(见接口文档表4-3)
*
* | SID(4) | DID(4) | BID(4) | Reserve(4) | Date(4) | Time(4) | DataLen(4) | DataField(N) |
*
* 大端字节序,28字节固定帧头
*/
@Data
public class D2dFrame {
public static final int HEADER_LENGTH = 28;
/** 信源地址 U32 */
private long sid;
/** 信宿地址 U32 */
private long did;
/** 信息类别 U32 */
private long bid;
/** 保留 U32,填0 */
private long reserve = 0;
/** 日期 U32 BCD,如20261101→0x20261101 */
private long date;
/** 时间 U32,单位0.1ms,从当日0时累积 */
private long time;
/** 数据长度 U32 */
private long dataLen;
/** 数据域 */
private byte[] data;
/** 以当前时间构造帧头(data域由外部设置)*/
public static D2dFrame create(long sid, long did, long bid) {
D2dFrame f = new D2dFrame();
f.sid = sid;
f.did = did;
f.bid = bid;
f.date = encodeDateBcd(LocalDate.now());
f.time = encodeTime(LocalTime.now());
return f;
}
/** 编码整帧为ByteBuf(大端字节序)*/
public ByteBuf encode() {
byte[] d = data != null ? data : new byte[0];
ByteBuf buf = Unpooled.buffer(HEADER_LENGTH + d.length);
buf.writeInt((int) sid);
buf.writeInt((int) did);
buf.writeInt((int) bid);
buf.writeInt(0); // reserve
buf.writeInt((int) date);
buf.writeInt((int) time);
buf.writeInt(d.length);
buf.writeBytes(d);
return buf;
}
/** 从ByteBuf解码帧(不移动readerIndex之外) */
public static D2dFrame decode(ByteBuf buf) {
if (buf.readableBytes() < HEADER_LENGTH) {
return null;
}
D2dFrame f = new D2dFrame();
f.sid = buf.readUnsignedInt();
f.did = buf.readUnsignedInt();
f.bid = buf.readUnsignedInt();
f.reserve = buf.readUnsignedInt();
f.date = buf.readUnsignedInt();
f.time = buf.readUnsignedInt();
f.dataLen = buf.readUnsignedInt();
int dLen = (int) Math.min(f.dataLen, buf.readableBytes());
f.data = new byte[dLen];
buf.readBytes(f.data);
return f;
}
// ── 日期/时间编码 ────────────────────────────────────────────────
/** 日期→4字节BCD: 2026-05-18 → 0x20260518 */
public static long encodeDateBcd(LocalDate d) {
int y = d.getYear(), m = d.getMonthValue(), day = d.getDayOfMonth();
return ((long)((y / 1000) & 0xF) << 28)
| ((long)((y / 100 % 10) & 0xF) << 24)
| ((long)((y / 10 % 10) & 0xF) << 20)
| ((long)( y % 10 & 0xF) << 16)
| ((long)((m / 10) & 0xF) << 12)
| ((long)( m % 10 & 0xF) << 8)
| ((long)((day/ 10) & 0xF) << 4)
| ((long)( day% 10 & 0xF));
}
/** BCD日期→可读字符串 */
public static String decodeDateBcd(long bcd) {
return String.format("%04X-%02X-%02X",
(bcd >> 16) & 0xFFFF,
(bcd >> 8) & 0xFF,
bcd & 0xFF);
}
/** 时间→U32(0.1ms since midnight) */
public static long encodeTime(LocalTime t) {
long ms = (t.getHour() * 3600L + t.getMinute() * 60L + t.getSecond()) * 1000L
+ t.getNano() / 1_000_000L;
return ms * 10L; // 0.1ms单位
}
/** 0.1ms→可读字符串 */
public static String decodeTime(long t01ms) {
long ms = t01ms / 10;
long s = ms / 1000; ms %= 1000;
long m = s / 60; s %= 60;
long h = m / 60; m %= 60;
return String.format("%02d:%02d:%02d.%03d", h, m, s, ms);
}
/** 格式化为可读摘要 */
public String summary() {
return String.format("SID=%08X DID=%08X BID=%08X Date=%s Time=%s DataLen=%d",
sid, did, bid, decodeDateBcd(date), decodeTime(time), dataLen);
}
}
\ No newline at end of file
package com.monitor.d2d.protocol;
import java.nio.charset.StandardCharsets;
/**
* 协议字节流读取工具(大端字节序)
*/
public class ProtoReader {
private final byte[] buf;
private int pos;
public ProtoReader(byte[] buf) {
this.buf = buf;
this.pos = 0;
}
public int remaining() {
return buf.length - pos;
}
public boolean hasRemaining(int n) {
return remaining() >= n;
}
public int readU8() {
return buf[pos++] & 0xFF;
}
public int readU16() {
int v = ((buf[pos] & 0xFF) << 8) | (buf[pos + 1] & 0xFF);
pos += 2;
return v;
}
public long readU32() {
long v = 0;
for (int i = 0; i < 4; i++) v = (v << 8) | (buf[pos++] & 0xFF);
return v;
}
public long readU64() {
long v = 0;
for (int i = 0; i < 8; i++) v = (v << 8) | (buf[pos++] & 0xFF);
return v;
}
public float readFloat() {
return Float.intBitsToFloat((int) readU32());
}
public double readDouble() {
return Double.longBitsToDouble(readU64());
}
/**
* 读固定长度字符串(ASCII,去除尾部0字节)
*/
public String readFixedString(int len) {
if (pos + len > buf.length) len = buf.length - pos;
int end = len;
while (end > 0 && buf[pos + end - 1] == 0) end--;
String s = new String(buf, pos, end, StandardCharsets.US_ASCII);
pos += len;
return s.trim();
}
public byte[] readBytes(int n) {
if (pos + n > buf.length) n = buf.length - pos;
byte[] b = new byte[n];
System.arraycopy(buf, pos, b, 0, n);
pos += n;
return b;
}
public void skip(int n) {
pos = Math.min(pos + n, buf.length);
}
}
\ No newline at end of file
package com.monitor.d2d.protocol;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.ByteOrder;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
/**
* 协议字节流构建工具(大端字节序)
*/
public class ProtoWriter {
private final ByteArrayOutputStream out = new ByteArrayOutputStream();
public ProtoWriter writeU8(int v) {
out.write(v & 0xFF);
return this;
}
public ProtoWriter writeU16(int v) {
out.write((v >> 8) & 0xFF);
out.write(v & 0xFF);
return this;
}
public ProtoWriter writeU32(long v) {
out.write((int)((v >> 24) & 0xFF));
out.write((int)((v >> 16) & 0xFF));
out.write((int)((v >> 8) & 0xFF));
out.write((int)( v & 0xFF));
return this;
}
public ProtoWriter writeU64(long v) {
for (int i = 56; i >= 0; i -= 8) {
out.write((int)((v >> i) & 0xFF));
}
return this;
}
public ProtoWriter writeFloat(float v) {
return writeU32(Float.floatToIntBits(v));
}
public ProtoWriter writeDouble(double v) {
return writeU64(Double.doubleToLongBits(v));
}
/**
* 写入固定长度字符串(ASCII,不足补0,超出截断)
*/
public ProtoWriter writeFixedString(String s, int len) {
byte[] bytes = new byte[len];
if (s != null && !s.isEmpty()) {
byte[] src = s.getBytes(StandardCharsets.US_ASCII);
System.arraycopy(src, 0, bytes, 0, Math.min(src.length, len));
}
try {
out.write(bytes);
} catch (IOException e) {
throw new RuntimeException(e);
}
return this;
}
public ProtoWriter writeBytes(byte[] b) {
try {
out.write(b);
} catch (IOException e) {
throw new RuntimeException(e);
}
return this;
}
public byte[] toBytes() {
return out.toByteArray();
}
/** CMD_ID写入:高字节0x54,低字节为命令编号 */
public ProtoWriter writeCmdId(short cmdId) {
return writeU16(cmdId & 0xFFFF);
}
}
\ No newline at end of file
package com.monitor.d2d.protocol.cmd5422;
import lombok.Data;
import java.util.ArrayList;
import java.util.List;
/**
* 5422H 返向数据任务控制 - 完整命令参数(文档表6-35)
*
* 字节流布局(大端):
* ┌─────────────────────┬──────┬──────────────────────────────────────────┐
* │ 字段 │ 类型 │ 说明 │
* ├─────────────────────┼──────┼──────────────────────────────────────────┤
* │ CMD_ID │ U16 │ 5422H │
* │ 任务控制 │ U8 │ 01=开始(部分) 02=开始(完整) 03=修改 04=停止 │
* │ 主备标识 │ U8 │ 0x0F=主机 0xF0=备机 │
* │ 跟踪接收计划号 │ C[9] │ 9位数字定长 │
* │ 终端站标识 │ C[20]│ │
* │ 中继星任务代号 │ U16 │ 如 0x411F │
* │ 中继卫星天线标识 │ U16 │ 如 0x000F=东天线 │
* │ 航天器标识 │ C[6] │ │
* │ 圈号 │ C[10]│ │
* │ 传输号 │ U32 │ │
* │ 捕跟开始时间 │ C[14]│ V6.11新增,预留 │
* │ 模式 │ U8 │ 1=实战(OP) 2=联试(TS) │
* ├─── 终端站端信息 ────┼──────┼──────────────────────────────────────────┤
* │ 应用层协议 │ U8 │ 01=FEP 05=PDXP 06=UDF │
* │ 传输层协议 │ U8 │ 50=UDP 60=UDP指定源 70=UDP任意源 90=TCP │
* │ 星间链路标识 │ C[5] │ │
* │ 初始接收总帧数 │ U64 │ V6.11新增,正常下发填0,切换时填原主机值 │
* │ 初始接收总数据量 │ U64 │ V6.11新增,正常下发填0,切换时填原主机值 │
* │ 链路数量 │ U8 │ │
* │ 链路i主备 │ U8 │ 0x00=不区分 0x0F=主 0xF0=备 │
* │ 链路i路由平面 │ U8 │ 0x00=第一路由 0x01=第二路由 │
* │ 链路i终端站IP │ U32 │ │
* │ 链路i终端站端口 │ U16 │ │
* │ 链路i数传IP │ U32 │ │
* │ 链路i数传端口 │ U16 │ │
* ├─── 对外端信息 ──────┼──────┼──────────────────────────────────────────┤
* │ 对外端数量 │ U8 │ 填0则无后续字段 │
* │ 目的中心标识 │ C[20]│ 每个对外端 │
* │ 应用层协议 │ U8 │ │
* │ 传输层协议 │ U8 │ │
* │ 初始发送总帧数 │ U64 │ V6.11新增,正常下发填0,切换时填原主机值 │
* │ 初始发送总数据量 │ U64 │ V6.11新增,正常下发填0,切换时填原主机值 │
* │ 链路数量 │ U8 │ │
* │ 链路j数传IP │ U32 │ │
* │ 链路j数传端口 │ U16 │ │
* │ 链路j对外端IP │ U32 │ │
* │ 链路j对外端端口 │ U16 │ │
* └─────────────────────┴──────┴──────────────────────────────────────────┘
*/
@Data
public class Cmd5422Request {
// ── 基本参数 ──────────────────────────────────────────────────────
/**
* 任务控制
* 0x01=开始(部分参数,终端站接入申请还未全收到)
* 0x02=开始(参数完整,所有方向均已接入)
* 0x03=修改(修改主备:本地主备 or 应用中心主备)
* 0x04=停止(按计划号+终端站+数据类型停止)
*/
private int taskControl = 0x02;
/**
* 主备标识(数传软件自身角色)
* 0x0F=主机(正常收发) 0xF0=备机(只建连接,不记录不转发)
*/
private int primaryFlag = 0x0F;
/** 跟踪接收计划号(9位数字) */
private String planNo = "";
/** 终端站标识(20字节) */
private String terminalStation = "";
/** 中继星任务代号(U16,如 0x411F) */
private int relayTaskCode = 0;
/** 中继卫星天线标识(U16,如 0x000F=东天线) */
private int relayAntenna = 0;
/** 航天器标识(6字节) */
private String spacecraft = "";
/** 圈号(10字节) */
private String orbitNo = "";
/** 传输号(U32) */
private long transferNo = 0;
/**
* 捕跟开始时间(14字节,预留)(V6.11新增)
*/
private String captureStartTime = "";
/**
* 模式
* 1=实战(OP) 2=联试(TS)
*/
private int mode = 1;
// ── 终端站端信息 ──────────────────────────────────────────────────
/** 终端站端通信参数(协议 + 星间链路ID + 链路列表) */
private TerminalStationInfo terminalInfo = new TerminalStationInfo();
// ── 对外端信息 ────────────────────────────────────────────────────
/**
* 对外发送端列表(应用中心列表)
* 填空 = 不对外发送数据(数传只收不转)
*/
private List<OuterEndpointInfo> outerEndpoints = new ArrayList<>();
// ── 工厂方法 ──────────────────────────────────────────────────────
/**
* 最简场景:单链路TCP,数传→单个应用中心
*
* @param planNo 计划号
* @param terminalStation 终端站标识
* @param relayTaskCode 中继星任务代号
* @param relayAntenna 天线标识
* @param spacecraft 航天器
* @param orbitNo 圈号
* @param transferNo 传输号
* @param terminalIp 终端站IP
* @param terminalPort 终端站端口
* @param dataTransIp 数传本端IP
* @param dataTransPort 数传本端端口
* @param centerIp 应用中心IP(对外端)
* @param centerPort 应用中心端口
* @param centerId 应用中心标识
*/
public static Cmd5422Request simpleTcp(
String planNo, String terminalStation,
int relayTaskCode, int relayAntenna,
String spacecraft, String orbitNo, long transferNo,
String terminalIp, int terminalPort,
String dataTransIp, int dataTransPort,
String centerIp, int centerPort, String centerId) {
Cmd5422Request req = new Cmd5422Request();
req.setTaskControl(0x02);
req.setPrimaryFlag(0x0F);
req.setPlanNo(planNo);
req.setTerminalStation(terminalStation);
req.setRelayTaskCode(relayTaskCode);
req.setRelayAntenna(relayAntenna);
req.setSpacecraft(spacecraft);
req.setOrbitNo(orbitNo);
req.setTransferNo(transferNo);
req.setMode(1);
// 终端站端:TCP单链路
req.setTerminalInfo(TerminalStationInfo.tcpSingle(terminalIp, terminalPort, dataTransIp, dataTransPort));
// 对外端:单中心
req.getOuterEndpoints().add(OuterEndpointInfo.tcpSingle(centerId, dataTransIp, dataTransPort, centerIp, centerPort));
return req;
}
/** 停止命令(只需计划号+终端站标识即可) */
public static Cmd5422Request stop(String planNo, String terminalStation,
int relayTaskCode, int relayAntenna,
String spacecraft, String orbitNo, long transferNo) {
Cmd5422Request req = new Cmd5422Request();
req.setTaskControl(0x04);
req.setPrimaryFlag(0x0F);
req.setPlanNo(planNo);
req.setTerminalStation(terminalStation);
req.setRelayTaskCode(relayTaskCode);
req.setRelayAntenna(relayAntenna);
req.setSpacecraft(spacecraft);
req.setOrbitNo(orbitNo);
req.setTransferNo(transferNo);
req.setMode(1);
// 停止时终端站端信息仍需填充(按文档协议要求保留字段位置)
req.setTerminalInfo(new TerminalStationInfo());
return req;
}
}
package com.monitor.d2d.protocol.cmd5422;
import lombok.Data;
import java.util.ArrayList;
import java.util.List;
/**
* 对外端信息(文档表6-35,Row30~Row38)
*
* 对应协议字段(每个对外端):
* 目的中心标识 Char[20]
* 应用层协议 U8 0x01=FEP 0x05=PDXP 0x06=UDF
* 传输层协议 U8 0x50=UDP 0x60=UDP指定源 0x70=UDP任意源 0x90=TCP
* 初始发送总帧数 U64 (V6.11新增,正常下发填0,切换时填原主机发送总帧数)
* 初始发送总数据量 U64 (V6.11新增,正常下发填0,切换时填原主机发送总数据量)
* 链路数量 U8
* 链路1..N 见 OuterEndpointLink
*/
@Data
public class OuterEndpointInfo {
/**
* 目的中心标识(20字节,对应接收数据的应用中心)
*/
private String centerId = "";
/**
* 应用层协议
* 0x01=FEP 0x05=PDXP 0x06=UDF
*/
private int appProtocol = 0x01;
/**
* 传输层协议
* 0x50=UDP单播 0x60=UDP指定源组播 0x70=UDP任意源组播 0x90=TCP
*/
private int transportProtocol = 0x90;
/**
* 初始发送总帧数(U64,V6.11新增)
* 监控正常下发时填0,切换(主备切换)时填写原主机发送总帧数
*/
private long initialSendFrames = 0;
/**
* 初始发送总数据量(U64,V6.11新增)
* 监控正常下发时填0,切换(主备切换)时填写原主机发送总数据量
*/
private long initialSendBytes = 0;
/**
* 链路列表(不同 linkID 可对应不同地址)
*/
private List<OuterEndpointLink> links = new ArrayList<>();
/** 快速添加链路 */
public OuterEndpointInfo addLink(OuterEndpointLink link) {
links.add(link);
return this;
}
/**
* 快捷方法:单中心单链路TCP(最常见场景)
*
* @param centerId 目的中心标识
* @param dataTransIp 数传本端IP
* @param dataTransPort 数传本端端口
* @param centerIp 应用中心IP
* @param centerPort 应用中心端口
*/
public static OuterEndpointInfo tcpSingle(String centerId,
String dataTransIp, int dataTransPort,
String centerIp, int centerPort) {
OuterEndpointInfo info = new OuterEndpointInfo();
info.setCenterId(centerId);
info.setAppProtocol(0x01);
info.setTransportProtocol(0x90);
info.addLink(OuterEndpointLink.of(dataTransIp, dataTransPort, centerIp, centerPort));
return info;
}
}
package com.monitor.d2d.protocol.cmd5422;
import lombok.Data;
/**
* 对外端 - 单条链路信息(文档表6-35,Row34~Row37)
*
* 对应协议字段(每条链路):
* 链路数传IP U32 数传软件绑定的本地IP(站节点做服务端,中心节点做客户端)
* 链路数传端口 U16 数传软件绑定的本地端口
* 链路对外端IP U32 发送目的IP(站节点无此字段,中心节点有)
* 链路对外端端口 U16 发送目的端口(站节点无此字段)
*
* 注:站节点(地面站)做服务端,无对外端IP/端口;
* 中心节点做客户端,有对外端IP/端口。
* 当前监控作为中心侧,包含完整4字段。
*/
@Data
public class OuterEndpointLink {
/** 数传软件本地绑定IP */
private long dataTransIp;
/** 数传软件本地绑定端口 */
private int dataTransPort;
/** 目标应用中心IP */
private long outerIp;
/** 目标应用中心端口 */
private int outerPort;
/** 便捷构造(字符串IP) */
public static OuterEndpointLink of(String dataTransIp, int dataTransPort,
String outerIp, int outerPort) {
OuterEndpointLink l = new OuterEndpointLink();
l.dataTransIp = TerminalLink.ipToU32(dataTransIp);
l.dataTransPort = dataTransPort;
l.outerIp = TerminalLink.ipToU32(outerIp);
l.outerPort = outerPort;
return l;
}
}
package com.monitor.d2d.protocol.cmd5422;
import lombok.Data;
/**
* 终端站端 - 单条链路信息(文档表6-35,Row16~Row27)
*
* 对应协议字段(每条链路):
* 链路主备 U8 0x00=不区分 0x0F=主机 0xF0=备机
* 路由平面 U8 0x00=第一路由 0x01=第二路由
* 终端站IP U32 点分十进制
* 终端站端口 U16
* 数传IP U32 点分十进制(数传软件绑定的本地IP)
* 数传端口 U16
*/
@Data
public class TerminalLink {
/** 主备标识: 0x00=不区分主备 0x0F=主机 0xF0=备机 */
private int primaryFlag = 0x00;
/** 路由平面: 0x00=第一路由 0x01=第二路由 */
private int routePlane = 0x00;
/** 终端站IP(U32,大端;可用 ipToU32 工具转换) */
private long terminalIp;
/** 终端站端口 */
private int terminalPort;
/** 数传软件本端IP(数传侧绑定地址) */
private long dataTransIp;
/** 数传软件本端端口 */
private int dataTransPort;
// ── 工具方法 ──────────────────────────────────────────────────────
/** "192.168.1.1" → U32(大端) */
public static long ipToU32(String ip) {
if (ip == null || ip.isEmpty()) return 0L;
String[] parts = ip.split("\\.");
if (parts.length != 4) return 0L;
long v = 0;
for (String p : parts) v = (v << 8) | (Integer.parseInt(p.trim()) & 0xFF);
return v;
}
/** U32 → "192.168.1.1" */
public static String u32ToIp(long v) {
return ((v >> 24) & 0xFF) + "." + ((v >> 16) & 0xFF) + "."
+ ((v >> 8) & 0xFF) + "." + (v & 0xFF);
}
/** 便捷构造:直接传字符串IP */
public static TerminalLink of(int primaryFlag, int routePlane,
String terminalIp, int terminalPort,
String dataTransIp, int dataTransPort) {
TerminalLink l = new TerminalLink();
l.primaryFlag = primaryFlag;
l.routePlane = routePlane;
l.terminalIp = ipToU32(terminalIp);
l.terminalPort = terminalPort;
l.dataTransIp = ipToU32(dataTransIp);
l.dataTransPort = dataTransPort;
return l;
}
}
package com.monitor.d2d.protocol.cmd5422;
import lombok.Data;
import java.util.ArrayList;
import java.util.List;
/**
* 终端站端信息(文档表6-35,Row12~Row28)
*
* 对应协议字段:
* 应用层协议 U8 0x01=FEP 0x05=PDXP 0x06=UDF
* 传输层协议 U8 0x50=UDP 0x60=UDP指定源 0x70=UDP任意源 0x90=TCP
* 星间链路标识 Char[5]
* 初始接收总帧数 U64 (V6.11新增,正常下发填0,切换时填原主机接收总帧数)
* 初始接收总数据量 U64 (V6.11新增,正常下发填0,切换时填原主机接收总数据量)
* 链路数量 U8 (最多4条:区分主备+双路由)
* 链路1..N 见 TerminalLink
*/
@Data
public class TerminalStationInfo {
/**
* 应用层协议
* 0x01=FEP 0x05=PDXP 0x06=UDF
*/
private int appProtocol = 0x01;
/**
* 传输层协议
* 0x50=UDP单播 0x60=UDP指定源组播 0x70=UDP任意源组播 0x90=TCP
*/
private int transportProtocol = 0x90;
/**
* 星间链路标识(5字节定长,不用时留空)
*/
private String linkId = "";
/**
* 初始接收总帧数(U64,V6.11新增)
* 监控正常下发时填0,切换(主备切换)时填写原主机接收总帧数
*/
private long initialRecvFrames = 0;
/**
* 初始接收总数据量(U64,V6.11新增)
* 监控正常下发时填0,切换(主备切换)时填写原主机接收总数据量
*/
private long initialRecvBytes = 0;
/**
* 链路列表(最多4条:主/备 × 第一/第二路由平面)
*/
private List<TerminalLink> links = new ArrayList<>();
// ── 工具方法 ──────────────────────────────────────────────────────
/** 快速添加一条链路 */
public TerminalStationInfo addLink(TerminalLink link) {
links.add(link);
return this;
}
/**
* 快捷方法:添加单条TCP链路(最常见场景)
*
* @param terminalIp 终端站IP
* @param terminalPort 终端站端口
* @param dataTransIp 数传本端IP
* @param dataTransPort 数传本端端口
*/
public static TerminalStationInfo tcpSingle(String terminalIp, int terminalPort,
String dataTransIp, int dataTransPort) {
TerminalStationInfo info = new TerminalStationInfo();
info.setAppProtocol(0x01); // FEP
info.setTransportProtocol(0x90); // TCP
info.addLink(TerminalLink.of(0x00, 0x00, terminalIp, terminalPort, dataTransIp, dataTransPort));
return info;
}
/**
* 快捷方法:添加主备双链路TCP(双机场景)
*/
public static TerminalStationInfo tcpDual(
String terminalPrimaryIp, int terminalPrimaryPort,
String terminalBackupIp, int terminalBackupPort,
String dataTransIp, int dataTransPort) {
TerminalStationInfo info = new TerminalStationInfo();
info.setAppProtocol(0x01);
info.setTransportProtocol(0x90);
info.addLink(TerminalLink.of(0x0F, 0x00, terminalPrimaryIp, terminalPrimaryPort, dataTransIp, dataTransPort));
info.addLink(TerminalLink.of(0xF0, 0x00, terminalBackupIp, terminalBackupPort, dataTransIp, dataTransPort));
return info;
}
}
package com.monitor.d2d.service;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.monitor.d2d.config.AppProperties;
import com.monitor.d2d.service.CommandService;
import com.monitor.d2d.netty.ReportEntry;
import com.monitor.d2d.service.ReportService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.ResponseEntity;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.web.bind.annotation.*;
import java.io.IOException;
import java.net.UnknownHostException;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
/**
* 监控接口 REST API
*/
@RestController
@RequestMapping("/api")
@CrossOrigin
@Slf4j
public class ApiController {
@Autowired private CommandService commandService;
@Autowired private ReportService reportService;
@Autowired private AppProperties props;
// ── 系统信息 ─────────────────────────────────────────────────────
@GetMapping("/info")
public Map<String, Object> info() {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sid", props.getSid());
m.put("did", props.getDid());
m.put("sendTarget", props.getSendHost() + ":" + props.getSendPort());
m.put("recvGroup", props.getRecvGroup() + ":" + props.getRecvPort());
m.put("queryInterval", props.getQueryInterval());
m.put("totalReports", reportService.getTotalCount());
m.put("serverTime", LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
return m;
}
// ── 状态查询令 ────────────────────────────────────────────────────
@PostMapping("/query")
public ResponseEntity<?> statusQuery() {
commandService.sendStatusQuery();
return ok("状态查询令已发送");
}
// ── 通用命令发送 ──────────────────────────────────────────────────
/**
* POST /api/cmd/5401 body: {"mode":2}
* POST /api/cmd/5404 body: {"control":1,"planNo":"000000001",...}
*/
@PostMapping("/cmd/{cmdId}")
public ResponseEntity<?> sendCommand(
@PathVariable String cmdId,
@RequestBody(required = false) Map<String, Object> params) throws UnknownHostException {
if (params == null) params = Collections.emptyMap();
commandService.dispatch(cmdId, params);
return ok("命令 [" + cmdId.toUpperCase() + "] 已发送");
}
// ── 上报记录 ──────────────────────────────────────────────────────
@GetMapping("/reports")
public List<Map<String, Object>> reports(
@RequestParam(defaultValue = "50") int limit) {
List<ReportEntry> entries = reportService.getRecent(Math.min(limit, 200));
List<Map<String, Object>> result = new ArrayList<>();
for (ReportEntry e : entries) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("time", e.getRecvTimeStr());
m.put("bidName", e.getBidName());
m.put("sid", String.format("%08X", e.getSid()));
m.put("did", String.format("%08X", e.getDid()));
m.put("bid", String.format("%08X", e.getBid()));
m.put("summary", e.getFrameSummary());
m.put("parsed", e.getParsed());
result.add(m);
}
return result;
}
@DeleteMapping("/reports")
public ResponseEntity<?> clearReports() {
reportService.clear();
return ok("上报记录已清空");
}
// ── 工具 ──────────────────────────────────────────────────────────
private ResponseEntity<Map<String, Object>> ok(String msg) {
Map<String, Object> r = new LinkedHashMap<>();
r.put("success", true);
r.put("message", msg);
r.put("time", LocalDateTime.now().format(DateTimeFormatter.ofPattern("HH:mm:ss.SSS")));
return ResponseEntity.ok(r);
}
private int num = 0;
// ──定时器──────────────────────────────────────────────────────────
@Scheduled(cron = "1 */30 * * * *")
public void autoSend5422() {
String localIP = "127.0.0.1";
String remoteIP = "127.0.0.1";
int concurrency = 4;
try {
for (int i = 0; i < concurrency; i++) {
int count = num * concurrency + i;
ObjectMapper mapper = new ObjectMapper();
String json = "{\"taskControl\":2,\"primaryFlag\":15,\"planNo\":\""+String.format("%09d", count)+"\",\"terminalStation\":\"TEST44621\",\"relayTaskCode\":16671,\"relayAntenna\":15,\"spacecraft\":\"TEST01\",\"orbitNo\":\""+String.format("%05d", count)+"\",\"transferNo\":1,\"captureStartTime\":\"20260701154300\",\"mode\":1,\"terminalAppProtocol\":5,\"terminalTransProtocol\":144,\"linkId\":\"LINK0\",\"initialRecvFrames\":0,\"initialRecvBytes\":0,\"terminalLinks\":[{\"primaryFlag\":15,\"routePlane\":0,\"terminalIp\":\""+remoteIP+"\",\"terminalPort\":27586,\"dataTransIp\":\""+localIP+"\",\"dataTransPort\":1001"+i+"}],\"outerEndpoints\":[{\"centerId\":\"TESTZX\",\"appProtocol\":0,\"initialSendFrames\":0,\"initialSendBytes\":0,\"transportProtocol\":144,\"links\":[{\"dataTransIp\":\""+localIP+"\",\"dataTransPort\":444"+i+",\"outerIp\":\""+remoteIP+"\",\"outerPort\":24584}]}]}";
Map<String, Object> parm = mapper.readValue(json, Map.class);
sendCommand("5422",parm);
}
log.info("执行次数: {}", num++);
} catch (IOException e) {
e.printStackTrace();
}
}
}
\ No newline at end of file
package com.monitor.d2d.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonMappingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.monitor.d2d.config.AppProperties;
import com.monitor.d2d.netty.UdpSender;
import com.monitor.d2d.protocol.*;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import java.io.IOException;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
/**
* 所有命令的构建与发送
* 对应接口文档第6章全部27条命令
*/
@Slf4j
@Service
public class CommandService {
@Autowired private AppProperties props;
@Autowired private UdpSender sender;
/** 状态查询序列号,从1开始,U16循环 */
private final AtomicInteger queryNo = new AtomicInteger(1);
// ────────────────────────────────────────────────────────────────
// 5.1 状态查询令 BID=0x0000F003
// ────────────────────────────────────────────────────────────────
@Scheduled(fixedDelayString = "${d2d.query-interval:10}000",
initialDelayString = "${d2d.query-interval:10}000")
public void autoStatusQuery() {
if (props.getQueryInterval() > 0) {
sendStatusQuery();
}
}
/**
* 手动触发状态查询
*/
public void sendStatusQuery() {
int no = queryNo.getAndUpdate(n -> n >= 0xFFFF ? 1 : n + 1);
ProtoWriter w = new ProtoWriter();
w.writeU16(no);
sendFrame(BidType.STATUS_QUERY, w.toBytes());
// log.info("→ 发送状态查询令 No={}", no);
}
// ────────────────────────────────────────────────────────────────
// 6.2.1 控制方式控制 5401H
// ────────────────────────────────────────────────────────────────
/**
* @param mode 0x01=分控 0x02=本控
*/
public void cmd5401ControlMode(int mode) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.CONTROL_MODE);
w.writeU8(mode);
sendCommand(w);
log.info("→ 控制方式控制 mode={}", mode == 1 ? "分控" : "本控");
}
// ────────────────────────────────────────────────────────────────
// 6.2.2 传输层协议控制 5402H
// ────────────────────────────────────────────────────────────────
/**
* @param protocol 0x01=FEP2.0/UDP 0x02=FEP2.0/TCP 0x03=FEP1.0/UDP 0x04=FEP1.0/TCP
*/
public void cmd5402TransportProtocol(int protocol) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.TRANSPORT_PROTOCOL);
w.writeU8(protocol);
sendCommand(w);
log.info("→ 传输层协议控制 protocol={}", protocol);
}
// ────────────────────────────────────────────────────────────────
// 6.2.3 存储路径控制 5403H
// ────────────────────────────────────────────────────────────────
public void cmd5403StoragePath(String path) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.STORAGE_PATH);
w.writeFixedString(path, 128);
sendCommand(w);
log.info("→ 存储路径控制 path={}", path);
}
// ────────────────────────────────────────────────────────────────
// 6.2.4 数据接收控制 5404H
// ────────────────────────────────────────────────────────────────
/**
* @param control 1=开始 2=停止
* @param planNo 跟踪接收计划号(9位数字)
* @param srcAddr 信源地址(20字节)
* @param spacecraft 航天器标识(6字节)
* @param orbitNo 圈号(10字节)
* @param channelNo 通道号(0xFFFFFFFF表示多通道)
* @param transferNo 传输号
* @param protocol 1=FEP 5=PDXP 6=TCP 7=FTP 8=PDXP/UDP接收 9=PDXP/UDP汇集
* @param mode 1=实战 2=联试
*/
public void cmd5404DataReceive(int control, String planNo, String srcAddr,
String spacecraft, String orbitNo, long channelNo,
long transferNo, int protocol, int mode) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.DATA_RECEIVE);
w.writeU8(control);
w.writeFixedString(planNo, 9);
w.writeFixedString(srcAddr, 20);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeU32(channelNo);
w.writeU32(transferNo);
w.writeU8(protocol);
w.writeU8(mode);
sendCommand(w);
log.info("→ 数据接收控制 control={} planNo={} spacecraft={}", control, planNo, spacecraft);
}
// ────────────────────────────────────────────────────────────────
// 6.2.5 数据接收结束告知 5405H
// ────────────────────────────────────────────────────────────────
/**
* @param endType 0x01=正常结束 0x02=异常结束
*/
public void cmd5405DataReceiveEnd(String planNo, String srcAddr, String spacecraft,
String orbitNo, long channelNo, long transferNo, int endType) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.DATA_RECEIVE_END);
w.writeFixedString(planNo, 9);
w.writeFixedString(srcAddr, 20);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeU32(channelNo);
w.writeU32(transferNo);
w.writeU8(endType);
sendCommand(w);
log.info("→ 数据接收结束告知 planNo={} endType={}", planNo, endType == 1 ? "正常" : "异常");
}
// ────────────────────────────────────────────────────────────────
// 6.2.6 数据发送控制 5406H
// ────────────────────────────────────────────────────────────────
/**
* 简化版(不挑通道、不挑编号、不挑虚拟信道),1个目标IP
*
* @param control 1=开始 2=停止 3=暂停 4=恢复 5=修改IP 6=断点续传 7=即时续传
* @param transferNo 传输号
* @param protocol 0=纯TCP流 1=FEP2.0 2=PDXP 3=其他
* @param maxRateKbps 最大速率
* @param targetIp 目标IP(U32)
* @param targetPort 目标端口
*/
public void cmd5406DataSend(int control, long transferNo, String planNo, String srcAddr,
String spacecraft, String orbitNo, int protocol, String dataCenter,
long maxRateKbps, String targetIp, long targetPort,
boolean pickChannel, long[] channels,
boolean pickVc, int vcMode, int[] vcs) throws UnknownHostException {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.DATA_SEND);
w.writeU8(control);
w.writeU32(transferNo);
w.writeFixedString(planNo, 9);
w.writeFixedString(srcAddr, 20);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeU8(protocol);
w.writeFixedString(dataCenter, 20);
w.writeU32(maxRateKbps);
// 地址个数=1
w.writeU8(1);
w.writeBytes(InetAddress.getByName(targetIp).getAddress());
w.writeU32(targetPort);
// 是否挑通道
w.writeU8(1);
w.writeU8(1);
w.writeU32(Integer.parseInt("311",16));
/*if (pickChannel && channels != null && channels.length > 0) {
w.writeU8(1);
w.writeU8(channels.length);
for (long ch : channels) w.writeU32(ch);
} else {
w.writeU8(0);
}*/
// 是否挑圈次内编号
w.writeU8(0);
// 是否挑虚拟信道
if (vcMode > 0 && vcs != null && vcs.length > 0) {
w.writeU8(vcMode); // 1=挑 2=剔除
w.writeU8(vcs.length);
for (int vc : vcs) w.writeU8(vc);
} else {
w.writeU8(0);
}
sendCommand(w);
log.info("→ 数据发送控制 control={} transferNo={} planNo={}", control, transferNo, planNo);
}
// ────────────────────────────────────────────────────────────────
// 6.2.7 速率调整控制 5407H
// ────────────────────────────────────────────────────────────────
public void cmd5407RateAdjust(long transferNo, String planNo, String srcAddr,
String spacecraft, String orbitNo, String dataCenter,
long maxRateKbps) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.RATE_ADJUST);
w.writeU32(transferNo);
w.writeFixedString(planNo, 9);
w.writeFixedString(srcAddr, 20);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeFixedString(dataCenter, 20);
w.writeU32(maxRateKbps);
sendCommand(w);
log.info("→ 速率调整控制 transferNo={} maxRate={}Kbps", transferNo, maxRateKbps);
}
// ────────────────────────────────────────────────────────────────
// 6.2.8 航天器调整控制 5408H
// ────────────────────────────────────────────────────────────────
public void cmd5408SpacecraftAdjust(int mid) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.SPACECRAFT_ADJUST);
w.writeU16(mid);
sendCommand(w);
log.info("→ 航天器调整控制 MID={}", String.format("%04X", mid));
}
// ────────────────────────────────────────────────────────────────
// 6.2.9 航天器删除控制 5409H
// ────────────────────────────────────────────────────────────────
public void cmd5409SpacecraftDelete(int mid) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.SPACECRAFT_DELETE);
w.writeU16(mid);
sendCommand(w);
log.info("→ 航天器删除控制 MID={}", String.format("%04X", mid));
}
// ────────────────────────────────────────────────────────────────
// 6.2.10 设备调整控制 5410H
// ────────────────────────────────────────────────────────────────
public void cmd5410DeviceAdjust(String deviceId) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.DEVICE_ADJUST);
w.writeFixedString(deviceId, 20);
sendCommand(w);
log.info("→ 设备调整控制 deviceId={}", deviceId);
}
// ────────────────────────────────────────────────────────────────
// 6.2.11 设备删除控制 5411H
// ────────────────────────────────────────────────────────────────
public void cmd5411DeviceDelete(String deviceId) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.DEVICE_DELETE);
w.writeFixedString(deviceId, 20);
sendCommand(w);
log.info("→ 设备删除控制 deviceId={}", deviceId);
}
// ────────────────────────────────────────────────────────────────
// 6.2.12 中心调整控制 5412H
// ────────────────────────────────────────────────────────────────
public void cmd5412CenterAdjust(String centerId) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.CENTER_ADJUST);
w.writeFixedString(centerId, 20);
sendCommand(w);
log.info("→ 中心调整控制 centerId={}", centerId);
}
// ────────────────────────────────────────────────────────────────
// 6.2.13 中心删除控制 5413H
// ────────────────────────────────────────────────────────────────
public void cmd5413CenterDelete(String centerId) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.CENTER_DELETE);
w.writeFixedString(centerId, 20);
sendCommand(w);
log.info("→ 中心删除控制 centerId={}", centerId);
}
// ────────────────────────────────────────────────────────────────
// 6.2.14 软件控制 5414H
// ────────────────────────────────────────────────────────────────
/**
* @param control 1=开启 2=关闭 3=维护
*/
public void cmd5414SoftwareControl(int control) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.SOFTWARE_CONTROL);
w.writeU8(control);
sendCommand(w);
log.info("→ 软件控制 control={}", control == 1 ? "开启" : control == 2 ? "关闭" : "维护");
}
// ────────────────────────────────────────────────────────────────
// 6.2.15 模块控制 5415H
// ────────────────────────────────────────────────────────────────
/**
* @param unitId 软件单元标识(见表6-28)
* @param control 0=启动 1=断开 2=维护 3=取消维护
*/
public void cmd5415ModuleControl(int unitId, int control) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.MODULE_CONTROL);
w.writeU8(unitId);
w.writeU8(control);
sendCommand(w);
log.info("→ 模块控制 unitId={} control={}", String.format("%02X", unitId), control);
}
// ────────────────────────────────────────────────────────────────
// 6.2.16 传输测试控制 5416H
// ────────────────────────────────────────────────────────────────
/**
* @param sendControl 1=开始 2=停止
* @param protocol 1=FEP2.0 2=PDXP
* @param testMode 测试数据选择
* @param maxRateKbps 最大速率
* @param ip 目标IP(U32)
* @param port 目标端口
*/
public void cmd5416TransmissionTest(int sendControl, int protocol, int testMode,
long maxRateKbps, long ip, long port) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.TRANSMISSION_TEST);
w.writeU8(sendControl);
w.writeU8(protocol);
w.writeU8(testMode);
w.writeU32(maxRateKbps);
w.writeU32(ip);
w.writeU32(port);
sendCommand(w);
log.info("→ 传输测试控制 control={} protocol={} ip={}", sendControl, protocol, ip);
}
// ────────────────────────────────────────────────────────────────
// 6.2.17 存储周期控制 5417H
// ────────────────────────────────────────────────────────────────
public void cmd5417StoragePeriod(int days, int percent) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.STORAGE_PERIOD);
w.writeU8(days);
w.writeU8(percent);
sendCommand(w);
log.info("→ 存储周期控制 days={} percent={}%", days, percent);
}
// ────────────────────────────────────────────────────────────────
// 6.2.18 数据删除控制 5418H
// ────────────────────────────────────────────────────────────────
public void cmd5418DataDelete(String deviceId, String spacecraft, String orbitNo,
long channelNo) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.DATA_DELETE);
w.writeFixedString(deviceId, 20);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeU32(channelNo);
w.writeU32(0); // 保留
sendCommand(w);
log.info("→ 数据删除控制 device={} spacecraft={} channel={}", deviceId, spacecraft,
String.format("%08X", channelNo));
}
// ────────────────────────────────────────────────────────────────
// 6.2.22 返向数据任务控制 5422H(完整版,对应文档表6-35)
// ────────────────────────────────────────────────────────────────
/**
* 发送返向数据任务控制命令(完整参数)
*
* @param req 完整命令参数,见 Cmd5422Request / TerminalStationInfo / OuterEndpointInfo
*
* 命令字节布局(大端):
* CMD_ID(U16) 任务控制(U8) 主备标识(U8) 计划号(C9) 终端站标识(C20)
* 中继星任务代号(U16) 中继卫星天线标识(U16) 航天器(C6) 圈号(C10)
* 传输号(U32) 捕跟开始时间(C14,V6.11新增) 模式(U8)
* ── 终端站端 ──
* 应用层协议(U8) 传输层协议(U8) 星间链路标识(C5)
* 初始接收总帧数(U64,V6.11新增) 初始接收总数据量(U64,V6.11新增) 链路数量(U8)
* [每条链路] 主备(U8) 路由平面(U8) 终端站IP(U32) 终端站端口(U16) 数传IP(U32) 数传端口(U16)
* ── 对外端 ──
* 对外端数量(U8)
* [每个对外端] 目的中心标识(C20) 应用层协议(U8) 传输层协议(U8)
* 初始发送总帧数(U64,V6.11新增) 初始发送总数据量(U64,V6.11新增) 链路数量(U8)
* [每条链路] 数传IP(U32) 数传端口(U16) 对外端IP(U32) 对外端端口(U16)
*/
public void cmd5422ReturnDataTask(com.monitor.d2d.protocol.cmd5422.Cmd5422Request req) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.RETURN_DATA_TASK);
// ── 基本参数 ─────────────────────────────────────────────────
w.writeU8(req.getTaskControl());
w.writeU8(req.getPrimaryFlag());
w.writeFixedString(req.getPlanNo(), 9);
w.writeFixedString(req.getTerminalStation(), 20);
w.writeU16(req.getRelayTaskCode());
w.writeU16(req.getRelayAntenna());
w.writeFixedString(req.getSpacecraft(), 6);
w.writeFixedString(req.getOrbitNo(), 10);
w.writeU32(req.getTransferNo());
w.writeFixedString(req.getCaptureStartTime(), 20); // 捕跟开始时间(V6.11新增)
w.writeU8(req.getMode());
// ── 终端站端信息 ─────────────────────────────────────────────
com.monitor.d2d.protocol.cmd5422.TerminalStationInfo ti = req.getTerminalInfo();
if (ti == null) ti = new com.monitor.d2d.protocol.cmd5422.TerminalStationInfo();
w.writeU8(ti.getAppProtocol());
w.writeU8(ti.getTransportProtocol());
w.writeFixedString(ti.getLinkId(), 5);
w.writeU64(ti.getInitialRecvFrames()); // 初始接收总帧数(V6.11新增)
w.writeU64(ti.getInitialRecvBytes()); // 初始接收总数据量(V6.11新增)
java.util.List<com.monitor.d2d.protocol.cmd5422.TerminalLink> tLinks =
ti.getLinks() != null ? ti.getLinks() : java.util.Collections.emptyList();
w.writeU8(tLinks.size());
for (com.monitor.d2d.protocol.cmd5422.TerminalLink tl : tLinks) {
w.writeU8(tl.getPrimaryFlag()); // 链路主备
w.writeU8(tl.getRoutePlane()); // 路由平面
w.writeU32(tl.getTerminalIp()); // 终端站IP
w.writeU16(tl.getTerminalPort()); // 终端站端口
w.writeU32(tl.getDataTransIp()); // 数传IP
w.writeU16(tl.getDataTransPort());// 数传端口
}
// ── 对外端信息 ───────────────────────────────────────────────
java.util.List<com.monitor.d2d.protocol.cmd5422.OuterEndpointInfo> outers =
req.getOuterEndpoints() != null ? req.getOuterEndpoints() : java.util.Collections.emptyList();
w.writeU8(outers.size());
for (com.monitor.d2d.protocol.cmd5422.OuterEndpointInfo oe : outers) {
w.writeFixedString(oe.getCenterId(), 20); // 目的中心标识
w.writeU8(oe.getAppProtocol()); // 应用层协议
w.writeU8(oe.getTransportProtocol()); // 传输层协议
w.writeU64(oe.getInitialSendFrames()); // 初始发送总帧数(V6.11新增)
w.writeU64(oe.getInitialSendBytes()); // 初始发送总数据量(V6.11新增)
java.util.List<com.monitor.d2d.protocol.cmd5422.OuterEndpointLink> oLinks =
oe.getLinks() != null ? oe.getLinks() : java.util.Collections.emptyList();
w.writeU8(oLinks.size());
for (com.monitor.d2d.protocol.cmd5422.OuterEndpointLink ol : oLinks) {
w.writeU32(ol.getDataTransIp()); // 数传本地IP
w.writeU16(ol.getDataTransPort()); // 数传本地端口
w.writeU32(ol.getOuterIp()); // 对外端IP(应用中心)
w.writeU16(ol.getOuterPort()); // 对外端端口
}
}
sendCommand(w);
log.info("→ 返向数据任务控制 taskControl=0x{} primaryFlag=0x{} planNo={} terminalStation={} 终端站链路数={} 对外端数={}",
Integer.toHexString(req.getTaskControl()).toUpperCase(),
Integer.toHexString(req.getPrimaryFlag()).toUpperCase(),
req.getPlanNo(), req.getTerminalStation(),
tLinks.size(), outers.size());
// 若命令中终端站传输层协议为UDP,视为"触发UDP参数":延迟5秒后启动PDXP模拟推送
handleUdpPdxpTrigger(req, ti, tLinks);
}
// ────────────────────────────────────────────────────────────────
// UDP参数触发:cmd5422终端站为UDP传输时,延迟5秒启动模拟PDXP推送
// 绑定本地地址=terminalLinks.terminalIp,目标地址=terminalLinks.dataTransIp:dataTransPort
// ────────────────────────────────────────────────────────────────
/** 传输层协议 U8:0x50=UDP单播 0x60=UDP指定源组播 0x70=UDP任意源组播 */
private boolean isUdpTransport(int transportProtocol) {
return transportProtocol == 0x50 || transportProtocol == 0x60 || transportProtocol == 0x70;
}
private void handleUdpPdxpTrigger(com.monitor.d2d.protocol.cmd5422.Cmd5422Request req,
com.monitor.d2d.protocol.cmd5422.TerminalStationInfo ti,
java.util.List<com.monitor.d2d.protocol.cmd5422.TerminalLink> tLinks) {
if (tLinks == null || tLinks.isEmpty()) return;
// 任务控制 0x04=停止:停止对应的模拟推送任务
boolean isStop = req.getTaskControl() == 0x04;
boolean udp = isUdpTransport(ti.getTransportProtocol());
for (com.monitor.d2d.protocol.cmd5422.TerminalLink tl : tLinks) {
String key = req.getPlanNo() + "_" + req.getTerminalStation()
+ "_" + tl.getTerminalIp() + "_" + tl.getDataTransIp() + "_" + tl.getDataTransPort();
if (isStop) {
com.monitor.d2d.netty.UdpPdxpPusher.stop(key);
continue;
}
if (!udp || tl.getDataTransPort() <= 0) continue;
String bindIp = com.monitor.d2d.protocol.cmd5422.TerminalLink.u32ToIp(tl.getTerminalIp());
String targetIp = com.monitor.d2d.protocol.cmd5422.TerminalLink.u32ToIp(tl.getDataTransIp());
com.monitor.d2d.netty.UdpPdxpPusher.triggerDelayedStart(key, bindIp, targetIp, tl.getDataTransPort());
}
}
/**
* 便捷重载:从 Map 参数中构造 Cmd5422Request 并发送
* 支持单链路简单场景,复杂场景请直接构造 Cmd5422Request 调用上面的方法。
*/
@SuppressWarnings("unchecked")
private void cmd5422FromMap(java.util.Map<String, Object> params) {
com.monitor.d2d.protocol.cmd5422.Cmd5422Request req = new com.monitor.d2d.protocol.cmd5422.Cmd5422Request();
req.setTaskControl(getInt(params, "taskControl", 0x02));
req.setPrimaryFlag(getInt(params, "primaryFlag", 0x0F));
req.setPlanNo(getString(params, "planNo", "000000001"));
req.setTerminalStation(getString(params, "terminalStation", ""));
req.setRelayTaskCode(getInt(params, "relayTaskCode", 0));
req.setRelayAntenna(getInt(params, "relayAntenna", 0));
req.setSpacecraft(getString(params, "spacecraft", ""));
req.setOrbitNo(getString(params, "orbitNo", ""));
req.setTransferNo(getLong(params, "transferNo", 1));
req.setCaptureStartTime(getString(params, "captureStartTime", ""));
req.setMode(getInt(params, "mode", 1));
// 终端站端信息
com.monitor.d2d.protocol.cmd5422.TerminalStationInfo ti = new com.monitor.d2d.protocol.cmd5422.TerminalStationInfo();
ti.setAppProtocol(getInt(params, "terminalAppProtocol", 0x01));
ti.setTransportProtocol(getInt(params, "terminalTransProtocol", 0x90));
ti.setLinkId(getString(params, "linkId", ""));
ti.setInitialRecvFrames(getLong(params, "initialRecvFrames", 0));
ti.setInitialRecvBytes(getLong(params, "initialRecvBytes", 0));
// 支持 terminalLinks 数组:[{primaryFlag,routePlane,terminalIp,terminalPort,dataTransIp,dataTransPort}]
Object rawLinks = params.get("terminalLinks");
if (rawLinks instanceof java.util.List) {
for (Object item : (java.util.List<?>) rawLinks) {
if (item instanceof java.util.Map) {
java.util.Map<String, Object> lm = (java.util.Map<String, Object>) item;
com.monitor.d2d.protocol.cmd5422.TerminalLink tl = new com.monitor.d2d.protocol.cmd5422.TerminalLink();
tl.setPrimaryFlag(getInt(lm, "primaryFlag", 0x00));
tl.setRoutePlane(getInt(lm, "routePlane", 0x00));
tl.setTerminalIp(com.monitor.d2d.protocol.cmd5422.TerminalLink.ipToU32(getString(lm, "terminalIp", "0.0.0.0")));
tl.setTerminalPort(getInt(lm, "terminalPort", 0));
tl.setDataTransIp(com.monitor.d2d.protocol.cmd5422.TerminalLink.ipToU32(getString(lm, "dataTransIp", "0.0.0.0")));
tl.setDataTransPort(getInt(lm, "dataTransPort", 0));
ti.getLinks().add(tl);
}
}
} else {
// 单链路简化参数
String terminalIp = getString(params, "terminalIp", "");
if (!terminalIp.isEmpty()) {
com.monitor.d2d.protocol.cmd5422.TerminalLink tl = new com.monitor.d2d.protocol.cmd5422.TerminalLink();
tl.setPrimaryFlag(getInt(params, "terminalPrimaryFlag", 0x00));
tl.setRoutePlane(getInt(params, "routePlane", 0x00));
tl.setTerminalIp(com.monitor.d2d.protocol.cmd5422.TerminalLink.ipToU32(terminalIp));
tl.setTerminalPort(getInt(params, "terminalPort", 0));
tl.setDataTransIp(com.monitor.d2d.protocol.cmd5422.TerminalLink.ipToU32(getString(params, "dataTransIp", "0.0.0.0")));
tl.setDataTransPort(getInt(params, "dataTransPort", 0));
ti.getLinks().add(tl);
}
}
req.setTerminalInfo(ti);
// 对外端信息:支持 outerEndpoints 数组
Object rawOuters = params.get("outerEndpoints");
if (rawOuters instanceof java.util.List) {
for (Object item : (java.util.List<?>) rawOuters) {
if (item instanceof java.util.Map) {
java.util.Map<String, Object> om = (java.util.Map<String, Object>) item;
com.monitor.d2d.protocol.cmd5422.OuterEndpointInfo oe = new com.monitor.d2d.protocol.cmd5422.OuterEndpointInfo();
oe.setCenterId(getString(om, "centerId", ""));
oe.setAppProtocol(getInt(om, "appProtocol", 0x01));
oe.setTransportProtocol(getInt(om, "transportProtocol", 0x90));
oe.setInitialSendFrames(getLong(om, "initialSendFrames", 0));
oe.setInitialSendBytes(getLong(om, "initialSendBytes", 0));
Object rawOLinks = om.get("links");
if (rawOLinks instanceof java.util.List) {
for (Object li : (java.util.List<?>) rawOLinks) {
if (li instanceof java.util.Map) {
java.util.Map<String, Object> lm = (java.util.Map<String, Object>) li;
com.monitor.d2d.protocol.cmd5422.OuterEndpointLink ol = new com.monitor.d2d.protocol.cmd5422.OuterEndpointLink();
ol.setDataTransIp(com.monitor.d2d.protocol.cmd5422.TerminalLink.ipToU32(getString(lm, "dataTransIp", "0.0.0.0")));
ol.setDataTransPort(getInt(lm, "dataTransPort", 0));
ol.setOuterIp(com.monitor.d2d.protocol.cmd5422.TerminalLink.ipToU32(getString(lm, "outerIp", "0.0.0.0")));
ol.setOuterPort(getInt(lm, "outerPort", 0));
oe.getLinks().add(ol);
}
}
}
req.getOuterEndpoints().add(oe);
}
}
}
cmd5422ReturnDataTask(req);
}
// ────────────────────────────────────────────────────────────────
// 6.2.23 中继卫星参数控制 5423H
// ────────────────────────────────────────────────────────────────
/**
* @param relayTaskCode 中继星任务代号
* @param controlType 1=新增或修改 2=删除
*/
public void cmd5423RelaySatParam(int relayTaskCode, int controlType) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.RELAY_SAT_PARAM);
w.writeU16(relayTaskCode);
w.writeU8(controlType);
sendCommand(w);
log.info("→ 中继卫星参数控制 taskCode={} controlType={}", String.format("%04X", relayTaskCode), controlType);
}
// ────────────────────────────────────────────────────────────────
// 6.2.24 快传任务周期调整控制 5424H
// ────────────────────────────────────────────────────────────────
public void cmd5424FastTransferPeriod(long periodSeconds) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.FAST_TRANSFER_PERIOD);
w.writeU32(periodSeconds);
w.writeU8(1); // 控制=修改
sendCommand(w);
log.info("→ 快传任务周期调整 period={}s", periodSeconds);
}
// ────────────────────────────────────────────────────────────────
// 6.2.25 快传任务调整控制 5425H
// ────────────────────────────────────────────────────────────────
/**
* @param taskRuleId 任务规则ID
* @param control 1=新增或修改 2=删除
*/
public void cmd5425FastTransferTask(long taskRuleId, int control) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.FAST_TRANSFER_TASK);
w.writeU32(taskRuleId);
w.writeU8(control);
sendCommand(w);
log.info("→ 快传任务调整控制 taskRuleId={} control={}", taskRuleId, control);
}
// ────────────────────────────────────────────────────────────────
// 6.2.26 天基误码率控制 5426H
// ────────────────────────────────────────────────────────────────
public void cmd5426SpaceBasedBer(String terminalStation, int relayTaskCode,
String spacecraft, String orbitNo, String linkId,
String startTime, String endTime, String rule) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.SPACE_BER);
w.writeFixedString(terminalStation, 20);
w.writeU16(relayTaskCode);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeFixedString(linkId, 5);
w.writeFixedString(startTime, 14);
w.writeFixedString(endTime, 14);
w.writeFixedString(rule, 20);
sendCommand(w);
log.info("→ 天基误码率控制 terminal={} spacecraft={}", terminalStation, spacecraft);
}
// ────────────────────────────────────────────────────────────────
// 6.2.27 天基任务速率调整控制 5427H
// ────────────────────────────────────────────────────────────────
public void cmd5427TaskAdjust(String planNo, String terminalStation, int relayTaskCode,
String spacecraft, String orbitNo, long transferNo,
String linkId, long speedKbps) {
ProtoWriter w = new ProtoWriter();
w.writeCmdId(CmdId.TASK_ADJUST);
w.writeFixedString(planNo, 9);
w.writeFixedString(terminalStation, 20);
w.writeU16(relayTaskCode);
w.writeFixedString(spacecraft, 6);
w.writeFixedString(orbitNo, 10);
w.writeU32(transferNo);
w.writeFixedString(linkId, 5);
w.writeU64(speedKbps);
sendCommand(w);
log.info("→ 天基任务速率调整 planNo={} speed={}Kbps", planNo, speedKbps);
}
// ────────────────────────────────────────────────────────────────
// 内部发送工具
// ────────────────────────────────────────────────────────────────
private void sendCommand(ProtoWriter w) {
sendFrame(BidType.PROCESS_CONTROL, w.toBytes());
}
private void sendFrame(long bid, byte[] data) {
D2dFrame frame = D2dFrame.create(props.parseSid(), props.parseDid(), bid);
frame.setData(data);
frame.setDataLen(data.length);
sender.send(frame);
}
// ────────────────────────────────────────────────────────────────
// Map参数入口(供REST API调用)
// ────────────────────────────────────────────────────────────────
public void dispatch(String cmdId, Map<String, Object> params) throws UnknownHostException {
switch (cmdId.toUpperCase()) {
case "STATUS_QUERY":
case "F003":
sendStatusQuery(); break;
case "5401":
cmd5401ControlMode(getInt(params, "mode", 2)); break;
case "5402":
cmd5402TransportProtocol(getInt(params, "protocol", 1)); break;
case "5403":
cmd5403StoragePath(getString(params, "path", "/data")); break;
case "5404":
cmd5404DataReceive(
getInt(params, "control", 1),
getString(params, "planNo", "000000001"),
getString(params, "srcAddr", ""),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
Integer.parseInt(getString(params, "channelNo", "0"),16),
getLong(params, "transferNo", 1),
getInt(params, "protocol", 1),
getInt(params, "mode", 1)
); break;
case "5405":
cmd5405DataReceiveEnd(
getString(params, "planNo", "000000001"),
getString(params, "srcAddr", ""),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
Integer.parseInt(getString(params, "channelNo", "0"),16),
getLong(params, "transferNo", 1),
getInt(params, "endType", 1)
); break;
case "5406":
cmd5406DataSend(
getInt(params, "control", 1),
getLong(params, "transferNo", 1),
getString(params, "planNo", "000000001"),
getString(params, "srcAddr", ""),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
getInt(params, "protocol", 1),
getString(params, "dataCenter", ""),
getLong(params, "maxRateKbps", 2000000),
getString(params, "targetIp", ""),
getLong(params, "targetPort", 0),
true, null, false, 0, null
); break;
case "5407":
cmd5407RateAdjust(
getLong(params, "transferNo", 1),
getString(params, "planNo", "000000001"),
getString(params, "srcAddr", ""),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
getString(params, "dataCenter", ""),
getLong(params, "maxRateKbps", 10000)
); break;
case "5408":
cmd5408SpacecraftAdjust(getInt(params, "mid", 0)); break;
case "5409":
cmd5409SpacecraftDelete(getInt(params, "mid", 0)); break;
case "5410":
cmd5410DeviceAdjust(getString(params, "deviceId", "")); break;
case "5411":
cmd5411DeviceDelete(getString(params, "deviceId", "")); break;
case "5412":
cmd5412CenterAdjust(getString(params, "centerId", "")); break;
case "5413":
cmd5413CenterDelete(getString(params, "centerId", "")); break;
case "5414":
cmd5414SoftwareControl(getInt(params, "control", 1)); break;
case "5415":
cmd5415ModuleControl(getInt(params, "unitId", 1), getInt(params, "control", 0)); break;
case "5416":
cmd5416TransmissionTest(
getInt(params, "sendControl", 1),
getInt(params, "protocol", 1),
getInt(params, "testMode", 1),
getLong(params, "maxRateKbps", 10000),
getLong(params, "ip", 0),
getLong(params, "port", 0)
); break;
case "5417":
cmd5417StoragePeriod(getInt(params, "days", 7), getInt(params, "percent", 80)); break;
case "5418":
cmd5418DataDelete(
getString(params, "deviceId", ""),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
getLong(params, "channelNo", 0xFFFFFFFFL)
); break;
case "5422":
cmd5422FromMap(params); break;
case "5423":
cmd5423RelaySatParam(getInt(params, "relayTaskCode", 0), getInt(params, "controlType", 1)); break;
case "5424":
cmd5424FastTransferPeriod(getLong(params, "period", 60)); break;
case "5425":
cmd5425FastTransferTask(getLong(params, "taskRuleId", 1), getInt(params, "control", 1)); break;
case "5426":
cmd5426SpaceBasedBer(
getString(params, "terminalStation", ""),
getInt(params, "relayTaskCode", 0),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
getString(params, "linkId", ""),
getString(params, "startTime", ""),
getString(params, "endTime", ""),
getString(params, "rule", "")
); break;
case "5427":
cmd5427TaskAdjust(
getString(params, "planNo", ""),
getString(params, "terminalStation", ""),
getInt(params, "relayTaskCode", 0),
getString(params, "spacecraft", ""),
getString(params, "orbitNo", ""),
getLong(params, "transferNo", 1),
getString(params, "linkId", ""),
getLong(params, "speedKbps", 10000)
); break;
default:
log.warn("未知命令: {}", cmdId);
}
}
private int getInt(Map<String, Object> m, String key, int def) {
Object v = m.get(key);
if (v == null) return def;
if (v instanceof Number) return ((Number)v).intValue();
try { return Integer.parseInt(v.toString()); } catch (Exception e) { return def; }
}
private long getLong(Map<String, Object> m, String key, long def) {
Object v = m.get(key);
if (v == null) return def;
if (v instanceof Number) return ((Number)v).longValue();
try { return Long.parseLong(v.toString()); } catch (Exception e) { return def; }
}
private String getString(Map<String, Object> m, String key, String def) {
Object v = m.get(key);
return v != null ? v.toString() : def;
}
}
\ No newline at end of file
package com.monitor.d2d.service;
import com.monitor.d2d.netty.ReportEntry;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Deque;
import java.util.List;
import java.util.concurrent.ConcurrentLinkedDeque;
import java.util.concurrent.atomic.AtomicLong;
/**
* 上报帧内存存储(固定容量,先进先出)
*/
@Service
public class ReportService {
private static final int MAX_SIZE = 500;
private final Deque<ReportEntry> store = new ConcurrentLinkedDeque<>();
private final AtomicLong totalCount = new AtomicLong(0);
public void add(ReportEntry entry) {
store.addFirst(entry);
totalCount.incrementAndGet();
while (store.size() > MAX_SIZE) {
store.pollLast();
}
}
/**
* 获取最近 n 条,按时间倒序(最新在前)
*/
public List<ReportEntry> getRecent(int n) {
List<ReportEntry> list = new ArrayList<>(store);
return list.subList(0, Math.min(n, list.size()));
}
public long getTotalCount() {
return totalCount.get();
}
public void clear() {
store.clear();
}
}
\ No newline at end of file
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment