From 21a698bb63663659e12d6b724197ef9840cc51f5 Mon Sep 17 00:00:00 2001 From: mive93 Date: Mon, 15 Apr 2019 11:38:54 +0200 Subject: [PATCH 1/2] added server and serialization --- CMakeLists.txt | 6 +- demo/demo/demo.cpp | 4 + demo/server/server_less_dummy.cpp | 283 ++++++++++++++++++++++++++++++ include/send.h | 68 ++++--- include/serialize.hpp | 38 ++++ 5 files changed, 371 insertions(+), 28 deletions(-) create mode 100644 demo/server/server_less_dummy.cpp create mode 100644 include/serialize.hpp diff --git a/CMakeLists.txt b/CMakeLists.txt index fef13a0..c960112 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -37,7 +37,7 @@ file(GLOB tkdnn_SRC "src/*.cpp") set(tkdnn_LIBS kernels ${CUDA_LIBRARIES} ${CUDA_CUBLAS_LIBRARIES} -lcudnn -lnvinfer ${OpenCV_LIBS}) set(CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS} -Wall -std=c++11") -include_directories(${CMAKE_CURRENT_SOURCE_DIR}/include ${CUDA_INCLUDE_DIRS} ${OPENCV_INCLUDE_DIRS} ${NVINFER_INCLUDES}) +include_directories(${CMAKE_CURRENT_SOURCE_DIR}/include ${CUDA_INCLUDE_DIRS} ${OPENCV_INCLUDE_DIRS} ${NVINFER_INCLUDES} "~/repos/cereal/include") add_library(tkDNN SHARED ${tkdnn_SRC}) target_link_libraries(tkDNN ${tkdnn_LIBS}) @@ -88,6 +88,10 @@ target_link_libraries(test_rtinference tkDNN) add_executable(yolo3_demo demo/demo/demo.cpp) target_link_libraries(yolo3_demo tkDNN) +add_executable(class_server demo/server/server_less_dummy.cpp) +target_link_libraries(class_server pthread tkDNN ) + + #install #if (CMAKE_INSTALL_PREFIX_INITIALIZED_TO_DEFAULT) # set (CMAKE_INSTALL_PREFIX "${CMAKE_BINARY_DIR}/install" diff --git a/demo/demo/demo.cpp b/demo/demo/demo.cpp index 98ebc49..e4255d9 100644 --- a/demo/demo/demo.cpp +++ b/demo/demo/demo.cpp @@ -77,7 +77,11 @@ int main(int argc, char *argv[]) { cap >> frame; if(!frame.data) { + usleep(1000000); + cap.open(input); + printf("cap reinitialize\n"); continue; + } // this will be resized to the net format diff --git a/demo/server/server_less_dummy.cpp b/demo/server/server_less_dummy.cpp new file mode 100644 index 0000000..53a3ebb --- /dev/null +++ b/demo/server/server_less_dummy.cpp @@ -0,0 +1,283 @@ +/* + C socket server example, handles multiple clients using threads +*/ + +#include +#include //strlen +#include //strlen +#include +#include +#include //inet_addr +#include //write +#include //for threading , link with lpthread +#include +#include +#include +#include "serialize.hpp" +#include + +sem_t semaphore; + +struct obj_coords +{ + float LAT; + float LONG; + float cl; +}; + +int write_coords_to_file(std::string buffer, FILE *f) +{ + + std::istringstream is(buffer); + cereal::PortableBinaryInputArchive retrieve(is); + Message m; + retrieve(m); + + int cam_id = m.cam_idx; + unsigned long long t_stamp_ms = m.t_stamp_ms; + int obj_n = m.num_objects; + + struct obj_coords *c = (struct obj_coords *)malloc(obj_n * sizeof(struct obj_coords)); + + int i; + for (i = 0; i < obj_n; i++) + { + c[i].LAT = m.objects.at(i).latitude; + c[i].LONG = m.objects.at(i).longitude; + c[i].cl = m.objects.at(i).category; + //m.objects.at(i).speed; + //m.objects.at(i).orientation; + } + + char *to_print = (char *)malloc(100000); + memset(to_print, 0, 100000); + + /*int obj_n; + int cam_id; + unsigned long long t_stamp_ms; + char type_of_m; + + int consumed_chars = 0; + char *shifted_chars = (char*)malloc(4000); + char *to_print = (char*)malloc(100000); + memset(to_print,0,100000); + + sscanf(buffer, "%c %lld %d %d ", &type_of_m, &t_stamp_ms, &cam_id, &obj_n); + sprintf(shifted_chars, "%c %lld %d %d ", type_of_m, t_stamp_ms, cam_id, obj_n); + consumed_chars += strlen(shifted_chars); + //printf("%c %lld %d %d\n", type_of_m, t_stamp_ms, cam_id, *obj_n); + + struct obj_coords *c = (struct obj_coords *)malloc(obj_n * sizeof(struct obj_coords)); + + int i; + for (i = 0; i < obj_n; i++) + { + sscanf(buffer + consumed_chars, "%f %f %f ", &c[i].LAT, &c[i].LONG, &c[i].cl); + sprintf(shifted_chars, "%.9f %.9f %.0f ", c[i].LAT, c[i].LONG, c[i].cl); + consumed_chars += strlen(shifted_chars); + //printf("%Lf %f %f \n", c[i].LAT, c[i].LONG, c[i].cl); + }*/ + + char *command = (char *)malloc(4000); + for (i = 0; i < obj_n; i++) + { + sprintf(command, "%d %lld %.0f %.9f %.9f %f %f\n", cam_id, t_stamp_ms, c[i].cl, c[i].LAT, c[i].LONG, 0.0f, 0.0f); + strcat(to_print, command); + } + + printf("%s", to_print); + + fprintf(f, "%s", to_print); + free(to_print); + free(command); + free(c); + + return obj_n; +} + +//the thread function +void *connection_handler(void *); + +int main(int argc, char *argv[]) +{ + const int path_size = 400; + //const int message_size = 100000; + char path[path_size]; + //int obj_n; + const char *basepath = "/data/class/"; + struct stat st = {0}; + + sem_init(&semaphore, 0, 1); + + int socket_desc, client_sock, c, *new_sock; + struct sockaddr_in server, client; + + time_t rawtime; + time(&rawtime); + struct tm *tm_struct = localtime(&rawtime); + int tm_hour = tm_struct->tm_hour; + int tm_yday = tm_struct->tm_yday; + + sprintf(path, "%s%d/", basepath, tm_yday); + if (stat(path, &st) == -1) + { + mkdir(path, 0700); + } + memset(path, 0, path_size); + + sprintf(path, "%s%d/%d/", basepath, tm_yday, tm_hour); + if (stat(path, &st) == -1) + { + mkdir(path, 0700); + } + memset(path, 0, path_size); + + //Create socket + socket_desc = socket(AF_INET, SOCK_STREAM, 0); + if (socket_desc == -1) + { + printf("Could not create socket"); + } + puts("Socket created"); + + //Prepare the sockaddr_in structure + server.sin_family = AF_INET; + server.sin_addr.s_addr = INADDR_ANY; + server.sin_port = htons(8888); + + //Bind + if (bind(socket_desc, (struct sockaddr *)&server, sizeof(server)) < 0) + { + //print the error message + perror("bind failed. Error"); + return 1; + } + puts("bind done"); + + //Listen + listen(socket_desc, 3); + + //Accept and incoming connection + puts("Waiting for incoming connections..."); + c = sizeof(struct sockaddr_in); + while ((client_sock = accept(socket_desc, (struct sockaddr *)&client, (socklen_t *)&c))) + { + puts("Connection accepted"); + + pthread_t sniffer_thread; + new_sock = (int *)malloc(1); + *new_sock = client_sock; + + if (pthread_create(&sniffer_thread, NULL, connection_handler, (void *)new_sock) < 0) + { + perror("could not create thread"); + return 1; + } + + //Now join the thread , so that we dont terminate before the thread + //pthread_join( sniffer_thread , NULL); + puts("Handler assigned"); + } + + if (client_sock < 0) + { + perror("accept failed"); + return 1; + } + + sem_destroy(&semaphore); + + return 0; +} + +/* + * This will handle connection for each client + * */ +void *connection_handler(void *socket_desc) +{ + time_t rawtime; + struct tm *tm_struct; + int tm_hour; + int tm_yday; + const int path_size = 400; + const int message_size = 100000; + char path[path_size]; + int obj_n; + const char *basepath = "/data/class/"; + struct stat st = {0}; + FILE *f; + int mess_hour, mess_min, mess_yday; + + //Get the socket descriptor + int sock = *(int *)socket_desc; + int read_size; + + void *client_message = (void*)malloc(message_size); + + /* //Send some messages to the client + message = "Greetings! I am your connection handler\n"; + write(sock, message, strlen(message)); */ + + //Receive a message from client + while ((read_size = recv(sock, client_message, message_size, 0)) > 0) + { + + std::string s((char *)client_message, message_size); + //std::cout<<"Message received: "<tm_min; + mess_hour = tm_struct->tm_hour; + mess_yday = tm_struct->tm_yday; + if (mess_yday != tm_yday) + { + tm_yday = mess_yday; + sprintf(path, "%s%d/", basepath, tm_yday); + if (stat(path, &st) == -1) + { + mkdir(path, 0700); + } + memset(path, 0, path_size); + } + if (mess_hour != tm_hour) + { + tm_hour = mess_hour; + sprintf(path, "%s%d/%d/", basepath, tm_yday, tm_hour); + if (stat(path, &st) == -1) + { + mkdir(path, 0700); + } + memset(path, 0, path_size); + } + + sprintf(path, "%s%d/%d/%d.txt", basepath, tm_yday, tm_hour, mess_min); + + //CRITICAL SECTION + sem_wait(&semaphore); + f = fopen(path, "a"); + memset(path, 0, path_size); + + obj_n = write_coords_to_file(s, f); + + fclose(f); + sem_post(&semaphore); + } + + if (read_size == 0) + { + puts("Client disconnected"); + fflush(stdout); + } + else if (read_size == -1) + { + perror("recv failed"); + } + + //Free the socket pointer + free(socket_desc); + + free(client_message); + + return 0; +} diff --git a/include/send.h b/include/send.h index 2ecf7aa..aa70d23 100644 --- a/include/send.h +++ b/include/send.h @@ -9,6 +9,8 @@ #include //inet_addr #include //write +#include "serialize.hpp" + struct obj_coords { float LAT; @@ -17,49 +19,57 @@ struct obj_coords }; -char *serialize_coords(struct obj_coords *c, int obj_n, int CAM_IDX) +void serialize_coords(struct obj_coords *c, int obj_n, int CAM_IDX, std::stringbuf* buf) { - - char *buffer = (char *)malloc(100000 * sizeof(char)); + + std::ostream os(buf); + cereal::PortableBinaryOutputArchive archive(os); struct timeval tv; gettimeofday(&tv, NULL); unsigned long long t_stamp_ms = (unsigned long long)(tv.tv_sec) * 1000 + (unsigned long long)(tv.tv_usec) / 1000; + + std::vector ruv; - sprintf(buffer, "m %lld %d %d ", t_stamp_ms, CAM_IDX, obj_n); int i; for (i = 0; i < obj_n; i++) - sprintf(buffer + strlen(buffer), "%.9f %.9f %.0f ", c[i].LAT, c[i].LONG, c[i].cl); + { + Road_User r{c[i].LAT,c[i].LONG,0,0,(int)c[i].cl}; + ruv.push_back(r); + } - //printf("%s\n", buffer); + Message m{CAM_IDX,t_stamp_ms,ruv.size(),ruv}; + archive(m); - return buffer; + + //std::cout<str()<str().length()<str().data(), message->str().length(), 0) < 0) { puts("Send failed"); socket_opened = 0; } - free(message); + delete message; + //free(message); //close(sock); return 1; } diff --git a/include/serialize.hpp b/include/serialize.hpp new file mode 100644 index 0000000..975c3ec --- /dev/null +++ b/include/serialize.hpp @@ -0,0 +1,38 @@ +#ifndef SERIALIZE_H +#define SERIALIZE_H + +#include +#include + +struct Road_User{ + float latitude; + float longitude; + uint8_t speed; + uint8_t orientation; + uint8_t category; + + template + void serialize(Archive & archive) + { + archive( latitude, longitude, speed, orientation, category ); + } + +}; + +struct Message{ + uint32_t cam_idx; + uint64_t t_stamp_ms; + uint16_t num_objects; + std::vector objects; + + template + void serialize(Archive & archive) + { + archive( cam_idx, t_stamp_ms, num_objects, objects ); + } +}; + + + + +#endif \ No newline at end of file From 40c67e85369296c7f58283f2256907cf00f09d89 Mon Sep 17 00:00:00 2001 From: mive93 Date: Mon, 15 Apr 2019 11:54:16 +0200 Subject: [PATCH 2/2] dla commented --- src/NetworkRT.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/NetworkRT.cpp b/src/NetworkRT.cpp index 93dd321..7ef7e75 100644 --- a/src/NetworkRT.cpp +++ b/src/NetworkRT.cpp @@ -35,7 +35,7 @@ NetworkRT::NetworkRT(Network *net, const char *name) { builderRT = createInferBuilder(loggerRT); std::cout<<"Float16 support: "<platformHasFastFp16()<<"\n"; std::cout<<"Int8 support: "<platformHasFastInt8()<<"\n"; - std::cout<<"DLAs: "<getNbDLACores()<<"\n"; + //std::cout<<"DLAs: "<getNbDLACores()<<"\n"; networkRT = builderRT->createNetwork(); if(!fileExist(name)) {