The command processing loop is working now. Ready for the real deal

This commit is contained in:
Maas-Maarten Zeeman
2011-10-21 00:03:49 +02:00
parent 3a0eb12fd8
commit 87318f8764
2 changed files with 104 additions and 37 deletions

View File

@@ -2,19 +2,17 @@
* Esqlite -- an erlang sqlite nif. * Esqlite -- an erlang sqlite nif.
*/ */
#include <stdio.h> /* for debugging */
#include <assert.h> #include <assert.h>
#include <erl_nif.h> #include <erl_nif.h>
#include <stdio.h> /* for debugging */
#include "queue.h" #include "queue.h"
#include "sqlite3.h" #include "sqlite3.h"
#define MAX_PATHNAME 512 /* unfortunately not in sqlite.h. */ #define MAX_PATHNAME 512 /* unfortunately not in sqlite.h. */
static ErlNifResourceType *esqlite_sqlite3_type = NULL; static ErlNifResourceType *esqlite_db_type = NULL;
static ERL_NIF_TERM _atom_ok;
static ERL_NIF_TERM _atom_error;
/* database connection context */ /* database connection context */
typedef struct { typedef struct {
@@ -27,6 +25,9 @@ typedef struct {
int alive; int alive;
} esqlite_db; } esqlite_db;
static ERL_NIF_TERM _atom_ok;
static ERL_NIF_TERM _atom_error;
typedef enum { typedef enum {
cmd_unknown, cmd_unknown,
cmd_open, cmd_open,
@@ -46,6 +47,29 @@ typedef struct {
} esqlite_command; } esqlite_command;
static ERL_NIF_TERM
make_atom(ErlNifEnv *env, const char *atom_name)
{
ERL_NIF_TERM atom;
if(enif_make_existing_atom(env, atom_name, &atom, ERL_NIF_LATIN1))
return atom;
return enif_make_atom(env, atom_name);
}
static ERL_NIF_TERM
make_ok_tuple(ErlNifEnv *env, ERL_NIF_TERM value)
{
return enif_make_tuple2(env, _atom_ok, value);
}
static ERL_NIF_TERM
make_error_tuple(ErlNifEnv *env, const char *reason)
{
return enif_make_tuple2(env, _atom_error, make_atom(env, reason));
}
static void static void
command_destroy(void *obj) command_destroy(void *obj)
{ {
@@ -82,8 +106,8 @@ descruct_esqlite_db(ErlNifEnv *env, void *arg)
{ {
esqlite_db *db = (esqlite_db *) arg; esqlite_db *db = (esqlite_db *) arg;
esqlite_command *cmd = command_create(); esqlite_command *cmd = command_create();
/* send the stop command */ /* Send the stop command */
cmd->type = cmd_stop; cmd->type = cmd_stop;
queue_push(db->commands, cmd); queue_push(db->commands, cmd);
queue_send(db->commands, cmd); queue_send(db->commands, cmd);
@@ -104,28 +128,29 @@ esqlite_db_run(void *arg)
while(1) { while(1) {
cmd = queue_pop(db->commands); cmd = queue_pop(db->commands);
/* We are stopping... */
if(cmd_stop == cmd->type) { if(cmd_stop == cmd->type) {
fprintf(stderr, "received stop\n");
command_destroy(cmd); command_destroy(cmd);
break; break;
} }
/* Evaluate the command */ /* Evaluate the command */
switch(cmd->type) { switch(cmd->type) {
cmd_open: case cmd_open:
/* do open */ fprintf(stderr, "received open\n");
break; break;
cmd_exec: case cmd_exec:
/* do exec */ fprintf(stderr, "received exec\n");
break; break;
cmd_close: case cmd_close:
/* do close */ fprintf(stderr, "received close\n");
break; break;
default: default:
assert(0 && "Invalid command"); assert(0 && "Invalid command");
} }
enif_send(NULL, &(cmd->pid), cmd->env, _atom_ok); /* TODO: A real implementation for a command */
enif_send(NULL, &cmd->pid, cmd->env, enif_make_tuple2(cmd->env, cmd->ref, _atom_ok));
command_destroy(cmd); command_destroy(cmd);
} }
@@ -133,55 +158,98 @@ esqlite_db_run(void *arg)
return NULL; return NULL;
} }
/* /*
* Open database. Expects utf-8 input * Start the processing thread
*
* Note the database is opened in a thread. New commands are send to
* the thread and when it finishes the result is send back. The reason
* for this is that we don't want to block the erlang scheduler of the
* calling function.
*/ */
static ERL_NIF_TERM static ERL_NIF_TERM
start_nif(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[]) start_nif(ErlNifEnv* env, int argc, const ERL_NIF_TERM argv[])
{ {
esqlite_db *esqldb; esqlite_db *esqldb;
ERL_NIF_TERM esqlite_db; ERL_NIF_TERM db;
/* initialize the resource */ /* Initialize the resource */
esqldb = enif_alloc_resource(esqlite_sqlite3_type, sizeof(esqlite_db)); esqldb = enif_alloc_resource(esqlite_db_type, sizeof(esqlite_db));
esqldb->db = NULL; esqldb->db = NULL;
/* Start the command processing thread */ /* Create command queue */
esqldb->commands = queue_create();
if(!esqldb->commands) {
enif_release_resource(esqldb);
return make_error_tuple(env, "command_queue_create_failed");
}
/* Start command processing thread */
esqldb->opts = enif_thread_opts_create("esqldb_thread_opts"); esqldb->opts = enif_thread_opts_create("esqldb_thread_opts");
if(enif_thread_create("", &esqldb->tid, esqlite_db_run, esqldb, esqldb->opts) != 0) { if(enif_thread_create("", &esqldb->tid, esqlite_db_run, esqldb, esqldb->opts) != 0) {
enif_release_resource(esqldb); enif_release_resource(esqldb);
return enif_make_tuple2(env, _atom_error, _atom_ok); return make_error_tuple(env, "thread_create_failed");
} }
/* We got the resource... now return it */ db = enif_make_resource(env, esqldb);
esqlite_db = enif_make_resource(env, esqldb);
enif_release_resource(esqldb); enif_release_resource(esqldb);
return enif_make_tuple2(env, _atom_ok, esqlite_db);
return make_ok_tuple(env, db);
} }
/*
* Open the database
*/
static ERL_NIF_TERM static ERL_NIF_TERM
esqlite_open_nif(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[]) esqlite_open_nif(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[])
{ {
esqlite_db *db;
esqlite_command *cmd = NULL;
ErlNifPid pid;
if(argc != 4)
return enif_make_badarg(env);
if(!enif_get_resource(env, argv[0], esqlite_db_type, (void **) &db))
return enif_make_badarg(env);
if(!enif_is_ref(env, argv[1]))
return make_error_tuple(env, "invalid_ref");
if(!enif_get_local_pid(env, argv[2], &pid))
return make_error_tuple(env, "invalid_pid");
cmd = command_create();
if(!cmd)
return make_error_tuple(env, "command_create_failed");
/* command */
cmd->type = cmd_open;
cmd->ref = enif_make_copy(cmd->env, argv[1]);
cmd->pid = pid;
/* todo add the filename */
if(!queue_push(db->commands, cmd))
return make_error_tuple(env, "command_push_failed");
return _atom_ok; return _atom_ok;
} }
/*
* Execute the sql statement
*/
static ERL_NIF_TERM static ERL_NIF_TERM
esqlite_exec_nif(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[]) esqlite_exec_nif(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[])
{ {
return _atom_ok; return _atom_ok;
} }
/*
* Close the database
*/
static ERL_NIF_TERM static ERL_NIF_TERM
esqlite_close_nif(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[]) esqlite_close_nif(ErlNifEnv *env, int argc, const ERL_NIF_TERM argv[])
{ {
return _atom_ok; return _atom_ok;
} }
/*
* Load the nif. Initialize some stuff and such
*/
static int static int
on_load(ErlNifEnv* env, void** priv, ERL_NIF_TERM info) on_load(ErlNifEnv* env, void** priv, ERL_NIF_TERM info)
{ {
@@ -190,18 +258,17 @@ on_load(ErlNifEnv* env, void** priv, ERL_NIF_TERM info)
if(!rt) if(!rt)
return -1; return -1;
esqlite_sqlite3_type = rt; esqlite_db_type = rt;
_atom_ok = enif_make_atom(env, "ok"); _atom_ok = make_atom(env, "ok");
_atom_error = enif_make_atom(env, "error"); _atom_error = make_atom(env, "error");
return 0; return 0;
} }
static ErlNifFunc nif_funcs[] = { static ErlNifFunc nif_funcs[] = {
{"esqlite_start", 0, start_nif}, {"esqlite_start", 0, start_nif},
{"esqlite_open", 2, esqlite_open_nif}, {"esqlite_open", 4, esqlite_open_nif},
{"esqlite_exec", 4, esqlite_exec_nif}, {"esqlite_exec", 4, esqlite_exec_nif},
{"esqlite_close", 3, esqlite_close_nif} {"esqlite_close", 3, esqlite_close_nif}
}; };

View File

@@ -22,10 +22,10 @@ open(Filename) ->
%% @doc Open a database connection %% @doc Open a database connection
%% %%
open(Filename, Timeout) -> open(Filename, Timeout) ->
Db = esqlite_start(), {ok, Db} = esqlite_start(),
Ref = make_ref(), Ref = make_ref(),
ok = esqlite_open(Db, Ref, Filename, self()), ok = esqlite_open(Db, Ref, self(), Filename),
case receive_answer(Ref, Timeout) of case receive_answer(Ref, Timeout) of
ok -> ok ->
{ok, Db}; {ok, Db};