use multiple queues

This commit is contained in:
pigeatgarlic 2024-04-04 13:50:21 -07:00
parent 43cc8acf12
commit 8ddf3e4f96
5 changed files with 71 additions and 55 deletions

View File

@ -22,12 +22,7 @@ using namespace std::literals;
typedef struct {
Packet audio[QUEUE_SIZE];
Packet video[QUEUE_SIZE];
int audio_order[QUEUE_SIZE];
int video_order[QUEUE_SIZE];
Queue queues[QueueType::Max];
Event events[EVENT_TYPE_MAX];
interprocess_mutex lock;
}SharedMemoryInternal;
@ -72,8 +67,8 @@ managed_shared_memory segment(create_only, random.c_str(), 2 * sizeof(SharedMemo
void
init_shared_memory(SharedMemory* memory){
for (int i = 0; i < QUEUE_SIZE; i++) {
memory->audio_order[i] = -1;
memory->video_order[i] = -1;
memory->queues[QueueType::Audio].order[i] = -1;
memory->queues[QueueType::Video].order[i] = -1;
}
for (int i = 0; i < EventType::EVENT_TYPE_MAX; i++)

View File

@ -50,13 +50,20 @@ typedef struct {
int read;
} Event;
enum QueueType {
Video,
Audio,
Microphone,
Max
};
typedef struct _Queue{
Packet array[QUEUE_SIZE];
int order[QUEUE_SIZE];
}Queue;
typedef struct {
Packet audio[QUEUE_SIZE];
Packet video[QUEUE_SIZE];
int audio_order[QUEUE_SIZE];
int video_order[QUEUE_SIZE];
Queue queues[QueueType::Max];
Event events[EVENT_TYPE_MAX];
}SharedMemory;

View File

@ -40,13 +40,20 @@ typedef struct {
int read;
} Event;
enum QueueType {
Video,
Audio,
Microphone,
Max
};
typedef struct _Queue{
Packet array[QUEUE_SIZE];
int order[QUEUE_SIZE];
}Queue;
typedef struct {
Packet audio[QUEUE_SIZE];
Packet video[QUEUE_SIZE];
int audio_order[QUEUE_SIZE];
int video_order[QUEUE_SIZE];
Queue queues[Max];
Event events[EVENT_TYPE_MAX];
}SharedMemory;
@ -80,9 +87,9 @@ type DataType int
func peek(memory *C.SharedMemory, media DataType) bool {
if media == video {
return memory.video_order[0] != -1
return memory.queues[C.Video].order[0] != -1
} else if media == audio {
return memory.audio_order[0] != -1
return memory.queues[C.Audio].order[0] != -1
}
panic(fmt.Errorf("unknown data type"))
@ -158,7 +165,7 @@ func main() {
panic(err)
}
pointer,_,err := obtain.Call(
pointer, _, err := obtain.Call(
uintptr(unsafe.Pointer(&buffer[0])),
)
if !errors.Is(err, windows.ERROR_SUCCESS) {
@ -170,28 +177,28 @@ func main() {
lock.Call(pointer)
defer unlock.Call(pointer)
block := memory.video[memory.video_order[0]]
block := memory.queues[C.Video].array[memory.queues[C.Video].order[0]]
fmt.Printf("video buffer %d\n", block.size)
for i := 0; i < C.QUEUE_SIZE-1; i++ {
memory.video_order[i] = memory.video_order[i+1]
memory.queues[C.Video].order[i] = memory.queues[C.Video].order[i+1]
}
memory.video_order[C.QUEUE_SIZE-1] = -1
memory.queues[C.Video].order[C.QUEUE_SIZE-1] = -1
}
handle_audio := func() {
lock.Call(pointer)
defer unlock.Call(pointer)
block := memory.audio[memory.audio_order[0]]
block := memory.queues[C.Audio].array[memory.queues[C.Audio].order[0]]
fmt.Printf("audio buffer %d\n", block.size)
for i := 0; i < C.QUEUE_SIZE-1; i++ {
memory.audio_order[i] = memory.audio_order[i+1]
memory.queues[C.Audio].order[i] = memory.queues[C.Audio].order[i+1]
}
memory.audio_order[C.QUEUE_SIZE-1] = -1
memory.queues[C.Audio].order[C.QUEUE_SIZE-1] = -1
}
go func() {

View File

@ -84,15 +84,15 @@ int find_available_slot(int* orders) {
void
push_audio_packet(SharedMemory* memory, void* data, int size){
// wait while queue is full
while (queue_size(memory->audio_order) == QUEUE_SIZE)
while (queue_size(memory->queues[QueueType::Audio].order) == QUEUE_SIZE)
std::this_thread::sleep_for(1ms);
scoped_lock<interprocess_mutex> lock(memory->lock);
int available = find_available_slot(memory->audio_order);
memory->audio_order[queue_size(memory->audio_order)] = available;
Packet* block = &memory->audio[available];
int available = find_available_slot(memory->queues[QueueType::Audio].order);
memory->queues[QueueType::Audio].order[queue_size(memory->queues[QueueType::Audio].order)] = available;
Packet* block = &memory->queues[QueueType::Audio].array[available];
memcpy(block->data,data,size);
block->size = size;
@ -104,14 +104,14 @@ push_video_packet(SharedMemory* memory,
int size,
VideoMetadata metadata){
// wait while queue is full
while (queue_size(memory->video_order) == QUEUE_SIZE)
while (queue_size(memory->queues[QueueType::Video].order) == QUEUE_SIZE)
std::this_thread::sleep_for(1ms);
scoped_lock<interprocess_mutex> lock(memory->lock);
int available = find_available_slot(memory->video_order);
memory->video_order[queue_size(memory->video_order)] = available;
Packet* block = &memory->video[available];
int available = find_available_slot(memory->queues[QueueType::Video].order);
memory->queues[QueueType::Video].order[queue_size(memory->queues[QueueType::Video].order)] = available;
Packet* block = &memory->queues[QueueType::Video].array[available];
memcpy(block->data,data,size);
block->size = size;
block->metadata = metadata;
@ -119,12 +119,12 @@ push_video_packet(SharedMemory* memory,
int
peek_video_packet(SharedMemory* memory){
return memory->video_order[0] != -1;
return memory->queues[QueueType::Video].order[0] != -1;
}
int
peek_audio_packet(SharedMemory* memory){
return memory->audio_order[0] != -1;
return memory->queues[QueueType::Audio].order[0] != -1;
}
void
@ -133,18 +133,18 @@ pop_audio_packet(SharedMemory* memory, void* data, int* size){
std::this_thread::sleep_for(1ms);
scoped_lock<interprocess_mutex> lock(memory->lock);
// std::cout << "Audio buffer size : " << queue_size(memory->audio_order) << "\n";
// std::cout << "Audio buffer size : " << queue_size(memory->queues[QueueType::Audio].order) << "\n";
int pop = memory->audio_order[0];
Packet *block = &memory->audio[pop];
int pop = memory->queues[QueueType::Audio].order[0];
Packet *block = &memory->queues[QueueType::Audio].array[pop];
memcpy(data,block->data,block->size);
*size = block->size;
// reorder
for (int i = 0; i < QUEUE_SIZE - 1; i++)
memory->audio_order[i] = memory->audio_order[i+1];
memory->queues[QueueType::Audio].order[i] = memory->queues[QueueType::Audio].order[i+1];
memory->audio_order[QUEUE_SIZE - 1] = -1;
memory->queues[QueueType::Audio].order[QUEUE_SIZE - 1] = -1;
}
@ -154,19 +154,19 @@ pop_video_packet(SharedMemory* memory, void* data, int* size){
std::this_thread::sleep_for(1ms);
scoped_lock<interprocess_mutex> lock(memory->lock);
// std::cout << "Video buffer size : " << queue_size(memory->video_order) << "\n";
// std::cout << "Video buffer size : " << queue_size(memory->queues[QueueType::Video].order) << "\n";
int pop = memory->video_order[0];
Packet *block = &memory->video[pop];
int pop = memory->queues[QueueType::Video].order[0];
Packet *block = &memory->queues[QueueType::Video].array[pop];
memcpy(data,block->data,block->size);
*size = block->size;
auto copy = block->metadata;
// reorder
for (int i = 0; i < QUEUE_SIZE - 1; i++)
memory->video_order[i] = memory->video_order[i+1];
memory->queues[QueueType::Video].order[i] = memory->queues[QueueType::Video].order[i+1];
memory->video_order[QUEUE_SIZE - 1] = -1;
memory->queues[QueueType::Video].order[QUEUE_SIZE - 1] = -1;
return copy;

View File

@ -53,13 +53,20 @@ typedef struct {
int read;
} Event;
enum QueueType {
Video,
Audio,
Microphone,
Max
};
typedef struct _Queue{
Packet array[QUEUE_SIZE];
int order[QUEUE_SIZE];
}Queue;
typedef struct {
Packet audio[QUEUE_SIZE];
Packet video[QUEUE_SIZE];
int audio_order[QUEUE_SIZE];
int video_order[QUEUE_SIZE];
Queue queues[QueueType::Max];
Event events[EVENT_TYPE_MAX];
interprocess_mutex lock;
}SharedMemory;