mirror of
https://github.com/latentPrion/libspinscale.git
synced 2026-08-12 23:48:22 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
21fa125d52 | ||
|
|
3156652257 |
@@ -87,6 +87,7 @@ add_library(spinscale SHARED
|
|||||||
src/puppetApplication.cpp
|
src/puppetApplication.cpp
|
||||||
src/runtime.cpp
|
src/runtime.cpp
|
||||||
src/callableTracer.cpp
|
src/callableTracer.cpp
|
||||||
|
src/multiOperationResultSet.cpp
|
||||||
)
|
)
|
||||||
|
|
||||||
set_target_properties(spinscale PROPERTIES
|
set_target_properties(spinscale PROPERTIES
|
||||||
|
|||||||
@@ -202,17 +202,19 @@ nursery.launch(
|
|||||||
nvc.checkAndRethrowException();
|
nvc.checkAndRethrowException();
|
||||||
});
|
});
|
||||||
|
|
||||||
nursery.requestCancelOnAll();
|
|
||||||
nursery.closeAdmission();
|
nursery.closeAdmission();
|
||||||
|
nursery.requestCancelOnAll();
|
||||||
nursery.syncAwaitAllSettlements(
|
nursery.syncAwaitAllSettlements(
|
||||||
sscl::ComponentThread::getSelf()->getIoContext());
|
sscl::ComponentThread::getSelf()->getIoContext());
|
||||||
```
|
```
|
||||||
|
|
||||||
Each slot owns a `SyncCancelerForAsyncWork`. `requestCancelOnAll()` only signals
|
Each slot owns a `SyncCancelerForAsyncWork`. `requestCancelOnAll()` only signals
|
||||||
cooperative stop; it does not destroy invokers. Invokers are retired when their
|
cooperative stop; it does not destroy invokers. Invokers are retired when their
|
||||||
completion callbacks run. Call `closeAdmission()` explicitly before
|
completion callbacks run. Call `closeAdmission()` before `requestCancelOnAll()`
|
||||||
`asyncAwaitAllSettlements()` or `syncAwaitAllSettlements()`; those APIs wait until
|
so no new work can be admitted after cancel begins, and call `closeAdmission()`
|
||||||
all slots have retired naturally and throw if admission is still open.
|
explicitly before `asyncAwaitAllSettlements()` or `syncAwaitAllSettlements()`;
|
||||||
|
those APIs wait until all slots have retired naturally and throw if admission is
|
||||||
|
still open.
|
||||||
|
|
||||||
`syncAwaitAllSettlements()` runs a nested `io_context` loop on the **calling
|
`syncAwaitAllSettlements()` runs a nested `io_context` loop on the **calling
|
||||||
thread** (it blocks in `run_one()` until every slot has retired). Pass the
|
thread** (it blocks in `run_one()` until every slot has retired). Pass the
|
||||||
|
|||||||
@@ -55,9 +55,6 @@ struct MemberInvoker : MemberInvokerBase
|
|||||||
* nursery member. The external submitter should add the complete flow to the
|
* nursery member. The external submitter should add the complete flow to the
|
||||||
* nursery and then return; the nursery owns that flow until the flow settles.
|
* nursery and then return; the nursery owns that flow until the flow settles.
|
||||||
*
|
*
|
||||||
* Call closeAdmission() explicitly before asyncAwaitAllSettlements() or
|
|
||||||
* syncAwaitAllSettlements().
|
|
||||||
*
|
|
||||||
* syncAwaitAllSettlements() runs a nested io_context loop on the calling
|
* syncAwaitAllSettlements() runs a nested io_context loop on the calling
|
||||||
* thread (AsynchronousBridge). Pass the calling thread's io_context —
|
* thread (AsynchronousBridge). Pass the calling thread's io_context —
|
||||||
* typically
|
* typically
|
||||||
@@ -240,6 +237,36 @@ public:
|
|||||||
s.rsrc.admissionOpen = true;
|
s.rsrc.admissionOpen = true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** EXPLANATION:
|
||||||
|
* Stopping a nursery: always closeAdmission() before
|
||||||
|
* requestCancelOnAll(). requestCancelOnAll() only marks currently
|
||||||
|
* ACTIVE_UNSETTLED slots; it does not refuse new leases. If cancel
|
||||||
|
* runs while admission is still open, a concurrent submitter can still
|
||||||
|
* getNewSlotLease() / launch() after cancel has fanned out, and that
|
||||||
|
* newly admitted work will not have been cancelled — it races past the
|
||||||
|
* stop wave and keeps the drain from reaching "all settled" until it
|
||||||
|
* finishes on its own (or a later cancel). Closing admission first
|
||||||
|
* seals the nursery so cancel applies to a fixed membership set, then
|
||||||
|
* drain with asyncAwaitAllSettlements() / syncAwaitAllSettlements()
|
||||||
|
* (those APIs also require admission already closed).
|
||||||
|
*
|
||||||
|
* Preferred stop stack for a daemon/service that enqueues request
|
||||||
|
* coroutines into the nursery: keep protocol "stop listening /
|
||||||
|
* disconnect / refuse new connections and requests" separate from
|
||||||
|
* protocol state destruction. Stop accepting at the protocol level
|
||||||
|
* first, then nursery closeAdmission(), then requestCancelOnAll(),
|
||||||
|
* then cancel any awaited I/O owned outside the cancelers, then drain,
|
||||||
|
* then destroy protocol state. That ordering stops new work at the
|
||||||
|
* source before admission is sealed.
|
||||||
|
*
|
||||||
|
* If the daemon/service cannot disconnect/stop listening separately
|
||||||
|
* from destruction, the spinscale-using embedding project must handle
|
||||||
|
* closed-admission failures when it tries to enqueue. For example,
|
||||||
|
* catch the "admission closed" throw around nursery.launch() (or in
|
||||||
|
* the factory that calls it) and emit a protocol-specific failure such
|
||||||
|
* as "connection failed" or "request timed out" instead of letting the
|
||||||
|
* exception escape the accept/request path unbounded.
|
||||||
|
*/
|
||||||
void closeAdmission()
|
void closeAdmission()
|
||||||
{
|
{
|
||||||
sscl::SpinLock::Guard guard(s.lock);
|
sscl::SpinLock::Guard guard(s.lock);
|
||||||
|
|||||||
@@ -4,6 +4,9 @@
|
|||||||
#include <exception>
|
#include <exception>
|
||||||
|
|
||||||
namespace sscl {
|
namespace sscl {
|
||||||
|
namespace co {
|
||||||
|
struct Group;
|
||||||
|
} // namespace co
|
||||||
|
|
||||||
/** Plain aggregate for fan-out / fan-in results returned from coroutines. */
|
/** Plain aggregate for fan-out / fan-in results returned from coroutines. */
|
||||||
struct MultiOperationResultSet
|
struct MultiOperationResultSet
|
||||||
@@ -38,9 +41,29 @@ struct MultiOperationResultSetWithException
|
|||||||
memberFailureException(memberFailureExceptionIn)
|
memberFailureException(memberFailureExceptionIn)
|
||||||
{}
|
{}
|
||||||
|
|
||||||
|
/** Summarize a settled Group into counts + aggregated member failure. */
|
||||||
|
explicit MultiOperationResultSetWithException(const co::Group &group);
|
||||||
|
|
||||||
bool hasMemberFailure() const
|
bool hasMemberFailure() const
|
||||||
{ return memberFailureException != nullptr; }
|
{ return memberFailureException != nullptr; }
|
||||||
|
|
||||||
|
/** Combine this result set with another phase's counts and exception. */
|
||||||
|
MultiOperationResultSetWithException mergeWith(
|
||||||
|
const MultiOperationResultSetWithException &other) const
|
||||||
|
{
|
||||||
|
std::exception_ptr memberFailure = memberFailureException;
|
||||||
|
if (!memberFailure && other.hasMemberFailure()) {
|
||||||
|
memberFailure = other.memberFailureException;
|
||||||
|
}
|
||||||
|
|
||||||
|
return MultiOperationResultSetWithException(
|
||||||
|
MultiOperationResultSet(
|
||||||
|
results.nTotal + other.results.nTotal,
|
||||||
|
results.nSucceeded + other.results.nSucceeded,
|
||||||
|
results.nFailed + other.results.nFailed),
|
||||||
|
memberFailure);
|
||||||
|
}
|
||||||
|
|
||||||
MultiOperationResultSet results;
|
MultiOperationResultSet results;
|
||||||
std::exception_ptr memberFailureException = nullptr;
|
std::exception_ptr memberFailureException = nullptr;
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -0,0 +1,35 @@
|
|||||||
|
#include <boostAsioLinkageFix.h>
|
||||||
|
|
||||||
|
#include <spinscale/multiOperationResultSet.h>
|
||||||
|
#include <spinscale/co/group.h>
|
||||||
|
|
||||||
|
namespace sscl {
|
||||||
|
|
||||||
|
MultiOperationResultSetWithException::MultiOperationResultSetWithException(
|
||||||
|
const co::Group &group)
|
||||||
|
{
|
||||||
|
unsigned int nSucceeded = 0;
|
||||||
|
unsigned int nFailed = 0;
|
||||||
|
using SettlementType = co::Group::SettlementDescriptor::TypeE;
|
||||||
|
|
||||||
|
for (const auto &desc : group.s.rsrc.settlements)
|
||||||
|
{
|
||||||
|
if (desc.type == SettlementType::EXCEPTION_THROWN) {
|
||||||
|
nFailed++;
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
nSucceeded++;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
results = MultiOperationResultSet(
|
||||||
|
static_cast<unsigned int>(group.s.rsrc.settlements.size()),
|
||||||
|
nSucceeded,
|
||||||
|
nFailed);
|
||||||
|
|
||||||
|
if (nFailed > 0) {
|
||||||
|
memberFailureException = group.captureAggregatedGroupExceptions();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
} // namespace sscl
|
||||||
Reference in New Issue
Block a user