Packages
rocksdb
0.24.0
3.1.1
3.1.0
3.0.0
2.6.2
2.6.1
retired
2.6.0
retired
2.5.0
2.4.1
2.4.0
2.3.0
2.2.0
2.1.0
2.0.0
1.9.0
1.8.0
1.7.0
1.6.0
1.5.1
1.5.0
1.4.0
1.3.2
1.3.1
1.3.0
1.2.0
1.1.1
1.1.0
1.0.0
0.26.2
0.26.1
0.26.0
0.25.0
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.0
0.21.0
0.20.1
0.20.0
0.19.0
0.18.0
0.17.0
0.16.0
0.15.0
0.14.0
0.13.1
0.13.0
0.12.0
0.11.0
0.10.0
0.9.1
0.9.0
0.8.2
0.8.1
0.8.0
0.7.1
0.7.0
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
RocksDB for Erlang
Current section
Files
Jump to
Current section
Files
c_src/erocksdb_iter.cc
// -------------------------------------------------------------------
// Copyright (c) 2016-2017 Benoit Chesneau. All Rights Reserved.
//
// This file is provided to you under the Apache License,
// Version 2.0 (the "License"); you may not use this file
// except in compliance with the License. You may obtain
// a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//
// -------------------------------------------------------------------
//
#include <mutex>
#include <vector>
#include <memory>
#include "rocksdb/db.h"
#include "rocksdb/comparator.h"
#include "rocksdb/write_batch.h"
#include "rocksdb/slice_transform.h"
#include "atoms.h"
#include "erl_nif.h"
#include "erocksdb_db.h"
#include "refobjects.h"
#include "util.h"
struct ItrBounds {
rocksdb::ReadOptions* read_options;
rocksdb::Slice *upper_bound_slice;
rocksdb::Slice *lower_bound_slice;
ItrBounds();
};
ItrBounds::ItrBounds()
: upper_bound_slice(nullptr),
lower_bound_slice(nullptr) {};
int
parse_iterator_options(
ErlNifEnv* env,
ErlNifEnv* itr_env,
ERL_NIF_TERM term,
rocksdb::ReadOptions& opts,
ItrBounds& bounds)
{
const ERL_NIF_TERM* option;
ERL_NIF_TERM head, tail;
int arity;
if(!enif_is_list(env, term))
return 0;
tail = term;
while(enif_get_list_cell(env, tail, &head, &tail))
{
if (enif_get_tuple(env, head, &arity, &option) && 2==arity)
{
if (option[0] == erocksdb::ATOM_VERIFY_CHECKSUMS)
opts.verify_checksums = (option[1] == erocksdb::ATOM_TRUE);
else if (option[0] == erocksdb::ATOM_FILL_CACHE)
opts.fill_cache = (option[1] == erocksdb::ATOM_TRUE);
else if (option[0] == erocksdb::ATOM_ITERATE_UPPER_BOUND)
{
ERL_NIF_TERM upper_bound = enif_make_copy(itr_env, option[1]);
ErlNifBinary upper_bound_bin;
if(!enif_inspect_binary(itr_env, upper_bound, &upper_bound_bin))
return 0;
bounds.upper_bound_slice = new rocksdb::Slice(
reinterpret_cast<char*>(upper_bound_bin.data),
upper_bound_bin.size);
opts.iterate_upper_bound = bounds.upper_bound_slice;
}
else if (option[0] == erocksdb::ATOM_ITERATE_LOWER_BOUND)
{
ERL_NIF_TERM lower_bound = enif_make_copy(itr_env, option[1]);
ErlNifBinary lower_bound_bin;
if(!enif_inspect_binary(itr_env, lower_bound, &lower_bound_bin))
return 0;
bounds.lower_bound_slice = new rocksdb::Slice(
reinterpret_cast<char*>(lower_bound_bin.data),
lower_bound_bin.size);
opts.iterate_lower_bound = bounds.lower_bound_slice;
}
else if (option[0] == erocksdb::ATOM_TAILING)
opts.tailing = (option[1] == erocksdb::ATOM_TRUE);
else if (option[0] == erocksdb::ATOM_TOTAL_ORDER_SEEK)
opts.total_order_seek = (option[1] == erocksdb::ATOM_TRUE);
else if (option[0] == erocksdb::ATOM_PREFIX_SAME_AS_START)
opts.prefix_same_as_start = (option[1] == erocksdb::ATOM_TRUE);
else if (option[0] == erocksdb::ATOM_SNAPSHOT)
{
erocksdb::ReferencePtr<erocksdb::SnapshotObject> snapshot_ptr;
snapshot_ptr.assign(erocksdb::SnapshotObject::RetrieveSnapshotObject(env, option[1]));
if(NULL==snapshot_ptr.get())
return 0;
opts.snapshot = snapshot_ptr->m_Snapshot;
}
}
}
return 1;
}
namespace erocksdb {
ERL_NIF_TERM
Iterator(
ErlNifEnv* env,
int argc,
const ERL_NIF_TERM argv[])
{
ReferencePtr<DbObject> db_ptr;
if(!enif_get_db(env, argv[0], &db_ptr))
return enif_make_badarg(env);
int i = 1;
if(argc==3) i = 2;
if(!enif_is_list(env, argv[i]))
return enif_make_badarg(env);
rocksdb::ReadOptions *opts = new rocksdb::ReadOptions;
ItrBounds bounds;
auto itr_env = std::make_shared<ErlEnvCtr>();
if (!parse_iterator_options(env, itr_env->env, argv[i], *opts, bounds))
{
delete opts;
return enif_make_badarg(env);
}
ItrObject * itr_ptr;
rocksdb::Iterator * iterator;
if(argc==3)
{
ReferencePtr<ColumnFamilyObject> cf_ptr;
if(!enif_get_cf(env, argv[1], &cf_ptr))
{
delete opts;
return enif_make_badarg(env);
}
iterator = db_ptr->m_Db->NewIterator(*opts, cf_ptr->m_ColumnFamily);
}
else
{
iterator = db_ptr->m_Db->NewIterator(*opts);
}
itr_ptr = ItrObject::CreateItrObject(db_ptr.get(), itr_env, iterator);
if(bounds.upper_bound_slice != nullptr)
{
itr_ptr->SetUpperBoundSlice(bounds.upper_bound_slice);
}
if(bounds.lower_bound_slice != nullptr)
{
itr_ptr->SetLowerBoundSlice(bounds.lower_bound_slice);
}
ERL_NIF_TERM result = enif_make_resource(env, itr_ptr);
// release reference created during CreateItrObject()
enif_release_resource(itr_ptr);
delete opts;
iterator = NULL;
return enif_make_tuple2(env, ATOM_OK, result);
} // erocksdb::Iterator
ERL_NIF_TERM
Iterators(
ErlNifEnv* env,
int argc,
const ERL_NIF_TERM argv[])
{
ReferencePtr<DbObject> db_ptr;
if(!enif_get_db(env, argv[0], &db_ptr))
return enif_make_badarg(env);
if(!enif_is_list(env, argv[1]) || !enif_is_list(env, argv[2]))
return enif_make_badarg(env);
rocksdb::ReadOptions *opts = new rocksdb::ReadOptions();
ItrBounds bounds;
auto itr_env = std::make_shared<ErlEnvCtr>();
if (!parse_iterator_options(env, itr_env->env, argv[2], *opts, bounds))
{
delete opts;
return enif_make_badarg(env);
}
std::vector<rocksdb::ColumnFamilyHandle*> column_families;
ERL_NIF_TERM head, tail = argv[1];
while(enif_get_list_cell(env, tail, &head, &tail))
{
ReferencePtr<ColumnFamilyObject> cf_ptr;
cf_ptr.assign(ColumnFamilyObject::RetrieveColumnFamilyObject(env, head));
ColumnFamilyObject* cf = cf_ptr.get();
column_families.push_back(cf->m_ColumnFamily);
}
std::vector<rocksdb::Iterator*> iterators;
db_ptr->m_Db->NewIterators(*opts, column_families, &iterators);
ERL_NIF_TERM result = enif_make_list(env, 0);
ItrObject * itr_ptr;
try {
for (size_t i = 0; i < iterators.size(); i++) {
itr_ptr = ItrObject::CreateItrObject(db_ptr.get(), itr_env, iterators[i]);
if(bounds.upper_bound_slice != nullptr)
{
itr_ptr->SetUpperBoundSlice(bounds.upper_bound_slice);
}
if(bounds.lower_bound_slice != nullptr)
{
itr_ptr->SetLowerBoundSlice(bounds.lower_bound_slice);
}
ERL_NIF_TERM itr_res = enif_make_resource(env, itr_ptr);
result = enif_make_list_cell(env, itr_res, result);
enif_release_resource(itr_ptr);
}
} catch (const std::exception& e) {
// pass through and return nullptr
}
opts=NULL;
ERL_NIF_TERM result_out;
enif_make_reverse_list(env, result, &result_out);
return enif_make_tuple2(env, erocksdb::ATOM_OK, result_out);
}
ERL_NIF_TERM
IteratorMove(
ErlNifEnv* env,
int argc,
const ERL_NIF_TERM argv[])
{
const ERL_NIF_TERM& itr_handle_ref = argv[0];
const ERL_NIF_TERM& action_or_target = argv[1];
ReferencePtr<ItrObject> itr_ptr;
itr_ptr.assign(ItrObject::RetrieveItrObject(env, itr_handle_ref));
if(NULL==itr_ptr.get())
{
return enif_make_badarg(env);
}
rocksdb::Iterator* itr = itr_ptr->m_Iterator;
rocksdb::Slice key;
if(enif_is_atom(env, action_or_target))
{
if(ATOM_FIRST == action_or_target) itr->SeekToFirst();
if(ATOM_LAST == action_or_target) itr->SeekToLast();
if(ATOM_NEXT == action_or_target) itr->Next();
if(ATOM_PREV == action_or_target) itr->Prev();
}
else if(enif_is_tuple(env, action_or_target))
{
int arity;
const ERL_NIF_TERM* seek;
if(enif_get_tuple(env, action_or_target, &arity, &seek) && 2==arity)
{
if(seek[0] == erocksdb::ATOM_SEEK_FOR_PREV)
{
if(!binary_to_slice(env, seek[1], &key))
return error_einval(env);
itr->SeekForPrev(key);
}
else if(seek[0] == erocksdb::ATOM_SEEK)
{
if(!binary_to_slice(env, seek[1], &key))
return error_einval(env);
itr->Seek(key);
}
else
{
return enif_make_badarg(env);
}
}
else
{
return enif_make_badarg(env);
}
}
else
{
if(!binary_to_slice(env, action_or_target, &key))
return error_einval(env);
itr->Seek(key);
}
if(!itr->Valid())
{
return enif_make_tuple2(env, ATOM_ERROR, ATOM_INVALID_ITERATOR);
}
rocksdb::Status status = itr->status();
if(!status.ok())
{
return error_tuple(env, ATOM_ERROR, status);
}
return enif_make_tuple3(env, ATOM_OK, slice_to_binary(env, itr->key()), slice_to_binary(env, itr->value()));
} // erocksdb::IteratorMove
ERL_NIF_TERM
IteratorRefresh(
ErlNifEnv* env,
int argc,
const ERL_NIF_TERM argv[])
{
const ERL_NIF_TERM& itr_handle_ref = argv[0];
ReferencePtr<ItrObject> itr_ptr;
itr_ptr.assign(ItrObject::RetrieveItrObject(env, itr_handle_ref));
if(NULL==itr_ptr.get())
{
return enif_make_badarg(env);
}
rocksdb::Iterator* itr = itr_ptr->m_Iterator;
itr->Refresh();
return(ATOM_OK);
} // erocksdb::IteratorRefresh
ERL_NIF_TERM
IteratorClose(
ErlNifEnv* env,
int argc,
const ERL_NIF_TERM argv[])
{
ItrObject * itr_ptr;
itr_ptr=ItrObject::RetrieveItrObject(env, argv[0], true);
if (NULL!=itr_ptr)
{
// set closing flag ... atomic likely unnecessary (but safer)
ErlRefObject::InitiateCloseRequest(itr_ptr);
itr_ptr=NULL;
return(ATOM_OK);
} // if
else
{
ERL_NIF_TERM bad_arg=enif_make_badarg(env);
return(bad_arg);
} // else
} // erocksdb:IteratorClose
}