Skip to content

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere
, '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

Gb 7250 add sse client subscription test - #10

Open
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test
Open

Gb 7250 add sse client subscription test#10
hackal wants to merge 3 commits into
mainfrom
gb-7250-add-sse-client-subscription-test

Conversation

@hackal

@hackalhackal commented Aug 13, 2024

Copy link
Copy Markdown
Contributor

Description by Cal

PR Description

This PR introduces a new feature for SSE client subscription testing. It includes configuration changes, dependency updates, and new test cases for SSE and multipart streaming in a federated GraphQL environment.

Key Issues

None

Files Changed

File: /.callstack.yml Added configuration for Callstack.ai PR review with modules for description, bug hunting, code suggestions, performance, and security.
File: /.github/workflows/callstack-reviewer.yml Created a GitHub Actions workflow for Callstack.ai PR review with inputs for configuration and commit SHAs.
File: /Cargo.lock Updated dependencies to include 'async-sse' and 'bytes'.
File: /Cargo.toml Added 'async-sse' and 'multipart-stream' as dependencies.
File: /engine/crates/gateway-core/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'async-trait', and 'multipart-stream'.
File: /engine/crates/integration-tests/Cargo.toml Updated dependencies to use workspace versions for 'async-sse', 'bytes', and 'multipart-stream'.
File: /engine/crates/integration-tests/docker-compose.yml Modified Docker Compose file to add a new service 'sse-subgraph' and changed restart policies for existing services.
File: /engine/crates/integration-tests/src/federation/mod.rs Refactored federation module to move request handling to a separate file.
File: /engine/crates/integration-tests/src/federation/request.rs Introduced a new module for handling GraphQL request execution and response processing.
File: /engine/crates/integration-tests/src/federation/request/stream.rs Added new module for handling multipart and SSE streaming requests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/mod.rs Created a new module for subscription tests, including multipart and SSE tests.
File: /engine/crates/integration-tests/tests/federation/subscriptions/sse.rs Added test cases for SSE subscriptions in a federated GraphQL environment.
# Description

Please include a summary of the change and which issue is fixed.
Please also include relevant motivation and context.

Type of change

  • 💔 Breaking
  • 🚀 Feature
  • 🐛 Fix
  • 🛠️ Tooling
  • 🔨 Refactoring
  • 🧪 Test
  • 📦 Dependency
  • 📖 Requires documentation update

Finistereand others added 3 commits August 6, 2024 22:24
We were already testing this in the cli tests, so not a huge win but
I mostly run integration-tests so well.

@callstackai-actioncallstackai-actionBot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

impl MultipartStreamRequest {
pub async fn collect<B>(self) -> B

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The collect method in both MultipartStreamRequest and SseStreamRequest is identical. Consider refactoring this method into a shared utility function to avoid code duplication.

use gateway_core::StreamingFormat;
use headers::HeaderMapExt;

pub struct MultipartStreamRequest(pub(super) super::ExecutionRequest);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider using a more descriptive name for the MultipartStreamRequest and SseStreamRequest structs to clearly indicate their purpose and usage.

where
B: Default + Extend<serde_json::Value>,
{
self.await.stream.collect().await

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
In both the MultipartStreamRequest::collect and SseStreamRequest::collect methods, the line self.await.stream.collect().await is incorrect because self is not a future and cannot be awaited. The correct approach would be to access the stream field of the ExecutionRequest and then collect it.

Suggested change
self.await.stream.collect().await
self.0.stream.collect().await

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), "-")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🐛 Bug
The use of unwrap in serde_json::from_slice(&result.unwrap().body).unwrap() and serde_json::from_slice(msg.data()).unwrap() can lead to panics if the result is an error or if the data is not valid JSON. It is important to handle these cases more gracefully to improve robustness and prevent potential panics. Consider handling the Err case properly to ensure the application can handle unexpected data without crashing.

Suggestion:

Suggested change
let stream = multipart_stream::parse(response.body.into_stream().map_ok(Into::into),"-")
Ok(async_sse::Event::Message(msg)) => serde_json::from_slice(msg.data()).unwrap_or_else(|_| serde_json::Value::String("Invalid JSON".into()))

@devslovecoffeedevslovecoffee left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PR Review Summary

This pull request has been reviewed. Please check the comments and suggestions provided.

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());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The unwrap() calls in the MultipartStreamRequest and SseStreamRequest implementations can lead to panics if the Result is an Err. Consider handling the error more gracefully, perhaps by using ? to propagate the error or by providing a default value.

let stream = async_sse::decode(stream.into_async_read())
.into_stream()
.try_take_while(|event| {
let take = if let async_sse::Event::Message(msg) = event {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
The use of string literals like "complete" and "Got retry?" in the SseStreamRequest implementation can be considered as excessive use of literals. Consider defining these as constants to improve maintainability and readability.

let mut headers = http::HeaderMap::new();

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion
Consider handling potential errors when using unwrap() to avoid panics. Specifically, unwrap() is used on HeaderName::from_str, HeaderValue::from_str, Box::pin, and serde_json::to_string_pretty. In each of these cases, if the operation fails, it could lead to a panic. Handling these errors more gracefully would improve the robustness of the code.

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@hackal@devslovecoffee@Finistere