diff --git a/.github/ci/spack-envs/gcc13_py312_mpich_h5_ad2/spack.yaml b/.github/ci/spack-envs/gcc13_py312_mpich_h5_ad2/spack.yaml index 6b3851858f..c6ba694d5b 100644 --- a/.github/ci/spack-envs/gcc13_py312_mpich_h5_ad2/spack.yaml +++ b/.github/ci/spack-envs/gcc13_py312_mpich_h5_ad2/spack.yaml @@ -6,7 +6,10 @@ # spack: specs: - - adios2@2.10 +mpi + # Pin the oldest ADIOS2 release that still lacks resetting of memory + # selections (added upstream in 2.11.0, backported to 2.10.1). This exercises + # the code path that rejects memory selections for such versions. + - adios2@2.10.0 +mpi - hdf5 +mpi - mpich diff --git a/CMakeLists.txt b/CMakeLists.txt index 66f9a53a6e..32d598090d 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -408,6 +408,7 @@ set(CORE_SOURCE src/Format.cpp src/Iteration.cpp src/IterationEncoding.cpp + src/LoadStoreChunk.cpp src/Mesh.cpp src/ParticlePatches.cpp src/ParticleSpecies.cpp @@ -419,6 +420,7 @@ set(CORE_SOURCE src/version.cpp src/auxiliary/Date.cpp src/auxiliary/Filesystem.cpp + src/auxiliary/Future.cpp src/auxiliary/JSON.cpp src/auxiliary/JSONMatcher.cpp src/auxiliary/Memory.cpp @@ -856,6 +858,11 @@ if(openPMD_BUILD_TESTING) test/Files_SerialIO/filebased_write_test.cpp test/Files_SerialIO/issue_1744_unique_ptrs_at_close_time.cpp test/Files_SerialIO/components_without_extent.cpp + test/Files_SerialIO/deferred_load_unrelated_flush.cpp + test/Files_SerialIO/container_buffer_size_1d_downsizes.cpp + test/Files_SerialIO/joined_dim_buffer_api.cpp + test/Files_SerialIO/memory_selection_rejected_before_flush.cpp + test/Files_SerialIO/memory_selection_old_adios2.cpp ) elseif(${test_name} STREQUAL "ParallelIO" AND openPMD_HAVE_MPI) list(APPEND ${out_list} diff --git a/docs/source/details/mpi.rst b/docs/source/details/mpi.rst index 38bdc2643d..3ea6fb19ad 100644 --- a/docs/source/details/mpi.rst +++ b/docs/source/details/mpi.rst @@ -31,6 +31,8 @@ Functionality Behavior Description ``::makeConstant`` [3]_ *backend-specific* declare, write ``::storeChunk`` [1]_ independent write ``::loadChunk`` independent read +``::prepareLoadStore()`` independent configure deferred I/O +``DeferredComputation`` [5]_ **collective** deferred store/load + flush ``::availableChunks`` [4]_ collective read, immediate result ============================ ================== ================================ @@ -49,6 +51,14 @@ Functionality Behavior Description .. [4] We usually open iterations delayed on first access. This first access is usually the ``flush()`` call after a ``storeChunk``/``loadChunk`` operation. If the first access is non-collective, an explicit, collective ``Iteration::open()`` can be used to have the files already open. Alternatively, iterations might be accessed for the first time by immediate operations such as ``::availableChunks()``. +.. [5] The experimental deferred I/O API (``RecordComponent::prepareLoadStore()`` followed by ``store()`` / ``load()``) returns handles of type ``auxiliary::DeferredComputation``. + By default, invoking such a handle via ``get()`` / ``operator()()`` performs the data operation *and* automatically flushes the underlying Series. + Since ``Series::flush()`` is collective, this makes invoking the handle a collective operation: every rank must invoke (or explicitly destroy) its handles in a consistent order, even if the number of ``store()`` / ``load()`` calls differs between ranks. + The per-``RecordComponent`` flush counter is used to skip a flush if the component has already been flushed by another component in the meantime, so a given invocation does not necessarily flush. + The returned handles are marked ``[[nodiscard]]`` and the automatic flush is only triggered by an explicit invocation, so handles cannot be dropped unnoticed. + Call ``unsafeNoAutomaticFlush()`` on the configuration object to disable the automatic flush and instead call ``Series::flush()`` explicitly and collectively at a suitable point. + The destructor of a still-valid handle also performs the automatic flush and therefore has the same collective semantics as ``get()``. + .. warning:: The openPMD-api will by default flush only those Iterations which are dirty, i.e. have been written to. diff --git a/docs/source/dev/design.rst b/docs/source/dev/design.rst index 6fb81071ce..e9b7d1f1b4 100644 --- a/docs/source/dev/design.rst +++ b/docs/source/dev/design.rst @@ -23,7 +23,7 @@ Therefore, enabling users to handle hierarchical, self-describing file formats w .. literalinclude:: IOTask.hpp :language: cpp - :lines: 50-81 + :lines: 57-91 Every task is designed to be a fully self-contained description of one such atomic operation. By describing a required minimal step of work (without any side-effect), these operations are the foundation of the unified handling mechanism across suitable file formats. The actual low-level exchange of data is implemented in ``IOHandlers``, one per file format (possibly two if handlingi MPI-parallel work is possible and requires different behaviour). diff --git a/docs/source/usage/workflow.rst b/docs/source/usage/workflow.rst index 8f79b3108e..8050841290 100644 --- a/docs/source/usage/workflow.rst +++ b/docs/source/usage/workflow.rst @@ -127,6 +127,15 @@ Flush points are triggered by: Flush point guarantees affect only the corresponding iteration. * Calling ``Writable::seriesFlush()`` or ``Attributable::seriesFlush()``. * The streaming API (i.e. ``Series.readIterations()`` and ``Series.writeIteration()``) automatically before accessing the next iteration. +* Invoking a handle returned by the experimental deferred I/O API (``RecordComponent::prepareLoadStore()`` followed by ``store()`` / ``load()``) via ``get()`` / ``operator()()``, unless automatic flushing was disabled via ``unsafeNoAutomaticFlush()``. + +.. note:: + + The automatic flush performed by the deferred I/O handles is MPI-collective (see :ref:`details-mpi`). + Each rank must invoke such a handle the same number of times and in the same order, *even if* the number of ``store()`` / ``load()`` calls differs per rank. + The flush is skipped if the pertaining ``RecordComponent`` has already been flushed by another component in the meantime (tracked via per-component flush counters), so a given invocation does not necessarily flush. + To avoid this coupling, disable the automatic flush with ``unsafeNoAutomaticFlush()`` and call ``Series::flush()`` explicitly at a collective point. + The handles are returned with the ``[[nodiscard]]`` attribute and the flush is only triggered by explicitly invoking them (or by destroying a still-valid handle), so it cannot be triggered accidentally. Attributes are (currently) unaffected by this: @@ -142,8 +151,8 @@ Attributes are (currently) unaffected by this: For user-guided selection of such implementations, ``Series::flush`` and ``Attributable::seriesFlush()`` take an optional JSON/TOML string as a parameter. See the section on :ref:`backend-specific configuration ` for details. -Deferred Data API Contract --------------------------- +Verbose Logging +--------------- A verbose debug log can optionally be printed to the standard error output by specifying the environment variable ``OPENPMD_VERBOSE=1``. Note that this functionality is at the current time still relatively basic. diff --git a/include/openPMD/Dataset.hpp b/include/openPMD/Dataset.hpp index e1f0058885..b8af286e14 100644 --- a/include/openPMD/Dataset.hpp +++ b/include/openPMD/Dataset.hpp @@ -34,6 +34,18 @@ namespace openPMD using Extent = std::vector; using Offset = std::vector; +/** Selection of a region of memory for storing chunks. + * + * Used to specify a non-contiguous memory region when storing + * data chunks. This allows writing data that is not contiguous + * in memory. + */ +struct MemorySelection +{ + Offset offset; + Extent extent; +}; + class Dataset { friend class RecordComponent; diff --git a/include/openPMD/Datatype.hpp b/include/openPMD/Datatype.hpp index 17cf6b67f4..3f667ae47f 100644 --- a/include/openPMD/Datatype.hpp +++ b/include/openPMD/Datatype.hpp @@ -420,6 +420,11 @@ inline size_t toBits(Datatype d) return toBytes(d) * CHAR_BIT; } +/** Check if a Datatype is a signed type + * + * @param d Datatype to test + * @return true if signed type (integer, floating point, complex), else false + */ constexpr bool isSigned(Datatype d); /** Compare if a Datatype is a vector type @@ -602,6 +607,13 @@ inline bool isSameFloatingPoint(Datatype d) return isSameFloatingPoint(d, determineDatatype()); } +/** Compare if two Datatypes are equivalent floating point types + * + * @param d1 First Datatype to compare + * @param d2 Second Datatype to compare + * @return true if both types are floating point and have same bitness, else + * false + */ inline bool isSameFloatingPoint(Datatype d1, Datatype d2) { // template @@ -629,6 +641,13 @@ inline bool isSameComplexFloatingPoint(Datatype d) return isSameComplexFloatingPoint(d, determineDatatype()); } +/** Compare if two Datatypes are equivalent complex floating point types + * + * @param d1 First Datatype to compare + * @param d2 Second Datatype to compare + * @return true if both types are complex floating point and have same bitness, + * else false + */ inline bool isSameComplexFloatingPoint(Datatype d1, Datatype d2) { // template @@ -656,6 +675,13 @@ inline bool isSameInteger(Datatype d) return isSameInteger(d, determineDatatype()); } +/** Compare if two Datatypes are equivalent integer types + * + * @param d1 First Datatype to compare + * @param d2 Second Datatype to compare + * @return true if both types are integers, same signedness and same bitness, + * else false + */ inline bool isSameInteger(Datatype d1, Datatype d2) { // template @@ -708,6 +734,13 @@ constexpr bool isChar(Datatype d) template constexpr bool isSameChar(Datatype d); +/** Compare if two Datatypes are equivalent char types + * + * @param d1 First Datatype to compare + * @param d2 Second Datatype to compare + * @return true if both types are chars with same signedness and size, else + * false + */ constexpr bool isSameChar(Datatype d1, Datatype d2); /** Comparison for two Datatypes @@ -715,6 +748,10 @@ constexpr bool isSameChar(Datatype d1, Datatype d2); * Besides returning true for the same types, identical implementations on * some platforms, e.g. if long and long long are the same or double and * long double will also return true. + * + * @param d First Datatype to compare + * @param e Second Datatype to compare + * @return true if the datatypes are equivalent */ constexpr bool isSame(openPMD::Datatype d, openPMD::Datatype e); @@ -726,15 +763,34 @@ constexpr bool isSame(openPMD::Datatype d, openPMD::Datatype e); */ Datatype basicDatatype(Datatype dt); +/** Convert a scalar Datatype to its vector variant + * + * @param dt Scalar Datatype to convert + * @return Vector Datatype (e.g., INT becomes VEC_INT) + */ Datatype toVectorType(Datatype dt); +/** Convert a Datatype to its string representation + * + * @param dt Datatype to convert + * @return String representation of the Datatype + */ std::string datatypeToString(Datatype dt); +/** Convert a string to a Datatype + * + * @param s String representation of a Datatype + * @return The corresponding Datatype + */ Datatype stringToDatatype(const std::string &s); -void warnWrongDtype(std::string const &key, Datatype store, Datatype request); - -std::ostream &operator<<(std::ostream &, openPMD::Datatype const &); +/** Stream operator for Datatype + * + * @param os Output stream + * @param dt Datatype to output + * @return Reference to the stream + */ +std::ostream &operator<<(std::ostream &os, openPMD::Datatype const &dt); template constexpr auto datatypeIndex() -> size_t diff --git a/include/openPMD/Datatype.tpp b/include/openPMD/Datatype.tpp index e35f2e26b6..b09351fd77 100644 --- a/include/openPMD/Datatype.tpp +++ b/include/openPMD/Datatype.tpp @@ -223,36 +223,52 @@ namespace detail template constexpr bool is_char_v = is_char::value; - template - inline bool isSameChar() + struct IsChar { - return - // both must be char types - is_char_v && is_char_v && - // both must have equivalent sign - std::is_signed_v == std::is_signed_v && - // both must have equivalent size - sizeof(T_Char1) == sizeof(T_Char2); + template + static constexpr bool call() + { + return is_char_v; + } + template + static constexpr bool call() + { + return false; + } + }; + + constexpr inline bool isChar(Datatype dtype) + { + return switchType(dtype); } - template - struct IsSameChar + struct DtypeSize { - template - static bool call() + template + static constexpr size_t call() { - return isSameChar(); + return sizeof(T); } - - static constexpr char const *errorMsg = "IsSameChar"; + static constexpr char const *errorMsg = "DtypeSize"; }; + constexpr inline size_t dtypeSize(Datatype dtype) + { + return switchType(dtype); + } } // namespace detail template constexpr inline bool isSameChar(Datatype d) { - return switchType>(d); + return isSameChar(d, determineDatatype()); +} + +constexpr bool isSameChar(Datatype d1, Datatype d2) +{ + return detail::isChar(d1) && detail::isChar(d2) && + isSigned(d1) == isSigned(d2) && + detail::dtypeSize(d1) == detail::dtypeSize(d2); } namespace detail @@ -285,11 +301,6 @@ constexpr inline bool isSigned(Datatype d) return switchType(d); } -constexpr inline bool isSameChar(Datatype d, Datatype e) -{ - return isChar(d) && isChar(e) && isSigned(d) == isSigned(e); -} - constexpr bool isSame(openPMD::Datatype const d, openPMD::Datatype const e) { return diff --git a/include/openPMD/IO/ADIOS/ADIOS2File.hpp b/include/openPMD/IO/ADIOS/ADIOS2File.hpp index d34cc8ebe5..66aa47e702 100644 --- a/include/openPMD/IO/ADIOS/ADIOS2File.hpp +++ b/include/openPMD/IO/ADIOS/ADIOS2File.hpp @@ -20,6 +20,7 @@ */ #pragma once +#include "openPMD/Dataset.hpp" #include "openPMD/IO/ADIOS/ADIOS2Auxiliary.hpp" #include "openPMD/IO/ADIOS/ADIOS2PreloadAttributes.hpp" #include "openPMD/IO/ADIOS/ADIOS2PreloadVariables.hpp" @@ -107,11 +108,14 @@ struct WriteDataset static void call(Params &&...); }; +/** Buffered put operation with unique pointer */ struct BufferedUniquePtrPut { std::string name; Offset offset; Extent extent; + /** Optional memory selection for non-contiguous memory regions */ + std::optional memorySelection; UniquePtrWithLambda data; Datatype dtype = Datatype::UNDEFINED; diff --git a/include/openPMD/IO/ADIOS/ADIOS2IOHandler.hpp b/include/openPMD/IO/ADIOS/ADIOS2IOHandler.hpp index 4316f5181f..218e733668 100644 --- a/include/openPMD/IO/ADIOS/ADIOS2IOHandler.hpp +++ b/include/openPMD/IO/ADIOS/ADIOS2IOHandler.hpp @@ -21,6 +21,7 @@ */ #pragma once +#include "openPMD/Dataset.hpp" #include "openPMD/Error.hpp" #include "openPMD/IO/ADIOS/ADIOS2Auxiliary.hpp" #include "openPMD/IO/ADIOS/ADIOS2FilePosition.hpp" @@ -509,6 +510,7 @@ class ADIOS2IOHandlerImpl adios2::Variable verifyDataset( Offset const &offset, Extent const &extent, + std::optional const &memorySelection, adios2::IO &IO, adios2::Engine &engine, std::string const &varName, @@ -622,6 +624,28 @@ class ADIOS2IOHandlerImpl var.SetSelection( {adios2::Dims(offset.begin(), offset.end()), adios2::Dims(extent.begin(), extent.end())}); + + if (memorySelection.has_value()) + { + if (!openPMD::CanTheMemorySelectionBeReset) + { + throw error::OperationUnsupportedInBackend( + "ADIOS2", + "Non-contiguous memory selections are not supported with " + "this version of ADIOS2 (upstream since 2.11.0, backported " + "to 2.10.1): A memory selection can not be reset once " + "specified, which would silently affect subsequent put " + "operations."); + } + var.SetMemorySelection( + {adios2::Dims( + memorySelection->offset.begin(), + memorySelection->offset.end()), + adios2::Dims( + memorySelection->extent.begin(), + memorySelection->extent.end())}); + } + return var; } @@ -942,7 +966,7 @@ class ADIOS2IOHandler : public AbstractIOHandler try { auto params = internal::defaultParsedFlushParams; - this->flush(params); + this->flush_impl(params); } catch (std::exception const &ex) { @@ -990,6 +1014,20 @@ class ADIOS2IOHandler : public AbstractIOHandler return true; } - std::future flush(internal::ParsedFlushParams &) override; + bool supportsMemorySelection() const override + { +#if openPMD_HAVE_ADIOS2 + /* + * A memory selection that cannot be reset would silently leak into + * subsequent store operations of the same variable. That ability was + * added upstream in ADIOS2 v2.11.0 and backported to v2.10.1. + */ + return openPMD::CanTheMemorySelectionBeReset; +#else + return false; +#endif + } + + std::future flush_impl(internal::ParsedFlushParams &) override; }; // ADIOS2IOHandler } // namespace openPMD diff --git a/include/openPMD/IO/ADIOS/macros.hpp b/include/openPMD/IO/ADIOS/macros.hpp index 8195d36e8a..d62b16f69e 100644 --- a/include/openPMD/IO/ADIOS/macros.hpp +++ b/include/openPMD/IO/ADIOS/macros.hpp @@ -26,6 +26,8 @@ #include +#include + #define openPMD_HAS_ADIOS_2_10 \ (ADIOS2_VERSION_MAJOR * 100 + ADIOS2_VERSION_MINOR >= 210) @@ -46,6 +48,40 @@ #define openPMD_HAVE_ADIOS2_BP5 0 #endif +namespace openPMD +{ +namespace detail +{ + /** Trait to check if SetMemorySelection can be called without arguments to + * clear a previously set memory selection. + * + * ADIOS2 gained this capability in v2.11.0; it was backported to v2.10.1. + * Before that, an empty box would fail the dimensionality and rank checks + * inside ADIOS2, so a memory selection could never be reset. + * + * @tparam Variable ADIOS2 variable type + */ + template + struct CanTheMemorySelectionBeReset + { + static constexpr bool value = false; + }; + + template + struct CanTheMemorySelectionBeReset< + Variable, + std::void_t().SetMemorySelection())>> + { + static constexpr bool value = true; + }; +} // namespace detail + +/** Whether ADIOS2 Variable supports resetting a memory selection via + * SetMemorySelection() without arguments (ADIOS2 >= 2.10.1) */ +constexpr bool CanTheMemorySelectionBeReset = + detail::CanTheMemorySelectionBeReset>::value; +} // namespace openPMD + #else #define openPMD_HAS_ADIOS_2_8 0 diff --git a/include/openPMD/IO/AbstractIOHandler.hpp b/include/openPMD/IO/AbstractIOHandler.hpp index b54b661c78..f10d432670 100644 --- a/include/openPMD/IO/AbstractIOHandler.hpp +++ b/include/openPMD/IO/AbstractIOHandler.hpp @@ -333,12 +333,21 @@ class AbstractIOHandler * @return Future indicating the completion state of the operation for * backends that decide to implement this operation asynchronously. */ - virtual std::future flush(internal::ParsedFlushParams &) = 0; + std::future flush(internal::ParsedFlushParams &); /** The currently used backend */ virtual std::string backendName() const = 0; virtual bool fullSupportForVariableBasedEncoding() const; + /** + * Whether the backend supports storing chunks with a non-contiguous + * memory selection. Backends that do not support this must reject such + * chunks in the frontend (RecordComponent::storeChunk_impl()) instead of + * letting the backend throw at flush time, which would clear the whole IO + * queue and corrupt the output. + */ + virtual bool supportsMemorySelection() const; + std::string directory; /* * Originally, the reason for distinguishing these two was that during @@ -377,6 +386,17 @@ class AbstractIOHandler IterationEncoding m_encoding = IterationEncoding::groupBased; OpenpmdStandard m_standard = auxiliary::parseStandard(getStandardDefault()); bool m_verify_homogeneous_extents = true; + +protected: + /** Implementation of flush operation for subclasses + * + * Do not call directly, use flush() wrapper instead. + * + * @param params Parsed flush parameters + * @return Future indicating completion state + */ + virtual std::future + flush_impl(internal::ParsedFlushParams ¶ms) = 0; }; // AbstractIOHandler } // namespace openPMD diff --git a/include/openPMD/IO/AbstractIOHandlerImpl.hpp b/include/openPMD/IO/AbstractIOHandlerImpl.hpp index fdb8af3599..2aa0e9373f 100644 --- a/include/openPMD/IO/AbstractIOHandlerImpl.hpp +++ b/include/openPMD/IO/AbstractIOHandlerImpl.hpp @@ -430,6 +430,9 @@ class AbstractIOHandlerImpl virtual void setWritten(Writable *, Parameter const ¶m); + virtual void increaseFlushCounter( + Writable *, Parameter const ¶m); + AbstractIOHandler *m_handler; bool m_verboseIOTasks = false; diff --git a/include/openPMD/IO/DummyIOHandler.hpp b/include/openPMD/IO/DummyIOHandler.hpp index 8abcf20990..100711c944 100644 --- a/include/openPMD/IO/DummyIOHandler.hpp +++ b/include/openPMD/IO/DummyIOHandler.hpp @@ -44,7 +44,7 @@ class DummyIOHandler : public AbstractIOHandler /** No-op consistent with the IOHandler interface to enable library use * without IO. */ - std::future flush(internal::ParsedFlushParams &) override; + std::future flush_impl(internal::ParsedFlushParams &) override; std::string backendName() const override; }; // DummyIOHandler } // namespace openPMD diff --git a/include/openPMD/IO/HDF5/HDF5IOHandler.hpp b/include/openPMD/IO/HDF5/HDF5IOHandler.hpp index 07b3978b87..5e2ddfd526 100644 --- a/include/openPMD/IO/HDF5/HDF5IOHandler.hpp +++ b/include/openPMD/IO/HDF5/HDF5IOHandler.hpp @@ -46,7 +46,7 @@ class HDF5IOHandler : public AbstractIOHandler return "HDF5"; } - std::future flush(internal::ParsedFlushParams &) override; + std::future flush_impl(internal::ParsedFlushParams &) override; private: std::unique_ptr m_impl; diff --git a/include/openPMD/IO/HDF5/ParallelHDF5IOHandler.hpp b/include/openPMD/IO/HDF5/ParallelHDF5IOHandler.hpp index abeb196b11..7b0c43129d 100644 --- a/include/openPMD/IO/HDF5/ParallelHDF5IOHandler.hpp +++ b/include/openPMD/IO/HDF5/ParallelHDF5IOHandler.hpp @@ -56,7 +56,7 @@ class ParallelHDF5IOHandler : public AbstractIOHandler return "MPI_HDF5"; } - std::future flush(internal::ParsedFlushParams &) override; + std::future flush_impl(internal::ParsedFlushParams &) override; private: std::unique_ptr m_impl; diff --git a/include/openPMD/IO/IOTask.hpp b/include/openPMD/IO/IOTask.hpp index a8efab571e..03673a1899 100644 --- a/include/openPMD/IO/IOTask.hpp +++ b/include/openPMD/IO/IOTask.hpp @@ -86,7 +86,8 @@ OPENPMDAPI_EXPORT_ENUM_CLASS(Operation){ AVAILABLE_CHUNKS, //!< Query chunks that can be loaded in a dataset DEREGISTER, //!< Inform the backend that an object has been deleted. TOUCH, //!< tell the backend that the file is to be considered active - SET_WRITTEN //!< tell backend to consider a file written / not written + SET_WRITTEN, //!< tell backend to consider a file written / not written + INCREASE_FLUSH_COUNTER //!< track if an object has been flushed }; // note: if you change the enum members here, please update // docs/source/dev/design.rst @@ -498,6 +499,8 @@ struct OPENPMDAPI_EXPORT Extent extent = {}; Offset offset = {}; + /** Optional memory selection for non-contiguous memory regions */ + std::optional memorySelection = std::nullopt; Datatype dtype = Datatype::UNDEFINED; auxiliary::WriteBuffer data; }; @@ -561,7 +564,9 @@ struct OPENPMDAPI_EXPORT } // in parameters - bool queryOnly = false; // query if the backend supports this + /** If true, only query if the backend supports buffer views without + * performing operation */ + bool queryOnly = false; Offset offset; Extent extent; Datatype dtype = Datatype::UNDEFINED; @@ -828,6 +833,27 @@ struct OPENPMDAPI_EXPORT bool target_status = false; }; +template <> +struct OPENPMDAPI_EXPORT + Parameter : public AbstractParameter +{ + explicit Parameter() = default; + + Parameter(Parameter const &) = default; + Parameter(Parameter &&) = default; + + Parameter &operator=(Parameter const &) = default; + Parameter &operator=(Parameter &&) = default; + + std::unique_ptr to_heap() && override + { + return std::make_unique>( + std::move(*this)); + } + + std::shared_ptr flush_counter; +}; + /** @brief Self-contained description of a single IO operation. * * Contained are diff --git a/include/openPMD/IO/InvalidatableFile.hpp b/include/openPMD/IO/InvalidatableFile.hpp index 89b763cda6..0fd7b82855 100644 --- a/include/openPMD/IO/InvalidatableFile.hpp +++ b/include/openPMD/IO/InvalidatableFile.hpp @@ -70,6 +70,11 @@ struct InvalidatableFile explicit operator bool() const; + /* + * + * Enables using InvalidatableFile in ordered containers like std::set + * for consistent ordering across parallel processes. + */ bool operator<(InvalidatableFile const &f) const; }; } // namespace openPMD diff --git a/include/openPMD/IO/JSON/JSONIOHandler.hpp b/include/openPMD/IO/JSON/JSONIOHandler.hpp index 07e797d4b3..db8b687ed9 100644 --- a/include/openPMD/IO/JSON/JSONIOHandler.hpp +++ b/include/openPMD/IO/JSON/JSONIOHandler.hpp @@ -59,7 +59,7 @@ class JSONIOHandler : public AbstractIOHandler return "JSON"; } - std::future flush(internal::ParsedFlushParams &) override; + std::future flush_impl(internal::ParsedFlushParams &) override; private: JSONIOHandlerImpl m_impl; diff --git a/include/openPMD/Iteration.hpp b/include/openPMD/Iteration.hpp index f1a4fd436b..6234e155b3 100644 --- a/include/openPMD/Iteration.hpp +++ b/include/openPMD/Iteration.hpp @@ -482,6 +482,12 @@ class Iteration namespace traits { + /** Generation policy for Iteration objects. + * + * This policy populates the cached iteration index when an Iteration + * is created or inserted into a Series, enabling constant-time lookup + * of the owning map entry. + */ template <> struct GenerationPolicy { diff --git a/include/openPMD/LoadStoreChunk.hpp b/include/openPMD/LoadStoreChunk.hpp new file mode 100644 index 0000000000..94ba1cabde --- /dev/null +++ b/include/openPMD/LoadStoreChunk.hpp @@ -0,0 +1,541 @@ +#pragma once + +#include "openPMD/Dataset.hpp" +#include "openPMD/auxiliary/Future.hpp" +#include "openPMD/auxiliary/Memory.hpp" +#include "openPMD/auxiliary/UniquePtr.hpp" + +// comment to prevent this include from being moved by clang-format +#include "openPMD/DatatypeMacros.hpp" + +#include +#include + +namespace openPMD +{ +class RecordComponent; +class ConfigureStoreChunkFromBuffer; +class ConfigureLoadStoreFromBuffer; +template +class DynamicMemoryView; +class Attributable; + +namespace internal +{ + /** Internal configuration for load/store operations without buffer. Default + * values for optionally specified parameters (offset, extent) must be + * computed to create this configuration struct. */ + struct LoadStoreConfig + { + Offset offset; + Extent extent; + }; + /** Internal configuration for load/store operations with buffer. Default + * values for optionally specified parameters (offset, extent) must be + * computed to create this configuration struct. MemorySelection remains + * optional even then. */ + struct LoadStoreConfigWithBuffer + { + Offset offset; + Extent extent; + std::optional memorySelection; + }; + +} // namespace internal + +namespace auxiliary::detail +{ +#define OPENPMD_ENUMERATE_TYPES(type) , std::shared_ptr + using shared_ptr_dataset_types = auxiliary::detail::variant_tail_t< + auxiliary::detail::bottom OPENPMD_FOREACH_DATASET_DATATYPE( + OPENPMD_ENUMERATE_TYPES)>; +#undef OPENPMD_ENUMERATE_TYPES +} // namespace auxiliary::detail + +/** Base class for configuring load/store chunk operations. + * + * Actual data members of `ConfigureLoadStore<>` and methods that don't + * depend on the ChildClass template parameter. By extracting the members to + * this struct, we can pass them around between different instances of the + * class template. Numbers of method instantiations can be reduced. + */ +class ConfigureLoadStore +{ + friend class openPMD::RecordComponent; + +protected: + ConfigureLoadStore(RecordComponent &); + std::shared_ptr m_rc; + + std::optional m_offset; + std::optional m_extent; + + bool m_unsafeNoAutomaticFlush = false; + +public: + ConfigureLoadStore(ConfigureLoadStore const &other); + ConfigureLoadStore &operator=(ConfigureLoadStore const &other); + ConfigureLoadStore(ConfigureLoadStore &&); + ConfigureLoadStore &operator=(ConfigureLoadStore &&); + +protected: + [[nodiscard]] auto dim() const -> uint8_t; + auto storeChunkConfig() -> internal::LoadStoreConfig; + + auto deferFlush(RecordComponent &); + + // The below methods return void. + // For chaining calls, they should return *this, but this class right + // here is going to be somewhere in the inheritance chain, and the final + // class should be returned. Could be solved more elegantly with CRT, + // but that blows up compile-time, so we make internal void functions + // and then repeat them in the final classes. + // (e.g. ConfigureLoadStoreFromBuffer::offset()) + + void offset_impl(Offset); + void extent_impl(Extent); + void unsafeNoAutomaticFlush_impl(); + + virtual auto getBufferSize() -> std::optional; + +private: + auto withSharedPtr_impl_mut(std::shared_ptr data, Datatype) + -> openPMD::ConfigureLoadStoreFromBuffer; + auto withSharedPtr_impl_const(std::shared_ptr data, Datatype) + -> openPMD::ConfigureStoreChunkFromBuffer; + auto withUniquePtr_impl_mut(UniquePtrWithLambda, Datatype) + -> openPMD::ConfigureStoreChunkFromBuffer; + auto withUniquePtr_impl_const(UniquePtrWithLambda, Datatype) + -> openPMD::ConfigureStoreChunkFromBuffer; + auto withRawPtr_impl_mut(void *data, Datatype) + -> openPMD::ConfigureLoadStoreFromBuffer; + auto withRawPtr_impl_const(void const *data, Datatype) + -> openPMD::ConfigureStoreChunkFromBuffer; + +public: + using this_t = ConfigureLoadStore; + + /** Retrieve the configured offset. If no offset has been specified, compute + * and store it now (default: full dataset selection). May be overwritten + * at a later point using offset(). + */ + auto computeOffset() -> Offset const &; + /** Retrieve the configured extent. If no extent has been specified, compute + * and store it now (default: full dataset selection). May be overwritten + * at a later point using extent(). + */ + auto computeExtent() -> Extent const &; + + // Configuration methods (always available) + + /** Set the offset within the dataset + * + * Optional. The operation will apply without offset by default (i.e. offset + * = (0, 0, ...)). + * + * @param offset Offset within the dataset + * @return Reference to this object for chaining + */ + auto offset(Offset offset) -> this_t & + { + offset_impl(std::move(offset)); + return *this; + } + /** Set the extent within the dataset + * + * Optional. The operation will apply to the entire dataset by default (i.e. + * operation extent = global dataset extent - operation offset). + * + * @param extent Extent within the dataset, counted from the offset + * @return Reference to this object for chaining + */ + auto extent(Extent extent) -> this_t & + { + extent_impl(std::move(extent)); + return *this; + } + /** Disable automatic flush after store operation + * + * The returned objects of type DeferredComputation will still return a + * buffer upon get() / operator()(), but these buffers are not guaranteed to + * be filled until explicitly flushing. + * By default, invoking get() / operator()() would also flush the underlying + * Series, an MPI-collective operation. Disabling that automatic flush + * removes the need for all ranks to invoke their handles in a consistent + * order; the user is then responsible for calling Series::flush() + * collectively at a suitable point. + * + * @return Reference to this object for chaining + */ + auto unsafeNoAutomaticFlush() -> this_t & + { + unsafeNoAutomaticFlush_impl(); + return *this; + } + + /* + * If the type is non-const, then the return type should be + * ConfigureLoadStoreFromBuffer, but if it is a const type, Load operations + * make no sense, so the return type should be + * ConfigureStoreChunkFromBuffer<>. + */ + template + using shared_ptr_return_type = std::conditional_t< + std::is_const_v, + ConfigureStoreChunkFromBuffer, + ConfigureLoadStoreFromBuffer>; + + /* + * As loading into unique pointer types makes no sense, the case is + * simpler for unique pointers. Just remove the array extents here. + * (Our interface wrappers still support const-type unique pointers, + * but the internal logic does not handle them separately.) + */ + template + using unique_ptr_return_type = openPMD::ConfigureStoreChunkFromBuffer; + + // Buffer specification methods (return specialized configurations) + template + [[nodiscard]] auto withSharedPtr(std::shared_ptr) + -> shared_ptr_return_type; + template + [[nodiscard]] auto withUniquePtr(UniquePtrWithLambda) + -> unique_ptr_return_type; + template + [[nodiscard]] auto withUniquePtr(std::unique_ptr) + -> unique_ptr_return_type; + template + [[nodiscard]] auto withRawPtr(T *data) -> shared_ptr_return_type; + /** Specify a contiguous container (std::vector, std::array, ...) as the + * buffer for the operation. + * + * The buffer size is inferred from the container and set automatically, so + * that the operation's extent is adjusted to the container when no explicit + * extent is given. + * + * @param data Contiguous container large enough for the selected data + * @return A buffer-specific configuration for the operation + */ + template + [[nodiscard]] auto withContiguousContainer(T_ContiguousContainer &data) + -> std::enable_if_t< + auxiliary::IsContiguousContainer_v, + shared_ptr_return_type>; + + // Enqueue methods (deferred execution) + template + [[nodiscard]] auto storeSpan() -> DynamicMemoryView; + // definition for this one is in RecordComponent.tpp since it needs the + // definition of class RecordComponent. + template + [[nodiscard]] auto storeSpan(F &&createBuffer) -> DynamicMemoryView; + + /** Load the chunk data, allocating a buffer + * + * The returned handle performs the load, and the automatic flush of the + * underlying Series, when invoked via get() / operator()(). That flush is + * an MPI-collective operation, so in parallel codes every rank must invoke + * (or explicitly destroy) its handles in a consistent order. Invoking the + * handle is an explicit, user-controlled action and the return value is + * marked [[nodiscard]], so handles cannot be dropped unnoticed. Use + * unsafeNoAutomaticFlush() to defer flushing to a later explicit + * Series::flush() instead. + * + * @return Deferred computation that performs the load when invoked + */ + template + [[nodiscard]] auto load() + -> auxiliary::DeferredComputation>; + + /** Type-erased version of load(). See load() for the collective semantics + * of the automatic flush. */ + [[nodiscard]] auto loadVariant() -> auxiliary::DeferredComputation< + auxiliary::detail::shared_ptr_dataset_types>; + + [[nodiscard]] auto getComponentHandle() const -> RecordComponent; +}; + +/** Configuration for storing chunks from a buffer. + * + * This class is used to configure a store chunk operation, where data is + * stored from a provided buffer into a dataset. + * This class is distinct from ConfigureLoadStoreFromBuffer, since reading + * data does not make sense on const / unique pointer types. This way, the type + * system will only allow read operations where they can actually run. + */ +class ConfigureStoreChunkFromBuffer : public ConfigureLoadStore +{ + friend class ConfigureLoadStore; + +protected: + auxiliary::WriteBuffer m_buffer; + Datatype m_datatype; + std::optional m_mem_select; + std::optional m_buffer_size; + + ConfigureStoreChunkFromBuffer( + auxiliary::WriteBuffer buffer, Datatype, ConfigureLoadStore &&); + + // The below methods return void. + // For chaining calls, they should return *this, but this class right + // here is going to be somewhere in the inheritance chain, and the final + // class should be returned. Could be solved more elegantly with CRT, + // but that blows up compile-time, so we make internal void functions + // and then repeat them in the final classes. + + /** Set memory selection for non-contiguous memory regions */ + void memorySelection_impl(MemorySelection); + + auto storeChunkConfig() -> internal::LoadStoreConfigWithBuffer; + + void bufferSize_impl(size_t); + + auto getBufferSize() -> std::optional override; + +public: + using this_t = ConfigureStoreChunkFromBuffer; + + // Configuration methods (always available) + + /** Set the offset within the dataset + * + * Optional. The operation will apply without offset by default (i.e. offset + * = (0, 0, ...)). + * + * @param offset Offset within the dataset + * @return Reference to this object for chaining + */ + auto offset(Offset offset) -> this_t & + { + offset_impl(std::move(offset)); + return *this; + } + + /** Set the extent within the dataset + * + * Optional. The operation will apply to the entire dataset by default (i.e. + * operation extent = global dataset extent - operation offset). + * + * @param extent Extent within the dataset, counted from the offset + * @return Reference to this object for chaining + */ + auto extent(Extent extent) -> this_t & + { + extent_impl(std::move(extent)); + return *this; + } + + /** Disable automatic flush after store operation + * + * The returned objects of type DeferredComputation will still return a + * buffer upon get() / operator()(), but these buffers are not guaranteed to + * be filled until explicitly flushing. + * By default, invoking get() / operator()() would also flush the underlying + * Series, an MPI-collective operation. Disabling that automatic flush + * removes the need for all ranks to invoke their handles in a consistent + * order; the user is then responsible for calling Series::flush() + * collectively at a suitable point. + * + * @return Reference to this object for chaining + */ + auto unsafeNoAutomaticFlush() -> this_t & + { + unsafeNoAutomaticFlush_impl(); + return *this; + } + + /** Set memory selection for non-contiguous memory regions + * + * Only supported with ADIOS2 >= 2.10.1 (the capability to reset a memory + * selection was added upstream in 2.11.0 and backported to 2.10.1). Older + * versions cannot reset a memory selection once it has been set, which + * would silently leak it into subsequent store operations of the same + * variable. Those versions reject memory selections with an error at store + * time. + * + * @param memorySelection Selection of memory region + * @return Reference to this object for chaining + */ + auto memorySelection(MemorySelection memorySelection) -> this_t & + { + memorySelection_impl(std::move(memorySelection)); + return *this; + } + + /** Set the number of elements the buffer can hold + * + * Optional. This tells the openPMD API the size of the buffer in elements. + * It is used to bound the operation's extent to the buffer when the extent + * is not set explicitly: for one-dimensional datasets, an otherwise full + * selection is shortened to the buffer size so that the buffer is not read + * from or written to beyond its bounds. + * It is set automatically when passing a contiguous container (see + * withContiguousContainer()) and can be set explicitly for raw pointers + * (see withRawPtr()), where the buffer size cannot be inferred. + * + * @param size Number of elements in the buffer + * @return Reference to this object for chaining + */ + auto bufferSize(size_t size) -> this_t & + { + bufferSize_impl(size); + return *this; + } + + // Enqueue method (deferred execution) + + /** Store the chunk data + * + * The returned handle performs the store, and the automatic flush of the + * underlying Series, when invoked via get() / operator()(). That flush is + * an MPI-collective operation, so in parallel codes every rank must invoke + * (or explicitly destroy) its handles in a consistent order. Invoking the + * handle is an explicit, user-controlled action and the return value is + * marked [[nodiscard]], so handles cannot be dropped unnoticed. Use + * unsafeNoAutomaticFlush() to defer flushing to a later explicit + * Series::flush() instead. + * + * @return Deferred computation that performs the store when invoked + */ + [[nodiscard]] auto store() -> auxiliary::DeferredComputation; + + /** This intentionally shadows the parent class's enqueueLoad methods in + * order to show a compile error when using load() on an object + * of this class. The parent method can still be accessed through + * typecasting if needed. + */ + template + auto load() + { + static_assert( + auxiliary::dependent_false_v, + "Cannot load chunk data into a buffer that is const or a " + "unique_ptr."); + } +}; + +/** Configuration for loading/storing chunks from/to a buffer. + * + * This class supports both loading and storing operations, allowing + * reading data into or writing data from a provided buffer. + */ +class ConfigureLoadStoreFromBuffer : public ConfigureStoreChunkFromBuffer +{ + friend class ConfigureLoadStore; + friend class RecordComponent; + + using ConfigureStoreChunkFromBuffer::ConfigureStoreChunkFromBuffer; + +public: + using this_t = ConfigureLoadStoreFromBuffer; + + // Configuration methods (always available) + + /** Set the offset within the dataset + * + * Optional. The operation will apply without offset by default (i.e. offset + * = (0, 0, ...)). + * + * @param offset Offset within the dataset + * @return Reference to this object for chaining + */ + auto offset(Offset offset) -> this_t & + { + offset_impl(std::move(offset)); + return *this; + } + + /** Set the extent within the dataset + * + * Optional. The operation will apply to the entire dataset by default (i.e. + * operation extent = global dataset extent - operation offset). + * + * @param extent Extent within the dataset, counted from the offset + * @return Reference to this object for chaining + */ + auto extent(Extent extent) -> this_t & + { + extent_impl(std::move(extent)); + return *this; + } + + /** Disable automatic flush after operation + * + * The returned objects of type DeferredComputation will still return a + * buffer upon get() / operator()(), but these buffers are not guaranteed to + * be filled until explicitly flushing. + * By default, invoking get() / operator()() would also flush the underlying + * Series, an MPI-collective operation. Disabling that automatic flush + * removes the need for all ranks to invoke their handles in a consistent + * order; the user is then responsible for calling Series::flush() + * collectively at a suitable point. + * + * @return Reference to this object for chaining + */ + auto unsafeNoAutomaticFlush() -> this_t & + { + unsafeNoAutomaticFlush_impl(); + return *this; + } + + /** Set memory selection for non-contiguous memory regions + * + * Only supported with ADIOS2 >= 2.10.1 (the capability to reset a memory + * selection was added upstream in 2.11.0 and backported to 2.10.1). Older + * versions cannot reset a memory selection once it has been set, which + * would silently leak it into subsequent store operations of the same + * variable. Those versions reject memory selections with an error at store + * time. + * + * @param memorySelection Selection of memory region + * @return Reference to this object for chaining + */ + auto memorySelection(MemorySelection memorySelection) -> this_t & + { + memorySelection_impl(std::move(memorySelection)); + return *this; + } + + /** Set the number of elements the buffer can hold + * + * Optional. This tells the openPMD API the size of the buffer in elements. + * It is used to bound the operation's extent to the buffer when the extent + * is not set explicitly: for one-dimensional datasets, an otherwise full + * selection is shortened to the buffer size so that data is not loaded + * beyond the buffer's bounds. + * It is set automatically when passing a contiguous container (see + * withContiguousContainer()) and can be set explicitly for raw pointers + * (see withRawPtr()), where the buffer size cannot be inferred. + * + * @param size Number of elements in the buffer + * @return Reference to this object for chaining + */ + auto bufferSize(size_t size) -> this_t & + { + bufferSize_impl(size); + return *this; + } + + // Enqueue method (deferred execution) + + /** Load the chunk data into the buffer + * + * The returned handle performs the load, and the automatic flush of the + * underlying Series, when invoked via get() / operator()(). That flush is + * an MPI-collective operation, so in parallel codes every rank must invoke + * (or explicitly destroy) its handles in a consistent order. Invoking the + * handle is an explicit, user-controlled action and the return value is + * marked [[nodiscard]], so handles cannot be dropped unnoticed. Use + * unsafeNoAutomaticFlush() to defer flushing to a later explicit + * Series::flush() instead. + * + * @return Deferred computation that performs the load when invoked + */ + [[nodiscard]] auto load() -> auxiliary::DeferredComputation; +}; + +} // namespace openPMD + +#include "openPMD/UndefDatatypeMacros.hpp" +// comment to prevent these includes from being moved by clang-format +#include "openPMD/LoadStoreChunk.tpp" diff --git a/include/openPMD/LoadStoreChunk.tpp b/include/openPMD/LoadStoreChunk.tpp new file mode 100644 index 0000000000..5d7ec83b0a --- /dev/null +++ b/include/openPMD/LoadStoreChunk.tpp @@ -0,0 +1,74 @@ +#pragma once + +#include "openPMD/LoadStoreChunk.hpp" + +namespace openPMD +{ +template +auto ConfigureLoadStore::withSharedPtr(std::shared_ptr data) + -> shared_ptr_return_type +{ + using T_decayed = std::remove_cv_t>; + constexpr auto dtype = determineDatatype(); + if constexpr (std::is_const_v) + { + return withSharedPtr_impl_const(data, dtype); + } + else + { + return withSharedPtr_impl_mut(data, dtype); + } +} + +template +auto ConfigureLoadStore::withUniquePtr(UniquePtrWithLambda data) + -> unique_ptr_return_type + +{ + using T_decayed = std::remove_cv_t>; + constexpr auto dtype = determineDatatype(); + if constexpr (std::is_const_v) + { + return withUniquePtr_impl_const( + std::move(data).template static_cast_(), dtype); + } + else + { + return withUniquePtr_impl_mut( + std::move(data).template static_cast_(), dtype); + } +} + +template +auto ConfigureLoadStore::withUniquePtr(std::unique_ptr data) + -> unique_ptr_return_type +{ + return withUniquePtr(UniquePtrWithLambda(std::move(data))); +} + +template +auto ConfigureLoadStore::withRawPtr(T *data) -> shared_ptr_return_type +{ + using T_decayed = std::remove_cv_t>; + constexpr auto dtype = determineDatatype(); + if constexpr (std::is_const_v) + { + return withRawPtr_impl_const(data, dtype); + } + else + { + return withRawPtr_impl_mut(data, dtype); + } +} + +template +auto ConfigureLoadStore::withContiguousContainer(T_ContiguousContainer &data) + -> std::enable_if_t< + auxiliary::IsContiguousContainer_v, + shared_ptr_return_type> +{ + auto res = withRawPtr(data.data()); + res.bufferSize(data.size()); + return res; +} +} // namespace openPMD diff --git a/include/openPMD/ParticleSpecies.hpp b/include/openPMD/ParticleSpecies.hpp index df6024db31..e436876b66 100644 --- a/include/openPMD/ParticleSpecies.hpp +++ b/include/openPMD/ParticleSpecies.hpp @@ -70,6 +70,11 @@ class ParticleSpecies namespace traits { + /** Generation policy for ParticleSpecies objects. + * + * Links particle patches to their parent hierarchy when a species is + * created. + */ template <> struct GenerationPolicy { diff --git a/include/openPMD/RecordComponent.hpp b/include/openPMD/RecordComponent.hpp index e8052d2f3e..ec66a8356d 100644 --- a/include/openPMD/RecordComponent.hpp +++ b/include/openPMD/RecordComponent.hpp @@ -22,6 +22,7 @@ #include "openPMD/Dataset.hpp" #include "openPMD/Datatype.hpp" +#include "openPMD/LoadStoreChunk.hpp" #include "openPMD/auxiliary/ShareRaw.hpp" #include "openPMD/auxiliary/TypeTraits.hpp" #include "openPMD/auxiliary/UniquePtr.hpp" @@ -30,9 +31,6 @@ #include "openPMD/backend/HierarchyVisitor.hpp" #include "openPMD/backend/scientific_defaults/ScientificDefaults.hpp" -// comment to prevent this include from being moved by clang-format -#include "openPMD/DatatypeMacros.hpp" - #include #include #include @@ -134,6 +132,12 @@ class RecordComponent friend T &internal::makeOwning(T &self, Series_type); friend class internal::ScientificDefaults; friend class Attributable; + friend class ConfigureLoadStore; + friend class ConfigureLoadStoreFromBuffer; + friend class ConfigureStoreChunkFromBuffer; + friend struct VisitorEnqueueLoadVariantWithoutFlush; + friend struct VisitorEnqueueLoadVariantWithFlush; + friend struct VisitorLoadVariant; public: enum class Allocation @@ -220,6 +224,16 @@ class RecordComponent */ bool empty() const; + /** Prepare a load/store chunk configuration object + * + * This is the entry point for the experimental new API for loading and + * storing chunks. It returns a ConfigureLoadStore object that can be used + * to specify offset, extent, and buffer for the operation. + * + * @return ConfigureLoadStore object for configuring the operation + */ + [[nodiscard]] ConfigureLoadStore prepareLoadStore(); + /** Load and allocate a chunk of data * * Set offset to {0u} and extent to {-1u} for full selection. @@ -230,11 +244,8 @@ class RecordComponent template std::shared_ptr loadChunk(Offset = {0u}, Extent = {-1u}); -#define OPENPMD_ENUMERATE_TYPES(type) , std::shared_ptr - using shared_ptr_dataset_types = auxiliary::detail::variant_tail_t< - auxiliary::detail::bottom OPENPMD_FOREACH_DATASET_DATATYPE( - OPENPMD_ENUMERATE_TYPES)>; -#undef OPENPMD_ENUMERATE_TYPES + using shared_ptr_dataset_types = + auxiliary::detail::shared_ptr_dataset_types; /** std::variant-based version of allocating loadChunk(Offset, Extent) * @@ -272,25 +283,6 @@ class RecordComponent template void loadChunk(std::shared_ptr data, Offset offset, Extent extent); - /** Load a chunk of data into pre-allocated memory, array version. - * - * @param data Preallocated, contiguous buffer, large enough to load the - * the requested data into it. - * The shared pointer must own and manage the buffer. - * Optimizations might be implemented based on this - * assumption (e.g. skipping the operation if the backend - * is the unique owner). - * The array-based overload helps avoid having to manually - * specify the delete[] destructor (C++17 feature). - * @param offset Offset within the dataset. Set to {0u} for full selection. - * @param extent Extent within the dataset, counted from the offset. - * Set to {-1u} for full selection. - * If offset is non-zero and extent is {-1u} the leftover - * extent in the record component will be selected. - */ - template - void loadChunk(std::shared_ptr data, Offset offset, Extent extent); - /** Load a chunk of data into pre-allocated memory, raw pointer version. * * @param data Preallocated, contiguous buffer, large enough to load the @@ -330,18 +322,6 @@ class RecordComponent template void storeChunk(std::shared_ptr data, Offset offset, Extent extent); - /** Store a chunk of data from a chunk of memory, array version. - * - * @param data Preallocated, contiguous buffer, large enough to read the - * the specified data from it. - * The array-based overload helps avoid having to manually - * specify the delete[] destructor (C++17 feature). - * @param offset Offset within the dataset. - * @param extent Extent within the dataset, counted from the offset. - */ - template - void storeChunk(std::shared_ptr data, Offset offset, Extent extent); - /** Store a chunk of data from a chunk of memory, unique pointer version. * * @param data Preallocated, contiguous buffer, large enough to read the @@ -505,8 +485,28 @@ class RecordComponent */ RecordComponent &makeEmpty(Dataset d); - void storeChunk( - auxiliary::WriteBuffer buffer, Datatype datatype, Offset o, Extent e); + void storeChunk_impl( + auxiliary::WriteBuffer buffer, + Datatype datatype, + internal::LoadStoreConfigWithBuffer); + + template + DynamicMemoryView storeChunkSpan_impl(internal::LoadStoreConfig); + template + DynamicMemoryView storeChunkSpanCreateBuffer_impl( + internal::LoadStoreConfig, F &&createBuffer); + + template + void loadChunk_impl( + std::shared_ptr const &, internal::LoadStoreConfigWithBuffer); + void loadChunk_impl( + std::shared_ptr const &, + Datatype, + internal::LoadStoreConfigWithBuffer); + template + std::shared_ptr loadChunkAllocate_impl(internal::LoadStoreConfig); + std::shared_ptr loadChunkAllocate_impl( + Datatype, size_t dtype_size, internal::LoadStoreConfig); // clang-format off OPENPMD_protected @@ -576,6 +576,4 @@ namespace internal } // namespace openPMD -#include "openPMD/UndefDatatypeMacros.hpp" -// comment to prevent these includes from being moved by clang-format #include "RecordComponent.tpp" diff --git a/include/openPMD/RecordComponent.tpp b/include/openPMD/RecordComponent.tpp index 951a520cfb..8c90e5f735 100644 --- a/include/openPMD/RecordComponent.tpp +++ b/include/openPMD/RecordComponent.tpp @@ -23,6 +23,7 @@ #include "openPMD/Datatype.hpp" #include "openPMD/Error.hpp" +#include "openPMD/LoadStoreChunk.hpp" #include "openPMD/RecordComponent.hpp" #include "openPMD/Span.hpp" #include "openPMD/auxiliary/Memory.hpp" @@ -32,6 +33,7 @@ #include "openPMD/backend/Attributable.hpp" #include +#include #include namespace openPMD @@ -41,8 +43,19 @@ template inline void RecordComponent::storeChunk(std::unique_ptr data, Offset o, Extent e) { - storeChunk( - UniquePtrWithLambda(std::move(data)), std::move(o), std::move(e)); + auto operation = prepareLoadStore(); + if (o.size() != 1u || o.at(0) != 0u) + { + operation.offset(std::move(o)); + } + if (e.size() != 1u || e.at(0) != -1u) + { + operation.extent(std::move(e)); + } + operation.withUniquePtr(std::move(data)) + .unsafeNoAutomaticFlush() + .store() + .get(); } template @@ -50,39 +63,48 @@ inline typename std::enable_if_t< auxiliary::IsContiguousContainer_v> RecordComponent::storeChunk(T_ContiguousContainer &data, Offset o, Extent e) { - uint8_t dim = getDimensionality(); + auto storeChunkConfig = prepareLoadStore(); - // default arguments - // offset = {0u}: expand to right dim {0u, 0u, ...} - Offset offset = o; - if (o.size() == 1u && o.at(0) == 0u) + // guard against default arguments + // we will take care of joined dimension handling later in computeOffset / + // computeExtent + if (o.size() != 1 || o.at(0) != 0u) { - if (joinedDimension().has_value()) - { - offset.clear(); - } - else if (dim > 1u) - { - offset = Offset(dim, 0u); - } + storeChunkConfig.offset(std::move(o)); + } + if (e.size() != 1 || e.at(0) != -1u) + { + storeChunkConfig.extent(std::move(e)); } - // extent = {-1u}: take full size - Extent extent(dim, 1u); - // avoid outsmarting the user: - // - stdlib data container implement 1D -> 1D chunk to write - if (e.size() == 1u && e.at(0) == -1u && dim == 1u) - extent.at(0) = data.size(); - else - extent = e; - - storeChunk(auxiliary::shareRaw(data.data()), offset, extent); + std::move(storeChunkConfig) + .withContiguousContainer(data) + .unsafeNoAutomaticFlush() + .store() + .get(); } template inline DynamicMemoryView RecordComponent::storeChunk(Offset o, Extent e, F &&createBuffer) { + auto operation = prepareLoadStore(); + if (o.size() != 1u || o.at(0) != 0u) + { + operation.offset(std::move(o)); + } + if (e.size() != 1u || e.at(0) != -1u) + { + operation.extent(std::move(e)); + } + return operation.storeSpan(std::forward(createBuffer)); +} + +template +inline DynamicMemoryView RecordComponent::storeChunkSpanCreateBuffer_impl( + internal::LoadStoreConfig cfg, F &&createBuffer) +{ + auto [o, e] = std::move(cfg); verifyChunk(o, e); size_t size = 1; @@ -193,4 +215,12 @@ inline auto RecordComponent::visit(Args &&...args) return switchDatasetType>( getDatatype(), *this, std::forward(args)...); } + +// definitions for LoadStoreChunk.hpp +template +auto ConfigureLoadStore::storeSpan(F &&createBuffer) -> DynamicMemoryView +{ + return m_rc->storeChunkSpanCreateBuffer_impl( + storeChunkConfig(), std::forward(createBuffer)); +} } // namespace openPMD diff --git a/include/openPMD/auxiliary/Defer.hpp b/include/openPMD/auxiliary/Defer.hpp index 8792a14ee0..7fccd906c0 100644 --- a/include/openPMD/auxiliary/Defer.hpp +++ b/include/openPMD/auxiliary/Defer.hpp @@ -6,6 +6,17 @@ namespace openPMD::auxiliary { +/** Defer wrapper + * + * Executes a functor when destroyed unless explicitly cancelled. + * Similar to Go's defer or C++'s experimental::scope_exit. + * + * Similar also to DeferredComputation under Future.hpp, but has another + * application scope (this: internal resource cleanup, that: public Future-like + * API) and is hence kept separate. + * + * @tparam F The functor type + */ template struct defer_type { @@ -51,8 +62,17 @@ struct defer_type auto operator=(defer_type const &) -> defer_type & = delete; }; +/** Type-erased defer wrapper for void functors */ using opaque_defer_type = defer_type>; +/** Create a defer wrapper + * + * Creates a defer wrapper that will execute the given functor when + * destroyed. + * + * @param functor The functor to execute on destruction + * @return A defer wrapper + */ template auto defer(F &&functor) -> defer_type> { diff --git a/include/openPMD/auxiliary/Future.hpp b/include/openPMD/auxiliary/Future.hpp new file mode 100644 index 0000000000..45f7f03083 --- /dev/null +++ b/include/openPMD/auxiliary/Future.hpp @@ -0,0 +1,157 @@ +#pragma once + +#include +#include +#include + +namespace openPMD::auxiliary::detail +{ +/** Internal helper for deferred computation - executes task once */ +template +struct OneTimeTask +{ + using task_type = std::function; + // Helper struct so we get auto-generated move constructor / assignment + // operator, but can still override constructors outside + struct Members + { + task_type m_task; + bool m_task_valid = true; + }; + Members members; + + static constexpr bool noexcept_move = + std::is_move_constructible_v && + std::is_move_assignable_v; + + explicit OneTimeTask(); + OneTimeTask(task_type); + + OneTimeTask(OneTimeTask &&) noexcept(noexcept_move); + OneTimeTask(OneTimeTask const &) = delete; + + auto operator=(OneTimeTask &&) noexcept(noexcept_move) -> OneTimeTask &; + auto operator=(OneTimeTask const &) -> OneTimeTask & = delete; + + auto operator()() -> T; +}; + +/** Internal helper for cached value storage. Used when the API requires + * creation of a DeferredComputation object, but there is not actually a + * computation to run. */ +template +struct CachedValue +{ + T val; +}; +template <> +struct CachedValue +{ + // this is silly +}; +} // namespace openPMD::auxiliary::detail + +namespace openPMD::auxiliary +{ +/** A computation that is deferred until explicitly invoked. + * + * This class wraps a callable, allowing lazy evaluation. + * The computation is performed once on first invocation, repeated invocation is + * an error. Check if the computation is still valid by calling valid(). + * + * Note: Some API operations may construct a DeferredComputation without any + * actual computation, instead emplacing a cached value. This is treated + * transparently to the user. In this case however, the object will not turn + * invalid upon invocation. + * + * @tparam T The return type of the computation + */ +template +class DeferredComputation +{ + using task_type = std::function; + using cached_type = std::conditional_t< + std::is_void_v, + // just something that is not void + detail::CachedValue, + T>; + std::variant, detail::CachedValue> m_task; + +public: + static constexpr bool noexcept_move = + std::is_move_constructible_v> && + std::is_move_assignable_v> && + std::is_move_constructible_v> && + std::is_move_assignable_v>; + /** Construct from a callable + * + * @param task The callable to execute + */ + DeferredComputation(task_type task); + /** Construct from a cached value + * + * @param val The pre-computed value + */ + DeferredComputation(cached_type val); + + explicit DeferredComputation(); + + DeferredComputation(DeferredComputation &&) noexcept(noexcept_move); + DeferredComputation(DeferredComputation const &) = delete; + + auto operator=(DeferredComputation &&) noexcept(noexcept_move) + -> DeferredComputation &; + auto operator=(DeferredComputation const &) + -> DeferredComputation & = delete; + + /** Destroy the computation, executing it if it has not been invoked yet. + * + * Computations returned by openPMD's data operations (see + * ConfigureLoadStore::store() and ConfigureLoadStore::load()) carry out + * the operation and, unless disabled via + * ConfigureLoadStore::unsafeNoAutomaticFlush(), the automatic flush of the + * underlying Series. That flush is an MPI-collective operation, hence + * destroying a still-valid handle without invoking it has the same + * collective semantics as calling get(). + * Exceptions thrown here cannot escape the destructor and are logged to + * standard error instead. + */ + ~DeferredComputation(); + + /** Get the result of the computation + * + * Invoking the computation executes the deferred data operation and, if + * automatic flushing was not disabled via + * ConfigureLoadStore::unsafeNoAutomaticFlush(), the automatic flush of the + * underlying Series. Since Series::flush() is MPI-collective, all ranks of + * a parallel Series must invoke (or explicitly destroy) their deferred + * computations in a consistent order. The handles returned by openPMD's + * store() / load() operations are marked [[nodiscard]] so that they cannot + * be dropped unnoticed. The flush itself is skipped if the pertaining + * RecordComponent has already been flushed by another component in the + * meantime (tracked via per-component flush counters), so a given + * invocation does not necessarily flush. + * + * @return The result of the computation + */ + auto get() -> T; + /** Invoke the computation + * + * Alias for get(). See get() for the collective semantics of the automatic + * flush. + * + * @return The result of the computation + */ + auto operator()() -> T; + + /** Discard the computation without executing it + */ + void invalidate() &&; + + /** Check if the computation is valid + * + * @return true if the computation has not been invalidated + */ + [[nodiscard]] auto valid() const noexcept -> bool; +}; +} // namespace openPMD::auxiliary diff --git a/include/openPMD/auxiliary/Memory.hpp b/include/openPMD/auxiliary/Memory.hpp index 6f8807b354..99b6b53d22 100644 --- a/include/openPMD/auxiliary/Memory.hpp +++ b/include/openPMD/auxiliary/Memory.hpp @@ -65,6 +65,7 @@ namespace auxiliary [[nodiscard]] auto release() -> UniquePtrWithLambda; }; using SharedPtr = std::shared_ptr; + using ReadSharedPtr = std::shared_ptr; /* * Use std::any publically since some compilers have trouble with * certain uses of std::variant, so hide it from them. @@ -73,17 +74,21 @@ namespace auxiliary */ std::any m_buffer; - WriteBuffer(); - WriteBuffer(std::shared_ptr ptr); - WriteBuffer(UniquePtrWithLambda ptr); + explicit WriteBuffer(); + // @todo implementation must distinguish const types + template + explicit WriteBuffer(std::shared_ptr ptr); + explicit WriteBuffer(UniquePtrWithLambda ptr); WriteBuffer(WriteBuffer &&) noexcept; WriteBuffer(WriteBuffer const &) = delete; WriteBuffer &operator=(WriteBuffer &&) noexcept; WriteBuffer &operator=(WriteBuffer const &) = delete; - WriteBuffer const &operator=(std::shared_ptr ptr); - WriteBuffer const &operator=(UniquePtrWithLambda ptr); + // @todo implementation must distinguish const types + template + WriteBuffer &operator=(std::shared_ptr const &ptr); + WriteBuffer &operator=(UniquePtrWithLambda ptr); void const *get() const; diff --git a/include/openPMD/auxiliary/Memory_internal.hpp b/include/openPMD/auxiliary/Memory_internal.hpp index bf3c8ccb4b..ee7134d9dc 100644 --- a/include/openPMD/auxiliary/Memory_internal.hpp +++ b/include/openPMD/auxiliary/Memory_internal.hpp @@ -25,6 +25,8 @@ namespace openPMD::auxiliary { // cannot use a unique_ptr inside a std::variant, so we represent it with this -using WriteBufferTypes = - std::variant; +using WriteBufferTypes = std::variant< + WriteBuffer::CopyableUniquePtr, + WriteBuffer::SharedPtr, + WriteBuffer::ReadSharedPtr>; } // namespace openPMD::auxiliary diff --git a/include/openPMD/auxiliary/UniquePtr.hpp b/include/openPMD/auxiliary/UniquePtr.hpp index 87f3261b45..84ea64db1d 100644 --- a/include/openPMD/auxiliary/UniquePtr.hpp +++ b/include/openPMD/auxiliary/UniquePtr.hpp @@ -24,6 +24,7 @@ #include #include #include +#include #include namespace openPMD @@ -176,10 +177,23 @@ template UniquePtrWithLambda UniquePtrWithLambda::static_cast_() && { using other_type = std::remove_extent_t; + auto original_ptr = this->release(); + auto original_ptr_casted = static_cast(original_ptr); return UniquePtrWithLambda{ - static_cast(this->release()), - [deleter = std::move(this->get_deleter())](other_type *ptr) { - deleter(static_cast(ptr)); + original_ptr_casted, + [deleter = std::move(this->get_deleter()), + original_ptr, + original_ptr_casted](other_type *casted_ptr) { + if (original_ptr_casted == casted_ptr) + { + deleter(original_ptr); + } + else + { + throw std::runtime_error( + "UniquePtrWithLambda: Pointer of casted pointer was " + "disassociated from its deleter."); + } }}; } } // namespace openPMD diff --git a/include/openPMD/backend/Attributable.hpp b/include/openPMD/backend/Attributable.hpp index b69e4e7a46..713b9b1c52 100644 --- a/include/openPMD/backend/Attributable.hpp +++ b/include/openPMD/backend/Attributable.hpp @@ -249,6 +249,7 @@ class Attributable friend struct internal::HomogenizeExtents; friend struct internal::ConfigAttribute; friend class internal::ScientificDefaults; + friend class ConfigureLoadStore; protected: // tag for internal constructor diff --git a/include/openPMD/backend/BaseRecordComponent.hpp b/include/openPMD/backend/BaseRecordComponent.hpp index af23e08468..f1303f2652 100644 --- a/include/openPMD/backend/BaseRecordComponent.hpp +++ b/include/openPMD/backend/BaseRecordComponent.hpp @@ -58,6 +58,15 @@ namespace internal */ bool m_datasetDefined = false; + /** Counter tracking the number of flush operations. This is later used + * to avoid repeated flushing in the DeferredComputation objects + * returned by the loadStoreChunk() API. (The counter is copied as a + * weak reference to the shared pointer, and the value is compared to + * the value upon enqueuing the operation. If the flush counter has + * proceeded past the old value, our operation has already been run.) */ + std::shared_ptr m_flushCounter = + std::make_shared(0); + BaseRecordComponentData(BaseRecordComponentData const &) = delete; BaseRecordComponentData(BaseRecordComponentData &&) = delete; BaseRecordComponentData & diff --git a/src/Datatype.cpp b/src/Datatype.cpp index 479286066c..ead2859c81 100644 --- a/src/Datatype.cpp +++ b/src/Datatype.cpp @@ -29,13 +29,6 @@ namespace openPMD { -void warnWrongDtype(std::string const &key, Datatype store, Datatype request) -{ - std::cerr << "Warning: Attribute '" << key << "' stored as " << store - << ", requested as " << request - << ". Casting unconditionally with possible loss of precision.\n"; -} - std::ostream &operator<<(std::ostream &os, openPMD::Datatype const &d) { using DT = openPMD::Datatype; diff --git a/src/IO/ADIOS/ADIOS2File.cpp b/src/IO/ADIOS/ADIOS2File.cpp index 0260836136..6507b49ff5 100644 --- a/src/IO/ADIOS/ADIOS2File.cpp +++ b/src/IO/ADIOS/ADIOS2File.cpp @@ -23,6 +23,7 @@ #include "openPMD/Error.hpp" #include "openPMD/IO/ADIOS/ADIOS2Auxiliary.hpp" #include "openPMD/IO/ADIOS/ADIOS2IOHandler.hpp" +#include "openPMD/IO/ADIOS/macros.hpp" #include "openPMD/IO/AbstractIOHandler.hpp" #include "openPMD/IterationEncoding.hpp" #include "openPMD/auxiliary/Environment.hpp" @@ -70,6 +71,7 @@ void DatasetReader::call( adios2::Variable var = impl->verifyDataset( bp.param.offset, bp.param.extent, + std::nullopt, IO, engine, bp.name, @@ -98,7 +100,9 @@ void WriteDataset::call(ADIOS2File &ba, detail::BufferedPut &bp) std::visit( [&](auto &&arg) { using ptr_type = std::decay_t; - if constexpr (std::is_same_v>) + if constexpr ( + std::is_same_v> || + std::is_same_v>) { auto ptr = static_cast(arg.get()); auto &engine = ba.getEngine(); @@ -106,6 +110,7 @@ void WriteDataset::call(ADIOS2File &ba, detail::BufferedPut &bp) adios2::Variable var = ba.m_impl->verifyDataset( bp.param.offset, bp.param.extent, + bp.param.memorySelection, ba.m_IO, engine, bp.name, @@ -113,6 +118,20 @@ void WriteDataset::call(ADIOS2File &ba, detail::BufferedPut &bp) ba.variables()); engine.Put(var, ptr); + if (bp.param.memorySelection.has_value()) + { + /* + * Reset the memory selection so that it does not leak into + * the next put of this variable. ADIOS2 versions that + * cannot reset a memory selection reject the selection + * already in verifyDataset(), so this branch is only ever + * reached with supporting versions. + */ + if constexpr (openPMD::CanTheMemorySelectionBeReset) + { + var.SetMemorySelection(); + } + } } else if constexpr ( std::is_same_v< @@ -123,6 +142,14 @@ void WriteDataset::call(ADIOS2File &ba, detail::BufferedPut &bp) bput.name = std::move(bp.name); bput.offset = std::move(bp.param.offset); bput.extent = std::move(bp.param.extent); + bput.memorySelection = std::move(bp.param.memorySelection); + /* + * Note: Moving is required here since it's a unique_ptr. + * std::forward<>() would theoretically work, but it + * requires the type parameter and we don't have that + * inside the lambda. + * (ptr_type does not work for this case). + */ bput.data = arg.release(); bput.dtype = bp.param.dtype; ba.m_uniquePtrPuts.push_back(std::move(bput)); @@ -170,12 +197,27 @@ struct RunUniquePtrPut adios2::Variable var = ba.m_impl->verifyDataset( bufferedPut.offset, bufferedPut.extent, + bufferedPut.memorySelection, ba.m_IO, engine, bufferedPut.name, std::nullopt, ba.variables()); engine.Put(var, ptr); + if (bufferedPut.memorySelection.has_value()) + { + /* + * Reset the memory selection so that it does not leak into the + * next put of this variable. ADIOS2 versions that cannot reset a + * memory selection reject the selection already in + * verifyDataset(), so this branch is only ever reached with + * supporting versions. + */ + if constexpr (openPMD::CanTheMemorySelectionBeReset) + { + var.SetMemorySelection(); + } + } } static constexpr char const *errorMsg = "RunUniquePtrPut"; diff --git a/src/IO/ADIOS/ADIOS2IOHandler.cpp b/src/IO/ADIOS/ADIOS2IOHandler.cpp index 036b4f6098..a2539e42cd 100644 --- a/src/IO/ADIOS/ADIOS2IOHandler.cpp +++ b/src/IO/ADIOS/ADIOS2IOHandler.cpp @@ -51,6 +51,7 @@ #include #include #include +#include #include #include #include @@ -1247,6 +1248,7 @@ namespace detail adios2::Variable variable = impl->verifyDataset( params.offset, params.extent, + std::nullopt, IO, engine, varName, @@ -2643,7 +2645,7 @@ ADIOS2IOHandler::ADIOS2IOHandler( {} std::future -ADIOS2IOHandler::flush(internal::ParsedFlushParams &flushParams) +ADIOS2IOHandler::flush_impl(internal::ParsedFlushParams &flushParams) { return m_impl.flush(flushParams); } @@ -2684,7 +2686,7 @@ ADIOS2IOHandler::ADIOS2IOHandler( std::move(initialize_from), std::move(path), at, std::move(config)) {} -std::future ADIOS2IOHandler::flush(internal::ParsedFlushParams &) +std::future ADIOS2IOHandler::flush_impl(internal::ParsedFlushParams &) { return std::future(); } diff --git a/src/IO/AbstractIOHandler.cpp b/src/IO/AbstractIOHandler.cpp index 564eebba01..27a8d49751 100644 --- a/src/IO/AbstractIOHandler.cpp +++ b/src/IO/AbstractIOHandler.cpp @@ -145,11 +145,26 @@ std::future AbstractIOHandler::flush(internal::FlushParams const ¶ms) return future; } +std::future AbstractIOHandler::flush(internal::ParsedFlushParams ¶ms) +{ + auto res = this->flush_impl(params); + if (!m_work.empty()) + { + throw error::Internal("flush() did not clear all work!"); + } + return res; +} + bool AbstractIOHandler::fullSupportForVariableBasedEncoding() const { return false; } +bool AbstractIOHandler::supportsMemorySelection() const +{ + return false; +} + #if openPMD_HAVE_MPI template <> AbstractIOHandler::AbstractIOHandler( diff --git a/src/IO/AbstractIOHandlerImpl.cpp b/src/IO/AbstractIOHandlerImpl.cpp index 153601677c..cc34aa0291 100644 --- a/src/IO/AbstractIOHandlerImpl.cpp +++ b/src/IO/AbstractIOHandlerImpl.cpp @@ -113,6 +113,7 @@ namespace case Operation::LIST_PATHS: case Operation::OPEN_PATH: case Operation::SET_WRITTEN: + case Operation::INCREASE_FLUSH_COUNTER: case Operation::CREATE_PATH: break; case Operation::CLOSE_PATH: @@ -340,10 +341,26 @@ std::future AbstractIOHandlerImpl::flush(FlushLevel l) i.writable->parent, "->", i.writable, - "] WRITE_DATASET, offset=", - [¶meter]() { return vec_as_string(parameter.offset); }, - ", extent=", - [¶meter]() { return vec_as_string(parameter.extent); }); + "] WRITE_DATASET: ", + [&]() { + std::stringstream stream; + stream << "offset: " << vec_as_string(parameter.offset) + << " extent: " << vec_as_string(parameter.extent) + << " mem-selection: "; + if (parameter.memorySelection.has_value()) + { + stream << vec_as_string( + parameter.memorySelection->offset) + << "--" + << vec_as_string( + parameter.memorySelection->extent); + } + else + { + stream << "NONE"; + } + return stream.str(); + }); writeDataset(i.writable, parameter); break; } @@ -531,6 +548,19 @@ std::future AbstractIOHandlerImpl::flush(FlushLevel l) setWritten(i.writable, parameter); break; } + case O::INCREASE_FLUSH_COUNTER: { + auto ¶meter = + deref_dynamic_cast>( + i.parameter.get()); + writeToStderr( + "[", + i.writable->parent, + "->", + i.writable, + "] INCREASE_FLUSH_COUNTER "); + increaseFlushCounter(i.writable, parameter); + break; + } } } catch (...) @@ -600,4 +630,10 @@ void AbstractIOHandlerImpl::setWritten( { w->written = param.target_status; } + +void AbstractIOHandlerImpl::increaseFlushCounter( + Writable *, Parameter const ¶m) +{ + ++*param.flush_counter; +} } // namespace openPMD diff --git a/src/IO/DummyIOHandler.cpp b/src/IO/DummyIOHandler.cpp index f3b4e155d2..52c83bb624 100644 --- a/src/IO/DummyIOHandler.cpp +++ b/src/IO/DummyIOHandler.cpp @@ -39,7 +39,7 @@ DummyIOHandler::DummyIOHandler(std::string path, Access at) void DummyIOHandler::enqueue(IOTask const &) {} -std::future DummyIOHandler::flush(internal::ParsedFlushParams &) +std::future DummyIOHandler::flush_impl(internal::ParsedFlushParams &) { return std::future(); } diff --git a/src/IO/HDF5/HDF5IOHandler.cpp b/src/IO/HDF5/HDF5IOHandler.cpp index 3096b2297c..a6bb015eb3 100644 --- a/src/IO/HDF5/HDF5IOHandler.cpp +++ b/src/IO/HDF5/HDF5IOHandler.cpp @@ -1918,6 +1918,12 @@ void HDF5IOHandlerImpl::writeDataset( "[HDF5] Writing into a dataset in a file opened as read only is " "not possible."); + if (parameters.memorySelection.has_value()) + { + throw error::OperationUnsupportedInBackend( + "HDF5", + "Non-contiguous memory selections not supported in HDF5 backend."); + } File file = requireFile("writeDataset", writable, /* checkParent = */ true); herr_t status; @@ -3608,7 +3614,7 @@ HDF5IOHandler::HDF5IOHandler( HDF5IOHandler::~HDF5IOHandler() = default; -std::future HDF5IOHandler::flush(internal::ParsedFlushParams ¶ms) +std::future HDF5IOHandler::flush_impl(internal::ParsedFlushParams ¶ms) { return m_impl->flush(params); } @@ -3627,7 +3633,7 @@ HDF5IOHandler::HDF5IOHandler( HDF5IOHandler::~HDF5IOHandler() = default; -std::future HDF5IOHandler::flush(internal::ParsedFlushParams &) +std::future HDF5IOHandler::flush_impl(internal::ParsedFlushParams &) { return std::future(); } diff --git a/src/IO/HDF5/ParallelHDF5IOHandler.cpp b/src/IO/HDF5/ParallelHDF5IOHandler.cpp index 7de4960feb..0cf8b396a3 100644 --- a/src/IO/HDF5/ParallelHDF5IOHandler.cpp +++ b/src/IO/HDF5/ParallelHDF5IOHandler.cpp @@ -76,7 +76,7 @@ ParallelHDF5IOHandler::ParallelHDF5IOHandler( ParallelHDF5IOHandler::~ParallelHDF5IOHandler() = default; std::future -ParallelHDF5IOHandler::flush(internal::ParsedFlushParams ¶ms) +ParallelHDF5IOHandler::flush_impl(internal::ParsedFlushParams ¶ms) { if (auto hdf5_config_it = params.backendConfig.json().find("hdf5"); hdf5_config_it != params.backendConfig.json().end()) @@ -462,7 +462,8 @@ ParallelHDF5IOHandler::ParallelHDF5IOHandler( ParallelHDF5IOHandler::~ParallelHDF5IOHandler() = default; -std::future ParallelHDF5IOHandler::flush(internal::ParsedFlushParams &) +std::future +ParallelHDF5IOHandler::flush_impl(internal::ParsedFlushParams &) { return std::future(); } diff --git a/src/IO/IOTask.cpp b/src/IO/IOTask.cpp index af12ff766c..87fb747d34 100644 --- a/src/IO/IOTask.cpp +++ b/src/IO/IOTask.cpp @@ -123,6 +123,9 @@ std::ostream &operator<<(std::ostream &os, Operation op) case Operation::SET_WRITTEN: os << "SET_WRITTEN"; break; + case Operation::INCREASE_FLUSH_COUNTER: + os << "INCREASE_FLUSH_COUNTER"; + break; } return os; } @@ -309,6 +312,9 @@ namespace internal case Operation::SET_WRITTEN: return "SET_WRITTEN"; break; + case Operation::INCREASE_FLUSH_COUNTER: + return "INCREASE_FLUSH_COUNTER"; + break; } return "unknown"; } diff --git a/src/IO/JSON/JSONIOHandler.cpp b/src/IO/JSON/JSONIOHandler.cpp index 9af9a06728..c26964e75c 100644 --- a/src/IO/JSON/JSONIOHandler.cpp +++ b/src/IO/JSON/JSONIOHandler.cpp @@ -53,7 +53,7 @@ JSONIOHandler::JSONIOHandler( {} #endif -std::future JSONIOHandler::flush(internal::ParsedFlushParams ¶ms) +std::future JSONIOHandler::flush_impl(internal::ParsedFlushParams ¶ms) { return m_impl.flush(params); } diff --git a/src/IO/JSON/JSONIOHandlerImpl.cpp b/src/IO/JSON/JSONIOHandlerImpl.cpp index d83a719d8b..18351c1341 100644 --- a/src/IO/JSON/JSONIOHandlerImpl.cpp +++ b/src/IO/JSON/JSONIOHandlerImpl.cpp @@ -1147,6 +1147,13 @@ void JSONIOHandlerImpl::writeDataset( access::write(m_handler->m_backendAccess), "[JSON] Cannot write data in read-only mode."); + if (parameters.memorySelection.has_value()) + { + throw error::OperationUnsupportedInBackend( + "JSON", + "Non-contiguous memory selections not supported in JSON backend."); + } + auto pos = setAndGetFilePosition(writable); auto file = refreshFileFromParent(writable); auto &j = obtainJsonContents(writable); diff --git a/src/LoadStoreChunk.cpp b/src/LoadStoreChunk.cpp new file mode 100644 index 0000000000..0356db0ad9 --- /dev/null +++ b/src/LoadStoreChunk.cpp @@ -0,0 +1,467 @@ + + +#include "openPMD/LoadStoreChunk.hpp" +#include "openPMD/Datatype.hpp" +#include "openPMD/Error.hpp" +#include "openPMD/RecordComponent.hpp" +#include "openPMD/Span.hpp" +#include "openPMD/auxiliary/Future.hpp" +#include "openPMD/auxiliary/Memory.hpp" +#include "openPMD/auxiliary/Memory_internal.hpp" +#include "openPMD/auxiliary/ShareRawInternal.hpp" +#include "openPMD/auxiliary/StringManip.hpp" +#include "openPMD/auxiliary/UniquePtr.hpp" + +// comment to keep clang-format from reordering +#include "openPMD/DatatypeMacros.hpp" +#include "openPMD/backend/Attributable.hpp" + +#include +#include +#include +#include + +namespace openPMD +{ +namespace +{ + template + auto asWriteBuffer(std::shared_ptr &&ptr) -> auxiliary::WriteBuffer + { + /* std::static_pointer_cast correctly reference-counts the pointer */ + return auxiliary::WriteBuffer( + std::static_pointer_cast(std::move(ptr))); + } + template + auto asWriteBuffer(UniquePtrWithLambda &&ptr) -> auxiliary::WriteBuffer + { + return auxiliary::WriteBuffer( + std::move(ptr).template static_cast_()); + } + + /* + * There is no backend support currently for const unique pointers. + * We support these mostly for providing a clean API to users that have such + * pointers and want to store from them, but there will be no + * backend-specific optimizations for such buffers as there are for + * non-const unique pointers. + */ + template + auto asWriteBuffer(UniquePtrWithLambda &&ptr) + -> auxiliary::WriteBuffer + { + auto raw_ptr = ptr.release(); + return asWriteBuffer( + std::shared_ptr{ + raw_ptr, + [deleter = std::move(ptr.get_deleter())]( + auto const *delete_me) { deleter(delete_me); }}); + } +} // namespace + +ConfigureLoadStore::ConfigureLoadStore(RecordComponent &rc) + : m_rc(std::make_shared(rc)) +{} +ConfigureLoadStore::ConfigureLoadStore(ConfigureLoadStore const &other) = + default; +ConfigureLoadStore & +ConfigureLoadStore::operator=(ConfigureLoadStore const &other) = default; +ConfigureLoadStore::ConfigureLoadStore(ConfigureLoadStore &&) = default; +ConfigureLoadStore & +ConfigureLoadStore::operator=(ConfigureLoadStore &&) = default; + +auto ConfigureLoadStore::dim() const -> uint8_t +{ + return m_rc->getDimensionality(); +} + +auto ConfigureLoadStore::storeChunkConfig() -> internal::LoadStoreConfig +{ + return internal::LoadStoreConfig{computeOffset(), computeExtent()}; +} + +auto ConfigureLoadStore::deferFlush(RecordComponent &attr) +{ + if (m_unsafeNoAutomaticFlush) + { + throw error::Internal( + "Configuring an automatic flush operating after configuring that " + "those should be switched off."); + } + auto ioHandler = attr.IOHandler(); + if (!ioHandler) + { + throw error::Internal( + "Cannot configure automatic flush: the underlying Series is " + "already closed."); + } + auto index = attr.get().m_flushCounter; + return [attr, + old_index = *index, + current_index = std::weak_ptr(index)]() mutable { + auto lock_current_index = current_index.lock(); + if (!lock_current_index || *lock_current_index > old_index) + { + return; + } + attr.seriesFlush(); + }; +} + +auto ConfigureLoadStore::computeOffset() -> Offset const & +{ + auto joined_dim = m_rc->joinedDimension(); + if (!m_offset.has_value()) + { + if (joined_dim.has_value()) + { + m_offset = std::make_optional(); + } + else + { + m_offset = std::make_optional(dim(), 0); + } + } + else if (joined_dim.has_value() && !m_offset->empty()) + { + throw error::WrongAPIUsage( + "Joined dimension: Must specify empty offset."); + } + return *m_offset; +} + +auto ConfigureLoadStore::computeExtent() -> Extent const & +{ + bool allow_downsizing = false; + if (!m_extent.has_value()) + { + m_extent = std::make_optional(m_rc->getExtent()); + if (m_offset.has_value()) + { + auto it_o = m_offset->begin(); + auto end_o = m_offset->end(); + auto it_e = m_extent->begin(); + auto end_e = m_extent->end(); + for (; it_o != end_o && it_e != end_e; ++it_e, ++it_o) + { + *it_e -= *it_o; + } + } + allow_downsizing = true; + } + if (auto buffer_size = getBufferSize(); buffer_size.has_value()) + { + size_t requestedExtent = std::accumulate( + m_extent->begin(), + m_extent->end(), + 1, + [](size_t l, size_t r) -> size_t { return r * l; }); + if (requestedExtent > *buffer_size) + { + if (m_extent->size() == 1 && allow_downsizing) + { + (*m_extent)[0] = *buffer_size; + } + else + { + std::stringstream error; + error << "Requesting to load/store a chunk of size " + << requestedExtent << " (n-dimensional extent is "; + auxiliary::write_vec_to_stream(error, *m_extent) + << ") to/from a buffer of size " << *buffer_size << "."; + throw error::WrongAPIUsage(error.str()); + } + } + } + return *m_extent; +} + +auto ConfigureLoadStore::withSharedPtr_impl_mut( + std::shared_ptr data, Datatype datatype) + -> openPMD::ConfigureLoadStoreFromBuffer +{ + if (!data) + { + throw std::runtime_error( + "Unallocated pointer passed during chunk store."); + } + return openPMD::ConfigureLoadStoreFromBuffer( + auxiliary::WriteBuffer(std::move(data)), datatype, {std::move(*this)}); +} +auto ConfigureLoadStore::withSharedPtr_impl_const( + std::shared_ptr data, Datatype datatype) + -> openPMD::ConfigureStoreChunkFromBuffer +{ + if (!data) + { + throw std::runtime_error( + "Unallocated pointer passed during chunk store."); + } + return openPMD::ConfigureStoreChunkFromBuffer( + auxiliary::WriteBuffer(std::move(data)), datatype, {std::move(*this)}); +} + +auto ConfigureLoadStore::withUniquePtr_impl_mut( + UniquePtrWithLambda data, Datatype dtype) + -> openPMD::ConfigureStoreChunkFromBuffer + +{ + if (!data) + { + throw std::runtime_error( + "Unallocated pointer passed during chunk store."); + } + + return openPMD::ConfigureStoreChunkFromBuffer( + auxiliary::WriteBuffer(std::move(data)), dtype, {std::move(*this)}); +} +auto ConfigureLoadStore::withUniquePtr_impl_const( + UniquePtrWithLambda data, Datatype dtype) + -> openPMD::ConfigureStoreChunkFromBuffer + +{ + if (!data) + { + throw std::runtime_error( + "Unallocated pointer passed during chunk store."); + } + + void const *raw_ptr = data.release(); + auto &deleter = data.get_deleter(); + return openPMD::ConfigureStoreChunkFromBuffer( + auxiliary::WriteBuffer( + std::shared_ptr( + raw_ptr, + [deleter_lambda = std::move(deleter)](auto const *p) { + deleter_lambda(p); + })), + dtype, + {std::move(*this)}); +} + +auto ConfigureLoadStore::withRawPtr_impl_mut(void *data, Datatype dtype) + -> openPMD::ConfigureLoadStoreFromBuffer +{ + if (!data) + { + throw std::runtime_error( + "Unallocated pointer passed during chunk store."); + } + return openPMD::ConfigureLoadStoreFromBuffer( + auxiliary::WriteBuffer(auxiliary::shareRaw(data)), + dtype, + {std::move(*this)}); +} + +auto ConfigureLoadStore::withRawPtr_impl_const(void const *data, Datatype dtype) + -> openPMD::ConfigureStoreChunkFromBuffer +{ + if (!data) + { + throw std::runtime_error( + "Unallocated pointer passed during chunk store."); + } + return openPMD::ConfigureStoreChunkFromBuffer( + auxiliary::WriteBuffer(auxiliary::shareRaw(data)), + dtype, + {std::move(*this)}); +} + +template +auto ConfigureLoadStore::storeSpan() -> DynamicMemoryView +{ + return m_rc->storeChunkSpan_impl(storeChunkConfig()); +} + +template +auto ConfigureLoadStore::load() + -> auxiliary::DeferredComputation> +{ + auto res = m_rc->loadChunkAllocate_impl(storeChunkConfig()); + if (m_unsafeNoAutomaticFlush) + { + return auxiliary::DeferredComputation>( + std::move(res)); + } + return auxiliary::DeferredComputation>( + [res_lambda = std::move(res), dflush = deferFlush(*m_rc)]() mutable { + dflush(); + return res_lambda; + }); +} + +struct VisitorEnqueueLoadVariantWithFlush +{ + template + static auto + call(RecordComponent &rc, internal::LoadStoreConfig cfg, F &&dflush) + -> auxiliary::DeferredComputation< + auxiliary::detail::shared_ptr_dataset_types> + { + auto res = rc.loadChunkAllocate_impl(std::move(cfg)); + return auxiliary::DeferredComputation< + auxiliary::detail::shared_ptr_dataset_types>( + [res_lambda = std::move(res), + dflush_lambda = std::forward(dflush)]() mutable + -> auxiliary::detail::shared_ptr_dataset_types { + dflush_lambda(); + return res_lambda; + }); + } +}; +struct VisitorEnqueueLoadVariantWithoutFlush +{ + template + static auto call(RecordComponent &rc, internal::LoadStoreConfig cfg) + -> auxiliary::DeferredComputation< + auxiliary::detail::shared_ptr_dataset_types> + { + auto res = rc.loadChunkAllocate_impl(std::move(cfg)); + return auxiliary::DeferredComputation< + auxiliary::detail::shared_ptr_dataset_types>(std::move(res)); + } +}; + +auto ConfigureLoadStore::loadVariant() -> auxiliary::DeferredComputation< + auxiliary::detail::shared_ptr_dataset_types> +{ + if (m_unsafeNoAutomaticFlush) + { + return m_rc->visit( + this->storeChunkConfig()); + } + else + { + return m_rc->visit( + this->storeChunkConfig(), deferFlush(*m_rc)); + } +} + +auto ConfigureLoadStore::getComponentHandle() const -> RecordComponent +{ + return *m_rc; +} + +struct VisitorLoadVariant +{ + template + static auto call(RecordComponent &rc, internal::LoadStoreConfig cfg) + -> auxiliary::detail::shared_ptr_dataset_types + { + return rc.loadChunkAllocate_impl(std::move(cfg)); + } +}; + +ConfigureStoreChunkFromBuffer::ConfigureStoreChunkFromBuffer( + auxiliary::WriteBuffer buffer, Datatype dt, ConfigureLoadStore &&core) + : ConfigureLoadStore(std::move(core)) + , m_buffer(std::move(buffer)) + , m_datatype(dt) +{} + +auto ConfigureStoreChunkFromBuffer::storeChunkConfig() + -> internal::LoadStoreConfigWithBuffer +{ + return internal::LoadStoreConfigWithBuffer{ + this->computeOffset(), this->computeExtent(), m_mem_select}; +} + +void ConfigureStoreChunkFromBuffer::bufferSize_impl(size_t size) +{ + m_buffer_size = size; +} + +auto ConfigureStoreChunkFromBuffer::getBufferSize() -> std::optional +{ + return m_buffer_size; +} + +auto ConfigureStoreChunkFromBuffer::store() + -> auxiliary::DeferredComputation +{ + this->m_rc->storeChunk_impl( + std::move(m_buffer), m_datatype, storeChunkConfig()); + if (m_unsafeNoAutomaticFlush) + { + return auxiliary::DeferredComputation( + auxiliary::detail::CachedValue()); + } + return auxiliary::DeferredComputation( + [dflush = deferFlush(*m_rc)]() mutable -> void { dflush(); }); +} + +auto ConfigureLoadStoreFromBuffer::load() + -> auxiliary::DeferredComputation +{ + auto *shared_ptr = std::get_if( + &this->m_buffer.as_variant()); + if (!shared_ptr) + { + throw std::runtime_error( + "ConfigureLoadStoreFromBuffer must be instantiated with a " + "non-const shared_ptr type."); + } + this->m_rc->loadChunk_impl( + *shared_ptr, m_datatype, this->storeChunkConfig()); + if (m_unsafeNoAutomaticFlush) + { + return auxiliary::DeferredComputation( + auxiliary::detail::CachedValue()); + } + return auxiliary::DeferredComputation( + [dflush = this->deferFlush(*this->m_rc)]() mutable -> void { + dflush(); + }); +} + +void ConfigureLoadStore::extent_impl(Extent extent) +{ + m_extent = std::make_optional(std::move(extent)); +} + +void ConfigureLoadStore::offset_impl(Offset offset) +{ + m_offset = std::make_optional(std::move(offset)); +} + +void ConfigureLoadStore::unsafeNoAutomaticFlush_impl() +{ + m_unsafeNoAutomaticFlush = true; +} + +auto ConfigureLoadStore::getBufferSize() -> std::optional +{ + return std::nullopt; +} + +void ConfigureStoreChunkFromBuffer::memorySelection_impl(MemorySelection sel) +{ + m_mem_select = std::make_optional(std::move(sel)); +} +// namespace core + +// need this for clang-tidy +#define OPENPMD_ARRAY(type) type[] +#define OPENPMD_POINTER(type) type * +#define OPENPMD_APPLY_TEMPLATE(template_, type) template_ + +#define INSTANTIATE_METHOD_TEMPLATES(dtype) \ + template auto ConfigureLoadStore::load() \ + -> auxiliary::DeferredComputation; +#define INSTANTIATE_METHOD_TEMPLATES_WITH_AND_WITHOUT_EXTENT(type) \ + INSTANTIATE_METHOD_TEMPLATES(type) \ + INSTANTIATE_METHOD_TEMPLATES(OPENPMD_ARRAY(type)) \ + template auto ConfigureLoadStore::storeSpan() -> DynamicMemoryView; + +OPENPMD_FOREACH_DATASET_DATATYPE( + INSTANTIATE_METHOD_TEMPLATES_WITH_AND_WITHOUT_EXTENT) + +#undef INSTANTIATE_METHOD_TEMPLATES +#undef INSTANTIATE_METHOD_TEMPLATES_WITH_AND_WITHOUT_EXTENT + +#undef INSTANTIATE_METHOD_TEMPLATES +#undef OPENPMD_ARRAY +#undef OPENPMD_POINTER +#undef OPENPMD_APPLY_TEMPLATE +} // namespace openPMD diff --git a/src/RecordComponent.cpp b/src/RecordComponent.cpp index 42738e1222..f8df85a242 100644 --- a/src/RecordComponent.cpp +++ b/src/RecordComponent.cpp @@ -24,9 +24,11 @@ #include "openPMD/Error.hpp" #include "openPMD/IO/AbstractIOHandler.hpp" #include "openPMD/IO/Format.hpp" +#include "openPMD/LoadStoreChunk.hpp" #include "openPMD/Series.hpp" #include "openPMD/auxiliary/Environment.hpp" #include "openPMD/auxiliary/Memory.hpp" +#include "openPMD/auxiliary/ShareRawInternal.hpp" #include "openPMD/auxiliary/StringManip.hpp" #include "openPMD/backend/Attributable.hpp" #include "openPMD/backend/BaseRecord.hpp" @@ -37,6 +39,9 @@ // comment so clang-format does not move this #include "openPMD/DatatypeMacros.hpp" +// comment +#include "openPMD/DatatypeMacros.hpp" + #include #include #include @@ -192,6 +197,73 @@ auto resource(T &t) -> attribute_types & return t.template resource(); } +ConfigureLoadStore RecordComponent::prepareLoadStore() +{ + return ConfigureLoadStore{*this}; +} + +namespace +{ +#if (defined(_LIBCPP_VERSION) && _LIBCPP_VERSION < 11000) || \ + (defined(__apple_build_version__) && __clang_major__ < 14) + template + auto createSpanBufferFallback(size_t size) -> UniquePtrWithLambda + { + return UniquePtrWithLambda{ + new T[size], [](auto *ptr) { delete[] ptr; }}; + } +#else + template + auto createSpanBufferFallback(size_t size) -> std::unique_ptr + { + return std::unique_ptr{new T[size]}; + } +#endif +} // namespace + +template +DynamicMemoryView +RecordComponent::storeChunkSpan_impl(internal::LoadStoreConfig cfg) +{ + return storeChunkSpanCreateBuffer_impl( + std::move(cfg), &createSpanBufferFallback); +} + +template +std::shared_ptr +RecordComponent::loadChunkAllocate_impl(internal::LoadStoreConfig cfg) +{ + using T = std::remove_cv_t>; + auto res = loadChunkAllocate_impl( + determineDatatype(), sizeof(T), std::move(cfg)); + return std::static_pointer_cast(res); +} + +std::shared_ptr RecordComponent::loadChunkAllocate_impl( + Datatype dtype, size_t dtype_size, internal::LoadStoreConfig cfg) +{ + auto [o, e] = std::move(cfg); + + size_t numPoints = 1; + for (auto val : e) + { + numPoints *= val; + } + + auto newData = + std::shared_ptr(new char[numPoints * dtype_size], [](void *p) { + delete[] (static_cast(p)); + }); + prepareLoadStore() + .offset(std::move(o)) + .extent(std::move(e)) + .withSharedPtr_impl_mut(newData, dtype) + .unsafeNoAutomaticFlush() + .load() + .get(); + return newData; +} + RecordComponent::RecordComponent() : BaseRecordComponent(NoInit()) { setData(std::make_shared()); @@ -421,11 +493,15 @@ void RecordComponent::flush( { return; } - if (access::readOnly(IOHandler()->m_frontendAccess)) + auto ioHandler = IOHandler(); + Parameter increaseFlushCounter; + increaseFlushCounter.flush_counter = rc.m_flushCounter; + if (access::readOnly(ioHandler->m_frontendAccess)) { + ioHandler->enqueue(IOTask(this, increaseFlushCounter)); while (!rc.m_chunks.empty()) { - IOHandler()->enqueue(rc.m_chunks.front()); + ioHandler->enqueue(rc.m_chunks.front()); rc.m_chunks.pop(); } } @@ -460,6 +536,8 @@ void RecordComponent::flush( return val == Dataset::JOINED_DIMENSION; }); }; + ioHandler->enqueue(IOTask(this, increaseFlushCounter)); + if (!written()) { if (constant()) @@ -468,7 +546,7 @@ void RecordComponent::flush( IterationEncoding::variableBased; Parameter pCreate; pCreate.path = name; - IOHandler()->enqueue(IOTask(this, pCreate)); + ioHandler->enqueue(IOTask(this, pCreate)); Parameter aWrite; aWrite.name = "value"; aWrite.dtype = rc.m_constantValue.dtype; @@ -478,7 +556,7 @@ void RecordComponent::flush( aWrite.changesOverSteps = Parameter< Operation::WRITE_ATT>::ChangesOverSteps::IfPossible; } - IOHandler()->enqueue(IOTask(this, aWrite)); + ioHandler->enqueue(IOTask(this, aWrite)); if (constant_component_write_shape()) { aWrite.name = "shape"; @@ -490,7 +568,7 @@ void RecordComponent::flush( aWrite.changesOverSteps = Parameter< Operation::WRITE_ATT>::ChangesOverSteps::IfPossible; } - IOHandler()->enqueue(IOTask(this, aWrite)); + ioHandler->enqueue(IOTask(this, aWrite)); } } else @@ -498,7 +576,7 @@ void RecordComponent::flush( Parameter dCreate( rc.m_dataset.value()); dCreate.name = name; - IOHandler()->enqueue(IOTask(this, dCreate)); + ioHandler->enqueue(IOTask(this, dCreate)); } } @@ -525,20 +603,20 @@ void RecordComponent::flush( aWrite.changesOverSteps = Parameter< Operation::WRITE_ATT>::ChangesOverSteps::IfPossible; } - IOHandler()->enqueue(IOTask(this, aWrite)); + ioHandler->enqueue(IOTask(this, aWrite)); } else { Parameter pExtend( rc.m_dataset.value().extent); - IOHandler()->enqueue(IOTask(this, std::move(pExtend))); + ioHandler->enqueue(IOTask(this, std::move(pExtend))); rc.m_hasBeenExtended = false; } } while (!rc.m_chunks.empty()) { - IOHandler()->enqueue(rc.m_chunks.front()); + ioHandler->enqueue(rc.m_chunks.front()); rc.m_chunks.pop(); } @@ -629,14 +707,61 @@ void RecordComponent::readBase() } } -void RecordComponent::storeChunk( - auxiliary::WriteBuffer buffer, Datatype dtype, Offset o, Extent e) +void RecordComponent::storeChunk_impl( + auxiliary::WriteBuffer buffer, + Datatype dtype, + internal::LoadStoreConfigWithBuffer cfg) { + auto [o, e, memorySelection] = std::move(cfg); verifyChunk(dtype, o, e); + if (memorySelection.has_value()) + { + auto const &mem_offset = memorySelection->offset; + auto const &mem_extent = memorySelection->extent; + if (mem_offset.size() != e.size() || mem_extent.size() != e.size()) + { + throw error::WrongAPIUsage( + "Memory selection: dimensionality of memory offset and memory " + "extent must match the chunk extent."); + } + for (size_t i = 0; i < e.size(); ++i) + { + if (mem_offset[i] + e[i] > mem_extent[i]) + { + throw error::WrongAPIUsage( + "Memory selection: memory offset + chunk extent exceeds " + "the memory extent (dimension " + + std::to_string(i) + ")."); + } + } + } Parameter dWrite; dWrite.offset = std::move(o); dWrite.extent = std::move(e); + if (memorySelection.has_value()) + { + if (joinedDimension().has_value()) + { + throw error::WrongAPIUsage( + "Memory selections are not supported for joined arrays."); + } + /* + * Reject unsupported memory selections here, while enqueueing the + * chunk, rather than letting the backend throw at flush time: an + * exception raised inside an IO task makes AbstractIOHandlerImpl::flush + * clear the entire IO queue, losing all other pending chunks and + * attributes and leaving an unreadable file behind. + */ + auto *ioHandler = IOHandler(); + if (ioHandler != nullptr && !ioHandler->supportsMemorySelection()) + { + throw error::OperationUnsupportedInBackend( + ioHandler->backendName(), + "Non-contiguous memory selections are not supported."); + } + } + dWrite.memorySelection = memorySelection; dWrite.dtype = dtype; /* std::static_pointer_cast correctly reference-counts the pointer */ dWrite.data = std::move(buffer); @@ -769,69 +894,82 @@ RecordComponent &RecordComponent::makeEmpty(uint8_t dimensions) template std::shared_ptr RecordComponent::loadChunk(Offset o, Extent e) { - uint8_t dim = getDimensionality(); + auto operation = prepareLoadStore(); // default arguments + // we will take care of joined dimension handling later in computeOffset / + // computeExtent // offset = {0u}: expand to right dim {0u, 0u, ...} - Offset offset = o; - if (o.size() == 1u && o.at(0) == 0u && dim > 1u) - offset = Offset(dim, 0u); + if (o.size() != 1u || o.at(0) != 0u) + { + operation.offset(std::move(o)); + } // extent = {-1u}: take full size - Extent extent(dim, 1u); - if (e.size() == 1u && e.at(0) == -1u) + if (e.size() != 1u || e.at(0) != -1u) { - extent = getExtent(); - for (uint8_t i = 0u; i < dim; ++i) - extent[i] -= offset[i]; + operation.extent(std::move(e)); } - else - extent = e; - uint64_t numPoints = 1u; - for (auto const &dimensionSize : extent) - numPoints *= dimensionSize; - -#if (defined(_LIBCPP_VERSION) && _LIBCPP_VERSION < 11000) || \ - (defined(__apple_build_version__) && __clang_major__ < 14) - auto newData = - std::shared_ptr(new T[numPoints], [](T *p) { delete[] p; }); - loadChunk(newData, offset, extent); - return newData; -#else - auto newData = std::shared_ptr[]>( - new std::remove_extent_t[numPoints]); - loadChunk(newData, offset, extent); - return std::static_pointer_cast(std::move(newData)); -#endif + return operation.unsafeNoAutomaticFlush().load().get(); } namespace detail { - template - struct do_convert + struct FillBuffer { - template - static std::optional call(Attribute &attr) + template + static void call( + void *target, + size_t numPoints, + RecordComponent const &component, + internal::RecordComponentData const &rc) { - if constexpr (std::is_convertible_v) + std::optional val = rc.m_constantValue.getOptional(); + + if (val.has_value()) { - return std::make_optional(attr.get()); + auto raw_ptr = static_cast(target); + std::fill(raw_ptr, raw_ptr + numPoints, *val); } else { - return std::nullopt; + std::string const data_type_str = + datatypeToString(component.getDatatype()); + std::string const requ_type_str = + datatypeToString(determineDatatype()); + std::string err_msg = + "Type conversion during chunk loading not possible! "; + err_msg += + "Data: " + data_type_str + "; Load as: " + requ_type_str; + throw error::WrongAPIUsage(err_msg); } } - static constexpr char const *errorMsg = "is_conversible"; + static constexpr char const *errorMsg = "FillBuffer"; }; } // namespace detail template -void RecordComponent::loadChunk(std::shared_ptr data, Offset o, Extent e) +void RecordComponent::loadChunk_impl( + std::shared_ptr const &data, internal::LoadStoreConfigWithBuffer cfg) +{ + loadChunk_impl( + std::static_pointer_cast(data), + determineDatatype>>(), + std::move(cfg)); +} + +void RecordComponent::loadChunk_impl( + std::shared_ptr const &data, + Datatype dtype_requested, + internal::LoadStoreConfigWithBuffer cfg) { - Datatype dtype = determineDatatype(data); + if (cfg.memorySelection.has_value()) + { + throw error::WrongAPIUsage( + "Unsupported: Memory selections in chunk loading."); + } /* * For constant components, we implement type conversion, so there is * a separate check further below. @@ -842,36 +980,28 @@ void RecordComponent::loadChunk(std::shared_ptr data, Offset o, Extent e) * * Attention: Do NOT use operator==(), doesnt work properly on Windows! */ - if (!isSame(dtype, getDatatype()) && !constant()) + if (!isSame(dtype_requested, getDatatype()) && !constant()) { std::string const data_type_str = datatypeToString(getDatatype()); - std::string const requ_type_str = - datatypeToString(determineDatatype()); + std::string const requ_type_str = datatypeToString(dtype_requested); std::string err_msg = "Type conversion during chunk loading not yet implemented! "; err_msg += "Data: " + data_type_str + "; Load as: " + requ_type_str; throw std::runtime_error(err_msg); } - uint8_t dim = getDimensionality(); + auto dim = getDimensionality(); + auto [offset, extent, memorySelection] = std::move(cfg); - // default arguments - // offset = {0u}: expand to right dim {0u, 0u, ...} - Offset offset = o; - if (o.size() == 1u && o.at(0) == 0u && dim > 1u) - offset = Offset(dim, 0u); - - // extent = {-1u}: take full size - Extent extent(dim, 1u); - if (e.size() == 1u && e.at(0) == -1u) + if (joinedDimension().has_value()) { - extent = getExtent(); - for (uint8_t i = 0u; i < dim; ++i) - extent[i] -= offset[i]; + throw error::WrongAPIUsage( + "Cannot load chunks of a joined array: the extent of the joined " + "dimension is not known during the write session. Close and " + "reopen the Series to read back the joined data."); } - else - extent = e; + Extent dse = getExtent(); if (extent.size() != dim || offset.size() != dim) { std::ostringstream oss; @@ -882,16 +1012,12 @@ void RecordComponent::loadChunk(std::shared_ptr data, Offset o, Extent e) << "do not match."; throw std::runtime_error(oss.str()); } - Extent dse = getExtent(); for (uint8_t i = 0; i < dim; ++i) if (dse[i] < offset[i] + extent[i]) throw std::runtime_error( "Chunk does not reside inside dataset (Dimension on index " + std::to_string(i) + ". DS: " + std::to_string(dse[i]) + " - Chunk: " + std::to_string(offset[i] + extent[i]) + ")"); - if (!data) - throw std::runtime_error( - "Unallocated pointer passed during chunk loading."); auto &rc = get(); if (constant()) @@ -900,25 +1026,8 @@ void RecordComponent::loadChunk(std::shared_ptr data, Offset o, Extent e) for (auto const &dimensionSize : extent) numPoints *= dimensionSize; - std::optional val = - switchNonVectorType>( - /* dt = */ getDatatype(), rc.m_constantValue); - - if (val.has_value()) - { - T *raw_ptr = data.get(); - std::fill(raw_ptr, raw_ptr + numPoints, *val); - } - else - { - std::string const data_type_str = datatypeToString(getDatatype()); - std::string const requ_type_str = - datatypeToString(determineDatatype()); - std::string err_msg = - "Type conversion during chunk loading not possible! "; - err_msg += "Data: " + data_type_str + "; Load as: " + requ_type_str; - throw error::WrongAPIUsage(err_msg); - } + switchDatasetType( + dtype_requested, data.get(), numPoints, *this, rc); } else { @@ -932,13 +1041,30 @@ void RecordComponent::loadChunk(std::shared_ptr data, Offset o, Extent e) } template -void RecordComponent::loadChunk( - std::shared_ptr ptr, Offset offset, Extent extent) +void RecordComponent::loadChunk(std::shared_ptr data, Offset o, Extent e) { - loadChunk( - std::static_pointer_cast(std::move(ptr)), - std::move(offset), - std::move(extent)); + // static_assert(!std::is_same_v, "EVIL"); + auto operation = prepareLoadStore(); + + // default arguments + // we will take care of joined dimension handling later in computeOffset / + // computeExtent + // offset = {0u}: expand to right dim {0u, 0u, ...} + if (o.size() != 1u || o.at(0) != 0u) + { + operation.offset(std::move(o)); + } + + // extent = {-1u}: take full size + if (e.size() != 1u || e.at(0) != -1u) + { + operation.extent(std::move(e)); + } + + operation.withSharedPtr(std::move(data)) + .unsafeNoAutomaticFlush() + .load() + .get(); } template @@ -950,62 +1076,71 @@ void RecordComponent::loadChunkRaw(T *ptr, Offset offset, Extent extent) template void RecordComponent::storeChunk(std::shared_ptr data, Offset o, Extent e) { - if (!data) - throw std::runtime_error( - "Unallocated pointer passed during chunk store."); - Datatype dtype = determineDatatype(data); - - /* std::static_pointer_cast correctly reference-counts the pointer */ - storeChunk( - auxiliary::WriteBuffer(std::static_pointer_cast(data)), - dtype, - std::move(o), - std::move(e)); + auto operation = prepareLoadStore(); + // default arguments + // we will take care of joined dimension handling later in computeOffset / + // computeExtent + if (o.size() != 1u || o.at(0) != 0u) + { + operation.offset(std::move(o)); + } + if (e.size() != 1u || e.at(0) != -1u) + { + operation.extent(std::move(e)); + } + operation.withSharedPtr(std::move(data)) + .unsafeNoAutomaticFlush() + .store() + .get(); } template void RecordComponent::storeChunk( UniquePtrWithLambda data, Offset o, Extent e) { - if (!data) - throw std::runtime_error( - "Unallocated pointer passed during chunk store."); - Datatype dtype = determineDatatype<>(data); - - storeChunk( - auxiliary::WriteBuffer{std::move(data).template static_cast_()}, - dtype, - std::move(o), - std::move(e)); -} - -template -void RecordComponent::storeChunk(std::shared_ptr data, Offset o, Extent e) -{ - storeChunk( - std::static_pointer_cast(std::move(data)), - std::move(o), - std::move(e)); + auto operation = prepareLoadStore(); + if (o.size() != 1u || o.at(0) != 0u) + { + operation.offset(std::move(o)); + } + if (e.size() != 1u || e.at(0) != -1u) + { + operation.extent(std::move(e)); + } + operation.withUniquePtr(std::move(data)) + .unsafeNoAutomaticFlush() + .store() + .get(); } template void RecordComponent::storeChunkRaw(T const *ptr, Offset offset, Extent extent) { - storeChunk(auxiliary::shareRaw(ptr), std::move(offset), std::move(extent)); + auto operation = prepareLoadStore(); + if (offset.size() != 1u || offset.at(0) != 0u) + { + operation.offset(std::move(offset)); + } + if (extent.size() != 1u || extent.at(0) != -1u) + { + operation.extent(std::move(extent)); + } + operation.withRawPtr(ptr).unsafeNoAutomaticFlush().store().get(); } template DynamicMemoryView RecordComponent::storeChunk(Offset offset, Extent extent) { - return storeChunk(std::move(offset), std::move(extent), [](size_t size) { -#if (defined(_LIBCPP_VERSION) && _LIBCPP_VERSION < 11000) || \ - (defined(__apple_build_version__) && __clang_major__ < 14) - return UniquePtrWithLambda{ - new T[size], [](auto *ptr) { delete[] ptr; }}; -#else - return std::unique_ptr{new T[size]}; -#endif - }); + auto operation = prepareLoadStore(); + if (offset.size() != 1u || offset.at(0) != 0u) + { + operation.offset(std::move(offset)); + } + if (extent.size() != 1u || extent.at(0) != -1u) + { + operation.extent(std::move(extent)); + } + return operation.storeSpan(); } template @@ -1019,10 +1154,6 @@ void RecordComponent::verifyChunk(Offset const &o, Extent const &e) const #define OPENPMD_ARRAY(type) type[] #define OPENPMD_INSTANTIATE_BASIC(type) \ - template void RecordComponent::loadChunk( \ - std::shared_ptr data, Offset o, Extent e); \ - template void RecordComponent::loadChunk( \ - std::shared_ptr data, Offset o, Extent e); \ template void RecordComponent::loadChunkRaw( \ OPENPMD_PTR(type) ptr, Offset offset, Extent extent); \ template void RecordComponent::verifyChunk( \ @@ -1030,21 +1161,28 @@ void RecordComponent::verifyChunk(Offset const &o, Extent const &e) const template DynamicMemoryView RecordComponent::storeChunk( \ Offset offset, Extent extent); \ template void RecordComponent::storeChunkRaw( \ - OPENPMD_PTR(type const) ptr, Offset offset, Extent extent); + OPENPMD_PTR(type const) ptr, Offset offset, Extent extent); \ + template DynamicMemoryView RecordComponent::storeChunkSpan_impl( \ + internal::LoadStoreConfig cfg); -#define OPENPMD_INSTANTIATE_CONST_AND_NONCONST(type) \ - template void RecordComponent::storeChunk( \ - std::shared_ptr data, Offset o, Extent e); \ - template void RecordComponent::storeChunk( \ - std::shared_ptr data, Offset o, Extent e); +#define OPENPMD_INSTANTIATE_CONST_AND_NONCONST(type) #define OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT(type) \ + template void RecordComponent::loadChunk( \ + std::shared_ptr data, Offset o, Extent e); \ template std::shared_ptr RecordComponent::loadChunk( \ Offset o, Extent e); \ template void RecordComponent::storeChunk( \ - UniquePtrWithLambda data, Offset o, Extent e); + UniquePtrWithLambda data, Offset o, Extent e); \ + template void RecordComponent::loadChunk_impl( \ + std::shared_ptr const &data, \ + internal::LoadStoreConfigWithBuffer cfg); \ + template std::shared_ptr RecordComponent::loadChunkAllocate_impl( \ + internal::LoadStoreConfig cfg); #define OPENPMD_INSTANTIATE_FULLMATRIX(type) \ + template void RecordComponent::storeChunk( \ + std::shared_ptr data, Offset o, Extent e); \ template RecordComponent &RecordComponent::makeConstant(type); \ template RecordComponent &RecordComponent::makeEmpty( \ uint8_t dimensions); diff --git a/src/Series.cpp b/src/Series.cpp index b08996c2c4..dd470d1c08 100644 --- a/src/Series.cpp +++ b/src/Series.cpp @@ -539,12 +539,13 @@ void Series::flushRankTable(FlushLevel l, Attributable &attributable) }; auto writeDataset = [&rank, &maxSize, this, &attributable]( - std::shared_ptr put, size_t num_lines = 1) { + std::shared_ptr const &put, + size_t num_lines = 1) { Parameter chunk; chunk.dtype = Datatype::CHAR; chunk.offset = {uint64_t(rank), 0}; chunk.extent = {num_lines, maxSize}; - chunk.data = std::move(put); + chunk.data = put; IOHandler()->enqueue(IOTask(&attributable, std::move(chunk))); }; @@ -584,7 +585,7 @@ void Series::flushRankTable(FlushLevel l, Attributable &attributable) * > } */ [asRawPtr](char *) { delete asRawPtr; }}; - writeDataset(std::move(put), /* num_lines = */ size); + writeDataset(put, /* num_lines = */ size); } // Must ensure that the Writable is consistently set to written on all @@ -602,7 +603,7 @@ void Series::flushRankTable(FlushLevel l, Attributable &attributable) new char[maxSize]{}, [](char const *ptr) { delete[] ptr; }}; std::copy_n(myRankInfo.c_str(), mySize, put.get()); - writeDataset(std::move(put)); + writeDataset(put); } std::string Series::particlesPath() const diff --git a/src/auxiliary/Future.cpp b/src/auxiliary/Future.cpp new file mode 100644 index 0000000000..39af555f9e --- /dev/null +++ b/src/auxiliary/Future.cpp @@ -0,0 +1,190 @@ +#include "openPMD/auxiliary/Future.hpp" +#include "openPMD/Error.hpp" +#include "openPMD/RecordComponent.hpp" + +#include +#include +#include + +// comment + +#include "openPMD/DatatypeMacros.hpp" + +namespace openPMD::auxiliary::detail +{ +template +OneTimeTask::OneTimeTask() = default; + +template +OneTimeTask::OneTimeTask(task_type task) : members{std::move(task)} +{} + +template +OneTimeTask::OneTimeTask(OneTimeTask &&other) noexcept(noexcept_move) + : members(std::move(other.members)) +{ + other.members.m_task_valid = false; +} + +template +auto OneTimeTask::operator=(OneTimeTask &&other) noexcept(noexcept_move) + -> OneTimeTask & +{ + this->members = std::move(other.members); + other.members.m_task_valid = false; + return *this; +} + +template +auto OneTimeTask::operator()() -> T +{ + if (!members.m_task_valid) + { + throw error::WrongAPIUsage( + "[DeferredComputation] No valid state. Probably already " + "computed."); + } + if (!members.m_task) + { + throw error::WrongAPIUsage( + "[DeferredComputation] No valid task was specified."); + } + members.m_task_valid = false; + if constexpr (std::is_void_v) + { + std::move(members.m_task)(); + members.m_task = {}; + } + else + { + auto res = std::move(members.m_task)(); + members.m_task = {}; // reset + return res; + } +} +} // namespace openPMD::auxiliary::detail + +namespace openPMD::auxiliary +{ + +template +DeferredComputation::DeferredComputation(task_type task) + : m_task(detail::OneTimeTask{std::move(task)}) +{} + +template +DeferredComputation::DeferredComputation(cached_type cached_val) + : m_task(detail::CachedValue{std::move(cached_val)}) +{} + +template +DeferredComputation::DeferredComputation() = default; + +template +DeferredComputation::DeferredComputation(DeferredComputation &&) noexcept( + noexcept_move) = default; + +template +auto DeferredComputation::operator=(DeferredComputation &&) noexcept( + noexcept_move) -> DeferredComputation & = default; + +template +DeferredComputation::~DeferredComputation() +{ + try + { + std::visit( + auxiliary::overloaded{ + [](detail::OneTimeTask &task) { + if (task.members.m_task_valid) + { + std::move(task)(); + } + }, + [](detail::CachedValue &) {}}, + this->m_task); + } + catch (std::exception const &e) + { + std::cerr << "[DeferredComputation] Error in destructor: '" << e.what() + << "'." << std::endl; + } + catch (...) + { + std::cerr << "[DeferredComputation] Unknown error in destructor." + << std::endl; + } +} + +template +auto DeferredComputation::get() -> T +{ + return std::visit( + auxiliary::overloaded{ + [](detail::OneTimeTask &task) -> T { return std::move(task)(); }, + [](detail::CachedValue &cached) -> T { return cached.val; }}, + this->m_task); +} + +template <> +auto DeferredComputation::get() -> void +{ + std::visit( + auxiliary::overloaded{ + [](detail::OneTimeTask &task) { std::move(task)(); }, + [](detail::CachedValue &) { return; }}, + this->m_task); +} + +template +auto DeferredComputation::operator()() -> T +{ + return get(); +} + +template +void DeferredComputation::invalidate() && +{ + std::visit( + auxiliary::overloaded{ + [](detail::OneTimeTask &task) { + task.members.m_task = {}; + task.members.m_task_valid = false; + }, + [](detail::CachedValue const &) {}}, + this->m_task); +} + +template +auto DeferredComputation::valid() const noexcept -> bool +{ + return std::visit( + auxiliary::overloaded{ + [](detail::OneTimeTask const &task) { + return task.members.m_task_valid; + }, + [](detail::CachedValue const &) { return true; }}, + this->m_task); +} + +template class DeferredComputation; +template class DeferredComputation; +template class DeferredComputation; // used in tests + +// need this for clang-tidy +#define OPENPMD_ARRAY(type) type[] +#define OPENPMD_APPLY_TEMPLATE(template_, type) template_ + +#define INSTANTIATE_FUTURE(dtype) \ + template class DeferredComputation; +#define INSTANTIATE_FUTURE_WITH_AND_WITHOUT_EXTENT(type) \ + INSTANTIATE_FUTURE(type) INSTANTIATE_FUTURE(OPENPMD_ARRAY(type)) +OPENPMD_FOREACH_NONVECTOR_DATATYPE(INSTANTIATE_FUTURE_WITH_AND_WITHOUT_EXTENT) +#undef INSTANTIATE_FUTURE +#undef INSTANTIATE_FUTURE_WITH_AND_WITHOUT_EXTENT +#undef OPENPMD_ARRAY +#undef OPENPMD_APPLY_TEMPLATE +} // namespace openPMD::auxiliary + +#include "openPMD/UndefDatatypeMacros.hpp" diff --git a/src/auxiliary/Memory.cpp b/src/auxiliary/Memory.cpp index c2a0f2aa0d..b002288a88 100644 --- a/src/auxiliary/Memory.cpp +++ b/src/auxiliary/Memory.cpp @@ -21,8 +21,10 @@ #include "openPMD/auxiliary/Memory.hpp" #include "openPMD/ChunkInfo.hpp" +#include "openPMD/Datatype.tpp" #include "openPMD/auxiliary/Memory_internal.hpp" #include "openPMD/auxiliary/UniquePtr.hpp" +#include "openPMD/backend/Variant_internal.hpp" #include #include @@ -193,8 +195,21 @@ auto WriteBuffer::CopyableUniquePtr::release() -> UniquePtrWithLambda WriteBuffer::WriteBuffer() : m_buffer(std::make_any()) {} -WriteBuffer::WriteBuffer(std::shared_ptr ptr) - : m_buffer(std::make_any(std::move(ptr))) +template +WriteBuffer::WriteBuffer(std::shared_ptr ptr) + : m_buffer //(std::make_any(std::move(ptr))) + ([&]() { + if constexpr (std::is_const_v) + { + return std::make_any( + std::static_pointer_cast(ptr)); + } + else + { + return std::make_any( + std::static_pointer_cast(ptr)); + } + }()) {} WriteBuffer::WriteBuffer(UniquePtrWithLambda ptr) : m_buffer( @@ -204,12 +219,22 @@ WriteBuffer::WriteBuffer(UniquePtrWithLambda ptr) WriteBuffer::WriteBuffer(WriteBuffer &&) noexcept = default; WriteBuffer &WriteBuffer::operator=(WriteBuffer &&) noexcept = default; -WriteBuffer const &WriteBuffer::operator=(std::shared_ptr ptr) +template +WriteBuffer &WriteBuffer::operator=(std::shared_ptr const &ptr) { - m_buffer = std::make_any(std::move(ptr)); + if constexpr (std::is_const_v) + { + m_buffer = std::make_any( + std::static_pointer_cast(ptr)); + } + else + { + m_buffer = std::make_any( + std::static_pointer_cast(ptr)); + } return *this; } -WriteBuffer const &WriteBuffer::operator=(UniquePtrWithLambda ptr) +WriteBuffer &WriteBuffer::operator=(UniquePtrWithLambda ptr) { m_buffer = std::make_any(CopyableUniquePtr(std::move(ptr))); @@ -226,4 +251,22 @@ void const *WriteBuffer::get() const }, as_variant()); } + +#define OPENPMD_INSTANTIATE(dtype) \ + template WriteBuffer::WriteBuffer(std::shared_ptr); \ + template WriteBuffer &WriteBuffer::operator=( \ + std::shared_ptr const &); + +#ifndef DOXYGEN_SHOULD_SKIP_THIS + +OPENPMD_FOREACH_DATASET_DATATYPE(OPENPMD_INSTANTIATE) +template WriteBuffer::WriteBuffer(std::shared_ptr); +template WriteBuffer &WriteBuffer::operator=(std::shared_ptr const &); +template WriteBuffer::WriteBuffer(std::shared_ptr); +template WriteBuffer & +WriteBuffer::operator=(std::shared_ptr const &); + +#endif /* DOXYGEN_SHOULD_SKIP_THIS */ + +#undef OPENPMD_INSTANTIATE } // namespace openPMD::auxiliary diff --git a/src/auxiliary/UniquePtr.cpp b/src/auxiliary/UniquePtr.cpp index 6625ca3a47..828c33a785 100644 --- a/src/auxiliary/UniquePtr.cpp +++ b/src/auxiliary/UniquePtr.cpp @@ -58,7 +58,7 @@ namespace auxiliary OPENPMD_FOREACH_DATASET_DATATYPE( OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT) - OPENPMD_INSTANTIATE(void) + OPENPMD_INSTANTIATE(void) OPENPMD_INSTANTIATE(void const) #undef OPENPMD_INSTANTIATE #undef OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT @@ -99,12 +99,16 @@ UniquePtrWithLambda::UniquePtrWithLambda( std::unique_ptr); #define OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT(type) \ - OPENPMD_INSTANTIATE(type) OPENPMD_INSTANTIATE(OPENPMD_ARRAY(type)) + OPENPMD_INSTANTIATE(type) \ + OPENPMD_INSTANTIATE(OPENPMD_ARRAY(type)) \ + OPENPMD_INSTANTIATE(type const) \ + OPENPMD_INSTANTIATE(OPENPMD_ARRAY(type const)) -OPENPMD_FOREACH_DATASET_DATATYPE(OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT) +OPENPMD_FOREACH_NONVECTOR_DATATYPE(OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT) // Instantiate this directly, do not instantiate the // `std::unique_ptr`-based constructor. template class UniquePtrWithLambda; +template class UniquePtrWithLambda; #undef OPENPMD_INSTANTIATE #undef OPENPMD_INSTANTIATE_WITH_AND_WITHOUT_EXTENT #undef OPENPMD_ARRAY diff --git a/src/binding/python/PatchRecordComponent.cpp b/src/binding/python/PatchRecordComponent.cpp index 5887ba7d03..c68cdbc16f 100644 --- a/src/binding/python/PatchRecordComponent.cpp +++ b/src/binding/python/PatchRecordComponent.cpp @@ -130,7 +130,9 @@ void init_PatchRecordComponent(py::module &m) switch (dtype) { case DT::BOOL: - return prc.store(idx, *static_cast(buf.ptr)); + throw std::runtime_error( + "make_constant: " + "Boolean type not supported!"); break; case DT::SHORT: return prc.store(idx, *static_cast(buf.ptr)); diff --git a/test/AuxiliaryTest.cpp b/test/AuxiliaryTest.cpp index ee0b029473..3612b13adb 100644 --- a/test/AuxiliaryTest.cpp +++ b/test/AuxiliaryTest.cpp @@ -19,6 +19,8 @@ * If not, see . */ // expose private and protected members for invasive testing +#include "openPMD/Error.hpp" +#include "openPMD/auxiliary/Future.hpp" #if openPMD_USE_INVASIVE_TESTS #define OPENPMD_private public: #define OPENPMD_protected public: @@ -538,3 +540,31 @@ TEST_CASE("filesystem_test", "[auxiliary]") REQUIRE(!remove_file("./nonexistent_file_in_cmake_bin_directory")); #endif } + +TEST_CASE("future_test", "[auxiliary]") +{ + using task_type = auxiliary::DeferredComputation; + size_t counter = 0; + + auto make_task = [&counter]() { + counter = 0; + return task_type{[&counter]() { + ++counter; + return "success"; + }}; + }; + + auto move_construct = make_task(); + task_type move_constructed(std::move(move_construct)); + REQUIRE(counter == 0); + REQUIRE(move_constructed() == "success"); + REQUIRE(counter == 1); + REQUIRE_THROWS_AS(move_constructed(), error::WrongAPIUsage); + + auto move_assign = make_task(); + task_type move_assigned = std::move(move_assign); + REQUIRE(counter == 0); + REQUIRE(move_assigned() == "success"); + REQUIRE(counter == 1); + REQUIRE_THROWS_AS(move_assigned(), error::WrongAPIUsage); +} diff --git a/test/CoreTest.cpp b/test/CoreTest.cpp index 821907b38e..3888e1dbe1 100644 --- a/test/CoreTest.cpp +++ b/test/CoreTest.cpp @@ -19,6 +19,7 @@ * If not, see . */ // expose private and protected members for invasive testing +#include #if openPMD_USE_INVASIVE_TESTS #define OPENPMD_private public: #define OPENPMD_protected public: @@ -1286,7 +1287,7 @@ TEST_CASE("use_count_test", "[core]") pprc.resetDataset(Dataset(determineDatatype(), {4})); pprc.store(0, static_cast(1)); REQUIRE( - std::get>( + std::get>( static_cast *>( pprc.get().m_chunks.front().parameter.get()) ->data.as_variant()) @@ -1785,6 +1786,24 @@ TEST_CASE("unique_ptr", "[core]") UniquePtrWithLambda arrptrFilled{new int[5]{}}; UniquePtrWithLambda arrptrFilledCustom{ new int[5]{}, [](int const *p) { delete[] p; }}; + + auto ptr3 = UniquePtrWithLambda(new int{5}).static_cast_(); + auto ptr4 = UniquePtrWithLambda(new int{6}).static_cast_(); + auto ptr5 = new int{7}; + ptr3.swap(ptr4); + REQUIRE(*reinterpret_cast(ptr3.get()) == 6); + REQUIRE(*reinterpret_cast(ptr4.get()) == 5); + // cannot hand a new pointer to a casted UniquePtrWithLambda + // so let's check that it throws + ptr3.reset(ptr5); + // (we could theoretically reinterpret_cast the new pointer back to the + // original type and apply the deleter, but let's not.) + REQUIRE_THROWS_AS(ptr3.get_deleter()(ptr3.get()), std::runtime_error); + // The pointer is now broken, so we need to replace the deleter + ptr3.get_deleter() = + auxiliary::CustomDelete([](void const *ptr_in_lambda) { + delete reinterpret_cast(ptr_in_lambda); + }); } TEST_CASE("scalar_and_vector", "[core]") diff --git a/test/Files_SerialIO/container_buffer_size_1d_downsizes.cpp b/test/Files_SerialIO/container_buffer_size_1d_downsizes.cpp new file mode 100644 index 0000000000..759585de9e --- /dev/null +++ b/test/Files_SerialIO/container_buffer_size_1d_downsizes.cpp @@ -0,0 +1,140 @@ +/* Copyright 2026 + * + * This file is part of openPMD-api. + * + * openPMD-api is free software: you can redistribute it and/or modify + * it under the terms of of either the GNU General Public License or + * the GNU Lesser General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * openPMD-api is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License and the GNU Lesser General Public License + * for more details. + * + * You should have received a copy of the GNU General Public License + * and the GNU Lesser General Public License along with openPMD-api. + * If not, see . + */ + +/* + * Tests for container buffer-size validation in the prepareLoadStore() / + * storeChunk() API. + * + * The size of a contiguous container buffer is remembered when the buffer is + * specified (withContiguousContainer) and checked against the final operation + * extent, which is only known at store() / load() time. This prevents + * out-of-bounds reads (store) and writes (load) when the container is too + * small for the selected data: + * - For N-D datasets (or an explicit extent) the selection cannot be + * downsized, so a too-small container throws error::WrongAPIUsage. + * - For a 1-D dataset with the default (full) extent the selection is + * downsized to the buffer size instead. + * + * Before the fix ("Remember buffer sizes"), the N-D default extent fell back + * to the full dataset extent, so the backend read prod(dataset extent) + * elements out of a small vector (heap buffer over-read / over-write). + */ +#include "SerialIOTests.hpp" + +#include + +#include +#include + +using namespace openPMD; + +namespace +{ +std::string const nd_name = "../samples/issue3_nd_buffer_size.json"; + +void write_nd_dataset() +{ + Series s(nd_name, Access::CREATE); + auto rc = s.iterations[0].meshes["E"]["x"]; + rc.resetDataset(Dataset(Datatype::INT, {4, 4})); + std::vector data{0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15}; + rc.prepareLoadStore().withContiguousContainer(data).store().get(); + s.flush(); +} +} // namespace + +TEST_CASE("container_buffer_size_too_small_throws", "[serial][json]") +{ + write_nd_dataset(); + Series s(nd_name, Access::READ_ONLY); + auto rc = s.iterations[0].meshes["E"]["x"]; + + // N-D load, default extent, container far too small. + { + std::vector small(2); + REQUIRE_THROWS_AS( + rc.prepareLoadStore().withContiguousContainer(small).load().get(), + error::WrongAPIUsage); + } + // N-D load, explicit extent that over-reads the container. + { + std::vector small(2); + REQUIRE_THROWS_AS( + rc.prepareLoadStore() + .withContiguousContainer(small) + .offset({0, 0}) + .extent({4, 4}) + .load() + .get(), + error::WrongAPIUsage); + } + // N-D store, default extent, container too small (over-read). + { + std::vector small(2, 7); + REQUIRE_THROWS_AS( + rc.prepareLoadStore() + .withContiguousContainer(small) + .offset({0, 0}) + .store(), + error::WrongAPIUsage); + } + s.flush(); +} + +TEST_CASE("container_buffer_size_1d_downsizes", "[serial][json]") +{ + // 1-D dataset, default (full) extent, container smaller than the dataset: + // the selection is downsized to the buffer size, no throw. + constexpr int N = 100; + std::string const name = "../samples/issue3_1d_downsize.json"; + { + Series write(name, Access::CREATE); + auto rc = write.iterations[0].meshes["E"]["x"]; + rc.resetDataset(Dataset(Datatype::INT, {N})); + std::vector v(4, 7); + rc.storeChunk(v, {0}); // no extent: defaults to full, downsized to 4 + write.flush(); + } + Series read(name, Access::READ_ONLY); + auto rc = read.iterations[0].meshes["E"]["x"]; + std::vector buf(4, -1); + rc.loadChunkRaw(buf.data(), {0}, {4}); + read.flush(); + REQUIRE(buf == std::vector(4, 7)); + + std::fill(buf.begin(), buf.end(), -1); + rc.prepareLoadStore().offset({0}).withContiguousContainer(buf).load().get(); + REQUIRE(buf == std::vector(4, 7)); +} + +TEST_CASE("container_buffer_size_exact_loads", "[serial][json]") +{ + // N-D load with a container that is exactly large enough: no throw. + write_nd_dataset(); + Series s(nd_name, Access::READ_ONLY); + auto rc = s.iterations[0].meshes["E"]["x"]; + std::vector buf(16, -1); + rc.prepareLoadStore().withContiguousContainer(buf).load().get(); + s.flush(); + REQUIRE( + buf == + std::vector{0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15}); +} diff --git a/test/Files_SerialIO/deferred_load_unrelated_flush.cpp b/test/Files_SerialIO/deferred_load_unrelated_flush.cpp new file mode 100644 index 0000000000..04e694470b --- /dev/null +++ b/test/Files_SerialIO/deferred_load_unrelated_flush.cpp @@ -0,0 +1,89 @@ +/* Copyright 2026 + * + * This file is part of openPMD-api. + * + * openPMD-api is free software: you can redistribute it and/or modify + * it under the terms of of either the GNU General Public License or + * the GNU Lesser General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * openPMD-api is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License and the GNU Lesser General Public License + * for more details. + * + * You should have received a copy of the GNU General Public License + * and the GNU Lesser General Public License along with openPMD-api. + * If not, see . + */ + +/* + * Regression test for the per-record-component flush counter. + * + * A pending deferred load must not be treated as complete just because an + * *unrelated* backend operation flushed the I/O handler. The counter that + * records "has this record component been flushed?" is per record component + * and only advances when the component's own chunks are flushed. + * + * Before the fix ("Track flush counter per record component") the counter + * lived on the I/O handler and *any* backend flush (here: a rank-table + * query) advanced it, so the deferred load's own flush was skipped and the + * buffer kept its initial contents. + */ +#include "SerialIOTests.hpp" + +#include + +#include +#include + +using namespace openPMD; + +TEST_CASE("deferred_load_unrelated_flush", "[serial][json]") +{ + constexpr int N = 4; + std::vector data{1, 2, 3, 4}; + std::string const name = "../samples/deferred_load_unrelated_flush.json"; + + { + Series write(name, Access::CREATE); + auto E_x = write.iterations[0].meshes["E"]["x"]; + E_x.resetDataset(Dataset(Datatype::INT, {N})); + E_x.storeChunk(data, {0}, {N}); + write.setRankTable("hostname"); + write.flush(); + } + + auto read_one = [&name, &data](bool const allocating) { + Series read(name, Access::READ_ONLY); + auto E_x = read.iterations[0].meshes["E"]["x"]; + if (allocating) + { + auto pending = E_x.prepareLoadStore().load(); + // Unrelated backend query: drains the I/O handler queue but does + // not flush `E_x`'s pending frontend chunk. + (void)read.rankTable(/* collective = */ false); + auto const loaded = pending.get(); + for (int i = 0; i < N; ++i) + { + REQUIRE(loaded.get()[i] == data[i]); + } + } + else + { + int values[N] = {-7, -7, -7, -7}; + auto pending = E_x.prepareLoadStore().withRawPtr(values).load(); + // Unrelated backend query, as above. + (void)read.rankTable(/* collective = */ false); + pending.get(); + for (int i = 0; i < N; ++i) + { + REQUIRE(values[i] == data[i]); + } + } + }; + read_one(/* allocating = */ false); + read_one(/* allocating = */ true); +} diff --git a/test/Files_SerialIO/joined_dim_buffer_api.cpp b/test/Files_SerialIO/joined_dim_buffer_api.cpp new file mode 100644 index 0000000000..96a1966693 --- /dev/null +++ b/test/Files_SerialIO/joined_dim_buffer_api.cpp @@ -0,0 +1,86 @@ +/* Copyright 2026 + * + * This file is part of openPMD-api. + * + * openPMD-api is free software: you can redistribute it and/or modify + * it under the terms of of either the GNU General Public License or + * the GNU Lesser General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * openPMD-api is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License and the GNU Lesser General Public License + * for more details. + * + * You should have received a copy of the GNU General Public License + * and the GNU Lesser General Public License along with openPMD-api. + * If not, see . + */ + +/* + * Buffer-based storeChunk / loadChunk on a joined dimension. + * + * A joined dimension's extent is only known once all writers have flushed and + * the data is read back (i.e. after a close + reopen). During the write + * session it is Dataset::JOINED_DIMENSION (max), so: + * + * - storing a chunk is meaningful (an empty offset {} means "append"; a {0} + * offset is treated as the default) and must pass the frontend checks; + * - loading a chunk is NOT meaningful (the total size is unknown, and the + * streaming engine cannot read it back mid-write) and must be rejected + * with a clear error. + * + * The argument handling is frontend validation, backend agnostic. A joined + * dataset can only be flushed on backends that support it (ADIOS2), so this is + * exercised on the JSON backend, where the deferred store is unsupported and + * only reported at flush time (not fatal). The load must throw before any + * backend operation is queued. + */ +#include "SerialIOTests.hpp" + +#include + +#include +#include +#include +#include +#include + +using namespace openPMD; + +TEST_CASE("joined_dim_buffer_api", "[serial][json]") +{ + using type = float; + constexpr size_t N = 8; + std::string const name = "../samples/joinedDimBufferApi.json"; + + std::vector data(N); + std::iota(data.begin(), data.end(), 0.f); + + Series s(name, Access::CREATE); + auto epx = s.iterations[0].particles["e"]["position"]["x"]; + epx.resetDataset(Dataset(Datatype::FLOAT, {Dataset::JOINED_DIMENSION})); + REQUIRE(epx.joinedDimension().has_value()); + + // store: empty offset {} (canonical) and {0} offset (treated as default) + // must pass the frontend checks. + epx.storeChunkRaw(data.data(), {}, {N}); + epx.storeChunkRaw(data.data(), {0}, {N}); + + // load: a joined array cannot be loaded during the write session (its + // extent is not known), so every buffer-based load overload must reject it. + std::vector buf(N, -1.f); + REQUIRE_THROWS_AS( + epx.loadChunkRaw(buf.data(), {0}, {-1u}), error::WrongAPIUsage); + REQUIRE_THROWS_AS( + epx.loadChunkRaw(buf.data(), {}, {-1u}), error::WrongAPIUsage); + { + std::shared_ptr sptr( + new type[N], [](const type *p) { delete[] p; }); + std::fill(sptr.get(), sptr.get() + N, -1.f); + REQUIRE_THROWS_AS( + epx.loadChunk(sptr, {0}, {-1u}), error::WrongAPIUsage); + } +} diff --git a/test/Files_SerialIO/memory_selection_old_adios2.cpp b/test/Files_SerialIO/memory_selection_old_adios2.cpp new file mode 100644 index 0000000000..9eca70073e --- /dev/null +++ b/test/Files_SerialIO/memory_selection_old_adios2.cpp @@ -0,0 +1,70 @@ +/* Copyright 2026 + * + * This file is part of openPMD-api. + * + * openPMD-api is free software: you can redistribute it and/or modify + * it under the terms of of either the GNU General Public License or + * the GNU Lesser General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * openPMD-api is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License and the GNU Lesser General Public License + * for more details. + * + * You should have received a copy of the GNU General Public License + * and the GNU Lesser General Public License along with openPMD-api. + * If not, see . + */ + +/* + * Memory selections require the ability to reset them, since a stale memory + * selection would otherwise silently leak into subsequent store operations of + * the same variable (see + * https://github.com/ornladios/ADIOS2/pull/4169). + * + * The reset capability was added upstream in ADIOS2 v2.11.0 and backported to + * v2.10.1. Older versions must reject memory selections instead. + */ +#include "SerialIOTests.hpp" + +#include "openPMD/Error.hpp" +#include "openPMD/IO/ADIOS/macros.hpp" + +#include + +#include +#include + +using namespace openPMD; + +#if openPMD_HAVE_ADIOS2 +TEST_CASE("memory_selection_old_adios2", "[serial][adios2]") +{ + std::string const name = "../samples/memorySelectionOldAdios2.bp"; + Series s(name, Access::CREATE, R"({"backend": "adios2"})"); + auto rc = s.iterations[0].meshes["E"]["x"]; + rc.resetDataset({Datatype::INT, {5, 5}}); + std::vector data(25, 1); + auto store = [&]() { + rc.prepareLoadStore() + .withContiguousContainer(data) + .offset({1, 1}) + .extent({2, 2}) + .memorySelection({{1, 1}, {5, 5}}) + .store() + .get(); + }; + if (CanTheMemorySelectionBeReset) + { + store(); + SUCCEED("Memory selections supported by this ADIOS2 version."); + } + else + { + REQUIRE_THROWS_AS(store(), error::OperationUnsupportedInBackend); + } +} +#endif diff --git a/test/Files_SerialIO/memory_selection_rejected_before_flush.cpp b/test/Files_SerialIO/memory_selection_rejected_before_flush.cpp new file mode 100644 index 0000000000..c5fdf82686 --- /dev/null +++ b/test/Files_SerialIO/memory_selection_rejected_before_flush.cpp @@ -0,0 +1,153 @@ +/* Copyright 2026 + * + * This file is part of openPMD-api. + * + * openPMD-api is free software: you can redistribute it and/or modify + * it under the terms of of either the GNU General Public License or + * the GNU Lesser General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * openPMD-api is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License and the GNU Lesser General Public License + * for more details. + * + * You should have received a copy of the GNU General Public License + * and the GNU Lesser General Public License along with openPMD-api. + * If not, see . + */ + +/* + * Bug: rejecting a memory selection inside the IO task (at flush time) + * corrupts the whole output file. + * + * The HDF5 and JSON backends do not support non-contiguous memory selections. + * They used to throw error::OperationUnsupportedInBackend from within + * writeDataset(), i.e. while flushing an IO task. AbstractIOHandlerImpl::flush + * reacts to an exception in an IO task by clearing the whole IO queue and + * rethrowing ("Clearing IO queue and passing on the exception"), so every + * other pending chunk and the file's root attributes were dropped. Subsequent + * flush() and close() then "succeeded", but the file could not be read back + * (AttributeNotFound: openPMD). + * + * The capability is known when the chunk is enqueued, so storeChunk_impl() + * now checks it up front and throws before anything is queued. This test + * verifies that: + * 1. the unsupported store throws immediately (not at flush time), and + * 2. a valid chunk stored in the same Series survives and can be read back. + * + * The ADIOS2 backend does support memory selections, so the corresponding + * store must not throw there. + */ + +#include "SerialIOTests.hpp" + +#include "openPMD/IO/ADIOS/macros.hpp" + +#include + +#include +#include + +using namespace openPMD; + +namespace +{ +constexpr size_t N = 4; + +#if openPMD_HAVE_ADIOS2 +bool adios2SupportsMemorySelection() +{ + return CanTheMemorySelectionBeReset; +} +#else +bool adios2SupportsMemorySelection() +{ + return false; +} +#endif + +struct WriteResult +{ + bool rejected; + bool backendSupportsMemorySelection; +}; + +/* + * Writes a valid 1D chunk for component A and then attempts a store with a + * non-contiguous memory selection on component B. + */ +WriteResult write_and_reject(std::string const &name) +{ + Series s(name, Access::CREATE); + // HDF5 and JSON never support memory selections. ADIOS2 only supports them + // if its version can reset a memory selection (>= 2.10.1). + bool const backendSupports = + s.backend() == "ADIOS2" && adios2SupportsMemorySelection(); + s.setAttribute("some_global", "attribute"); + + // Valid chunk that must not be lost. + auto A = s.iterations[0].meshes["A"]["x"]; + A.resetDataset(Dataset(Datatype::INT, {N})); + std::vector aval{0, 1, 2, 3}; + A.storeChunk(aval, {0}, {N}); + + // Memory selection on another component. The selection (2x2) is a + // sub-block of the in-memory buffer layout (NxN), i.e. non-contiguous. + auto B = s.iterations[0].meshes["B"]["y"]; + B.resetDataset(Dataset(Datatype::INT, {N, N})); + std::vector bval(N * N, 7); + bool rejected = false; + try + { + B.prepareLoadStore() + .withContiguousContainer(bval) + .offset({0, 0}) + .extent({2, 2}) + .memorySelection({{1, 1}, {N, N}}) + .unsafeNoAutomaticFlush() + .store() + .get(); + } + catch (error::OperationUnsupportedInBackend const &) + { + rejected = true; + } + + // The Series must still be usable and the valid chunk must be written. + s.flush(); + s.close(); + return {rejected, backendSupports}; +} +} // namespace + +TEST_CASE("memory_selection_rejected_before_flush", "[serial]") +{ + auto allExtensions = getFileExtensions(); + for (auto const &ext : allExtensions) + { + if (ext == "sst" || ext == "ssc" || ext == "bp5" || ext == "toml") + { + continue; + } + std::string const name = + std::string("../samples/memory_selection_rejected.") + ext; + WriteResult const result = write_and_reject(name); + REQUIRE(result.rejected == !result.backendSupportsMemorySelection); + + // Whether or not the memory selection was rejected, the valid chunk + // and the root attribute must be readable. Before the fix, the + // backend-side throw at flush time cleared the IO queue and left an + // unreadable file behind. + Series read(name, Access::READ_ONLY); + REQUIRE( + read.getAttribute("some_global").get() == "attribute"); + std::vector aval(N, -1); + read.iterations[0].meshes["A"]["x"].loadChunkRaw(aval.data(), {0}, {N}); + read.flush(); + REQUIRE(aval == std::vector{0, 1, 2, 3}); + read.close(); + } +} diff --git a/test/ParallelIOTest.cpp b/test/ParallelIOTest.cpp index d32d31253b..9f516070b7 100644 --- a/test/ParallelIOTest.cpp +++ b/test/ParallelIOTest.cpp @@ -473,18 +473,73 @@ void available_chunks_test(std::string const &file_ending) } )END"; - std::vector data{2, 4, 6, 8}; + std::vector xdata{2, 4, 6, 8}; + std::vector ydata{0, 0, 0, 0, 0, // + 0, 1, 2, 3, 0, // + 0, 4, 5, 6, 0, // + 0, 7, 8, 9, 0, // + 0, 0, 0, 0, 0}; + std::vector ydata_firstandlastrow{-1, -1, -1}; + // Equivalent contiguous representation of the 3x3 block with values 1..9 + // for ADIOS2 versions that do not support memory selections. + std::vector ydata_block{1, 2, 3, 4, 5, 6, 7, 8, 9}; { Series write(name, Access::CREATE, MPI_COMM_WORLD, parameters.str()); Iteration it0 = write.iterations[0]; auto E_x = it0.meshes["E"]["x"]; E_x.resetDataset({Datatype::INT, {mpi_size, 4}}); - E_x.storeChunk(data, {mpi_rank, 0}, {1, 4}); + E_x.storeChunk(xdata, {mpi_rank, 0}, {1, 4}); + auto E_y = it0.meshes["E"]["y"]; + E_y.resetDataset({Datatype::INT, {5, 3ul * mpi_size}}); + E_y.prepareLoadStore() + .withContiguousContainer(ydata_firstandlastrow) + .offset({0, 3ul * mpi_rank}) + .extent({1, 3}) + .unsafeNoAutomaticFlush() + .store() + .get(); + // Memory selections can only be reset in ADIOS2 >= 2.10.1 + // (https://github.com/ornladios/ADIOS2/pull/4169). Older versions + // reject them, so use an equivalent contiguous buffer there. + if constexpr (CanTheMemorySelectionBeReset) + { + // Take the 3x3 block (values 1..9) out of the 5x5 buffer `ydata` + // and store it non-contiguously. + E_y.prepareLoadStore() + .offset({1, 3ul * mpi_rank}) + .extent({3, 3}) + .withContiguousContainer(ydata) + .memorySelection({{1, 1}, {5, 5}}) + .unsafeNoAutomaticFlush() + .store() + .get(); + } + else + { + E_y.prepareLoadStore() + .withContiguousContainer(ydata_block) + .offset({1, 3ul * mpi_rank}) + .extent({3, 3}) + .unsafeNoAutomaticFlush() + .store() + .get(); + } + E_y.prepareLoadStore() + .withContiguousContainer(ydata_firstandlastrow) + .offset({4, 3ul * mpi_rank}) + .extent({1, 3}) + .unsafeNoAutomaticFlush() + .store() + .get(); it0.close(); } { - Series read(name, Access::READ_ONLY, MPI_COMM_WORLD); + Series read( + name, + Access::READ_ONLY, + MPI_COMM_WORLD, + R"({"verify_homogeneous_extents": false})"); Iteration it0 = read.iterations[0]; auto E_x = it0.meshes["E"]["x"]; ChunkTable table = E_x.availableChunks(); @@ -511,6 +566,35 @@ void available_chunks_test(std::string const &file_ending) { REQUIRE(ranks[i] == i); } + + auto E_y = it0.meshes["E"]["y"]; + auto width = E_y.getExtent()[1]; + auto first_row_deferred = + E_y.prepareLoadStore().extent({1, width}).load(); + auto middle_rows_deferred = E_y.prepareLoadStore() + .offset({1, 0}) + .extent({3, width}) + .load(); + auto last_row_deferred = + E_y.prepareLoadStore().offset({4, 0}).load(); + + auto first_row = first_row_deferred.get(); + auto middle_rows = middle_rows_deferred.get(); + auto last_row = last_row_deferred.get(); + + for (auto row : {&first_row, &last_row}) + { + for (size_t i = 0; i < width; ++i) + { + REQUIRE(row->get()[i] == -1); + } + } + for (size_t i = 0; i < width * 3; ++i) + { + size_t row = i / width; + int required_value = row * 3 + (i % 3) + 1; + REQUIRE(middle_rows.get()[i] == required_value); + } } } diff --git a/test/SerialIOTest.cpp b/test/SerialIOTest.cpp index 8637389d0f..26be8ef0b5 100644 --- a/test/SerialIOTest.cpp +++ b/test/SerialIOTest.cpp @@ -945,7 +945,13 @@ inline void constant_scalar(std::string const &file_ending) new unsigned int[6], [](unsigned int const *p) { delete[] p; }); unsigned int e{0}; std::generate(E.get(), E.get() + 6, [&e] { return e++; }); - E_y.storeChunk(std::move(E), {0, 0, 0}, {1, 2, 3}); + // check that const-type unique pointers work in the builder pattern + E_y.prepareLoadStore() + .extent({1, 2, 3}) + .withUniquePtr(std::move(E).static_cast_()) + .unsafeNoAutomaticFlush() + .store() + .get(); // store a number of predefined attributes in E Mesh &E_mesh = s.snapshots()[1].meshes["E"]; @@ -1756,13 +1762,17 @@ inline void write_test( auto opaqueTypeDataset = rc.visit(); auto variantTypeDataset = rc.loadChunkVariant(); + auto variantTypeDataset2 = rc.prepareLoadStore().loadVariant().get(); rc.seriesFlush(); - std::visit( - [](auto &&shared_ptr) { - std::cout << "First value in loaded chunk: '" << shared_ptr.get()[0] - << '\'' << std::endl; - }, - variantTypeDataset); + for (auto ptr : {&variantTypeDataset, &variantTypeDataset2}) + { + std::visit( + [](auto &&shared_ptr) { + std::cout << "First value in loaded chunk: '" + << shared_ptr.get()[0] << '\'' << std::endl; + }, + *ptr); + } #ifndef _WIN32 if (test_rank_table)