Working with output
This commit is contained in:
@@ -11,10 +11,10 @@ static ErlDrvEntry basic_driver_entry = {
|
|||||||
"sqlite3_drv", /* the name of the driver */
|
"sqlite3_drv", /* the name of the driver */
|
||||||
NULL, /* finish */
|
NULL, /* finish */
|
||||||
NULL, /* handle */
|
NULL, /* handle */
|
||||||
NULL, /* control */
|
control, /* control */
|
||||||
NULL, /* timeout */
|
NULL, /* timeout */
|
||||||
outputv, /* outputv (defined below) */
|
NULL, /* outputv (defined below) */
|
||||||
ready_async, /* ready_async */
|
NULL, /* ready_async */
|
||||||
NULL, /* flush */
|
NULL, /* flush */
|
||||||
NULL, /* call */
|
NULL, /* call */
|
||||||
NULL, /* event */
|
NULL, /* event */
|
||||||
@@ -60,29 +60,18 @@ static void stop(ErlDrvData handle) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Handle input from Erlang VM
|
// Handle input from Erlang VM
|
||||||
static void outputv(ErlDrvData handle, ErlIOVec *ev) {
|
static int control(ErlDrvData drv_data, unsigned int command, char *buf,
|
||||||
sqlite3_drv_t* driver_data = (sqlite3_drv_t*) handle;
|
int len, char **rbuf, int rlen) {
|
||||||
ErlDrvBinary* data = ev->binv[1];
|
sqlite3_drv_t* driver_data = (sqlite3_drv_t*) drv_data;
|
||||||
|
|
||||||
int command = data->orig_bytes[0];
|
|
||||||
|
|
||||||
|
|
||||||
switch(command) {
|
switch(command) {
|
||||||
case CMD_SQL_EXEC:
|
case CMD_SQL_EXEC:
|
||||||
sql_exec(driver_data, ev);
|
sql_exec(driver_data, buf, len);
|
||||||
break;
|
break;
|
||||||
//
|
|
||||||
// // case CMD_GET:
|
|
||||||
// // get(driver_data, ev);
|
|
||||||
// // break;
|
|
||||||
// //
|
|
||||||
// // case CMD_DEL:
|
|
||||||
// // del(driver_data, ev);
|
|
||||||
// // break;
|
|
||||||
//
|
|
||||||
default:
|
default:
|
||||||
unkown(driver_data, ev);
|
unknown(driver_data, buf, len);
|
||||||
}
|
}
|
||||||
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data)
|
static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data)
|
||||||
@@ -90,202 +79,143 @@ static void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data)
|
|||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
static int callback(void *data, int argc, char **argv, char **azColName)
|
static inline int return_error(sqlite3_drv_t *drv, const char *error) {
|
||||||
{
|
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
||||||
sqlite3_drv_t *drv = (sqlite3_drv_t *)data;
|
ERL_DRV_STRING, (ErlDrvTermData)error, strlen(error),
|
||||||
|
ERL_DRV_TUPLE, 2};
|
||||||
|
|
||||||
|
return driver_output_term(drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
||||||
|
if (error) {
|
||||||
|
// sqlite3_free((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;
|
||||||
|
char *rest = NULL;
|
||||||
|
sqlite3_stmt *statement;
|
||||||
|
|
||||||
|
double *floats = NULL;
|
||||||
|
int float_count = 0;
|
||||||
|
|
||||||
|
ErlDrvBinary **binaries = NULL;
|
||||||
|
int binaries_count = 0;
|
||||||
int i;
|
int i;
|
||||||
|
|
||||||
// ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
|
||||||
// ERL_DRV_STRING, error_reason, strlen(error_reason),
|
|
||||||
// ERL_DRV_TUPLE, 2};
|
|
||||||
//
|
|
||||||
// driver_output_term(drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
|
||||||
|
|
||||||
|
fprintf(stderr, "Exec: %.*s\n", command_size, command);
|
||||||
|
|
||||||
|
result = sqlite3_prepare_v2(drv->db, command, command_size, &statement, (const char **)&rest);
|
||||||
// record_list = malloc(argc * sizeof(ETERM *));
|
if(result != SQLITE_OK) {
|
||||||
|
return return_error(drv, sqlite3_errmsg(drv->db));
|
||||||
fprintf(stderr, "runs %d\n", argc);
|
|
||||||
for (i = 0; i < argc; i++) {
|
|
||||||
fprintf(stderr, "%s = %s\n", azColName[i], argv[i] ? argv[i] : "NULL");
|
|
||||||
}
|
}
|
||||||
fprintf(stderr, "\n");
|
|
||||||
fflush(stderr);
|
|
||||||
|
|
||||||
// result = erl_cons(erl_mk_tuple(record_list, argc), result);
|
column_count = sqlite3_column_count(statement);
|
||||||
|
dataset = NULL;
|
||||||
|
|
||||||
// free(record_list);
|
fprintf(stderr, "Going to read some rows with %d columns\n", column_count);
|
||||||
return 0;
|
while ((next_row = sqlite3_step(statement)) == SQLITE_ROW) {
|
||||||
}
|
if (row_count == 0) {
|
||||||
|
term_count = 2 + column_count*2 + 2 + 2;
|
||||||
|
dataset = realloc(dataset, sizeof(*dataset) * term_count);
|
||||||
static void sql_exec(sqlite3_drv_t *drv, ErlIOVec *ev) {
|
dataset[0] = ERL_DRV_ATOM;
|
||||||
|
dataset[1] = driver_mk_atom("columns");
|
||||||
ErlDrvBinary* input = ev->binv[1];
|
for (i = 0; i < column_count; i++) {
|
||||||
char *command = input->orig_bytes + 1;
|
dataset[2 + (i*2)] = ERL_DRV_ATOM;
|
||||||
int command_size = input->orig_size - 1;
|
fprintf(stderr, "Column: %s\n", sqlite3_column_name(statement, i));
|
||||||
int status;
|
dataset[2 + (i*2) + 1] = driver_mk_atom((char *)sqlite3_column_name(statement, i));
|
||||||
char *error = NULL;
|
}
|
||||||
|
dataset[2 + column_count*2] = ERL_DRV_LIST;
|
||||||
fprintf(stderr, "Exec: %*s\n", command_size, command);
|
dataset[2 + column_count*2 + 1] = column_count;
|
||||||
|
dataset[2 + column_count*2 + 2] = ERL_DRV_TUPLE;
|
||||||
status = sqlite3_exec(drv->db, command, callback, drv, &error);
|
dataset[2 + column_count*2 + 3] = 2;
|
||||||
|
|
||||||
if(status != SQLITE_OK) {
|
|
||||||
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
|
||||||
ERL_DRV_STRING, error, strlen(error),
|
|
||||||
ERL_DRV_TUPLE, 2};
|
|
||||||
|
|
||||||
driver_output_term(drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
|
||||||
}
|
|
||||||
if (error) {
|
|
||||||
sqlite3_free(error);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#if 0
|
|
||||||
|
|
||||||
// Retrieve a record from the database, if it exists
|
|
||||||
static void get(sqlite3_drv_t *bdb_drv, ErlIOVec *ev) {
|
|
||||||
ErlDrvBinary* input = ev->binv[1];
|
|
||||||
ErlDrvBinary *output_bytes;
|
|
||||||
char *bytes = input->orig_bytes;
|
|
||||||
char *key_bytes = bytes+1;
|
|
||||||
|
|
||||||
DB *db = bdb_drv->db;
|
|
||||||
DBT key;
|
|
||||||
DBT value;
|
|
||||||
int status;
|
|
||||||
|
|
||||||
bzero(&key, sizeof(DBT));
|
|
||||||
bzero(&value, sizeof(DBT));
|
|
||||||
|
|
||||||
key.data = key_bytes;
|
|
||||||
key.size = KEY_SIZE;
|
|
||||||
|
|
||||||
// Have BerkeleyDB allocate memory big enough to store the value
|
|
||||||
value.flags = DB_DBT_MALLOC; // Don't forget to free it later
|
|
||||||
|
|
||||||
// Retrieve the record
|
|
||||||
status = db->get(db, NULL, &key, &value, 0);
|
|
||||||
|
|
||||||
if(status == 0) {
|
|
||||||
// Get went OK
|
|
||||||
|
|
||||||
// Copy the record value to an output structure to return to Erlang VM
|
|
||||||
output_bytes = driver_alloc_binary(value.size);
|
|
||||||
output_bytes->orig_size = value.size;
|
|
||||||
memcpy(output_bytes->orig_bytes, value.data, value.size);
|
|
||||||
free(value.data);
|
|
||||||
|
|
||||||
// TODO:Figure out if we can somehow use this original memory without recopying a la:
|
|
||||||
//binary->orig_bytes = (char *)&data.data;
|
|
||||||
|
|
||||||
// Returns tuple {ok, Data}
|
|
||||||
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("ok"),
|
|
||||||
ERL_DRV_BINARY, (ErlDrvTermData) output_bytes, output_bytes->orig_size, 0,
|
|
||||||
ERL_DRV_TUPLE, 2};
|
|
||||||
|
|
||||||
driver_output_term(bdb_drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
|
||||||
driver_free_binary(output_bytes);
|
|
||||||
} else {
|
|
||||||
// there was an error
|
|
||||||
char *error_reason;
|
|
||||||
|
|
||||||
switch(status) {
|
|
||||||
case DB_LOCK_DEADLOCK:
|
|
||||||
error_reason = "deadlock";
|
|
||||||
break;
|
|
||||||
case DB_SECONDARY_BAD:
|
|
||||||
error_reason = "bad_secondary_index";
|
|
||||||
break;
|
|
||||||
case ENOMEM:
|
|
||||||
error_reason = "insufficient_memory";
|
|
||||||
break;
|
|
||||||
case EINVAL:
|
|
||||||
error_reason = "bad_flag";
|
|
||||||
break;
|
|
||||||
case DB_RUNRECOVERY:
|
|
||||||
error_reason = "run_recovery";
|
|
||||||
break;
|
|
||||||
default:
|
|
||||||
error_reason = "unknown";
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Return tuple {error, Reason}
|
for (i = 0; i < column_count; i++) {
|
||||||
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
fprintf(stderr, "Column %d type: %d\n", i, sqlite3_column_type(statement, i));
|
||||||
ERL_DRV_ATOM, driver_mk_atom(error_reason),
|
switch (sqlite3_column_type(statement, i)) {
|
||||||
ERL_DRV_TUPLE, 2};
|
case SQLITE_INTEGER: {
|
||||||
driver_output_term(bdb_drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
term_count += 2;
|
||||||
}
|
dataset = realloc(dataset, sizeof(*dataset) * term_count);
|
||||||
}
|
dataset[term_count - 2] = ERL_DRV_INT;
|
||||||
|
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_BLOB:
|
||||||
|
case SQLITE_TEXT: {
|
||||||
|
int bytes = sqlite3_column_bytes(statement, i);
|
||||||
|
binaries_count++;
|
||||||
|
binaries = realloc(binaries, sizeof(*binaries) * binaries_count);
|
||||||
|
binaries[binaries_count - 1] = driver_alloc_binary(bytes);
|
||||||
|
binaries[binaries_count - 1]->orig_size = bytes;
|
||||||
|
memcpy(binaries[binaries_count - 1]->orig_bytes, sqlite3_column_blob(statement, i), bytes);
|
||||||
|
|
||||||
// Delete a record from the database
|
term_count += 2;
|
||||||
static void del(sqlite3_drv_t *bdb_drv, ErlIOVec *ev) {
|
dataset = realloc(dataset, sizeof(*dataset) * term_count);
|
||||||
ErlDrvBinary* data = ev->binv[1];
|
dataset[term_count - 2] = ERL_DRV_FLOAT;
|
||||||
char *bytes = data->orig_bytes;
|
dataset[term_count - 1] = (ErlDrvTermData)&floats[float_count - 1];
|
||||||
char *key_bytes = bytes+1;
|
break;
|
||||||
|
}
|
||||||
DB *db = bdb_drv->db;
|
case SQLITE_NULL: {
|
||||||
DBT key;
|
break;
|
||||||
int status;
|
}
|
||||||
|
}
|
||||||
bzero(&key, sizeof(DBT));
|
|
||||||
|
|
||||||
key.data = key_bytes;
|
|
||||||
key.size = KEY_SIZE;
|
|
||||||
|
|
||||||
status = db->del(db, NULL, &key, 0);
|
|
||||||
db->sync(db, 0);
|
|
||||||
|
|
||||||
if(status == 0) {
|
|
||||||
// Delete went OK, return atom 'ok'
|
|
||||||
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("ok")};
|
|
||||||
|
|
||||||
driver_output_term(bdb_drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
|
||||||
|
|
||||||
} else {
|
|
||||||
// There was an error
|
|
||||||
char *error_reason;
|
|
||||||
|
|
||||||
switch(status) {
|
|
||||||
case DB_NOTFOUND:
|
|
||||||
error_reason = "not_found";
|
|
||||||
break;
|
|
||||||
case DB_LOCK_DEADLOCK:
|
|
||||||
error_reason = "deadlock";
|
|
||||||
break;
|
|
||||||
case DB_SECONDARY_BAD:
|
|
||||||
error_reason = "bad_secondary_index";
|
|
||||||
break;
|
|
||||||
case EINVAL:
|
|
||||||
error_reason = "bad_flag";
|
|
||||||
break;
|
|
||||||
case EACCES:
|
|
||||||
error_reason = "readonly";
|
|
||||||
break;
|
|
||||||
case DB_RUNRECOVERY:
|
|
||||||
error_reason = "run_recovery";
|
|
||||||
break;
|
|
||||||
default:
|
|
||||||
error_reason = "unknown";
|
|
||||||
}
|
}
|
||||||
|
term_count += 2;
|
||||||
|
dataset = realloc(dataset, sizeof(*dataset) * term_count);
|
||||||
|
dataset[term_count - 2] = ERL_DRV_TUPLE;
|
||||||
|
dataset[term_count - 1] = column_count;
|
||||||
|
|
||||||
// Return tuple {error, Reason}
|
row_count++;
|
||||||
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
|
||||||
ERL_DRV_ATOM, driver_mk_atom(error_reason),
|
|
||||||
ERL_DRV_TUPLE, 2};
|
|
||||||
|
|
||||||
driver_output_term(bdb_drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (next_row == SQLITE_BUSY) {
|
||||||
|
sqlite3_finalize(statement);
|
||||||
|
return return_error(drv, "SQLite3 database is busy");
|
||||||
|
}
|
||||||
|
|
||||||
|
term_count += 2;
|
||||||
|
dataset = realloc(dataset, sizeof(*dataset) * term_count);
|
||||||
|
dataset[term_count - 2] = ERL_DRV_LIST;
|
||||||
|
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);
|
||||||
}
|
}
|
||||||
#endif
|
|
||||||
|
|
||||||
// Unkown Command
|
// Unkown Command
|
||||||
static void unkown(sqlite3_drv_t *bdb_drv, ErlIOVec *ev) {
|
static int unknown(sqlite3_drv_t *drv, char *command, int command_size) {
|
||||||
// Return {error, unkown_command}
|
// Return {error, unkown_command}
|
||||||
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
ErlDrvTermData spec[] = {ERL_DRV_ATOM, driver_mk_atom("error"),
|
||||||
ERL_DRV_ATOM, driver_mk_atom("uknown_command"),
|
ERL_DRV_ATOM, driver_mk_atom("uknown_command"),
|
||||||
ERL_DRV_TUPLE, 2};
|
ERL_DRV_TUPLE, 2};
|
||||||
driver_output_term(bdb_drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
return driver_output_term(drv->port, spec, sizeof(spec) / sizeof(spec[0]));
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -27,9 +27,10 @@ typedef struct _bdb_drv_t {
|
|||||||
|
|
||||||
static ErlDrvData start(ErlDrvPort port, char* cmd);
|
static ErlDrvData start(ErlDrvPort port, char* cmd);
|
||||||
static void stop(ErlDrvData handle);
|
static void stop(ErlDrvData handle);
|
||||||
static void outputv(ErlDrvData handle, ErlIOVec *ev);
|
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 void ready_async(ErlDrvData drv_data, ErlDrvThreadData thread_data);
|
||||||
static void sql_exec(sqlite3_drv_t *drv, ErlIOVec *ev);
|
static int sql_exec(sqlite3_drv_t *drv, char *buf, int len);
|
||||||
// static void get(bdb_drv_t *bdb_drv, ErlIOVec *ev);
|
// static void get(bdb_drv_t *bdb_drv, ErlIOVec *ev);
|
||||||
// static void del(bdb_drv_t *bdb_drv, ErlIOVec *ev);
|
// static void del(bdb_drv_t *bdb_drv, ErlIOVec *ev);
|
||||||
static void unkown(sqlite3_drv_t *bdb_drv, ErlIOVec *ev);
|
static int unknown(sqlite3_drv_t *bdb_drv, char *buf, int len);
|
||||||
|
|||||||
@@ -331,7 +331,7 @@ init(Options) ->
|
|||||||
SearchDir = filename:join([filename:dirname(code:which(?MODULE)), "..", "ebin"]),
|
SearchDir = filename:join([filename:dirname(code:which(?MODULE)), "..", "ebin"]),
|
||||||
case erl_ddll:load(SearchDir, atom_to_list(?DRIVER_NAME)) of
|
case erl_ddll:load(SearchDir, atom_to_list(?DRIVER_NAME)) of
|
||||||
ok ->
|
ok ->
|
||||||
Port = open_port({spawn, ?DRIVER_NAME}, [{packet, 2}, binary]),
|
Port = open_port({spawn, ?DRIVER_NAME}, [binary]),
|
||||||
{ok, #state{port = Port, ops = Options}};
|
{ok, #state{port = Port, ops = Options}};
|
||||||
{error, Error} ->
|
{error, Error} ->
|
||||||
io:format("Error loading ~p: ~p", [?DRIVER_NAME, erl_ddll:format_error(Error)]),
|
io:format("Error loading ~p: ~p", [?DRIVER_NAME, erl_ddll:format_error(Error)]),
|
||||||
@@ -456,18 +456,25 @@ create_cmd(Dbase) ->
|
|||||||
|
|
||||||
wait_result(Port) ->
|
wait_result(Port) ->
|
||||||
receive
|
receive
|
||||||
{Port, {data, Data}} when is_binary(Data) ->
|
{Port, Reply} ->
|
||||||
List = binary_to_term(Data),
|
io:format("Reply: ~p~n", [Reply]),
|
||||||
if is_list(List) ->
|
Reply;
|
||||||
lists:reverse(List);
|
% List = binary_to_term(Data),
|
||||||
true -> List
|
% if is_list(List) ->
|
||||||
end;
|
% lists:reverse(List);
|
||||||
_ ->
|
% true -> List
|
||||||
ok
|
% end;
|
||||||
|
{error, Reason} ->
|
||||||
|
io:format("Error: ~p~n", [Reason]),
|
||||||
|
{error, Reason};
|
||||||
|
_Else ->
|
||||||
|
io:format("Else: ~p~n", [_Else]),
|
||||||
|
_Else
|
||||||
end.
|
end.
|
||||||
|
|
||||||
exec(Port, {sql_exec, Cmd}) ->
|
exec(Port, {sql_exec, Cmd}) ->
|
||||||
port_command(Port, <<?SQL_EXEC_COMMAND, (list_to_binary(Cmd))/binary>>),
|
port_control(Port, ?SQL_EXEC_COMMAND, list_to_binary(Cmd)),
|
||||||
|
io:format("Returned from driver, waiting for reply~n"),
|
||||||
wait_result(Port).
|
wait_result(Port).
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user