// sqs-consumer — SQS poller → Fission function invoker // // Env vars: // SQS_ENDPOINT — URL SQS сервиса // SQS_QUEUE_NAME — имя очереди (обязательно) // SQS_ACCESS_KEY — AccessKeyId тенанта // SQS_SECRET_KEY — SecretAccessKey тенанта // SQS_REGION — регион (default: us-east-1) // FUNCTION_URL — полный URL функции // ROUTER_USERNAME — логин для /auth/login роутера (optional) // ROUTER_PASSWORD — пароль для /auth/login роутера (optional) // ROUTER_LOGIN_URL — URL /auth/login (default: выводится из FUNCTION_URL) // AUTH_TOKEN — статичный Bearer токен (если ROUTER_USERNAME не задан) // POLL_INTERVAL — интервал поллинга в секундах (default: 5) // MAX_MESSAGES — макс. сообщений за раз (default: 1) // MAX_RETRIES — попыток вызова функции перед skip (default: 3) // // 2026-05-12 package main import ( "bytes" "context" "encoding/json" "fmt" "io" "log/slog" "net/http" "os" "strconv" "strings" "sync" "time" "github.com/aws/aws-sdk-go-v2/aws" awsconfig "github.com/aws/aws-sdk-go-v2/config" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/sqs" ) // tokenManager управляет JWT токеном с автообновлением. type tokenManager struct { mu sync.Mutex client *http.Client loginURL string username string password string staticToken string cached string expiresAt time.Time } func newTokenManager(client *http.Client, functionURL, loginURL, username, password, staticToken string) *tokenManager { if loginURL == "" && username != "" { if i := strings.Index(functionURL, "://"); i >= 0 { rest := functionURL[i+3:] if j := strings.Index(rest, "/"); j >= 0 { loginURL = functionURL[:i+3] + rest[:j] + "/auth/login" } else { loginURL = functionURL + "/auth/login" } } } return &tokenManager{ client: client, loginURL: loginURL, username: username, password: password, staticToken: staticToken, } } func (tm *tokenManager) getToken(log *slog.Logger) string { tm.mu.Lock() defer tm.mu.Unlock() if tm.username == "" { return tm.staticToken } if tm.cached != "" && time.Now().Add(30*time.Second).Before(tm.expiresAt) { return tm.cached } token, exp := tm.login(log) if token != "" { tm.cached = token tm.expiresAt = exp log.Info("router token refreshed", "expiresAt", exp) } return tm.cached } func (tm *tokenManager) login(log *slog.Logger) (string, time.Time) { body, _ := json.Marshal(map[string]string{"username": tm.username, "password": tm.password}) resp, err := tm.client.Post(tm.loginURL, "application/json", bytes.NewReader(body)) if err != nil { log.Error("router login failed", "err", err) return "", time.Time{} } defer resp.Body.Close() if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated { log.Error("router login non-2xx", "status", resp.StatusCode) return "", time.Time{} } var result struct { AccessToken string `json:"accesstoken"` } if err := json.NewDecoder(resp.Body).Decode(&result); err != nil || result.AccessToken == "" { log.Error("router login decode failed", "err", err) return "", time.Time{} } return result.AccessToken, time.Now().Add(5 * time.Minute) } func (tm *tokenManager) invalidate() { tm.mu.Lock() defer tm.mu.Unlock() tm.cached = "" } func main() { log := slog.New(slog.NewJSONHandler(os.Stdout, nil)) endpoint := getenv("SQS_ENDPOINT", "http://shared-sqs.shared-sqs.svc.cluster.local:4100") queueName := mustenv("SQS_QUEUE_NAME") accessKey := mustenv("SQS_ACCESS_KEY") secretKey := mustenv("SQS_SECRET_KEY") region := getenv("SQS_REGION", "us-east-1") functionURL := mustenv("FUNCTION_URL") routerUser := getenv("ROUTER_USERNAME", "") routerPass := getenv("ROUTER_PASSWORD", "") routerLoginURL := getenv("ROUTER_LOGIN_URL", "") staticToken := getenv("AUTH_TOKEN", "") pollSec := parseInt(getenv("POLL_INTERVAL", "5"), 5) maxMsg := int32(parseInt(getenv("MAX_MESSAGES", "1"), 1)) maxRetries := parseInt(getenv("MAX_RETRIES", "3"), 3) customResolver := aws.EndpointResolverWithOptionsFunc( func(service, reg string, opts ...interface{}) (aws.Endpoint, error) { return aws.Endpoint{URL: endpoint, HostnameImmutable: true}, nil }, ) cfg, err := awsconfig.LoadDefaultConfig(context.Background(), awsconfig.WithRegion(region), awsconfig.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKey, secretKey, "")), awsconfig.WithEndpointResolverWithOptions(customResolver), ) if err != nil { log.Error("failed to create AWS config", "err", err) os.Exit(1) } client := sqs.NewFromConfig(cfg) ctx := context.Background() urlResult, err := client.GetQueueUrl(ctx, &sqs.GetQueueUrlInput{QueueName: aws.String(queueName)}) if err != nil { log.Error("GetQueueUrl failed", "queue", queueName, "err", err) os.Exit(1) } queueURL := aws.ToString(urlResult.QueueUrl) log.Info("sqs-consumer started", "queue", queueName, "endpoint", endpoint, "functionURL", functionURL) httpClient := &http.Client{Timeout: 30 * time.Second} tm := newTokenManager(httpClient, functionURL, routerLoginURL, routerUser, routerPass, staticToken) ticker := time.NewTicker(time.Duration(pollSec) * time.Second) defer ticker.Stop() for range ticker.C { poll(ctx, log, client, httpClient, tm, queueURL, functionURL, maxMsg, maxRetries) } } func poll(ctx context.Context, log *slog.Logger, sqsClient *sqs.Client, httpClient *http.Client, tm *tokenManager, queueURL, functionURL string, maxMsg int32, maxRetries int) { result, err := sqsClient.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{ QueueUrl: aws.String(queueURL), MaxNumberOfMessages: maxMsg, WaitTimeSeconds: 5, }) if err != nil { log.Error("SQS ReceiveMessage failed", "err", err) return } for _, msg := range result.Messages { body := aws.ToString(msg.Body) receipt := aws.ToString(msg.ReceiptHandle) if invokeWithRetry(log, httpClient, tm, functionURL, body, maxRetries) { if _, err := sqsClient.DeleteMessage(ctx, &sqs.DeleteMessageInput{ QueueUrl: aws.String(queueURL), ReceiptHandle: aws.String(receipt), }); err != nil { log.Error("DeleteMessage failed", "err", err) } else { log.Info("message processed", "msgId", aws.ToString(msg.MessageId)) } } } } func invokeWithRetry(log *slog.Logger, client *http.Client, tm *tokenManager, url, body string, maxRetries int) bool { for attempt := 1; attempt <= maxRetries; attempt++ { token := tm.getToken(log) req, _ := http.NewRequest(http.MethodPost, url, bytes.NewBufferString(body)) req.Header.Set("Content-Type", "application/json") if token != "" { req.Header.Set("Authorization", "Bearer "+token) } resp, err := client.Do(req) if err != nil { log.Warn("function invoke error", "attempt", attempt, "err", err) time.Sleep(time.Duration(attempt) * time.Second) continue } io.Copy(io.Discard, resp.Body) resp.Body.Close() if resp.StatusCode >= 200 && resp.StatusCode < 300 { return true } if resp.StatusCode == http.StatusUnauthorized { log.Warn("got 401, refreshing token", "attempt", attempt) tm.invalidate() } else { log.Warn("function returned non-2xx", "attempt", attempt, "status", resp.StatusCode) } time.Sleep(time.Duration(attempt) * time.Second) } log.Error("all retries exhausted, skipping message", "url", url) return false } func getenv(key, def string) string { if v := os.Getenv(key); v != "" { return v } return def } func mustenv(key string) string { v := os.Getenv(key) if v == "" { fmt.Fprintf(os.Stderr, "ERROR: env var %s is required\n", key) os.Exit(1) } return v } func parseInt(s string, def int) int { n, err := strconv.Atoi(s) if err != nil || n <= 0 { return def } return n }