SocketStream
This commit is contained in:
@@ -0,0 +1 @@
|
||||
#include <c_Socket.h>
|
||||
@@ -0,0 +1,40 @@
|
||||
#ifndef INCLUDED_C_SOCKET_H
|
||||
#define INCLUDED_C_SOCKET_H
|
||||
|
||||
#ifndef INCLUDED_C_TYPES_H
|
||||
#include <c_Types.h>
|
||||
#endif /*INCLUDED_C_TYPES_H*/
|
||||
|
||||
|
||||
/* ------------------------------------------------------------------------------------------------------------------ */
|
||||
/* */
|
||||
|
||||
|
||||
#if defined(_WIN32) || defined(_WIN64)
|
||||
#include <winsock2.h>
|
||||
typedef SOCKET c_socket_t;
|
||||
#define C_INVALID_SOCKET INVALID_SOCKET
|
||||
#define C_SOCKET_ERROR SOCKET_ERROR
|
||||
#define c_Socket_Close(x) closesocket(x)
|
||||
|
||||
C_STATIC_FORCE_INLINE
|
||||
void c_Socket_Init(void) {
|
||||
WSADATA wsa; WSAStartup(MAKEWORD(2,2), &wsa);
|
||||
}
|
||||
|
||||
C_STATIC_FORCE_INLINE
|
||||
void c_Socket_Destroy(void) {
|
||||
WSACleanup();
|
||||
}
|
||||
#else
|
||||
typedef int c_socket_t;
|
||||
#define C_INVALID_SOCKET (-1)
|
||||
#define C_SOCKET_ERROR (-1)
|
||||
#define c_Socket_Close(x) close(x)
|
||||
#define c_Socket_Init()
|
||||
#define c_Socket_Destroy()
|
||||
#endif
|
||||
|
||||
|
||||
|
||||
#endif /*INCLUDED_C_SOCKET_H*/
|
||||
@@ -0,0 +1,125 @@
|
||||
#include <c_SocketStream.h>
|
||||
|
||||
#if defined(_WIN32) || defined(_WIN64)
|
||||
#include <winsock2.h>
|
||||
#else
|
||||
#include <sys/socket.h>
|
||||
#include <unistd.h>
|
||||
#endif
|
||||
|
||||
/* ================================================================================================================== */
|
||||
/* [套接字输入流子类]: c_SocketInStream 内部具体化上下文 */
|
||||
|
||||
typedef struct {
|
||||
c_InStream_t base;
|
||||
c_socket_t sock;
|
||||
c_Allocator_t allocator;
|
||||
} c_SocketInImpl_t;
|
||||
|
||||
static c_err_t _SocketIn_Read(c_InStream_t* self, void* buf, c_size_t len, c_size_t* bytes_read) {
|
||||
c_SocketInImpl_t* impl = (c_SocketInImpl_t*)self;
|
||||
if (impl->sock == C_INVALID_SOCKET || !buf || len == 0) return C_ERR_PARAM;
|
||||
|
||||
/* 执行底层网络字节流套接字拦截 */
|
||||
#if defined(_WIN32) || defined(_WIN64)
|
||||
int n = recv(impl->sock, (char*)buf, (int)len, 0);
|
||||
#else
|
||||
ssize_t n = recv(impl->sock, buf, (size_t)len, 0);
|
||||
#endif
|
||||
|
||||
if (n == 0) {
|
||||
if (bytes_read) *bytes_read = 0;
|
||||
return C_ERR_OUTOFBOUND; /* 对端优雅关闭连接 */
|
||||
}
|
||||
if (n == C_SOCKET_ERROR) {
|
||||
if (bytes_read) *bytes_read = 0;
|
||||
return C_ERR_FAIL; /* 网络链路异常断开 */
|
||||
}
|
||||
|
||||
if (bytes_read) *bytes_read = (c_size_t)n;
|
||||
return ( (c_size_t)n == len ) ? C_SUCCESS : C_ERR_OK; /* 部分读取或完全读取皆视为合规 */
|
||||
}
|
||||
|
||||
static void _SocketIn_Destroy(c_InStream_t* self) {
|
||||
c_SocketInImpl_t* impl = (c_SocketInImpl_t*)self;
|
||||
c_Allocator_t alloc = impl->allocator;
|
||||
c_Allocator_Free(&alloc, impl);
|
||||
}
|
||||
|
||||
static const c_InStreamVtbl_t g_SocketInVtbl = { _SocketIn_Read, _SocketIn_Destroy };
|
||||
|
||||
c_err_t c_SocketInStream_Create(c_InStream_t** out_stream, c_socket_t sock, c_Allocator_t* allocator) {
|
||||
if (!out_stream || sock == C_INVALID_SOCKET) return C_ERR_PARAM;
|
||||
c_Allocator_t alloc = allocator ? *allocator : c_DefaultAllocator;
|
||||
|
||||
c_SocketInImpl_t* impl = (c_SocketInImpl_t*)c_Allocator_Calloc(&alloc, 1, sizeof(c_SocketInImpl_t));
|
||||
if (!impl) return C_ERR_NOMEM;
|
||||
|
||||
impl->base.vtbl = &g_SocketInVtbl;
|
||||
impl->sock = sock;
|
||||
impl->allocator = alloc;
|
||||
*out_stream = &impl->base;
|
||||
return C_SUCCESS;
|
||||
}
|
||||
|
||||
/* ================================================================================================================== */
|
||||
/* [套接字输出流子类]: c_SocketOutStream 内部具体化上下文 */
|
||||
|
||||
typedef struct {
|
||||
c_OutStream_t base;
|
||||
c_socket_t sock;
|
||||
c_Allocator_t allocator;
|
||||
} c_SocketOutImpl_t;
|
||||
|
||||
static c_err_t _SocketOut_Write(c_OutStream_t* self, const void* buf, c_size_t len, c_size_t* bytes_written) {
|
||||
c_SocketOutImpl_t* impl = (c_SocketOutImpl_t*)self;
|
||||
if (impl->sock == C_INVALID_SOCKET || !buf || len == 0) return C_ERR_PARAM;
|
||||
|
||||
#if defined(_WIN32) || defined(_WIN64)
|
||||
int n = send(impl->sock, (const char*)buf, (int)len, 0);
|
||||
#else
|
||||
/* MSG_NOSIGNAL 能够有效防止 Unix 系统在往断开的连接写入时触发致命的 SIGPIPE 信号崩溃进程 */
|
||||
#ifdef MSG_NOSIGNAL
|
||||
ssize_t n = send(impl->sock, buf, (size_t)len, MSG_NOSIGNAL);
|
||||
#else
|
||||
ssize_t n = send(impl->sock, buf, (size_t)len, 0);
|
||||
#endif
|
||||
#endif
|
||||
|
||||
if (n == C_SOCKET_ERROR) {
|
||||
if (bytes_written) *bytes_written = 0;
|
||||
return C_ERR_FAIL;
|
||||
}
|
||||
|
||||
if (bytes_written) *bytes_written = (c_size_t)n;
|
||||
return ( (c_size_t)n == len ) ? C_SUCCESS : C_ERR_FAIL;
|
||||
}
|
||||
|
||||
static c_err_t _SocketOut_Flush(c_OutStream_t* self) {
|
||||
/* TCP 套接字在操作系统内核层通常是实时刷新或通过 Nagle 算法控制,上层虚函数表在此可直接视为合规 */
|
||||
(void)self;
|
||||
return C_SUCCESS;
|
||||
}
|
||||
|
||||
static void _SocketOut_Destroy(c_OutStream_t* self) {
|
||||
c_SocketOutImpl_t* impl = (c_SocketOutImpl_t*)self;
|
||||
c_Allocator_t alloc = impl->allocator;
|
||||
c_Allocator_Free(&alloc, impl);
|
||||
}
|
||||
|
||||
static const c_OutStreamVtbl_t g_SocketOutVtbl = { _SocketOut_Write, _SocketOut_Flush, _SocketOut_Destroy };
|
||||
|
||||
c_err_t c_SocketOutStream_Create(c_OutStream_t** out_stream, c_socket_t sock, c_Allocator_t* allocator) {
|
||||
if (!out_stream || sock == C_INVALID_SOCKET) return C_ERR_PARAM;
|
||||
c_Allocator_t alloc = allocator ? *allocator : c_DefaultAllocator;
|
||||
|
||||
c_SocketOutImpl_t* impl = (c_SocketOutImpl_t*)c_Allocator_Calloc(&alloc, 1, sizeof(c_SocketOutImpl_t));
|
||||
if (!impl) return C_ERR_NOMEM;
|
||||
|
||||
impl->base.vtbl = &g_SocketOutVtbl;
|
||||
impl->sock = sock;
|
||||
impl->allocator = alloc;
|
||||
*out_stream = &impl->base;
|
||||
return C_SUCCESS;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
#ifndef INCLUDED_C_SOCKETSTREAM_H
|
||||
#define INCLUDED_C_SOCKETSTREAM_H
|
||||
|
||||
#ifndef INCLUDED_C_SOCKET_H
|
||||
#include <c_Socket.h>
|
||||
#endif /*INCLUDED_C_SOCKET_H*/
|
||||
|
||||
|
||||
#ifndef INCLUDED_C_INSTREAM_H
|
||||
#include <c_InStream.h>
|
||||
#endif /*INCLUDED_C_INSTREAM_H*/
|
||||
|
||||
|
||||
#ifndef INCLUDED_C_OUTSTREAM_H
|
||||
#include <c_OutStream.h>
|
||||
#endif /*INCLUDED_C_OUTSTREAM_H*/
|
||||
|
||||
|
||||
#ifndef INCLUDED_C_ALLOCATOR_H
|
||||
#include <c_Allocator.h>
|
||||
#endif /*INCLUDED_C_ALLOCATOR_H*/
|
||||
|
||||
|
||||
/* ------------------------------------------------------------------------------------------------------------------ */
|
||||
/* */
|
||||
|
||||
/**
|
||||
* @brief 实例化一个绑定到网络套接字的 polymorphic 输入流读取器
|
||||
* @param out_stream 成功后接收 c_InStream_t 抽象接口指针
|
||||
* @param sock 有效的且已建立连接的底层套接字句柄描述符
|
||||
*/
|
||||
c_err_t c_SocketInStream_Create(c_InStream_t** out_stream, c_socket_t sock, c_Allocator_t* allocator);
|
||||
|
||||
/**
|
||||
* @brief 实例化一个绑定到网络套接字的 polymorphic 输出流写入器
|
||||
* @param out_stream 成功后接收 c_OutStream_t 抽象接口指针
|
||||
* @param sock 有效的且已建立连接的底层套接字句柄描述符
|
||||
*/
|
||||
c_err_t c_SocketOutStream_Create(c_OutStream_t** out_stream, c_socket_t sock, c_Allocator_t* allocator);
|
||||
|
||||
|
||||
|
||||
#endif /*INCLUDED_C_SOCKETSTREAM_H*/
|
||||
@@ -0,0 +1,81 @@
|
||||
#include "c_Test.h"
|
||||
#include "c_SocketStream.h"
|
||||
#include "c_BinaryOutStream.h"
|
||||
#include "c_BinaryInStream.h"
|
||||
|
||||
/* 模拟简单的跨平台 Socket Pair 创建逻辑 */
|
||||
static void CreateLocalSocketPair(c_socket_t* listener, c_socket_t* client, c_socket_t* server) {
|
||||
c_Socket_Init();
|
||||
*listener = socket(AF_INET, SOCK_STREAM, 0);
|
||||
struct sockaddr_in addr;
|
||||
memset(&addr, 0, sizeof(addr));
|
||||
addr.sin_family = AF_INET;
|
||||
addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK);
|
||||
addr.sin_port = 0; /* 让 OS 自适应分派空闲端口 */
|
||||
|
||||
bind(*listener, (struct sockaddr*)&addr, sizeof(addr));
|
||||
listen(*listener, 1);
|
||||
|
||||
int addr_len = sizeof(addr);
|
||||
getsockname(*listener, (struct sockaddr*)&addr, &addr_len);
|
||||
|
||||
*client = socket(AF_INET, SOCK_STREAM, 0);
|
||||
connect(*client, (struct sockaddr*)&addr, sizeof(addr));
|
||||
*server = accept(*listener, NULL, NULL);
|
||||
}
|
||||
|
||||
TEST_CASE(test_polymorphic_socket_streams_network_pipeline) {
|
||||
c_socket_t listener, client_sock, server_sock;
|
||||
CreateLocalSocketPair(&listener, &client_sock, &server_sock);
|
||||
|
||||
/* -------------------------------------------------------------------------------------------------------------- */
|
||||
/* 发送方逻辑:构建网络流并将二进制打包器拼装到虚函数表抽象句柄上 */
|
||||
c_OutStream_t* net_out = NULL;
|
||||
c_err_t err = c_SocketOutStream_Create(&net_out, client_sock, NULL);
|
||||
ASSERT_INT_EQ(C_SUCCESS, err);
|
||||
|
||||
c_BinaryOutStream_t bin_out;
|
||||
c_BinaryOutStream_Init(&bin_out, net_out);
|
||||
|
||||
/* 写入不常规宽度的非对齐 bit 流数据 */
|
||||
c_BinaryOutStream_WriteBits(&bin_out, 0x05, 3); /* 3 bits */
|
||||
c_BinaryOutStream_WriteUInt32(&bin_out, 0xABCDEF12); /* 32 bits */
|
||||
c_BinaryOutStream_Flush(&bin_out);
|
||||
|
||||
/* -------------------------------------------------------------------------------------------------------------- */
|
||||
/* 接收方逻辑:接收网络多态流并解包验证 */
|
||||
c_InStream_t* net_in = NULL;
|
||||
err = c_SocketInStream_Create(&net_in, server_sock, NULL);
|
||||
ASSERT_INT_EQ(C_SUCCESS, err);
|
||||
|
||||
c_BinaryInStream_t bin_in;
|
||||
c_BinaryInStream_Init(&bin_in, net_in);
|
||||
|
||||
uint32_t val_bits = 0;
|
||||
err = c_BinaryInStream_ReadBits(&bin_in, 3, &val_bits);
|
||||
ASSERT_INT_EQ(C_SUCCESS, err);
|
||||
ASSERT_LL_EQ(0x05, val_bits);
|
||||
|
||||
uint32_t val_u32 = 0;
|
||||
err = c_BinaryInStream_ReadUInt32(&bin_in, &val_u32);
|
||||
ASSERT_INT_EQ(C_SUCCESS, err);
|
||||
ASSERT_LL_EQ(0xABCDEF12, val_u32);
|
||||
|
||||
/* -------------------------------------------------------------------------------------------------------------- */
|
||||
/* 清理所有多态虚函数组件及底层网络套接字句柄描述符 */
|
||||
c_OutStream_Destroy(net_out);
|
||||
c_InStream_Destroy(net_in);
|
||||
|
||||
c_Socket_Close(client_sock);
|
||||
c_Socket_Close(server_sock);
|
||||
c_Socket_Close(listener);
|
||||
|
||||
c_Socket_Destroy();
|
||||
}
|
||||
|
||||
int main(void) {
|
||||
TEST_START(Polymorphic_SocketStream_Network_Pipeline_Suite);
|
||||
RUN_TEST(test_polymorphic_socket_streams_network_pipeline);
|
||||
TEST_REPORT();
|
||||
RETURN_TEST_STATUS;
|
||||
}
|
||||
Reference in New Issue
Block a user