/
t3
/
cli
Обзор
Документация
Войти
/
t3
/
cli
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
internal/engine/bridge/s3_cache.go
720 строк
23 KB
Ivan
move s3 to the root
09 июл 2026, 13:51
09 июл 2026, 13:51
c214eff
Код
Авторство
О чём код?
package bridge import ( "bytes" "context" "crypto/sha256" "crypto/tls" "encoding/binary" "encoding/hex" "encoding/json" "errors" "fmt" "io" "net/http" "net/url" "os" "path" "path/filepath" "sort" "strings" "sync" "time" "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/service/s3" smithyhttp "github.com/aws/smithy-go/transport/http" ) // diskEntry хранит информацию о файле на диске для LRU-эвикшена. type diskEntry struct { size int64 modTime time.Time } // S3Cache управляет кэшированием S3-файлов. // Кэш имеет два уровня: // - in-memory: map[s3URL]*CachedFile для быстрого доступа // - disk: файлы в cacheDir, загруженные между запусками // // После завершения agent_init() кэш запечатывается (sealed), // и новые файлы нельзя добавить — можно только читать существующие. type S3Cache struct { mu sync.RWMutex sealed bool files map[string]*CachedFile // key = s3 URL maxAge time.Duration cacheDir string // путь к папке дискового кэша (пусто = отключён) cacheMaxSize int64 // макс. размер дискового кэша в байтах diskIndex map[string]*diskEntry // key = хэш URL → entry currentSize int64 // текущий размер дискового кэша } // NewS3Cache создает новый S3-кэш. // maxAge — максимальное время жизни кэшированного файла до проверки ETag. // cacheDir — путь к папке для дискового кэша (пустая строка = без дискового кэша). // cacheMaxSize — максимальный размер дискового кэша в байтах (игнорируется при cacheDir == ""). func NewS3Cache(maxAge time.Duration, cacheDir string, cacheMaxSize int64) *S3Cache { if maxAge <= 0 { maxAge = 5 * time.Minute //nolint:mnd } c := &S3Cache{ files: make(map[string]*CachedFile), maxAge: maxAge, cacheDir: cacheDir, cacheMaxSize: cacheMaxSize, diskIndex: make(map[string]*diskEntry), } // Если указана папка кэша — инициализируем дисковый кэш if cacheDir != "" { c.initDiskCache() } return c } // initDiskCache сканирует cacheDir и строит diskIndex. func (c *S3Cache) initDiskCache() { // Создаём папку, если не существует if err := os.MkdirAll(c.cacheDir, 0755); err != nil { // Не фатально — просто не будет дискового кэша return } entries, err := os.ReadDir(c.cacheDir) if err != nil { return } for _, entry := range entries { if entry.IsDir() { continue } info, err := entry.Info() if err != nil { continue } hash := entry.Name() if len(hash) != 64 { // SHA256 hex = 64 символа continue } c.diskIndex[hash] = &diskEntry{ size: info.Size(), modTime: info.ModTime(), } c.currentSize += info.Size() } } // cacheFilePath возвращает путь к файлу на диске для заданного URL. func (c *S3Cache) cacheFilePath(s3URL string) string { hash := sha256.Sum256([]byte(s3URL)) return filepath.Join(c.cacheDir, hex.EncodeToString(hash[:])) } // Sealed возвращает true, если кэш запечатан (agent_init завершен). func (c *S3Cache) Sealed() bool { c.mu.RLock() defer c.mu.RUnlock() return c.sealed } // Seal запечатывает кэш. После вызова Seal() новые файлы нельзя добавить. func (c *S3Cache) Seal() { c.mu.Lock() defer c.mu.Unlock() c.sealed = true } // Get возвращает файл из кэша по URL. // Если файла нет в кэше и кэш запечатан — возвращает ошибку. // Если кэш не запечатан (agent_init еще выполняется), скачивает файл. // Между запусками проверяет закэшированный на диск файл через ETag (If-None-Match). func (c *S3Cache) Get(ctx context.Context, s3URL string, s3Cfg *S3Config) (*CachedFile, error) { c.mu.RLock() sealed := c.sealed f, ok := c.files[s3URL] c.mu.RUnlock() if ok { // Файл есть в in-memory кэше — проверяем ETag, если expired return c.handleCacheHit(ctx, s3URL, s3Cfg, f) } // Файла нет в памяти — пробуем загрузить с диска (если есть дисковый кэш) if c.cacheDir != "" { diskFile, err := c.loadFromDisk(ctx, s3URL, s3Cfg) if err == nil { // Успешно загрузили с диска — сохраняем в память c.mu.Lock() c.files[s3URL] = diskFile c.mu.Unlock() // Проверяем ETag, если expired return c.handleCacheHit(ctx, s3URL, s3Cfg, diskFile) } } if sealed { return nil, fmt.Errorf("file '%s' was not preloaded in agent_init(). Add t3.s3.download() to your agent_init() function", s3URL) } // Фаза 1: реальная загрузка (вызывается из agent_init) return c.downloadAndCache(ctx, s3URL, s3Cfg) } // handleCacheHit обрабатывает ситуацию, когда файл найден в кэше (in-memory или с диска). // Проверяет ETag, если файл expired. func (c *S3Cache) handleCacheHit(ctx context.Context, s3URL string, s3Cfg *S3Config, f *CachedFile) (*CachedFile, error) { if f.IsExpired(c.maxAge) { c.mu.Lock() // Проверяем повторно под локом if f.IsExpired(c.maxAge) { updated, err := c.checkETag(ctx, s3URL, s3Cfg, f) if err != nil { c.mu.Unlock() return nil, fmt.Errorf("failed to check etag for %s: %w", s3URL, err) } if updated != nil { f = updated c.files[s3URL] = f // Сохраняем обновлённый файл на диск if c.cacheDir != "" { if saveErr := c.saveToDisk(f); saveErr != nil { // Не фатально — кэш в памяти работает } } } } c.mu.Unlock() } return f, nil } // downloadAndCache скачивает файл с S3 и сохраняет в кэш (память + диск). func (c *S3Cache) downloadAndCache(ctx context.Context, s3URL string, s3Cfg *S3Config) (*CachedFile, error) { c.mu.Lock() defer c.mu.Unlock() // Проверяем повторно — другой goroutine мог уже загрузить if f, ok := c.files[s3URL]; ok { return f, nil } f, err := c.downloadFromS3(ctx, s3URL, s3Cfg) if err != nil { return nil, err } c.files[s3URL] = f // Сохраняем на диск, если папка кэша указана if c.cacheDir != "" { if saveErr := c.saveToDisk(f); saveErr != nil { // Не фатально } } return f, nil } // loadFromDisk загружает файл с диска, проверяет ETag (If-None-Match). // Если файл изменился на S3 — скачивает заново. // Если файла нет на диске — возвращает ошибку. func (c *S3Cache) loadFromDisk(ctx context.Context, s3URL string, s3Cfg *S3Config) (*CachedFile, error) { diskPath := c.cacheFilePath(s3URL) // Проверяем, есть ли файл на диске if _, err := os.Stat(diskPath); os.IsNotExist(err) { return nil, fmt.Errorf("file not in disk cache: %s", s3URL) } // Читаем файл с диска data, err := os.ReadFile(diskPath) if err != nil { return nil, fmt.Errorf("failed to read disk cache for %s: %w", s3URL, err) } // Парсим бинарный формат: [4 байта длина ETag] [ETag строка] [8 байт Updated] [RawBytes] if len(data) < 4 { return nil, fmt.Errorf("corrupt disk cache for %s: too short", s3URL) } etagLen := int(binary.LittleEndian.Uint32(data[:4])) if len(data) < 4+etagLen+8 { return nil, fmt.Errorf("corrupt disk cache for %s: truncated etag or missing updated timestamp", s3URL) } etag := string(data[4 : 4+etagLen]) // Пропускаем 8 байт Updated (не используется — при загрузке с диска всегда expired) _ = int64(binary.LittleEndian.Uint64(data[4+etagLen : 4+etagLen+8])) rawBytes := data[4+etagLen+8:] // Преобразуем в CachedFile content := string(rawBytes) lines := strings.Split(content, "\n") fileName := path.Base(extractKeyFromURL(s3URL)) // Пробуем парсить как JSON var parsedJSON interface{} if err := json.Unmarshal(rawBytes, &parsedJSON); err == nil { // валидный JSON } else { parsedJSON = nil } cf := &CachedFile{ ETag: etag, URL: s3URL, FileName: fileName, RawBytes: rawBytes, Content: content, Lines: lines, ParsedJSON: parsedJSON, Updated: time.Time{}, // всегда expired → принудительная ETag-проверка из handleCacheHit } return cf, nil } // extractKeyFromURL извлекает ключ (путь) из s3 URL. func extractKeyFromURL(rawURL string) string { trimmed := strings.TrimPrefix(rawURL, "s3://") parts := strings.SplitN(trimmed, "/", 3) if len(parts) == 3 { return parts[2] } if len(parts) == 2 { return parts[1] } return rawURL } // saveToDisk сохраняет CachedFile на диск. // Проверяет лимит размера кэша и при необходимости удаляет старые файлы (LRU). func (c *S3Cache) saveToDisk(f *CachedFile) error { if c.cacheDir == "" { return nil } diskPath := c.cacheFilePath(f.URL) hash := filepath.Base(diskPath) // Сериализуем: [4 байта длина ETag] [ETag] [8 байт Updated] [RawBytes] etagBytes := []byte(f.ETag) buf := new(bytes.Buffer) etagLen := make([]byte, 4) binary.LittleEndian.PutUint32(etagLen, uint32(len(etagBytes))) buf.Write(etagLen) buf.Write(etagBytes) updatedNano := make([]byte, 8) binary.LittleEndian.PutUint64(updatedNano, uint64(f.Updated.UnixNano())) buf.Write(updatedNano) buf.Write(f.RawBytes) fileSize := int64(buf.Len()) // Проверяем, не превысит ли новый файл лимит oldSize := int64(0) if existing, ok := c.diskIndex[hash]; ok { oldSize = existing.size } // Если новый файл больше старого — нужно освободить место neededSpace := fileSize - oldSize if c.cacheMaxSize > 0 && c.currentSize+neededSpace > c.cacheMaxSize { c.evictLRU(neededSpace) } // Записываем файл if err := os.WriteFile(diskPath, buf.Bytes(), 0644); err != nil { return fmt.Errorf("failed to write disk cache for %s: %w", f.URL, err) } // Обновляем индекс if oldSize > 0 { c.currentSize -= oldSize } c.currentSize += fileSize c.diskIndex[hash] = &diskEntry{ size: fileSize, modTime: time.Now(), } return nil } // evictLRU удаляет самые старые файлы из дискового кэша, пока не освободится // нужное количество байт. func (c *S3Cache) evictLRU(neededBytes int64) { if len(c.diskIndex) == 0 { return } // Сортируем по modTime (самые старые первыми) type hashEntry struct { hash string mod time.Time size int64 } entries := make([]hashEntry, 0, len(c.diskIndex)) for hash, entry := range c.diskIndex { entries = append(entries, hashEntry{hash: hash, mod: entry.modTime, size: entry.size}) } sort.Slice(entries, func(i, j int) bool { return entries[i].mod.Before(entries[j].mod) }) freed := int64(0) for _, entry := range entries { if freed >= neededBytes { break } diskPath := filepath.Join(c.cacheDir, entry.hash) if err := os.Remove(diskPath); err == nil { freed += entry.size delete(c.diskIndex, entry.hash) c.currentSize -= entry.size } } } // updateDiskTimestamp обновляет mtime файла на диске (при 304 Not Modified). func (c *S3Cache) updateDiskTimestamp(s3URL string) { if c.cacheDir == "" { return } diskPath := c.cacheFilePath(s3URL) hash := filepath.Base(diskPath) // Обновляем mtime файла now := time.Now() if err := os.Chtimes(diskPath, now, now); err != nil { return } // Обновляем индекс if entry, ok := c.diskIndex[hash]; ok { entry.modTime = now } } // readDiskEntry возвращает размер файла на диске для LRU-расчётов. func (c *S3Cache) readDiskEntry(s3URL string) (*diskEntry, bool) { diskPath := c.cacheFilePath(s3URL) hash := filepath.Base(diskPath) entry, ok := c.diskIndex[hash] return entry, ok } // resolveEndpoint возвращает endpoint для S3-клиента. // Приоритет: 1. parsed.Endpoint (из URL), 2. s3Cfg.Endpoint (из конфига). // Учитывает s3Cfg.UseSSL для коррекции схемы. // Если задан s3Cfg.BasePath, он добавляется к endpoint. func resolveEndpoint(parsedEndpoint string, s3Cfg *S3Config) string { endpoint := parsedEndpoint if endpoint == "" && s3Cfg != nil && s3Cfg.Endpoint != "" { endpoint = s3Cfg.Endpoint } // Если endpoint не содержит схему, добавляем http:// if !strings.HasPrefix(endpoint, "http://") && !strings.HasPrefix(endpoint, "https://") { endpoint = "http://" + endpoint } // Корректируем схему в соответствии с UseSSL if s3Cfg != nil && s3Cfg.UseSSL { endpoint = strings.Replace(endpoint, "http://", "https://", 1) } // Добавляем basePath (префикс пути API, e.g. /api/v1) if s3Cfg != nil && s3Cfg.BasePath != "" { basePath := s3Cfg.BasePath if !strings.HasPrefix(basePath, "/") { basePath = "/" + basePath } // Убираем концевой слэш у endpoint, если есть endpoint = strings.TrimRight(endpoint, "/") + basePath } return endpoint } // downloadFromS3 выполняет реальную загрузку файла с S3. func (c *S3Cache) downloadFromS3(ctx context.Context, rawURL string, s3Cfg *S3Config) (*CachedFile, error) { parsed, err := ParseS3URL(rawURL) if err != nil { return nil, err } endpoint := resolveEndpoint(parsed.Endpoint, s3Cfg) client, err := c.newS3Client(ctx, endpoint, s3Cfg) if err != nil { return nil, fmt.Errorf("failed to create s3 client: %w", err) } // Скачиваем объект downloadCtx, cancel := context.WithTimeout(ctx, 30*time.Second) //nolint:mnd defer cancel() result, err := client.GetObject(downloadCtx, &s3.GetObjectInput{ Bucket: aws.String(parsed.Bucket), Key: aws.String(parsed.Key), }) if err != nil { return nil, fmt.Errorf("failed to download %s: %w", rawURL, err) } defer result.Body.Close() // Получаем ETag etag := "" if result.ETag != nil { etag = strings.Trim(*result.ETag, "\"") } // Читаем содержимое rawBytes, err := io.ReadAll(result.Body) if err != nil { return nil, fmt.Errorf("failed to read s3 object body: %w", err) } content := string(rawBytes) lines := strings.Split(content, "\n") // Извлекаем имя файла из URL (последний сегмент key) fileName := path.Base(parsed.Key) // Пробуем парсить как JSON для метода .json() var parsedJSON interface{} if err := json.Unmarshal(rawBytes, &parsedJSON); err == nil { // валидный JSON — сохраняем } else { parsedJSON = nil } return &CachedFile{ ETag: etag, URL: rawURL, FileName: fileName, RawBytes: rawBytes, Content: content, Lines: lines, ParsedJSON: parsedJSON, Updated: time.Now(), }, nil } // checkETag проверяет ETag файла на S3 с помощью If-None-Match. // - 304 Not Modified: обновляем Updated, возвращаем nil (файл не изменился) // - 200 OK с тем же ETag: сервер не поддержал If-None-Match, но файл не изменился. // Прокидываем тело (drain), обновляем Updated, возвращаем nil. // - 200 OK с новым ETag: читаем тело из этого же ответа (без второго запроса), // создаём новый CachedFile и возвращаем его. // - Если ETag пустой (первый запуск): полная перекачка через downloadFromS3. func (c *S3Cache) checkETag(ctx context.Context, rawURL string, s3Cfg *S3Config, current *CachedFile) (*CachedFile, error) { if current.ETag == "" { return c.downloadFromS3(ctx, rawURL, s3Cfg) } parsed, err := ParseS3URL(rawURL) if err != nil { return nil, err } endpoint := resolveEndpoint(parsed.Endpoint, s3Cfg) client, err := c.newS3Client(ctx, endpoint, s3Cfg) if err != nil { return nil, fmt.Errorf("failed to create s3 client: %w", err) } // GET с If-None-Match getCtx, cancel := context.WithTimeout(ctx, 10*time.Second) //nolint:mnd defer cancel() result, err := client.GetObject(getCtx, &s3.GetObjectInput{ Bucket: aws.String(parsed.Bucket), Key: aws.String(parsed.Key), IfNoneMatch: aws.String(`"` + current.ETag + `"`), }) if err != nil { // 304 Not Modified — ошибка в SDK var respErr *smithyhttp.ResponseError if errors.As(err, &respErr) && respErr.HTTPStatusCode() == http.StatusNotModified { current.Updated = time.Now() if c.cacheDir != "" { c.updateDiskTimestamp(rawURL) } return nil, nil } return nil, fmt.Errorf("etag check failed for s3://%s/%s: %w", parsed.Bucket, parsed.Key, err) } // 200 OK — читаем ETag из ответа defer result.Body.Close() serverETag := "" if result.ETag != nil { serverETag = strings.Trim(*result.ETag, "\"") } // Если ETag не изменился — сервер не поддержал If-None-Match, // но файл тот же. Дренируем тело и обновляем timestamp. if serverETag == current.ETag { // Дренируем тело, чтобы вернуть соединение в пул _, _ = io.Copy(io.Discard, result.Body) current.Updated = time.Now() if c.cacheDir != "" { c.updateDiskTimestamp(rawURL) } return nil, nil } // ETag изменился — читаем тело из этого же ответа (без второго запроса) rawBytes, err := io.ReadAll(result.Body) if err != nil { return nil, fmt.Errorf("failed to read s3 object body for %s: %w", rawURL, err) } content := string(rawBytes) lines := strings.Split(content, "\n") fileName := path.Base(parsed.Key) var parsedJSON interface{} if err := json.Unmarshal(rawBytes, &parsedJSON); err == nil { // валидный JSON } else { parsedJSON = nil } return &CachedFile{ ETag: serverETag, URL: rawURL, FileName: fileName, RawBytes: rawBytes, Content: content, Lines: lines, ParsedJSON: parsedJSON, Updated: time.Now(), }, nil } // newS3Client создает S3-клиент для заданного endpoint. func (c *S3Cache) newS3Client(ctx context.Context, endpoint string, s3Cfg *S3Config) (*s3.Client, error) { if s3Cfg == nil { return nil, fmt.Errorf("s3 config is nil: provide s3 configuration in setup_data or ensure the 's3' key exists in the setup() return value") } // Создаем custom resolver для S3-совместимых хранилищ resolver := aws.EndpointResolverWithOptionsFunc(func(service, region string, options ...interface{}) (aws.Endpoint, error) { return aws.Endpoint{ URL: endpoint, HostnameImmutable: true, }, nil }) // Для публичных бакетов используем anonymous credentials, // иначе — статические credentials из конфигурации. var credsProvider aws.CredentialsProvider if s3Cfg.AccessKey == "" && s3Cfg.SecretKey == "" { credsProvider = aws.AnonymousCredentials{} } else { credsProvider = credentials.NewStaticCredentialsProvider(s3Cfg.AccessKey, s3Cfg.SecretKey, "") } opts := []func(*config.LoadOptions) error{ config.WithRegion(s3Cfg.Region), config.WithCredentialsProvider(credsProvider), config.WithEndpointResolverWithOptions(resolver), } if s3Cfg.InsecureSkipVerify { tr := &http.Transport{ TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, //nolint:gosec } opts = append(opts, config.WithHTTPClient(&http.Client{Transport: tr})) } cfg, err := config.LoadDefaultConfig(ctx, opts...) if err != nil { return nil, err } client := s3.NewFromConfig(cfg, func(o *s3.Options) { o.UsePathStyle = true // MinIO и S3-совместимые хранилища используют path-style // Отключаем проверку контрольных сумм для ответов от самописных S3-совместимых // хранилищ, которые их не поддерживают (подавляет WARN: "Response has no supported checksum") o.ResponseChecksumValidation = aws.ResponseChecksumValidationWhenRequired }) return client, nil } // ParseS3URL парсит s3 URL. // Поддерживаются два формата: // - s3://endpoint/bucket/key (endpoint из URL) // - s3://bucket/key (endpoint пустой, заполняется вызывающим кодом из s3Cfg.Endpoint) func ParseS3URL(rawURL string) (*S3URL, error) { const prefix = "s3://" if !strings.HasPrefix(rawURL, prefix) { return nil, fmt.Errorf("invalid s3 url: must start with 's3://', got %q", rawURL) } // Убираем s3:// и разделяем trimmed := rawURL[len(prefix):] parts := strings.SplitN(trimmed, "/", 3) if len(parts) < 2 { return nil, fmt.Errorf("invalid s3 url: expected s3://bucket/key or s3://endpoint/bucket/key, got %q", rawURL) } if len(parts) == 2 { // Формат: s3://bucket/key — endpoint будет заполнен из s3Cfg.Endpoint bucket := parts[0] key := parts[1] if bucket == "" { return nil, fmt.Errorf("invalid s3 url: bucket is empty in %q", rawURL) } if key == "" { return nil, fmt.Errorf("invalid s3 url: key is empty in %q", rawURL) } return &S3URL{ Endpoint: "", // заполняется вызывающим кодом Bucket: bucket, Key: key, }, nil } // Формат: s3://endpoint/bucket/key endpointRaw := parts[0] bucket := parts[1] key := parts[2] // Если endpoint не содержит схему, добавляем http:// endpoint := endpointRaw if !strings.HasPrefix(endpoint, "http://") && !strings.HasPrefix(endpoint, "https://") { endpoint = "http://" + endpoint } // Проверяем, что URL валидный _, err := url.Parse(endpoint) if err != nil { return nil, fmt.Errorf("invalid s3 endpoint %q: %w", endpoint, err) } if bucket == "" { return nil, fmt.Errorf("invalid s3 url: bucket is empty in %q", rawURL) } if key == "" { return nil, fmt.Errorf("invalid s3 url: key is empty in %q", rawURL) } return &S3URL{ Endpoint: endpoint, Bucket: bucket, Key: key, }, nil }