-
Notifications
You must be signed in to change notification settings - Fork 0
API Shuttle
𧬠API ⺠Shuttle
π Shuttles are neat little fairies that specialise in doing a thing and telling you how it went. You interact with them by simply .Ask()ing.
Shuttle is a pattern for performing an operation and reporting its outcome β save/load, network calls, path computation, moving the player, anything with a success/failure result.
-
Latency is not what defines a Shuttle.
Request β perform β reportis the shape; slowness is incidental. Most Shuttles are async because most interesting operations are, but a synchronous one is exactly as legitimate β see Synchronous Shuttles. -
The Shuttle is where the work happens. It is expected to touch the engine, hit the disk, call the network. Performing is its job, not a side effect it should feel bad about β
Response<T>carriesWasSuccessful, and that field is meaningless unless something was actually attempted (where engine work lives). - For live state rather than an operation (player health, game settings), prefer
[HasState]matter types instead. The contrast there is state versus operation, not sync versus async. - Shuttle operates on two specialised Matter types (
RequestandResponse<T>), it is a round-trip pair so you need to define both.
/*
* ππ§ Since the two types are useless apart, it helps to keep them in one .cs file.
* SaveGameRR or SaveGameRequestResponse, whichever you prefer.
*/
class SaveGameRequest : Request { }
class SaveGameResponse : Response<SaveGameRequest>
{
public SaveGameResponse(SaveGameRequest request, bool wasSuccessful)
// The `: base(request, wasSuccessful)` constructor call is required.
: base(request, wasSuccessful) { }
}- The response carries a reference back to the original request.
- That's what lets
Askroute reply to the correct caller (viaIsRespondingTo) when multiple requests of same type are in flight. - It's also how the request is recorded as response's cause automatically, correct even across an async boundary.
- That's what lets
π Register a handler with Shuttle which will process Requests and provide Responses to them.
- Most Shuttles cross an async boundary. When yours does, wrap the operation so that unsubscription can cancel it:
-
Observable.FromAsync(ct => ...)for aTask-based API β one Task, one value, token wired for you. -
Observable.Create<T>(async (observer, ct) => ...)when you need several async steps or more than one emission. -
Observable.Create<T>(observer => ...)for a callback-style API β return a disposable that really cancels, neverDisposable.Empty.
-
- Note there is no manual circumstance stamping here β the response's instance carries
req, and rzeka records it as the cause for you.- You only stamp circumstances on a Shuttle response when it has other causes apart from the request β e.g. ambient state you pulled in via
Scry(see below Multi-context Response using Scry). - The request itself is always handled automatically β stamping it too is harmless, rzeka dedupes it.
- You only stamp circumstances on a Shuttle response when it has other causes apart from the request β e.g. ambient state you pulled in via
- Example:
Q += rzeka.Shuttle<SaveGameRequest, SaveGameResponse>(
this,
reqs => reqs.SelectMany(req =>
Observable.Create<SaveGameResponse>(observer =>
{
var cts = new CancellationTokenSource();
_saveSystem.SaveAsync(cts.Token, success =>
{
observer.OnNext(new SaveGameResponse(req, success));
observer.OnCompleted();
});
// Runs when this inner sequence is disposed. With the SelectMany above
// that means one thing only: the spell itself being torn down. An operator
// that drops inner sequences - Switch, Take, TakeUntil - would trigger it
// too, but there is none in this chain, and none on the Ask side can reach
// in here. Cancellation does not cross the river.
return Disposable.Create(() => { cts.Cancel(); cts.Dispose(); });
})
)
);π When the operation has no latency, the Shuttle is the same shape without the wrapping β the work happens in the chain, and the response reports that it happened.
Q += rzeka.Shuttle<MovePlayerRequest, MovePlayerResponse>(
this,
reqs => reqs.Select(req =>
{
// The response below asserts the player HAS been moved.
// That is only true because this ran first.
_playerBody.GlobalPosition = req.Position;
_playerBody.Velocity = req.LinearVelocity;
_cameraRig.SetRotationAngles(req.Pitch, req.Yaw);
return new MovePlayerResponse(req, true);
})
);ππΉ If you find yourself writing a Loom that performs work and then reports it, that is a Shuttle typed out by hand. Reach for the Shuttle instead β you get request/response correlation for free, and the effect stops needing an apology.
ππ§΅ The main-thread hop happens before rzeka publishes the response, not before your lambda. Synchronous work taken straight off reqs runs on whatever thread plucked the request β normally the main thread, but that is a property of your callers, not a guarantee from the Shuttle. If a request can arrive from a worker, put .ObserveOn(rzeka.MainThread) in front of the engine work (ingredient or consequence?).
π Use Ask extension method on rzeka to send Requests to Shuttles.
- It sends your request into the river.
- Returns an observable that emits only the response to your specific request β not responses to other concurrent requests of the same type or even ones coming from the same spell where this
Askis defined. - When Ask(ing) you also have to manually stamp the request matter circumstances with
.WithCircumstances(...circumstances).
π Ask completes after the response. One request, one verdict β it takes the first correlated response and closes, which disposes the Weave it registered. You do not need .Take(1).
π𧨠Build the Ask inside the lambda, never outside it. Each call constructs a fresh observable, which is what lets repeated triggers each get their own round trip. Hoisting one out and re-subscribing re-plucks the same request instance into the river. todo: is this really worth mentioning, isn't it a ridiculous idea, that of course if an ask is hoisted out with a single specific request and not the new coming in, it will
// β a new Ask per trigger
evts.SelectMany(evt => rzeka.Ask<Req, Res>(this, new Req(evt)))
// β one Ask, re-subscribed - the same request goes into the river again
var ask = rzeka.Ask<Req, Res>(this, req);
evts.SelectMany(_ => ask)π Flattening Asks: an Ask is one round trip, so inside a spell it is always flattened into the chain.
-
.SelectMany(evt => rzeka.Ask(...))- requests overlap, responses arrive as they complete -
.Select(evt => rzeka.Ask(...)).Concat()- requests queue up, one at a time, in order -
.Select(evt => rzeka.Ask(...)).Switch()- a new trigger discards the response still in flight
ππΉ A Shuttle can physically answer one request more than once β nothing stops its lambda emitting repeatedly β but Ask takes the first verdict and closes, because that is what Response<T> means. Recurring side-channel data (load progress, partial results) belongs in its own matter type, where it can be named honestly and consumed by whoever cares; squeezed through a response it inherits a WasSuccessful flag that means nothing until the last one.
π𧨠Switch on the Ask side drops the response. It does not cancel the work.
-
Asksubscribes a Weave to the response stream and then plucks the request into the river. Disposing that subscription tears down the Weave β the listener. The request is already in the river by then, and the Shuttle answering it is a separate, long-lived spell, so nothing reaches back into the work in flight. -
This holds even if the Shuttle wrapped its work in
Observable.FromAsync(ct => ...)and honours the token. The river is a hard cancellation boundary; a token cannot cross it. Do not let the presence of actconvince you the work stopped. - What you get is correct latest-wins semantics at the consumer: the stale response never enters your chain. What you pay is the full network and CPU cost of the request you abandoned.
- If that cost is unacceptable, the cancellation has to live inside the Shuttle β see Cancelling inside a Shuttle.
// Inside a Loom: on level completion, save then show results
Q += rzeka.Loom<LevelCompletedEvent, GameSaved>(
this,
levelCompletedEvent => levelCompletedEvent.SelectMany(evt =>
rzeka.Ask<SaveGameRequest, SaveGameResponse>(
this,
new SaveGameRequest(evt.Score, evt.Timestamp) // instantiating your request
.WithCircumstances(evt)) // manually stamped Ask request's circumstances
.Select(saveResponse => new GameSaved(evt.Score, save.WasSuccessful)
.WithCircumstances(evt, saveResponse))); // if you omit this stamp only 'evt' will be its circumstance
);π Ask responses are invisible to Loom's auto-tracking.
- Loom only auto-stamps matter from its declared input slots - the
LevelCompletedEventhere. - The save response comes from inside the lambda, so without the manual stamp it silently vanishes from the causal graph. See circumstance rules.
π Don't nest Asks.
- Every Ask nested inside a lambda is another cause auto-tracking cannot see, so each hop needs another own manual stamp β messy β and in Eris the whole chain collapses into a single spell occurrence.
- This makes you lose the step-by-step story.
- The solution is to decompose into separate Looms and Shuttles.
π The only place a CancellationToken actually bites is inside the Shuttle's own chain, where Switch owns the inner sequence it is dropping.
Q += rzeka.Shuttle<SearchRequest, SearchResponse>(
this,
reqs => reqs
.Select(req => Observable.FromAsync(ct => _index.SearchAsync(req.Term, ct))
.Select(hits => new SearchResponse(req, hits, true)))
.Switch() // a newer request cancels the older one's token for real
);π𧨠This breaks the Shuttle's contract, and rzeka will not warn you.
- A Shuttle is meant to be total: every request gets a response, success or failure.
Switchcancels the superseded request without answering it. - The caller that sent it is left with an
Askthat never emits, never completes and never errors β the "AnAskthat never returns" symptom, self-inflicted. - Only reach for it when every caller is the same spell re-asking, so the abandoned caller is one you already replaced.
ππ§ If callers are independent, make losing explicit instead β cancel the work and answer:
reqs => reqs.SelectMany(req =>
Observable.FromAsync(ct => _index.SearchAsync(req.Term, ct))
.Select(hits => new SearchResponse(req, hits, true))
// A newer request ends this one early...
.TakeUntil(reqs)
// ...and this guarantees the caller still hears back.
.DefaultIfEmpty(new SearchResponse(req, null, false))
);ππΉ If you want latest-wins but the work is genuinely cancellable and nobody else asks, a Loom with .Switch() avoids this problem entirely β it keeps correct causality and really cancels, because it never crosses the river. See Async Operations.
When the response depends on more than just the request, pull the additional matter in with Scry β sampled when the request arrives, before the async work starts.
ππΉ Only Scry [HasState] types as ambient context.
-
WithLatestFromMattersilently drops the trigger if a context stream has not emitted yet β in a Shuttle that means the request vanishes, no response is ever produced, and the Ask caller waits forever with no error. -
[HasState]matter always carries a current value, so this cannot happen.
π Example:
Q += rzeka.Shuttle<LoadSceneRequest, LoadSceneResponse>(
this,
reqs => reqs
/*
* The request stays the trigger, Scry'd matter is read as ambient context.
* The context is sampled at trigger time and captured in the closure,
* so it cannot drift while the async work runs.
*/
.WithLatestFromMatter(rzeka.Scry<GameState>(), rzeka.Scry<Settings>())
.SelectMany(ctx =>
{
var (req, state, settings) = ctx;
return LoadSceneThreaded(req, state, settings) // your Observable.Create
.Select(scene => new LoadSceneResponse(req, scene, true)
// Stamp only the extra Scry'd context. The request is recorded
// automatically, so the response ends up with [request, state, settings].
.WithCircumstances<LoadSceneResponse>(state, settings))
.Catch<LoadSceneResponse, Exception>(ex =>
{
rzeka.Whisper(ex, req);
// The failure response needs the same stamps as the success one
return Observable.Return(new LoadSceneResponse(req, null, false)
.WithCircumstances<LoadSceneResponse>(state, settings));
});
})
);π Keep .Catch inside the SelectMany.
- Scoped there, a failure ends one request and the Shuttle lives on.
- Moved outside, the first failure kills the whole spell and every later request goes unanswered.
See also: Scry · Loom · Async Operations · circumstance rules · 𧬠API overview