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¶
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¶
function HasQueueOwnership¶
function StopEngine¶
Stops the processing engine and joins its in-flight subtask threads so no result write or publish outlives this call.
function setMirrorResultCallback¶
Set callback for mirroring results from other nodes
function setBitswap¶
Set bitswap instance for data availability checks
function SetMembershipFilter¶
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¶
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¶
Return: Progress percentage (0.0 to 100.0)
Get current processing progress
Updated on 2026-10-06 at 13:34:20 +0000