1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
// Types re-exported from crate root
use futures::StreamExt as _;
use futures::channel::mpsc;
use crate::jsonrpc::protocol_compat::ProtocolCompat;
use crate::jsonrpc::{OutgoingMessage, PendingReplies, RawJsonRpcMessage, TransportFrame};
use crate::schema::v1::RequestId;
pub type OutgoingMessageTx = mpsc::UnboundedSender<OutgoingMessage>;
pub(crate) fn send_raw_message(
tx: &OutgoingMessageTx,
message: OutgoingMessage,
) -> Result<(), crate::Error> {
tracing::debug!(?message, ?tx, "send_raw_message");
tx.unbounded_send(message)
.map_err(crate::util::internal_error)
}
/// Outgoing protocol actor: Converts application-level OutgoingMessage to protocol-level RawJsonRpcMessage.
///
/// This actor handles JSON-RPC protocol semantics:
/// - Verifies that outgoing requests still have pending response registrations
/// - Converts OutgoingMessage variants to RawJsonRpcMessage
///
/// This is the protocol layer - it has no knowledge of how messages are transported.
pub(super) async fn outgoing_protocol_actor(
mut outgoing_rx: mpsc::UnboundedReceiver<OutgoingMessage>,
pending_replies: PendingReplies,
transport_tx: mpsc::UnboundedSender<TransportFrame>,
protocol_compat: ProtocolCompat,
foreground_done: super::SharedCompletionSignal,
) -> Result<(), crate::Error> {
let mut drain_waiters = Vec::new();
while let Some(message) = outgoing_rx.next().await {
tracing::debug!(?message, "outgoing_protocol_actor");
// Create the message to be sent over the transport
let (json_rpc_message, destination) = match message {
OutgoingMessage::CloseAfterDraining { done } => {
// Reject later sends while preserving every message that was
// already accepted into this receiver's buffer.
outgoing_rx.close();
drain_waiters.push(done);
continue;
}
OutgoingMessage::BatchDispatchComplete { completion } => {
if let Some(frame) = completion.complete() {
transport_tx
.unbounded_send(frame)
.map_err(crate::Error::into_internal_error)?;
}
continue;
}
OutgoingMessage::BatchHandlerAttemptComplete { destination } => {
if let Some(frame) = destination.finish_handler_attempt() {
transport_tx
.unbounded_send(frame)
.map_err(crate::Error::into_internal_error)?;
}
continue;
}
OutgoingMessage::AbandonedBatchResponse {
id,
method,
destination,
} => {
tracing::warn!(
?id,
%method,
"Completing abandoned JSON-RPC batch request with Internal Error"
);
let fallback = protocol_compat.outgoing_response_to(
&id,
&method,
Err(crate::Error::internal_error().data(format!(
"request handler dropped its responder for `{method}`"
))),
);
let fallback = RawJsonRpcMessage::response(id, fallback);
if let Some(frame) = destination.abandon(fallback) {
transport_tx
.unbounded_send(frame)
.map_err(crate::Error::into_internal_error)?;
}
continue;
}
OutgoingMessage::Request {
id,
method,
untyped,
remote_style,
readiness,
} => {
// Requests register their response destination synchronously
// before entering this queue. EOF removes that registration,
// so skip work that can no longer receive a response.
if !pending_replies.contains(&id) {
continue;
}
if let Some(readiness) = readiness {
// A route-installation gate may depend on an application
// handler that never returns. Foreground success cancels
// unresolved gates, but ready gates still publish in FIFO
// order; shutdown must never bypass route readiness.
// Poll a fresh clone for each gate: a Shared handle that
// returned Ready cannot itself be polled a second time.
let result =
match futures::future::select(Box::pin(readiness), foreground_done.clone())
.await
{
futures::future::Either::Left((result, _)) => result,
futures::future::Either::Right(((), _)) => {
Err(crate::Error::internal_error()
.data("foreground completed before outgoing request readiness"))
}
};
if let Err(error) = result {
tracing::warn!(
?id,
%method,
?error,
"Outgoing request readiness failed"
);
if let Some(pending_reply) = pending_replies.remove(&id) {
pending_reply.fail(error);
}
continue;
}
}
if !pending_replies.contains(&id) {
continue;
}
let request = match protocol_compat
.outgoing_message(untyped, remote_style)
.and_then(|untyped| remote_style.transform_outgoing_message(untyped))
.and_then(|untyped| untyped.into_raw_jsonrpc_message(Some(id.clone())))
{
Ok(request) => request,
Err(error) => {
tracing::warn!(?id, %method, ?error, "Failed to prepare outgoing request");
if let Some(pending_reply) = pending_replies.remove(&id) {
pending_reply.fail(error);
}
continue;
}
};
if !pending_replies.contains(&id) {
continue;
}
if let Err(error) = transport_tx.unbounded_send(TransportFrame::Single(request)) {
let error = crate::Error::into_internal_error(error);
if let Some(pending_reply) = pending_replies.remove(&id) {
pending_reply.fail(error.clone());
}
return Err(error);
}
continue;
}
OutgoingMessage::Notification { untyped } => {
let messages = match protocol_compat.outgoing_notification(untyped) {
Ok(messages) => messages,
Err(error) => {
tracing::warn!(
?error,
"Dropping outgoing notification after preparation failed"
);
continue;
}
};
for untyped in messages {
let message = match untyped.into_raw_jsonrpc_message(None) {
Ok(message) => message,
Err(error) => {
tracing::warn!(
?error,
"Dropping outgoing notification after serialization failed"
);
continue;
}
};
transport_tx
.unbounded_send(TransportFrame::Single(message))
.map_err(crate::Error::into_internal_error)?;
}
continue;
}
OutgoingMessage::Response {
id,
method,
response,
destination,
} => match protocol_compat.outgoing_response_to(&id, &method, response) {
Ok(value) => {
tracing::debug!(?id, "Sending success response");
(RawJsonRpcMessage::response(id, Ok(value)), destination)
}
Err(error) => {
tracing::warn!(?id, %method, ?error, "Sending error response");
(RawJsonRpcMessage::response(id, Err(error)), destination)
}
},
OutgoingMessage::UncorrelatedErrorResponse { error, destination } => {
// JSON-RPC reports parse/invalid-request errors with id null when
// they cannot be correlated to a specific request.
(
RawJsonRpcMessage::response(RequestId::Null, Err(error)),
destination,
)
}
};
if let Some(frame) = destination.complete(json_rpc_message) {
transport_tx
.unbounded_send(frame)
.map_err(crate::Error::into_internal_error)?;
}
}
// Closing the raw queue lets the transport actor finish all buffered
// writes. The caller separately awaits that transport future before
// treating the drain as complete.
drop(transport_tx);
for done in drain_waiters {
let _ = done.send(());
}
Ok(())
}