diff --git a/src/node/rpc/node_frontend.h b/src/node/rpc/node_frontend.h index 882b39eb438..ba420468271 100644 --- a/src/node/rpc/node_frontend.h +++ b/src/node/rpc/node_frontend.h @@ -467,320 +467,1405 @@ namespace ccf } } - public: - NodeEndpoints(NetworkState& network_, ccf::AbstractNodeContext& context_) : - CommonEndpointRegistry(get_actor_prefix(ActorsType::nodes), context_), - network(network_), - node_operation(*context_.get_subsystem()) - { - openapi_info.title = "CCF Public Node API"; - openapi_info.description = - "This API provides public, uncredentialed access to service and node " - "state."; - openapi_info.document_version = "5.0.8"; - } - - // NOLINTNEXTLINE(readability-function-cognitive-complexity) - void init_handlers() override + auto accept( + ccf::endpoints::EndpointContext& args, const nlohmann::json& params) { - CommonEndpointRegistry::init_handlers(); + const auto in = params.get(); - const auto self_cert_auth_policy = - std::make_shared(this->context); + // Not part of network => Internal error + if ( + !this->node_operation.is_part_of_network() && + !this->node_operation.is_part_of_public_network() && + !this->node_operation.is_reading_private_ledger()) + { + const std::string payload = + "Target node should be part of network to accept new nodes."; + LOG_INFO_FMT("Join request rejected: {}", payload); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + payload); + } - auto accept = [this](auto& args, const nlohmann::json& params) { - const auto in = params.get(); + // No service => Internal error + auto* service = args.tx.rw(this->network.service); + auto active_service = service->get(); + if (!active_service.has_value()) + { + const std::string payload = + "No service is available to accept new node."; + LOG_INFO_FMT("Join request rejected: {}", payload); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + payload); + } - // Not part of network => Internal error - if ( - !this->node_operation.is_part_of_network() && - !this->node_operation.is_part_of_public_network() && - !this->node_operation.is_reading_private_ledger()) + auto* current_consensus = get_consensus(); + const auto should_redirect_to_primary = + current_consensus != nullptr && !this->node_operation.can_replicate(); + auto redirect_to_primary = [&]() { + auto primary_id = current_consensus->primary(); + if (primary_id.has_value()) { + const auto address = node::get_redirect_address_for_node( + args, args.tx, primary_id.value()); + if (!address.has_value()) + { + LOG_INFO_FMT( + "Join request rejected: no redirect address for " + "primary {}", + primary_id.value()); + return already_populated_response(); + } + + args.rpc_ctx->set_response_header( + http::headers::LOCATION, + fmt::format("https://{}/node/join", address.value())); + const std::string payload = - "Target node should be part of network to accept new nodes."; - LOG_INFO_FMT("Join request rejected: {}", payload); + "Node is not primary; cannot handle write"; + LOG_INFO_FMT( + "Join request redirected to primary {} at {}: {}", + primary_id.value(), + address.value(), + payload); return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, + HTTP_STATUS_PERMANENT_REDIRECT, + ccf::errors::NodeCannotHandleRequest, payload); } - // No service => Internal error - auto service = args.tx.rw(this->network.service); - auto active_service = service->get(); - if (!active_service.has_value()) + const std::string payload = "Primary unknown"; + LOG_INFO_FMT("Join request rejected: {}", payload); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + payload); + }; + + auto* nodes = args.tx.ro(network.nodes); + + // If already joined => return equivalent response + auto existing_node_info = check_node_exists( + args.tx, args.rpc_ctx->get_session_context()->caller_cert); + if (existing_node_info.has_value()) + { + JoinNetworkNodeToNode::Out rep; + + // If the node already exists, return network secrets if is already + // trusted. Otherwise, only return its status + auto node_info = nodes->get(existing_node_info->node_id); + auto node_status = + node_info->status; // NOLINT(bugprone-unchecked-optional-access) + rep.node_status = node_status; + rep.node_id = existing_node_info->node_id; + if (node_status == NodeStatus::TRUSTED) { - const std::string payload = - "No service is available to accept new node."; - LOG_INFO_FMT("Join request rejected: {}", payload); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - payload); + rep.network_info = JoinNetworkNodeToNode::Out::NetworkInfo( + node_operation.is_part_of_public_network(), + node_operation.get_last_recovered_signed_idx(), + this->network.ledger_secrets->get( + args.tx, existing_node_info->ledger_secret_seqno), + *this->network.identity, + active_service->status, + existing_node_info->endorsed_certificate, + node_operation.get_cose_signatures_config()); + + LOG_DEBUG_FMT( + "Join request accepted: {} already marked as TRUSTED", + existing_node_info->node_id); + return make_success(rep); } - auto* current_consensus = get_consensus(); - const auto should_redirect_to_primary = - current_consensus != nullptr && !this->node_operation.can_replicate(); - auto redirect_to_primary = [&]() { - auto primary_id = current_consensus->primary(); - if (primary_id.has_value()) + if (node_status == NodeStatus::PENDING) + { + const auto pending_node_timeout = get_pending_node_timeout(); + if (!pending_node_timeout.has_value()) + { + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "NodeConfiguration subsystem is not available"); + } + + if (pending_node_timeout.value() > std::chrono::milliseconds::zero()) { - const auto address = node::get_redirect_address_for_node( - args, args.tx, primary_id.value()); - if (!address.has_value()) + if (should_redirect_to_primary) { - LOG_INFO_FMT( - "Join request rejected: no redirect address for " - "primary {}", - primary_id.value()); - return already_populated_response(); + return redirect_to_primary(); } - args.rpc_ctx->set_response_header( - http::headers::LOCATION, - fmt::format("https://{}/node/join", address.value())); + // NOLINTNEXTLINE(bugprone-unchecked-optional-access) + node_info->pending_last_seen = current_time_ms(); + args.tx.rw(network.nodes) + ->put( + existing_node_info->node_id, + // NOLINTNEXTLINE(bugprone-unchecked-optional-access) + node_info.value()); + } - const std::string payload = - "Node is not primary; cannot handle write"; - LOG_INFO_FMT( - "Join request redirected to primary {} at {}: {}", - primary_id.value(), - address.value(), - payload); - return make_error( - HTTP_STATUS_PERMANENT_REDIRECT, - ccf::errors::NodeCannotHandleRequest, - payload); + // Only return node status and ID + LOG_DEBUG_FMT( + "Join request accepted: {} already marked as PENDING", + existing_node_info->node_id); + return make_success(rep); + } + + const std::string payload = fmt::format( + "Joining node is not in expected state ({}).", node_status); + LOG_INFO_FMT("Join request rejected: {}", payload); + return make_error( + HTTP_STATUS_BAD_REQUEST, ccf::errors::InvalidNodeState, payload); + } + + // Not the primary => Redirect if possible to primary + if (should_redirect_to_primary) + { + return redirect_to_primary(); + } + + // Joiner's snapshot too old => StartupSeqnoIsOld + // (causes joiner to fetch a more recent snapshot) + // + // The joiner always wants to use the most recent snapshot. + // However this will result in the joiner chasing the primary if + // snapshot production period ~= snapshot fetching delay + // + // So we have hysteresis in the fetching constraint: + // If the joiner has already fetched a snapshot: joiner seqno > startup + // snapshot seqno Otherwise: joiner seqno > latest snapshot on disk + // seqno + auto this_startup_seqno = + this->node_operation.get_startup_snapshot_seqno(); + ccf::kv::Version required_seqno = this_startup_seqno; + // If the joiner does not enable fetching, or is a legacy node, + // join_fetch_count is unset and we should use the required bound to + // prevent it chasing the primary. + // Otherwise if this is the first request, use the preferred bound + bool using_preferred_bound = + (in.join_fetch_count.has_value() && in.join_fetch_count.value() == 0); + if (using_preferred_bound) + { + auto node_configuration_subsystem = + this->context.get_subsystem(); + if (node_configuration_subsystem != nullptr) + { + const auto& snapshots_config = + node_configuration_subsystem->get().node_config.snapshots; + const auto latest_committed_snapshot = + snapshots::find_latest_committed_snapshot_in_directory( + snapshots_config.directory); + if (latest_committed_snapshot.has_value()) + { + const auto latest_snapshot_seqno = + snapshots::get_snapshot_idx_from_file_name( + latest_committed_snapshot->filename().string()); + required_seqno = std::max( + required_seqno, + static_cast(latest_snapshot_seqno)); } + } + } + if ( + in.startup_seqno.has_value() && + in.startup_seqno.value() < required_seqno) + { + // Make sure that the joiner's snapshot is more recent than this + // node's snapshot. Otherwise, the joiner may not be given all the + // ledger secrets required to replay historical transactions. + const std::string payload = fmt::format( + "Node requested to join from seqno {} which is older than this " + "node {} {}. A snapshot at least as recent as {} must " + "be used instead.", + in.startup_seqno.value(), + using_preferred_bound ? "latest_on_disk_seqno" : "startup_seqno", + required_seqno, + required_seqno); + LOG_INFO_FMT("Join request rejected: {}", payload); + return make_error( + HTTP_STATUS_BAD_REQUEST, ccf::errors::StartupSeqnoIsOld, payload); + } - const std::string payload = "Primary unknown"; - LOG_INFO_FMT("Join request rejected: {}", payload); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - payload); - }; + auto joining_node_status = NodeStatus::PENDING; + // If the service is opening, new nodes are trusted straight away + if ( + active_service->status == ServiceStatus::OPENING || + active_service->status == ServiceStatus::RECOVERING) + { + joining_node_status = NodeStatus::TRUSTED; + } + + return add_node( + args.tx, + args.rpc_ctx->get_session_context()->caller_cert, + in, + joining_node_status, + active_service->status); + } + + auto remove_expired_pending( + ccf::endpoints::EndpointContext& ctx, nlohmann::json&& /*params*/) + { + const auto pending_node_timeout = get_pending_node_timeout(); + if (!pending_node_timeout.has_value()) + { + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "NodeConfiguration subsystem is not available"); + } + + if (pending_node_timeout.value() <= std::chrono::milliseconds::zero()) + { + return make_success(true); + } + + const auto now = current_time_ms(); + auto* nodes = ctx.tx.rw(network.nodes); + std::map untimestamped_pending_nodes; + std::vector expired_pending_nodes; + nodes->foreach([&](const auto& node_id, const auto& node_info) { + if (node_info.status != NodeStatus::PENDING) + { + return true; + } + + if ( + !node_info.pending_last_seen.has_value() || + node_info.pending_last_seen.value() < 0 || + node_info.pending_last_seen.value() > now) + { + auto updated_node_info = node_info; + updated_node_info.pending_last_seen = now; + untimestamped_pending_nodes.emplace( + node_id, std::move(updated_node_info)); + } + else if ( + now - node_info.pending_last_seen.value() >= + pending_node_timeout.value().count()) + { + expired_pending_nodes.push_back(node_id); + } + + return true; + }); + + for (const auto& [node_id, node_info] : untimestamped_pending_nodes) + { + nodes->put(node_id, node_info); + } + + for (const auto& node_id : expired_pending_nodes) + { + LOG_INFO_FMT("Removing expired Pending node {}", node_id); + InternalTablesAccess::remove_node(ctx.tx, node_id); + } - auto nodes = args.tx.ro(network.nodes); + return make_success(true); + } + + auto set_retired_committed( + ccf::endpoints::EndpointContext& ctx, nlohmann::json&& /*params*/) + { + auto* nodes = ctx.tx.rw(network.nodes); + nodes->foreach([&nodes](const auto& node_id, auto node_info) { + auto gc_node = nodes->get_globally_committed(node_id); + if ( + gc_node.has_value() && gc_node->status == ccf::NodeStatus::RETIRED && + !node_info.retired_committed) + { + // Set retired_committed on nodes for which RETIRED status + // has been committed. + node_info.retired_committed = true; + nodes->put(node_id, node_info); + + LOG_DEBUG_FMT("Setting retired_committed on node {}", node_id); + } + return true; + }); + + return make_success(); + } + + auto get_state( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + GetState::Out result; + auto [s, rts, lrs] = this->node_operation.state(); + result.node_id = this->context.get_node_id(); + result.state = s; + result.recovery_target_seqno = rts; + result.last_recovered_seqno = lrs; + result.startup_seqno = this->node_operation.get_startup_snapshot_seqno(); + + // Read last signed seqno from both raw and COSE signature tables + auto* signatures = args.tx.ro(Tables::SIGNATURES); + auto sig = signatures->get(); + + ccf::kv::Version raw_seqno = 0; + if (sig.has_value()) + { + raw_seqno = sig.value().seqno; + } - // If already joined => return equivalent response - auto existing_node_info = check_node_exists( - args.tx, args.rpc_ctx->get_session_context()->caller_cert); - if (existing_node_info.has_value()) + ccf::kv::Version cose_seqno = 0; + auto* cose_signatures = + args.tx.ro(Tables::COSE_SIGNATURES); + auto cose_sig = cose_signatures->get(ccf::IdentityType::CLASSICAL); + if (cose_sig.has_value() && !cose_sig->empty()) + { + auto receipt = ccf::cose::decode_ccf_receipt(cose_sig.value(), false); + auto txid = ccf::TxID::from_str(receipt.phdr.ccf.txid); + if (!txid.has_value()) + { + throw std::logic_error(fmt::format( + "Failed to parse txid from COSE signature: {}", + receipt.phdr.ccf.txid)); + } + cose_seqno = txid->seqno; + } + + result.last_signed_seqno = std::max(raw_seqno, cose_seqno); + + auto node_configuration_subsystem = + this->context.get_subsystem(); + if (!node_configuration_subsystem) + { + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "NodeConfigurationSubsystem is not available"); + } + result.stop_notice = + node_configuration_subsystem->has_received_stop_notice(); + + return make_success(result); + } + + auto get_quote( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + QuoteInfo node_quote_info; + const auto result = get_quote_for_this_node_v1(args.tx, node_quote_info); + if (result == ApiResult::OK) + { + Quote q; + q.node_id = context.get_node_id(); + q.raw = node_quote_info.quote; + q.endorsements = node_quote_info.endorsements; + q.format = node_quote_info.format; + q.uvm_endorsements = node_quote_info.uvm_endorsements; + + auto* nodes = args.tx.ro(network.nodes); + auto node_info = nodes->get(context.get_node_id()); + if (node_info.has_value() && node_info->code_digest.has_value()) { - JoinNetworkNodeToNode::Out rep; - - // If the node already exists, return network secrets if is already - // trusted. Otherwise, only return its status - auto node_info = nodes->get(existing_node_info->node_id); - auto node_status = node_info->status; - rep.node_status = node_status; - rep.node_id = existing_node_info->node_id; - if (node_status == NodeStatus::TRUSTED) + q.measurement = node_info->code_digest.value(); + } + else + { + auto measurement = + AttestationProvider::get_measurement(node_quote_info); + if (measurement.has_value()) + { + q.measurement = measurement.value().hex_str(); + } + else { - rep.network_info = JoinNetworkNodeToNode::Out::NetworkInfo( - node_operation.is_part_of_public_network(), - node_operation.get_last_recovered_signed_idx(), - this->network.ledger_secrets->get( - args.tx, existing_node_info->ledger_secret_seqno), - *this->network.identity, - active_service->status, - existing_node_info->endorsed_certificate, - node_operation.get_cose_signatures_config()); - - LOG_DEBUG_FMT( - "Join request accepted: {} already marked as TRUSTED", - existing_node_info->node_id); - return make_success(rep); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InvalidQuote, + "Failed to extract code id from node quote."); } + } + + return make_success(q); + } + + if (result == ApiResult::NotFound) + { + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "Could not find node quote."); + } + + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + fmt::format("Error code: {}", ccf::api_result_to_str(result))); + } + + auto get_quotes( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + GetQuotes::Out result; - if (node_status == NodeStatus::PENDING) + auto* nodes = args.tx.ro(network.nodes); + nodes->foreach( + ["es = result.quotes](const auto& node_id, const auto& node_info) { + if (node_info.status == ccf::NodeStatus::TRUSTED) { - const auto pending_node_timeout = get_pending_node_timeout(); - if (!pending_node_timeout.has_value()) + Quote q; + q.node_id = node_id; + q.raw = node_info.quote_info.quote; + q.endorsements = node_info.quote_info.endorsements; + q.format = node_info.quote_info.format; + q.uvm_endorsements = node_info.quote_info.uvm_endorsements; + + if (node_info.code_digest.has_value()) { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "NodeConfiguration subsystem is not available"); + q.measurement = node_info.code_digest.value(); } - - if ( - pending_node_timeout.value() > std::chrono::milliseconds::zero()) + else { - if (should_redirect_to_primary) + auto measurement = + AttestationProvider::get_measurement(node_info.quote_info); + if (measurement.has_value()) { - return redirect_to_primary(); + q.measurement = measurement.value().hex_str(); } - - node_info->pending_last_seen = current_time_ms(); - args.tx.rw(network.nodes) - ->put(existing_node_info->node_id, node_info.value()); } + quotes.emplace_back(q); + } + return true; + }); + + return make_success(result); + } + + auto get_attestations( + ccf::endpoints::ReadOnlyEndpointContext& args, nlohmann::json&& params) + { + auto res = get_quotes(args, std::move(params)); + const auto* body = std::get_if(&res); + if (body != nullptr) + { + auto result = nlohmann::json::object(); + result["attestations"] = (*body)["quotes"]; + return make_success(result); + } + + return res; + } + + auto network_status( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + GetNetworkInfo::Out out; + auto* service = args.tx.ro(network.service); + auto service_state = service->get(); + if (service_state.has_value()) + { + const auto& service_value = service_state.value(); + out.service_status = service_value.status; + out.service_certificate = service_value.cert; + out.recovery_count = service_value.recovery_count.value_or(0); + out.service_data = service_value.service_data; + out.current_service_create_txid = + service_value.current_service_create_txid; + auto* current_consensus = get_consensus(); + if (current_consensus != nullptr) + { + out.current_view = current_consensus->get_view(); + auto primary_id = current_consensus->primary(); + if (primary_id.has_value()) + { + out.primary_id = primary_id.value(); + } + } + return make_success(out); + } + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "Service state not available."); + } + + static auto service_previous_identity( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + auto* psi_handle = args.tx.ro( + ccf::Tables::PREVIOUS_SERVICE_IDENTITY); + const auto psi = psi_handle->get(); + if (psi.has_value()) + { + GetServicePreviousIdentity::Out out; + out.previous_service_identity = psi.value(); + return make_success(out); + } + + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "This service is not a recovery of a previous service."); + } + + auto get_nodes( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + const auto parsed_query = + http::parse_query(args.rpc_ctx->get_request_query()); + + std::string error_string; // Ignored - all params are optional + const auto host = http::get_query_value_opt( + parsed_query, "host", error_string); + const auto port = http::get_query_value_opt( + parsed_query, "port", error_string); + const auto status_str = http::get_query_value_opt( + parsed_query, "status", error_string); + + std::optional status; + if (status_str.has_value()) + { + // Convert the query argument to a JSON string, try to parse it as + // a NodeStatus, return an error if this doesn't work + try + { + status = nlohmann::json(status_str.value()).get(); + } + catch (const ccf::JsonParseError& e) + { + return ccf::make_error( + HTTP_STATUS_BAD_REQUEST, + ccf::errors::InvalidQueryParameterValue, + fmt::format( + "Query parameter '{}' is not a valid node status", + status_str.value())); + } + } + + GetNodes::Out out; + + auto* nodes = args.tx.ro(this->network.nodes); + auto* current_consensus = get_consensus(); + nodes->foreach([host, port, status, &out, nodes, current_consensus]( + const NodeId& nid, const NodeInfo& ni) { + if (status.has_value() && status.value() != ni.status) + { + return true; + } + + // Match on any interface + bool is_matched = false; + for (auto const& interface : ni.rpc_interfaces) + { + const auto& [pub_host, pub_port] = + split_net_address(interface.second.published_address); - // Only return node status and ID - LOG_DEBUG_FMT( - "Join request accepted: {} already marked as PENDING", - existing_node_info->node_id); - return make_success(rep); + if ( + (!host.has_value() || host.value() == pub_host) && + (!port.has_value() || port.value() == pub_port)) + { + is_matched = true; + break; } + } + + if (!is_matched) + { + return true; + } + + bool is_primary = false; + if (current_consensus != nullptr) + { + is_primary = current_consensus->primary() == nid; + } + + out.nodes.push_back( + {nid, + ni.status, + is_primary, + ni.rpc_interfaces, + ni.node_data, + nodes->get_version_of_previous_write(nid).value_or(0)}); + return true; + }); + + return make_success(out); + } + + auto get_removable_nodes( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + GetNodes::Out out; + + auto* nodes = args.tx.ro(this->network.nodes); + nodes->foreach( + [&out, nodes](const NodeId& node_id, const NodeInfo& /*ni*/) { + // Only nodes whose retire_committed status is committed can be + // safely removed, because any primary elected from here on would + // consider them retired, and would consequently not need their + // input in any quorum. We must therefore read the KV at its + // globally committed watermark, for the purpose of this RPC. Since + // this transaction does not perform a write, it is safe to do this. + auto node = nodes->get_globally_committed(node_id); + if ( + node.has_value() && node->status == ccf::NodeStatus::RETIRED && + node->retired_committed) + { + out.nodes.push_back( + {node_id, + node->status, + false /* is_primary */, + node->rpc_interfaces, + node->node_data, + nodes->get_version_of_previous_write(node_id).value_or(0)}); + } + return true; + }); + + return make_success(out); + } + + auto delete_retired_committed_node( + ccf::endpoints::EndpointContext& args, nlohmann::json&& /*params*/) + { + GetNodes::Out out; + + std::string node_id; + std::string error; + if (!get_path_param( + args.rpc_ctx->get_request_path_params(), "node_id", node_id, error)) + { + return make_error( + HTTP_STATUS_BAD_REQUEST, ccf::errors::InvalidResourceName, error); + } + + auto* nodes = args.tx.rw(this->network.nodes); + if (!nodes->has(node_id)) + { + return make_error( + HTTP_STATUS_NOT_FOUND, ccf::errors::ResourceNotFound, "No such node"); + } + + // A node's retirement is only complete when the + // transition of retired_committed is itself committed, + // i.e. when the next eligible primary is guaranteed to + // be aware the retirement is committed. + // As a result, the handler must check node info at the + // current committed level, rather than at the end of the + // local suffix. + // While this transaction does execute a write, it specifically + // deletes the value it reads from. It is therefore safe to + // execute on the basis of a potentially stale read-set, + // which get_globally_committed() typically produces. + auto node = nodes->get_globally_committed(node_id); + if ( + node.has_value() && node->status == ccf::NodeStatus::RETIRED && + node->retired_committed) + { + InternalTablesAccess::remove_node(args.tx, node_id); + } + else + { + return make_error( + HTTP_STATUS_BAD_REQUEST, + ccf::errors::NodeNotRetiredCommitted, + "Node is not completely retired"); + } + + return make_success(true); + } + + auto get_self_signed_certificate( + ccf::endpoints::CommandEndpointContext& /*args*/, + nlohmann::json&& /*params*/) + { + return SelfSignedNodeCertificateInfo{ + this->node_operation.get_self_signed_node_certificate()}; + } + + auto get_node_info( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + std::string node_id; + std::string error; + if (!get_path_param( + args.rpc_ctx->get_request_path_params(), "node_id", node_id, error)) + { + return make_error( + HTTP_STATUS_BAD_REQUEST, ccf::errors::InvalidResourceName, error); + } + + auto* nodes = args.tx.ro(this->network.nodes); + auto info = nodes->get(node_id); + + if (!info) + { + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "Node not found"); + } + + bool is_primary = false; + auto* current_consensus = get_consensus(); + if (current_consensus != nullptr) + { + auto primary = current_consensus->primary(); + if (primary.has_value() && primary.value() == node_id) + { + is_primary = true; + } + } + auto& ni = info.value(); + return make_success(GetNode::Out{ + node_id, + ni.status, + is_primary, + ni.rpc_interfaces, + ni.node_data, + nodes->get_version_of_previous_write(node_id).value_or(0)}); + } + + auto get_self_node( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + auto node_id = this->context.get_node_id(); + auto* nodes = args.tx.ro(this->network.nodes); + auto info = nodes->get(node_id); + + bool is_primary = false; + auto* current_consensus = get_consensus(); + if (current_consensus != nullptr) + { + auto primary = current_consensus->primary(); + if (primary.has_value() && primary.value() == node_id) + { + is_primary = true; + } + } + + if (info.has_value()) + { + // Answers from the KV are preferred, as they are more up-to-date, + // especially status and node_data. + auto& ni = info.value(); + return make_success(GetNode::Out{ + node_id, + ni.status, + is_primary, + ni.rpc_interfaces, + ni.node_data, + nodes->get_version_of_previous_write(node_id).value_or(0)}); + } + + // If the node isn't in its KV yet, fall back to configuration + auto node_configuration_subsystem = + this->context.get_subsystem(); + if (!node_configuration_subsystem) + { + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "NodeConfigurationSubsystem is not available"); + } + const auto& node_startup_config = + node_configuration_subsystem->get().node_config; + return make_success(GetNode::Out{ + node_id, + ccf::NodeStatus::PENDING, + is_primary, + node_startup_config.network.rpc_interfaces, + node_configuration_subsystem->get().node_data, + 0}); + } + + auto get_primary_node( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + auto* current_consensus = get_consensus(); + if (current_consensus != nullptr) + { + auto primary_id = current_consensus->primary(); + if (!primary_id.has_value()) + { + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "Primary unknown"); + } + + auto* nodes = args.tx.ro(this->network.nodes); + auto info = nodes->get(primary_id.value()); + if (!info) + { + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "Node not found"); + } + + auto& ni = info.value(); + return make_success(GetNode::Out{ + primary_id.value(), + ni.status, + true, + ni.rpc_interfaces, + ni.node_data, + nodes->get_version_of_previous_write(primary_id.value()) + .value_or(0)}); + } + + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "No configured consensus"); + } + + auto head_primary(ccf::endpoints::ReadOnlyEndpointContext& args) + { + if (this->node_operation.can_replicate()) + { + args.rpc_ctx->set_response_status(HTTP_STATUS_OK); + } + else + { + auto* current_consensus = get_consensus(); + if (current_consensus == nullptr) + { + args.rpc_ctx->set_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "Consensus not initialised"); + return; + } + + auto primary_id = current_consensus->primary(); + if (!primary_id.has_value()) + { + args.rpc_ctx->set_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "Primary unknown"); + return; + } + + const auto address = node::get_redirect_address_for_node( + args, args.tx, primary_id.value()); + if (!address.has_value()) + { + return; + } + + args.rpc_ctx->set_response_header( + http::headers::LOCATION, + fmt::format("https://{}/node/primary", address.value())); + args.rpc_ctx->set_response_status(HTTP_STATUS_PERMANENT_REDIRECT); + } + } + + auto get_primary(ccf::endpoints::ReadOnlyEndpointContext& args) + { + if (this->node_operation.can_replicate()) + { + args.rpc_ctx->set_response_status(HTTP_STATUS_OK); + return; + } + + args.rpc_ctx->set_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "Node is not primary"); + } + + auto get_backup(ccf::endpoints::ReadOnlyEndpointContext& args) + { + if (!this->node_operation.can_replicate()) + { + args.rpc_ctx->set_response_status(HTTP_STATUS_OK); + return; + } + + args.rpc_ctx->set_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "Node is not backup"); + } + + auto consensus_config( + ccf::endpoints::CommandEndpointContext& /*args*/, + nlohmann::json&& /*params*/) + { + // Query node for configurations, separate current from pending + auto* current_consensus = get_consensus(); + if (current_consensus != nullptr) + { + auto cfg = current_consensus->get_latest_configuration(); + ConsensusConfig cc; + for (auto& [nid, ninfo] : cfg) + { + cc.emplace( + nid.value(), + ConsensusNodeConfig{ + fmt::format("{}:{}", ninfo.hostname, ninfo.port)}); + } + return make_success(cc); + } + + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "No configured consensus"); + } + + auto consensus_state( + ccf::endpoints::CommandEndpointContext& /*args*/, + nlohmann::json&& /*params*/) + { + auto* current_consensus = get_consensus(); + if (current_consensus != nullptr) + { + return make_success( + ConsensusConfigDetails{current_consensus->get_details()}); + } + + return make_error( + HTTP_STATUS_NOT_FOUND, + ccf::errors::ResourceNotFound, + "No configured consensus"); + } + + auto node_metrics(ccf::endpoints::CommandEndpointContext& args) + { + NodeMetrics nm; + nm.sessions = node_operation.get_session_metrics(); + + args.rpc_ctx->set_response_status(HTTP_STATUS_OK); + args.rpc_ctx->set_response_header( + http::headers::CONTENT_TYPE, http::headervalues::contenttype::JSON); + args.rpc_ctx->set_response_body(nlohmann::json(nm).dump()); + } + + auto js_metrics( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& /*params*/) + { + auto* bytecode_map = args.tx.ro(this->network.modules_quickjs_bytecode); + auto* version_val = args.tx.ro(this->network.modules_quickjs_version); + uint64_t bytecode_size = 0; + bytecode_map->foreach( + [&bytecode_size](const auto&, const auto& bytecode) { + bytecode_size += bytecode.size(); + return true; + }); + auto* js_engine_map = args.tx.ro(this->network.js_engine); + JavaScriptMetrics m; + m.bytecode_size = bytecode_size; + m.bytecode_used = version_val->get() == std::string(ccf::quickjs_version); + + auto options = js_engine_map->get().value_or(ccf::JSRuntimeOptions{}); + m.max_stack_size = options.max_stack_bytes; + m.max_heap_size = options.max_heap_bytes; + m.max_execution_time = options.max_execution_time_ms; + m.max_cached_interpreters = options.max_cached_interpreters; + + return m; + } + + static auto version( + ccf::endpoints::CommandEndpointContext& /*args*/, + nlohmann::json&& /*params*/) + { + GetVersion::Out result; + result.ccf_version = ccf::ccf_version; + result.quickjs_version = ccf::quickjs_version; + result.unsafe = false; + + return make_success(result); + } + + auto create(ccf::endpoints::EndpointContext& ctx, nlohmann::json&& params) + { + LOG_INFO_FMT("Processing create RPC"); + + bool recovering = node_operation.is_reading_public_ledger(); + + // This endpoint can only be called once, directly from the starting + // node for the genesis or end of public recovery transaction to + // initialise the service + if (!node_operation.is_in_initialised_state() && !recovering) + { + return make_error( + HTTP_STATUS_FORBIDDEN, + ccf::errors::InternalError, + "Node is not in initial state."); + } + + const auto in = params.get(); + + if (InternalTablesAccess::is_service_created(ctx.tx, in.service_cert)) + { + return make_error( + HTTP_STATUS_FORBIDDEN, + ccf::errors::InternalError, + "Service is already created."); + } + + InternalTablesAccess::create_service( + ctx.tx, in.service_cert, in.create_txid, in.service_data, recovering); + + if (recovering) + { + // Recovery starts with a fresh consensus configuration, so previous + // service nodes can be removed immediately. + InternalTablesAccess::remove_previous_service_nodes(ctx.tx); + } + + // Genesis transaction (i.e. not after recovery) + if (in.genesis_info.has_value()) + { + // Note that it is acceptable to start a network without any member + // having a recovery share. The service will check that at least one + // recovery member is added before the service is opened. + for (const auto& info : in.genesis_info->members) + { + InternalTablesAccess::add_member(ctx.tx, info); + } + + InternalTablesAccess::init_configuration( + ctx.tx, in.genesis_info->service_configuration); + InternalTablesAccess::set_constitution( + ctx.tx, in.genesis_info->constitution); + } + else + { + // On recovery, force a new ledger chunk + auto* tx_ = dynamic_cast(&ctx.tx); + if (tx_ == nullptr) + { + throw std::logic_error("Could not cast tx to CommittableTx"); + } + tx_->set_tx_flag( + ccf::kv::CommittableTx::TxFlag::LEDGER_CHUNK_BEFORE_THIS_TX); + } + + auto* endorsed_certificates = + ctx.tx.rw(network.node_endorsed_certificates); + endorsed_certificates->put(in.node_id, in.node_endorsed_certificate); + + NodeInfo node_info = { + in.node_info_network, + {in.quote_info}, + in.public_encryption_key, + NodeStatus::TRUSTED, + std::nullopt, + in.measurement.hex_str(), + in.certificate_signing_request, + in.public_key, + in.node_data}; + InternalTablesAccess::add_node(ctx.tx, in.node_id, node_info); + + if (in.sealing_recovery_data.has_value()) + { + const auto& [sealing_keys, sealing_recovery_name] = + in.sealing_recovery_data.value(); + auto* sealed_recovery_keys = + ctx.tx.rw(Tables::SEALED_RECOVERY_KEYS); + sealed_recovery_keys->put(in.node_id, sealing_keys); + + auto* local_sealing_node_id_map = + ctx.tx.rw(Tables::SEALING_RECOVERY_NAMES); + local_sealing_node_id_map->put(sealing_recovery_name, in.node_id); + } + + node_operation.shuffle_sealed_shares(ctx.tx); + + if ( + in.quote_info.format != QuoteFormat::amd_sev_snp_v1 || + !in.snp_uvm_endorsements.has_value()) + { + // For improved serviceability on SNP, do not record trusted + // measurements if UVM endorsements are available + InternalTablesAccess::trust_node_measurement( + ctx.tx, in.measurement, in.quote_info.format); + } + + switch (in.quote_info.format) + { + case QuoteFormat::insecure_virtual: + { + auto host_data = AttestationProvider::get_host_data(in.quote_info); + if (host_data.has_value()) + { + InternalTablesAccess::trust_node_virtual_host_data( + ctx.tx, host_data.value()); + } + else + { + LOG_FAIL_FMT("Unable to extract host data from virtual quote"); + } + break; + } + + case QuoteFormat::amd_sev_snp_v1: + { + auto host_data = + // NOLINTNEXTLINE(bugprone-unchecked-optional-access) + AttestationProvider::get_host_data(in.quote_info).value(); + InternalTablesAccess::trust_node_snp_host_data( + ctx.tx, host_data, in.snp_security_policy); + + InternalTablesAccess::trust_node_uvm_endorsements( + ctx.tx, in.snp_uvm_endorsements, recovering); + + auto attestation = + // NOLINTNEXTLINE(bugprone-unchecked-optional-access) + AttestationProvider::get_snp_attestation_report(in.quote_info) + .value(); + InternalTablesAccess::trust_node_snp_tcb_version( + ctx.tx, attestation, recovering); + break; + } + case QuoteFormat::oe_sgx_v1: + { + break; + } + } + + std::optional digest = + ccf::get_create_tx_claims_digest(ctx.tx); + if (digest.has_value()) + { + auto digest_value = digest.value(); + ctx.rpc_ctx->set_claims_digest(std::move(digest_value)); + } + + this->node_operation.recovery_decision_protocol().reset_state(ctx.tx); + this->node_operation.recovery_decision_protocol().try_start( + ctx.tx, recovering); + + LOG_INFO_FMT("Created service"); + return make_success(true); + } + + auto refresh_jwt_keys( + ccf::endpoints::EndpointContext& ctx, nlohmann::json&& body) + { + // All errors are server errors since the client is the server. + + auto* current_consensus = get_consensus(); + auto primary_id = current_consensus->primary(); + if (!primary_id.has_value()) + { + LOG_FAIL_FMT("JWT key auto-refresh: primary unknown"); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "Primary is unknown"); + } + + const auto& sig_auth_ident = ctx.get_caller(); + if (primary_id.value() != sig_auth_ident.node_id) + { + LOG_FAIL_FMT( + "JWT key auto-refresh: request does not originate from primary"); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "Request does not originate from primary."); + } + + SetJwtPublicSigningKeys parsed; + try + { + parsed = body.get(); + } + catch (const ccf::JsonParseError& e) + { + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "Unable to parse body."); + } + + auto* issuers = ctx.tx.ro(this->network.jwt_issuers); + auto issuer_metadata_ = issuers->get(parsed.issuer); + if (!issuer_metadata_.has_value()) + { + LOG_FAIL_FMT( + "JWT key auto-refresh: {} is not a valid issuer", parsed.issuer); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + fmt::format("{} is not a valid issuer.", parsed.issuer)); + } + auto& issuer_metadata = issuer_metadata_.value(); + + if (!issuer_metadata.auto_refresh) + { + LOG_FAIL_FMT( + "JWT key auto-refresh: {} does not have auto_refresh enabled", + parsed.issuer); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + fmt::format("{} does not have auto_refresh enabled.", parsed.issuer)); + } + + if (!set_jwt_public_signing_keys( + ctx.tx, + "", + parsed.issuer, + issuer_metadata, + parsed.jwks)) + { + LOG_FAIL_FMT( + "JWT key auto-refresh: error while storing signing keys for issuer " + "{}", + parsed.issuer); + return make_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + fmt::format( + "Error while storing signing keys for issuer {}.", parsed.issuer)); + } + + return make_success(true); + } + + auto get_jwt_metrics( + ccf::endpoints::ReadOnlyEndpointContext& /*args*/, + const nlohmann::json& /*params*/) + { + JWTRefreshMetrics metrics; + { + ccf::ds::MutexGuard guard(jwt_refresh_metrics_lock); + metrics = jwt_refresh_metrics; + } + return make_success(metrics); + } + + auto service_config_handler( + ccf::endpoints::EndpointContext& args, const nlohmann::json& /*params*/) + { + return make_success(args.tx.ro(network.config)->get()); + } + + auto list_indexing_strategies( + ccf::endpoints::EndpointContext& /*args*/, + const nlohmann::json& /*params*/) + { + return make_success(this->context.get_indexing_strategies().describe()); + } + + auto get_ready_app(ccf::endpoints::CommandEndpointContext& ctx) + { + auto node_configuration_subsystem = + this->context.get_subsystem(); + if (!node_configuration_subsystem) + { + ctx.rpc_ctx->set_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "NodeConfigurationSubsystem is not available"); + return; + } + if ( + !node_configuration_subsystem->has_received_stop_notice() && + this->node_operation.is_part_of_network() && + this->node_operation.is_user_frontend_open()) + { + ctx.rpc_ctx->set_response_status(HTTP_STATUS_NO_CONTENT); + } + else + { + ctx.rpc_ctx->set_response_status(HTTP_STATUS_SERVICE_UNAVAILABLE); + } + } + + auto get_ready_gov(ccf::endpoints::CommandEndpointContext& ctx) + { + auto node_configuration_subsystem = + this->context.get_subsystem(); + if (!node_configuration_subsystem) + { + ctx.rpc_ctx->set_error( + HTTP_STATUS_INTERNAL_SERVER_ERROR, + ccf::errors::InternalError, + "NodeConfigurationSubsystem is not available"); + return; + } + if ( + !node_configuration_subsystem->has_received_stop_notice() && + this->node_operation.is_accessible_to_members() && + this->node_operation.is_member_frontend_open()) + { + ctx.rpc_ctx->set_response_status(HTTP_STATUS_NO_CONTENT); + } + else + { + ctx.rpc_ctx->set_response_status(HTTP_STATUS_SERVICE_UNAVAILABLE); + } + } + + auto create_snapshot( + ccf::endpoints::EndpointContext& args, nlohmann::json&& /*params*/) + { + auto* snapshot_create = + args.tx.rw(ccf::Tables::SNAPSHOT_CREATE); + snapshot_create->touch(); + this->node_operation.trigger_snapshot(args.tx); + return make_success(); + } - const std::string payload = fmt::format( - "Joining node is not in expected state ({}).", node_status); - LOG_INFO_FMT("Join request rejected: {}", payload); - return make_error( - HTTP_STATUS_BAD_REQUEST, ccf::errors::InvalidNodeState, payload); - } + auto historical_cache_info( + [[maybe_unused]] ccf::endpoints::ReadOnlyEndpointContext& args, + [[maybe_unused]] nlohmann::json&& /*params*/) + { + GetHistoricalCacheInfo::Out result{}; + result.estimated_size = + this->context.get_historical_state().get_estimated_store_cache_size(); + return make_success(result); + } - // Not the primary => Redirect if possible to primary - if (should_redirect_to_primary) - { - return redirect_to_primary(); - } + public: + NodeEndpoints(NetworkState& network_, ccf::AbstractNodeContext& context_) : + CommonEndpointRegistry(get_actor_prefix(ActorsType::nodes), context_), + network(network_), + node_operation(*context_.get_subsystem()) + { + openapi_info.title = "CCF Public Node API"; + openapi_info.description = + "This API provides public, uncredentialed access to service and node " + "state."; + openapi_info.document_version = "5.0.8"; + } - // Joiner's snapshot too old => StartupSeqnoIsOld - // (causes joiner to fetch a more recent snapshot) - // - // The joiner always wants to use the most recent snapshot. - // However this will result in the joiner chasing the primary if - // snapshot production period ~= snapshot fetching delay - // - // So we have hysteresis in the fetching constraint: - // If the joiner has already fetched a snapshot: joiner seqno > startup - // snapshot seqno Otherwise: joiner seqno > latest snapshot on disk - // seqno - auto this_startup_seqno = - this->node_operation.get_startup_snapshot_seqno(); - ccf::kv::Version required_seqno = this_startup_seqno; - // If the joiner does not enable fetching, or is a legacy node, - // join_fetch_count is unset and we should use the required bound to - // prevent it chasing the primary. - // Otherwise if this is the first request, use the preferred bound - bool using_preferred_bound = - (in.join_fetch_count.has_value() && in.join_fetch_count.value() == 0); - if (using_preferred_bound) - { - auto node_configuration_subsystem = - this->context.get_subsystem(); - if (node_configuration_subsystem != nullptr) - { - const auto& snapshots_config = - node_configuration_subsystem->get().node_config.snapshots; - const auto latest_committed_snapshot = - snapshots::find_latest_committed_snapshot_in_directory( - snapshots_config.directory); - if (latest_committed_snapshot.has_value()) - { - const auto latest_snapshot_seqno = - snapshots::get_snapshot_idx_from_file_name( - latest_committed_snapshot->filename().string()); - required_seqno = std::max( - required_seqno, - static_cast(latest_snapshot_seqno)); - } - } - } - if ( - in.startup_seqno.has_value() && - in.startup_seqno.value() < required_seqno) - { - // Make sure that the joiner's snapshot is more recent than this - // node's snapshot. Otherwise, the joiner may not be given all the - // ledger secrets required to replay historical transactions. - const std::string payload = fmt::format( - "Node requested to join from seqno {} which is older than this " - "node {} {}. A snapshot at least as recent as {} must " - "be used instead.", - in.startup_seqno.value(), - using_preferred_bound ? "latest_on_disk_seqno" : "startup_seqno", - required_seqno, - required_seqno); - LOG_INFO_FMT("Join request rejected: {}", payload); - return make_error( - HTTP_STATUS_BAD_REQUEST, ccf::errors::StartupSeqnoIsOld, payload); - } + void init_handlers() override + { + CommonEndpointRegistry::init_handlers(); - auto joining_node_status = NodeStatus::PENDING; - // If the service is opening, new nodes are trusted straight away - if ( - active_service->status == ServiceStatus::OPENING || - active_service->status == ServiceStatus::RECOVERING) - { - joining_node_status = NodeStatus::TRUSTED; - } + const auto self_cert_auth_policy = + std::make_shared(this->context); - return add_node( - args.tx, - args.rpc_ctx->get_session_context()->caller_cert, - in, - joining_node_status, - active_service->status); + auto accept = [this]( + ccf::endpoints::EndpointContext& args, + const nlohmann::json& params) { + return this->accept(args, params); }; make_endpoint("/join", HTTP_POST, json_adapter(accept), no_auth_required) .set_forwarding_required(endpoints::ForwardingRequired::Never) .set_openapi_hidden(true) .install(); - auto remove_expired_pending = [this](auto& ctx, nlohmann::json&&) { - const auto pending_node_timeout = get_pending_node_timeout(); - if (!pending_node_timeout.has_value()) - { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "NodeConfiguration subsystem is not available"); - } - - if (pending_node_timeout.value() <= std::chrono::milliseconds::zero()) - { - return make_success(true); - } - - const auto now = current_time_ms(); - auto nodes = ctx.tx.rw(network.nodes); - std::map untimestamped_pending_nodes; - std::vector expired_pending_nodes; - nodes->foreach([&](const auto& node_id, const auto& node_info) { - if (node_info.status != NodeStatus::PENDING) - { - return true; - } - - if ( - !node_info.pending_last_seen.has_value() || - node_info.pending_last_seen.value() < 0 || - node_info.pending_last_seen.value() > now) - { - auto updated_node_info = node_info; - updated_node_info.pending_last_seen = now; - untimestamped_pending_nodes.emplace( - node_id, std::move(updated_node_info)); - } - else if ( - now - node_info.pending_last_seen.value() >= - pending_node_timeout.value().count()) - { - expired_pending_nodes.push_back(node_id); - } - - return true; - }); - - for (const auto& [node_id, node_info] : untimestamped_pending_nodes) - { - nodes->put(node_id, node_info); - } - - for (const auto& node_id : expired_pending_nodes) - { - LOG_INFO_FMT("Removing expired Pending node {}", node_id); - InternalTablesAccess::remove_node(ctx.tx, node_id); - } - - return make_success(true); - }; + auto remove_expired_pending = + [this](ccf::endpoints::EndpointContext& ctx, nlohmann::json&& json) { + return this->remove_expired_pending(ctx, std::move(json)); + }; make_endpoint( "network/nodes/remove_expired_pending", HTTP_POST, @@ -790,27 +1875,10 @@ namespace ccf .set_openapi_hidden(true) .install(); - auto set_retired_committed = [this](auto& ctx, nlohmann::json&&) { - auto nodes = ctx.tx.rw(network.nodes); - nodes->foreach([&nodes](const auto& node_id, auto node_info) { - auto gc_node = nodes->get_globally_committed(node_id); - if ( - gc_node.has_value() && - gc_node->status == ccf::NodeStatus::RETIRED && - !node_info.retired_committed) - { - // Set retired_committed on nodes for which RETIRED status - // has been committed. - node_info.retired_committed = true; - nodes->put(node_id, node_info); - - LOG_DEBUG_FMT("Setting retired_committed on node {}", node_id); - } - return true; - }); - - return make_success(); - }; + auto set_retired_committed = + [this](ccf::endpoints::EndpointContext& ctx, nlohmann::json&& json) { + return this->set_retired_committed(ctx, std::move(json)); + }; make_endpoint( "network/nodes/set_retired_committed", HTTP_POST, @@ -819,58 +1887,10 @@ namespace ccf .set_openapi_hidden(true) .install(); - auto get_state = [this](auto& args, nlohmann::json&&) { - GetState::Out result; - auto [s, rts, lrs] = this->node_operation.state(); - result.node_id = this->context.get_node_id(); - result.state = s; - result.recovery_target_seqno = rts; - result.last_recovered_seqno = lrs; - result.startup_seqno = - this->node_operation.get_startup_snapshot_seqno(); - - // Read last signed seqno from both raw and COSE signature tables - auto signatures = args.tx.template ro(Tables::SIGNATURES); - auto sig = signatures->get(); - - ccf::kv::Version raw_seqno = 0; - if (sig.has_value()) - { - raw_seqno = sig.value().seqno; - } - - ccf::kv::Version cose_seqno = 0; - auto cose_signatures = - args.tx.template ro(Tables::COSE_SIGNATURES); - auto cose_sig = cose_signatures->get(ccf::IdentityType::CLASSICAL); - if (cose_sig.has_value() && !cose_sig->empty()) - { - auto receipt = ccf::cose::decode_ccf_receipt(cose_sig.value(), false); - auto txid = ccf::TxID::from_str(receipt.phdr.ccf.txid); - if (!txid.has_value()) - { - throw std::logic_error(fmt::format( - "Failed to parse txid from COSE signature: {}", - receipt.phdr.ccf.txid)); - } - cose_seqno = txid->seqno; - } - - result.last_signed_seqno = std::max(raw_seqno, cose_seqno); - - auto node_configuration_subsystem = - this->context.get_subsystem(); - if (!node_configuration_subsystem) - { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "NodeConfigurationSubsystem is not available"); - } - result.stop_notice = - node_configuration_subsystem->has_received_stop_notice(); - - return make_success(result); + auto get_state = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_state(args, std::move(json)); }; make_read_only_endpoint( "/state", HTTP_GET, json_read_only_adapter(get_state), no_auth_required) @@ -878,57 +1898,10 @@ namespace ccf .set_forwarding_required(endpoints::ForwardingRequired::Never) .install(); - auto get_quote = [this](auto& args, nlohmann::json&&) { - QuoteInfo node_quote_info; - const auto result = - get_quote_for_this_node_v1(args.tx, node_quote_info); - if (result == ApiResult::OK) - { - Quote q; - q.node_id = context.get_node_id(); - q.raw = node_quote_info.quote; - q.endorsements = node_quote_info.endorsements; - q.format = node_quote_info.format; - q.uvm_endorsements = node_quote_info.uvm_endorsements; - - auto nodes = args.tx.ro(network.nodes); - auto node_info = nodes->get(context.get_node_id()); - if (node_info.has_value() && node_info->code_digest.has_value()) - { - q.measurement = node_info->code_digest.value(); - } - else - { - auto measurement = - AttestationProvider::get_measurement(node_quote_info); - if (measurement.has_value()) - { - q.measurement = measurement.value().hex_str(); - } - else - { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InvalidQuote, - "Failed to extract code id from node quote."); - } - } - - return make_success(q); - } - - if (result == ApiResult::NotFound) - { - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "Could not find node quote."); - } - - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - fmt::format("Error code: {}", ccf::api_result_to_str(result))); + auto get_quote = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_quote(args, std::move(json)); }; make_read_only_endpoint( "/quotes/self", @@ -947,40 +1920,10 @@ namespace ccf .set_forwarding_required(endpoints::ForwardingRequired::Never) .install(); - auto get_quotes = [this](auto& args, nlohmann::json&&) { - GetQuotes::Out result; - - auto nodes = args.tx.ro(network.nodes); - nodes->foreach(["es = result.quotes]( - const auto& node_id, const auto& node_info) { - if (node_info.status == ccf::NodeStatus::TRUSTED) - { - Quote q; - q.node_id = node_id; - q.raw = node_info.quote_info.quote; - q.endorsements = node_info.quote_info.endorsements; - q.format = node_info.quote_info.format; - q.uvm_endorsements = node_info.quote_info.uvm_endorsements; - - if (node_info.code_digest.has_value()) - { - q.measurement = node_info.code_digest.value(); - } - else - { - auto measurement = - AttestationProvider::get_measurement(node_info.quote_info); - if (measurement.has_value()) - { - q.measurement = measurement.value().hex_str(); - } - } - quotes.emplace_back(q); - } - return true; - }); - - return make_success(result); + auto get_quotes = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_quotes(args, std::move(json)); }; make_read_only_endpoint( "/quotes", @@ -990,19 +1933,11 @@ namespace ccf .set_auto_schema() .install(); - auto get_attestations = - [get_quotes](auto& args, nlohmann::json&& params) { - auto res = get_quotes(args, std::move(params)); - const auto* body = std::get_if(&res); - if (body != nullptr) - { - auto result = nlohmann::json::object(); - result["attestations"] = (*body)["quotes"]; - return make_success(result); - } - - return res; - }; + auto get_attestations = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& params) { + return this->get_attestations(args, std::move(params)); + }; make_read_only_endpoint( "/attestations", HTTP_GET, @@ -1011,35 +1946,10 @@ namespace ccf .set_auto_schema() .install(); - auto network_status = [this](auto& args, nlohmann::json&&) { - GetNetworkInfo::Out out; - auto service = args.tx.ro(network.service); - auto service_state = service->get(); - if (service_state.has_value()) - { - const auto& service_value = service_state.value(); - out.service_status = service_value.status; - out.service_certificate = service_value.cert; - out.recovery_count = service_value.recovery_count.value_or(0); - out.service_data = service_value.service_data; - out.current_service_create_txid = - service_value.current_service_create_txid; - auto* current_consensus = get_consensus(); - if (current_consensus != nullptr) - { - out.current_view = current_consensus->get_view(); - auto primary_id = current_consensus->primary(); - if (primary_id.has_value()) - { - out.primary_id = primary_id.value(); - } - } - return make_success(out); - } - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "Service state not available."); + auto network_status = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->network_status(args, std::move(json)); }; make_read_only_endpoint( "/network", @@ -1049,22 +1959,13 @@ namespace ccf .set_auto_schema() .install(); - auto service_previous_identity = [](auto& args, nlohmann::json&&) { - auto psi_handle = args.tx.template ro( - ccf::Tables::PREVIOUS_SERVICE_IDENTITY); - const auto psi = psi_handle->get(); - if (psi.has_value()) - { - GetServicePreviousIdentity::Out out; - out.previous_service_identity = psi.value(); - return make_success(out); - } - - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "This service is not a recovery of a previous service."); - }; + auto service_previous_identity = + []( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return NodeEndpoints::service_previous_identity( + args, std::move(json)); + }; make_read_only_endpoint( "/service/previous_identity", HTTP_GET, @@ -1073,87 +1974,10 @@ namespace ccf .set_auto_schema() .install(); - auto get_nodes = [this](auto& args, nlohmann::json&&) { - const auto parsed_query = - http::parse_query(args.rpc_ctx->get_request_query()); - - std::string error_string; // Ignored - all params are optional - const auto host = http::get_query_value_opt( - parsed_query, "host", error_string); - const auto port = http::get_query_value_opt( - parsed_query, "port", error_string); - const auto status_str = http::get_query_value_opt( - parsed_query, "status", error_string); - - std::optional status; - if (status_str.has_value()) - { - // Convert the query argument to a JSON string, try to parse it as - // a NodeStatus, return an error if this doesn't work - try - { - status = nlohmann::json(status_str.value()).get(); - } - catch (const ccf::JsonParseError& e) - { - return ccf::make_error( - HTTP_STATUS_BAD_REQUEST, - ccf::errors::InvalidQueryParameterValue, - fmt::format( - "Query parameter '{}' is not a valid node status", - status_str.value())); - } - } - - GetNodes::Out out; - - auto nodes = args.tx.ro(this->network.nodes); - auto* current_consensus = get_consensus(); - nodes->foreach([host, port, status, &out, nodes, current_consensus]( - const NodeId& nid, const NodeInfo& ni) { - if (status.has_value() && status.value() != ni.status) - { - return true; - } - - // Match on any interface - bool is_matched = false; - for (auto const& interface : ni.rpc_interfaces) - { - const auto& [pub_host, pub_port] = - split_net_address(interface.second.published_address); - - if ( - (!host.has_value() || host.value() == pub_host) && - (!port.has_value() || port.value() == pub_port)) - { - is_matched = true; - break; - } - } - - if (!is_matched) - { - return true; - } - - bool is_primary = false; - if (current_consensus != nullptr) - { - is_primary = current_consensus->primary() == nid; - } - - out.nodes.push_back( - {nid, - ni.status, - is_primary, - ni.rpc_interfaces, - ni.node_data, - nodes->get_version_of_previous_write(nid).value_or(0)}); - return true; - }); - - return make_success(out); + auto get_nodes = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_nodes(args, std::move(json)); }; make_read_only_endpoint( "/network/nodes", @@ -1169,36 +1993,12 @@ namespace ccf "status", ccf::endpoints::OptionalParameter) .install(); - auto get_removable_nodes = [this](auto& args, nlohmann::json&&) { - GetNodes::Out out; - - auto nodes = args.tx.ro(this->network.nodes); - nodes->foreach( - [&out, nodes](const NodeId& node_id, const NodeInfo& /*ni*/) { - // Only nodes whose retire_committed status is committed can be - // safely removed, because any primary elected from here on would - // consider them retired, and would consequently not need their - // input in any quorum. We must therefore read the KV at its - // globally committed watermark, for the purpose of this RPC. Since - // this transaction does not perform a write, it is safe to do this. - auto node = nodes->get_globally_committed(node_id); - if ( - node.has_value() && node->status == ccf::NodeStatus::RETIRED && - node->retired_committed) - { - out.nodes.push_back( - {node_id, - node->status, - false /* is_primary */, - node->rpc_interfaces, - node->node_data, - nodes->get_version_of_previous_write(node_id).value_or(0)}); - } - return true; - }); - - return make_success(out); - }; + auto get_removable_nodes = + [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_removable_nodes(args, std::move(json)); + }; make_read_only_endpoint( "/network/removable_nodes", @@ -1209,57 +2009,8 @@ namespace ccf .install(); auto delete_retired_committed_node = - [this](auto& args, nlohmann::json&&) { - GetNodes::Out out; - - std::string node_id; - std::string error; - if (!get_path_param( - args.rpc_ctx->get_request_path_params(), - "node_id", - node_id, - error)) - { - return make_error( - HTTP_STATUS_BAD_REQUEST, ccf::errors::InvalidResourceName, error); - } - - auto nodes = args.tx.rw(this->network.nodes); - if (!nodes->has(node_id)) - { - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "No such node"); - } - - // A node's retirement is only complete when the - // transition of retired_committed is itself committed, - // i.e. when the next eligible primary is guaranteed to - // be aware the retirement is committed. - // As a result, the handler must check node info at the - // current committed level, rather than at the end of the - // local suffix. - // While this transaction does execute a write, it specifically - // deletes the value it reads from. It is therefore safe to - // execute on the basis of a potentially stale read-set, - // which get_globally_committed() typically produces. - auto node = nodes->get_globally_committed(node_id); - if ( - node.has_value() && node->status == ccf::NodeStatus::RETIRED && - node->retired_committed) - { - InternalTablesAccess::remove_node(args.tx, node_id); - } - else - { - return make_error( - HTTP_STATUS_BAD_REQUEST, - ccf::errors::NodeNotRetiredCommitted, - "Node is not completely retired"); - } - - return make_success(true); + [this](ccf::endpoints::EndpointContext& args, nlohmann::json&& json) { + return this->delete_retired_committed_node(args, std::move(json)); }; make_endpoint( @@ -1271,61 +2022,23 @@ namespace ccf .install(); auto get_self_signed_certificate = - [this](auto& /*args*/, nlohmann::json&&) { - return SelfSignedNodeCertificateInfo{ - this->node_operation.get_self_signed_node_certificate()}; + [this]( + ccf::endpoints::CommandEndpointContext& args, nlohmann::json&& json) { + return this->get_self_signed_certificate(args, std::move(json)); }; make_command_endpoint( "/self_signed_certificate", HTTP_GET, json_command_adapter(get_self_signed_certificate), no_auth_required) - .set_forwarding_required(endpoints::ForwardingRequired::Never) - .set_auto_schema() - .install(); - - auto get_node_info = [this](auto& args, nlohmann::json&&) { - std::string node_id; - std::string error; - if (!get_path_param( - args.rpc_ctx->get_request_path_params(), - "node_id", - node_id, - error)) - { - return make_error( - HTTP_STATUS_BAD_REQUEST, ccf::errors::InvalidResourceName, error); - } - - auto nodes = args.tx.ro(this->network.nodes); - auto info = nodes->get(node_id); - - if (!info) - { - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "Node not found"); - } - - bool is_primary = false; - auto* current_consensus = get_consensus(); - if (current_consensus != nullptr) - { - auto primary = current_consensus->primary(); - if (primary.has_value() && primary.value() == node_id) - { - is_primary = true; - } - } - auto& ni = info.value(); - return make_success(GetNode::Out{ - node_id, - ni.status, - is_primary, - ni.rpc_interfaces, - ni.node_data, - nodes->get_version_of_previous_write(node_id).value_or(0)}); + .set_forwarding_required(endpoints::ForwardingRequired::Never) + .set_auto_schema() + .install(); + + auto get_node_info = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_node_info(args, std::move(json)); }; make_read_only_endpoint( "/network/nodes/{node_id}", @@ -1335,55 +2048,10 @@ namespace ccf .set_auto_schema() .install(); - auto get_self_node = [this](auto& args, nlohmann::json&&) { - auto node_id = this->context.get_node_id(); - auto nodes = args.tx.ro(this->network.nodes); - auto info = nodes->get(node_id); - - bool is_primary = false; - auto* current_consensus = get_consensus(); - if (current_consensus != nullptr) - { - auto primary = current_consensus->primary(); - if (primary.has_value() && primary.value() == node_id) - { - is_primary = true; - } - } - - if (info.has_value()) - { - // Answers from the KV are preferred, as they are more up-to-date, - // especially status and node_data. - auto& ni = info.value(); - return make_success(GetNode::Out{ - node_id, - ni.status, - is_primary, - ni.rpc_interfaces, - ni.node_data, - nodes->get_version_of_previous_write(node_id).value_or(0)}); - } - - // If the node isn't in its KV yet, fall back to configuration - auto node_configuration_subsystem = - this->context.get_subsystem(); - if (!node_configuration_subsystem) - { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "NodeConfigurationSubsystem is not available"); - } - const auto& node_startup_config = - node_configuration_subsystem->get().node_config; - return make_success(GetNode::Out{ - node_id, - ccf::NodeStatus::PENDING, - is_primary, - node_startup_config.network.rpc_interfaces, - node_configuration_subsystem->get().node_data, - 0}); + auto get_self_node = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_self_node(args, std::move(json)); }; make_read_only_endpoint( "/network/nodes/self", @@ -1394,44 +2062,10 @@ namespace ccf .set_forwarding_required(endpoints::ForwardingRequired::Never) .install(); - auto get_primary_node = [this](auto& args, nlohmann::json&&) { - auto* current_consensus = get_consensus(); - if (current_consensus != nullptr) - { - auto primary_id = current_consensus->primary(); - if (!primary_id.has_value()) - { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "Primary unknown"); - } - - auto nodes = args.tx.ro(this->network.nodes); - auto info = nodes->get(primary_id.value()); - if (!info) - { - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "Node not found"); - } - - auto& ni = info.value(); - return make_success(GetNode::Out{ - primary_id.value(), - ni.status, - true, - ni.rpc_interfaces, - ni.node_data, - nodes->get_version_of_previous_write(primary_id.value()) - .value_or(0)}); - } - - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "No configured consensus"); + auto get_primary_node = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->get_primary_node(args, std::move(json)); }; make_read_only_endpoint( "/network/nodes/primary", @@ -1441,46 +2075,10 @@ namespace ccf .set_auto_schema() .install(); - auto head_primary = [this](auto& args) { - if (this->node_operation.can_replicate()) - { - args.rpc_ctx->set_response_status(HTTP_STATUS_OK); - } - else - { - auto* current_consensus = get_consensus(); - if (current_consensus == nullptr) - { - args.rpc_ctx->set_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "Consensus not initialised"); - return; - } - - auto primary_id = current_consensus->primary(); - if (!primary_id.has_value()) - { - args.rpc_ctx->set_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "Primary unknown"); - return; - } - - const auto address = node::get_redirect_address_for_node( - args, args.tx, primary_id.value()); - if (!address.has_value()) - { - return; - } - - args.rpc_ctx->set_response_header( - http::headers::LOCATION, - fmt::format("https://{}/node/primary", address.value())); - args.rpc_ctx->set_response_status(HTTP_STATUS_PERMANENT_REDIRECT); - } - }; + auto head_primary = + [this](ccf::endpoints::ReadOnlyEndpointContext& args) { + this->head_primary(args); + }; make_read_only_endpoint( "/primary", HTTP_HEAD, head_primary, no_auth_required) .set_forwarding_required(endpoints::ForwardingRequired::Never) @@ -1489,60 +2087,25 @@ namespace ccf "Redirect to the current primary node.") .install(); - auto get_primary = [this](auto& args) { - if (this->node_operation.can_replicate()) - { - args.rpc_ctx->set_response_status(HTTP_STATUS_OK); - return; - } - - args.rpc_ctx->set_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "Node is not primary"); + auto get_primary = [this](ccf::endpoints::ReadOnlyEndpointContext& args) { + this->get_primary(args); }; make_read_only_endpoint( "/primary", HTTP_GET, get_primary, no_auth_required) .set_forwarding_required(endpoints::ForwardingRequired::Never) .install(); - auto get_backup = [this](auto& args) { - if (!this->node_operation.can_replicate()) - { - args.rpc_ctx->set_response_status(HTTP_STATUS_OK); - return; - } - - args.rpc_ctx->set_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "Node is not backup"); + auto get_backup = [this](ccf::endpoints::ReadOnlyEndpointContext& args) { + this->get_backup(args); }; make_read_only_endpoint("/backup", HTTP_GET, get_backup, no_auth_required) .set_forwarding_required(endpoints::ForwardingRequired::Never) .install(); - auto consensus_config = [this](auto& /*args*/, nlohmann::json&&) { - // Query node for configurations, separate current from pending - auto* current_consensus = get_consensus(); - if (current_consensus != nullptr) - { - auto cfg = current_consensus->get_latest_configuration(); - ConsensusConfig cc; - for (auto& [nid, ninfo] : cfg) - { - cc.emplace( - nid.value(), - ConsensusNodeConfig{ - fmt::format("{}:{}", ninfo.hostname, ninfo.port)}); - } - return make_success(cc); - } - - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "No configured consensus"); + auto consensus_config = [this]( + ccf::endpoints::CommandEndpointContext& args, + nlohmann::json&& json) { + return this->consensus_config(args, std::move(json)); }; make_command_endpoint( @@ -1554,18 +2117,10 @@ namespace ccf .set_auto_schema() .install(); - auto consensus_state = [this](auto& /*args*/, nlohmann::json&&) { - auto* current_consensus = get_consensus(); - if (current_consensus != nullptr) - { - return make_success( - ConsensusConfigDetails{current_consensus->get_details()}); - } - - return make_error( - HTTP_STATUS_NOT_FOUND, - ccf::errors::ResourceNotFound, - "No configured consensus"); + auto consensus_state = [this]( + ccf::endpoints::CommandEndpointContext& args, + nlohmann::json&& json) { + return this->consensus_state(args, std::move(json)); }; make_command_endpoint( @@ -1577,14 +2132,8 @@ namespace ccf .set_auto_schema() .install(); - auto node_metrics = [this](auto& args) { - NodeMetrics nm; - nm.sessions = node_operation.get_session_metrics(); - - args.rpc_ctx->set_response_status(HTTP_STATUS_OK); - args.rpc_ctx->set_response_header( - http::headers::CONTENT_TYPE, http::headervalues::contenttype::JSON); - args.rpc_ctx->set_response_body(nlohmann::json(nm).dump()); + auto node_metrics = [this](ccf::endpoints::CommandEndpointContext& args) { + this->node_metrics(args); }; make_command_endpoint( @@ -1593,28 +2142,10 @@ namespace ccf .set_auto_schema() .install(); - auto js_metrics = [this](auto& args, nlohmann::json&&) { - auto bytecode_map = args.tx.ro(this->network.modules_quickjs_bytecode); - auto version_val = args.tx.ro(this->network.modules_quickjs_version); - uint64_t bytecode_size = 0; - bytecode_map->foreach( - [&bytecode_size](const auto&, const auto& bytecode) { - bytecode_size += bytecode.size(); - return true; - }); - auto js_engine_map = args.tx.ro(this->network.js_engine); - JavaScriptMetrics m; - m.bytecode_size = bytecode_size; - m.bytecode_used = - version_val->get() == std::string(ccf::quickjs_version); - - auto options = js_engine_map->get().value_or(ccf::JSRuntimeOptions{}); - m.max_stack_size = options.max_stack_bytes; - m.max_heap_size = options.max_heap_bytes; - m.max_execution_time = options.max_execution_time_ms; - m.max_cached_interpreters = options.max_cached_interpreters; - - return m; + auto js_metrics = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->js_metrics(args, std::move(json)); }; make_read_only_endpoint( @@ -1625,13 +2156,10 @@ namespace ccf .set_auto_schema() .install(); - auto version = [](auto&, nlohmann::json&&) { - GetVersion::Out result; - result.ccf_version = ccf::ccf_version; - result.quickjs_version = ccf::quickjs_version; - result.unsafe = false; - - return make_success(result); + auto version = []( + ccf::endpoints::CommandEndpointContext& args, + nlohmann::json&& json) { + return NodeEndpoints::version(args, std::move(json)); }; make_command_endpoint( @@ -1640,258 +2168,20 @@ namespace ccf .set_auto_schema() .install(); - auto create = [this](auto& ctx, nlohmann::json&& params) { - LOG_INFO_FMT("Processing create RPC"); - - bool recovering = node_operation.is_reading_public_ledger(); - - // This endpoint can only be called once, directly from the starting - // node for the genesis or end of public recovery transaction to - // initialise the service - if (!node_operation.is_in_initialised_state() && !recovering) - { - return make_error( - HTTP_STATUS_FORBIDDEN, - ccf::errors::InternalError, - "Node is not in initial state."); - } - - const auto in = params.get(); - - if (InternalTablesAccess::is_service_created(ctx.tx, in.service_cert)) - { - return make_error( - HTTP_STATUS_FORBIDDEN, - ccf::errors::InternalError, - "Service is already created."); - } - - InternalTablesAccess::create_service( - ctx.tx, in.service_cert, in.create_txid, in.service_data, recovering); - - if (recovering) - { - // Recovery starts with a fresh consensus configuration, so previous - // service nodes can be removed immediately. - InternalTablesAccess::remove_previous_service_nodes(ctx.tx); - } - - // Genesis transaction (i.e. not after recovery) - if (in.genesis_info.has_value()) - { - // Note that it is acceptable to start a network without any member - // having a recovery share. The service will check that at least one - // recovery member is added before the service is opened. - for (const auto& info : in.genesis_info->members) - { - InternalTablesAccess::add_member(ctx.tx, info); - } - - InternalTablesAccess::init_configuration( - ctx.tx, in.genesis_info->service_configuration); - InternalTablesAccess::set_constitution( - ctx.tx, in.genesis_info->constitution); - } - else - { - // On recovery, force a new ledger chunk - auto* tx_ = static_cast(&ctx.tx); - if (tx_ == nullptr) - { - throw std::logic_error("Could not cast tx to CommittableTx"); - } - tx_->set_tx_flag( - ccf::kv::CommittableTx::TxFlag::LEDGER_CHUNK_BEFORE_THIS_TX); - } - - auto endorsed_certificates = - ctx.tx.rw(network.node_endorsed_certificates); - endorsed_certificates->put(in.node_id, in.node_endorsed_certificate); - - NodeInfo node_info = { - in.node_info_network, - {in.quote_info}, - in.public_encryption_key, - NodeStatus::TRUSTED, - std::nullopt, - in.measurement.hex_str(), - in.certificate_signing_request, - in.public_key, - in.node_data}; - InternalTablesAccess::add_node(ctx.tx, in.node_id, node_info); - - if (in.sealing_recovery_data.has_value()) - { - const auto& [sealing_keys, sealing_recovery_name] = - in.sealing_recovery_data.value(); - auto* sealed_recovery_keys = ctx.tx.template rw( - Tables::SEALED_RECOVERY_KEYS); - sealed_recovery_keys->put(in.node_id, sealing_keys); - - auto* local_sealing_node_id_map = - ctx.tx.template rw( - Tables::SEALING_RECOVERY_NAMES); - local_sealing_node_id_map->put(sealing_recovery_name, in.node_id); - } - - node_operation.shuffle_sealed_shares(ctx.tx); - - if ( - in.quote_info.format != QuoteFormat::amd_sev_snp_v1 || - !in.snp_uvm_endorsements.has_value()) - { - // For improved serviceability on SNP, do not record trusted - // measurements if UVM endorsements are available - InternalTablesAccess::trust_node_measurement( - ctx.tx, in.measurement, in.quote_info.format); - } - - switch (in.quote_info.format) - { - case QuoteFormat::insecure_virtual: - { - auto host_data = AttestationProvider::get_host_data(in.quote_info); - if (host_data.has_value()) - { - InternalTablesAccess::trust_node_virtual_host_data( - ctx.tx, host_data.value()); - } - else - { - LOG_FAIL_FMT("Unable to extract host data from virtual quote"); - } - break; - } - - case QuoteFormat::amd_sev_snp_v1: - { - auto host_data = - AttestationProvider::get_host_data(in.quote_info).value(); - InternalTablesAccess::trust_node_snp_host_data( - ctx.tx, host_data, in.snp_security_policy); - - InternalTablesAccess::trust_node_uvm_endorsements( - ctx.tx, in.snp_uvm_endorsements, recovering); - - auto attestation = - AttestationProvider::get_snp_attestation_report(in.quote_info) - .value(); - InternalTablesAccess::trust_node_snp_tcb_version( - ctx.tx, attestation, recovering); - break; - } - case QuoteFormat::oe_sgx_v1: - { - break; - } - } - - std::optional digest = - ccf::get_create_tx_claims_digest(ctx.tx); - if (digest.has_value()) - { - auto digest_value = digest.value(); - ctx.rpc_ctx->set_claims_digest(std::move(digest_value)); - } - - this->node_operation.recovery_decision_protocol().reset_state(ctx.tx); - this->node_operation.recovery_decision_protocol().try_start( - ctx.tx, recovering); - - LOG_INFO_FMT("Created service"); - return make_success(true); - }; + auto create = + [this](ccf::endpoints::EndpointContext& ctx, nlohmann::json&& params) { + return this->create(ctx, std::move(params)); + }; make_endpoint( "/create", HTTP_POST, json_adapter(create), {self_cert_auth_policy}) .set_openapi_hidden(true) .install(); // Only called from node. See node_state.h. - auto refresh_jwt_keys = [this](auto& ctx, nlohmann::json&& body) { - // All errors are server errors since the client is the server. - - auto* current_consensus = get_consensus(); - auto primary_id = current_consensus->primary(); - if (!primary_id.has_value()) - { - LOG_FAIL_FMT("JWT key auto-refresh: primary unknown"); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "Primary is unknown"); - } - - const auto& sig_auth_ident = - ctx.template get_caller(); - if (primary_id.value() != sig_auth_ident.node_id) - { - LOG_FAIL_FMT( - "JWT key auto-refresh: request does not originate from primary"); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "Request does not originate from primary."); - } - - SetJwtPublicSigningKeys parsed; - try - { - parsed = body.get(); - } - catch (const ccf::JsonParseError& e) - { - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "Unable to parse body."); - } - - auto issuers = ctx.tx.ro(this->network.jwt_issuers); - auto issuer_metadata_ = issuers->get(parsed.issuer); - if (!issuer_metadata_.has_value()) - { - LOG_FAIL_FMT( - "JWT key auto-refresh: {} is not a valid issuer", parsed.issuer); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - fmt::format("{} is not a valid issuer.", parsed.issuer)); - } - auto& issuer_metadata = issuer_metadata_.value(); - - if (!issuer_metadata.auto_refresh) - { - LOG_FAIL_FMT( - "JWT key auto-refresh: {} does not have auto_refresh enabled", - parsed.issuer); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - fmt::format( - "{} does not have auto_refresh enabled.", parsed.issuer)); - } - - if (!set_jwt_public_signing_keys( - ctx.tx, - "", - parsed.issuer, - issuer_metadata, - parsed.jwks)) - { - LOG_FAIL_FMT( - "JWT key auto-refresh: error while storing signing keys for issuer " - "{}", - parsed.issuer); - return make_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - fmt::format( - "Error while storing signing keys for issuer {}.", - parsed.issuer)); - } - - return make_success(true); - }; + auto refresh_jwt_keys = + [this](ccf::endpoints::EndpointContext& ctx, nlohmann::json&& body) { + return this->refresh_jwt_keys(ctx, std::move(body)); + }; make_endpoint( "/jwt_keys/refresh", HTTP_POST, @@ -1900,15 +2190,11 @@ namespace ccf .set_openapi_hidden(true) .install(); - auto get_jwt_metrics = - [this](auto& /*args*/, const nlohmann::json& /*params*/) { - JWTRefreshMetrics metrics; - { - ccf::ds::MutexGuard guard(jwt_refresh_metrics_lock); - metrics = jwt_refresh_metrics; - } - return make_success(metrics); - }; + auto get_jwt_metrics = [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + const nlohmann::json& params) { + return this->get_jwt_metrics(args, params); + }; make_read_only_endpoint( "/jwt_keys/refresh/metrics", HTTP_GET, @@ -1917,10 +2203,11 @@ namespace ccf .set_auto_schema() .install(); - auto service_config_handler = - [this](auto& args, const nlohmann::json& /*params*/) { - return make_success(args.tx.ro(network.config)->get()); - }; + auto service_config_handler = [this]( + ccf::endpoints::EndpointContext& args, + const nlohmann::json& params) { + return this->service_config_handler(args, params); + }; make_endpoint( "/service/configuration", HTTP_GET, @@ -1931,9 +2218,9 @@ namespace ccf .install(); auto list_indexing_strategies = [this]( - auto& /*args*/, - const nlohmann::json& /*params*/) { - return make_success(this->context.get_indexing_strategies().describe()); + ccf::endpoints::EndpointContext& args, + const nlohmann::json& params) { + return this->list_indexing_strategies(args, params); }; make_endpoint( @@ -1946,28 +2233,7 @@ namespace ccf .install(); auto get_ready_app = [this](ccf::endpoints::CommandEndpointContext& ctx) { - auto node_configuration_subsystem = - this->context.get_subsystem(); - if (!node_configuration_subsystem) - { - ctx.rpc_ctx->set_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "NodeConfigurationSubsystem is not available"); - return; - } - if ( - !node_configuration_subsystem->has_received_stop_notice() && - this->node_operation.is_part_of_network() && - this->node_operation.is_user_frontend_open()) - { - ctx.rpc_ctx->set_response_status(HTTP_STATUS_NO_CONTENT); - } - else - { - ctx.rpc_ctx->set_response_status(HTTP_STATUS_SERVICE_UNAVAILABLE); - } - return; + this->get_ready_app(ctx); }; make_command_endpoint( "/ready/app", HTTP_GET, get_ready_app, no_auth_required) @@ -1979,28 +2245,7 @@ namespace ccf .install(); auto get_ready_gov = [this](ccf::endpoints::CommandEndpointContext& ctx) { - auto node_configuration_subsystem = - this->context.get_subsystem(); - if (!node_configuration_subsystem) - { - ctx.rpc_ctx->set_error( - HTTP_STATUS_INTERNAL_SERVER_ERROR, - ccf::errors::InternalError, - "NodeConfigurationSubsystem is not available"); - return; - } - if ( - !node_configuration_subsystem->has_received_stop_notice() && - this->node_operation.is_accessible_to_members() && - this->node_operation.is_member_frontend_open()) - { - ctx.rpc_ctx->set_response_status(HTTP_STATUS_NO_CONTENT); - } - else - { - ctx.rpc_ctx->set_response_status(HTTP_STATUS_SERVICE_UNAVAILABLE); - } - return; + this->get_ready_gov(ctx); }; make_command_endpoint( "/ready/gov", HTTP_GET, get_ready_gov, no_auth_required) @@ -2011,13 +2256,10 @@ namespace ccf .set_forwarding_required(endpoints::ForwardingRequired::Never) .install(); - auto create_snapshot = [this](auto& args, nlohmann::json&&) { - auto* snapshot_create = args.tx.template rw( - ccf::Tables::SNAPSHOT_CREATE); - snapshot_create->touch(); - this->node_operation.trigger_snapshot(args.tx); - return make_success(); - }; + auto create_snapshot = + [this](ccf::endpoints::EndpointContext& args, nlohmann::json&& json) { + return this->create_snapshot(args, std::move(json)); + }; make_endpoint( "/snapshot:create", HTTP_POST, @@ -2032,14 +2274,12 @@ namespace ccf ccf::node::init_file_serving_handlers(*this, context); - auto historical_cache_info = [this]( - [[maybe_unused]] auto& args, - [[maybe_unused]] nlohmann::json&&) { - GetHistoricalCacheInfo::Out result{}; - result.estimated_size = - this->context.get_historical_state().get_estimated_store_cache_size(); - return make_success(result); - }; + auto historical_cache_info = + [this]( + ccf::endpoints::ReadOnlyEndpointContext& args, + nlohmann::json&& json) { + return this->historical_cache_info(args, std::move(json)); + }; make_read_only_endpoint( "/historical_cache", HTTP_GET,