Skip to content

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

~ProcessingServiceImpl()

function StartProcessing

void StartProcessing(
    const std::string & processingGridChannelId
)

function StopProcessing

void StopProcessing()

function GetProcessingNodesCount

size_t GetProcessingNodesCount() const

function SetChannelListRequestTimeout

void SetChannelListRequestTimeout(
    boost::posix_time::time_duration channelListRequestTimeout
)

function GetProcessingStatus

ProcessingStatus GetProcessingStatus() const

function setMirrorResultCallback

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.

function setBitswap

void setBitswap(
    std::shared_ptr< sgns::ipfs_bitswap::Bitswap > bitswap
)

Set bitswap instance propagated to all processing nodes for data availability checks.

function SetMembershipFilter

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.

function SetGossipSigningKey

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.


Updated on 2026-10-07 at 18:58:59 +0000