sgns::processing::ProcessingServiceImpl¶
#include <processing_service.hpp>
Inherits from std::enable_shared_from_this< ProcessingServiceImpl >
Public Classes¶
| Name | |
|---|---|
| struct | ProcessingStatus |
Public Types¶
| Name | |
|---|---|
| enum class uint8_t | Status |
Public Functions¶
| Name | |
|---|---|
| ProcessingServiceImpl(std::shared_ptr< ipfs_pubsub::GossipPubSub > gossipPubSub, size_t maximalNodesCount, std::shared_ptr< SubTaskEnqueuer > subTaskEnqueuer, std::shared_ptr< SubTaskResultStorage > subTaskResultStorage, std::shared_ptr< ProcessingCore > processingCore, std::function< void(const std::string &subTaskQueueId, const SGProcessing::TaskResult &taskresult)> userCallbackSuccess, std::function< void(const std::string &subTaskQueueId)> userCallbackError, std::string node_address) Constructs a processing service with user callbacks. |
|
| ~ProcessingServiceImpl() | |
| void | StartProcessing(const std::string & processingGridChannelId) |
| void | StopProcessing() |
| size_t | GetProcessingNodesCount() const |
| void | SetChannelListRequestTimeout(boost::posix_time::time_duration channelListRequestTimeout) |
| ProcessingStatus | GetProcessingStatus() const |
| void | setMirrorResultCallback(std::function< void(const std::string &)> callback) Set callback for mirroring processing results. When set, results with mirror_result=true trigger a fetch. |
| void | setBitswap(std::shared_ptr< sgns::ipfs_bitswap::Bitswap > bitswap) Set bitswap instance propagated to all processing nodes for data availability checks. |
| void | SetMembershipFilter(sgns::networkregistry::MembershipFilter filter) Set membership filter enforced at every processing-path message handler (grid, results, and processing queue channels). Empty filter = public pass-through. The filter is snapshotted BEFORE node creation at both creation sites and passed INTO ProcessingNode::New, where it is installed BEFORE any subscription goes live – before the queue-channel Listen() and before the results-channel CreateResultsChannel/ConnectToSubTaskQueue – so there is no creation-time enrollment window (T-15-13-06, delivered by the pre-subscription install; CR-G02a closed). Set-time propagation (this call) refreshes existing nodes. |
| void | SetGossipSigningKey(std::shared_ptr< const libp2p::crypto::KeyPair > key) Set the gossip host keypair used to SEAL private-network processing-channel publishes and authenticate inbound ones (CR-G01). Under a set membership filter every publish is sealed (sgns::base::SealGossipPayload) and every inbound message must open a verifiable envelope whose embedded key derives the from-field PeerId (sgns::base::OpenGossipPayload) BEFORE the membership predicate runs. Filter set + no key = publishes fail closed. Propagates to all existing processing nodes and, symmetric with SetMembershipFilter, is snapshotted BEFORE node creation and passed INTO ProcessingNode::New to be installed BEFORE any subscription goes live. No filter -> raw, byte-identical. |
Public Types Documentation¶
enum Status¶
| Enumerator | Value | Description |
|---|---|---|
| DISABLED | ||
| IDLE | ||
| PROCESSING |
Public Functions Documentation¶
function ProcessingServiceImpl¶
ProcessingServiceImpl(
std::shared_ptr< ipfs_pubsub::GossipPubSub > gossipPubSub,
size_t maximalNodesCount,
std::shared_ptr< SubTaskEnqueuer > subTaskEnqueuer,
std::shared_ptr< SubTaskResultStorage > subTaskResultStorage,
std::shared_ptr< ProcessingCore > processingCore,
std::function< void(const std::string &subTaskQueueId, const SGProcessing::TaskResult &taskresult)> userCallbackSuccess,
std::function< void(const std::string &subTaskQueueId)> userCallbackError,
std::string node_address
)
Constructs a processing service with user callbacks.
Parameters:
- gossipPubSub PubSub service.
- maximalNodesCount Max number of processing nodes handled by the service.
- subTaskEnqueuer Subtask enqueuer used to dispatch tasks.
- subTaskResultStorage Storage for subtask results.
- processingCore Processing core used to execute subtasks.
- userCallbackSuccess Callback invoked on successful task completion.
- userCallbackError Callback invoked on task error.
- node_address Local node address used in coordination.
function ~ProcessingServiceImpl¶
function StartProcessing¶
function StopProcessing¶
function GetProcessingNodesCount¶
function SetChannelListRequestTimeout¶
function GetProcessingStatus¶
function setMirrorResultCallback¶
Set callback for mirroring processing results. When set, results with mirror_result=true trigger a fetch.
function setBitswap¶
Set bitswap instance propagated to all processing nodes for data availability checks.
function SetMembershipFilter¶
Set membership filter enforced at every processing-path message handler (grid, results, and processing queue channels). Empty filter = public pass-through. The filter is snapshotted BEFORE node creation at both creation sites and passed INTO ProcessingNode::New, where it is installed BEFORE any subscription goes live – before the queue-channel Listen() and before the results-channel CreateResultsChannel/ConnectToSubTaskQueue – so there is no creation-time enrollment window (T-15-13-06, delivered by the pre-subscription install; CR-G02a closed). Set-time propagation (this call) refreshes existing nodes.
function SetGossipSigningKey¶
Set the gossip host keypair used to SEAL private-network processing-channel publishes and authenticate inbound ones (CR-G01). Under a set membership filter every publish is sealed (sgns::base::SealGossipPayload) and every inbound message must open a verifiable envelope whose embedded key derives the from-field PeerId (sgns::base::OpenGossipPayload) BEFORE the membership predicate runs. Filter set + no key = publishes fail closed. Propagates to all existing processing nodes and, symmetric with SetMembershipFilter, is snapshotted BEFORE node creation and passed INTO ProcessingNode::New to be installed BEFORE any subscription goes live. No filter -> raw, byte-identical.
Updated on 2026-10-07 at 18:58:59 +0000