Skip to content
Closed
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions config/das.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
{
"atomdb": {
"type": "redismongodb",
"protected": true,
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
"composite_type_enabled": true,
"redis": {
"image": "redis:7.2.3-alpine",
Expand Down
6 changes: 3 additions & 3 deletions src/agents/atomdb_broker/AtomDBProxy.cc
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ void AtomDBProxy::add_atoms_callback(const vector<string>& tokens) {
for (auto& atom : atoms) {
buffer.push_back(atom.get());
}
this->atomdb->add_atoms(buffer, false, true);
this->atomdb->add_atoms(buffer, "", false, true);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
} catch (const exception& e) {
LOG_ERROR("Error processing batch: " << e.what());
}
Expand All @@ -180,7 +180,7 @@ void AtomDBProxy::delete_atoms_callback(const vector<string>& args) {
}
vector<string> handles(args.begin(), args.end() - 1);
bool delete_link_targets = args.back() == "1";
uint deleted_count = this->atomdb->delete_atoms(handles, delete_link_targets);
uint deleted_count = this->atomdb->delete_atoms(handles, "", delete_link_targets);
LOG_INFO("Deleted " << deleted_count << " atoms");
} catch (const exception& e) {
LOG_ERROR("Error processing delete_atoms command: " << e.what());
Expand Down Expand Up @@ -215,7 +215,7 @@ void AtomDBProxy::process_atom_batches() {
this->pending_atoms_count -= atoms.size();
lock.unlock();
auto job = [this, atoms = std::move(atoms)]() {
this->atomdb->add_atoms(atoms, false, true);
this->atomdb->add_atoms(atoms, "", false, true);
for (auto& atom : atoms) {
delete atom;
}
Expand Down
12 changes: 6 additions & 6 deletions src/agents/evolution/QueryEvolutionProcessor.cc
Original file line number Diff line number Diff line change
Expand Up @@ -497,12 +497,12 @@ string QueryEvolutionProcessor::answer_to_string_2(shared_ptr<QueryAnswer> answe
vector<string> path_link = {" -> ", " -> "};
bool first = true;
for (string& handle : answer->get_path_vector(i)) {
auto link = db->get_link(handle);
auto link = db->get_link(handle, "");
if ((link == nullptr) || (link->arity() != 3)) {
return "Invalid link: " + handle;
}
auto target1 = db->get_link(link->targets[1]);
auto target2 = db->get_link(link->targets[2]);
auto target1 = db->get_link(link->targets[1], "");
auto target2 = db->get_link(link->targets[2], "");
if ((target1 == nullptr) || (target2 == nullptr)) {
return "Invalid link: " + link->to_string();
}
Expand Down Expand Up @@ -533,9 +533,9 @@ string QueryEvolutionProcessor::answer_to_string_1(shared_ptr<QueryAnswer> answe
string path_link = " -> ";
bool first = true;
for (string& handle : answer->get_path_vector(0)) {
auto link = db->get_link(handle);
auto target1 = db->get_link(link->targets[1]);
auto target2 = db->get_link(link->targets[2]);
auto link = db->get_link(handle, "");
auto target1 = db->get_link(link->targets[1], "");
auto target2 = db->get_link(link->targets[2], "");
if (first) {
first = false;
path = target1->metta_representation(*(this->decoder)) + path_link;
Expand Down
4 changes: 2 additions & 2 deletions src/agents/evolution/fitness_functions/CountLetterFunction.cc
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,9 @@ float CountLetterFunction::eval(shared_ptr<QueryAnswer> query_answer) {
shared_ptr<Link> sentence_link;
shared_ptr<Node> sentence_name_node;
string handle = query_answer->assignment.get(VARIABLE_NAME);
sentence_link = this->db->get_link(handle);
sentence_link = this->db->get_link(handle, "");
handle = sentence_link->targets[1];
sentence_name_node = this->db->get_node(handle);
sentence_name_node = this->db->get_node(handle, "");
string sentence_name = sentence_name_node->name;
unsigned int count = 0;
unsigned int sentence_length = 0;
Expand Down
6 changes: 3 additions & 3 deletions src/agents/link_creation_agent/EquivalenceProcessor.cc
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,8 @@ bool EquivalenceProcessor::link_exists(const string& handle1, const string& hand
vector<string> targets_c2_c1 = {equivalence_node.handle(), handle2, handle1};
shared_ptr<Link> link_c1_c2 = make_shared<Link>("Expression", targets_c1_c2);
shared_ptr<Link> link_c2_c1 = make_shared<Link>("Expression", targets_c2_c1);
return AtomDBSingleton::get_instance()->link_exists(link_c1_c2->handle()) &&
AtomDBSingleton::get_instance()->link_exists(link_c2_c1->handle());
return AtomDBSingleton::get_instance()->link_exists(link_c1_c2->handle(), "") &&
AtomDBSingleton::get_instance()->link_exists(link_c2_c1->handle(), "");
}

static vector<string> build_equivalence_query(const string& handle) {
Expand Down Expand Up @@ -89,7 +89,7 @@ vector<shared_ptr<Link>> EquivalenceProcessor::process_query(shared_ptr<QueryAns
vector<shared_ptr<Link>> result;
Node equivalence_node("Symbol", "Equivalence");
try {
AtomDBSingleton::get_instance()->add_node(&equivalence_node);
AtomDBSingleton::get_instance()->add_node(&equivalence_node, "");
} catch (const std::exception& e) {
LOG_ERROR("Failed to add node to AtomDB: " << e.what());
}
Expand Down
6 changes: 3 additions & 3 deletions src/agents/link_creation_agent/ImplicationProcessor.cc
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,8 @@ bool ImplicationProcessor::link_exists(const string& handle1, const string& hand
vector<string> targets_p2_p1 = {implication_node.handle(), handle2, handle1};
shared_ptr<Link> p1_link = make_shared<Link>("Expression", targets_p1_p2);
shared_ptr<Link> p2_link = make_shared<Link>("Expression", targets_p2_p1);
return AtomDBSingleton::get_instance()->link_exists(p1_link->handle()) &&
AtomDBSingleton::get_instance()->link_exists(p2_link->handle());
return AtomDBSingleton::get_instance()->link_exists(p1_link->handle(), "") &&
AtomDBSingleton::get_instance()->link_exists(p2_link->handle(), "");
}

vector<shared_ptr<Link>> ImplicationProcessor::process_query(shared_ptr<QueryAnswer> query_answer,
Expand Down Expand Up @@ -113,7 +113,7 @@ vector<shared_ptr<Link>> ImplicationProcessor::process_query(shared_ptr<QueryAns
vector<shared_ptr<Link>> result;
Node implication_node("Symbol", "Implication");
try {
AtomDBSingleton::get_instance()->add_node(&implication_node);
AtomDBSingleton::get_instance()->add_node(&implication_node, "");
} catch (const std::exception& e) {
LOG_ERROR("Failed to add node to AtomDB: " << e.what());
}
Expand Down
8 changes: 4 additions & 4 deletions src/agents/link_creation_agent/LinkCreationService.cc
Original file line number Diff line number Diff line change
Expand Up @@ -126,14 +126,14 @@ void LinkCreationService::set_timeout(int timeout) { this->timeout = timeout; }

static void add_or_update_link(shared_ptr<Link> link) {
auto db_instance = AtomDBSingleton::get_instance();
if (!db_instance->link_exists(link->handle())) {
if (!db_instance->link_exists(link->handle(), "")) {
LOG_INFO("Adding link to AtomDB: " << link->to_string());
db_instance->add_link(link.get());
db_instance->add_link(link.get(), "");
} else {
LOG_INFO("Updating link in AtomDB: " << link->to_string());
auto old_link = db_instance->get_atom(link->handle());
db_instance->delete_link(link->handle(), false);
db_instance->add_link(link.get());
db_instance->delete_link(link->handle(), "", false);
db_instance->add_link(link.get(), "");
}
}

Expand Down
4 changes: 2 additions & 2 deletions src/agents/link_creation_agent/MettaTemplateProcessor.cc
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ static void create_missing_atoms_in_atomdb(shared_ptr<MettaParserActions> parser
for (const auto& element : parser_actions->handle_to_atom) {
if (dynamic_pointer_cast<Node>(element.second) != nullptr) {
try {
atomdb->add_node(dynamic_pointer_cast<Node>(element.second).get(), false);
atomdb->add_node(dynamic_pointer_cast<Node>(element.second).get(), "", false);
LOG_DEBUG("Node added to AtomDB: " << element.second->to_string());
} catch (const std::exception& e) {
LOG_ERROR("Error adding node to AtomDB: " << e.what());
Expand Down Expand Up @@ -68,7 +68,7 @@ static void create_missing_atoms_in_atomdb(shared_ptr<MettaParserActions> parser
RAISE_ERROR("Parsed atom is not a Link for metta expression: " + metta_expression_cp);
continue;
}
atomdb->add_link(dynamic_pointer_cast<Link>(link).get(), false);
atomdb->add_link(dynamic_pointer_cast<Link>(link).get(), "", false);
LOG_DEBUG("Link added to AtomDB: " << metta_expression_cp);
} catch (const std::exception& e) {
LOG_ERROR("Error adding link to AtomDB: " << e.what());
Expand Down
4 changes: 2 additions & 2 deletions src/agents/query_engine/query_element/LinkTemplate.cc
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ void LinkTemplate::processor_method(shared_ptr<StoppableThread> monitor) {
handles = LinkTemplate::fetched_links_cache().get(link_schema_handle);
} else {
LOG_INFO("Fetching " + link_schema_handle + " from AtomDB");
handles = db->query_for_pattern(this->link_schema);
handles = db->query_for_pattern(this->link_schema, "");
if (this->use_cache) {
LinkTemplate::fetched_links_cache().set(link_schema_handle, handles);
}
Expand Down Expand Up @@ -220,7 +220,7 @@ void LinkTemplate::processor_method(shared_ptr<StoppableThread> monitor) {
pending = 0;
} else {
if (tagged_handle.second > 0 || !this->positive_importance_flag) {
if (db->allow_nested_indexing()) {
if (db->allow_nested_indexing("")) {
if ((this->attention_focus_strictness == 0.0) ||
(this->attention_focus_strictness == 1.0)) {
this->source_element->add_handle(
Expand Down
111 changes: 74 additions & 37 deletions src/atomdb/AtomDB.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,56 +20,93 @@ class AtomDB : public HandleDecoder {
AtomDB() = default;
virtual ~AtomDB() = default;

virtual bool allow_nested_indexing() = 0;
virtual bool allow_nested_indexing(const string& public_key) = 0;
virtual bool composite_type_enabled() const = 0;

virtual shared_ptr<Atom> get_atom(const string& handle) = 0; // HandleDecoder interface
virtual shared_ptr<Node> get_node(const string& handle) = 0;
virtual shared_ptr<Link> get_link(const string& handle) = 0;

virtual vector<shared_ptr<Atom>> get_matching_atoms(bool is_toplevel, Atom& key) = 0;

virtual shared_ptr<atomdb_api_types::HandleSet> query_for_pattern(const LinkSchema& link_schema) = 0;
virtual shared_ptr<atomdb_api_types::HandleList> query_for_targets(const string& handle) = 0;
virtual shared_ptr<atomdb_api_types::HandleSet> query_for_incoming_set(const string& handle) = 0;

virtual bool atom_exists(const string& handle) = 0;
virtual bool node_exists(const string& handle) = 0;
virtual bool link_exists(const string& handle) = 0;

virtual set<string> atoms_exist(const vector<string>& handles) = 0;
virtual set<string> nodes_exist(const vector<string>& handles) = 0;
virtual set<string> links_exist(const vector<string>& handles) = 0;

virtual string add_atom(const atoms::Atom* atom, bool throw_if_exists = false) = 0;
virtual string add_node(const atoms::Node* node, bool throw_if_exists = false) = 0;
virtual string add_link(const atoms::Link* link, bool throw_if_exists = false) = 0;
/**
* @brief Reports whether this backend points to a protected database.
*/
virtual bool is_protected() const = 0;

/**
* HandleDecoder requires get_atom(handle) without public_key. Existing callers use that interface,
* so this forwards to get_atom(handle, "").
*/
shared_ptr<Atom> get_atom(const string& handle) override { return get_atom(handle, ""); }

virtual shared_ptr<Atom> get_atom(const string& handle, const string& public_key) = 0;
virtual shared_ptr<Node> get_node(const string& handle, const string& public_key) = 0;
virtual shared_ptr<Link> get_link(const string& handle, const string& public_key) = 0;

virtual vector<shared_ptr<Atom>> get_matching_atoms(bool is_toplevel,
Atom& key,
const string& public_key) = 0;

virtual shared_ptr<atomdb_api_types::HandleSet> query_for_pattern(const LinkSchema& link_schema,
const string& public_key) = 0;
virtual shared_ptr<atomdb_api_types::HandleList> query_for_targets(const string& handle,
const string& public_key) = 0;
virtual shared_ptr<atomdb_api_types::HandleSet> query_for_incoming_set(const string& handle,
const string& public_key) = 0;

virtual bool atom_exists(const string& handle, const string& public_key) = 0;
virtual bool node_exists(const string& handle, const string& public_key) = 0;
virtual bool link_exists(const string& handle, const string& public_key) = 0;

virtual set<string> atoms_exist(const vector<string>& handles, const string& public_key) = 0;
virtual set<string> nodes_exist(const vector<string>& handles, const string& public_key) = 0;
virtual set<string> links_exist(const vector<string>& handles, const string& public_key) = 0;

virtual string add_atom(const atoms::Atom* atom,
const string& public_key,
bool throw_if_exists = false) = 0;
virtual string add_node(const atoms::Node* node,
const string& public_key,
bool throw_if_exists = false) = 0;
virtual string add_link(const atoms::Link* link,
const string& public_key,
bool throw_if_exists = false) = 0;

virtual vector<string> add_atoms(const vector<atoms::Atom*>& atoms,
const string& public_key,
bool throw_if_exists = false,
bool is_transactional = false) = 0;
virtual vector<string> add_nodes(const vector<atoms::Node*>& nodes,
const string& public_key,
bool throw_if_exists = false,
bool is_transactional = false) = 0;
virtual vector<string> add_links(const vector<atoms::Link*>& links,
const string& public_key,
bool throw_if_exists = false,
bool is_transactional = false) = 0;

virtual bool delete_atom(const string& handle, bool delete_link_targets = false) = 0;
virtual bool delete_node(const string& handle, bool delete_link_targets = false) = 0;
virtual bool delete_link(const string& handle, bool delete_link_targets = false) = 0;

virtual uint delete_atoms(const vector<string>& handles, bool delete_link_targets = false) = 0;
virtual uint delete_nodes(const vector<string>& handles, bool delete_link_targets = false) = 0;
virtual uint delete_links(const vector<string>& handles, bool delete_link_targets = false) = 0;

virtual void re_index_patterns(bool flush_patterns = true) = 0;

virtual size_t node_count() const = 0;
virtual size_t link_count() const = 0;
virtual size_t atom_count() const = 0;

bool empty() const { return atom_count() == 0; }
virtual bool delete_atom(const string& handle,
const string& public_key,
bool delete_link_targets = false) = 0;
virtual bool delete_node(const string& handle,
const string& public_key,
bool delete_link_targets = false) = 0;
virtual bool delete_link(const string& handle,
const string& public_key,
bool delete_link_targets = false) = 0;

virtual uint delete_atoms(const vector<string>& handles,
const string& public_key,
bool delete_link_targets = false) = 0;
virtual uint delete_nodes(const vector<string>& handles,
const string& public_key,
bool delete_link_targets = false) = 0;
virtual uint delete_links(const vector<string>& handles,
const string& public_key,
bool delete_link_targets = false) = 0;

virtual void re_index_patterns(const string& public_key, bool flush_patterns = true) = 0;

virtual size_t node_count(const string& public_key) const = 0;
virtual size_t link_count(const string& public_key) const = 0;
virtual size_t atom_count(const string& public_key) const = 0;

bool empty(const string& public_key) const { return atom_count(public_key) == 0; }
};

} // namespace atomdb
41 changes: 25 additions & 16 deletions src/atomdb/AtomDBSingleton.cc
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

#include "AdapterDB.h"
#include "MorkDB.h"
#include "ProtectedAtomDB.h"
#include "RedisMongoDB.h"
#include "RemoteAtomDB.h"
#include "Utils.h"
Expand All @@ -19,24 +20,32 @@ void AtomDBSingleton::init(const JsonConfig& atomdb_config) {
if (AtomDBSingleton::initialized) {
RAISE_ERROR(
"AtomDBSingleton already initialized. AtomDBSingleton::init() should be called only once.");
}

shared_ptr<AtomDB> atomdb;
auto atomdb_type = atomdb_config.at_path("type").get_or<string>("");

if (atomdb_type == "morkdb") {
atomdb = shared_ptr<AtomDB>(new MorkDB("", atomdb_config));
} else if (atomdb_type == "redismongodb") {
atomdb = shared_ptr<AtomDB>(new RedisMongoDB("", false, atomdb_config));
} else if (atomdb_type == "remotedb") {
auto remote_peers_config =
atomdb_config.at_path("remote_peers").get_or<JsonConfig>(JsonConfig());
atomdb = shared_ptr<AtomDB>(new RemoteAtomDB(remote_peers_config));
} else if (atomdb_type == "adapterdb") {
atomdb = shared_ptr<AtomDB>(new AdapterDB(atomdb_config));
} else {
RAISE_ERROR("Invalid AtomDB type: " + atomdb_type);
}

if (atomdb->is_protected()) {
AtomDBSingleton::atom_db = shared_ptr<AtomDB>(new ProtectedAtomDB(atomdb, atomdb_config));
} else {
auto atomdb_type = atomdb_config.at_path("type").get_or<string>("");
if (atomdb_type == "morkdb") {
AtomDBSingleton::atom_db = shared_ptr<AtomDB>(new MorkDB("", atomdb_config));
} else if (atomdb_type == "redismongodb") {
AtomDBSingleton::atom_db = shared_ptr<AtomDB>(new RedisMongoDB("", false, atomdb_config));
} else if (atomdb_type == "remotedb") {
auto remote_peers_config =
atomdb_config.at_path("remote_peers").get_or<JsonConfig>(JsonConfig());
AtomDBSingleton::atom_db = shared_ptr<AtomDB>(new RemoteAtomDB(remote_peers_config));
} else if (atomdb_type == "adapterdb") {
AtomDBSingleton::atom_db = shared_ptr<AtomDB>(new AdapterDB(atomdb_config));
} else {
RAISE_ERROR("Invalid AtomDB type: " + atomdb_type);
}

AtomDBSingleton::initialized = true;
AtomDBSingleton::atom_db = atomdb;
}

AtomDBSingleton::initialized = true;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

shared_ptr<AtomDB> AtomDBSingleton::get_instance() {
Expand Down
2 changes: 2 additions & 0 deletions src/atomdb/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ cc_library(
":atomdb_singleton",
":atomdbutils",
"//atomdb/adapterdb:adapterdb_lib",
"//atomdb/auth:protected_atomdb_lib",
"//atomdb/inmemorydb:inmemorydb_lib",
"//atomdb/morkdb:morkdb_lib",
"//atomdb/redis_mongodb:redis_mongodb_lib",
Expand Down Expand Up @@ -57,6 +58,7 @@ cc_library(
deps = [
"//atomdb:atomdb_api_types",
"//atomdb/adapterdb",
"//atomdb/auth:protected_atomdb_lib",
"//atomdb/morkdb",
"//atomdb/redis_mongodb",
"//atomdb/remotedb:remotedb_lib",
Expand Down
Loading
Loading