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
1 change: 0 additions & 1 deletion .mockery.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ packages:
interfaces:
RegistryHandler:
ImageHandler:
DAGCheck:
DockerRegistryAPI:
github.com/astronomer/astro-cli/internal/platform/astro/clients/astrov1:
config:
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ astro local check Validate this project's DAGs without starting Airflow (--s
astro local upgrade airflow [version] Move the project to a new Airflow, the newest of its generation if none is given (--with-otto hands the rest to Otto)
```

`astro start`, `astro stop`, and `astro logs` work as shorthand for the `local` versions. The cloud commands — `astro login`, `astro deploy`, `astro deployment`, `astro workspace` — are unchanged from v1.
`astro start`, `astro stop`, and `astro logs` work as shorthand for the `local` versions. The cloud commands — `astro login`, `astro deployment`, `astro workspace` — work as they did in v1. `astro deploy` ships a project with a `pyproject.toml`: a project made by Astro CLI 1.x is converted with `astro init` first, or keeps deploying with Astro CLI 1.x.

Every command takes `--output json` for scripting; streaming commands like `logs` emit one JSON object per line.

Expand Down
3 changes: 0 additions & 3 deletions airflow/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,3 @@ const (
QuayBaseImageName = "quay.io/astronomer"
AstroImageRegistryBaseImageName = "astrocrpublic.azurecr.io"
)

// DefaultTestPath is the DAG integrity test a deploy's parse step runs.
const DefaultTestPath = ".astro/test_dag_integrity_default.py"
73 changes: 0 additions & 73 deletions airflow/container.go
Original file line number Diff line number Diff line change
@@ -1,26 +1,11 @@
package airflow

import (
"crypto/md5" //nolint:gosec // reviewed; not a new risk in this shell code
"fmt"
"regexp"
"strings"

"github.com/docker/docker/client"
"github.com/pkg/errors"

"github.com/astronomer/astro-cli/airflow/types"
"github.com/astronomer/astro-cli/config"
"github.com/astronomer/astro-cli/pkg/fileutil"
)

// DAGCheck checks a project's DAGs in its image, the parse and
// pytest steps of `astro deploy`.
type DAGCheck interface {
Pytest(pytestFile, customImageName, deployImageName, pytestArgsString string, buildSecrets []string) (string, error)
Parse(customImageName, deployImageName string, buildSecrets []string) error
}

// RegistryHandler defines methods require to handle all operations with registry
type RegistryHandler interface {
Login(username, token string) error
Expand All @@ -32,74 +17,16 @@ type ImageHandler interface {
Push(remoteImage, username, token string, getImageRepoSha bool) (string, error)
GetLabel(altImageName, labelName string) (string, error)
TagLocalImage(localImage string) error
Pytest(pytestFile, airflowHome, envFile, testHomeDirectory string, pytestArgs []string, htmlReport bool, config types.ImageBuildConfig) (string, error)
}

type DockerRegistryAPI interface {
client.APIClient
}

func DAGCheckInit(airflowHome, envFile, dockerfile, projectName string) (DAGCheck, error) {
return NewDAGChecker(airflowHome, envFile, dockerfile, projectName)
}

func RegistryHandlerInit(registry string) (RegistryHandler, error) {
return DockerRegistryInit(registry)
}

func ImageHandlerInit(image string) ImageHandler {
return DockerImageInit(image)
}

// ProjectNameUnique creates a reasonably unique project name based on the hashed
// path of the project. This prevents collisions of projects with identical dir names
// in different paths. ie (~/dev/project1 vs ~/prod/project1)
func ProjectNameUnique() (string, error) {
projectName := config.CFG.ProjectName.GetString()

pwd, err := fileutil.GetWorkingDir()
if err != nil {
return "", errors.Wrap(err, "error retrieving working directory")
}

// #nosec
b := md5.Sum([]byte(pwd))
s := fmt.Sprintf("%x", b[:])

return sanitizeImageName(projectName + "_" + s[0:6]), nil
}

// validImageName matches docker's grammar for a single image-name path
// component: alphanumeric runs joined by single "." / "_" / "-" separators (or
// a double "__"), starting and ending with an alphanumeric.
var validImageName = regexp.MustCompile(`^[a-z0-9]+(?:(?:[._]|__|-+)[a-z0-9]+)*$`)

// sanitizeImageName turns an arbitrary project name into a name docker will
// accept as an image tag. Names docker already accepts are returned unchanged,
// so cached image tags keep their names. Anything else is lowercased, has every
// run of non-alphanumeric characters collapsed to a single "-", and its leading
// and trailing separators trimmed. The result is always a non-empty valid name.
func sanitizeImageName(s string) string {
s = strings.ToLower(s)
if validImageName.MatchString(s) {
return s
}

var b strings.Builder
prevSep := false
for _, r := range s {
if (r >= 'a' && r <= 'z') || (r >= '0' && r <= '9') {
b.WriteRune(r)
prevSep = false
} else if !prevSep {
b.WriteByte('-')
prevSep = true
}
}

out := strings.Trim(b.String(), "-")
if out == "" {
out = "project"
}
return out
}
117 changes: 0 additions & 117 deletions airflow/docker.go
Original file line number Diff line number Diff line change
@@ -1,123 +1,6 @@
package airflow

import (
"fmt"
"strconv"
"strings"

"github.com/pkg/errors"

airflowTypes "github.com/astronomer/astro-cli/airflow/types"
"github.com/astronomer/astro-cli/pkg/ansi"
"github.com/astronomer/astro-cli/pkg/util"
)

const (
RuntimeImageLabel = "io.astronomer.docker.runtime.version"
pytestDirectory = "tests"
componentName = "airflow"
)

// DAGChecker runs a project's DAG checks in its image: the parse and
// pytest steps of a 1.x project's deploy. It is what is left of the 1.x
// `astro dev` container handler, which drove docker compose.
type DAGChecker struct {
airflowHome string
envFile string
dockerfile string
imageHandler ImageHandler
}

func NewDAGChecker(airflowHome, envFile, dockerfile, imageName string) (*DAGChecker, error) {
if imageName == "" {
// Get project name from config
projectName, err := ProjectNameUnique()
if err != nil {
return nil, fmt.Errorf("error retrieving working directory: %w", err)
}
imageName = projectName
}

return &DAGChecker{
airflowHome: airflowHome,
envFile: envFile,
dockerfile: dockerfile,
imageHandler: DockerImageInit(ImageName(imageName, "latest")),
}, nil
}

// Pytest creates and runs a container containing the users airflow image, requirments, packages, and volumes(DAGs folder, etc...)
// These containers runs pytest on a specified pytest file (pytestFile). A deploy's --pytest and --parse use it
func (d *DAGChecker) Pytest(pytestFile, customImageName, deployImageName, pytestArgsString string, buildSecrets []string) (string, error) {
// deployImageName may be provided to the function if it is being used in the deploy command
if deployImageName == "" {
// build image
if customImageName == "" {
err := d.imageHandler.Build(d.dockerfile, buildSecrets, airflowTypes.ImageBuildConfig{Path: d.airflowHome})
if err != nil {
return "", err
}
} else {
// skip build if an customImageName is passed
err := d.imageHandler.TagLocalImage(customImageName)
if err != nil {
return "", err
}
}
}

// determine pytest args and file
pytestArgs := strings.Fields(pytestArgsString)

// Determine pytest file
if pytestFile != DefaultTestPath {
if !strings.Contains(pytestFile, pytestDirectory) {
pytestFile = pytestDirectory + "/" + pytestFile
} else if pytestFile == "" {
pytestFile = pytestDirectory + "/"
}
}

// run pytests
exitCode, err := d.imageHandler.Pytest(pytestFile, d.airflowHome, d.envFile, "", pytestArgs, false, airflowTypes.ImageBuildConfig{Path: d.airflowHome})
if err != nil {
return exitCode, err
}
if code, convErr := strconv.Atoi(exitCode); convErr == nil && code == 0 { // exit code 0 means the pytests passed
return "", nil
}
return exitCode, errors.New("something went wrong while Pytesting your Dags")
}

func (d *DAGChecker) Parse(customImageName, deployImageName string, buildSecrets []string) error {
// check for file
path := d.airflowHome + "/" + DefaultTestPath

fileExist, err := util.Exists(path)
if err != nil {
return err
}
if !fileExist {
// Only a 1.x project deploys with --parse, and Astro CLI 1.x's
// `astro dev init` wrote this file into it; v2 writes it nowhere. A
// project without it has never had a parse check to run, so the deploy
// goes on, as it always has, but says so.
fmt.Println("\nSkipping the DAG parse check: it runs " + path + ", which this project does not have. " +
"Astro CLI 1.x's `astro dev init` created that file; add it back to the project to run the check.")

return nil
}

fmt.Println("Checking your Dags for errors…")

pytestFile := DefaultTestPath
exitCode, err := d.Pytest(pytestFile, customImageName, deployImageName, "", buildSecrets)
if err != nil {
if code, convErr := strconv.Atoi(exitCode); convErr == nil && code == 1 { // exit code 1 means tests failed
return errors.New("See above for errors detected in your Dags")
}
return errors.Wrap(err, "something went wrong while parsing your Dags")
}
fmt.Println(ansi.Green("✔") + " No errors detected in your Dags ")
return err
}
98 changes: 0 additions & 98 deletions airflow/docker_image.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@ import (
"github.com/astronomer/astro-cli/config"
"github.com/astronomer/astro-cli/pkg/logger"
"github.com/astronomer/astro-cli/pkg/spinner"
"github.com/astronomer/astro-cli/pkg/util"
)

const (
Expand Down Expand Up @@ -182,103 +181,6 @@ func (d *DockerImage) Build(dockerfilePath string, buildSecrets []string, buildC
return nil
}

func (d *DockerImage) Pytest(pytestFile, airflowHome, envFile, testHomeDirectory string, pytestArgs []string, htmlReport bool, buildConfig airflowTypes.ImageBuildConfig) (string, error) {
// delete container
containerRuntime, err := runtimes.GetContainerRuntimeBinary()
if err != nil {
return "", err
}
err = cmdExec(containerRuntime, nil, nil, "rm", "astro-pytest")
if err != nil {
logger.Debug(err)
}
// Change to location of Dockerfile
err = os.Chdir(buildConfig.Path)
if err != nil {
return "", err
}
args := []string{
"create",
"-i",
"--name",
"astro-pytest",
}
fileExist, err := util.Exists(airflowHome + "/" + envFile)
if err != nil {
return "", err
}
if fileExist {
args = append(args, []string{"--env-file", envFile}...)
}
args = append(args, []string{d.imageName, "pytest", pytestFile}...)
args = append(args, pytestArgs...)
// run pytest image
var stdout, stderr io.Writer
if logger.IsLevelEnabled(logrus.WarnLevel) {
stdout = os.Stdout
stderr = os.Stderr
} else {
stdout = nil
stderr = nil
}

// create pytest container
err = cmdExec(containerRuntime, stdout, stderr, args...)
if err != nil {
return "", err
}

// Copy host directories into the container using docker cp.
// This ensures fresh files from the host are used (not stale from
// image build cache) and works with remote Docker daemons (CI).
copyDirs := []string{"dags", "tests", "plugins", "include", ".astro"}
for _, dir := range copyDirs {
srcPath := airflowHome + "/" + dir
if exists, _ := util.Exists(srcPath); !exists { //nolint:errcheck // treated as absent on error
continue
}
docErr := cmdExec(containerRuntime, stdout, stderr, "cp", srcPath, "astro-pytest:/usr/local/airflow/")
if docErr != nil {
return "", docErr
}
}

// start pytest container
err = cmdExec(containerRuntime, stdout, stderr, []string{"start", "astro-pytest", "-a"}...)
if err != nil {
logger.Debugf("Error starting pytest container: %s", err.Error())
}

// get exit code
args = []string{
"inspect",
"astro-pytest",
"--format={{.State.ExitCode}}",
}
var outb bytes.Buffer
inspectErr := cmdExec(containerRuntime, &outb, stderr, args...)
if inspectErr != nil {
logger.Debug(inspectErr)
}

if htmlReport {
// Copy the dag-test-report.html file from the container to the destination folder
cpErr := cmdExec(containerRuntime, nil, stderr, "cp", "astro-pytest:/usr/local/airflow/dag-test-report.html", "./"+testHomeDirectory)
if cpErr != nil {
logger.Debugf("Error copying dag-test-report.html file from the pytest container: %s", cpErr.Error())
}
}

// delete container
rmErr := cmdExec(containerRuntime, nil, stderr, "rm", "astro-pytest")
if rmErr != nil {
logger.Debugf("Error removing the astro-pytest container: %s", rmErr.Error())
}

// trim the trailing newline so consumers get a clean integer string to parse
return strings.TrimSpace(outb.String()), err
}

func (d *DockerImage) Push(remoteImage, username, token string, getImageRepoSha bool) (string, error) {
containerRuntime, err := runtimes.GetContainerRuntimeBinary()
if err != nil {
Expand Down
Loading
Loading