Skip to content

processing/processing_validation_core.cpp

Source file of core implementation of processing task results validation. More...

Namespaces

Name
sgns
sgns::processing

Functions

Name
OUTCOME_CPP_DEFINE_CATEGORY_3(sgns::processing , ProcessingValidationCore::Error , e )

Detailed Description

Source file of core implementation of processing task results validation.

Date: 2022-05-08 creativeid00

Note: This was mostly rewritten by Henrique A. Klein ([email protected]) and Justin Church ([email protected])

Functions Documentation

function OUTCOME_CPP_DEFINE_CATEGORY_3

OUTCOME_CPP_DEFINE_CATEGORY_3(
    sgns::processing ,
    ProcessingValidationCore::Error ,
    e 
)

Source code

#include <optional>
#include <unordered_set>
#include "processing_validation_core.hpp"
#include "processing/processing_subtask_queue.hpp"
#include "processing/processing_subtask_queue_channel.hpp"
#include "util/diff_utils.hpp"
#include "base/hexutil.hpp"

OUTCOME_CPP_DEFINE_CATEGORY_3( sgns::processing, ProcessingValidationCore::Error, e )
{
    using ValidationError = sgns::processing::ProcessingValidationCore::Error;
    switch ( e )
    {
        case ValidationError::NO_RESULTS_FOR_SUBTASK:
            return "Subtask was finalized with no results";
        case ValidationError::WRONG_RESULT_HASHES_LENGTH:
            return "The hashes length doesn't match the chunks to process length";
        case ValidationError::DUPLICATE_CHUNK_RESULT_HASH:
            return "A duplicate chunk result hash was found";
        case ValidationError::EMPTY_CHUNK_RESULT_HASH:
            return "Empty chunk result hash was found";
        case ValidationError::MISSING_CHUNK_RESULT:
            return "Missing chunk result found";
        case ValidationError::INVALID_CHUNK_RESULT_HASH:
            return "The chunk result hash is invalid";
        case ValidationError::SUBTASK_ID_MISMATCH:
            return "The subtask id doesn't match the result id";
        case ValidationError::INVALID_RESULTS_BATCH:
            return "The results batch is invalid";
        case ValidationError::CHUNK_HASH_MISMATCH_UNTOLERATED:
            return "Chunk hash mismatch could not be resolved as a tolerant match";
        case ValidationError::INVALID_PAYOUT_METADATA:
            return "The result payout metadata is missing or out of range";
    }
    return "Unknown error";
}

namespace sgns::processing
{
    outcome::result<void> ProcessingValidationCore::ValidateResults(
        const SGProcessing::SubTaskCollection                                            &subTasks,
        const std::map<std::string, SGProcessing::SubTaskResult>                         &results,
        std::set<std::string>                                                             &invalidSubTaskIds,
        const std::vector<sgns::Parameter>                                               *jobParameters,
        std::function<outcome::result<std::vector<uint8_t>>( const std::string &outputUri )> fetchOutputData )
    {
        std::optional<std::error_code> error;

        // Compare result hashes for each chunk
        // If a chunk hashes didn't match each other add the all subtasks with invalid hashes to VALID ITEMS LIST
        // Keyed on chunkKey -> {subtaskId -> ChunkContribution}: each subtask's contribution stays
        // independently addressable (fixes the old concatenation bug, which appended every contributing
        // subtask's chunk-hash bytes onto one shared buffer instead of comparing them).
        std::map<std::string, std::map<std::string, ChunkContribution>> chunksBySubtask;

        for ( int itemIdx = 0; itemIdx < subTasks.items_size(); ++itemIdx )
        {
            const auto &subTask  = subTasks.items( itemIdx );
            auto        itResult = results.find( subTask.subtaskid() );
            if ( itResult != results.end() )
            {
                if ( itResult->second.chunk_hashes_size() != subTask.chunkstoprocess_size() )
                {
                    m_logger->error( "WRONG_RESULT_HASHES_LENGTH {}: {} {}",
                                     subTask.subtaskid(),
                                     itResult->second.chunk_hashes_size(),
                                     subTask.chunkstoprocess_size() );
                    invalidSubTaskIds.insert( subTask.subtaskid() );
                    if ( !error )
                    {
                        error = make_error_code( Error::WRONG_RESULT_HASHES_LENGTH );
                    }
                }
                else
                {
                    for ( int chunkIdx = 0; chunkIdx < subTask.chunkstoprocess_size(); ++chunkIdx )
                    {
                        chunksBySubtask[subTask.chunkstoprocess( chunkIdx ).SerializeAsString()][subTask.subtaskid()] =
                            ChunkContribution{ itResult->second.chunk_hashes( chunkIdx ),
                                               chunkIdx,
                                               subTask.chunkstoprocess_size() };
                    }
                }
            }
            else
            {
                // Since all subtasks are processed a result should be found for all of them
                m_logger->error( "NO_RESULTS_FOUND {} on ", subTask.subtaskid() );
                invalidSubTaskIds.insert( subTask.subtaskid() );
                if ( !error )
                {
                    error = make_error_code( Error::NO_RESULTS_FOR_SUBTASK );
                }
            }
        }

        // Genuine cross-subtask comparison pass: for each chunk key contributed to by 2+ subtasks, either
        // every contribution agrees (trivial match) or a real mismatch exists that must be resolved --
        // via the tolerance fallback (XNODE-02) -- before being declared a genuine divergence (XNODE-01b).
        // Runs before the per-subtask CheckSubTaskResultHashes loop below so an already-invalidated
        // subtask (NO_RESULTS_FOR_SUBTASK/WRONG_RESULT_HASHES_LENGTH) is not double-processed there.
        for ( const auto &[chunkKey, contributions] : chunksBySubtask )
        {
            if ( contributions.size() < 2 )
            {
                continue; // nothing to compare -- only one subtask contributed this chunk
            }

            const auto &firstHash = contributions.begin()->second.hashBytes;
            bool        allMatch  = true;
            for ( const auto &[subtaskId, contribution] : contributions )
            {
                if ( contribution.hashBytes != firstHash )
                {
                    allMatch = false;
                    break;
                }
            }

            if ( allMatch )
            {
                continue; // genuine match (SC2)
            }

            // Genuine hash mismatch -- the exact case the old concatenation bug silently missed.
            if ( !AttemptToleranceFallback( contributions, results, jobParameters, fetchOutputData ) )
            {
                for ( const auto &[subtaskId, contribution] : contributions )
                {
                    invalidSubTaskIds.insert( subtaskId );
                }
                if ( !error )
                {
                    error = make_error_code( Error::CHUNK_HASH_MISMATCH_UNTOLERATED );
                }
            }
        }

        for ( int itemIdx = 0; itemIdx < subTasks.items_size(); ++itemIdx )
        {
            const auto &subTask = subTasks.items( itemIdx );
            if ( invalidSubTaskIds.find( subTask.subtaskid() ) != invalidSubTaskIds.end() )
            {
                m_logger->trace( "Subtask already invalid {}, no need to check chunk hashes ", subTask.subtaskid() );
                continue;
            }

            auto subtaskCheck = CheckSubTaskResultHashes( subTask, chunksBySubtask );
            if ( subtaskCheck.has_failure() )
            {
                invalidSubTaskIds.insert( subTask.subtaskid() );
                if ( !error )
                {
                    error = subtaskCheck.error();
                }
            }
        }

        if ( error )
        {
            return outcome::failure( *error );
        }
        return outcome::success();
    }

    outcome::result<void> ProcessingValidationCore::ValidateIndividualResult(
        const SGProcessing::SubTask       &subTask,
        const SGProcessing::SubTaskResult &result ) const
    {
        // Check 1: Verify subtask IDs match
        if ( subTask.subtaskid() != result.subtaskid() )
        {
            m_logger->error( "SUBTASK_ID_MISMATCH: expected {}, got {}", subTask.subtaskid(), result.subtaskid() );
            return outcome::failure( Error::SUBTASK_ID_MISMATCH );
        }

        // Check 2: Verify hash count matches chunk count
        if ( result.chunk_hashes_size() != subTask.chunkstoprocess_size() )
        {
            m_logger->error( "WRONG_RESULT_HASHES_LENGTH {}: {} {}",
                             subTask.subtaskid(),
                             result.chunk_hashes_size(),
                             subTask.chunkstoprocess_size() );
            return outcome::failure( Error::WRONG_RESULT_HASHES_LENGTH );
        }

        // Check 3: Verify no duplicate hashes
        std::unordered_set<std::string> encounteredHashes;
        for ( int chunkIdx = 0; chunkIdx < result.chunk_hashes_size(); ++chunkIdx )
        {
            const std::string &chunkHash = result.chunk_hashes( chunkIdx );

            if ( !encounteredHashes.insert( chunkHash ).second )
            {
                const auto &chunk = subTask.chunkstoprocess( chunkIdx );
                m_logger->error( "DUPLICATE_CHUNK_RESULT_HASH [{}, {}]", subTask.subtaskid(), chunk.chunkid() );
                return outcome::failure( Error::DUPLICATE_CHUNK_RESULT_HASH );
            }

            // Check 4: Verify hash is not empty
            if ( chunkHash.empty() )
            {
                const auto &chunk = subTask.chunkstoprocess( chunkIdx );
                m_logger->error( "EMPTY_CHUNK_RESULT_HASH [{}, {}]", subTask.subtaskid(), chunk.chunkid() );
                return outcome::failure( Error::EMPTY_CHUNK_RESULT_HASH );
            }
        }

        // Check 5: Verify the payout metadata this peer reported for itself and for the developer
        // of the app that ran the work. A payout cannot be built without it, so reject the result
        // here rather than let it block the escrow release later.
        if ( !base::IsHexAddress( result.node_address() ) )
        {
            m_logger->error( "INVALID_PAYOUT_METADATA {}: peer address is not an account address",
                             subTask.subtaskid() );
            return outcome::failure( Error::INVALID_PAYOUT_METADATA );
        }
        if ( result.developer_address().empty() )
        {
            m_logger->error( "INVALID_PAYOUT_METADATA {}: no developer address", subTask.subtaskid() );
            return outcome::failure( Error::INVALID_PAYOUT_METADATA );
        }
        // Checked on the raw bytes, not via TokenID, which left-pads anything shorter than 32.
        if ( result.token_id().size() != TOKEN_ID_BYTES )
        {
            m_logger->error( "INVALID_PAYOUT_METADATA {}: token id is {} bytes, expected {}",
                             subTask.subtaskid(),
                             result.token_id().size(),
                             TOKEN_ID_BYTES );
            return outcome::failure( Error::INVALID_PAYOUT_METADATA );
        }
        if ( result.developer_cut() > DEVELOPER_CUT_SCALE )
        {
            m_logger->error( "INVALID_PAYOUT_METADATA {}: developer cut {} exceeds {}",
                             subTask.subtaskid(),
                             result.developer_cut(),
                             DEVELOPER_CUT_SCALE );
            return outcome::failure( Error::INVALID_PAYOUT_METADATA );
        }

        return outcome::success();
    }

    outcome::result<void> ProcessingValidationCore::CheckSubTaskResultHashes(
        const SGProcessing::SubTask                                            &subTask,
        const std::map<std::string, std::map<std::string, ChunkContribution>> &chunksBySubtask ) const
    {
        std::unordered_set<std::string> encounteredHashes;
        for ( int chunkIdx = 0; chunkIdx < subTask.chunkstoprocess_size(); ++chunkIdx )
        {
            const auto &chunk = subTask.chunkstoprocess( chunkIdx );
            auto        it    = chunksBySubtask.find( chunk.SerializeAsString() );
            if ( it != chunksBySubtask.end() )
            {
                auto contributionIt = it->second.find( subTask.subtaskid() );
                if ( contributionIt == it->second.end() )
                {
                    // Should not happen -- this subtask's own contribution was just inserted in the
                    // build loop above -- but guard defensively rather than dereference blindly.
                    m_logger->error( "NO_CHUNK_RESULT_FOUND [{}, {}]", subTask.subtaskid(), chunk.chunkid() );
                    return outcome::failure( Error::MISSING_CHUNK_RESULT );
                }

                const std::string &chunkHash = contributionIt->second.hashBytes;
                if ( !encounteredHashes.insert( chunkHash ).second )
                {
                    m_logger->error( "INVALID_CHUNK_RESULT_HASH [{}, {}]", subTask.subtaskid(), chunk.chunkid() );
                    return outcome::failure( Error::INVALID_CHUNK_RESULT_HASH );
                }
            }
            else
            {
                m_logger->error( "NO_CHUNK_RESULT_FOUND [{}, {}]", subTask.subtaskid(), chunk.chunkid() );
                return outcome::failure( Error::MISSING_CHUNK_RESULT );
            }
        }
        return outcome::success();
    }

    bool ProcessingValidationCore::AttemptToleranceFallback(
        const std::map<std::string, ChunkContribution>                                  &contributions,
        const std::map<std::string, SGProcessing::SubTaskResult>                         &results,
        const std::vector<sgns::Parameter>                                               *jobParameters,
        const std::function<outcome::result<std::vector<uint8_t>>( const std::string & )> &fetchOutputData ) const
    {
        // D-02: fail closed when no fetch mechanism is available -- no fetch attempted at all.
        if ( !fetchOutputData )
        {
            return false;
        }

        // Fetch and slice out this chunk's byte range from each contributing subtask's output blob.
        std::vector<std::vector<uint8_t>> slicedChunks;
        slicedChunks.reserve( contributions.size() );

        for ( const auto &[subtaskId, contribution] : contributions )
        {
            auto resultIt = results.find( subtaskId );
            if ( resultIt == results.end() )
            {
                // Should not happen -- contributions are only built from results actually present --
                // but fail closed rather than dereference blindly.
                return false;
            }

            const std::string &outputUri = resultIt->second.ipfs_results_data_id();
            if ( outputUri.empty() )
            {
                // Pitfall 2: no fetchable data for this subtask -- fail closed, not a crash.
                return false;
            }

            auto fetchResult = fetchOutputData( outputUri );
            if ( !fetchResult )
            {
                m_logger->debug( "AttemptToleranceFallback: fetch failed for {}: {}",
                                  outputUri,
                                  fetchResult.error().message() );
                return false;
            }

            const std::vector<uint8_t> &blob = fetchResult.value();

            if ( contribution.totalChunksForSubtask <= 1 )
            {
                // The entire blob IS the chunk -- exact/primary case.
                slicedChunks.push_back( blob );
            }
            else if ( blob.size() % static_cast<size_t>( contribution.totalChunksForSubtask ) == 0 )
            {
                // Uniform-division slicing -- single-channel-only case. Multi-channel outputs are NOT
                // exactly sliceable this way (channel-major layout); that case is a documented, accepted
                // scope limit, not attempted here.
                const size_t sliceSize  = blob.size() / static_cast<size_t>( contribution.totalChunksForSubtask );
                const size_t sliceStart = static_cast<size_t>( contribution.chunkIdx ) * sliceSize;
                slicedChunks.emplace_back( blob.begin() + static_cast<long>( sliceStart ),
                                           blob.begin() + static_cast<long>( sliceStart + sliceSize ) );
            }
            else
            {
                // Cannot slice safely (e.g. overlapping-window subtask) -- fail closed rather than
                // attempt an inexact slice that could silently compare the wrong bytes (Pitfall 3).
                return false;
            }
        }

        // Determine comparison mode and diff every subsequent sliced buffer against the first one's.
        const auto elementType = sgprocmanagerdiff::ResolveChunkElementTypeHint( jobParameters );

        for ( size_t idx = 1; idx < slicedChunks.size(); ++idx )
        {
            sgprocmanagerdiff::ElementDiffStats stats;
            bool                                withinTolerance = false;

            if ( elementType == sgprocmanagerdiff::ChunkElementType::UINT8 )
            {
                withinTolerance =
                    sgprocmanagerdiff::IsByteChunkWithinTolerance( slicedChunks[0], slicedChunks[idx], jobParameters, stats );
            }
            else
            {
                withinTolerance =
                    sgprocmanagerdiff::IsFloatChunkWithinTolerance( slicedChunks[0], slicedChunks[idx], jobParameters, stats );
            }

            m_logger->debug( "AttemptToleranceFallback: maxAbsDelta={} maxRelDelta={} withinTolerance={}",
                              stats.maxAbsDelta,
                              stats.maxRelDelta,
                              withinTolerance );

            if ( !withinTolerance )
            {
                return false;
            }
        }

        return true;
    }

}

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