test data sending
This commit is contained in:
@@ -57,6 +57,15 @@ void on_data(TcpClient* client) {
|
||||
}
|
||||
}
|
||||
|
||||
// TODO; Move inside check, this is redundant
|
||||
if (capabilitiesCount > 0) {
|
||||
// Send ACKNOWLEDGE; Might embed which packet is being acknowledged later
|
||||
char ackPacket[1];
|
||||
ackPacket[0] = PACKET_TYPE_ACKNOWLEDGE;
|
||||
send(client->clientFd, ackPacket, 1, 0);
|
||||
printf("Sent ACKNOWLEDGE to client %u\n", client->clientId);
|
||||
}
|
||||
|
||||
break;
|
||||
case PACKET_TYPE_TASK_REQUEST:
|
||||
printf("Received TASK_REQUEST packet from client %u\n", client->clientId);
|
||||
@@ -74,10 +83,18 @@ void on_data(TcpClient* client) {
|
||||
break;
|
||||
}
|
||||
|
||||
int assigned = 0;
|
||||
|
||||
printf("Task count: %d", (int)DynArr_size(taskQueue.tasks));
|
||||
|
||||
// Find a task that matches capabilities and assign it
|
||||
for (size_t i = 0; i < DynArr_size(&taskQueue.tasks); i++) {
|
||||
task_t* task = (task_t*)DynArr_at(&taskQueue.tasks, i);
|
||||
int assigned = 0;
|
||||
for (size_t i = 0; i < DynArr_size(taskQueue.tasks); i++) {
|
||||
//printf("Iterated!\n");
|
||||
task_t* task = (task_t*)DynArr_at(taskQueue.tasks, i);
|
||||
|
||||
if (!task) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (task->state == TASK_PENDING) {
|
||||
// For simplicity, assign the first pending task with capabilities
|
||||
@@ -97,7 +114,36 @@ void on_data(TcpClient* client) {
|
||||
DynArr_push_back(packet, &packetType);
|
||||
|
||||
// Push back the task struct
|
||||
DynArr_push_back(packet, task);
|
||||
for (size_t k = 0; k < sizeof(task_t); k++) {
|
||||
DynArr_push_back(packet, ((char*)task) + k);
|
||||
}
|
||||
|
||||
// Special sequence to indicate start of args
|
||||
char sequence[] = "ARGSTART";
|
||||
for (size_t k = 0; k < sizeof(sequence); k++) {
|
||||
DynArr_push_back(packet, &sequence[k]);
|
||||
}
|
||||
|
||||
// Push back args
|
||||
for (size_t k = 0; k < DynArr_size(task->args); k++) {
|
||||
task_arg_t* arg = (task_arg_t*)DynArr_at(task->args, k);
|
||||
for (size_t m = 0; m < sizeof(task_arg_t); m++) {
|
||||
DynArr_push_back(packet, ((char*)arg) + m);
|
||||
}
|
||||
// Push back arg data
|
||||
for (size_t m = 0; m < DynArr_size(arg->data); m++) {
|
||||
char* dataPtr = (char*)DynArr_at(arg->data, m);
|
||||
for (size_t n = 0; n < sizeof(char); n++) {
|
||||
DynArr_push_back(packet, dataPtr + n);
|
||||
}
|
||||
}
|
||||
|
||||
// Push back next arg indicator
|
||||
char nextArgIndicator[] = "ARGNEXT";
|
||||
for (size_t n = 0; n < sizeof(nextArgIndicator); n++) {
|
||||
DynArr_push_back(packet, &nextArgIndicator[n]);
|
||||
}
|
||||
}
|
||||
|
||||
send(client->clientFd, packet->data, DynArr_size(packet) * sizeof(uint8_t), 0);
|
||||
printf("Assigned task %u to client %u\n", task->taskId, client->clientId);
|
||||
@@ -131,6 +177,7 @@ void on_disconnect(TcpClient* client) {
|
||||
}
|
||||
|
||||
int main(void) {
|
||||
srand(time(NULL));
|
||||
signal(SIGINT, signalHandler);
|
||||
|
||||
TaskQueue_Init(&taskQueue);
|
||||
@@ -139,8 +186,15 @@ int main(void) {
|
||||
task_t task1;
|
||||
Task_Create(&task1);
|
||||
task1.state = TASK_PENDING;
|
||||
|
||||
task1.args = DYNARR_CREATE(task_arg_t, 1);
|
||||
strcpy(task1.binary, "test_computation");
|
||||
task_arg_t* blk = DynArr_push_back(task1.args, &(task_arg_t){ .name = "input_data", .data = DYNARR_CREATE(char, 1) });
|
||||
DynArr_push_back(blk->data, "12"); // Compute the square of 12
|
||||
task1.taskId = 1;
|
||||
|
||||
TaskQueue_AddTask(&taskQueue, &task1);
|
||||
//Task_DestroyArgs(&task1);
|
||||
|
||||
printf("Task added!\n");
|
||||
|
||||
TcpServer* svr = TcpServer_Create();
|
||||
|
||||
Reference in New Issue
Block a user