fix index

This commit is contained in:
pigeatgarlic 2025-02-13 16:56:33 +00:00
parent 123bd31e5c
commit f8c2fa81ec
4 changed files with 40 additions and 118 deletions

View File

@ -1,31 +1,5 @@
#define QUEUE_SIZE 128
#ifdef _WIN32
#define QUEUE_SIZE 16
#define PACKET_SIZE 512 * 1024
#else
#define PACKET_SIZE 1024 * 1024
#endif
enum QueueType {
Video,
Audio,
QueueMax
};
typedef enum _EventType {
Pointer,
Bitrate,
Framerate,
Idr,
Hdr,
Stop,
EventMax
} EventType;
typedef struct {
int is_idr;
long long duration;
}PacketMetadata;
typedef struct {
int active;
@ -33,7 +7,6 @@ typedef struct {
int env_width, env_height;
int width, height;
// Offset x and y coordinates of the client
float client_offsetX, client_offsetY;
float offsetX, offsetY;
@ -41,32 +14,16 @@ typedef struct {
}QueueMetadata;
typedef struct {
int size;
PacketMetadata metadata;
int size;
char data[PACKET_SIZE];
} Packet;
typedef enum _DataType {
HDR_INFO,
NUMBER,
STRING,
} DataType;
typedef struct {
int read;
DataType type;
int data_size;
int value_number;
char value_raw[PACKET_SIZE];
} Event;
typedef struct _Queue{
int index;
int inindex;
int outindex;
void* handle;
QueueMetadata metadata;
Event events[EventMax];
Packet array[QUEUE_SIZE];
}Queue;
typedef struct {
Queue queues[QueueMax];
}SharedMemory;
Packet incoming[QUEUE_SIZE];
Packet outcoming[QUEUE_SIZE];
}Queue;

View File

@ -14,43 +14,6 @@
#pragma comment(lib, "user32.lib")
#define BUF_SIZE 256
void
push_packet(Queue* queue,
void* data,
int size,
PacketMetadata metadata){
// wait while queue is full
auto new_index = queue->index + 1;
auto real_index = new_index % QUEUE_SIZE;
Packet* block = &queue->array[real_index];
memcpy(block->data,data,size);
block->size = size;
block->metadata = metadata;
//always update index after write data
queue->index = new_index;
}
void
raise_event(Queue* queue, EventType type, Event event){
event.read = false;
memcpy(&queue->events[type],&event,sizeof(Event));
}
int
peek_event(Queue* memory, EventType type){
return !memory->events[type].read;
}
Event
pop_event(Queue* queue, EventType type){
queue->events[type].read = true;
BOOST_LOG(debug) << "Receive event " << type << ", value: "<< queue->events[type].value_number;
return queue->events[type];
}
void*
@ -82,8 +45,8 @@ map_file(char* name)
return pBuf;
}
SharedMemory*
Queue*
init_shared_memory(char* data){
SharedMemory* memory = (SharedMemory*)map_file(data);
Queue* memory = (Queue*)map_file(data);
return memory;
}

View File

@ -10,17 +10,5 @@
#include <smemory.h>
SharedMemory*
init_shared_memory(char* data);
void
push_packet(Queue* memory, void* data, int size, PacketMetadata metadata);
void
raise_event(Queue* memory, EventType type, Event event);
int
peek_event(Queue* memory, EventType type);
Event
pop_event(Queue* memory, EventType type);
Queue*
init_shared_memory(char* data);

View File

@ -33,6 +33,19 @@ enum StatusCode {
NORMAL_EXIT,
NO_ENCODER_AVAILABLE = 77
};
enum QueueType {
Video,
Audio,
};
enum EventType {
Pointer,
Bitrate,
Framerate,
Idr,
Hdr,
Stop,
EventMax
};
using namespace std::literals;
using namespace boost::asio::ip;
@ -157,7 +170,7 @@ main(int argc, char *argv[]) {
} else if(queuetype == QueueType::Video && video::probe_encoders()) {
BOOST_LOG(error) << "Video failed to find working encoder"sv;
return StatusCode::NO_ENCODER_AVAILABLE;
} else if (memory == nullptr) {
} else if (queue == nullptr) {
BOOST_LOG(error) << "Failed to find shared memory"sv;
return StatusCode::NO_ENCODER_AVAILABLE;
}
@ -228,8 +241,6 @@ main(int argc, char *argv[]) {
auto last_timestamp = std::chrono::high_resolution_clock::now().time_since_epoch().count();
bool first_video_packet = true;
auto buffer = (char*)malloc(5*1024*1024);
uint32_t index = 0;
while (!process_shutdown_event->peek() && !local_shutdown->peek()) {
if (queue_type == QueueType::Video) {
do {
@ -244,12 +255,14 @@ main(int argc, char *argv[]) {
size_t size = packet->data_size();
auto duration = uint32_t(timestamp - last_timestamp);;
memcpy(buffer,&index,sizeof(uint32_t));
memcpy(buffer + sizeof(uint32_t),&duration,sizeof(uint32_t));
memcpy(buffer + sizeof(uint32_t) + sizeof(uint32_t),ptr,size);
if (queue->inindex >= QUEUE_SIZE)
queue->inindex = 0;
memcpy(&queue->incoming[queue->inindex],&queue->inindex,sizeof(uint32_t));
memcpy(&queue->incoming[queue->inindex] + sizeof(uint32_t),&duration,sizeof(uint32_t));
memcpy(&queue->incoming[queue->inindex] + sizeof(uint32_t) + sizeof(uint32_t),ptr,size);
queue->incoming[queue->inindex].size = size + sizeof(uint32_t) + sizeof(uint32_t);
queue->inindex++;
last_timestamp = timestamp;
index++;
} while (video_packets->peek());
} else if (queue_type == QueueType::Audio) {
do {
@ -259,12 +272,14 @@ main(int argc, char *argv[]) {
size_t size = packet->second.size();
auto duration = uint32_t(timestamp - last_timestamp);;
memcpy(buffer,&index,sizeof(uint32_t));
memcpy(buffer + sizeof(uint32_t),&duration,sizeof(uint32_t));
memcpy(buffer + sizeof(uint32_t) + sizeof(uint32_t),ptr,size);
if (queue->inindex >= QUEUE_SIZE)
queue->inindex = 0;
memcpy(&queue->incoming[queue->inindex],&queue->inindex,sizeof(uint32_t));
memcpy(&queue->incoming[queue->inindex] + sizeof(uint32_t),&duration,sizeof(uint32_t));
memcpy(&queue->incoming[queue->inindex] + sizeof(uint32_t) + sizeof(uint32_t),ptr,size);
queue->incoming[queue->inindex].size = size + sizeof(uint32_t) + sizeof(uint32_t);
queue->inindex++;
last_timestamp = timestamp;
index++;
} while (audio_packets->peek());
}
}
@ -308,7 +323,6 @@ main(int argc, char *argv[]) {
};
auto queue = &memory->queues[queuetype];
BOOST_LOG(info) << "Starting capture on channel " << queuetype;
if (queuetype == QueueType::Video) {
auto capture = std::thread{video_capture,mail,target,0};