Skip to content
Merged
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
2 changes: 1 addition & 1 deletion app/graphql/subscription_triggers.rb
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ module SubscriptionTriggers
def self.execution_result(execution_result)
SagittariusSchema.subscriptions.trigger(
:namespaces_projects_flows_execution_result,
{ execution_identifier: execution_result.execution_identifier },
{ flow_id: execution_result.flow.to_global_id },
execution_result,
context: { visibility_profile: :execution }
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,28 +7,18 @@ module Flows
class ExecutionResult < BaseSubscription
description 'Subscription to asynchronously receive an execution result'

argument :execution_identifier,
type: GraphQL::Types::String,
argument :flow_id,
type: Types::GlobalIdType[Flow],
required: true,
description: 'Execution identifier of the triggered execution'
description: 'Id of the flow to receive execution results for'

field :execution_result,
type: Types::ExecutionResultType,
null: true,
description: 'The execution result of the relevant execution'

def subscribe(execution_identifier:)
result = ::ExecutionResult.find_by(execution_identifier: execution_identifier, created_at: 1.hour.ago..)

if result.present?
unsubscribe({ execution_result: result })
else
:no_response
end
end
description: 'The most recent execution result of the relevant flow'

def update(*)
unsubscribe({ execution_result: object })
{ execution_result: object }
end
end
end
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,10 @@ Subscription to asynchronously receive an execution result

| Name | Type | Description |
|------|------|-------------|
| `executionIdentifier` | [`String!`](../scalar/string.md) | Execution identifier of the triggered execution |
| `flowId` | [`FlowID!`](../scalar/flowid.md) | Id of the flow to receive execution results for |

## Fields

| Name | Type | Description |
|------|------|-------------|
| `executionResult` | [`ExecutionResult`](../object/executionresult.md) | The execution result of the relevant execution |
| `executionResult` | [`ExecutionResult`](../object/executionresult.md) | The most recent execution result of the relevant flow |
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,11 @@
let(:user) { create(:user) }
let(:token) { "Session #{authorization_token(user)}" }
let(:flow) { create(:flow) }
let(:execution_identifier) { 'existing-execution' }

let(:subscription_query) do
<<~GQL
subscription($executionIdentifier: String!) {
namespacesProjectsFlowsExecutionResult(executionIdentifier: $executionIdentifier) {
subscription($flowId: FlowID!) {
namespacesProjectsFlowsExecutionResult(flowId: $flowId) {
executionResult {
success
nodeResults {
Expand All @@ -42,29 +41,47 @@
subscribe(token: token)
end

context 'when the execution result already exists' do
before do
result = create(
:execution_result,
flow: flow,
execution_identifier: execution_identifier,
success: { 'done' => true }
)
node_result = create(:execution_node_result, execution_result: result)
create(:execution_parameter_result, execution_node_result: node_result)
context 'when subscribing' do
it 'does not deliver an execution result in the initial subscription response' do
perform :execute, query: subscription_query, variables: { flowId: flow.to_global_id.to_s }

execution_result = transmissions.last.dig('result', 'data', 'namespacesProjectsFlowsExecutionResult')
expect(execution_result).to be_nil
end
end

context 'when a new execution result is persisted for the flow after subscribing' do
it 'streams the result to the subscriber' do
perform :execute, query: subscription_query, variables: { flowId: flow.to_global_id.to_s }

result = create(:execution_result, flow: flow, success: { 'first' => true })
SubscriptionTriggers.execution_result(result)

it 'immediately delivers the result in the initial subscription response' do
perform :execute, query: subscription_query, variables: { executionIdentifier: execution_identifier }
first_transmission = transmissions.last
execution_result = first_transmission.dig('result', 'data', 'namespacesProjectsFlowsExecutionResult',
'executionResult')
expect(execution_result['success']).to eq({ 'first' => true })

other_result = create(:execution_result, flow: flow, success: { 'second' => true })
SubscriptionTriggers.execution_result(other_result)

second_transmission = transmissions.last
second_execution_result = second_transmission.dig('result', 'data', 'namespacesProjectsFlowsExecutionResult',
'executionResult')
expect(second_execution_result['success']).to eq({ 'second' => true })
end
end

result = transmissions.last
context 'when a result for a different flow is triggered' do
it 'does not deliver the result to the subscriber' do
perform :execute, query: subscription_query, variables: { flowId: flow.to_global_id.to_s }
transmission_count_before_trigger = transmissions.count

execution_result = result.dig('result', 'data', 'namespacesProjectsFlowsExecutionResult', 'executionResult')
expect(execution_result['success']).to eq({ 'done' => true })
other_flow = create(:flow)
result = create(:execution_result, flow: other_flow, success: { 'done' => true })
SubscriptionTriggers.execution_result(result)

execution_node_result = execution_result.dig('nodeResults', 'nodes', 0)
expect(execution_node_result['success']).to eq({ 'node' => 'ok' })
expect(execution_node_result.dig('parameterResults', 0, 'value')).to eq({ 'parameter' => 'ok' })
expect(transmissions.count).to eq(transmission_count_before_trigger)
end
end
end
Original file line number Diff line number Diff line change
Expand Up @@ -379,13 +379,13 @@

perform :execute,
query: <<~GQL,
subscription($executionIdentifier: String!) {
namespacesProjectsFlowsExecutionResult(executionIdentifier: $executionIdentifier) {
subscription($flowId: FlowID!) {
namespacesProjectsFlowsExecutionResult(flowId: $flowId) {
executionResult { success }
}
}
GQL
variables: { executionIdentifier: 'execution-identifier' }
variables: { flowId: flow.to_global_id.to_s }
end

it 'delivers the execution result to subscribers without visibility profile error' do
Expand Down