network chunk v4
This commit is contained in:
@@ -4,7 +4,6 @@ import { MessageHandler } from '../message_handler';
|
|||||||
import { SocketCommunicatorBase } from './socket_communicator_base';
|
import { SocketCommunicatorBase } from './socket_communicator_base';
|
||||||
import { OperationHandler } from '../operations_base/operation_handler';
|
import { OperationHandler } from '../operations_base/operation_handler';
|
||||||
|
|
||||||
const END_OF_MESSAGE = '<EOM>'; // Define a unique marker for end of message
|
|
||||||
const CHUNK_SIZE = 1024; // Define chunk size
|
const CHUNK_SIZE = 1024; // Define chunk size
|
||||||
|
|
||||||
export class TcpServerCommunicator extends SocketCommunicatorBase {
|
export class TcpServerCommunicator extends SocketCommunicatorBase {
|
||||||
@@ -13,7 +12,6 @@ export class TcpServerCommunicator extends SocketCommunicatorBase {
|
|||||||
private publicKey: string | null;
|
private publicKey: string | null;
|
||||||
private aesKey: Buffer | null;
|
private aesKey: Buffer | null;
|
||||||
private aesIv: Buffer | null;
|
private aesIv: Buffer | null;
|
||||||
private messageBuffer: string;
|
|
||||||
private chunkBuffers: { [messageId: string]: string[] }; // Buffer for reassembling chunks
|
private chunkBuffers: { [messageId: string]: string[] }; // Buffer for reassembling chunks
|
||||||
|
|
||||||
constructor(socket: Socket, ip: string, port: number, operationHandler: OperationHandler) {
|
constructor(socket: Socket, ip: string, port: number, operationHandler: OperationHandler) {
|
||||||
@@ -23,7 +21,6 @@ export class TcpServerCommunicator extends SocketCommunicatorBase {
|
|||||||
this.publicKey = null; // Client public key will be set later
|
this.publicKey = null; // Client public key will be set later
|
||||||
this.aesKey = null;
|
this.aesKey = null;
|
||||||
this.aesIv = null;
|
this.aesIv = null;
|
||||||
this.messageBuffer = ''; // Initialize the message buffer
|
|
||||||
this.chunkBuffers = {}; // Buffer for reassembling incoming messages
|
this.chunkBuffers = {}; // Buffer for reassembling incoming messages
|
||||||
this.generateKeyPair(); // Generate RSA key pair for encryption
|
this.generateKeyPair(); // Generate RSA key pair for encryption
|
||||||
}
|
}
|
||||||
@@ -97,37 +94,37 @@ export class TcpServerCommunicator extends SocketCommunicatorBase {
|
|||||||
|
|
||||||
// Handle incoming chunks of data
|
// Handle incoming chunks of data
|
||||||
async handleIncomingChunk(data: Buffer): Promise<void> {
|
async handleIncomingChunk(data: Buffer): Promise<void> {
|
||||||
const incomingMessage = data.toString();
|
const incomingMessage = data.toString().trim();
|
||||||
this.messageBuffer += incomingMessage;
|
|
||||||
|
|
||||||
if (this.messageBuffer.includes(END_OF_MESSAGE)) {
|
// Extract header and chunk content from incoming data
|
||||||
const messages = this.messageBuffer.split(END_OF_MESSAGE);
|
const [headerJson, chunkContent] = incomingMessage.split('|');
|
||||||
|
|
||||||
for (let i = 0; i < messages.length - 1; i++) {
|
|
||||||
const completeMessage = messages[i];
|
|
||||||
if (completeMessage) {
|
|
||||||
this.processCompleteMessage(completeMessage);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
this.messageBuffer = messages[messages.length - 1];
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private processCompleteMessage(completeMessage: string): void {
|
|
||||||
const [headerJson, chunkContent] = completeMessage.split('|');
|
|
||||||
const header = JSON.parse(headerJson);
|
const header = JSON.parse(headerJson);
|
||||||
|
|
||||||
|
// Initialize chunk array if this is the first chunk for this messageId
|
||||||
if (!this.chunkBuffers[header.messageId]) {
|
if (!this.chunkBuffers[header.messageId]) {
|
||||||
this.chunkBuffers[header.messageId] = new Array(header.totalChunks);
|
this.chunkBuffers[header.messageId] = new Array(header.totalChunks);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Place the chunk in the correct position in the chunk buffer
|
||||||
this.chunkBuffers[header.messageId][header.sequenceNumber] = chunkContent;
|
this.chunkBuffers[header.messageId][header.sequenceNumber] = chunkContent;
|
||||||
|
if(header.sequenceNumber === header.totalChunks) {
|
||||||
|
await this.processCompleteMessage();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (this.chunkBuffers[header.messageId].every((chunk) => chunk !== undefined)) {
|
// Process the complete message when EOM is received
|
||||||
const fullMessage = this.chunkBuffers[header.messageId].join('');
|
private async processCompleteMessage(): Promise<void> {
|
||||||
this.handleIncomingMessage(fullMessage);
|
for (const [messageId, chunks] of Object.entries(this.chunkBuffers)) {
|
||||||
delete this.chunkBuffers[header.messageId];
|
if (chunks.every((chunk) => chunk !== undefined)) {
|
||||||
|
// Join all chunks to form the full message
|
||||||
|
const fullMessage = chunks.join('');
|
||||||
|
|
||||||
|
// Handle the completed and possibly decrypted message
|
||||||
|
await this.handleIncomingMessage(fullMessage);
|
||||||
|
|
||||||
|
// Clean up buffer after processing
|
||||||
|
delete this.chunkBuffers[messageId];
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -154,29 +151,25 @@ export class TcpServerCommunicator extends SocketCommunicatorBase {
|
|||||||
const chunk = outgoingMessage.slice(i * CHUNK_SIZE, (i + 1) * CHUNK_SIZE);
|
const chunk = outgoingMessage.slice(i * CHUNK_SIZE, (i + 1) * CHUNK_SIZE);
|
||||||
const chunkHeader = JSON.stringify({
|
const chunkHeader = JSON.stringify({
|
||||||
messageId,
|
messageId,
|
||||||
sequenceNumber: i,
|
sequenceNumber: i + 1,
|
||||||
totalChunks,
|
totalChunks,
|
||||||
});
|
});
|
||||||
|
|
||||||
const chunkWithHeader = `${chunkHeader}|${chunk}`;
|
const chunkWithHeader = `${chunkHeader}|${chunk}`;
|
||||||
|
|
||||||
await this.writeToSocket(chunkWithHeader);
|
await this.writeToSocket(chunkWithHeader);
|
||||||
|
|
||||||
if (i === totalChunks - 1) {
|
|
||||||
await this.writeToSocket(END_OF_MESSAGE);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Handle incoming message (decrypt with AES if available)
|
// Handle incoming message (decrypt with AES if available)
|
||||||
handleIncomingMessage(incomingMessage: string): void {
|
async handleIncomingMessage(incomingMessage: string): Promise<void> {
|
||||||
let messageToProcess = incomingMessage;
|
let messageToProcess = incomingMessage;
|
||||||
|
|
||||||
if (this.aesKey && this.aesIv) {
|
if (this.aesKey && this.aesIv) {
|
||||||
messageToProcess = this.decryptWithAes(incomingMessage);
|
messageToProcess = this.decryptWithAes(incomingMessage);
|
||||||
}
|
}
|
||||||
|
|
||||||
this.handlerResult = this.operationHandler.handleOperation(messageToProcess);
|
this.handlerResult = await this.operationHandler.handleOperation(messageToProcess);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Write message to socket
|
// Write message to socket
|
||||||
|
|||||||
@@ -11,7 +11,6 @@ export class TcpClientCommunicator extends SocketCommunicatorBase {
|
|||||||
private readonly socket: Socket;
|
private readonly socket: Socket;
|
||||||
private aesKey: Buffer | null;
|
private aesKey: Buffer | null;
|
||||||
private aesIv: Buffer | null;
|
private aesIv: Buffer | null;
|
||||||
private messageBuffer: string;
|
|
||||||
private serverPublicKey: string | null;
|
private serverPublicKey: string | null;
|
||||||
private isAesKeySetFlag: boolean;
|
private isAesKeySetFlag: boolean;
|
||||||
private chunkBuffers: { [messageId: string]: string[] }; // Buffer for reassembling chunks
|
private chunkBuffers: { [messageId: string]: string[] }; // Buffer for reassembling chunks
|
||||||
@@ -21,8 +20,7 @@ export class TcpClientCommunicator extends SocketCommunicatorBase {
|
|||||||
this.socket = socket;
|
this.socket = socket;
|
||||||
this.aesKey = null;
|
this.aesKey = null;
|
||||||
this.aesIv = null;
|
this.aesIv = null;
|
||||||
this.serverPublicKey = null;
|
this.serverPublicKey = null;c
|
||||||
this.messageBuffer = ''; // Buffer for message reassembly
|
|
||||||
this.isAesKeySetFlag = false;
|
this.isAesKeySetFlag = false;
|
||||||
this.chunkBuffers = {}; // Initialize chunk buffer for reassembling messages
|
this.chunkBuffers = {}; // Initialize chunk buffer for reassembling messages
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,7 +12,6 @@ export class TcpServerCommunicator extends SocketCommunicatorBase {
|
|||||||
private publicKey: string | null;
|
private publicKey: string | null;
|
||||||
private aesKey: Buffer | null;
|
private aesKey: Buffer | null;
|
||||||
private aesIv: Buffer | null;
|
private aesIv: Buffer | null;
|
||||||
private messageBuffer: string;
|
|
||||||
private chunkBuffers: { [messageId: string]: string[] }; // Buffer for reassembling chunks
|
private chunkBuffers: { [messageId: string]: string[] }; // Buffer for reassembling chunks
|
||||||
|
|
||||||
constructor(socket: Socket, ip: string, port: number, operationHandler: OperationHandler) {
|
constructor(socket: Socket, ip: string, port: number, operationHandler: OperationHandler) {
|
||||||
@@ -22,7 +21,6 @@ export class TcpServerCommunicator extends SocketCommunicatorBase {
|
|||||||
this.publicKey = null; // Client public key will be set later
|
this.publicKey = null; // Client public key will be set later
|
||||||
this.aesKey = null;
|
this.aesKey = null;
|
||||||
this.aesIv = null;
|
this.aesIv = null;
|
||||||
this.messageBuffer = ''; // Initialize the message buffer
|
|
||||||
this.chunkBuffers = {}; // Buffer for reassembling incoming messages
|
this.chunkBuffers = {}; // Buffer for reassembling incoming messages
|
||||||
this.generateKeyPair(); // Generate RSA key pair for encryption
|
this.generateKeyPair(); // Generate RSA key pair for encryption
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user