cmvr-es/cmvr-es/service/grpc/include/grpc_command_transaction.h
xtkuang f4be2ffaaa feat(safety): unify device admission and recovery
Add the DeviceManager-owned safety coordinator, shared sensor/control policies, command ledger, service guards, generalized StopAll, and RecoverSafetyState. Preserve device-side hardware checks and AUBO hardware E-stop release reconciliation while keeping software E-stop independently latched.
2026-08-17 08:34:44 +08:00

212 lines
7.9 KiB
C++

#pragma once
#include <functional>
#include <memory>
#include <optional>
#include <string>
#include <grpcpp/server_context.h>
#include <grpcpp/support/status.h>
#include <google/protobuf/message.h>
#include "cmvr/api/common.pb.h"
#include "manager/safety/include/safety_coordinator.h"
#include "service/grpc/include/grpc_security.h"
namespace cmvr::service {
grpc::Status grpcStatusForSafetyReason(
safety::SafetyReason reason,
const std::string& detail = {});
struct GrpcStreamingSafetyOpen {
std::string full_method_name;
std::string device_id;
std::string session_id;
std::string expected_service_instance_id;
std::optional<std::uint64_t> expected_device_generation;
std::uint64_t authority_generation{0};
safety::SafetyClock::time_point deadline{
safety::SafetyClock::time_point::max()};
};
// Binds a long-lived control stream to one coordinator permit. Stream
// protocols retain their own sequence, watchdog, and control-lease rules;
// this object owns the safety epoch/device-generation checks shared by all of
// them. It deliberately does not use the unary idempotency ledger.
class GrpcStreamingSafetySession final {
public:
GrpcStreamingSafetySession(
safety::SafetyCoordinator& coordinator,
const GrpcRequestContext& request_context,
GrpcStreamingSafetyOpen open);
GrpcStreamingSafetySession(
GrpcStreamingSafetySession&&) noexcept = default;
GrpcStreamingSafetySession& operator=(
GrpcStreamingSafetySession&&) noexcept = default;
GrpcStreamingSafetySession(
const GrpcStreamingSafetySession&) = delete;
GrpcStreamingSafetySession& operator=(
const GrpcStreamingSafetySession&) = delete;
bool admitted() const noexcept { return permit_.has_value(); }
const grpc::Status& status() const noexcept { return status_; }
const safety::AdmissionDecision& admissionDecision() const noexcept
{
return admission_decision_;
}
bool revalidate();
safety::DispatchGuard beginDispatch();
std::uint64_t safetyEpoch() const noexcept;
std::uint64_t deviceGeneration() const noexcept;
std::uint64_t authorityGeneration() const noexcept;
private:
void reject_(safety::SafetyReason reason, std::string detail);
safety::SafetyCoordinator* coordinator_{nullptr};
std::optional<safety::AdmissionPermit> permit_;
safety::AdmissionDecision admission_decision_;
grpc::Status status_;
};
// Owns one unary command from identity reservation through the final hardware
// dispatch fence. Legacy/Shadow calls without a command ID still use admission,
// but deliberately remain outside the idempotency ledger for wire compatibility.
class GrpcCommandTransaction final {
public:
GrpcCommandTransaction(
safety::SafetyCoordinator& coordinator,
GrpcRequestContext request_context,
GrpcMethodPolicy method_policy,
const google::protobuf::Message& request,
google::protobuf::Message& response);
~GrpcCommandTransaction() noexcept;
GrpcCommandTransaction(GrpcCommandTransaction&& other) noexcept;
GrpcCommandTransaction& operator=(
GrpcCommandTransaction&& other) noexcept;
GrpcCommandTransaction(const GrpcCommandTransaction&) = delete;
GrpcCommandTransaction& operator=(const GrpcCommandTransaction&) = delete;
bool shouldExecute() const noexcept { return should_execute_; }
const grpc::Status& status() const noexcept { return status_; }
const safety::AdmissionDecision& admissionDecision() const noexcept
{
return admission_decision_;
}
// Must be called immediately before the first driver/SDK mutation. The
// returned guard remains owned by this transaction until finish().
bool beginDispatch();
// Long-running unary commands may submit more than one hardware command.
// Revalidate the original permit between submissions, then hold the
// returned guard only around one driver/SDK mutation.
bool revalidate();
safety::DispatchGuard beginScopedDispatch();
// Internal mitigation for a command-owned activity. This obtains a fresh
// Stop-lane permit, so an expired/revoked Actuate permit cannot suppress a
// physical stop.
safety::DispatchGuard beginSafetyStopDispatch();
const grpc::Status& dispatchStatus() const noexcept
{
return dispatch_status_;
}
grpc::Status finish(
grpc::Status operation_status,
safety::SafetyReason reason = safety::SafetyReason::None,
std::optional<safety::CommandLifecycle> lifecycle = std::nullopt);
grpc::Status finishException(std::string detail) noexcept;
const std::string& deviceId() const noexcept { return device_id_; }
const std::string& commandId() const noexcept { return command_id_; }
std::uint64_t safetyEpoch() const noexcept
{
return admission_decision_.safety_epoch;
}
std::uint64_t deviceGeneration() const noexcept
{
return admission_decision_.device_generation;
}
private:
void initialize_(const google::protobuf::Message& request);
void rejectBeforeDispatch_(
safety::SafetyReason reason,
std::string detail,
grpc::Status status,
bool complete_reserved_record);
bool restoreOutcome_(const safety::CommandOutcome& outcome);
bool completeLedger_(
safety::CommandLifecycle lifecycle,
safety::SafetyReason reason,
const std::string& detail,
bool hardware_submission_possible) noexcept;
void populateFeedback_(
bool success,
safety::SafetyReason reason,
safety::CommandLifecycle lifecycle,
const std::string& detail);
void abandon_() noexcept;
safety::SafetyCoordinator* coordinator_{nullptr};
GrpcRequestContext request_context_;
GrpcMethodPolicy method_policy_;
google::protobuf::Message* response_{nullptr};
safety::CommandLedger::Ticket ledger_ticket_;
std::optional<safety::AdmissionPermit> permit_;
std::optional<safety::DispatchGuard> dispatch_guard_;
safety::AdmissionDecision admission_decision_;
grpc::Status status_;
grpc::Status dispatch_status_;
std::string device_id_;
std::string command_id_;
std::string payload_hash_;
safety::SafetyClock::time_point deadline_{
safety::SafetyClock::time_point::max()};
bool should_execute_{false};
bool owns_ledger_record_{false};
bool dispatch_started_{false};
bool completed_{false};
std::uint64_t safety_stop_sequence_{0};
};
using GrpcUnaryCommandOperation =
std::function<grpc::Status(GrpcCommandTransaction&)>;
grpc::Status executeRegisteredGrpcCommand(
const std::shared_ptr<GrpcSecurityGateway>& gateway,
grpc::ServerContext* server_context,
safety::SafetyCoordinator& coordinator,
const std::string& full_method_name,
const google::protobuf::Message* request,
google::protobuf::Message* response,
GrpcUnaryCommandOperation operation);
// For a legacy RPC whose validated request selects one of several fixed
// server-side intents. Gateway authorization still uses the registered method
// policy; the supplied policy only narrows Coordinator admission after the
// server has parsed the request (for example PTZ START versus STOP).
grpc::Status executeServerDerivedGrpcCommand(
const std::shared_ptr<GrpcSecurityGateway>& gateway,
grpc::ServerContext* server_context,
safety::SafetyCoordinator& coordinator,
const std::string& full_method_name,
GrpcMethodPolicy effective_policy,
const google::protobuf::Message* request,
google::protobuf::Message* response,
GrpcUnaryCommandOperation operation);
// Stable across processes and protobuf map iteration order. Transport identity
// fields and client timestamps are excluded; device generation remains part of
// the semantic payload.
std::string deterministicGrpcPayloadHash(
const std::string& full_method_name,
const google::protobuf::Message& request);
} // namespace cmvr::service