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
use std::{borrow::Cow, pin::Pin};
use futures_util::stream::{Stream, StreamExt};
use crate::{
parser::types::{Selection, TypeCondition},
registry,
registry::Registry,
Context, ContextSelectionSet, PathSegment, Response, ServerError, ServerResult,
};
pub trait SubscriptionType: Send + Sync {
fn type_name() -> Cow<'static, str>;
fn qualified_type_name() -> String {
format!("{}!", Self::type_name())
}
fn create_type_info(registry: &mut registry::Registry) -> String;
#[doc(hidden)]
fn is_empty() -> bool {
false
}
#[doc(hidden)]
fn create_field_stream<'a>(
&'a self,
ctx: &'a Context<'_>,
) -> Option<Pin<Box<dyn Stream<Item = Response> + Send + 'a>>>;
}
type BoxFieldStream<'a> = Pin<Box<dyn Stream<Item = Response> + 'a + Send>>;
pub(crate) fn collect_subscription_streams<'a, T: SubscriptionType + 'static>(
ctx: &ContextSelectionSet<'a>,
root: &'a T,
streams: &mut Vec<BoxFieldStream<'a>>,
) -> ServerResult<()> {
for selection in &ctx.item.node.items {
match &selection.node {
Selection::Field(field) => streams.push(Box::pin({
let ctx = ctx.clone();
async_stream::stream! {
let ctx = ctx.with_field(field);
let field_name = ctx.item.node.response_key().node.clone();
let stream = root.create_field_stream(&ctx);
if let Some(mut stream) = stream {
while let Some(resp) = stream.next().await {
yield resp;
}
} else {
let err = ServerError::new(format!(r#"Cannot query field "{}" on type "{}"."#, field_name, T::type_name()), Some(ctx.item.pos))
.with_path(vec![PathSegment::Field(field_name.to_string())]);
yield Response::from_errors(vec![err]);
}
}
})),
Selection::FragmentSpread(fragment_spread) => {
if let Some(fragment) = ctx
.query_env
.fragments
.get(&fragment_spread.node.fragment_name.node)
{
collect_subscription_streams(
&ctx.with_selection_set(&fragment.node.selection_set),
root,
streams,
)?;
}
}
Selection::InlineFragment(inline_fragment) => {
if let Some(TypeCondition { on: name }) = inline_fragment
.node
.type_condition
.as_ref()
.map(|v| &v.node)
{
if name.node.as_str() == T::type_name() {
collect_subscription_streams(
&ctx.with_selection_set(&inline_fragment.node.selection_set),
root,
streams,
)?;
}
} else {
collect_subscription_streams(
&ctx.with_selection_set(&inline_fragment.node.selection_set),
root,
streams,
)?;
}
}
}
}
Ok(())
}
impl<T: SubscriptionType> SubscriptionType for &T {
fn type_name() -> Cow<'static, str> {
T::type_name()
}
fn create_type_info(registry: &mut Registry) -> String {
T::create_type_info(registry)
}
fn create_field_stream<'a>(
&'a self,
ctx: &'a Context<'_>,
) -> Option<Pin<Box<dyn Stream<Item = Response> + Send + 'a>>> {
T::create_field_stream(*self, ctx)
}
}