- audiorecords вместо transcribe_jobs: приложения (texts, structures, recognitions, record_events, topics) живут своими коллекциями, ссылки на исходник и на приведённую копию перестали переставляться - рубеж называет достигнутое, отказ стал признаком остановки с причиной, а сторожей стало двое: число отказов и время в рубеже - воркеры потеряли специализацию, их число задаётся [pipeline] workers, шаг выбирается по рубежу, а захват отдаёт идентификатор и признак захвата
128 lines
4.4 KiB
Go
128 lines
4.4 KiB
Go
package yandex
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"strings"
|
|
|
|
"github.com/aws/aws-sdk-go-v2/aws"
|
|
"github.com/aws/aws-sdk-go-v2/config"
|
|
"github.com/aws/aws-sdk-go-v2/credentials"
|
|
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
|
|
"github.com/aws/aws-sdk-go-v2/service/s3"
|
|
s3types "github.com/aws/aws-sdk-go-v2/service/s3/types"
|
|
"github.com/aws/smithy-go"
|
|
)
|
|
|
|
type s3Config struct {
|
|
Region string
|
|
AccessKey string
|
|
SecretKey string
|
|
BucketName string
|
|
Endpoint string
|
|
}
|
|
|
|
type yandexS3Service struct {
|
|
client *s3.Client
|
|
uploader *manager.Uploader
|
|
bucketName string
|
|
endpoint string
|
|
}
|
|
|
|
func newYandexS3Service(cfg s3Config) (*yandexS3Service, error) {
|
|
if cfg.Region == "" || cfg.AccessKey == "" || cfg.SecretKey == "" || cfg.BucketName == "" {
|
|
return nil, fmt.Errorf("missing required S3 configuration parameters")
|
|
}
|
|
|
|
// Создаем конфигурацию
|
|
awsCfg, err := config.LoadDefaultConfig(context.Background(),
|
|
config.WithRegion(cfg.Region),
|
|
config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(cfg.AccessKey, cfg.SecretKey, "")),
|
|
)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to load AWS config: %w", err)
|
|
}
|
|
|
|
// Создаем клиент S3
|
|
var client *s3.Client
|
|
if cfg.Endpoint != "" {
|
|
// Кастомный endpoint (например, для MinIO)
|
|
client = s3.NewFromConfig(awsCfg, func(o *s3.Options) {
|
|
o.BaseEndpoint = aws.String(cfg.Endpoint)
|
|
o.UsePathStyle = true
|
|
})
|
|
} else {
|
|
// Стандартный AWS S3
|
|
client = s3.NewFromConfig(awsCfg)
|
|
}
|
|
|
|
uploader := manager.NewUploader(client)
|
|
|
|
return &yandexS3Service{
|
|
client: client,
|
|
uploader: uploader,
|
|
bucketName: cfg.BucketName,
|
|
endpoint: cfg.Endpoint,
|
|
}, nil
|
|
}
|
|
|
|
func (s *yandexS3Service) uploadFile(ctx context.Context, file io.Reader, fileName string) error {
|
|
_, err := s.uploader.Upload(ctx, &s3.PutObjectInput{
|
|
Bucket: aws.String(s.bucketName),
|
|
Key: aws.String(fileName),
|
|
Body: file,
|
|
})
|
|
if err != nil {
|
|
// Отказ SDK несёт полный URL объекта, то есть имя файла в хранилище, а
|
|
// оно — последняя часть ссылки на скачивание: цепочка `%w` уехала бы в
|
|
// журнал вместе с ключом. Наружу идёт класс отказа и только он — по
|
|
// нему «ключи отозваны» отличимо от «бакета нет» и от «сети нет», а
|
|
// адреса в коде отказа SDK не бывает.
|
|
var apiErr smithy.APIError
|
|
if errors.As(err, &apiErr) {
|
|
return fmt.Errorf("failed to upload file to S3: %s", apiErr.ErrorCode())
|
|
}
|
|
return errors.New("failed to upload file to S3")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// objectExists отвечает, лежит ли объект нужного размера.
|
|
//
|
|
// По нему шаг решает, повторять ли заливку: повтор её бесплатен, но дорог по
|
|
// времени на многочасовой записи. Сверка идёт по присутствию и длине, а не по
|
|
// отпечатку содержимого: признак целостности у составного объекта не равен
|
|
// отпечатку, и сверка хешем расходилась бы на всякой большой записи.
|
|
func (s *yandexS3Service) objectExists(ctx context.Context, objectKey string, size int64) (bool, error) {
|
|
out, err := s.client.HeadObject(ctx, &s3.HeadObjectInput{
|
|
Bucket: aws.String(s.bucketName),
|
|
Key: aws.String(objectKey),
|
|
})
|
|
if err != nil {
|
|
var notFound *s3types.NotFound
|
|
if errors.As(err, ¬Found) {
|
|
return false, nil
|
|
}
|
|
// Отказ SDK несёт полный URL объекта, а он — ключ к чужому аудио: наружу
|
|
// идёт класс отказа и только он.
|
|
var apiErr smithy.APIError
|
|
if errors.As(err, &apiErr) {
|
|
if apiErr.ErrorCode() == "NotFound" || apiErr.ErrorCode() == "NoSuchKey" {
|
|
return false, nil
|
|
}
|
|
return false, fmt.Errorf("failed to head object in S3: %s", apiErr.ErrorCode())
|
|
}
|
|
return false, errors.New("failed to head object in S3")
|
|
}
|
|
|
|
return out.ContentLength != nil && *out.ContentLength == size, nil
|
|
}
|
|
|
|
func (s *yandexS3Service) fileUrl(fileName string) string {
|
|
endpoint := strings.TrimRight(s.endpoint, "/")
|
|
return fmt.Sprintf("%s/%s/%s", endpoint, s.bucketName, fileName)
|
|
}
|