12#include "GaudiKernel/ClassID.h"
13#include "GaudiKernel/FileIncident.h"
14#include "GaudiKernel/GenericAddress.h"
15#include "GaudiKernel/IIncidentSvc.h"
16#include "GaudiKernel/IOpaqueAddress.h"
44 return StatusCode::FAILURE;
62 getPoolSvc()->setShareMode(
true);
67 incSvc->addListener(
this,
"StoreCleared", pri);
70 return this->AthenaPoolCnvSvc::initialize();
92 return this->AthenaPoolCnvSvc::finalize();
96 const std::string& openMode) {
97 return AthenaPoolCnvSvc::connectOutput(outputConnectionSpec, openMode);
102 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
105 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
106 return(StatusCode::FAILURE);
111 return(StatusCode::SUCCESS);
115 ATH_MSG_DEBUG(std::format(
"connectOutput SKIPPED for metadata-only server: {}", outputConnectionSpec));
116 return(StatusCode::SUCCESS);
120 return(StatusCode::SUCCESS);
126 std::size_t apend = outputConnectionSpec.find(
'[');
127 if (apend != std::string::npos) {
128 outputConnection += outputConnectionSpec.substr(apend);
130 if (outputConnectionSpec.find(
"[PoolContainerPrefix=" +
m_metadataContainerProp.value() +
"]") != std::string::npos) {
131 return AthenaPoolCnvSvc::connectOutput(outputConnection,
"APPEND");
133 return AthenaPoolCnvSvc::connectOutput(outputConnection);
139 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
143 if (!this->
cleanUp(outputConnection).isSuccess()) {
145 return(StatusCode::FAILURE);
147 return(StatusCode::SUCCESS);
150 ATH_MSG_DEBUG(
"commitOutput SKIPPED for uninitialized server.");
151 return(StatusCode::SUCCESS);
153 std::map<void*, RootType> commitCache;
154 std::string fileName;
157 const char* placementStr =
nullptr;
160 if (
sc.isSuccess() && placementStr !=
nullptr && strlen(placementStr) > 6 && num > 0) {
161 const char * matchedChars = strstr(placementStr,
"[FILE=");
163 ATH_MSG_ERROR(std::format(
"No matching filename in {}", placementStr));
166 fileName = matchedChars;
167 fileName = fileName.substr(6, fileName.find(
']') - 6);
169 ATH_MSG_ERROR(std::format(
"Failed to connectOutput for {}", fileName));
173 bool dataHeaderSeen =
false;
174 std::string dataHeaderID;
176 std::string objName =
"ALL";
177 if (useDetailChronoStat()) {
178 objName = placementStr;
183 std::string_view pStr = placementStr;
184 std::string::size_type cpos = pStr.find (
"[CONT=");
185 if (cpos == std::string::npos) {
186 ATH_MSG_ERROR(std::format(
"No CONT field in placement string: {}", pStr));
187 return StatusCode::FAILURE;
189 std::string tokenStr (pStr.substr(0, cpos));
190 std::string contName (pStr.substr(cpos, std::string::npos));
191 std::string::size_type cl1 = contName.find(
']');
192 if (cl1 == std::string::npos) {
193 ATH_MSG_ERROR(std::format(
"Missing close bracket after CONT field in placement string: {}", pStr));
194 return StatusCode::FAILURE;
196 tokenStr.append(contName, cl1 + 1);
197 contName = contName.substr(6, cl1 - 6);
199 std::string::size_type ppos = pStr.find (
"[PNAME=");
200 if (ppos == std::string::npos) {
201 ATH_MSG_ERROR(std::format(
"No PNAME field in placement string: {}", pStr));
202 return StatusCode::FAILURE;
204 std::string className (pStr.substr(ppos, std::string::npos));
205 std::string::size_type cl2 = className.find(
']');
206 if (cl2 == std::string::npos) {
207 ATH_MSG_ERROR(std::format(
"Missing close bracket after PNAME field in placement string: {}", pStr));
208 return StatusCode::FAILURE;
210 className = className.substr(7, cl2 - 7);
213 const std::string numStr = std::to_string(num);
215 bool foundContainer =
false;
216 std::size_t opPos = contName.find(
'(');
218 foundContainer =
true;
221 if (contName.compare(0, opPos, item) == 0){
222 foundContainer =
true;
228 if (len > 0 && foundContainer && contName[len] ==
'(' ) {
235 memName, {}, memName,
236 "BeginInputMemFile",
"EndInputMemFile");
243 sc = metadataSvc->shmProxy(std::format(
"{}[NUM={}]", pStr, numStr));
244 if (
sc.isRecoverable()) {
246 }
else if (
sc.isFailure()) {
257 if( m_oneDataHeaderForm.value() ) {
258 auto placementWithSwn = [&] {
return std::format(
"{}[SWN={}]", placementStr, num); };
259 if( className ==
"DataHeaderForm_p6" ) {
262 "", placementWithSwn());
263 DHcnv->updateRepRefs(&address,
static_cast<DataObject*
>(obj)).ignore();
269 if (token ==
nullptr) {
273 tokenStr = token->toString();
275 if( className ==
"DataHeader_p6" ) {
278 tokenStr, placementWithSwn());
279 if (!DHcnv->updateRep(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
284 if (className !=
"Token" && className !=
"DataHeaderForm_p6" && !classDesc.
IsFundamental()) {
285 commitCache.insert(std::pair<void*, RootType>(obj, classDesc));
287 placementStr =
nullptr;
291 placement.
fromString(placementStr); placementStr =
nullptr;
293 if (token ==
nullptr) {
297 tokenStr = token->toString();
298 if (className ==
"DataHeader_p6") {
303 if (!DHcnv->updateRep(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
307 dataHeaderSeen =
true;
314 dataHeaderID = std::format(
"{}/{}/{}", token->contID(), numStr, token->dbID().toString());
315 }
else if (dataHeaderSeen) {
316 dataHeaderSeen =
false;
319 if (className ==
"DataHeaderForm_p6") {
322 tokenStr, dataHeaderID);
323 if (!DHcnv->updateRepRefs(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
324 ATH_MSG_ERROR(
"Failed updateRepRefs for obj = " << tokenStr);
329 GenericAddress address(0, 0,
"", dataHeaderID);
330 if (!DHcnv->updateRepRefs(&address,
nullptr).isSuccess()) {
336 if (className !=
"Token" && className !=
"DataHeaderForm_p6" && !classDesc.
IsFundamental()) {
337 commitCache.insert(std::pair<void*, RootType>(obj, classDesc));
344 while (
sc.isRecoverable()) {
347 if (!
sc.isSuccess()) {
353 while (
sc.isRecoverable()) {
356 if (
sc.isFailure()) {
361 if (dataHeaderSeen) {
363 GenericAddress address(0, 0,
"", std::move(dataHeaderID));
364 if (!DHcnv->updateRepRefs(&address,
nullptr).isSuccess()) {
369 placementStr =
nullptr;
370 }
else if (
sc.isSuccess() && placementStr !=
nullptr && strncmp(placementStr,
"stop", 4) == 0) {
371 return(StatusCode::RECOVERABLE);
372 }
else if (
sc.isRecoverable() || num == -1) {
373 return(StatusCode::RECOVERABLE);
375 if (
sc.isFailure() || fileName.empty()) {
380 memName, {}, memName,
381 "BeginInputMemFile",
"EndInputMemFile");
383 if (
sc.isFailure()) {
384 ATH_MSG_INFO(
"All SharedWriter clients stopped - exiting");
388 return(StatusCode::FAILURE);
392 ATH_MSG_DEBUG(std::format(
"commitOutput SKIPPED for metadata-only server: {}", outputConnectionSpec));
393 return(StatusCode::SUCCESS);
395 if (outputConnection.empty()) {
396 outputConnection = std::move(fileName);
398 outputConnection = outputConnectionSpec;
403 std::size_t
merge = outputConnection.find(
"?pmerge=");
404 const std::string baseOutputConnection = outputConnection.substr(0,
merge);
413 StatusCode status = AthenaPoolCnvSvc::commitOutput(outputConnection, doCommit);
414 for (
auto& [ptr, rootType] : commitCache) {
415 rootType.Destruct(ptr);
422 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
425 return(StatusCode::SUCCESS);
431 return(StatusCode::SUCCESS);
440 return AthenaPoolCnvSvc::disconnectOutput(outputConnectionSpec +
m_streamPortString.value());
447 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
451 Token* token =
nullptr;
455 std::string placementStr = placement->
toString();
456 placementStr +=
"[PNAME=";
457 placementStr += classDesc.
Name();
461 while (
sc.isRecoverable()) {
465 if (!
sc.isSuccess()) {
470 const void* buffer =
nullptr;
471 std::size_t nbytes = 0;
473 if (classDesc.
Name() ==
"Token") {
474 nbytes = strlen(
static_cast<const char*
>(obj)) + 1;
478 nbytes = classDesc.
SizeOf();
486 while (
sc.isRecoverable()) {
490 if (own) {
delete []
static_cast<const char*
>(buffer); }
492 if (!
sc.isSuccess()) {
493 ATH_MSG_ERROR(
"Could not share object for: " << placementStr);
498 if (info !=
nullptr) {
500 ATH_MSG_ERROR(
"Could not share dynamic aux store for: " << placementStr);
510 const char* tokenStr =
nullptr;
513 while (
sc.isRecoverable()) {
517 if (!
sc.isSuccess()) {
521 if (!strcmp(tokenStr,
"ABORT")) {
528 tempToken->
fromString(tokenStr); tokenStr =
nullptr;
530 token = tempToken; tempToken =
nullptr;
534 ATH_MSG_DEBUG(
"registerForWrite SKIPPED for uninitialized server, Placement = " << placement->
toString());
537 token = tempToken; tempToken =
nullptr;
543 token = getPoolSvc()->registerForWrite(placement, obj, classDesc);
548 token = AthenaPoolCnvSvc::registerForWrite(placement, obj, classDesc);
557 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
562 int num = token->
oid().first;
564 void* buffer =
nullptr;
565 std::size_t nbytes = 0;
567 while (
sc.isRecoverable()) {
570 if (!
sc.isSuccess()) {
578 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
581 std::string className = token->
auxString();
582 className = className.substr(className.find(
"[PNAME="));
583 className = className.substr(7, className.find(
']') - 7);
585 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
590 while (
sc.isRecoverable() && nbytes > 0) {
593 if (
sc.isSuccess() && nbytes > 0) {
595 classId.
fromString(
static_cast<const char*
>(buffer));
597 const std::string contName = std::string(
static_cast<const char*
>(buffer));
600 if (info !=
nullptr) {
601 if (
m_auxDynTool->hasAuxStore(contName, info->clazz().Class() ) && !
m_auxOutput->receiveStore(info->clazz().Class(), obj, num).isSuccess()) {
618 void* buffer =
nullptr;
619 std::size_t nbytes = 0;
620 StatusCode
sc = StatusCode::FAILURE;
625 while (
sc.isRecoverable()) {
630 if (!
sc.isSuccess()) {
638 while (
sc.isRecoverable() && nbytes > 0) {
641 if (
sc.isSuccess() && nbytes > 0) {
643 classId.
fromString(
static_cast<const char*
>(buffer));
645 const std::string contName = std::string(
static_cast<const char*
>(buffer));
648 if (info !=
nullptr) {
649 if (!
m_auxInput->receiveStore(info->clazz().Class(), obj).isSuccess()) {
660 AthenaPoolCnvSvc::setObjPtr(obj, token);
666 const std::string* par,
667 const unsigned long* ip,
668 IOpaqueAddress*& refpAddress) {
674 addressToken.
setDb(par[0].substr(4));
678 void* buffer =
nullptr;
679 std::size_t nbytes = 0;
681 while (
sc.isRecoverable()) {
685 if (!
sc.isSuccess()) {
687 return(StatusCode::FAILURE);
689 auto token = std::make_unique<Token>();
690 token->fromString(
static_cast<const char*
>(buffer)); buffer =
nullptr;
697 return(StatusCode::SUCCESS);
700 return(StatusCode::RECOVERABLE);
703 return AthenaPoolCnvSvc::createAddress(svcType, clid, par, ip, refpAddress);
709 const std::string& refAddress,
710 IOpaqueAddress*& refpAddress) {
711 return AthenaPoolCnvSvc::createAddress(svcType, clid, refAddress, refpAddress);
715 auto pos = connection.find(
"?pmerge=");
716 std::string conn = (pos == std::string::npos) ? connection : connection.substr(0, pos);
717 return AthenaPoolCnvSvc::cleanUp(conn);
731 return(StatusCode::FAILURE);
734 m_persSvcPerOutput.setValue(
false);
735 return(StatusCode::SUCCESS);
737 return(StatusCode::RECOVERABLE);
740 return(StatusCode::RECOVERABLE);
749 std::string streamPortSuffix;
752 return(StatusCode::FAILURE);
755 ATH_MSG_DEBUG(
"makeClient: Setting conversion service port suffix to " << streamPortSuffix);
760 return(StatusCode::SUCCESS);
763 std::string dummyStr;
769 return(StatusCode::FAILURE);
771 const char* tokenStr =
nullptr;
774 if (
sc.isSuccess() && tokenStr !=
nullptr && strlen(tokenStr) > 0 && num > 0) {
775 ATH_MSG_DEBUG(
"readData: " << tokenStr <<
", for client: " << num);
782 token.
fromString(tokenStr); tokenStr =
nullptr;
784 std::string objName =
"ALL";
785 if (useDetailChronoStat()) {
793 void* buffer =
nullptr;
794 std::size_t nbytes = 0;
797 while (
sc.isRecoverable()) {
800 delete []
static_cast<char*
>(buffer); buffer =
nullptr;
801 if (!
sc.isSuccess()) {
803 return(StatusCode::FAILURE);
806 if (info !=
nullptr) {
809 return(StatusCode::FAILURE);
815 return(StatusCode::FAILURE);
818 std::string returnToken;
820 if( metadataToken ) {
821 returnToken = metadataToken->
toString();
822 metadataToken->
release(); metadataToken =
nullptr;
830 return(StatusCode::FAILURE);
833 return(StatusCode::RECOVERABLE);
835 return(StatusCode::SUCCESS);
840 getPoolSvc()->commitCatalog();
841 getPoolSvc()->startCatalog();
842 return(StatusCode::SUCCESS);
851 StatusCode
sc = StatusCode::SUCCESS;
852 while (
sc.isSuccess()) {
858 while (
sc.isRecoverable()) {
862 return StatusCode::FAILURE;
873 base_class(name, pSvcLocator) {
#define ATH_CHECK
Evaluate an expression and check for errors.
#define ATH_MSG_VERBOSE(x)
#define ATH_MSG_WARNING(x)
This file contains the class definition for the AthenaPoolSharedIOCnvSvc class.
uint32_t CLID
The Class ID type.
This file contains the class definition for the Placement class (migrated from POOL).
This file contains the class definition for the TokenAddress class.
This file contains the class definition for the Token class (migrated from POOL).
std::map< std::string, int > m_fileCommitCounter
Force SharedWriter to flush data to output file at given intervals, needed by parallel compression.
virtual StatusCode commitOutput(const std::string &outputConnectionSpec, bool doCommit) override
Implementation of IConversionSvc: Commit pending output.
Gaudi::Property< std::string > m_streamPortString
Extension to use ROOT TMemFile for event data, "?pmerge=<host>:<port>".
virtual StatusCode finalize() override
Required of all Gaudi Services.
AthenaPoolSharedIOCnvSvc(const std::string &name, ISvcLocator *pSvcLocator)
Standard Service Constructor.
virtual StatusCode cleanUp(const std::string &connection) override
Implement cleanUp to call all registered IAthenaPoolCleanUp cleanUp() function.
virtual void setObjPtr(void *&obj, const Token *token) override
bool m_streamServerActive
StatusCode createAddress(long svcType, const CLID &clid, const std::string *par, const unsigned long *ip, IOpaqueAddress *&refpAddress) override
Create a Generic address using explicit arguments to identify a single object.
std::unique_ptr< RootAuxDynIO::IAuxDynShare > m_auxOutput
StatusCode abortSharedWrClients(int client_n)
Send abort to SharedWriter clients if the server quits on error.
virtual StatusCode makeServer(int num) override
Make this a server.
virtual void handle(const Incident &incident) override
Implementation of IIncidentListener: Handle for EndEvent incidence.
Gaudi::Property< std::string > m_metadataContainerProp
For SharedWriter: To use MetadataSvc to merge data placed in a certain container.
virtual StatusCode disconnectOutput(const std::string &outputConnectionSpec) override
Disconnect to the output connection.
virtual StatusCode readData() override
Read the next data object.
virtual StatusCode commitCatalog() override
Commit Catalog.
virtual StatusCode connectOutput(const std::string &outputConnectionSpec, const std::string &openMode) override
Implementation of IConversionSvc: Connect to the output connection specification with open mode.
std::unique_ptr< RootAuxDynIO::IAuxDynShare > m_auxInput
virtual ~AthenaPoolSharedIOCnvSvc()
Destructor.
Gaudi::Property< std::vector< std::string > > m_metadataContainersAug
Gaudi::Property< bool > m_parallelCompression
Use Athena Object sharing for metadata only, event data is collected and send via ROOT TMemFile.
ToolHandle< IAthenaIPCTool > m_outputStreamingTool
Gaudi::Property< std::map< std::string, int > > m_fileFlushSetting
Gaudi::Property< int > m_streamingTechnology
Use Streaming for selected technologies only.
virtual StatusCode initialize() override
Required of all Gaudi Services.
ToolHandle< IAthenaIPCTool > m_inputStreamingTool
ServiceHandle< IAthenaSerializeSvc > m_serializeSvc
virtual Token * registerForWrite(Placement *placement, const void *obj, const RootType &classDesc) override
Gaudi::Property< int > m_makeStreamingToolClient
Make this instance a Streaming Client during first connect/write automatically.
virtual StatusCode makeClient(int num) override
Make this a client.
std::unique_ptr< RootAuxDynIO::IFactoryTool > m_auxDynTool
This class provides a encapsulation of a GUID/UUID/CLSID/IID data structure (128 bit number).
constexpr void fromString(std::string_view s)
Automatic conversion from string representation.
static const Guid & null() noexcept
NULL-Guid: static class method.
constexpr void toString(std::span< char, StrLen > buf, bool uppercase=true) const noexcept
Automatic conversion to string representation.
This class holds all the necessary information to guide the writing of an object in a physical place.
const std::string & auxString() const
Access auxiliary string.
const std::string & containerName() const
Access container name.
Placement & setFileName(const std::string &fileName)
Set file name.
const std::string toString() const
Retrieve the string representation of the placement.
Placement & setTechnology(int technology)
Set technology type.
const std::string & fileName() const
Access file name.
Placement & fromString(const std::string &from)
Build from the string representation of a placement.
int technology() const
Access technology type.
static TScopeAdapter ByNameNoQuiet(const std::string &name, Bool_t load=kTRUE)
Bool_t IsFundamental() const
std::string Name(unsigned int mod=Reflex::SCOPED) const
void Destruct(void *place) const
This class provides a Generic Transient Address for POOL tokens.
This class provides a token that identifies in a unique way objects on the persistent storage.
const std::string & auxString() const
Access auxiliary string.
Token & setCont(const std::string &cnt)
Set container name.
Token & setDb(const Guid &db)
Set database name.
const std::string & contID() const
Access container identifier.
const Guid & classID() const
Access database identifier.
Token & setClassID(const Guid &cl_id)
Access database identifier.
const std::string toString() const
Retrieve the string representation of the token.
int technology() const
Access technology type.
int release()
Release token: Decrease reference count and eventually delete.
Token & setOid(const OID_t &oid)
Set object identifier.
const OID_t & oid() const
Access object identifier.
Token & fromString(const std::string_view from)
Build from the string representation of a token.
const Guid & dbID() const
Access database identifier.
Token & setAuxString(std::string &&auxString)
Set auxiliary string.
static const TypeH forGuid(const Guid &info)
Access classes by Guid.
static Guid guid(const TypeH &id)
Determine Guid (normalized string form) from reflection type.
Definition of class DbTypeInfo.
static const DbTypeInfo * create(const std::string &cl_name)
Create type information using name.
static DbType getType(const std::string &name)
Access known storage type object by name.
static const DbType POOL_StorageType
static constexpr CLID ID()