|
|
@@ -1,5 +1,6 @@
|
|
|
#include "database_grpc_impl.hpp"
|
|
|
#include "persistence/wal.hpp"
|
|
|
+#include "views/projection.hpp"
|
|
|
|
|
|
#include <spdlog/spdlog.h>
|
|
|
#include <nlohmann/json.hpp>
|
|
|
@@ -16,12 +17,14 @@ DatabaseGrpcImpl::DatabaseGrpcImpl(
|
|
|
PersistenceManager& persistence,
|
|
|
EventManager& events,
|
|
|
FileManager& files,
|
|
|
- EncryptionManager& encryption
|
|
|
+ EncryptionManager& encryption,
|
|
|
+ ViewManager& view_manager
|
|
|
) : store_(store)
|
|
|
, persistence_(persistence)
|
|
|
, events_(events)
|
|
|
, files_(files)
|
|
|
, encryption_(encryption)
|
|
|
+ , view_manager_(view_manager)
|
|
|
{
|
|
|
}
|
|
|
|
|
|
@@ -32,6 +35,10 @@ grpc::Status DatabaseGrpcImpl::Insert(
|
|
|
const pb::InsertRequest* request,
|
|
|
pb::InsertResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INVALID_ARGUMENT,
|
|
|
+ "cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ }
|
|
|
try {
|
|
|
nlohmann::json data = nlohmann::json::parse(
|
|
|
request->data().begin(), request->data().end()
|
|
|
@@ -82,7 +89,15 @@ grpc::Status DatabaseGrpcImpl::Get(
|
|
|
const pb::GetRequest* request,
|
|
|
pb::GetResponse* response
|
|
|
) {
|
|
|
- auto doc = store_.get(request->collection(), request->id());
|
|
|
+ // Resolve view to real collection (if applicable)
|
|
|
+ std::string targetCollection = request->collection();
|
|
|
+ std::optional<ViewInfo> view;
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ view = view_manager_.getView(request->collection());
|
|
|
+ targetCollection = view->collection;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto doc = store_.get(targetCollection, request->id());
|
|
|
if (!doc) {
|
|
|
response->set_found(false);
|
|
|
return grpc::Status::OK;
|
|
|
@@ -91,7 +106,26 @@ grpc::Status DatabaseGrpcImpl::Get(
|
|
|
// Decrypt sensitive fields
|
|
|
encryption_.decryptSensitiveFields(*doc);
|
|
|
|
|
|
- *response->mutable_document() = toProto(*doc);
|
|
|
+ // Apply view's where filters — if doc doesn't match, treat as not found
|
|
|
+ if (view && !view->where.empty()) {
|
|
|
+ if (!store_.matchesFilters(*doc, view->where)) {
|
|
|
+ response->set_found(false);
|
|
|
+ return grpc::Status::OK;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ pb::Document protoDoc = toProto(*doc);
|
|
|
+
|
|
|
+ // Apply view projection on the response data
|
|
|
+ if (view) {
|
|
|
+ nlohmann::json docJson = nlohmann::json::parse(protoDoc.data());
|
|
|
+ docJson = applyProjection(docJson, view->include, view->exclude);
|
|
|
+ protoDoc.set_data(docJson.dump());
|
|
|
+ // Preserve view name as the observed collection
|
|
|
+ protoDoc.set_collection(request->collection());
|
|
|
+ }
|
|
|
+
|
|
|
+ *response->mutable_document() = std::move(protoDoc);
|
|
|
response->set_found(true);
|
|
|
|
|
|
return grpc::Status::OK;
|
|
|
@@ -102,6 +136,11 @@ grpc::Status DatabaseGrpcImpl::Update(
|
|
|
const pb::UpdateRequest* request,
|
|
|
pb::UpdateResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ response->set_success(false);
|
|
|
+ response->set_error("cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ return grpc::Status::OK;
|
|
|
+ }
|
|
|
try {
|
|
|
nlohmann::json data = nlohmann::json::parse(
|
|
|
request->data().begin(), request->data().end()
|
|
|
@@ -171,6 +210,11 @@ grpc::Status DatabaseGrpcImpl::PatchDocument(
|
|
|
const pb::PatchDocumentRequest* request,
|
|
|
pb::PatchDocumentResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ response->set_success(false);
|
|
|
+ response->set_error("cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ return grpc::Status::OK;
|
|
|
+ }
|
|
|
try {
|
|
|
nlohmann::json patch = nlohmann::json::parse(
|
|
|
request->patch_json().begin(), request->patch_json().end()
|
|
|
@@ -216,6 +260,10 @@ grpc::Status DatabaseGrpcImpl::Upsert(
|
|
|
const pb::UpsertRequest* request,
|
|
|
pb::UpsertResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INVALID_ARGUMENT,
|
|
|
+ "cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ }
|
|
|
try {
|
|
|
nlohmann::json data = nlohmann::json::parse(
|
|
|
request->data().begin(), request->data().end()
|
|
|
@@ -283,6 +331,10 @@ grpc::Status DatabaseGrpcImpl::Delete(
|
|
|
const pb::DeleteRequest* request,
|
|
|
pb::DeleteResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INVALID_ARGUMENT,
|
|
|
+ "cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ }
|
|
|
// Check if document has a vector before deleting (remove() erases it from memory)
|
|
|
bool hadVector = false;
|
|
|
auto* vectors = store_.getCollectionVectors(request->collection());
|
|
|
@@ -399,6 +451,10 @@ grpc::Status DatabaseGrpcImpl::BatchInsert(
|
|
|
const pb::BatchInsertRequest* request,
|
|
|
pb::BatchInsertResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INVALID_ARGUMENT,
|
|
|
+ "cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ }
|
|
|
uint32_t successCount = 0;
|
|
|
uint32_t errorCount = 0;
|
|
|
std::string actor = request->actor();
|
|
|
@@ -460,6 +516,10 @@ grpc::Status DatabaseGrpcImpl::BatchDelete(
|
|
|
const pb::BatchDeleteRequest* request,
|
|
|
pb::BatchDeleteResponse* response
|
|
|
) {
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INVALID_ARGUMENT,
|
|
|
+ "cannot write to view '" + request->collection() + "': views are read-only");
|
|
|
+ }
|
|
|
std::vector<std::string> ids(request->ids().begin(), request->ids().end());
|
|
|
uint64_t deleted = store_.bulkDelete(request->collection(), ids);
|
|
|
response->set_deleted_count(deleted);
|
|
|
@@ -473,8 +533,29 @@ grpc::Status DatabaseGrpcImpl::Find(
|
|
|
const pb::FindRequest* request,
|
|
|
pb::FindResponse* response
|
|
|
) {
|
|
|
+ // Resolve view to real collection (if applicable)
|
|
|
+ std::string targetCollection = request->collection();
|
|
|
+ std::optional<ViewInfo> view;
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ view = view_manager_.getView(request->collection());
|
|
|
+ targetCollection = view->collection;
|
|
|
+ }
|
|
|
+
|
|
|
Query query = fromProtoQuery(*request);
|
|
|
- QueryResult result = store_.find(request->collection(), query);
|
|
|
+
|
|
|
+ if (view) {
|
|
|
+ // AND-merge: append view's where filters to the caller's filters
|
|
|
+ for (const auto& f : view->where) {
|
|
|
+ query.filters.push_back(f);
|
|
|
+ }
|
|
|
+
|
|
|
+ // Caller wins for sort: only apply view's default_sort if caller didn't specify
|
|
|
+ if (!query.sort && view->defaultSort) {
|
|
|
+ query.sort = view->defaultSort;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ QueryResult result = store_.find(targetCollection, query);
|
|
|
|
|
|
response->set_total_count(result.totalCount);
|
|
|
response->set_has_more(result.hasMore);
|
|
|
@@ -488,7 +569,15 @@ grpc::Status DatabaseGrpcImpl::Find(
|
|
|
|
|
|
for (auto& doc : result.documents) {
|
|
|
encryption_.decryptSensitiveFields(doc);
|
|
|
- *response->add_documents() = toProto(doc);
|
|
|
+ pb::Document protoDoc = toProto(doc);
|
|
|
+ if (view) {
|
|
|
+ nlohmann::json docJson = nlohmann::json::parse(protoDoc.data());
|
|
|
+ docJson = applyProjection(docJson, view->include, view->exclude);
|
|
|
+ protoDoc.set_data(docJson.dump());
|
|
|
+ // Preserve view name as the observed collection
|
|
|
+ protoDoc.set_collection(request->collection());
|
|
|
+ }
|
|
|
+ *response->add_documents() = std::move(protoDoc);
|
|
|
}
|
|
|
|
|
|
return grpc::Status::OK;
|
|
|
@@ -499,8 +588,24 @@ grpc::Status DatabaseGrpcImpl::Count(
|
|
|
const pb::CountRequest* request,
|
|
|
pb::CountResponse* response
|
|
|
) {
|
|
|
- auto filters = fromProtoFilters(request->filters());
|
|
|
- uint64_t count = store_.count(request->collection(), filters);
|
|
|
+ // Resolve view to real collection (if applicable)
|
|
|
+ std::string targetCollection = request->collection();
|
|
|
+ std::vector<Filter> filters;
|
|
|
+ if (view_manager_.isView(request->collection())) {
|
|
|
+ auto view = view_manager_.getView(request->collection());
|
|
|
+ targetCollection = view->collection;
|
|
|
+ // AND-merge view's where into the count filters
|
|
|
+ for (const auto& f : view->where) {
|
|
|
+ filters.push_back(f);
|
|
|
+ }
|
|
|
+ }
|
|
|
+ // Append caller's filters (AND-merge; order doesn't matter for AND)
|
|
|
+ auto callerFilters = fromProtoFilters(request->filters());
|
|
|
+ for (auto& f : callerFilters) {
|
|
|
+ filters.push_back(std::move(f));
|
|
|
+ }
|
|
|
+
|
|
|
+ uint64_t count = store_.count(targetCollection, filters);
|
|
|
response->set_count(count);
|
|
|
return grpc::Status::OK;
|
|
|
}
|
|
|
@@ -1029,6 +1134,130 @@ Query DatabaseGrpcImpl::fromProtoQuery(const pb::FindRequest& request) {
|
|
|
return query;
|
|
|
}
|
|
|
|
|
|
+// ===== View Operations =====
|
|
|
+
|
|
|
+grpc::Status DatabaseGrpcImpl::CreateView(
|
|
|
+ grpc::ServerContext* /*context*/,
|
|
|
+ const pb::CreateViewRequest* request,
|
|
|
+ pb::CreateViewResponse* response
|
|
|
+) {
|
|
|
+ try {
|
|
|
+ ViewInfo v;
|
|
|
+ v.name = request->name();
|
|
|
+ v.collection = request->collection();
|
|
|
+ for (const auto& p : request->include()) v.include.push_back(p);
|
|
|
+ for (const auto& p : request->exclude()) v.exclude.push_back(p);
|
|
|
+
|
|
|
+ // Filters
|
|
|
+ for (const auto& pf : request->where()) {
|
|
|
+ Filter f;
|
|
|
+ f.field = pf.field();
|
|
|
+ f.op = static_cast<FilterOp>(pf.op() - 1); // Adjust for UNSPECIFIED
|
|
|
+ if (!pf.value().empty()) {
|
|
|
+ try {
|
|
|
+ f.value = nlohmann::json::parse(pf.value());
|
|
|
+ } catch (...) {
|
|
|
+ f.value = pf.value(); // fall back to string if not JSON
|
|
|
+ }
|
|
|
+ }
|
|
|
+ v.where.push_back(f);
|
|
|
+ }
|
|
|
+
|
|
|
+ // Default sort
|
|
|
+ if (request->has_default_sort()) {
|
|
|
+ Sort s;
|
|
|
+ s.field = request->default_sort().field();
|
|
|
+ s.descending = request->default_sort().descending();
|
|
|
+ v.defaultSort = s;
|
|
|
+ }
|
|
|
+
|
|
|
+ std::string err;
|
|
|
+ bool ok = view_manager_.createView(v, err);
|
|
|
+ response->set_success(ok);
|
|
|
+ if (!ok) response->set_error(err);
|
|
|
+ return grpc::Status::OK;
|
|
|
+ } catch (const std::exception& e) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INTERNAL, e.what());
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+grpc::Status DatabaseGrpcImpl::DropView(
|
|
|
+ grpc::ServerContext* /*context*/,
|
|
|
+ const pb::DropViewRequest* request,
|
|
|
+ pb::DropViewResponse* response
|
|
|
+) {
|
|
|
+ try {
|
|
|
+ std::string err;
|
|
|
+ bool ok = view_manager_.dropView(request->name(), err);
|
|
|
+ response->set_success(ok);
|
|
|
+ if (!ok) response->set_error(err);
|
|
|
+ return grpc::Status::OK;
|
|
|
+ } catch (const std::exception& e) {
|
|
|
+ return grpc::Status(grpc::StatusCode::INTERNAL, e.what());
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+grpc::Status DatabaseGrpcImpl::ListViews(
|
|
|
+ grpc::ServerContext* /*context*/,
|
|
|
+ const pb::ListViewsRequest* /*request*/,
|
|
|
+ pb::ListViewsResponse* response
|
|
|
+) {
|
|
|
+ auto views = view_manager_.listViews();
|
|
|
+ for (const auto& v : views) {
|
|
|
+ auto* out = response->add_views();
|
|
|
+ out->set_name(v.name);
|
|
|
+ out->set_collection(v.collection);
|
|
|
+ for (const auto& p : v.include) out->add_include(p);
|
|
|
+ for (const auto& p : v.exclude) out->add_exclude(p);
|
|
|
+ for (const auto& f : v.where) {
|
|
|
+ auto* pbf = out->add_where();
|
|
|
+ pbf->set_field(f.field);
|
|
|
+ pbf->set_op(static_cast<pb::FilterOp>(static_cast<int>(f.op) + 1));
|
|
|
+ pbf->set_value(f.value.dump());
|
|
|
+ }
|
|
|
+ if (v.defaultSort) {
|
|
|
+ auto* s = out->mutable_default_sort();
|
|
|
+ s->set_field(v.defaultSort->field);
|
|
|
+ s->set_descending(v.defaultSort->descending);
|
|
|
+ }
|
|
|
+ out->set_created_at(v.createdAt);
|
|
|
+ out->set_updated_at(v.updatedAt);
|
|
|
+ }
|
|
|
+ return grpc::Status::OK;
|
|
|
+}
|
|
|
+
|
|
|
+grpc::Status DatabaseGrpcImpl::GetViewInfo(
|
|
|
+ grpc::ServerContext* /*context*/,
|
|
|
+ const pb::GetViewInfoRequest* request,
|
|
|
+ pb::GetViewInfoResponse* response
|
|
|
+) {
|
|
|
+ auto v = view_manager_.getView(request->name());
|
|
|
+ if (!v) {
|
|
|
+ response->set_found(false);
|
|
|
+ return grpc::Status::OK;
|
|
|
+ }
|
|
|
+ auto* out = response->mutable_view();
|
|
|
+ out->set_name(v->name);
|
|
|
+ out->set_collection(v->collection);
|
|
|
+ for (const auto& p : v->include) out->add_include(p);
|
|
|
+ for (const auto& p : v->exclude) out->add_exclude(p);
|
|
|
+ for (const auto& f : v->where) {
|
|
|
+ auto* pbf = out->add_where();
|
|
|
+ pbf->set_field(f.field);
|
|
|
+ pbf->set_op(static_cast<pb::FilterOp>(static_cast<int>(f.op) + 1));
|
|
|
+ pbf->set_value(f.value.dump());
|
|
|
+ }
|
|
|
+ if (v->defaultSort) {
|
|
|
+ auto* s = out->mutable_default_sort();
|
|
|
+ s->set_field(v->defaultSort->field);
|
|
|
+ s->set_descending(v->defaultSort->descending);
|
|
|
+ }
|
|
|
+ out->set_created_at(v->createdAt);
|
|
|
+ out->set_updated_at(v->updatedAt);
|
|
|
+ response->set_found(true);
|
|
|
+ return grpc::Status::OK;
|
|
|
+}
|
|
|
+
|
|
|
// ===== DatabaseReplicationGrpcImpl =====
|
|
|
|
|
|
DatabaseReplicationGrpcImpl::DatabaseReplicationGrpcImpl(
|