forked from DrFaust92/terraform-provider-airflow
-
Notifications
You must be signed in to change notification settings - Fork 0
/
resource_dag.go
121 lines (106 loc) · 2.73 KB
/
resource_dag.go
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
package main
import (
"fmt"
"github.com/apache/airflow-client-go/airflow"
"github.com/hashicorp/terraform-plugin-sdk/v2/helper/schema"
)
func resourceDag() *schema.Resource {
return &schema.Resource{
Create: resourceDagUpdate,
Read: resourceDagRead,
Update: resourceDagUpdate,
Delete: resourceDagDelete,
Importer: &schema.ResourceImporter{
StateContext: schema.ImportStatePassthroughContext,
},
Schema: map[string]*schema.Schema{
"dag_id": {
Type: schema.TypeString,
Required: true,
ForceNew: true,
},
"description": {
Type: schema.TypeString,
Computed: true,
},
"delete_dag": {
Type: schema.TypeBool,
Optional: true,
Default: false,
},
"file_token": {
Type: schema.TypeString,
Computed: true,
},
"fileloc": {
Type: schema.TypeString,
Computed: true,
},
"is_active": {
Type: schema.TypeBool,
Computed: true,
},
"is_paused": {
Type: schema.TypeBool,
Required: true,
},
"is_subdag": {
Type: schema.TypeBool,
Computed: true,
},
"root_dag_id": {
Type: schema.TypeString,
Computed: true,
},
},
}
}
func resourceDagUpdate(d *schema.ResourceData, m interface{}) error {
pcfg := m.(ProviderConfig)
client := pcfg.ApiClient
dagId := d.Get("dag_id").(string)
dagApi := client.DAGApi
dag := *airflow.NewDAG()
dag.SetIsPaused(d.Get("is_paused").(bool))
_, res, err := dagApi.PatchDag(pcfg.AuthContext, dagId).DAG(dag).Execute()
if res.StatusCode != 200 {
return fmt.Errorf("failed to update DAG `%s` from Airflow: %w", dagId, err)
}
d.SetId(dagId)
return resourceDagRead(d, m)
}
func resourceDagRead(d *schema.ResourceData, m interface{}) error {
pcfg := m.(ProviderConfig)
client := pcfg.ApiClient
DAG, resp, err := client.DAGApi.GetDag(pcfg.AuthContext, d.Id()).Execute()
if resp != nil && resp.StatusCode == 404 {
d.SetId("")
return nil
}
if resp.StatusCode != 200 {
return fmt.Errorf("failed to get DAG `%s` from Airflow: %w", d.Id(), err)
}
d.Set("dag_id", DAG.DagId)
d.Set("is_paused", DAG.IsPaused.Get())
d.Set("is_active", DAG.IsActive.Get())
d.Set("is_subdag", DAG.IsSubdag)
d.Set("description", DAG.Description.Get())
d.Set("file_token", DAG.FileToken)
d.Set("fileloc", DAG.Fileloc)
d.Set("root_dag_id", DAG.RootDagId.Get())
return nil
}
func resourceDagDelete(d *schema.ResourceData, m interface{}) error {
pcfg := m.(ProviderConfig)
client := pcfg.ApiClient.DAGApi
if d.Get("delete_dag").(bool) {
resp, err := client.DeleteDag(pcfg.AuthContext, d.Id()).Execute()
if err != nil {
return fmt.Errorf("failed to delete DAG `%s` from Airflow: %w", d.Id(), err)
}
if resp != nil && resp.StatusCode == 404 {
return nil
}
}
return nil
}