Skip to content

sgns::processing::ProcessingNode

Node for distributed computation. More...

#include <processing_node.hpp>

Inherits from std::enable_shared_from_this< ProcessingNode >

Public Functions

Name
std::shared_ptr< ProcessingNode > New(std::shared_ptr< ipfs_pubsub::GossipPubSub > gossipPubSub, std::shared_ptr< SubTaskResultStorage > subTaskResultStorage, std::shared_ptr< ProcessingCore > processingCore, std::function< void(const SGProcessing::TaskResult &)> taskResultProcessingSink, std::function< void(const std::string &)> processingErrorSink, std::function< void(void)> processingDoneSink, std::string node_id, const std::string & processingQueueChannelId, std::list< SGProcessing::SubTask > subTasks ={}, std::chrono::milliseconds msSubscriptionWaitingDuration =std::chrono::milliseconds(2000), std::chrono::seconds ttl =std::chrono::minutes(2), sgns::networkregistry::MembershipFilter membershipFilter ={}, std::shared_ptr< const libp2p::crypto::KeyPair > gossipSigningKey ={})
Creates a processing node instance.
~ProcessingNode()
bool HasQueueOwnership() const
void StopEngine()
Stops the processing engine and joins its in-flight subtask threads so no result write or publish outlives this call.
void setMirrorResultCallback(std::function< void(const std::string &)> callback)
void setBitswap(std::shared_ptr< sgns::ipfs_bitswap::Bitswap > bitswap)
void SetMembershipFilter(sgns::networkregistry::MembershipFilter filter)
void SetGossipSigningKey(std::shared_ptr< const libp2p::crypto::KeyPair > key)
float GetProgress() const

Detailed Description

class sgns::processing::ProcessingNode;

Node for distributed computation.

Coordinates subtask queue ownership, processing, and result publication.

Public Functions Documentation

function New

static std::shared_ptr< ProcessingNode > New(
    std::shared_ptr< ipfs_pubsub::GossipPubSub > gossipPubSub,
    std::shared_ptr< SubTaskResultStorage > subTaskResultStorage,
    std::shared_ptr< ProcessingCore > processingCore,
    std::function< void(const SGProcessing::TaskResult &)> taskResultProcessingSink,
    std::function< void(const std::string &)> processingErrorSink,
    std::function< void(void)> processingDoneSink,
    std::string node_id,
    const std::string & processingQueueChannelId,
    std::list< SGProcessing::SubTask > subTasks ={},
    std::chrono::milliseconds msSubscriptionWaitingDuration =std::chrono::milliseconds(2000),
    std::chrono::seconds ttl =std::chrono::minutes(2),
    sgns::networkregistry::MembershipFilter membershipFilter ={},
    std::shared_ptr< const libp2p::crypto::KeyPair > gossipSigningKey ={}
)

Creates a processing node instance.

Parameters:

  • gossipPubSub PubSub service for queue coordination.
  • subTaskResultStorage Storage for subtask results.
  • processingCore Processing core to execute subtasks.
  • taskResultProcessingSink Callback for task result processing.
  • processingErrorSink Callback for processing errors.
  • processingDoneSink Callback when processing is done.
  • node_id Identifier of the processing node.
  • processingQueueChannelId Queue channel identifier.
  • subTasks Optional initial subtask list.
  • msSubscriptionWaitingDuration Wait duration for queue subscription.
  • ttl Time-to-live for node ownership.
  • membershipFilter Membership gate installed BEFORE any subscription goes live: on the queue channel before Listen() and on the results accessor before CreateResultsChannel/ConnectToSubTaskQueue (CR-G02a – the creation-site snapshot eliminates the enrollment window). Empty (public node) -> nothing installed, byte-identical.
  • gossipSigningKey Gossip host keypair installed beside the filter at the same pre-subscription points (CR-G01 symmetry: sealed publishes + authenticated ingest from the first message).

function ~ProcessingNode

~ProcessingNode()

function HasQueueOwnership

bool HasQueueOwnership() const

function StopEngine

void StopEngine()

Stops the processing engine and joins its in-flight subtask threads so no result write or publish outlives this call.

function setMirrorResultCallback

void setMirrorResultCallback(
    std::function< void(const std::string &)> callback
)

Set callback for mirroring results from other nodes

function setBitswap

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

Set bitswap instance for data availability checks

function SetMembershipFilter

void SetMembershipFilter(
    sgns::networkregistry::MembershipFilter filter
)

Set membership filter forwarded to this node's results channel (queue accessor) and processing queue channel — gates non-member senders at both handlers.

function SetGossipSigningKey

void SetGossipSigningKey(
    std::shared_ptr< const libp2p::crypto::KeyPair > key
)

Set gossip signing key forwarded to this node's results channel and processing queue channel — seals private-network publishes and authenticates inbound envelopes at both handlers (CR-G01).

function GetProgress

float GetProgress() const

Return: Progress percentage (0.0 to 100.0)

Get current processing progress


Updated on 2026-10-06 at 13:34:20 +0000