diff --git a/app/graphql/subscription_triggers.rb b/app/graphql/subscription_triggers.rb index 37c6f4af9..be63287b0 100644 --- a/app/graphql/subscription_triggers.rb +++ b/app/graphql/subscription_triggers.rb @@ -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 } ) diff --git a/app/graphql/subscriptions/namespaces/projects/flows/execution_result.rb b/app/graphql/subscriptions/namespaces/projects/flows/execution_result.rb index 251f75795..45f0d5344 100644 --- a/app/graphql/subscriptions/namespaces/projects/flows/execution_result.rb +++ b/app/graphql/subscriptions/namespaces/projects/flows/execution_result.rb @@ -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 diff --git a/docs/graphql/subscription/namespacesprojectsflowsexecutionresult.md b/docs/graphql/subscription/namespacesprojectsflowsexecutionresult.md index 31540501a..b5e6add3e 100644 --- a/docs/graphql/subscription/namespacesprojectsflowsexecutionresult.md +++ b/docs/graphql/subscription/namespacesprojectsflowsexecutionresult.md @@ -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 | diff --git a/spec/requests/graphql/subscription/namespaces/projects/flows/execution_result_spec.rb b/spec/requests/graphql/subscription/namespaces/projects/flows/execution_result_spec.rb index 55f1c7863..97238a889 100644 --- a/spec/requests/graphql/subscription/namespaces/projects/flows/execution_result_spec.rb +++ b/spec/requests/graphql/subscription/namespaces/projects/flows/execution_result_spec.rb @@ -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 { @@ -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 diff --git a/spec/services/namespaces/projects/flows/persist_execution_result_service_spec.rb b/spec/services/namespaces/projects/flows/persist_execution_result_service_spec.rb index f96339e0c..2986ca578 100644 --- a/spec/services/namespaces/projects/flows/persist_execution_result_service_spec.rb +++ b/spec/services/namespaces/projects/flows/persist_execution_result_service_spec.rb @@ -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