Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 42 additions & 31 deletions settings/settings.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,11 +47,12 @@ var (
const (
airflowConnectionList = "airflow connections list -o yaml"
ariflowPoolsList = "airflow pools list -o yaml"
airflowConnExport = "airflow connections export tmp.connections --file-format env"
airflowConnExport = "airflow connections export /tmp/connections.env --file-format env"
airflowVarExport = "airflow variables export tmp.var"
catVarFile = "cat tmp.var"
rmVarFile = "rm tmp.var"
catConnFile = "cat tmp.connections"
catConnFile = "cat /tmp/connections.env"
rmConnFile = "rm /tmp/connections.env"
configReadErrorMsg = "Error reading Airflow Settings file. Connections, Variables, and Pools were not loaded please check your Settings file syntax: %s\n"
noColorString = "[\u001B\u009B][[\\]()#;?]*(?:(?:(?:[a-zA-Z\\d]*(?:;[a-zA-Z\\d]*)*)?\u0007)|(?:(?:\\d{1,4}(?:;\\d{0,4})*)?[\\dA-PRZcf-ntqry=><~]))"
)
Expand Down Expand Up @@ -481,49 +482,59 @@ func EnvExportVariables(id, envFile string) error {
return errors.New("variable export unsuccessful")
}

func EnvExportConnections(id, envFile string) error {
func EnvExportConnections(id, envFile string) (err error) {
// Airflow command to export connections to env uris
out, err := execAirflowCommand(id, airflowConnExport)
if err != nil {
return fmt.Errorf("error exporting connections: %w", err)
}
logger.Debugf("Env Export Connections logs:\n%s", out)

if strings.Contains(out, "successfully") {
// get connections from file craeted by airflow command
out, err = execAirflowCommand(id, catConnFile)
if err != nil {
return fmt.Errorf("error reading connections file: %w", err)
}
if !strings.Contains(out, "successfully") {
return errors.New("connection export unsuccessful")
}

vars := strings.Split(out, "\n")
// add connections to the env file; connection URIs contain passwords, so keep it owner-only
f, err := os.OpenFile(envFile, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o600) //nolint:mnd
if err != nil {
return errors.Wrap(err, "Writing connections to file unsuccessful")
// The command above wrote connection URIs, including plaintext passwords, to a
// temp file inside the container. Remove it once we're done reading it, even if
// something below fails, so secrets don't linger in the container.
defer func() {
if _, rmErr := execAirflowCommand(id, rmConnFile); rmErr != nil {
rmErr = fmt.Errorf("error removing connections file: %w", rmErr)
if err == nil {
err = rmErr
} else {
fmt.Println(rmErr)
}
}
}()

defer f.Close()
// get connections from file created by airflow command
out, err = execAirflowCommand(id, catConnFile)
if err != nil {
return fmt.Errorf("error reading connections file: %w", err)
}

for i := range vars {
varSplit := strings.SplitN(vars[i], "=", 2) //nolint:mnd
if len(varSplit) > 1 {
fmt.Println("Exporting Connection: " + varSplit[0])
_, err := f.WriteString("\nAIRFLOW_CONN_" + strings.ToUpper(varSplit[0]) + "=" + varSplit[1])
if err != nil {
fmt.Printf("error adding connection %s to file: %s\n", varSplit[0], err.Error())
}
vars := strings.Split(out, "\n")
// add connections to the env file; connection URIs contain passwords, so keep it owner-only
f, ferr := os.OpenFile(envFile, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o600) //nolint:mnd
if ferr != nil {
return errors.Wrap(ferr, "Writing connections to file unsuccessful")
}

defer f.Close()

for i := range vars {
varSplit := strings.SplitN(vars[i], "=", 2) //nolint:mnd
if len(varSplit) > 1 {
fmt.Println("Exporting Connection: " + varSplit[0])
_, err := f.WriteString("\nAIRFLOW_CONN_" + strings.ToUpper(varSplit[0]) + "=" + varSplit[1])
if err != nil {
fmt.Printf("error adding connection %s to file: %s\n", varSplit[0], err.Error())
}
}
fmt.Println("Aiflow connections successfully export to the file " + envFile + "\n")
rmCmd := "rm tmp.connection"
_, err = execAirflowCommand(id, rmCmd)
if err != nil {
return fmt.Errorf("error removing connections file: %w", err)
}
return nil
}
return errors.New("connection export unsuccessful")
fmt.Println("Airflow connections successfully exported to the file " + envFile + "\n")
return nil
}

func Export(id, settingsFile string, version uint64, connections, variables, pools bool) error {
Expand Down
49 changes: 49 additions & 0 deletions settings/settings_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -503,6 +503,55 @@ func (s *Suite) TestEnvExport() {
})
}

func (s *Suite) TestEnvExportConnectionsCleansUpTempFile() {
s.Run("export, read, and remove all target the same in-container temp path", func() {
exportedPath := strings.TrimPrefix(catConnFile, "cat ")
s.True(strings.HasPrefix(exportedPath, "/tmp/"), "connections should be exported under /tmp inside the container, got %q", exportedPath)
s.Contains(airflowConnExport, exportedPath, "export command should write to the same path the cat command reads")
s.Equal(exportedPath, strings.TrimPrefix(rmConnFile, "rm "), "remove command should target the same path the file was exported to")
})

s.Run("removes the temp file on success", func() {
var commands []string
execAirflowCommand = func(id, airflowCommand string) (string, error) {
commands = append(commands, airflowCommand)
switch airflowCommand {
case airflowConnExport:
return "1 connections successfully exported", nil
case catConnFile:
return "local_postgres=postgres://username:password@example.db.example.com:5432/schema", nil
default:
return "", nil
}
}

err := EnvExportConnections("id", "testfiles/test.env")
s.NoError(err)
s.Contains(commands, rmConnFile)
_ = fileutil.WriteStringToFile("testfiles/test.env", "")
})

s.Run("still removes the temp file when reading the connections back fails", func() {
var commands []string
execAirflowCommand = func(id, airflowCommand string) (string, error) {
commands = append(commands, airflowCommand)
switch airflowCommand {
case airflowConnExport:
return "1 connections successfully exported", nil
case catConnFile:
return "", fmt.Errorf("boom")
default:
return "", nil
}
}

err := EnvExportConnections("id", "testfiles/test.env")
s.Error(err)
s.Contains(err.Error(), "error reading connections file")
s.Contains(commands, rmConnFile, "the temp file should still be removed even though the read failed")
})
}

func (s *Suite) TestExport() {
s.Run("success", func() {
WorkingPath = "./testfiles/"
Expand Down
Loading