@@ -25,6 +25,7 @@ defmodule Cascade.DagLoader do
2525 require Logger
2626
2727 alias Cascade.Workflows
28+ alias Cascade.Events
2829 alias Cascade.DagLoader . { LocalSource , S3Source , Validator }
2930
3031 @ default_scan_interval 30_000 # 30 seconds
@@ -219,8 +220,14 @@ defmodule Cascade.DagLoader do
219220 :ok
220221
221222 dag ->
222- Workflows . update_dag ( dag , % { enabled: false } )
223- Logger . info ( " Disabled DAG: #{ dag_name } " )
223+ case Workflows . update_dag ( dag , % { enabled: false } ) do
224+ { :ok , _updated_dag } ->
225+ Logger . info ( " Disabled DAG: #{ dag_name } " )
226+ broadcast_dag_event ( :deleted , dag_name )
227+
228+ { :error , _reason } ->
229+ :ok
230+ end
224231 end
225232 end )
226233
@@ -287,14 +294,39 @@ defmodule Cascade.DagLoader do
287294 enabled: Map . get ( definition , "enabled" , true )
288295 }
289296
290- case Workflows . get_dag_by_name ( name ) do
291- nil ->
292- # Create new DAG
293- Workflows . create_dag ( attrs )
297+ result =
298+ case Workflows . get_dag_by_name ( name ) do
299+ nil ->
300+ # Create new DAG
301+ case Workflows . create_dag ( attrs ) do
302+ { :ok , _dag } = success ->
303+ broadcast_dag_event ( :created , name )
304+ success
305+
306+ error ->
307+ error
308+ end
309+
310+ existing_dag ->
311+ # Update existing DAG
312+ case Workflows . update_dag ( existing_dag , attrs ) do
313+ { :ok , _dag } = success ->
314+ broadcast_dag_event ( :updated , name )
315+ success
316+
317+ error ->
318+ error
319+ end
320+ end
294321
295- existing_dag ->
296- # Update existing DAG
297- Workflows . update_dag ( existing_dag , attrs )
298- end
322+ result
323+ end
324+
325+ defp broadcast_dag_event ( action , dag_name ) do
326+ Phoenix.PubSub . broadcast (
327+ Cascade.PubSub ,
328+ Events . dag_updates_topic ( ) ,
329+ { :dag_event , % { action: action , dag_name: dag_name } }
330+ )
299331 end
300332end
0 commit comments