Skip to content

Commit d5cd91e

Browse files
committed
revise how cancel_transaction works
1 parent 9b579c0 commit d5cd91e

11 files changed

Lines changed: 163 additions & 135 deletions

File tree

‎ChangeLog‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -20,12 +20,12 @@ Improvements:
2020

2121
API Changes:
2222
- New Engine method:
23-
- Engine::cancel_transaction() cancel the currently ongoing transaction
24-
performed by this engine. This method has to be called by an external
25-
actor not involved in the transaction. It will make the transaction raise
26-
a DTLMod::TransactionCancelledException that must be caught if you want
27-
the publishers and subscribers to survive to the cancellation of this
28-
transaction.
23+
- Engine::cancel_transaction(unsigned int transaction_id) cancels a
24+
specific transaction performed by this engine. This method has to be
25+
called by an external actor not involved in the transaction. It will make
26+
the transaction raise a DTLMod::TransactionCanceledException that must be
27+
caught if you want the publishers and subscribers to survive to the
28+
cancelation of this transaction.
2929
- New Stream helper methods:
3030
- Stream::get_engine_type() and Stream::get_transport_method() respectively
3131
return the enum value of the engine type and transport method for a

‎include/dtlmod/DTLException.hpp‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ DECLARE_DTLMOD_EXCEPTION(UnknownCompressionOptionException, "Unknown Compression
7070
DECLARE_DTLMOD_EXCEPTION(InconsistentCompressionRatioException, "Inconsistent Compression ratio");
7171
DECLARE_DTLMOD_EXCEPTION(SubscriberSideCompressionException, "Compression can only be applied on the publisher side");
7272

73-
DECLARE_DTLMOD_EXCEPTION(TransactionCancelledException, "Transaction cancelled");
73+
DECLARE_DTLMOD_EXCEPTION(TransactioncanceledException, "Transaction canceled");
7474

7575
} // namespace dtlmod
7676

‎include/dtlmod/Engine.hpp‎

Lines changed: 13 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ class Engine {
5757
std::weak_ptr<Stream> stream_;
5858

5959
bool pub_ever_present_ = false;
60-
std::atomic<bool> cancelled_{false};
60+
std::atomic<unsigned int> canceled_transaction_id_{0};
6161

6262
ActorRegistry publishers_;
6363

@@ -77,8 +77,9 @@ class Engine {
7777
[[nodiscard]] const sg4::ActivitySet& get_sub_transaction() const noexcept { return sub_transaction_; }
7878
[[nodiscard]] sg4::ActivitySet& get_sub_transaction() noexcept { return sub_transaction_; }
7979

80-
// Protected virtual method for derived classes to implement
81-
[[nodiscard]] virtual unsigned int get_current_transaction_impl() const noexcept = 0;
80+
// Protected virtual methods for derived classes to implement
81+
[[nodiscard]] virtual unsigned int get_current_transaction_impl() const noexcept = 0;
82+
[[nodiscard]] virtual unsigned int get_current_sub_transaction_impl() const noexcept = 0;
8283

8384
// Protected methods for derived classes only
8485
void close_stream() const;
@@ -93,7 +94,11 @@ class Engine {
9394

9495
[[nodiscard]] bool pub_ever_present() const noexcept { return pub_ever_present_; }
9596

96-
[[nodiscard]] bool is_cancelled() const noexcept { return cancelled_; }
97+
[[nodiscard]] bool is_canceled() const noexcept { return canceled_transaction_id_ != 0; }
98+
[[nodiscard]] bool is_transaction_canceled(unsigned int tx_id) const noexcept
99+
{
100+
return canceled_transaction_id_ == tx_id;
101+
}
97102

98103
// Pure virtual methods for derived classes to implement
99104
virtual void create_transport(const Transport::Method& transport_method) = 0;
@@ -144,9 +149,11 @@ class Engine {
144149
/// @return The id of the ongoing transaction.
145150
[[nodiscard]] unsigned int get_current_transaction() const noexcept { return get_current_transaction_impl(); }
146151

147-
/// @brief Cancel all in-flight activities of the current transaction, unblocking publishers and subscribers.
152+
/// @brief Cancel all in-flight activities of a specific transaction, unblocking publishers and subscribers.
153+
/// @param transaction_id The id of the transaction to cancel. If both sides have already moved past this
154+
/// transaction, the call is a no-op to avoid accidentally cancelling a subsequent transaction.
148155
/// @note Must be called from an external actor not participating in the transaction.
149-
void cancel_transaction();
156+
void cancel_transaction(unsigned int transaction_id);
150157

151158
/// @brief Close the Engine associated to a Stream.
152159
void close();

‎include/dtlmod/FileEngine.hpp‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,10 @@ class FileEngine : public Engine {
5555
{
5656
return current_pub_transaction_id_;
5757
}
58+
[[nodiscard]] unsigned int get_current_sub_transaction_impl() const noexcept override
59+
{
60+
return current_sub_transaction_id_;
61+
}
5862

5963
protected:
6064
[[nodiscard]] std::shared_ptr<FileTransport> get_file_transport() const;

‎include/dtlmod/StagingEngine.hpp‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,10 @@ class StagingEngine : public Engine {
4444
{
4545
return current_pub_transaction_id_;
4646
}
47+
[[nodiscard]] unsigned int get_current_sub_transaction_impl() const noexcept override
48+
{
49+
return current_sub_transaction_id_;
50+
}
4751

4852
protected:
4953
[[nodiscard]] std::shared_ptr<StagingTransport> get_staging_transport() const;

‎src/Engine.cpp‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -86,9 +86,15 @@ void Engine::close()
8686
publishers_.contains(sg4::Actor::self()) ? pub_close() : sub_close();
8787
}
8888

89-
void Engine::cancel_transaction()
89+
void Engine::cancel_transaction(unsigned int transaction_id)
9090
{
91-
cancelled_ = true;
91+
// No-op if both sides have already moved past the target transaction: cancelling now would
92+
// affect the wrong (subsequent) transaction.
93+
if (get_current_transaction_impl() > transaction_id && get_current_sub_transaction_impl() > transaction_id)
94+
return;
95+
// transaction_id == 0 means no transaction has started yet; treat as T1 so all checks against
96+
// (current + 1 >= 1) still fire correctly.
97+
canceled_transaction_id_.store(transaction_id == 0 ? 1 : transaction_id);
9298
cancel_activities();
9399
}
94100

‎src/FileEngine.cpp‎

Lines changed: 20 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -100,8 +100,8 @@ std::string FileEngine::get_path_to_dataset() const
100100

101101
void FileEngine::begin_pub_transaction()
102102
{
103-
if (is_cancelled())
104-
throw TransactionCancelledException(XBT_THROW_POINT);
103+
if (is_transaction_canceled(current_pub_transaction_id_ + 1))
104+
throw TransactioncanceledException(XBT_THROW_POINT);
105105

106106
auto self = sg4::Actor::self();
107107

@@ -115,12 +115,12 @@ void FileEngine::begin_pub_transaction()
115115
// Wait for the completion of the Publish activities from the previous transaction
116116
XBT_DEBUG("Wait for the completion of %u publish activities from the previous transaction",
117117
file_pub_transaction_[self].size());
118-
while (!is_cancelled() && file_pub_transaction_[self].size() > 0) {
118+
while (!is_canceled() && file_pub_transaction_[self].size() > 0) {
119119
std::unique_lock lock(*(get_publishers().get_mutex()));
120120
pub_activities_completed_->wait(lock);
121121
}
122-
if (is_cancelled())
123-
throw TransactionCancelledException(XBT_THROW_POINT);
122+
if (is_transaction_canceled(current_pub_transaction_id_))
123+
throw TransactioncanceledException(XBT_THROW_POINT);
124124
XBT_DEBUG("All on-flight publish activities are completed. Proceed with the current transaction.");
125125
get_file_transport()->clear_to_write_in_transaction(self);
126126
}
@@ -169,7 +169,7 @@ void FileEngine::pub_close()
169169

170170
XBT_DEBUG("[%s] Wait for the completion of %u publish activities from the previous transaction", get_cname(),
171171
file_pub_transaction_[self].size());
172-
while (!is_cancelled() && file_pub_transaction_[self].size() > 0) {
172+
while (!is_transaction_canceled(current_pub_transaction_id_) && file_pub_transaction_[self].size() > 0) {
173173
std::unique_lock lock(*(get_publishers().get_mutex()));
174174
pub_activities_completed_->wait(lock);
175175
}
@@ -191,8 +191,8 @@ void FileEngine::pub_close()
191191

192192
void FileEngine::begin_sub_transaction()
193193
{
194-
if (is_cancelled())
195-
throw TransactionCancelledException(XBT_THROW_POINT);
194+
if (is_transaction_canceled(current_sub_transaction_id_ + 1))
195+
throw TransactioncanceledException(XBT_THROW_POINT);
196196

197197
// Only one subscriber has to do this
198198
if (!sub_transaction_in_progress_) {
@@ -204,12 +204,15 @@ void FileEngine::begin_sub_transaction()
204204
// We have publishers on that stream, wait for them to complete a transaction first
205205
if (not get_publishers().is_empty()) {
206206
std::unique_lock lock(*get_subscribers().get_mutex());
207-
while (!is_cancelled() && completed_pub_transaction_id_ < current_sub_transaction_id_) {
207+
while (!is_transaction_canceled(current_sub_transaction_id_) &&
208+
completed_pub_transaction_id_ < current_sub_transaction_id_) {
208209
XBT_DEBUG("Wait for publishers to end the transaction I need");
209210
pub_transaction_completed_->wait(lock);
210211
}
211-
if (is_cancelled())
212-
throw TransactionCancelledException(XBT_THROW_POINT);
212+
if (is_transaction_canceled(current_sub_transaction_id_)) {
213+
sub_transaction_in_progress_ = false;
214+
throw TransactioncanceledException(XBT_THROW_POINT);
215+
}
213216
XBT_DEBUG("Publishers stored metadata for that transaction, proceed");
214217
}
215218
}
@@ -221,17 +224,17 @@ void FileEngine::end_sub_transaction()
221224

222225
// The files subscribers need to read may not have been fully written. Wait to be notified completion of the publish
223226
// activities
224-
if (!is_cancelled() && current_sub_transaction_id_ == current_pub_transaction_id_ &&
225-
not get_publishers().is_empty()) {
227+
if (!is_transaction_canceled(current_sub_transaction_id_) &&
228+
current_sub_transaction_id_ == current_pub_transaction_id_ && not get_publishers().is_empty()) {
226229
XBT_DEBUG("Wait for the completion of publish activities from the current transaction");
227230
pub_activities_completed_->wait(std::unique_lock(*get_subscribers().get_mutex()));
228231
XBT_DEBUG("All on-flight publish activities are completed. Proceed with the subscribe activities.");
229232
}
230-
if (is_cancelled()) {
233+
if (is_transaction_canceled(current_sub_transaction_id_)) {
231234
transport->close_sub_files(self);
232235
transport->clear_to_read_in_transaction(self);
233236
sub_transaction_in_progress_ = false;
234-
throw TransactionCancelledException(XBT_THROW_POINT);
237+
throw TransactioncanceledException(XBT_THROW_POINT);
235238
}
236239

237240
// Subscriber get the list of files and size to read that has been build during the get() operations
@@ -243,12 +246,12 @@ void FileEngine::end_sub_transaction()
243246

244247
XBT_DEBUG("Wait for the %d subscribe activities for the transaction", file_sub_transaction_[self].size());
245248
file_sub_transaction_[self].wait_all();
246-
if (is_cancelled()) {
249+
if (is_transaction_canceled(current_sub_transaction_id_)) {
247250
file_sub_transaction_[self].clear();
248251
transport->close_sub_files(self);
249252
transport->clear_to_read_in_transaction(self);
250253
sub_transaction_in_progress_ = false;
251-
throw TransactionCancelledException(XBT_THROW_POINT);
254+
throw TransactioncanceledException(XBT_THROW_POINT);
252255
}
253256
file_sub_transaction_[self].clear();
254257
// Close files opened in this transaction

‎src/StagingEngine.cpp‎

Lines changed: 23 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -55,8 +55,8 @@ std::shared_ptr<StagingTransport> StagingEngine::get_staging_transport() const
5555

5656
void StagingEngine::begin_pub_transaction()
5757
{
58-
if (is_cancelled())
59-
throw TransactionCancelledException(XBT_THROW_POINT);
58+
if (is_transaction_canceled(current_pub_transaction_id_ + 1))
59+
throw TransactioncanceledException(XBT_THROW_POINT);
6060

6161
if (!pub_transaction_in_progress_) {
6262
pub_transaction_in_progress_ = true;
@@ -78,24 +78,24 @@ void StagingEngine::begin_pub_transaction()
7878
try {
7979
get_pub_transaction().wait_all();
8080
} catch (const simgrid::NetworkFailureException&) {
81-
if (!is_cancelled())
81+
if (!is_canceled())
8282
throw;
8383
}
8484
XBT_DEBUG("All on-flight publish activities are completed. Proceed with the current transaction.");
8585
XBT_DEBUG("%u sub activities pending", get_sub_transaction().size());
8686
get_pub_transaction().clear();
87-
if (is_cancelled())
88-
throw TransactionCancelledException(XBT_THROW_POINT);
87+
if (is_transaction_canceled(current_pub_transaction_id_))
88+
throw TransactioncanceledException(XBT_THROW_POINT);
8989
}
9090

9191
// Then we wait for all subscribers to be at the same transaction
92-
while (!is_cancelled() &&
92+
while (!is_transaction_canceled(current_pub_transaction_id_) &&
9393
(get_subscribers().is_empty() || current_pub_transaction_id_ > current_sub_transaction_id_)) {
9494
XBT_DEBUG("Wait for subscribers");
9595
sub_transaction_started_->wait(lock);
9696
}
97-
if (is_cancelled())
98-
throw TransactionCancelledException(XBT_THROW_POINT);
97+
if (is_transaction_canceled(current_pub_transaction_id_))
98+
throw TransactioncanceledException(XBT_THROW_POINT);
9999
// Publisher has been notified by subscribers, it can proceed with the transaction
100100
}
101101

@@ -146,16 +146,16 @@ void StagingEngine::pub_close()
146146

147147
void StagingEngine::begin_sub_transaction()
148148
{
149-
if (is_cancelled())
150-
throw TransactionCancelledException(XBT_THROW_POINT);
149+
if (is_transaction_canceled(current_sub_transaction_id_ + 1))
150+
throw TransactioncanceledException(XBT_THROW_POINT);
151151

152152
if (current_sub_transaction_id_ == 0) { // This is the first transaction
153153
// Wait for at least one publisher to start a tran
154154
std::unique_lock lock(*get_subscribers().get_mutex());
155-
while (!is_cancelled() && current_pub_transaction_id_ == 0)
155+
while (!is_transaction_canceled(current_sub_transaction_id_ + 1) && current_pub_transaction_id_ == 0)
156156
first_pub_transaction_started_->wait(lock);
157-
if (is_cancelled())
158-
throw TransactionCancelledException(XBT_THROW_POINT);
157+
if (is_transaction_canceled(current_sub_transaction_id_ + 1))
158+
throw TransactioncanceledException(XBT_THROW_POINT);
159159
XBT_DEBUG("Publishers have started a transaction, create rendez-vous points");
160160
// We now know the number of publishers, subscriber can create mailboxes/mqs with publishers
161161
get_staging_transport()->create_rendez_vous_points();
@@ -178,10 +178,13 @@ void StagingEngine::begin_sub_transaction()
178178
}
179179

180180
std::unique_lock lock(*get_subscribers().get_mutex());
181-
while (!is_cancelled() && completed_pub_transaction_id_ < current_sub_transaction_id_)
181+
while (!is_transaction_canceled(current_sub_transaction_id_) &&
182+
completed_pub_transaction_id_ < current_sub_transaction_id_)
182183
pub_transaction_completed_->wait(lock);
183-
if (is_cancelled())
184-
throw TransactionCancelledException(XBT_THROW_POINT);
184+
if (is_transaction_canceled(current_sub_transaction_id_)) {
185+
sub_transaction_in_progress_ = false;
186+
throw TransactioncanceledException(XBT_THROW_POINT);
187+
}
185188
}
186189

187190
void StagingEngine::end_sub_transaction()
@@ -195,19 +198,19 @@ void StagingEngine::end_sub_transaction()
195198
try {
196199
get_sub_transaction().wait_all();
197200
} catch (const simgrid::CancelException&) {
198-
if (!is_cancelled())
201+
if (!is_canceled())
199202
throw;
200203
get_sub_transaction().clear();
201204
sub_transaction_in_progress_ = false;
202205
num_subscribers_starting_--;
203-
throw TransactionCancelledException(XBT_THROW_POINT);
206+
throw TransactioncanceledException(XBT_THROW_POINT);
204207
} catch (const simgrid::NetworkFailureException&) {
205-
if (!is_cancelled())
208+
if (!is_canceled())
206209
throw;
207210
get_sub_transaction().clear();
208211
sub_transaction_in_progress_ = false;
209212
num_subscribers_starting_--;
210-
throw TransactionCancelledException(XBT_THROW_POINT);
213+
throw TransactioncanceledException(XBT_THROW_POINT);
211214
}
212215
XBT_DEBUG("All on-flight subscribe activities are completed. Proceed with the current transaction.");
213216
get_sub_transaction().clear();

‎src/bindings/python/dtlmod_python.cpp‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ PYBIND11_MODULE(dtlmod, m)
8888
py::register_exception<dtlmod::InconsistentCompressionRatioException>(m, "InconsistentCompressionRatioException");
8989
py::register_exception<dtlmod::SubscriberSideCompressionException>(m, "SubscriberSideCompressionException");
9090

91-
py::register_exception<dtlmod::TransactionCancelledException>(m, "TransactionCancelledException");
91+
py::register_exception<dtlmod::TransactioncanceledException>(m, "TransactioncanceledException");
9292

9393
/* Class Engine */
9494
py::class_<Engine, std::shared_ptr<Engine>> engine(
@@ -108,7 +108,8 @@ PYBIND11_MODULE(dtlmod, m)
108108
.def_property_readonly("current_transaction", &Engine::get_current_transaction,
109109
"The id of the current transaction on this Engine (read-only)")
110110
.def("cancel_transaction", &Engine::cancel_transaction, py::call_guard<py::gil_scoped_release>(),
111-
"Cancel all in-flight activities of the current transaction (must be called from an external actor)")
111+
py::arg("transaction_id"),
112+
"Cancel all in-flight activities of a specific transaction (must be called from an external actor)")
112113
.def("close", &Engine::close, py::call_guard<py::gil_scoped_release>(), "Close this Engine");
113114

114115
py::enum_<Engine::Type>(engine, "Type", "The type of Engine")

0 commit comments

Comments
 (0)