From de53ecd1c02f3eef83c6af38cfe53db0e8628ce1 Mon Sep 17 00:00:00 2001 From: Max Lapshin Date: Wed, 14 Oct 2009 01:57:04 +0400 Subject: [PATCH] now can in right way reply to empty resultsets --- priv/sqlite3_drv.c | 206 ++++++++++++++++++++++++++++----------------- priv/sqlite3_drv.h | 21 +++-- 2 files changed, 146 insertions(+), 81 deletions(-) diff --git a/priv/sqlite3_drv.c b/priv/sqlite3_drv.c index 47b417a..441d30d 100644 --- a/priv/sqlite3_drv.c +++ b/priv/sqlite3_drv.c @@ -14,7 +14,7 @@ static ErlDrvEntry basic_driver_entry = { control, /* control */ NULL, /* timeout */ NULL, /* outputv (defined below) */ - NULL, /* ready_async */ + ready_async, /* ready_async */ NULL, /* flush */ NULL, /* call */ NULL, /* event */ @@ -46,6 +46,7 @@ static ErlDrvData start(ErlDrvPort port, char* cmd) { // Set the state for the driver retval->port = port; retval->db = db; + retval->key = 42; //FIXME: Just a magic number, make real key return (ErlDrvData) retval; } @@ -74,10 +75,6 @@ static int control(ErlDrvData drv_data, unsigned int command, char *buf, return 0; } -static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data) -{ - -} static inline int return_error(sqlite3_drv_t *drv, const char *error) { ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"), @@ -92,13 +89,58 @@ static inline int return_error(sqlite3_drv_t *drv, const char *error) { static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) { - ErlDrvTermData *dataset; - - int result, next_row, column_count, row_count = 0, term_count = 0; - char *error = NULL; + int result, next_row; char *rest = NULL; sqlite3_stmt *statement; + fprintf(stderr, "Exec: %.*s\n", command_size, command); + result = sqlite3_prepare_v2(drv->db, command, command_size, &statement, (const char **)&rest); + if(result != SQLITE_OK) { + return return_error(drv, sqlite3_errmsg(drv->db)); + } + + async_sqlite3_command *async_command = (async_sqlite3_command *)calloc(1, sizeof(async_sqlite3_command)); + async_command->driver_data = drv; + async_command->statement = statement; + + // fprintf(stderr, "Driver async: %p\n", async_command->statement); + + result = driver_async(drv->port, &drv->key, sql_exec_async, async_command, sql_free_async); + return 0; +} + +static void sql_free_async(void *_async_command) +{ + int i; + async_sqlite3_command *async_command = (async_sqlite3_command *)_async_command; + free(async_command->dataset); + + if (async_command->floats) { + free(async_command->floats); + } + for (i = 0; i < async_command->binaries_count; i++) { + driver_free_binary(async_command->binaries[i]); + } + if(async_command->binaries) { + free(async_command->binaries); + } + sqlite3_finalize(async_command->statement); + free(async_command); +} + + +static void sql_exec_async(void *_async_command) { + async_sqlite3_command *async_command = (async_sqlite3_command *)_async_command; + ErlDrvTermData *dataset = async_command->dataset; + int term_count = async_command->term_count; + int row_count = async_command->row_count; + sqlite3_drv_t *drv = async_command->driver_data; + + int result, next_row, column_count; + char *error = NULL; + char *rest = NULL; + sqlite3_stmt *statement = async_command->statement; + double *floats = NULL; int float_count = 0; @@ -106,13 +148,6 @@ static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) { int binaries_count = 0; int i; - - // fprintf(stderr, "Exec: %.*s\n", command_size, command); - - result = sqlite3_prepare_v2(drv->db, command, command_size, &statement, (const char **)&rest); - if(result != SQLITE_OK) { - return return_error(drv, sqlite3_errmsg(drv->db)); - } column_count = sqlite3_column_count(statement); dataset = NULL; @@ -122,28 +157,29 @@ static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) { dataset[term_count - 2] = ERL_DRV_PORT; dataset[term_count - 1] = driver_mk_port(drv->port); - - while ((next_row = sqlite3_step(statement)) == SQLITE_ROW) { - if (row_count == 0) { - int base = term_count; - term_count += 2 + column_count*2 + 1 + 2 + 2 + 2; - dataset = realloc(dataset, sizeof(*dataset) * term_count); - dataset[base] = ERL_DRV_ATOM; - dataset[base + 1] = driver_mk_atom("columns"); - for (i = 0; i < column_count; i++) { - dataset[base + 2 + (i*2)] = ERL_DRV_ATOM; - // fprintf(stderr, "Column: %s\n", sqlite3_column_name(statement, i)); - dataset[base + 2 + (i*2) + 1] = driver_mk_atom((char *)sqlite3_column_name(statement, i)); - } - dataset[base + 2 + column_count*2 + 0] = ERL_DRV_NIL; - dataset[base + 2 + column_count*2 + 1] = ERL_DRV_LIST; - dataset[base + 2 + column_count*2 + 2] = column_count + 1; - dataset[base + 2 + column_count*2 + 3] = ERL_DRV_TUPLE; - dataset[base + 2 + column_count*2 + 4] = 2; - - dataset[base + 2 + column_count*2 + 5] = ERL_DRV_ATOM; - dataset[base + 2 + column_count*2 + 6] = driver_mk_atom("rows"); + if (column_count > 0) { + int base = term_count; + term_count += 2 + column_count*2 + 1 + 2 + 2 + 2; + dataset = realloc(dataset, sizeof(*dataset) * term_count); + dataset[base] = ERL_DRV_ATOM; + dataset[base + 1] = driver_mk_atom("columns"); + for (i = 0; i < column_count; i++) { + dataset[base + 2 + (i*2)] = ERL_DRV_ATOM; + // fprintf(stderr, "Column: %s\n", sqlite3_column_name(statement, i)); + dataset[base + 2 + (i*2) + 1] = driver_mk_atom((char *)sqlite3_column_name(statement, i)); } + dataset[base + 2 + column_count*2 + 0] = ERL_DRV_NIL; + dataset[base + 2 + column_count*2 + 1] = ERL_DRV_LIST; + dataset[base + 2 + column_count*2 + 2] = column_count + 1; + dataset[base + 2 + column_count*2 + 3] = ERL_DRV_TUPLE; + dataset[base + 2 + column_count*2 + 4] = 2; + + dataset[base + 2 + column_count*2 + 5] = ERL_DRV_ATOM; + dataset[base + 2 + column_count*2 + 6] = driver_mk_atom("rows"); + } + + + while (column_count > 0 && (next_row = sqlite3_step(statement)) == SQLITE_ROW) { for (i = 0; i < column_count; i++) { // fprintf(stderr, "Column %d type: %d\n", i, sqlite3_column_type(statement, i)); @@ -155,17 +191,17 @@ static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) { dataset[term_count - 1] = sqlite3_column_int(statement, i); break; } - // case SQLITE_FLOAT: { - // float_count++; - // floats = realloc(floats, sizeof(double) * float_count); - // floats[float_count - 1] = sqlite3_column_double(statement, i); - // - // term_count += 2; - // dataset = realloc(dataset, sizeof(*dataset) * term_count); - // dataset[term_count - 2] = ERL_DRV_FLOAT; - // dataset[term_count - 1] = (ErlDrvTermData)&floats[float_count - 1]; - // break; - // } + case SQLITE_FLOAT: { + float_count++; + floats = realloc(floats, sizeof(double) * float_count); + floats[float_count - 1] = sqlite3_column_double(statement, i); + + term_count += 2; + dataset = realloc(dataset, sizeof(*dataset) * term_count); + dataset[term_count - 2] = ERL_DRV_FLOAT; + dataset[term_count - 1] = (ErlDrvTermData)&floats[float_count - 1]; + break; + } case SQLITE_BLOB: case SQLITE_TEXT: { int bytes = sqlite3_column_bytes(statement, i); @@ -198,49 +234,67 @@ static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) { if (next_row == SQLITE_BUSY) { sqlite3_finalize(statement); - return return_error(drv, "SQLite3 database is busy"); + return_error(drv, "SQLite3 database is busy"); } - term_count += 3; - dataset = realloc(dataset, sizeof(*dataset) * term_count); - dataset[term_count - 3] = ERL_DRV_NIL; - dataset[term_count - 2] = ERL_DRV_LIST; - dataset[term_count - 1] = row_count + 1; + if (column_count > 0) { + term_count += 3; + dataset = realloc(dataset, sizeof(*dataset) * term_count); + dataset[term_count - 3] = ERL_DRV_NIL; + dataset[term_count - 2] = ERL_DRV_LIST; + dataset[term_count - 1] = row_count + 1; - term_count += 2; - dataset = realloc(dataset, sizeof(*dataset) * term_count); - dataset[term_count - 2] = ERL_DRV_TUPLE; - dataset[term_count - 1] = 2; + term_count += 2; + dataset = realloc(dataset, sizeof(*dataset) * term_count); + dataset[term_count - 2] = ERL_DRV_TUPLE; + dataset[term_count - 1] = 2; - term_count += 3; - dataset = realloc(dataset, sizeof(*dataset) * term_count); - dataset[term_count - 3] = ERL_DRV_NIL; - dataset[term_count - 2] = ERL_DRV_LIST; - dataset[term_count - 1] = 3; + term_count += 3; + dataset = realloc(dataset, sizeof(*dataset) * term_count); + dataset[term_count - 3] = ERL_DRV_NIL; + dataset[term_count - 2] = ERL_DRV_LIST; + dataset[term_count - 1] = 3; + } else { + + sqlite3_last_insert_rowid(drv->db); + term_count += 4; + dataset = realloc(dataset, sizeof(*dataset) * term_count); + dataset[term_count - 4] = ERL_DRV_ATOM; + dataset[term_count - 3] = driver_mk_atom("ok"); + dataset[term_count - 2] = ERL_DRV_TUPLE; + dataset[term_count - 1] = 1; + } + term_count += 2; dataset = realloc(dataset, sizeof(*dataset) * term_count); dataset[term_count - 2] = ERL_DRV_TUPLE; dataset[term_count - 1] = 2; + - int res = driver_output_term(drv->port, dataset, term_count); - // fprintf(stderr, "Total term count: %d, rows count: %d (%d)\n", term_count, row_count, res); - free(dataset); - if (floats) { - free(floats); - } - for (i = 0; i < binaries_count; i++) { - driver_free_binary(binaries[i]); - } - if(binaries) { - free(binaries); - } - sqlite3_finalize(statement); + async_command->dataset = dataset; + async_command->term_count = term_count; + async_command->row_count = row_count; + async_command->floats = floats; + async_command->binaries = binaries; + async_command->binaries_count = binaries_count; + fprintf(stderr, "Total term count: %p %d, rows count: %dx%d\n", statement, term_count, column_count, row_count); } +static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data) +{ + async_sqlite3_command *async_command = (async_sqlite3_command *)thread_data; + sqlite3_drv_t *drv = async_command->driver_data; + + int res = driver_output_term(drv->port, async_command->dataset, async_command->term_count); + fprintf(stderr, "Total term count: %p %d, rows count: %d (%d)\n", async_command->statement, async_command->term_count, async_command->row_count, res); + sql_free_async(async_command); +} + + // Unkown Command static int unknown(sqlite3_drv_t *drv, char *command, int command_size) { // Return {error, unkown_command} diff --git a/priv/sqlite3_drv.h b/priv/sqlite3_drv.h index 5eb645b..62e44e3 100644 --- a/priv/sqlite3_drv.h +++ b/priv/sqlite3_drv.h @@ -18,19 +18,30 @@ #define KEY_SIZE 20 // Define struct to hold state across calls -typedef struct _bdb_drv_t { +typedef struct sqlite3_drv_t { ErlDrvPort port; - + unsigned int key; struct sqlite3 *db; } sqlite3_drv_t; +typedef struct async_sqlite3_command { + sqlite3_drv_t *driver_data; + sqlite3_stmt *statement; + ErlDrvTermData *dataset; + int term_count; + int row_count; + double *floats; + int binaries_count; + ErlDrvBinary **binaries; +} async_sqlite3_command; + static ErlDrvData start(ErlDrvPort port, char* cmd); static void stop(ErlDrvData handle); static int control(ErlDrvData drv_data, unsigned int command, char *buf, int len, char **rbuf, int rlen); -static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data); static int sql_exec(sqlite3_drv_t *drv, char *buf, int len); -// static void get(bdb_drv_t *bdb_drv, ErlIOVec *ev); -// static void del(bdb_drv_t *bdb_drv, ErlIOVec *ev); +static void sql_exec_async(void *async_command); +static void sql_free_async(void *async_command); +static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data); static int unknown(sqlite3_drv_t *bdb_drv, char *buf, int len);