@@ -113,38 +113,44 @@ void on_communication_requested(uv_async_t* s)
113113}
114114
115115DataProcessingDevice::DataProcessingDevice (RunningDeviceRef ref, ServiceRegistry& registry, ProcessingPolicies& policies)
116- : mSpec {registry.get <RunningWorkflowInfo const >(ServiceRegistry::threadSalt ()).devices [ref.index ]},
117- mState {registry.get <DeviceState>(ServiceRegistry::threadSalt ())},
116+ : mSpec {registry.get <RunningWorkflowInfo const >(ServiceRegistry::globalDeviceSalt ()).devices [ref.index ]},
117+ mState {registry.get <DeviceState>(ServiceRegistry::globalDeviceSalt ())},
118118 mInit {mSpec .algorithm .onInit },
119119 mStatefulProcess {nullptr },
120120 mStatelessProcess {mSpec .algorithm .onProcess },
121121 mError {mSpec .algorithm .onError },
122122 mConfigRegistry {nullptr },
123123 mServiceRegistry {registry},
124124 mProcessingPolicies {policies},
125- mQuotaEvaluator {registry.get <ComputingQuotaEvaluator>(ServiceRegistry::threadSalt ())}
125+ mQuotaEvaluator {registry.get <ComputingQuotaEvaluator>(ServiceRegistry::globalDeviceSalt ())}
126126{
127127
128128 // / FIXME: move erro handling to a service?
129129 if (mError != nullptr ) {
130130 mErrorHandling = [&errorCallback = mError ,
131131 &serviceRegistry = mServiceRegistry ](RuntimeErrorRef e, InputRecord& record) {
132132 ZoneScopedN (" Error handling" );
133+ // / FIXME: we should pass the salt in, so that the message
134+ // / can access information which were stored in the stream.
135+ ServiceRegistryRef ref{serviceRegistry, ServiceRegistry::globalDeviceSalt ()};
133136 auto & err = error_from_ref (e);
134137 LOGP (error, " Exception caught: {} " , err.what );
135138 demangled_backtrace_symbols (err.backtrace , err.maxBacktrace , STDERR_FILENO );
136- serviceRegistry .get <DataProcessingStats>(ServiceRegistry::threadSalt () ).exceptionCount ++;
137- ErrorContext errorContext{record, serviceRegistry , e};
139+ ref .get <DataProcessingStats>().exceptionCount ++;
140+ ErrorContext errorContext{record, ref , e};
138141 errorCallback (errorContext);
139142 };
140143 } else {
141144 mErrorHandling = [&errorPolicy = mProcessingPolicies .error ,
142145 &serviceRegistry = mServiceRegistry ](RuntimeErrorRef e, InputRecord& record) {
143146 ZoneScopedN (" Error handling" );
144147 auto & err = error_from_ref (e);
148+ // / FIXME: we should pass the salt in, so that the message
149+ // / can access information which were stored in the stream.
145150 LOGP (error, " Exception caught: {} " , err.what );
151+ ServiceRegistryRef ref{serviceRegistry, ServiceRegistry::globalDeviceSalt ()};
146152 demangled_backtrace_symbols (err.backtrace , err.maxBacktrace , STDERR_FILENO );
147- serviceRegistry .get <DataProcessingStats>(ServiceRegistry::threadSalt () ).exceptionCount ++;
153+ ref .get <DataProcessingStats>().exceptionCount ++;
148154 switch (errorPolicy) {
149155 case TerminationPolicy::QUIT :
150156 throw e;
@@ -155,9 +161,10 @@ DataProcessingDevice::DataProcessingDevice(RunningDeviceRef ref, ServiceRegistry
155161 }
156162
157163 std::function<void (const fair::mq::State)> stateWatcher = [this , ®istry = mServiceRegistry ](const fair::mq::State state) -> void {
158- auto & deviceState = registry.get <DeviceState>(ServiceRegistry::threadSalt ());
159- auto & control = registry.get <ControlService>(ServiceRegistry::threadSalt ());
160- auto & callbacks = registry.get <CallbackService>(ServiceRegistry::threadSalt ());
164+ auto ref = ServiceRegistryRef{registry, ServiceRegistry::globalDeviceSalt ()};
165+ auto & deviceState = ref.get <DeviceState>();
166+ auto & control = ref.get <ControlService>();
167+ auto & callbacks = ref.get <CallbackService>();
161168 control.notifyDeviceState (fair::mq::GetStateName (state));
162169 callbacks (CallbackService::Id::DeviceStateChanged, registry, state);
163170
@@ -356,7 +363,7 @@ void DataProcessingDevice::Init()
356363 str = entry.second .get_value <std::string>();
357364 }
358365 std::string configString = fmt::format (" [CONFIG] {}={} 1 {}" , entry.first , str, configStore->provenance (entry.first .c_str ())).c_str ();
359- mServiceRegistry .get <DriverClient>(ServiceRegistry::threadSalt ()).tell (configString.c_str ());
366+ mServiceRegistry .get <DriverClient>(ServiceRegistry::globalDeviceSalt ()).tell (configString.c_str ());
360367 }
361368
362369 mConfigRegistry = std::make_unique<ConfigParamRegistry>(std::move (configStore));
@@ -386,7 +393,7 @@ void DataProcessingDevice::Init()
386393 // Invoke the callback policy for this device.
387394 if (mSpec .callbacksPolicy .policy != nullptr ) {
388395 InitContext initContext{*mConfigRegistry , mServiceRegistry };
389- mSpec .callbacksPolicy .policy (mServiceRegistry .get <CallbackService>(ServiceRegistry::threadSalt ()), initContext);
396+ mSpec .callbacksPolicy .policy (mServiceRegistry .get <CallbackService>(ServiceRegistry::globalDeviceSalt ()), initContext);
390397 }
391398}
392399
@@ -864,7 +871,7 @@ void DataProcessingDevice::InitTask()
864871 // more a ServiceRegistry::globalDataProcessorSalt(N) where
865872 // N is the number of the multiplexed data processor.
866873 // We will get there.
867- this ->fillContext (mServiceRegistry .get <DataProcessorContext>(ServiceRegistry::threadSalt ()), deviceContext);
874+ this ->fillContext (mServiceRegistry .get <DataProcessorContext>(ServiceRegistry::globalDeviceSalt ()), deviceContext);
868875
869876 // / We now run an event loop also in InitTask. This is needed to:
870877 // / * Make sure region registration callbacks are invoked
0 commit comments