diff --git a/app/multitenant/dynamo_collector.go b/app/multitenant/dynamo_collector.go index 9c0896614..8c6e34b8f 100644 --- a/app/multitenant/dynamo_collector.go +++ b/app/multitenant/dynamo_collector.go @@ -9,7 +9,6 @@ import ( log "github.com/Sirupsen/logrus" "github.com/aws/aws-sdk-go/aws" - "github.com/aws/aws-sdk-go/aws/credentials" "github.com/aws/aws-sdk-go/aws/session" "github.com/aws/aws-sdk-go/service/dynamodb" "github.com/ugorji/go/codec" @@ -39,12 +38,9 @@ type dynamoDBCollector struct { // NewDynamoDBCollector the reaper of souls // https://github.com/aws/aws-sdk-go/wiki/common-examples -func NewDynamoDBCollector(url, region string, creds *credentials.Credentials, userIDer UserIDer) DynamoDBCollector { +func NewDynamoDBCollector(config *aws.Config, userIDer UserIDer) DynamoDBCollector { return &dynamoDBCollector{ - db: dynamodb.New(session.New(aws.NewConfig(). - WithEndpoint(url). - WithRegion(region). - WithCredentials(creds))), + db: dynamodb.New(session.New(config)), userIDer: userIDer, } } diff --git a/app/multitenant/sqs_control_router.go b/app/multitenant/sqs_control_router.go index 6358dd6e9..8f064997f 100644 --- a/app/multitenant/sqs_control_router.go +++ b/app/multitenant/sqs_control_router.go @@ -10,7 +10,6 @@ import ( log "github.com/Sirupsen/logrus" "github.com/aws/aws-sdk-go/aws" - "github.com/aws/aws-sdk-go/aws/credentials" "github.com/aws/aws-sdk-go/aws/session" "github.com/aws/aws-sdk-go/service/sqs" "golang.org/x/net/context" @@ -51,12 +50,9 @@ type sqsResponseMessage struct { } // NewSQSControlRouter the harbinger of death -func NewSQSControlRouter(url, region string, creds *credentials.Credentials, userIDer UserIDer) app.ControlRouter { +func NewSQSControlRouter(config *aws.Config, userIDer UserIDer) app.ControlRouter { result := &sqsControlRouter{ - service: sqs.New(session.New(aws.NewConfig(). - WithEndpoint(url). - WithRegion(region). - WithCredentials(creds))), + service: sqs.New(session.New(config)), responseQueueURL: nil, userIDer: userIDer, responses: map[string]chan xfer.Response{}, diff --git a/experimental/multitenant/run.sh b/experimental/multitenant/run.sh index e2488c0c6..25fd5ed4e 100755 --- a/experimental/multitenant/run.sh +++ b/experimental/multitenant/run.sh @@ -43,20 +43,16 @@ start_container 1 progrium/consul consul -p 8400:8400 -p 8500:8500 -p 8600:53/ud # These are the micro services common_args="--no-probe --app.weave.addr= --app.http.address=:80" -aws_args="--app.aws.region=us-east-1 --app.aws.id=abc --app.aws.secret=123 --app.aws.token=xyz" -start_container 2 weaveworks/scope collection -- ${common_args} --app.collector=dynamodb \ - ${aws_args} \ - --app.aws.dynamodb=http://dynamodb.weave.local:8000 \ +start_container 2 weaveworks/scope collection -- ${common_args} \ + --app.collector=dynamodb://abc:123@dynamodb.weave.local:8000 \ --app.aws.create.tables=true -start_container 2 weaveworks/scope query -- ${common_args} --app.collector=dynamodb \ - ${aws_args} \ - --app.aws.dynamodb=http://dynamodb.weave.local:8000 -start_container 2 weaveworks/scope controls -- ${common_args} --app.control.router=sqs \ - ${aws_args} \ - --app.aws.sqs=http://sqs.weave.local:9324 -start_container 2 weaveworks/scope pipes -- ${common_args} --app.pipe.router=consul \ - --app.consul.addr=consul.weave.local:8500 --app.consul.inf=ethwe \ - --app.consul.prefix=pipes/ +start_container 2 weaveworks/scope query -- ${common_args} \ + --app.collector=dynamodb://abc:123@dynamodb.weave.local:8000 +start_container 2 weaveworks/scope controls -- ${common_args} \ + --app.control.router=sqs://abc:123@sqs.weave.local:9324 +start_container 2 weaveworks/scope pipes -- ${common_args} \ + --app.pipe.router=consul://consul.weave.local:8500/pipes/ \ + --app.consul.inf=ethwe # And we bring it all together with a reverse proxy start_container 1 weaveworks/scope-frontend frontend --add-host=dns.weave.local:$(weave docker-bridge-ip) --publish=4040:80 diff --git a/prog/app.go b/prog/app.go index b80c7e74b..937d4c04c 100644 --- a/prog/app.go +++ b/prog/app.go @@ -2,13 +2,17 @@ package main import ( "flag" + "fmt" "math/rand" "net/http" _ "net/http/pprof" + "net/url" "strconv" + "strings" "time" log "github.com/Sirupsen/logrus" + "github.com/aws/aws-sdk-go/aws" "github.com/aws/aws-sdk-go/aws/credentials" "github.com/gorilla/mux" "github.com/weaveworks/go-checkpoint" @@ -30,6 +34,79 @@ func router(collector app.Collector, controlRouter app.ControlRouter, pipeRouter return app.TopologyHandler(collector, router, http.FileServer(FS(false))) } +func awsConfigFromURL(url *url.URL) *aws.Config { + password, _ := url.User.Password() + creds := credentials.NewStaticCredentials(url.User.Username(), password, "") + config := aws.NewConfig().WithCredentials(creds) + if strings.Contains(url.Host, ".") { + config = config.WithEndpoint(fmt.Sprintf("http://%s", url.Host)).WithRegion("dummy") + } else { + config = config.WithRegion(url.Host) + } + return config +} + +func collectorFactory(userIDer multitenant.UserIDer, collectorURL string, window time.Duration, createTables bool) (app.Collector, error) { + if collectorURL == "local" { + return app.NewCollector(window), nil + } + + parsed, err := url.Parse(collectorURL) + if err != nil { + return nil, err + } + + if parsed.Scheme == "dynamodb" { + dynamoCollector := multitenant.NewDynamoDBCollector(awsConfigFromURL(parsed), userIDer) + if createTables { + if err := dynamoCollector.CreateTables(); err != nil { + return nil, err + } + } + return dynamoCollector, nil + } + + return nil, fmt.Errorf("Invalid collector '%s'", collectorURL) +} + +func controlRouterFactory(userIDer multitenant.UserIDer, controlRouterURL string) (app.ControlRouter, error) { + if controlRouterURL == "local" { + return app.NewLocalControlRouter(), nil + } + + parsed, err := url.Parse(controlRouterURL) + if err != nil { + return nil, err + } + + if parsed.Scheme == "sqs" { + return multitenant.NewSQSControlRouter(awsConfigFromURL(parsed), userIDer), nil + } + + return nil, fmt.Errorf("Invalid control router '%s'", controlRouterURL) +} + +func pipeRouterFactory(userIDer multitenant.UserIDer, pipeRouterURL, consulInf string) (app.PipeRouter, error) { + if pipeRouterURL == "local" { + return app.NewLocalPipeRouter(), nil + } + + parsed, err := url.Parse(pipeRouterURL) + if err != nil { + return nil, err + } + + if parsed.Scheme == "consul" { + consulClient, err := multitenant.NewConsulClient(parsed.Host) + if err != nil { + return nil, err + } + return multitenant.NewConsulPipeRouter(consulClient, strings.TrimPrefix(parsed.Path, "/"), consulInf, userIDer) + } + + return nil, fmt.Errorf("Invalid pipe router '%s'", pipeRouterURL) +} + // Main runs the app func appMain() { var ( @@ -43,79 +120,39 @@ func appMain() { containerName = flag.String("container.name", app.DefaultContainerName, "Name of this container (to lookup container ID)") dockerEndpoint = flag.String("docker", app.DefaultDockerEndpoint, "Location of docker endpoint (to lookup container ID)") - collectorType = flag.String("collector", "local", "Collector to use (local of dynamodb)") - controlRouterType = flag.String("control.router", "local", "Control router to use (local or sqs)") - pipeRouterType = flag.String("pipe.router", "local", "Pipe router to use (local)") - userIDHeader = flag.String("userid.header", "", "HTTP header to use as userid") + collectorURL = flag.String("collector", "local", "Collector to use (local of dynamodb)") + controlRouterURL = flag.String("control.router", "local", "Control router to use (local or sqs)") + pipeRouterURL = flag.String("pipe.router", "local", "Pipe router to use (local)") + userIDHeader = flag.String("userid.header", "", "HTTP header to use as userid") - awsDynamoDB = flag.String("aws.dynamodb", "", "URL of DynamoDB instance") awsCreateTables = flag.Bool("aws.create.tables", false, "Create the tables in DynamoDB") - awsSQS = flag.String("aws.sqs", "", "URL of SQS instance") - awsRegion = flag.String("aws.region", "", "AWS Region") - awsID = flag.String("aws.id", "", "AWS Account ID") - awsSecret = flag.String("aws.secret", "", "AWS Account Secret") - awsToken = flag.String("aws.token", "", "AWS Account Token") - - consulPrefix = flag.String("consul.prefix", "", "Prefix for keys in consul") - consulAddr = flag.String("consul.addr", "", "Address of consul instance") - consulInf = flag.String("consul.inf", "", "The interface who's address I should advertise myself under in consul") + consulInf = flag.String("consul.inf", "", "The interface who's address I should advertise myself under in consul") ) flag.Parse() setLogLevel(*logLevel) setLogFormatter(*logPrefix) - // Do we need a user IDer? - var userIDer = multitenant.NoopUserIDer + userIDer := multitenant.NoopUserIDer if *userIDHeader != "" { userIDer = multitenant.UserIDHeader(*userIDHeader) } - // Create a collector - var collector app.Collector - if *collectorType == "local" { - collector = app.NewCollector(*window) - } else if *collectorType == "dynamodb" { - creds := credentials.NewStaticCredentials(*awsID, *awsSecret, *awsToken) - dynamoCollector := multitenant.NewDynamoDBCollector(*awsDynamoDB, *awsRegion, creds, userIDer) - collector = dynamoCollector - if *awsCreateTables { - if err := dynamoCollector.CreateTables(); err != nil { - log.Fatalf("Error createing DynamoDB tables: %v", err) - } - } - } else { - log.Fatalf("Invalid collector '%s'", *collectorType) + collector, err := collectorFactory(userIDer, *collectorURL, *window, *awsCreateTables) + if err != nil { + log.Fatalf("Error creating collector: %v", err) return } - // Create a control router - var controlRouter app.ControlRouter - if *controlRouterType == "local" { - controlRouter = app.NewLocalControlRouter() - } else if *controlRouterType == "sqs" { - creds := credentials.NewStaticCredentials(*awsID, *awsSecret, *awsToken) - controlRouter = multitenant.NewSQSControlRouter(*awsSQS, *awsRegion, creds, userIDer) - } else { - log.Fatalf("Invalid control router '%s'", *controlRouterType) + controlRouter, err := controlRouterFactory(userIDer, *controlRouterURL) + if err != nil { + log.Fatalf("Error creating control router: %v", err) return } - // Create a pipe router - var pipeRouter app.PipeRouter - if *pipeRouterType == "local" { - pipeRouter = app.NewLocalPipeRouter() - } else if *pipeRouterType == "consul" { - consulClient, err := multitenant.NewConsulClient(*consulAddr) - if err != nil { - log.Fatalf("Error createing consul client: %v", err) - } - pipeRouter, err = multitenant.NewConsulPipeRouter(consulClient, *consulPrefix, *consulInf, userIDer) - if err != nil { - log.Fatalf("Error createing consul pipe router: %v", err) - } - } else { - log.Fatalf("Invalid pipe router '%s'", *pipeRouterType) + pipeRouter, err := pipeRouterFactory(userIDer, *pipeRouterURL, *consulInf) + if err != nil { + log.Fatalf("Error creating pipe router: %v", err) return }