Skip to content
Closed
Show file tree
Hide file tree
Changes from all 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
Expand Up @@ -19,6 +19,7 @@
"endpoint": "localhost:40021",
"username": "admin",
"password": "admin",
"seed_protected": false,
"cluster": false,
"cluster_secret_key": "8UDJSgpUCaVOTQG",
"nodes": [
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