Files
transcriber/main.go
T
av 3a2da3004b приём и чтение записей сведены к одному контракту приложения
- адреса приложения переехали в своё пространство `/app/`, опрос готовности
  убран целиком: рубеж и причину остановки владелец узнаёт карточкой записи,
  текст — отдельным адресом названного вида
- заведена единая точка отображения доменной ошибки и слой, приводящий к той же
  форме отказы библиотеки: тело несёт машиночитаемый код рядом с сообщением
- у записи появились имя файла отправителя, длительность и размер своими
  колонками, а у ленты владельца — свой индекс: без него страница сканировала
  весь архив сервиса
2026-08-15 13:51:23 +03:00

338 lines
14 KiB
Go

package main
import (
"context"
"errors"
"flag"
"fmt"
"log/slog"
"net/http"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"time"
ffmpegconv "git.vakhrushev.me/av/transcriber/internal/adapter/converter/ffmpeg"
ffmpegmv "git.vakhrushev.me/av/transcriber/internal/adapter/metaviewer/ffmpeg"
"git.vakhrushev.me/av/transcriber/internal/adapter/recognizer/yandex"
pbrepo "git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase"
"git.vakhrushev.me/av/transcriber/internal/config"
httpcontroller "git.vakhrushev.me/av/transcriber/internal/controller/http"
"git.vakhrushev.me/av/transcriber/internal/controller/worker"
"git.vakhrushev.me/av/transcriber/internal/metrics"
"git.vakhrushev.me/av/transcriber/internal/service"
"github.com/joho/godotenv"
"github.com/pocketbase/pocketbase/apis"
"github.com/pocketbase/pocketbase/core"
"github.com/prometheus/client_golang/prometheus/promhttp"
"git.vakhrushev.me/av/transcriber/internal/clock"
)
func main() {
// Создаем структурированный логгер
logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{
Level: slog.LevelInfo,
}))
slog.SetDefault(logger)
// Parse command line flags
configPath := flag.String("c", "config.toml", "Path to config file")
flag.StringVar(configPath, "config", "config.toml", "Path to config file (alias for -c)")
flag.Parse()
// Load configuration
cfg, err := config.LoadConfig(*configPath)
if err != nil {
logger.Error("Unable to load configuration", "config_path", *configPath, "error", err)
os.Exit(1)
} else {
logger.Info("Configuration loaded successfully", "config_path", *configPath)
}
// Незаполненный вход роняет старт: подняться с молча выключенным входом
// значит остаться открытым наружу, и узнать об этом было бы неоткуда.
if err := cfg.Auth.Validate(); err != nil {
logger.Error("Unable to start with incomplete login settings", "error", err)
os.Exit(1)
}
// Числа конвейера проверяются здесь же: ноль воркеров — объявленный режим, а
// отрицательное число и нулевой предел простоя — опечатка, и подниматься с
// ней значит остановить всякую запись первым же захватом.
if err := cfg.Pipeline.Validate(); err != nil {
logger.Error("Unable to start with incorrect pipeline settings", "error", err)
os.Exit(1)
}
// Загружаем переменные окружения из .env файла
if err := godotenv.Load(); err != nil {
logger.Warn("Warning: .env file not found, using system environment variables")
}
// Хранилище поднимается библиотекой, а не её набором команд: разбор флагов
// и мягкая остановка остаются нашими. Схему накатывает Serve — он гоняет
// непринятые шаги прежде, чем поднять сервер.
storage, err := pbrepo.New(cfg.Storage.DataDir)
if err != nil {
logger.Error("Failed to open storage", "error", err)
os.Exit(1)
}
defer func() {
if err := storage.ResetBootstrapState(); err != nil {
logger.Error("Failed to close storage", "error", err)
}
}()
pbrepo.BindPanelRules(storage)
// Создаем репозитории
recordRepo := pbrepo.NewAudioRecordRepository(storage)
repos := service.Repositories{
Records: recordRepo,
Files: pbrepo.NewFileRepository(storage),
Texts: pbrepo.NewTextRepository(storage),
Structures: pbrepo.NewStructureRepository(storage),
Recognitions: pbrepo.NewRecognitionRepository(storage),
Events: pbrepo.NewRecordEventRepository(storage),
}
// Создаем адаптеры
metaviewer := ffmpegmv.NewFfmpegMetaViewer()
converter := ffmpegconv.NewFfmpegConverter()
recognizer, err := yandex.NewYandexAudioRecognizerService(yandex.YandexAudioRecognizerConfig{
Region: cfg.Yandex.ObjStorageRegion,
AccessKey: cfg.Yandex.ObjStorageAccessKey,
SecretKey: cfg.Yandex.ObjStorageSecretKey,
BucketName: cfg.Yandex.ObjStorageBucketName,
Endpoint: cfg.Yandex.ObjStorageEndpoint,
ApiKey: cfg.Yandex.SpeechKitAPIKey,
FolderID: cfg.Yandex.FolderID,
})
if err != nil {
logger.Error("failed to create audio recognizer", "error", err)
os.Exit(1)
}
// Отдавать отказ закрытия некому — процесс заканчивается, — поэтому он идёт
// в журнал владельца. Что он означает: gRPC-клиент отдаёт здесь отказ лишь
// при повторном закрытии, то есть запись говорит о нашей ошибке, а не о
// недоступности Yandex.
defer func() {
if err := recognizer.Close(); err != nil {
logger.Error("failed to close audio recognizer", "error", err)
}
}()
// Создаем сервисы
transcribeService := service.NewTranscribeService(
repos,
metaviewer,
converter,
recognizer,
cfg.Pipeline.StuckLimits(),
logger,
)
// Создаем контекст для graceful shutdown
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
// Создаем WaitGroup для ожидания завершения всех воркеров
var wg sync.WaitGroup
// Пул одинаковых воркеров: специализации у них нет, шаг выбирается по рубежу
// самой записи. Число приходит настройкой, ноль — законное значение.
pool := worker.NewPool(cfg.Pipeline.Workers, transcribeService.RunStep, logger)
wg.Add(1)
go func() {
defer wg.Done()
pool.Start(ctx)
}()
// Вход у сервиса один — приём по HTTP, — и метка ставится только ему. Метки
// убранного входа Telegram здесь нет намеренно: ноль читался бы как поломка,
// а признак существует ради того дня, когда входов снова станет больше.
metrics.IntakeUpGauge.WithLabelValues("http").Set(1)
// Наши маршруты живут на роутере хранилища: панель отдаётся тем же портом,
// и второму серверу на нём взяться неоткуда.
appHandler := httpcontroller.NewAppHandler(recordRepo, repos.Texts, repos.Structures, transcribeService, logger)
authHandler := httpcontroller.NewAuthHandler(storage, httpcontroller.AuthHandlerConfig{
AuthURL: cfg.Auth.AuthURL,
RedirectURL: cfg.Auth.RedirectURL,
ClientID: cfg.Auth.ClientID,
SecureCookie: cfg.Auth.SecureCookie,
}, logger)
// Сервер приезжает каналом, а не общей переменной: хук исполняется в
// горутине сервера, а читает его горутина остановки, и связи «произошло
// раньше» между ними иначе нет.
srvCh := make(chan *http.Server, 1)
storage.OnServe().BindFunc(func(se *core.ServeEvent) error {
// Шесть часов записи по медленному каналу переживают любой фиксированный
// таймаут чтения, а умолчание хранилища — пять минут. Стойкость к
// целенаправленной нагрузке объявлена вне модели угроз проекта.
se.Server.ReadTimeout = 0
srvCh <- se.Server
// Журнал входящих запросов вернулся своим слоем: вместе с gin ушёл
// `sloggin`, а хранилище пишет запросы в свою таблицу, которой в
// журнале контейнера не видно. Поля — те, что просит конвенция.
se.Router.BindFunc(func(e *core.RequestEvent) error {
start := clock.Start()
err := e.Next()
level := slog.LevelInfo
if e.Request.URL.Path == "/health" || e.Request.URL.Path == "/metrics" {
// Опрос здоровья и метрик идёт постоянно и полезного не несёт.
level = slog.LevelDebug
}
logger.Log(e.Request.Context(), level, "Incoming request",
"http.method", e.Request.Method,
"http.route", journalRoute(e.Request.URL.Path),
"http.status_code", e.Status(),
"duration_ms", time.Since(start).Milliseconds(),
"transport", "http")
return err
})
// Настройки провайдера приводятся к конфигу при каждом подъёме:
// применённый шаг схемы не переписывается, и секрет, положенный
// однажды шагом, не пережил бы ротации.
if err := pbrepo.ApplyProviderSettings(storage, pbrepo.ProviderSettings{
AuthURL: cfg.Auth.AuthURL,
TokenURL: cfg.Auth.TokenURL,
UserInfoURL: cfg.Auth.UserInfoURL,
ClientID: cfg.Auth.ClientID,
ClientSecret: cfg.Auth.ClientSecret,
}); err != nil {
return fmt.Errorf("failed to apply provider settings: %w", err)
}
// Своё правило ограничителя частоты под корень приложения. Правило
// хранилища настроено на его собственный корень и наших адресов больше не
// покрывает: вместе с переездом приложения ограничитель перестал бы
// существовать для него вовсе, и заметить это было бы нечем.
if err := httpcontroller.ApplyAppRateLimit(storage); err != nil {
return fmt.Errorf("failed to apply app rate limit: %w", err)
}
authHandler.Register(se.Router)
appHandler.Register(se.Router)
se.Router.GET("/health", func(e *core.RequestEvent) error {
return e.JSON(http.StatusOK, map[string]string{
"status": "ok",
"message": "Transcriber service is running",
})
})
se.Router.GET("/metrics", func(e *core.RequestEvent) error {
promhttp.Handler().ServeHTTP(e.Response, e.Request)
return nil
})
return se.Next()
})
// Запускаем HTTP сервер в отдельной горутине
serveErr := make(chan error, 1)
wg.Add(1)
go func() {
defer wg.Done()
logger.Info("Starting HTTP server", "port", cfg.Server.Port)
err := apis.Serve(storage, apis.ServeConfig{
HttpAddr: fmt.Sprintf(":%d", cfg.Server.Port),
ShowStartBanner: false,
})
if err != nil && !errors.Is(err, http.ErrServerClosed) {
logger.Error("HTTP server error", "error", err)
serveErr <- err
}
}()
// Настраиваем обработку сигналов для graceful shutdown
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
logger.Info("Transcriber service started", "pipeline_workers", pool.Size())
logger.Info("Press Ctrl+C to stop...")
// Ждем сигнал завершения либо отказ сервера
select {
case <-sigChan:
logger.Info("Received shutdown signal, initiating graceful shutdown...")
case <-serveErr:
logger.Error("HTTP server stopped unexpectedly, shutting down")
}
// Создаем контекст с таймаутом для graceful shutdown HTTP сервера
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), time.Duration(cfg.Server.ShutdownTimeout)*time.Second)
defer shutdownCancel()
// Останавливаем HTTP сервер
select {
case srv := <-srvCh:
logger.Info("Shutting down HTTP server...")
if err := srv.Shutdown(shutdownCtx); err != nil {
logger.Error("HTTP server forced to shutdown", "error", err)
} else {
logger.Info("HTTP server stopped gracefully")
}
default:
logger.Info("HTTP server was not started, nothing to shut down")
}
// Отменяем контекст для остановки воркеров
cancel()
// Создаем канал для уведомления о завершении всех воркеров
done := make(chan struct{})
go func() {
wg.Wait()
close(done)
}()
// Ждем завершения всех воркеров или таймаута
select {
case <-done:
logger.Info("All workers stopped gracefully")
case <-time.After(time.Duration(cfg.Server.ForceShutdownTimeout) * time.Second):
logger.Warn("Timeout reached, forcing shutdown")
}
logger.Info("Transcriber service stopped")
}
// filesPathPrefix — начало пути, которым хранилище отдаёт файл записи. Последний
// сегмент такого пути и есть имя файла в хранилище.
const filesPathPrefix = "/api/files/"
// journalRoute готовит путь запроса к записи в журнал.
//
// Инвариант проекта запрещает имени файла в хранилище попадать в журнал: имя —
// последняя часть ссылки `/api/files/...`, и строка журнала вместе с
// идентификатором записи собирала бы ссылку целиком. Слой журнала пишет путь
// всякого запроса, поэтому имя срезается здесь — иначе оно уезжало бы в
// собранные логи при каждом скачивании записи.
//
// Срезается только имя: маршрут остаётся различимым, и наблюдаемость от этого не
// теряется.
func journalRoute(path string) string {
if !strings.HasPrefix(path, filesPathPrefix) {
return path
}
cut := strings.LastIndex(path, "/")
if cut < len(filesPathPrefix) {
return path
}
return path[:cut+1] + "<имя>"
}