Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); Gb 7250 add sse client subscription test by hackal · Pull Request #10 · CallstackAI/grafbase-clone · GitHub
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 39 additions & 0 deletions .callstack.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
pr_review:
# Default: true
auto_run: true
modules:
# Automatically create a description summarizing the changes in pull request.
description:
enabled: true
diagram: false

# Find potential bugs in pull request changes or related files.
bug_hunter:
enabled: true
# Include fixes to possible bugs.
suggestions: true

# Suggest improvements to added code.
code_suggestions:
enabled: true

# Suggest changes to follow defined code conventions.
code_conventions:
enabled: false
# Describe your code conventions in plain text.
conventions: |
E.g. Exported variables, functions, classes and methods should be defined before private.


# Point out any typos or grammatical errors in variable names, texts, comments.
grammar:
enabled: false

# Suggest performance improvements to added code.
performance:
enabled: true

# Find potential security issues in added code.
security:
enabled: true

29 changes: 29 additions & 0 deletions .github/workflows/callstack-reviewer.yml
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
name: Callstack.ai PR Review

on:
workflow_dispatch:
inputs:
config:
type: string
description: "config for reviewer"
required: true
head:
type: string
description: "head commit sha"
required: true
base:
type: string
description: "base commit sha"
required: false

jobs:
callstack_pr_review_job:
runs-on: ubuntu-latest
steps:
- name: Review PR
uses: callstackai/action@main
with:
config: ${{ inputs.config }}
head: ${{ inputs.head }}
export: /code/chats.json

2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -71,6 +71,7 @@ async-graphql = "7.0.3"
async-graphql-axum = "7.0.3"
async-graphql-parser = "7.0.3"
async-graphql-value = "7.0.3"
async-sse = "5"
async-trait = "0.1.80"
axum = { version = "0.7.5", default-features = false }
axum-server = { version = "0.6", default-features = false }
Expand DownExpand Up@@ -102,6 +103,7 @@ internment = { version = "0.8", features = ["serde", "arc"] }
itertools = "0.13.0"
jsonwebtoken = "9.3.0"
governor = "0.6"
multipart-stream = "0.1.2"
num-traits = "0.2.18"
once_cell = "1.19.0"
openidconnect = "4.0.0-alpha.1"
Expand Down
6 changes: 3 additions & 3 deletions engine/crates/gateway-core/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,8 +19,8 @@ workspace = true
[dependencies]
async-graphql.workspace = true
async-runtime.workspace = true
async-sse = "5.1.0"
async-trait = "0.1.80"
async-sse.workspace = true
async-trait.workspace = true
blake3.workspace = true
bytes.workspace = true
common-types.workspace = true
Expand All@@ -36,7 +36,7 @@ headers.workspace = true
http.workspace = true
mediatype = "0.19.18"
mime = "0.3.17"
multipart-stream = "0.1.2"
multipart-stream.workspace = true
operation-normalizer = { path = "../operation-normalizer" }
partial-caching.workspace = true
registry-for-cache.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion engine/crates/integration-tests/Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,7 +12,9 @@ async-graphql-parser.workspace = true
async-graphql.workspace = true
async-once-cell = "0.5.3"
async-runtime.workspace = true
async-sse.workspace = true
async-trait.workspace = true
bytes.workspace = true
crossbeam-queue = "0.3"
cynic.workspace = true
cynic-introspection.workspace = true
Expand All@@ -32,7 +34,7 @@ headers.workspace = true
http.workspace = true
indoc = "2.0.5"
insta.workspace = true
multipart-stream = "0.1.2"
multipart-stream.workspace = true
names = "0.14.1-dev"
openidconnect.workspace = true
reqwest.workspace = true
Expand Down
17 changes: 12 additions & 5 deletions engine/crates/integration-tests/docker-compose.yml
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,16 @@
version: '3'
services:
sse-subgraph:
restart: unless-stopped
build:
context: ./data/sse-subgraph
ports:
- '4092:4092'

# MongoDB
data-api:
image: grafbase/mongodb-data-api:latest
restart: always
restart: unless-stopped
environment:
MONGODB_DATABASE_URL: 'mongodb://grafbase:grafbase@mongodb:27017'
ports:
Expand All@@ -15,7 +22,7 @@ services:

mongodb:
image: mongo:latest
restart: always
restart: unless-stopped
environment:
MONGO_INITDB_ROOT_USERNAME: 'grafbase'
MONGO_INITDB_ROOT_PASSWORD: 'grafbase'
Expand All@@ -28,7 +35,7 @@ services:
# Postgres
postgres:
image: postgres:16
restart: always
restart: unless-stopped
command: postgres -c 'max_connections=1000'
environment:
POSTGRES_PASSWORD: 'grafbase'
Expand DownExpand Up@@ -61,7 +68,7 @@ services:
environment:
DSN: 'sqlite:///var/lib/sqlite/db.sqlite?_fk=true'
URLS_SELF_ISSUER: 'http://127.0.0.1:4444'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand DownExpand Up@@ -104,7 +111,7 @@ services:
URLS_SELF_ISSUER: 'http://127.0.0.1:4454'
SERVE_PUBLIC_PORT: '4454'
SERVE_ADMIN_PORT: '4455'
restart: always
restart: unless-stopped
depends_on:
- hydra-migrate
networks:
Expand Down
192 changes: 3 additions & 189 deletions engine/crates/integration-tests/src/federation/mod.rs
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,11 @@
mod builder;
mod request;

use std::{
any::TypeId,
borrow::Cow,
collections::HashMap,
future::IntoFuture,
ops::{Deref, DerefMut},
str::FromStr,
sync::Arc,
};
use std::{any::TypeId, collections::HashMap, sync::Arc};

pub use builder::*;
use engine::{BatchRequest, Variables};
use engine_v2::{HttpGraphqlResponse, HttpGraphqlResponseBody};
use futures::{future::BoxFuture, stream::BoxStream, StreamExt, TryStreamExt};
use gateway_core::StreamingFormat;
use graphql_mocks::{MockGraphQlServer, ReceivedRequest};
use headers::HeaderMapExt;
use http::{header::Entry, HeaderName, HeaderValue};
use serde::de::Error;
pub use request::*;

use crate::engine_v1::GraphQlRequest;

Expand DownExpand Up@@ -69,176 +56,3 @@ impl TestEngineV2 {
.collect()
}
}

#[must_use]
pub struct ExecutionRequest {
request: GraphQlRequest,
#[allow(dead_code)]
headers: Vec<(String, String)>,
engine: Arc<engine_v2::Engine<TestRuntime>>,
}

impl ExecutionRequest {
pub fn by_client(self, name: &'static str, version: &'static str) -> Self {
self.header("x-grafbase-client-name", name)
.header("x-grafbase-client-version", version)
}

/// Adds a header into the request
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}

pub fn variables(mut self, variables: impl serde::Serialize) -> Self {
self.request.variables = Some(Variables::from_json(
serde_json::to_value(variables).expect("variables to be serializable"),
));
self
}

pub fn extensions(mut self, extensions: impl serde::Serialize) -> Self {
self.request.extensions =
serde_json::from_value(serde_json::to_value(extensions).expect("extensions to be serializable"))
.expect("extensions to be deserializable");
self
}

fn http_headers(&self) -> http::HeaderMap {
let mut headers = http::HeaderMap::new();

for (key, value) in &self.headers {
let key = HeaderName::from_str(key).unwrap();
let value = HeaderValue::from_str(value).unwrap();

if let Entry::Occupied(mut e) = headers.entry(key.clone()) {
e.append(value);
} else {
headers.insert(key, value);
}
}

headers
}

pub fn into_multipart_stream(self) -> MultipartStreamRequest {
MultipartStreamRequest(self)
}
}

impl IntoFuture for ExecutionRequest {
type Output = GraphqlResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let headers = self.http_headers();
let request = BatchRequest::Single(self.request.into_engine_request());
Box::pin(async move { self.engine.execute(headers, request).await.try_into().unwrap() })
}
}

pub struct MultipartStreamRequest(ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B
where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await
}
}

impl IntoFuture for MultipartStreamRequest {
type Output = GraphqlStreamingResponse;

type IntoFuture = BoxFuture<'static, Self::Output>;

fn into_future(self) -> Self::IntoFuture {
let mut headers = self.0.http_headers();
headers.typed_insert(StreamingFormat::IncrementalDelivery);
let request = BatchRequest::Single(self.0.request.into_engine_request());
Box::pin(async move {
let response = self.0.engine.execute(headers, request).await;
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into), "-")
.map(|result| serde_json::from_slice(&result.unwrap().body).unwrap());
GraphqlStreamingResponse {
stream: Box::pin(stream),
headers: response.headers,
}
})
}
}

pub struct GraphqlStreamingResponse {
pub stream: BoxStream<'static, serde_json::Value>,
pub headers: http::HeaderMap,
}

#[derive(serde::Serialize, Debug)]
pub struct GraphqlResponse {
#[serde(flatten)]
pub body: serde_json::Value,
#[serde(skip)]
pub headers: http::HeaderMap,
}

impl TryFrom<HttpGraphqlResponse> for GraphqlResponse {
type Error = serde_json::Error;

fn try_from(response: HttpGraphqlResponse) -> Result<Self, Self::Error> {
Ok(GraphqlResponse {
body: match response.body {
HttpGraphqlResponseBody::Bytes(bytes) => serde_json::from_slice(bytes.as_ref())?,
HttpGraphqlResponseBody::Stream(_) => {
return Err(serde_json::Error::custom("Unexpected stream response body"))?
}
},
headers: response.headers,
})
}
}

impl std::fmt::Display for GraphqlResponse {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", serde_json::to_string_pretty(&self.body).unwrap())
}
}

impl Deref for GraphqlResponse {
type Target = serde_json::Value;

fn deref(&self) -> &Self::Target {
&self.body
}
}

impl DerefMut for GraphqlResponse {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.body
}
}

impl GraphqlResponse {
pub fn into_value(self) -> serde_json::Value {
self.body
}

#[track_caller]
pub fn into_data(self) -> serde_json::Value {
assert!(self.errors().is_empty(), "{self:#?}");

match self.body {
serde_json::Value::Object(mut value) => value.remove("data"),
_ => None,
}
.unwrap_or_default()
}

pub fn errors(&self) -> Cow<'_, Vec<serde_json::Value>> {
self.body["errors"]
.as_array()
.map(Cow::Borrowed)
.unwrap_or_else(|| Cow::Owned(Vec::new()))
}
}
Loading