-
Notifications
You must be signed in to change notification settings - Fork 95
Expand file tree
/
Copy pathconnectors_controller.rb
More file actions
215 lines (193 loc) 路 6.49 KB
/
Copy pathconnectors_controller.rb
File metadata and controls
215 lines (193 loc) 路 6.49 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
# frozen_string_literal: true
module Api
module V1
# rubocop:disable Metrics/ClassLength
class ConnectorsController < ApplicationController
include Connectors
include AuditLogger
include ResourceLinkBuilder
before_action :set_connector, only: %i[show update destroy discover query_source execute_model]
# TODO: Enable this once we have query validation implemented for all the connectors
# before_action :validate_query, only: %i[query_source]
# TODO: Enable this for ai_ml sources
before_action :validate_catalog, only: %i[query_source execute_model]
after_action :event_logger
after_action :create_audit_log, only: %i[create update destroy]
def index
@connectors = current_workspace.connectors
authorize @connectors
@connectors = @connectors.send(params[:type].downcase) if params[:type]
@connectors = @connectors.send(params[:category].downcase) if params[:category]
@connectors = @connectors.send(params[:sub_category].downcase) if params[:sub_category]
@connectors = @connectors.page(params[:page] || 1).per(params[:per_page])
render json: @connectors, status: :ok
end
def show
authorize @connector
render json: @connector, status: :ok
end
def create
authorize current_workspace, policy_class: ConnectorPolicy
result = CreateConnector.call(
workspace: current_workspace,
connector_params:
)
if result.success?
@connector = result.connector
@audit_resource = @connector.name
@resource_id = @connector.id
@payload = connector_params
render json: @connector, status: :created
else
render_error(
message: result.error || "Connector creation failed",
status: :unprocessable_content,
details: result.connector ? format_errors(result.connector) : nil
)
end
end
def update
authorize @connector
result = UpdateConnector.call(
connector: @connector,
connector_params:
)
if result.success?
@connector = result.connector
@audit_resource = @connector.name
@payload = connector_params
render json: @connector, status: :ok
else
render_error(
message: "Connector update failed",
status: :unprocessable_content,
details: format_errors(result.connector)
)
end
end
def destroy
authorize @connector
@action = "delete"
@audit_resource = @connector.name
@connector.destroy!
head :no_content
end
def discover
authorize @connector
result = DiscoverConnector.call(
connector: @connector,
refresh: params[:refresh]
)
if result.success?
@catalog = result.catalog
render json: @catalog, status: :ok
else
render_error(
message: "Discover catalog failed",
status: :unprocessable_content,
details: format_errors(result.catalog)
)
end
end
def query_source
authorize @connector
if @connector.source?
result = QuerySource.call(
connector: @connector,
query: params[:query],
limit: params[:limit] || 50
)
if result.success?
@records = result.records.map(&:record).map(&:data)
render json: { data: @records }, status: :ok
else
render_error(
message: result["error"],
status: :unprocessable_content
)
end
else
render_error(
message: "Connector is not a source",
status: :unprocessable_content
)
end
end
def execute_model
authorize @connector
if @connector.source?
result = ExecuteModel.call(
connector: @connector,
payload: params[:payload]
)
if result.success?
@records = result.records.map(&:record).map(&:data)
render json: { data: @records }, status: :ok
else
render_error(
message: result["error"],
status: :unprocessable_content
)
end
else
render_error(
message: "Connector is not a source",
status: :unprocessable_content
)
end
end
private
def set_connector
@connector = current_workspace.connectors.find(params[:id])
rescue ActiveRecord::RecordNotFound
render_error(
message: "Connector not found",
status: :not_found
)
end
def validate_catalog
return unless @connector.ai_model?
return if @connector.catalog.present?
render_error(
message: "Catalog is not present for the connector",
status: :unprocessable_content
)
end
def validate_query
Utils::QueryValidator.validate_query(@connector.connector_query_type, params[:query])
rescue StandardError => e
render_error(
message: "Query validation failed: #{e.message}",
status: :unprocessable_content
)
end
def create_audit_log
resource_id = @resource_id || params[:id]
resource_link = @action == "delete" ? nil : build_link!(resource: @connector, resource_id:)
audit!(action: @action, resource_id:, resource: @audit_resource, payload: @payload, resource_link:)
end
def connector_params
params.require(:connector).permit(:workspace_id,
:connector_type,
:connector_name, :name, :description, :query_type,
configuration: {})
end
<<<<<<< HEAD
=======
def filter_params
{
workspace: current_workspace,
type: params[:type],
category: params[:category],
sub_category: params[:sub_category],
exclude_sub_category: params[:exclude_sub_category],
provider: params[:provider],
page: params[:page],
per_page: params[:per_page]
}
end
>>>>>>> f6d71341a (fix(CE): added an extra filter in connectors api for excluding on basis of sub category (#1893))
end
# rubocop:enable Metrics/ClassLength
end
end