Skip to content
Merged
4 changes: 2 additions & 2 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, true);
} catch (const exception& e) {
LOG_ERROR("Error processing batch: " << 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, true);
for (auto& atom : atoms) {
delete atom;
}
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());
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());
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
36 changes: 25 additions & 11 deletions src/atomdb/AtomDB.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include "AtomDBAPITypes.h"
#include "HandleDecoder.h"
#include "LinkSchema.h"
#include "Merger.h"
#include "Properties.h"

using namespace std;
Expand Down Expand Up @@ -41,19 +42,32 @@ class AtomDB : public HandleDecoder {
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;

virtual vector<string> add_atoms(const vector<atoms::Atom*>& atoms,
bool throw_if_exists = false,
bool is_transactional = false) = 0;
/**
* Add methods take an optional Merger.
* - merger == NULL: upsert (insert if missing, replace if present)
* - merger != NULL and atom exists: merge into a working copy, then persist only if
* merge() returns true; if merge() returns false, stored state is left unchanged
* and the returned handle (or batch slot) is ""
* - Use &ThrowIfExistsMerger::instance() to reject duplicates by throwing
* (former throw_if_exists=true); earlier items in a batch may already be applied
* - Use &SkipIfExistsMerger::instance() to soft-skip duplicates (merge returns false,
* returned handle/slot is "", batch continues)
* - Soft merge() returns false does not abort the rest of a batch
* - Merge-enabled adds assume a single writer per handle (no concurrent RMW atomicity)
*/
virtual string add_atom(const atoms::Atom* atom, const atoms::Merger* merger = NULL) = 0;
virtual string add_node(const atoms::Node* node, const atoms::Merger* merger = NULL) = 0;
virtual string add_link(const atoms::Link* link, const atoms::Merger* merger = NULL) = 0;

virtual vector<string> add_atoms(const vector<atoms::Atom*>& atom_list,
bool is_transactional = false,
const atoms::Merger* merger = NULL) = 0;
virtual vector<string> add_nodes(const vector<atoms::Node*>& nodes,
bool throw_if_exists = false,
bool is_transactional = false) = 0;
bool is_transactional = false,
const atoms::Merger* merger = NULL) = 0;
virtual vector<string> add_links(const vector<atoms::Link*>& links,
bool throw_if_exists = false,
bool is_transactional = false) = 0;
bool is_transactional = false,
const atoms::Merger* merger = NULL) = 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;
Expand Down
1 change: 1 addition & 0 deletions src/atomdb/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ cc_library(
includes = ["."],
deps = [
":atomdb_api_types",
"//commons/atoms:atoms_lib",
],
)

Expand Down
32 changes: 16 additions & 16 deletions src/atomdb/adapterdb/AdapterDB.cc
Original file line number Diff line number Diff line change
Expand Up @@ -137,40 +137,40 @@ set<string> AdapterDB::links_exist(const vector<string>& handles) {
return this->atomdb_backend->links_exist(handles);
}

string AdapterDB::add_atom(const atoms::Atom* atom, bool throw_if_exists) {
string AdapterDB::add_atom(const atoms::Atom* atom, const atoms::Merger* merger) {
this->ensure_backend_ready();
return this->atomdb_backend->add_atom(atom, throw_if_exists);
return this->atomdb_backend->add_atom(atom, merger);
}

string AdapterDB::add_node(const atoms::Node* node, bool throw_if_exists) {
string AdapterDB::add_node(const atoms::Node* node, const atoms::Merger* merger) {
this->ensure_backend_ready();
return this->atomdb_backend->add_node(node, throw_if_exists);
return this->atomdb_backend->add_node(node, merger);
}

string AdapterDB::add_link(const atoms::Link* link, bool throw_if_exists) {
string AdapterDB::add_link(const atoms::Link* link, const atoms::Merger* merger) {
this->ensure_backend_ready();
return this->atomdb_backend->add_link(link, throw_if_exists);
return this->atomdb_backend->add_link(link, merger);
}

vector<string> AdapterDB::add_atoms(const vector<atoms::Atom*>& atoms,
bool throw_if_exists,
bool is_transactional) {
vector<string> AdapterDB::add_atoms(const vector<atoms::Atom*>& atom_list,
bool is_transactional,
const atoms::Merger* merger) {
this->ensure_backend_ready();
return this->atomdb_backend->add_atoms(atoms, throw_if_exists, is_transactional);
return this->atomdb_backend->add_atoms(atom_list, is_transactional, merger);
}

vector<string> AdapterDB::add_nodes(const vector<atoms::Node*>& nodes,
bool throw_if_exists,
bool is_transactional) {
bool is_transactional,
const atoms::Merger* merger) {
this->ensure_backend_ready();
return this->atomdb_backend->add_nodes(nodes, throw_if_exists, is_transactional);
return this->atomdb_backend->add_nodes(nodes, is_transactional, merger);
}

vector<string> AdapterDB::add_links(const vector<atoms::Link*>& links,
bool throw_if_exists,
bool is_transactional) {
bool is_transactional,
const atoms::Merger* merger) {
this->ensure_backend_ready();
return this->atomdb_backend->add_links(links, throw_if_exists, is_transactional);
return this->atomdb_backend->add_links(links, is_transactional, merger);
}

bool AdapterDB::delete_atom(const string& handle, bool delete_link_targets) {
Expand Down
20 changes: 10 additions & 10 deletions src/atomdb/adapterdb/AdapterDB.h
Original file line number Diff line number Diff line change
Expand Up @@ -82,19 +82,19 @@ class AdapterDB : public AtomDB {
set<string> nodes_exist(const vector<string>& handles) override;
set<string> links_exist(const vector<string>& handles) override;

string add_atom(const atoms::Atom* atom, bool throw_if_exists = false) override;
string add_node(const atoms::Node* node, bool throw_if_exists = false) override;
string add_link(const atoms::Link* link, bool throw_if_exists = false) override;
string add_atom(const atoms::Atom* atom, const atoms::Merger* merger = NULL) override;
string add_node(const atoms::Node* node, const atoms::Merger* merger = NULL) override;
string add_link(const atoms::Link* link, const atoms::Merger* merger = NULL) override;

vector<string> add_atoms(const vector<atoms::Atom*>& atoms,
bool throw_if_exists = false,
bool is_transactional = false) override;
vector<string> add_atoms(const vector<atoms::Atom*>& atom_list,
bool is_transactional = false,
const atoms::Merger* merger = NULL) override;
vector<string> add_nodes(const vector<atoms::Node*>& nodes,
bool throw_if_exists = false,
bool is_transactional = false) override;
bool is_transactional = false,
const atoms::Merger* merger = NULL) override;
vector<string> add_links(const vector<atoms::Link*>& links,
bool throw_if_exists = false,
bool is_transactional = false) override;
bool is_transactional = false,
const atoms::Merger* merger = NULL) override;

bool delete_atom(const string& handle, bool delete_link_targets = false) override;
bool delete_node(const string& handle, bool delete_link_targets = false) override;
Expand Down
138 changes: 61 additions & 77 deletions src/atomdb/inmemorydb/InMemoryDB.cc
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include "InMemoryDBAPITypes.h"
#include "Link.h"
#include "LinkSchema.h"
#include "Merger.h"
#include "Node.h"
#include "Utils.h"

Expand Down Expand Up @@ -289,142 +290,125 @@ set<string> InMemoryDB::links_exist(const vector<string>& handles) {
return existing;
}

string InMemoryDB::add_atom(const atoms::Atom* atom, bool throw_if_exists) {
string InMemoryDB::add_atom(const atoms::Atom* atom, const atoms::Merger* merger) {
if (atom->arity() == 0) {
return add_node(dynamic_cast<const atoms::Node*>(atom), throw_if_exists);
return add_node(dynamic_cast<const atoms::Node*>(atom), merger);
} else {
return add_link(dynamic_cast<const atoms::Link*>(atom), throw_if_exists);
return add_link(dynamic_cast<const atoms::Link*>(atom), merger);
}
}

string InMemoryDB::add_node(const atoms::Node* node, bool throw_if_exists) {
string InMemoryDB::add_node(const atoms::Node* node, const atoms::Merger* merger) {
string handle = node->handle();

if (throw_if_exists && this->node_exists(handle)) {
RAISE_ERROR("Node already exists: " + handle);
return "";
}

// Check if already exists
auto existing = atoms_trie_->lookup(handle);
if (existing != NULL && !throw_if_exists) {
return handle; // Already exists, return handle
if ((existing == NULL) || (merger == NULL)) {
// Insert or upsert/replace — HandleTrie insert calls AtomTrieValue::merge,
// which deletes the previous Atom (if any) and takes ownership of the new one.
Node* cloned_node = new Node(*node);
atoms_trie_->insert(handle, new AtomTrieValue(cloned_node));
return handle;
}

// Merge a copy; persist only when merge() returns true.
auto* atom_trie_value = dynamic_cast<AtomTrieValue*>(existing);
unique_ptr<Node> working(new Node(*dynamic_cast<Node*>(atom_trie_value->get_atom())));
if (!merger->merge(working.get(), node)) {
return "";
}

// Clone the node to store in trie
Node* cloned_node = new Node(*node);
auto atom_trie_value = new AtomTrieValue(cloned_node);
atoms_trie_->insert(handle, atom_trie_value);
atoms_trie_->insert(handle, new AtomTrieValue(working.release()));

return handle;
}

string InMemoryDB::add_link(const atoms::Link* link, bool throw_if_exists) {
string InMemoryDB::add_link(const atoms::Link* link, const atoms::Merger* merger) {
vector<Link*> links = {const_cast<atoms::Link*>(link)};
auto handles = this->add_links(links, throw_if_exists, false);
auto handles = this->add_links(links, false, merger);
return handles.empty() ? "" : handles[0];
}

vector<string> InMemoryDB::add_atoms(const vector<atoms::Atom*>& atoms,
bool throw_if_exists,
bool is_transactional) {
if (atoms.empty()) {
vector<string> InMemoryDB::add_atoms(const vector<atoms::Atom*>& atom_list,
bool is_transactional,
const atoms::Merger* merger) {
if (atom_list.empty()) {
return {};
}

vector<Node*> nodes;
vector<Link*> links;
for (const auto& atom : atoms) {
for (const auto& atom : atom_list) {
LOG_DEBUG("Adding atom: " + atom->to_string());
if (atom->arity() == 0) {
nodes.push_back(dynamic_cast<atoms::Node*>(atom));
} else {
links.push_back(dynamic_cast<atoms::Link*>(atom));
}
}
auto node_handles = this->add_nodes(nodes, throw_if_exists, is_transactional);
auto link_handles = this->add_links(links, throw_if_exists, is_transactional);
auto node_handles = this->add_nodes(nodes, is_transactional, merger);
auto link_handles = this->add_links(links, is_transactional, merger);

node_handles.insert(node_handles.end(), link_handles.begin(), link_handles.end());
return node_handles;
}

vector<string> InMemoryDB::add_nodes(const vector<atoms::Node*>& nodes,
bool throw_if_exists,
bool is_transactional) {
bool is_transactional,
const atoms::Merger* merger) {
if (nodes.empty()) {
return {};
}

vector<string> handles;
handles.reserve(nodes.size());
for (const auto& node : nodes) {
handles.push_back(node->handle());
handles.push_back(this->add_node(node, merger));
}

if (throw_if_exists) {
auto existing_handles = this->nodes_exist(handles);
if (!existing_handles.empty()) {
vector<string> existing_handles_vector(existing_handles.begin(), existing_handles.end());
RAISE_ERROR("Failed to insert nodes, some nodes already exist: " +
Utils::join(existing_handles_vector, ','));
return {};
}
}

for (const auto& node : nodes) {
handles.push_back(this->add_node(node, throw_if_exists));
}

return handles;
}

vector<string> InMemoryDB::add_links(const vector<atoms::Link*>& links,
bool throw_if_exists,
bool is_transactional) {
bool is_transactional,
const atoms::Merger* merger) {
if (links.empty()) {
return {};
}

if (throw_if_exists) {
vector<string> handles;
for (const auto& link : links) {
handles.push_back(link->handle());
}
auto existing_handles = this->links_exist(handles);
if (!existing_handles.empty()) {
vector<string> existing_handles_vector(existing_handles.begin(), existing_handles.end());
RAISE_ERROR("Failed to insert links, some links already exist: " +
Utils::join(existing_handles_vector, ','));
return {};
}
}

vector<string> handles;
handles.reserve(links.size());

for (const auto& link : links) {
string link_handle = link->handle();
handles.push_back(link_handle);

// Check if already exists
auto existing = atoms_trie_->lookup(link_handle);
if (existing == NULL || !throw_if_exists) {
if (existing == NULL) {
// Clone the link to store in trie
Link* cloned_link = new Link(*link);
auto atom_trie_value = new AtomTrieValue(cloned_link);
atoms_trie_->insert(link_handle, atom_trie_value);
if ((existing == NULL) || (merger == NULL)) {
// Insert or upsert/replace — AtomTrieValue::merge frees the previous Atom.
Link* cloned_link = new Link(*link);
atoms_trie_->insert(link_handle, new AtomTrieValue(cloned_link));
} else {
// Merge a copy; persist only when merge() returns true.
// On failure, skip incoming-set/pattern updates — the link was already
// indexed by whichever add created it.
auto* atom_trie_value = dynamic_cast<AtomTrieValue*>(existing);
unique_ptr<Link> working(new Link(*dynamic_cast<Link*>(atom_trie_value->get_atom())));
if (!merger->merge(working.get(), link)) {
handles.push_back("");
continue;
}
atoms_trie_->insert(link_handle, new AtomTrieValue(working.release()));
}

// Update incoming sets for each target
for (const auto& target_handle : link->targets) {
this->add_incoming_set(target_handle, link_handle);
}
// Update incoming sets for each target
for (const auto& target_handle : link->targets) {
this->add_incoming_set(target_handle, link_handle);
}

// Index pattern
auto pattern_handles = this->match_pattern_index_schema(link);
for (const auto& pattern_handle : pattern_handles) {
this->add_pattern(pattern_handle, link_handle);
}
// Index pattern
auto pattern_handles = this->match_pattern_index_schema(link);
for (const auto& pattern_handle : pattern_handles) {
this->add_pattern(pattern_handle, link_handle);
}

handles.push_back(link_handle);
}

return handles;
Expand Down
Loading
Loading