diff --git a/api/go.mod b/api/go.mod index 1f3a6463..4dfabc62 100644 --- a/api/go.mod +++ b/api/go.mod @@ -3,7 +3,6 @@ module github.com/USACE/cumulus-api/api go 1.26.6 require ( - github.com/USACE/go-simple-asyncer v0.0.0-20201015223104-446ae10887a8 github.com/aws/aws-sdk-go-v2 v1.42.0 github.com/btcsuite/btcutil v1.0.2 github.com/georgysavva/scany/v2 v2.1.4 @@ -21,7 +20,6 @@ require ( ) require ( - github.com/aws/aws-sdk-go v1.35.7 // indirect github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.13 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.19.23 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.29 // indirect @@ -38,8 +36,6 @@ require ( github.com/aws/aws-sdk-go-v2/service/sts v1.43.3 // indirect github.com/aws/smithy-go v1.27.2 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect - github.com/jmespath/go-jmespath v0.4.0 // indirect - github.com/streadway/amqp v1.0.0 // indirect golang.org/x/sync v0.21.0 // indirect ) diff --git a/api/go.sum b/api/go.sum index 68f9e7db..09035768 100644 --- a/api/go.sum +++ b/api/go.sum @@ -1,10 +1,6 @@ filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= -github.com/USACE/go-simple-asyncer v0.0.0-20201015223104-446ae10887a8 h1:ZiPYYp2OgNoJ489M5jiyqt23LV3yQ9w+ZiqiswQLWP8= -github.com/USACE/go-simple-asyncer v0.0.0-20201015223104-446ae10887a8/go.mod h1:2Zftz61ghmwOivA7LUAWt8rUV1CXmHuKYYvK2bFSxMY= github.com/aead/siphash v1.0.1/go.mod h1:Nywa3cDsYNNK3gaciGTWPwHt0wlpNV15vwmswBAUSII= -github.com/aws/aws-sdk-go v1.35.7 h1:FHMhVhyc/9jljgFAcGkQDYjpC9btM0B8VfkLBfctdNE= -github.com/aws/aws-sdk-go v1.35.7/go.mod h1:tlPOdRjfxPBpNIwqDj61rmsnA85v9jc0Ps9+muhnW+k= github.com/aws/aws-sdk-go-v2 v1.42.0 h1:XvXMJTkFQtpBKIWZnmr9ZEOc2InWM2yldjXEJ/bymhA= github.com/aws/aws-sdk-go-v2 v1.42.0/go.mod h1:27+ACypSLljLAEKsCYOmrjKh83vuTRkuAe9Uv/3A4bg= github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.13 h1:p1BBrg/Hhp6uK7zpejeI8QFXHJeC/mynzi04Sl03k9g= @@ -89,10 +85,6 @@ github.com/jackc/pgx/v5 v5.10.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QII github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/jessevdk/go-flags v0.0.0-20141203071132-1679536dcc89/go.mod h1:4FA24M0QyGHXBuZZK/XkWh8h0e1EYbRYJSGM75WSRxI= -github.com/jmespath/go-jmespath v0.4.0 h1:BEgLn5cpjn8UN1mAw4NjwDrS35OdebyEtFe+9YPoQUg= -github.com/jmespath/go-jmespath v0.4.0/go.mod h1:T8mJZnbsbmF+m6zOOFylbeCJqk5+pHWvzYPziyZiYoo= -github.com/jmespath/go-jmespath/internal/testify v1.5.1 h1:shLQSRRSCCPj3f2gpwzGwWFoC7ycTf1rcQZHOlsJ6N8= -github.com/jmespath/go-jmespath/internal/testify v1.5.1/go.mod h1:L3OGu8Wl2/fWfCI6z80xFu9LTZmf1ZRjMHUOPmWr69U= github.com/jmoiron/sqlx v1.4.0 h1:1PLqN7S1UYp5t4SrVVnt4nUVNemrDAtxlulVe+Qgm3o= github.com/jmoiron/sqlx v1.4.0/go.mod h1:ZrZ7UsYB/weZdl2Bxg6jCRO9c3YHl8r3ahlKmRT4JLY= github.com/jrick/logrotate v1.0.0/go.mod h1:LNinyqDIJnpAur+b8yyulnQw/wDuN1+BYKlTRt3OuAQ= @@ -136,8 +128,6 @@ github.com/prometheus/common v0.68.1 h1:omjRRl4QP4komogpXuhfeOiisQg7xdy8VM1UY+pS github.com/prometheus/common v0.68.1/go.mod h1:ZzL3f6u94qUxh9p+tJTrF+FvBS1XXbbRAZCQkytAL0Y= github.com/prometheus/procfs v0.20.1 h1:XwbrGOIplXW/AU3YhIhLODXMJYyC1isLFfYCsTEycfc= github.com/prometheus/procfs v0.20.1/go.mod h1:o9EMBZGRyvDrSPH1RqdxhojkuXstoe4UlK79eF5TGGo= -github.com/streadway/amqp v1.0.0 h1:kuuDrUJFZL1QYL9hUNuCxNObNzB0bV/ZG5jV3RWAQgo= -github.com/streadway/amqp v1.0.0/go.mod h1:AZpEONHx3DKn8O/DFsRAY58/XVQiIPMTMB1SddzLXVw= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.5.0 h1:1zr/of2m5FGMsad5YfcqgdqdWrIhu+EBEJRhR1U7z/c= github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= @@ -160,7 +150,6 @@ golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto= golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio= golang.org/x/net v0.0.0-20180906233101-161cd47e91fd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20200202094626-16171245cfb2/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= @@ -182,8 +171,6 @@ gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8 gopkg.in/fsnotify.v1 v1.4.7/go.mod h1:Tz8NjZHkW78fSQdbUxIjBTcgA1z1m8ZHf0WmKUhAMys= gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7/go.mod h1:dt/ZhP58zS4L8KSrWDmTeBkI65Dw0HsyUHuEVlX15mw= gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.2.8 h1:obN1ZagJSUGI0Ek/LBmuj4SNLPfIny3KsKFopxRdj10= -gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/api/handlers/statistics.go b/api/handlers/statistics.go deleted file mode 100644 index 7861ad36..00000000 --- a/api/handlers/statistics.go +++ /dev/null @@ -1,26 +0,0 @@ -package handlers - -import ( - "net/http" - - "github.com/jmoiron/sqlx" - "github.com/labstack/echo/v4" - - "github.com/USACE/cumulus-api/api/models" - - "github.com/USACE/go-simple-asyncer/asyncer" - - // SQL Interface - _ "github.com/jackc/pgx/v5" -) - -// DoStatistics triggers statistics -func DoStatistics(db *sqlx.DB, ae asyncer.Asyncer) echo.HandlerFunc { - return func(c echo.Context) error { - err := models.DoStatistics(db, ae) - if err != nil { - return c.String(http.StatusInternalServerError, err.Error()) - } - return c.NoContent(http.StatusCreated) - } -} diff --git a/api/models/statistics.go b/api/models/statistics.go deleted file mode 100644 index 5c9b80f3..00000000 --- a/api/models/statistics.go +++ /dev/null @@ -1,31 +0,0 @@ -package models - -import ( - "encoding/json" - - "github.com/USACE/go-simple-asyncer/asyncer" - - "github.com/jmoiron/sqlx" -) - -// DoStatistics creates a new acquisition and triggers acquisition functions -func DoStatistics(db *sqlx.DB, ae asyncer.Asyncer) error { - - // Check if Cron should run for each acquirable - payload, err := json.Marshal( - map[string]string{ - "s3_bucket": "cwbi-data-develop", - "s3_key": "cumulus/ncep_mrms_v12_MultiSensor_QPE_01H_Pass1/MRMS_MultiSensor_QPE_01H_Pass1_00.00_20201106-130000.tif", - "vector": "redrivernorth.shp", - }, - ) - if err != nil { - return err - } - // Async Call - if err := ae.CallAsync(payload); err != nil { - return err - } - - return nil -} diff --git a/async_listener/dispatch/dispatch.go b/async_listener/dispatch/dispatch.go new file mode 100644 index 00000000..7ddbcffa --- /dev/null +++ b/async_listener/dispatch/dispatch.go @@ -0,0 +1,130 @@ +// Package dispatch sends job payloads to the queues the Python workers consume. +// +// It replaces github.com/USACE/go-simple-asyncer, which pinned aws-sdk-go v1 +// and so carried CVE-2020-8911 and CVE-2020-8912 -- neither of which was ever +// patched on the v1 line. +// +// The wire contract is deliberately unchanged: a plain SQS message whose body +// is the JSON payload. That is what async_geoprocess/worker.py and +// async_packager/packager.py read via boto3 receive_messages, so nothing on the +// consumer side has to change. +package dispatch + +import ( + "context" + "fmt" + "log" + "net/url" + "path" + "strings" + "sync" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/sqs" +) + +// Sender delivers a payload to a worker queue. +type Sender interface { + Name() string + Send(ctx context.Context, payload []byte) error +} + +// New builds a Sender from the same ASYNC_ENGINE_* / ASYNC_ENGINE_*_TARGET +// environment contract go-simple-asyncer used, so no infrastructure config +// changes with this swap: +// +// engine "AWSSQS", target "" -> real SQS, URL resolved via GetQueueUrl +// engine "AWSSQS", target "local/" -> SQS-compatible endpoint (ElasticMQ locally) +// anything else -> logging no-op, as MockAsyncer was +// +// Engines "AWSSNS", "AWSLAMBDA" and "AMQP" existed upstream but are not +// implemented here: every ASYNC_ENGINE_* value in this repo is AWSSQS or MOCK. +// They fall through to the no-op rather than failing silently at send time, and +// New reports that clearly so a misconfiguration is visible at startup. +func New(ctx context.Context, engine, target string) (Sender, error) { + if !strings.EqualFold(engine, "AWSSQS") { + if engine != "" && !strings.EqualFold(engine, "MOCK") { + log.Printf( + "dispatch: engine %q is not implemented; falling back to a no-op sender. "+ + "Only AWSSQS performs real delivery.", engine, + ) + } + return noop{engine: engine, target: target}, nil + } + + awsCfg, err := config.LoadDefaultConfig(ctx) + if err != nil { + return nil, fmt.Errorf("dispatch: loading AWS config: %w", err) + } + + // "local/" points at an SQS-compatible endpoint (ElasticMQ in + // docker-compose). The queue URL is known up front, so no lookup is needed. + if len(target) > 6 && strings.EqualFold(target[:6], "local/") { + raw := target[6:] + u, err := url.Parse(raw) + if err != nil { + return nil, fmt.Errorf("dispatch: parsing local target %q: %w", raw, err) + } + endpoint := fmt.Sprintf("%s://%s", u.Scheme, u.Host) + _, queueName := path.Split(u.Path) + client := sqs.NewFromConfig(awsCfg, func(o *sqs.Options) { + o.BaseEndpoint = aws.String(endpoint) + }) + return &sqsSender{client: client, queueURL: raw, queueName: queueName}, nil + } + + return &sqsSender{client: sqs.NewFromConfig(awsCfg), queueName: target}, nil +} + +type sqsSender struct { + client *sqs.Client + queueName string + + // queueURL is known up front for local targets. For real SQS it is resolved + // on first send and cached. Resolution is deliberately lazy so that a queue + // which is not reachable at boot does not stop the listener from starting -- + // go-simple-asyncer resolved per message, which had the same tolerance but + // cost an extra API call on every dispatch. + mu sync.Mutex + queueURL string +} + +func (s *sqsSender) Name() string { return "AWSSQS" } + +func (s *sqsSender) url(ctx context.Context) (string, error) { + s.mu.Lock() + defer s.mu.Unlock() + if s.queueURL != "" { + return s.queueURL, nil + } + out, err := s.client.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{QueueName: aws.String(s.queueName)}) + if err != nil { + return "", fmt.Errorf("dispatch: resolving queue %q: %w", s.queueName, err) + } + s.queueURL = aws.ToString(out.QueueUrl) + return s.queueURL, nil +} + +func (s *sqsSender) Send(ctx context.Context, payload []byte) error { + queueURL, err := s.url(ctx) + if err != nil { + return err + } + if _, err := s.client.SendMessage(ctx, &sqs.SendMessageInput{ + QueueUrl: aws.String(queueURL), + MessageBody: aws.String(string(payload)), + }); err != nil { + return fmt.Errorf("dispatch: sending to %q: %w", s.queueName, err) + } + return nil +} + +type noop struct{ engine, target string } + +func (n noop) Name() string { return "NOOP" } + +func (n noop) Send(_ context.Context, payload []byte) error { + log.Printf("dispatch: no-op sender (engine=%q target=%q); payload: %s", n.engine, n.target, payload) + return nil +} diff --git a/async_listener/go.mod b/async_listener/go.mod index b7f1ab33..5ca99b0d 100644 --- a/async_listener/go.mod +++ b/async_listener/go.mod @@ -3,13 +3,24 @@ module github.com/USACE/cumulus-api/listener go 1.26.6 require ( - github.com/USACE/go-simple-asyncer v0.0.0-20201015223104-446ae10887a8 + github.com/aws/aws-sdk-go-v2 v1.43.7 + github.com/aws/aws-sdk-go-v2/config v1.32.38 + github.com/aws/aws-sdk-go-v2/service/sqs v1.46.7 github.com/kelseyhightower/envconfig v1.4.0 github.com/lib/pq v1.10.9 ) require ( - github.com/aws/aws-sdk-go v1.55.5 // indirect - github.com/jmespath/go-jmespath v0.4.0 // indirect - github.com/streadway/amqp v1.1.0 // indirect + github.com/aws/aws-sdk-go-v2/credentials v1.19.37 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.38 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.38 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.38 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.39 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.17 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.38 // indirect + github.com/aws/aws-sdk-go-v2/service/signin v1.5.7 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.33.7 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.7 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.45.7 // indirect + github.com/aws/smithy-go v1.27.8 // indirect ) diff --git a/async_listener/go.sum b/async_listener/go.sum index 5853750f..ad57cb7c 100644 --- a/async_listener/go.sum +++ b/async_listener/go.sum @@ -1,29 +1,34 @@ -github.com/USACE/go-simple-asyncer v0.0.0-20201015223104-446ae10887a8 h1:ZiPYYp2OgNoJ489M5jiyqt23LV3yQ9w+ZiqiswQLWP8= -github.com/USACE/go-simple-asyncer v0.0.0-20201015223104-446ae10887a8/go.mod h1:2Zftz61ghmwOivA7LUAWt8rUV1CXmHuKYYvK2bFSxMY= -github.com/aws/aws-sdk-go v1.35.7/go.mod h1:tlPOdRjfxPBpNIwqDj61rmsnA85v9jc0Ps9+muhnW+k= -github.com/aws/aws-sdk-go v1.55.5 h1:KKUZBfBoyqy5d3swXyiC7Q76ic40rYcbqH7qjh59kzU= -github.com/aws/aws-sdk-go v1.55.5/go.mod h1:eRwEWoyTWFMVYVQzKMNHWP5/RV4xIUGMQfXQHfHkpNU= -github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8= -github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/jmespath/go-jmespath v0.4.0 h1:BEgLn5cpjn8UN1mAw4NjwDrS35OdebyEtFe+9YPoQUg= -github.com/jmespath/go-jmespath v0.4.0/go.mod h1:T8mJZnbsbmF+m6zOOFylbeCJqk5+pHWvzYPziyZiYoo= -github.com/jmespath/go-jmespath/internal/testify v1.5.1 h1:shLQSRRSCCPj3f2gpwzGwWFoC7ycTf1rcQZHOlsJ6N8= -github.com/jmespath/go-jmespath/internal/testify v1.5.1/go.mod h1:L3OGu8Wl2/fWfCI6z80xFu9LTZmf1ZRjMHUOPmWr69U= +github.com/aws/aws-sdk-go-v2 v1.43.7 h1:msCzvkeYJA9ehbV8mRRmkZLo/zJg/+yDVLNtflg83hQ= +github.com/aws/aws-sdk-go-v2 v1.43.7/go.mod h1:tXpPM+v0D1lndmga+HqqLDIzUFJlEeR21aspVklHF00= +github.com/aws/aws-sdk-go-v2/config v1.32.38 h1:n4yPHBjtQ3BrIIUyk0/LAqf/BL2iv0Tw6XZcMRzM0ps= +github.com/aws/aws-sdk-go-v2/config v1.32.38/go.mod h1:dencYsOS1R7rBy8zehCvwBYzdxxL4Q/nRK7In03wjN8= +github.com/aws/aws-sdk-go-v2/credentials v1.19.37 h1:FJ8Iz4/xISMB/rwLlgfWujfGDFWr0oneQgtA6KPcYLY= +github.com/aws/aws-sdk-go-v2/credentials v1.19.37/go.mod h1:Q6pWOgVUp49x4g5QVi29wHofUoICnZ+Zq4jHbRN/7ec= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.38 h1:Nqo2jU1wz5rnBM9XQyXfVD1RP8txkbP3EDx8hR/hbCE= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.38/go.mod h1:PzJFHhjR2vWFKHe8HmY5Lxhvwyxnr5MERtk0nDxWNbk= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.38 h1:MBMg0zJ6i4TkAJ0dVFLKKn2cOkY6FkicmUDM67BRr6g= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.38/go.mod h1:9MWuJbyiUyj6eA7W1/zm1zuePDPSB3g+xcgRQeMWsXc= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.38 h1:lHm4jPf3k1Lz5ZWc+Vcn3MKVwym+26kWCba9FkJ4f0Y= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.38/go.mod h1:Rn+P2XR+FbyZzjmWKjg/KUZNxmGfr5oZwh5jQiE+CzI= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.39 h1:vo4xvMRs/F6h1E52qsgLqCQgWIQXgIJUauG6rlZEh4U= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.39/go.mod h1:jB03R1ij/A+OE2e1dz6vgj076gd7vlYcfstAzj3HcnU= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.17 h1:OvYZOB3qA6zvfdRFiRFRzVSiElMYrz3GdntkXZxlp1o= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.17/go.mod h1:JgR/2Ew50ACfIWau1oeMRX59tMtC0kM+PYQGEaT04cY= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.38 h1:H/5TI1jqaHsNoDQ60UwvPvJBg4GURkinXI3Qga29t2w= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.38/go.mod h1:PTVFf+XH++7NJOky+RLBYQx0QA5NcaeEYFQ2fsi0nwo= +github.com/aws/aws-sdk-go-v2/service/signin v1.5.7 h1:YcczQ6zNH/ojIzD/ikDrO+RfW06wmdMp18d4NH5hXY4= +github.com/aws/aws-sdk-go-v2/service/signin v1.5.7/go.mod h1:nl9RVnb9ulgAYzOkjLq1NyFxmWcnH2maCUEuOdESy98= +github.com/aws/aws-sdk-go-v2/service/sqs v1.46.7 h1:1A1rBhfNSJJ8JSxxWzbY2UbaI9wOu1G3FvjKsuhvUdg= +github.com/aws/aws-sdk-go-v2/service/sqs v1.46.7/go.mod h1:StIEARuthBzD6irPINvOKymBc4hE/QVZeMacPE1olE0= +github.com/aws/aws-sdk-go-v2/service/sso v1.33.7 h1:P+bMNiA93gyuYT3Oh+4dWtvrnGcu2bd9Uy5hRJM8BNo= +github.com/aws/aws-sdk-go-v2/service/sso v1.33.7/go.mod h1:zy+397isDFLvleg9H18Zq2MGzMso7uKyJyzR7DWSgFk= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.7 h1:WWkehGZ4nWtOKLMy0yi8+RqzzVqAGe60hGaxwF06JAw= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.7/go.mod h1:T8AI4SbQYm9ybcVmki2T3n7Qg1g3kfWoeQlNwNYOyO8= +github.com/aws/aws-sdk-go-v2/service/sts v1.45.7 h1:yU/9y2r7s9kSUPbHXbpQTa4LA8kt+CMgpu1OBrhx8p4= +github.com/aws/aws-sdk-go-v2/service/sts v1.45.7/go.mod h1:0lQTDEBArMevQXpxu443LVGjKxxEeSsSnrw9n8YiTMg= +github.com/aws/smithy-go v1.27.8 h1:FR0dxZfIlV7Z8eh2iHfIofdunw382XsDV3Mxt9nUvRY= +github.com/aws/smithy-go v1.27.8/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/kelseyhightower/envconfig v1.4.0 h1:Im6hONhd3pLkfDFsbRgu68RDNkGF1r3dvMUtDTo2cv8= github.com/kelseyhightower/envconfig v1.4.0/go.mod h1:cccZRl6mQpaq41TPp5QxidR+Sa3axMbJDNb//FQX6Gg= github.com/lib/pq v1.10.9 h1:YXG7RB+JIjhP29X+OtkiDnYaXQwpS4JEWq7dtCCRUEw= github.com/lib/pq v1.10.9/go.mod h1:AlVN5x4E4T544tWzH6hKfbfQvm3HdbOxrmggDNAPY9o= -github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= -github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/streadway/amqp v1.0.0/go.mod h1:AZpEONHx3DKn8O/DFsRAY58/XVQiIPMTMB1SddzLXVw= -github.com/streadway/amqp v1.1.0 h1:py12iX8XSyI7aN/3dUT8DFIDJazNJsVJdxNVEpnQTZM= -github.com/streadway/amqp v1.1.0/go.mod h1:WYSrTEYHOXHd0nwFeUXAe2G2hRnQT+deZJJf88uS9Bg= -github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= -golang.org/x/net v0.0.0-20200202094626-16171245cfb2/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/yaml.v2 v2.2.8 h1:obN1ZagJSUGI0Ek/LBmuj4SNLPfIny3KsKFopxRdj10= -gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= diff --git a/async_listener/listener/main.go b/async_listener/listener/main.go index 976a49f8..8caf53ab 100644 --- a/async_listener/listener/main.go +++ b/async_listener/listener/main.go @@ -1,12 +1,13 @@ package main import ( + "context" "encoding/json" "fmt" "log" "time" - "github.com/USACE/go-simple-asyncer/asyncer" + "github.com/USACE/cumulus-api/listener/dispatch" "github.com/kelseyhightower/envconfig" "github.com/lib/pq" ) @@ -61,11 +62,17 @@ func (c Config) maxReconn() time.Duration { return d } -// NewAsyncNotificationHandler handles dependency injection of asyncer.Asyncer -func NewAsyncNotificationHandler(a asyncer.Asyncer) NotificationHandler { +// sendTimeout bounds a single dispatch. go-simple-asyncer had no timeout at +// all, so a hung send leaked its goroutine for the life of the process. +const sendTimeout = 30 * time.Second + +// NewAsyncNotificationHandler handles dependency injection of dispatch.Sender +func NewAsyncNotificationHandler(s dispatch.Sender) NotificationHandler { return func(d string) error { - if err := a.CallAsync([]byte(d)); err != nil { - fmt.Println("Error calling async") + ctx, cancel := context.WithTimeout(context.Background(), sendTimeout) + defer cancel() + if err := s.Send(ctx, []byte(d)); err != nil { + fmt.Println("Error dispatching to worker queue") fmt.Println(err.Error()) return err } @@ -100,6 +107,8 @@ func reportProblem(eq pq.ListenerEventType, err error) { func main() { + ctx := context.Background() + var cfg Config if err := envconfig.Process("cumulus", &cfg); err != nil { log.Fatal(err.Error()) @@ -113,26 +122,26 @@ func main() { } // downloadAsyncer defines async engine used to package DSS files for download - downloadAsyncer, err := asyncer.NewAsyncer(asyncer.Config{Engine: cfg.AsyncEnginePackager, Target: cfg.AsyncEnginePackagerTarget}) + downloadSender, err := dispatch.New(ctx, cfg.AsyncEnginePackager, cfg.AsyncEnginePackagerTarget) if err != nil { log.Fatal(err.Error()) } - d := NewAsyncNotificationHandler(downloadAsyncer) + d := NewAsyncNotificationHandler(downloadSender) // acquirablefileAsyncer defines async engine for processing new acquirable files - geoprocessAsyncer, err := asyncer.NewAsyncer(asyncer.Config{Engine: cfg.AsyncEngineGeoprocess, Target: cfg.AsyncEngineGeoprocessTarget}) + geoprocessSender, err := dispatch.New(ctx, cfg.AsyncEngineGeoprocess, cfg.AsyncEngineGeoprocessTarget) if err != nil { log.Fatal(err.Error()) } - g := NewAsyncNotificationHandler(geoprocessAsyncer) + g := NewAsyncNotificationHandler(geoprocessSender) // statisticsAsyncer defines async engine for computing raster statistics - // statisticsAsyncer, err := asyncer.NewAsyncer(asyncer.Config{Engine: cfg.AsyncEngineStatistics, Target: cfg.AsyncEngineStatisticsTarget}) + // statisticsSender, err := dispatch.New(ctx, cfg.AsyncEngineStatistics, cfg.AsyncEngineStatisticsTarget) // if err != nil { // log.Fatal(err.Error()) // } // // AddstatisticsAsyncer to map of registered handlers - // handlers["notify_statistics"] = NewAsyncNotificationHandler(statisticsAsyncer) + // handlers["notify_statistics"] = NewAsyncNotificationHandler(statisticsSender) // Map of handlers handlers := map[string]NotificationHandler{