Skip to content

Commit 3072fe8

Browse files
feat: add show detailed trigger history (wraft#176)
* feat: add show detailed trigger history * feat: add pipeline_id and trigger_Id to form entry response * chore: update swagger doc * chore: update dependencies in mix.exs and mix.lock --------- Co-authored-by: Salsabeel <[email protected]>
1 parent f176c63 commit 3072fe8

13 files changed

Lines changed: 240 additions & 61 deletions

File tree

lib/wraft_doc/forms/forms.ex

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -543,8 +543,8 @@ defmodule WraftDoc.Forms do
543543
end)
544544
|> Repo.transaction()
545545
|> case do
546-
{:ok, _} ->
547-
:ok
546+
{:ok, %{trigger_history: trigger_history}} ->
547+
{:ok, %{pipeline_id: pipeline_id, trigger_id: trigger_history.id}}
548548

549549
{:error, step, error, _} ->
550550
Logger.error("Pipeline failed in step #{inspect(step)}", error: error)

lib/wraft_doc/pipelines/trigger_histories/trigger_histories.ex

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -123,4 +123,19 @@ defmodule WraftDoc.Pipelines.TriggerHistories do
123123
end
124124

125125
def create_pipeline_job(_, _), do: nil
126+
127+
@doc """
128+
Get a trigger history by id.
129+
"""
130+
@spec get_trigger_history(User.t(), Ecto.UUID.t()) :: TriggerHistory.t() | nil
131+
def get_trigger_history(%User{current_org_id: org_id}, <<_::288>> = id) do
132+
TriggerHistory
133+
|> join(:inner, [t], p in Pipeline, on: t.pipeline_id == p.id)
134+
|> where([t, p], p.organisation_id == ^org_id)
135+
|> where([t], t.id == ^id)
136+
|> preload([:creator])
137+
|> Repo.one()
138+
end
139+
140+
def get_trigger_history(_, _), do: nil
126141
end

lib/wraft_doc/pipelines/trigger_histories/trigger_history.ex

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ defmodule WraftDoc.Pipelines.TriggerHistories.TriggerHistory do
1111
only: [
1212
:id,
1313
:data,
14+
:response,
1415
:error,
1516
:state,
1617
:pipeline_id,
@@ -24,6 +25,7 @@ defmodule WraftDoc.Pipelines.TriggerHistories.TriggerHistory do
2425

2526
schema "trigger_history" do
2627
field(:data, :map)
28+
field(:response, :map)
2729
field(:error, :map, default: %{})
2830
field(:state, :integer)
2931
field(:start_time, :naive_datetime)
@@ -73,8 +75,8 @@ defmodule WraftDoc.Pipelines.TriggerHistories.TriggerHistory do
7375

7476
def update_changeset(%TriggerHistory{} = trigger, attrs \\ %{}) do
7577
trigger
76-
|> cast(attrs, [:error, :state, :start_time, :zip_file])
77-
|> validate_required([:state, :error])
78+
|> cast(attrs, [:response, :error, :state, :start_time, :zip_file])
79+
|> validate_required([:state])
7880
end
7981

8082
def trigger_end_changeset(%TriggerHistory{} = trigger, attrs \\ %{}) do

lib/wraft_doc/webhooks/event_trigger.ex

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -145,7 +145,7 @@ defmodule WraftDoc.Webhooks.EventTrigger do
145145
@doc """
146146
Trigger pipeline.completed event when a pipeline execution completes successfully.
147147
"""
148-
@spec trigger_pipeline_completed(TriggerHistory.t(), map()) :: :ok
148+
@spec trigger_pipeline_completed(TriggerHistory.t(), map()) :: map() | :ok
149149
def trigger_pipeline_completed(
150150
%TriggerHistory{pipeline: %{organisation_id: org_id}} = trigger_history,
151151
pipeline_result
@@ -156,6 +156,8 @@ defmodule WraftDoc.Webhooks.EventTrigger do
156156
Logger.info(
157157
"Triggered pipeline.completed webhook for pipeline id #{trigger_history.pipeline_id}"
158158
)
159+
160+
payload.pipeline
159161
end
160162

161163
def trigger_pipeline_completed(%TriggerHistory{} = trigger_history, _pipeline_result) do
@@ -169,7 +171,7 @@ defmodule WraftDoc.Webhooks.EventTrigger do
169171
@doc """
170172
Trigger pipeline.failed event when a pipeline execution fails.
171173
"""
172-
@spec trigger_pipeline_failed(TriggerHistory.t(), map()) :: :ok
174+
@spec trigger_pipeline_failed(TriggerHistory.t(), map()) :: map() | :ok
173175
def trigger_pipeline_failed(
174176
%TriggerHistory{pipeline: %{organisation_id: org_id}} = trigger_history,
175177
error_data
@@ -180,6 +182,8 @@ defmodule WraftDoc.Webhooks.EventTrigger do
180182
Logger.info(
181183
"Triggered pipeline.failed webhook for pipeline id #{trigger_history.pipeline_id}"
182184
)
185+
186+
payload.pipeline
183187
end
184188

185189
def trigger_pipeline_failed(%TriggerHistory{} = trigger_history, _error_data) do
@@ -193,7 +197,7 @@ defmodule WraftDoc.Webhooks.EventTrigger do
193197
@doc """
194198
Trigger pipeline.partially_completed event when a pipeline execution partially completes.
195199
"""
196-
@spec trigger_pipeline_partially_completed(TriggerHistory.t(), map()) :: :ok
200+
@spec trigger_pipeline_partially_completed(TriggerHistory.t(), map()) :: map() | :ok
197201
def trigger_pipeline_partially_completed(
198202
%TriggerHistory{pipeline: %{organisation_id: org_id}} = trigger_history,
199203
pipeline_result
@@ -204,6 +208,8 @@ defmodule WraftDoc.Webhooks.EventTrigger do
204208
Logger.info(
205209
"Triggered pipeline.partially_completed webhook for pipeline id #{trigger_history.pipeline_id}"
206210
)
211+
212+
payload.pipeline
207213
end
208214

209215
def trigger_pipeline_partially_completed(%TriggerHistory{} = trigger_history, _pipeline_result) do

lib/wraft_doc/workers/bulk_worker.ex

Lines changed: 92 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -117,13 +117,22 @@ defmodule WraftDoc.Workers.BulkWorker do
117117
stage: stage
118118
}
119119

120-
trigger =
121-
update_trigger_history_state_and_error(trigger, state, error_data)
122-
123120
# Trigger webhook for pipeline failure
124121
Task.start(fn ->
125122
trigger = Repo.preload(trigger, :pipeline)
126-
EventTrigger.trigger_pipeline_failed(trigger, error_data)
123+
124+
payload =
125+
trigger
126+
|> EventTrigger.trigger_pipeline_failed(error_data)
127+
|> case do
128+
payload when is_map(payload) ->
129+
payload
130+
131+
_ ->
132+
%{}
133+
end
134+
135+
update_trigger_history_state_and_error(trigger, state, error_data, payload)
127136
end)
128137

129138
Logger.error("Form mapping not complete. Pipeline execution failed.")
@@ -143,13 +152,22 @@ defmodule WraftDoc.Workers.BulkWorker do
143152
stage: stage
144153
}
145154

146-
trigger =
147-
update_trigger_history_state_and_error(trigger, state, error_data)
148-
149155
# Trigger webhook for pipeline failure
150156
Task.start(fn ->
151157
trigger = Repo.preload(trigger, :pipeline)
152-
EventTrigger.trigger_pipeline_failed(trigger, error_data)
158+
159+
payload =
160+
trigger
161+
|> EventTrigger.trigger_pipeline_failed(error_data)
162+
|> case do
163+
payload when is_map(payload) ->
164+
payload
165+
166+
_ ->
167+
%{}
168+
end
169+
170+
update_trigger_history_state_and_error(trigger, state, error_data, payload)
153171
end)
154172

155173
Logger.error("Pipeline not found. Pipeline execution failed.")
@@ -169,13 +187,21 @@ defmodule WraftDoc.Workers.BulkWorker do
169187
stage: stage
170188
}
171189

172-
trigger =
173-
update_trigger_history_state_and_error(trigger, state, error_data)
174-
175190
# Trigger webhook for pipeline failure
176191
Task.start(fn ->
177-
trigger = Repo.preload(trigger, :pipeline)
178-
EventTrigger.trigger_pipeline_failed(trigger, error_data)
192+
payload =
193+
trigger
194+
|> Repo.preload(:pipeline)
195+
|> EventTrigger.trigger_pipeline_failed(error_data)
196+
|> case do
197+
payload when is_map(payload) ->
198+
payload
199+
200+
_ ->
201+
%{}
202+
end
203+
204+
update_trigger_history_state_and_error(trigger, state, error_data, payload)
179205
end)
180206

181207
Logger.error("Instance creation failed. Pipeline execution failed.")
@@ -195,13 +221,22 @@ defmodule WraftDoc.Workers.BulkWorker do
195221
stage: stage
196222
}
197223

198-
trigger =
199-
update_trigger_history_state_and_error(trigger, state, error_data)
200-
201224
# Trigger webhook for pipeline failure
202225
Task.start(fn ->
203226
trigger = Repo.preload(trigger, :pipeline)
204-
EventTrigger.trigger_pipeline_failed(trigger, error_data)
227+
228+
payload =
229+
trigger
230+
|> EventTrigger.trigger_pipeline_failed(error_data)
231+
|> case do
232+
payload when is_map(payload) ->
233+
payload
234+
235+
_ ->
236+
%{}
237+
end
238+
239+
update_trigger_history_state_and_error(trigger, state, error_data, payload)
205240
end)
206241

207242
Logger.error("Instance creation failed. Pipeline execution failed.")
@@ -221,12 +256,21 @@ defmodule WraftDoc.Workers.BulkWorker do
221256
stage: stage
222257
}
223258

224-
trigger =
225-
update_trigger_history_state_and_error(trigger, state, error_data)
226-
227259
Task.start(fn ->
228260
trigger = Repo.preload(trigger, :pipeline)
229-
EventTrigger.trigger_pipeline_failed(trigger, error_data)
261+
262+
payload =
263+
trigger
264+
|> EventTrigger.trigger_pipeline_failed(error_data)
265+
|> case do
266+
payload when is_map(payload) ->
267+
payload
268+
269+
_ ->
270+
%{}
271+
end
272+
273+
update_trigger_history_state_and_error(trigger, state, error_data, payload)
230274
end)
231275

232276
Logger.error("Invalid JSON error. Pipeline execution failed.")
@@ -238,7 +282,6 @@ defmodule WraftDoc.Workers.BulkWorker do
238282
current_user
239283
) do
240284
state = TriggerHistory.states()[:success]
241-
trigger = update_trigger_history(trigger, %{state: state, zip_file: zip_file})
242285

243286
instances = Map.get(result, :instances, [])
244287

@@ -249,7 +292,20 @@ defmodule WraftDoc.Workers.BulkWorker do
249292

250293
Task.start(fn ->
251294
trigger = Repo.preload(trigger, :pipeline)
252-
EventTrigger.trigger_pipeline_completed(trigger, pipeline_result)
295+
296+
payload =
297+
trigger
298+
|> EventTrigger.trigger_pipeline_completed(pipeline_result)
299+
|> case do
300+
payload when is_map(payload) -> payload
301+
_ -> %{}
302+
end
303+
304+
update_trigger_history(trigger, %{
305+
state: state,
306+
response: payload,
307+
zip_file: zip_file
308+
})
253309
end)
254310

255311
Task.start(fn ->
@@ -310,11 +366,18 @@ defmodule WraftDoc.Workers.BulkWorker do
310366
documents: successful_instances
311367
}
312368

313-
trigger = update_trigger_history_state_and_error(trigger, state, error_data_for_db)
314-
315369
Task.start(fn ->
316370
trigger_with_pipeline = Repo.preload(trigger, :pipeline)
317-
EventTrigger.trigger_pipeline_partially_completed(trigger_with_pipeline, pipeline_result)
371+
372+
payload =
373+
trigger_with_pipeline
374+
|> EventTrigger.trigger_pipeline_partially_completed(pipeline_result)
375+
|> case do
376+
payload when is_map(payload) -> payload
377+
_ -> %{}
378+
end
379+
380+
update_trigger_history_state_and_error(trigger, state, error_data_for_db, payload)
318381
end)
319382

320383
Task.start(fn ->
@@ -353,12 +416,12 @@ defmodule WraftDoc.Workers.BulkWorker do
353416
end
354417

355418
# Update state and error of a trigger history
356-
@spec update_trigger_history_state_and_error(TriggerHistory.t(), integer, map) ::
419+
@spec update_trigger_history_state_and_error(TriggerHistory.t(), integer(), map(), map()) ::
357420
TriggerHistory.t()
358-
defp update_trigger_history_state_and_error(trigger, state, error) do
421+
defp update_trigger_history_state_and_error(trigger, state, error, response \\ %{}) do
359422
failure_time = DateTime.to_iso8601(Timex.now())
360423
error = Map.put(error, :failure_time, failure_time)
361-
update_trigger_history(trigger, %{state: state, error: error})
424+
update_trigger_history(trigger, %{state: state, error: error, response: response})
362425
end
363426

364427
# Update trigger history on start

lib/wraft_doc_web/controllers/form_entry_controller.ex

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,8 @@ defmodule WraftDocWeb.Api.V1.FormEntryController do
7373
properties do
7474
form_id(:string, "Form ID")
7575
id(:string, "Form Entry ID")
76+
trigger_id(:string, "Trigger ID")
77+
pipeline_id(:string, "Pipeline ID")
7678
inserted_at(:string, "When was the form entry inserted", format: "ISO-8601")
7779
status(:string, "Status of the form entry")
7880
updated_at(:string, "When was the form entry last updated", format: "ISO-8601")
@@ -87,6 +89,8 @@ defmodule WraftDocWeb.Api.V1.FormEntryController do
8789
},
8890
form_id: "aa18afe1-3383-4653-bc0e-505ec3bbfc19",
8991
id: "f507ca98-9848-49af-89f8-a21f12202ec0",
92+
pipeline_id: "12345678-9abc-def0-1234-56789abcdef0",
93+
trigger_id: "af2cf1c6-f342-4042-8425-6346e9fd6c44",
9094
inserted_at: "2024-04-17T07:10:17",
9195
status: "draft",
9296
updated_at: "2024-04-17T07:10:17",
@@ -194,8 +198,11 @@ defmodule WraftDocWeb.Api.V1.FormEntryController do
194198
%FormPipeline{} <- Forms.get_form_pipeline(form, pipeline_id),
195199
{:ok, %FormEntry{data: data} = form_entry} <-
196200
Forms.create_form_entry(current_user, form, params),
197-
:ok <- Forms.trigger_pipeline(current_user, pipeline_id, data, 0) do
198-
render(conn, "form_entry.json", %{form_entry: form_entry})
201+
{:ok, trigger_response} <- Forms.trigger_pipeline(current_user, pipeline_id, data, 0) do
202+
render(conn, "form_entry.json", %{
203+
form_entry: form_entry,
204+
trigger_response: trigger_response
205+
})
199206
end
200207
end
201208

0 commit comments

Comments
 (0)