diff --git a/src/benchmarks/activeRulesBenchmark.h b/src/benchmarks/activeRulesBenchmark.h index 6f5e161..2d9fd2b 100644 --- a/src/benchmarks/activeRulesBenchmark.h +++ b/src/benchmarks/activeRulesBenchmark.h @@ -1,9 +1,8 @@ #include -#include // Required for timeBeginPeriod +#include #include #include #include -#include #ifdef DIST #include "embedDB.h" @@ -23,6 +22,8 @@ #ifdef ARDUINO +#include "Arduino.h" + #if defined(MEMBOARD) && STORAGE_TYPE == 1 #include "dataflashFileInterface.h" @@ -30,6 +31,11 @@ #endif #include "SDFileInterface.h" +#include "query-interface/activeRules.h" +#define FILE_TYPE SD_FILE +#define fopen sd_fopen +#define fread sd_fread +#define fclose sd_fclose #define getFileInterface getSDInterface #define setupFile setupSDFile #define tearDownFile tearDownSDFile @@ -38,7 +44,7 @@ #define INDEX_PATH "indexFile.bin" #else - +#define FILE_TYPE FILE #include "desktopFileInterface.h" #include "query-interface/activeRules.h" #define DATA_PATH "build/artifacts/dataFile.bin" @@ -46,24 +52,36 @@ #endif -#define NUM_INSERTIONS 2000 +#define NUM_INSERTIONS 1000 embedDBState* init_state(); embedDBSchema* createSchema(); void GTcallback(void* aggregateValue, void* currentValue, void* context); +int callbacks = 0; + // Get current time in nanoseconds uint64_t get_nanoseconds() { +#ifdef ARDUINO + return (uint64_t)micros() * 1000ULL; +#else +#if defined(TIME_UTC) struct timespec ts; - clock_gettime(CLOCK_MONOTONIC, &ts); - return (uint64_t)ts.tv_sec * 1e9 + ts.tv_nsec; + if (timespec_get(&ts, TIME_UTC) == TIME_UTC) { + return (uint64_t)ts.tv_sec * 1000000000ULL + (uint64_t)ts.tv_nsec; + } +#endif + return (uint64_t)((double)clock() * (1000000000.0 / CLOCKS_PER_SEC)); +#endif } // Callback function for active rule void GTcallback(void* aggregateValue, void* currentValue, void* context) { - FILE* perfLog = (FILE*)context; + // FILE_TYPE* perfLog = (FILE_TYPE*)context; uint64_t callbackTime = get_nanoseconds(); - fprintf(perfLog, "%llu,CALLBACK,%f\n", callbackTime, *(float*)aggregateValue); + // fprintf(perfLog, "%llu,CALLBACK,%f\n", (unsigned long long)callbackTime, *(float*)aggregateValue); + // printf("%llu,CALLBACK,%f\n", (unsigned long long)callbackTime, *(float*)aggregateValue); + callbacks++; } int8_t groupFunctionLocal(const void* lastRecord, const void* record) { @@ -72,7 +90,8 @@ int8_t groupFunctionLocal(const void* lastRecord, const void* record) { embedDBOperator* createOperatorLocal(embedDBState* state, embedDBSchema* schema, void*** allocatedValues, uint32_t key) { embedDBIterator* it = (embedDBIterator*)malloc(sizeof(embedDBIterator)); - uint32_t minKeyVal = key - (1000 - 1); + uint32_t numRecords = 1; + uint32_t minKeyVal = key - (numRecords - 1); uint32_t* minKeyPtr = (uint32_t*)malloc(sizeof(uint32_t)); *minKeyPtr = minKeyVal; it->minKey = minKeyPtr; @@ -118,23 +137,28 @@ void GetAvgLocal(embedDBState* state, embedDBSchema* schema, uint32_t key, float free(allocatedValues[i]); } free(allocatedValues); - if (avg > 0) { + if (avg < 0) { GTcallback(&avg, ¤tVal, context); } } uint32_t activeRulesBenchmark() { + printf("Active Rules Benchmark\n"); embedDBState* state = init_state(); embedDBPrintInit(state); embedDBSchema* schema = createSchema(); // Create active rule activeRule* activeRuleGT = createActiveRule(schema, NULL); + uint32_t numRecords = 10; + float minVal = 0.0f; activeRuleGT->IF(activeRuleGT, 1, GET_AVG) - ->ofLast(activeRuleGT, (void*)&(uint32_t){1000}) - ->is(activeRuleGT, GreaterThan, (void*)&(float){0}) + ->ofLast(activeRuleGT, (void*)&numRecords) + // ->is(activeRuleGT, GreaterThan, (void*)&minVal) + ->is(activeRuleGT, LessThan, (void*)&minVal) ->then(activeRuleGT, GTcallback); + printf("Window size: %u\n", numRecords); state->rules = (activeRule**)malloc(sizeof(activeRule*)); state->rules[0] = activeRuleGT; state->numRules = 1; @@ -142,81 +166,79 @@ uint32_t activeRulesBenchmark() { srand(12345); // Fixed seed for reproducibility // Open performance log file - FILE* perfLog = fopen("C:/Users/richa/OneDrive/Documents/influxdb/embeddb_perf_new.csv", "w"); + // FILE_TYPE perfLog = fopen("C:/tmp/EmbedDB-Vs-InfluxDB/embeddb_perf_new.csv", "w"); // FILE* advancedPerfLog = fopen("C:/Users/richa/OneDrive/Documents/influxdb/embeddb_advanced_perf.csv", "w"); // fprintf(advancedPerfLog, "timestamp,event,temperature,latency\n"); - fprintf(perfLog, "timestamp,event,temperature,latency\n"); + // fprintf(perfLog, "timestamp,event,temperature,latency\n"); + printf("timestamp,event,temperature,latency\n"); // Set callback context to the log file - state->rules[0]->context = perfLog; - timeBeginPeriod(1); + // state->rules[0]->context = perfLog; uint32_t j = 0; // Insert without active query first - for (int i = 0; i < NUM_INSERTIONS; i++) { + printf("Initial inserts\n"); + void* dataPtr = malloc(state->dataSize); + for (int i = 0; i < 1000; i++) { uint64_t timestamp = get_nanoseconds(); - float temperature = 15 + (float)rand() / RAND_MAX * 15; // Random temperature between 15°C and 30°C + float temperature = -5 + (float)rand() / RAND_MAX * 10; // Random temperature between -5°C and 5°C - LARGE_INTEGER start, end, freq; - QueryPerformanceFrequency(&freq); - QueryPerformanceCounter(&start); + uint64_t start = get_nanoseconds(); - void* dataPtr = malloc(state->dataSize); *((float*)dataPtr) = temperature; int8_t result = embedDBPut(state, &j, dataPtr); + (void)result; - QueryPerformanceCounter(&end); - int insertTime = (double)(end.QuadPart - start.QuadPart) / freq.QuadPart * 1e9; // Convert to nanoseconds + uint64_t end = get_nanoseconds(); + uint64_t insertTime = end - start; // Log insertion event // fprintf(advancedPerfLog, "%llu,INSERT,%f,%i\n", timestamp, temperature, insertTime); - fprintf(perfLog, "%llu,INSERT,%f,%i\n", timestamp, temperature, insertTime); + // fprintf(perfLog, "%llu,INSERT,%f,%llu\n", (unsigned long long)timestamp, temperature, (unsigned long long)insertTime); + if (i % 100 == 0) + printf("%llu,INSERT,%f,%llu\n", (unsigned long long)timestamp, temperature, (unsigned long long)insertTime); - free(dataPtr); j++; } - state->rules[0]->enabled = true; // Enable the rule for subsequent insertions + printf("Test inserts\n"); + // state->rules[0]->enabled = true; // Enable the rule for subsequent insertions uint64_t startTime = get_nanoseconds(); for (int i = 0; i < NUM_INSERTIONS; i++) { uint64_t timestamp = get_nanoseconds(); - float temperature = 15 + (float)rand() / RAND_MAX * 15; // Random temperature between 15°C and 30°C - - LARGE_INTEGER start, end, freq; - QueryPerformanceFrequency(&freq); - QueryPerformanceCounter(&start); + float temperature = -5 + (float)rand() / RAND_MAX * 10; // Random temperature between -5°C and 5°C - void* dataPtr = malloc(state->dataSize); + uint64_t start = get_nanoseconds(); + // void* dataPtr = malloc(state->dataSize); *((float*)dataPtr) = temperature; // using j instead of timestamp ensures same number of records queried each time independent of changing insert speed int8_t result = embedDBPut(state, &j, dataPtr); + // Uncomment this if want to test performance of advanced query without active rule callback + GetAvgLocal(state, schema, j, temperature, NULL); - // QueryPerformanceFrequency(&freq); - // QueryPerformanceCounter(&start); - // GetAvgLocal(state, schema, j, temperature, advancedPerfLog); - - QueryPerformanceCounter(&end); - int insertTime = (double)(end.QuadPart - start.QuadPart) / freq.QuadPart * 1e9; // Convert to nanoseconds + uint64_t end = get_nanoseconds(); + uint64_t insertTime = end - start; // Log insertion event // fprintf(advancedPerfLog, "%llu,INSERT,%f,%i\n", timestamp, temperature, insertTime); - fprintf(perfLog, "%llu,INSERT,%f,%i\n", timestamp, temperature, insertTime); + // fprintf(perfLog, "%llu,INSERT,%f,%llu\n", (unsigned long long)timestamp, temperature, (unsigned long long)insertTime); + if (i % 100000 == 0) + printf("%llu,INSERT,%f,%llu\n", (unsigned long long)timestamp, temperature, (unsigned long long)insertTime); - free(dataPtr); j++; } - timeEndPeriod(1); + free(dataPtr); uint64_t endTime = get_nanoseconds(); // Calculate throughput double totalTime = (double)(endTime - startTime) / 1e9; // Convert to seconds double throughput = NUM_INSERTIONS / totalTime; - printf("Throughput: %f insertions/second\n", throughput); + printf("Throughput: %f insertions/second Time: %f Records: %d Callbacks: %d\n", throughput, totalTime, NUM_INSERTIONS, callbacks); // Clean up - fclose(perfLog); + // fclose(perfLog); // fclose(advancedPerfLog); return 0; } diff --git a/src/desktopMain.c b/src/desktopMain.c index d4ae86a..614126f 100644 --- a/src/desktopMain.c +++ b/src/desktopMain.c @@ -18,9 +18,9 @@ #elif WHICH_PROGRAM == 3 #include "benchmarks/queryInterfaceBenchmark.h" #elif WHICH_PROGRAM == 4 -#include "benchmarks/sortBenchmark.h" -#elif WHICH_PROGRAM == 5 #include "benchmarks/activeRulesBenchmark.h" +#elif WHICH_PROGRAM == 5 +#include "benchmarks/sortBenchmark.h" #endif int main() { @@ -33,9 +33,9 @@ int main() { #elif WHICH_PROGRAM == 3 return advancedQueryExample(); #elif WHICH_PROGRAM == 4 - return sortQueryBenchmark(); -#elif WHICH_PROGRAM == 5 return activeRulesBenchmark(); +#elif WHICH_PROGRAM == 5 + return sortQueryBenchmark(); #endif } diff --git a/src/dueMain.cpp b/src/dueMain.cpp index 82d7312..dab1cc4 100644 --- a/src/dueMain.cpp +++ b/src/dueMain.cpp @@ -70,6 +70,8 @@ static ArduinoOutStream cout(Serial); #elif WHICH_PROGRAM == 3 #include "benchmarks/queryInterfaceBenchmark.h" #elif WHICH_PROGRAM == 4 +#include "benchmarks/activeRulesBenchmark.h" +#elif WHICH_PROGRAM == 5 #include "benchmarks/sortBenchmark.h" #endif @@ -102,7 +104,7 @@ void setup() { sd.ls("/", LS_R); } - init_sdcard((void *)&sd); + init_sdcard((void*)&sd); #if WHICH_PROGRAM == 0 embedDBExample(); #elif WHICH_PROGRAM == 1 @@ -112,6 +114,8 @@ void setup() { #elif WHICH_PROGRAM == 3 advancedQueryExample(); #elif WHICH_PROGRAM == 4 + activeRulesBenchmark(); +#elif WHICH_PROGRAM == 5 sortQueryBenchmark(); #endif } diff --git a/src/query-interface/advancedQueries.c b/src/query-interface/advancedQueries.c index db544dc..3c3663c 100644 --- a/src/query-interface/advancedQueries.c +++ b/src/query-interface/advancedQueries.c @@ -517,7 +517,6 @@ int8_t nextOrderBy(embedDBOperator* op) { void closeOrderBy(embedDBOperator* op) { op->input->close(op->input); - op->input = NULL; embedDBFreeSchema(&op->schema); closeSort(((sortData*)op->state)->fileIterator); @@ -704,7 +703,6 @@ int8_t nextAggregate(embedDBOperator* op) { void closeAggregate(embedDBOperator* op) { op->input->close(op->input); - op->input = NULL; embedDBFreeSchema(&op->schema); free(((struct aggregateInfo*)op->state)->lastRecordBuffer); free(op->state);