helics 3.7.0
Loading...
Searching...
No Matches
FederateState.hpp
1/*
2Copyright (c) 2017-2026,
3Battelle Memorial Institute; Lawrence Livermore National Security, LLC; Alliance for Energy
4Innovation LLC. See the top-level NOTICE for additional details. All rights reserved.
5SPDX-License-Identifier: BSD-3-Clause
6*/
7#pragma once
8
9#include "../common/GuardedTypes.hpp"
10#include "ActionMessage.hpp"
11#include "BasicHandleInfo.hpp"
12#include "CoreTypes.hpp"
13#include "InterfaceInfo.hpp"
14#include "core-data.hpp"
15#include "gmlc/containers/BlockingQueue.hpp"
16#include "helicsTime.hpp"
17
18#include <algorithm>
19#include <atomic>
20#include <chrono>
21#include <cstdint>
22#include <deque>
23#include <map>
24#include <memory>
25#include <optional>
26#include <string>
27#include <thread>
28#include <tuple>
29#include <utility>
30#include <vector>
31
32namespace helics {
33class SubscriptionInfo;
34class PublicationInfo;
35class EndpointInfo;
36class FilterInfo;
37class CommonCore;
38class CoreFederateInfo;
39
40class TimeCoordinator;
41class MessageTimer;
42class LogManager;
43
44constexpr Time startupTime = Time::minVal();
45constexpr Time initialTime{-1000000.0};
46
48enum class TimeSynchronizationMethod : uint8_t { DISTRIBUTED = 0, GLOBAL = 1, ASYNC = 2 };
49
52 public:
54 FederateState(const std::string& fedName, const CoreFederateInfo& fedInfo);
55 // the destructor is defined so some classes linked with unique ptrs don't have to be defined in
56 // the header
58 FederateState(const FederateState&) = delete;
59 FederateState& operator=(const FederateState&) = delete;
62
63 private:
64 const std::string name;
66 std::unique_ptr<TimeCoordinator> timeCoord;
67
68 public:
70 std::atomic<GlobalFederateId> global_id;
71 private:
73 std::atomic<FederateStates> state{FederateStates::CREATED};
74 bool only_transmit_on_change{false};
76 bool realtime{false};
77 bool observer{false};
78 bool reentrant{false};
79 bool mSourceOnly{false};
80 bool mCallbackBased{false};
82 bool strict_input_type_checking{false};
83 bool ignore_unit_mismatch{false};
85 bool mSlowResponding{false};
87 bool mAllowRemoteControl{true};
88 InterfaceInfo interfaceInformation;
89 std::unique_ptr<LogManager> mLogManager;
90 int maxLogLevel{HELICS_LOG_LEVEL_NO_PRINT};
91
92 public:
93 std::atomic<bool> init_transmitted{false};
95 int indexGroup{0};
96
97 private:
99 bool wait_for_current_time{false};
101 bool ignore_time_mismatch_warnings{false};
103 bool mProfilerActive{false};
106 bool mLocalProfileCapture{false};
107 int errorCode{0};
108 CommonCore* mParent{nullptr};
109 std::string errorString;
111 decltype(std::chrono::steady_clock::now()) start_clock_time;
112 Time rt_lag{timeZero};
113 Time rt_lead{timeZero};
114 Time grantTimeOutPeriod{timeZero};
115 std::uint64_t queuedValueBytes{0};
116 std::uint64_t valueBufferWarningLimit{100ULL * 1024ULL * 1024ULL};
117 std::uint64_t nextValueBufferWarning{100ULL * 1024ULL * 1024ULL};
118 std::int32_t realTimeTimerIndex{-1};
119 std::int32_t grantTimeoutTimeIndex{-1};
120 public:
122 std::atomic<bool> initRequested{false};
123 // temporary
124 std::atomic<bool> requestingMode{false};
125
126 std::atomic<bool> initIterating{false};
127
128 private:
129 bool iterating{false};
130 bool timeGranted_mode{false};
132 bool terminate_on_error{false};
136 TimeSynchronizationMethod timeMethod{TimeSynchronizationMethod::DISTRIBUTED};
138 std::uint32_t mGrantCount{0}; // this is intended to allow wrapping
140 std::shared_ptr<MessageTimer> mTimer;
142 gmlc::containers::BlockingQueue<ActionMessage> queue;
144 gmlc::containers::BlockingQueue<std::pair<std::string, std::string>> commandQueue;
146 std::atomic<uint16_t> interfaceFlags{0};
148 std::map<GlobalFederateId, std::deque<ActionMessage>> delayQueues;
149 std::vector<InterfaceHandle> events;
150 std::vector<InterfaceHandle> eventMessages;
151 std::vector<GlobalFederateId> delayedFederates;
152 Time time_granted{startupTime};
153 Time allowed_send_time{startupTime};
154 Time minimumReceiveTime{startupTime};
155
156#if __cplusplus >= 201703L
157 mutable std::atomic_flag processing{};
158#else
159 mutable std::atomic_flag processing = ATOMIC_FLAG_INIT;
160#endif
161
163 std::vector<std::function<std::string(std::string_view)>> queryCallbacks;
164 std::shared_ptr<FederateOperator> fedCallbacks;
165 std::vector<std::pair<std::string, std::string>> tags;
166 std::atomic<bool> queueProcessing{false};
168 Time nextValueTime() const;
170 Time nextMessageTime() const;
171
173 void setState(FederateStates newState);
174
176 bool messageShouldBeDelayed(const ActionMessage& cmd) const noexcept;
178 void addFederateToDelay(GlobalFederateId gid);
180 void generateConfig(nlohmann::json& base) const;
181
182 public:
184 void reset(const CoreFederateInfo& fedInfo);
186 const std::string& getIdentifier() const { return name; }
188 FederateStates getState() const;
190 InterfaceInfo& interfaces() { return interfaceInformation; }
192 const InterfaceInfo& interfaces() const { return interfaceInformation; }
193
195 uint64_t getQueueSize(InterfaceHandle hid) const;
198 uint64_t getQueueSize() const;
202 int32_t getCurrentIteration() const;
206 std::unique_ptr<Message> receive(InterfaceHandle hid);
209 std::unique_ptr<Message> receiveAny(InterfaceHandle& hid);
213 const std::shared_ptr<const SmallBuffer>& getValue(InterfaceHandle handle,
214 uint32_t* inputIndex);
215
220 const std::vector<std::shared_ptr<const SmallBuffer>>& getAllValues(InterfaceHandle handle);
221
223 std::pair<SmallBuffer, Time> getPublishedValue(InterfaceHandle handle);
225 void setParent(CommonCore* coreObject) { mParent = coreObject; }
230 void setProperties(const ActionMessage& cmd);
232 void setInterfaceProperty(const ActionMessage& cmd);
234 void setProperty(int timeProperty, Time propertyVal);
236 void setProperty(int intProperty, int propertyVal);
238 void setOptionFlag(int optionFlag, bool value);
240 Time getTimeProperty(int timeProperty) const;
242 bool getOptionFlag(int optionFlag) const;
244 int32_t getHandleOption(InterfaceHandle handle, char iType, int32_t option) const;
246 uint16_t getInterfaceFlags() const { return interfaceFlags.load(); }
248 int getIntegerProperty(int intProperty) const;
250 int publicationCount() const;
252 int endpointCount() const;
254 int inputCount() const;
256 void spinlock() const
257 {
258 while (processing.test_and_set()) {
259 ; // spin
260 }
261 }
263 void sleeplock() const
264 {
265 if (!processing.test_and_set()) {
266 return;
267 }
268 // spin for 10000 tries
269 for (int ii = 0; ii < 10000; ++ii) {
270 if (!processing.test_and_set()) {
271 return;
272 }
273 }
274 while (processing.test_and_set()) {
275 std::this_thread::yield();
276 }
277 }
279 void lock() { sleeplock(); }
280
282 bool try_lock() const { return !processing.test_and_set(); }
284 void unlock() const { processing.clear(); }
286 int loggingLevel() const;
287
289 void setTag(std::string_view tag, std::string_view value);
291 const std::string& getTag(std::string_view tag) const;
293 const std::pair<std::string, std::string>& getTagByIndex(size_t index) const
294 {
295 return tags[index];
296 }
298 auto tagCount() const { return tags.size(); }
300 bool isCallbackFederate() const { return mCallbackBased; }
301
302 private:
312 MessageProcessingResult processQueue() noexcept;
313
323 MessageProcessingResult processDelayQueue() noexcept;
324
326 checkProcResult(std::tuple<FederateStates, MessageProcessingResult, bool>& proc_result,
327 ActionMessage& cmd);
331 MessageProcessingResult processActionMessage(ActionMessage& cmd);
332
335 void processDataConnectionMessage(ActionMessage& cmd);
336
339 void processDataMessage(ActionMessage& cmd);
341 void timeoutCheck(ActionMessage& cmd);
343 void processLoggingMessage(ActionMessage& cmd);
347 void fillEventVectorUpTo(Time currentTime);
351 void fillEventVectorInclusive(Time currentTime);
355 void fillEventVectorNextIteration(Time currentTime);
357 void updateQueuedValueBytes();
359 void checkValueBufferWarning();
361 void addDependency(GlobalFederateId fedToDependOn);
363 void addDependent(GlobalFederateId fedThatDependsOnThis);
365 void resetDependency(GlobalFederateId gid);
366
368 int checkInterfaces();
370 std::string processQueryActual(std::string_view query) const;
374 void generateProfilingMessage(bool enterHelicsCode);
376 void generateProfilingMarker();
378 void updateMaxLogLevel();
379
381 void callbackProcessing() noexcept;
382 void callbackReturnResult(FederateStates lastState,
384 FederateStates newState) noexcept;
385 void initCallbackProcessing();
386 void execCallbackProcessing(IterationResult result);
388 void updateDataForExecEntry(MessageProcessingResult result, IterationRequest iterate);
390 void updateDataForTimeReturn(MessageProcessingResult result,
391 Time nextTime,
392 IterationRequest iterate);
393
394 public:
396 Time grantedTime() const { return time_granted; }
398 Time nextAllowedSendTime() const { return allowed_send_time; }
401 const std::vector<InterfaceHandle>& getEvents() const;
404 std::vector<GlobalFederateId> getDependencies() const;
407 std::vector<GlobalFederateId> getDependents() const;
409 const std::string& lastErrorString() const { return errorString; }
411 int lastErrorCode() const noexcept { return errorCode; }
413 void setCoreObject(CommonCore* parent);
414 // the next 5 functions are the processing functions that actually process the queue
430 iteration_time enterExecutingMode(IterationRequest iterate, bool sendRequest = false);
438 iteration_time requestTime(Time nextTime, IterationRequest iterate, bool sendRequest = false);
442 std::vector<GlobalHandle> getSubscribers(InterfaceHandle handle);
443
447 std::vector<std::pair<GlobalHandle, std::string_view>>
449
456 void finalize();
458 void processCommunications(std::chrono::milliseconds period);
460 void addAction(const ActionMessage& action);
462 void addAction(ActionMessage&& action);
464 std::optional<ActionMessage> processPostTerminationAction(const ActionMessage& action);
465
468
476 void logMessage(int level,
477 std::string_view logMessageSource,
478 std::string_view message,
479 bool fromRemote = false) const;
480
485 void setLogger(std::function<void(int, std::string_view, std::string_view)> logFunction);
486
489 void setCallbackOperator(std::shared_ptr<FederateOperator> fed)
490 {
491 fedCallbacks = std::move(fed);
492 }
493
497 void setQueryCallback(std::function<std::string(std::string_view)> queryCallbackFunction,
498 int order)
499 {
500 order = std::clamp(order, 1, 10);
501
502 if (static_cast<int>(queryCallbacks.size()) < order) {
503 queryCallbacks.resize(order);
504 }
505 queryCallbacks[order - 1] = std::move(queryCallbackFunction);
506 }
512 std::string processQuery(std::string_view query, bool force_ordering = false) const;
520 bool checkAndSetValue(InterfaceHandle pub_id, const char* data, uint64_t len);
521
523 void routeMessage(const ActionMessage& msg);
524
526 void routeMessage(ActionMessage&& msg);
529 InterfaceHandle handle,
530 std::string_view key,
531 std::string_view type,
532 std::string_view units,
533 uint16_t flags);
537 void sendCommand(ActionMessage& command);
538
540 std::pair<std::string, std::string> getCommand();
542 std::pair<std::string, std::string> waitCommand();
543};
544
545} // namespace helics
Definition ActionMessage.hpp:30
Definition CommonCore.hpp:74
Definition CoreFederateInfo.hpp:16
Definition FederateState.hpp:51
void setQueryCallback(std::function< std::string(std::string_view)> queryCallbackFunction, int order)
Definition FederateState.hpp:497
const std::string & getTag(std::string_view tag) const
Definition FederateState.cpp:2947
void setInterfaceProperty(const ActionMessage &cmd)
Definition FederateState.cpp:2029
std::atomic< bool > init_transmitted
Definition FederateState.hpp:93
std::pair< SmallBuffer, Time > getPublishedValue(InterfaceHandle handle)
Definition FederateState.cpp:350
void routeMessage(const ActionMessage &msg)
Definition FederateState.cpp:359
void spinlock() const
Definition FederateState.hpp:256
auto tagCount() const
Definition FederateState.hpp:298
std::atomic< bool > initRequested
Definition FederateState.hpp:122
void reset(const CoreFederateInfo &fedInfo)
Definition FederateState.cpp:163
std::optional< ActionMessage > processPostTerminationAction(const ActionMessage &action)
Definition FederateState.cpp:488
const std::vector< std::shared_ptr< const SmallBuffer > > & getAllValues(InterfaceHandle handle)
Definition FederateState.cpp:345
bool isCallbackFederate() const
Definition FederateState.hpp:300
void setCoreObject(CommonCore *parent)
Definition FederateState.cpp:2487
int inputCount() const
Definition FederateState.cpp:2407
std::vector< GlobalFederateId > getDependents() const
Definition FederateState.cpp:2417
const std::string & lastErrorString() const
Definition FederateState.hpp:409
std::pair< std::string, std::string > getCommand()
Definition FederateState.cpp:2671
void forceProcessMessage(ActionMessage &action)
Definition FederateState.cpp:501
iteration_time enterExecutingMode(IterationRequest iterate, bool sendRequest=false)
Definition FederateState.cpp:578
std::unique_ptr< Message > receive(InterfaceHandle hid)
Definition FederateState.cpp:302
const std::string & getIdentifier() const
Definition FederateState.hpp:186
LocalFederateId local_id
id code for the local federate descriptor
Definition FederateState.hpp:69
void closeInterface(InterfaceHandle handle, InterfaceType type)
Definition FederateState.cpp:445
void setProperties(const ActionMessage &cmd)
Definition FederateState.cpp:1988
Time nextAllowedSendTime() const
Definition FederateState.hpp:398
uint16_t getInterfaceFlags() const
Definition FederateState.hpp:246
MessageProcessingResult genericUnspecifiedQueueProcess(bool busyReturn)
Definition FederateState.cpp:946
void unlock() const
Definition FederateState.hpp:284
std::vector< GlobalFederateId > getDependencies() const
Definition FederateState.cpp:2412
int publicationCount() const
Definition FederateState.cpp:2397
void setOptionFlag(int optionFlag, bool value)
Definition FederateState.cpp:2181
void addAction(const ActionMessage &action)
Definition FederateState.cpp:389
void finalize()
Definition FederateState.cpp:985
std::vector< std::pair< GlobalHandle, std::string_view > > getMessageDestinations(InterfaceHandle handle)
Definition FederateState.cpp:692
IterationResult waitSetup()
Definition FederateState.cpp:511
uint64_t getQueueSize() const
Definition FederateState.cpp:287
int indexGroup
storage for index group location (this only matters on construction so can be public)
Definition FederateState.hpp:95
bool getOptionFlag(int optionFlag) const
Definition FederateState.cpp:2316
void logMessage(int level, std::string_view logMessageSource, std::string_view message, bool fromRemote=false) const
Definition FederateState.cpp:2494
const std::vector< InterfaceHandle > & getEvents() const
Definition FederateState.cpp:1047
std::pair< std::string, std::string > waitCommand()
Definition FederateState.cpp:2686
std::unique_ptr< Message > receiveAny(InterfaceHandle &hid)
Definition FederateState.cpp:311
FederateState(const FederateState &)=delete
FederateStates getState() const
Definition FederateState.cpp:241
iteration_time requestTime(Time nextTime, IterationRequest iterate, bool sendRequest=false)
Definition FederateState.cpp:701
int lastErrorCode() const noexcept
Definition FederateState.hpp:411
void sleeplock() const
Definition FederateState.hpp:263
int32_t getCurrentIteration() const
Definition FederateState.cpp:246
void setCallbackOperator(std::shared_ptr< FederateOperator > fed)
Definition FederateState.hpp:489
const std::shared_ptr< const SmallBuffer > & getValue(InterfaceHandle handle, uint32_t *inputIndex)
Definition FederateState.cpp:338
void createInterface(InterfaceType htype, InterfaceHandle handle, std::string_view key, std::string_view type, std::string_view units, uint16_t flags)
Definition FederateState.cpp:409
const InterfaceInfo & interfaces() const
Definition FederateState.hpp:192
void sendCommand(ActionMessage &command)
Definition FederateState.cpp:2552
void setParent(CommonCore *coreObject)
Definition FederateState.hpp:225
void setProperty(int timeProperty, Time propertyVal)
Definition FederateState.cpp:2095
IterationResult enterInitializingMode(IterationRequest request)
Definition FederateState.cpp:540
void processCommunications(std::chrono::milliseconds period)
Definition FederateState.cpp:1021
int getIntegerProperty(int intProperty) const
Definition FederateState.cpp:2377
void lock()
Definition FederateState.hpp:279
int loggingLevel() const
Definition FederateState.cpp:2928
std::vector< GlobalHandle > getSubscribers(InterfaceHandle handle)
Definition FederateState.cpp:678
int32_t getHandleOption(InterfaceHandle handle, char iType, int32_t option) const
Definition FederateState.cpp:2361
Time getTimeProperty(int timeProperty) const
Definition FederateState.cpp:2301
InterfaceInfo & interfaces()
Definition FederateState.hpp:190
bool try_lock() const
Definition FederateState.hpp:282
void setLogger(std::function< void(int, std::string_view, std::string_view)> logFunction)
Definition FederateState.cpp:296
Time grantedTime() const
Definition FederateState.hpp:396
const std::pair< std::string, std::string > & getTagByIndex(size_t index) const
Definition FederateState.hpp:293
bool checkAndSetValue(InterfaceHandle pub_id, const char *data, uint64_t len)
Definition FederateState.cpp:251
std::string processQuery(std::string_view query, bool force_ordering=false) const
Definition FederateState.cpp:2905
void setTag(std::string_view tag, std::string_view value)
Definition FederateState.cpp:2933
std::atomic< GlobalFederateId > global_id
global id code, default to invalid
Definition FederateState.hpp:70
int endpointCount() const
Definition FederateState.cpp:2402
Definition GlobalFederateId.hpp:75
Definition LocalFederateId.hpp:65
Definition InterfaceInfo.hpp:26
Definition LocalFederateId.hpp:22
@ HELICS_LOG_LEVEL_NO_PRINT
Definition helics_enums.h:206
the main namespace for the helics co-simulation library User functions will be in the helics namespac...
Definition AsyncFedCallInfo.hpp:14
IterationResult
Definition CoreTypes.hpp:94
FederateStates
Definition CoreTypes.hpp:21
MessageProcessingResult
Definition CoreTypes.hpp:75
constexpr Time timeZero
Definition helicsTime.hpp:31
TimeSynchronizationMethod
enumeration of possible time coordination methods
Definition FederateState.hpp:48
IterationRequest
Definition CoreTypes.hpp:103
@ NO_ITERATIONS
indicator that the iterations have completed
InterfaceType
Definition CoreTypes.hpp:112
TimeRepresentation< count_time< 9 > > Time
Definition helicsTime.hpp:27
STL namespace.
Definition helicsTime.hpp:46