now can in right way reply to empty resultsets

This commit is contained in:
Max Lapshin
2009-10-14 01:57:04 +04:00
parent be08fe62f3
commit de53ecd1c0
2 changed files with 146 additions and 81 deletions

View File

@@ -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;
@@ -107,13 +149,6 @@ static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) {
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,14 +234,38 @@ 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 += 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);
@@ -213,34 +273,28 @@ static int sql_exec(sqlite3_drv_t *drv, char *command, int command_size) {
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 += 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}

View File

@@ -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);