точки входа переехали в cmd/, заведена заглушка OIDC для локального входа
- main.go и journal_route_test.go переехали в cmd/transcriber без правок содержимого; образ собирает ./cmd/transcriber поимённо - cmd/oidcstub отвечает на /authorize, /token и /userinfo, проверок не делает и слушает петлевой адрес: войти без Authelia стало чем - ступень сборки приложения переехала на node:24 с alpine — musl ждёт ответа на AAAA, которого нет, и npm ci висел вместо отказа
This commit is contained in:
@@ -0,0 +1,171 @@
|
||||
// Команда oidcstub — подставной провайдер OIDC для локального запуска.
|
||||
//
|
||||
// Сервис держит вход через внешнего провайдера, и без него в приложение не
|
||||
// попасть вовсе: секция `[auth]` обязательна, а выдуманные адреса разбираются
|
||||
// как ссылки, но никуда не ведут. Настоящая Authelia — отдельный сервис, который
|
||||
// надо развернуть и настроить, а ключи боевого провайдера на машине разработчика
|
||||
// лежать не должны.
|
||||
//
|
||||
// Заглушка отвечает на `/authorize`, `/token` и `/userinfo`, и этого хватает:
|
||||
// `user_info_url` — обязательный ключ конфига, а при заполненном адресе сведений
|
||||
// библиотека хранилища берёт их обычным запросом с предъявленным токеном и
|
||||
// `id_token` не смотрит вовсе. Поэтому здесь нет ни ключей подписи, ни документа
|
||||
// обнаружения.
|
||||
//
|
||||
// Проверок она не делает никаких — ни секрета клиента, ни проверочного кода
|
||||
// PKCE, ни выданного токена. Вход у неё один и заранее известный: она нужна,
|
||||
// чтобы дойти до куки сессии, а не чтобы изображать провайдера. По той же
|
||||
// причине слушает она только петлевой адрес.
|
||||
//
|
||||
// Запуск:
|
||||
//
|
||||
// go run ./cmd/oidcstub
|
||||
// go run ./cmd/oidcstub -sub local-2
|
||||
//
|
||||
// Другое значение `-sub` заводит второго вошедшего, и почта переезжает вместе с
|
||||
// ним: без своей почты второй вход отвергается вовсе. Хранилище ищет учётную
|
||||
// запись сперва по признаку провайдера, а не найдя — по адресу почты, и на
|
||||
// найденной записи признак уникален. Второй `sub` при общей почте пришёл бы к той
|
||||
// же записи и получил бы отказ уникальности — вход отвечал бы 401, а причина
|
||||
// осталась бы в журнале хранилища строкой про связь.
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// authCode — код, который заглушка отдаёт на возврате. Значение постоянное:
|
||||
// одноразовость кода держит провайдер, а здесь её изображать не для кого.
|
||||
const authCode = "local-code"
|
||||
|
||||
// accessToken — то, что уезжает в обмен и приходит обратно заголовком запроса
|
||||
// сведений о вошедшем. Заглушка его не сверяет.
|
||||
const accessToken = "local-token"
|
||||
|
||||
// readHeaderTimeout — потолок чтения заголовков. Заглушка стоит на петлевом
|
||||
// адресе, но сервер без единого срока держит подвисшее соединение вечно.
|
||||
const readHeaderTimeout = 5 * time.Second
|
||||
|
||||
func main() {
|
||||
port := flag.Int("port", 9000, "Порт заглушки")
|
||||
sub := flag.String("sub", "local-1", "Идентификатор вошедшего: разные значения дают разных владельцев записей")
|
||||
name := flag.String("name", "Локальный", "Имя вошедшего")
|
||||
email := flag.String("email", "", "Почта вошедшего; пустое значение даёт <sub>@example.com")
|
||||
flag.Parse()
|
||||
|
||||
// Почта выводится из `sub`, а не стоит своим умолчанием: общая почта у двух
|
||||
// разных `sub` приводит второй вход к первой записи, а признак провайдера на
|
||||
// ней уникален — второй вошедший получал бы 401 вместо своей учётной записи.
|
||||
// Обещание «разные `-sub` дают разных владельцев» держится этой строкой.
|
||||
if *email == "" {
|
||||
*email = *sub + "@example.com"
|
||||
}
|
||||
|
||||
logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{
|
||||
Level: slog.LevelInfo,
|
||||
}))
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/authorize", authorize(logger))
|
||||
mux.HandleFunc("/token", token(logger))
|
||||
mux.HandleFunc("/userinfo", userInfo(logger, *sub, *name, *email))
|
||||
|
||||
// Петлевой адрес, а не все — заглушка пускает кого угодно кем угодно, и в
|
||||
// чужой сети это дыра, а не удобство.
|
||||
addr := fmt.Sprintf("127.0.0.1:%d", *port)
|
||||
server := &http.Server{
|
||||
Addr: addr,
|
||||
Handler: mux,
|
||||
ReadHeaderTimeout: readHeaderTimeout,
|
||||
}
|
||||
|
||||
logger.Info("oidc stub started", "addr", addr, "sub", *sub)
|
||||
|
||||
if err := server.ListenAndServe(); err != nil {
|
||||
logger.Error("oidc stub stopped", "error", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
// authorize уводит браузер обратно на адрес возврата с готовым кодом.
|
||||
//
|
||||
// Настоящий провайдер спросил бы здесь имя и пароль; заглушка не спрашивает
|
||||
// ничего и возвращает сразу — в этом вся её работа.
|
||||
func authorize(logger *slog.Logger) http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
query := r.URL.Query()
|
||||
|
||||
target, err := url.Parse(query.Get("redirect_uri"))
|
||||
if err != nil || target.Host == "" {
|
||||
logger.Error("authorize rejected", "reason", "redirect_uri is missing or malformed")
|
||||
http.Error(w, "redirect_uri не задан или не разбирается", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
// Состояние возвращается тем же значением, каким пришло: сервис сверяет
|
||||
// его со своей кукой и без совпадения отвергает возврат.
|
||||
back := target.Query()
|
||||
back.Set("code", authCode)
|
||||
if state := query.Get("state"); state != "" {
|
||||
back.Set("state", state)
|
||||
}
|
||||
target.RawQuery = back.Encode()
|
||||
|
||||
logger.Info("authorize passed", "redirect_host", target.Host)
|
||||
|
||||
http.Redirect(w, r, target.String(), http.StatusFound)
|
||||
}
|
||||
}
|
||||
|
||||
// token отдаёт токен в обмен на код. Ни код, ни секрет клиента, ни проверочный
|
||||
// код PKCE не сверяются.
|
||||
func token(logger *slog.Logger) http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
http.Error(w, "обмен идёт методом POST", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
|
||||
// `id_token` не отдаётся намеренно: подписать его нечем, а библиотека
|
||||
// хранилища его и не смотрит — `user_info_url` стоит в перечне
|
||||
// обязательных ключей конфига, и без него сервис не поднимается вовсе.
|
||||
writeJSON(logger, w, map[string]any{
|
||||
"access_token": accessToken,
|
||||
"token_type": "Bearer",
|
||||
"expires_in": 3600,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// userInfo отдаёт сведения о вошедшем.
|
||||
//
|
||||
// `email_verified` обязано быть истинным: без него библиотека адрес почты в
|
||||
// учётную запись не запишет, и запись заведётся без него.
|
||||
func userInfo(logger *slog.Logger, sub, name, email string) http.HandlerFunc {
|
||||
return func(w http.ResponseWriter, _ *http.Request) {
|
||||
writeJSON(logger, w, map[string]any{
|
||||
"sub": sub,
|
||||
"name": name,
|
||||
"preferred_username": sub,
|
||||
"email": email,
|
||||
"email_verified": true,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// writeJSON отвечает разметкой JSON. Отказ записи идёт в журнал: ответ к этому
|
||||
// моменту уже начат, и сказать о нём спрашивающему нечем.
|
||||
func writeJSON(logger *slog.Logger, w http.ResponseWriter, body map[string]any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
|
||||
if err := json.NewEncoder(w).Encode(body); err != nil {
|
||||
logger.Error("response write failed", "error", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,93 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
|
||||
httpcontroller "git.vakhrushev.me/av/transcriber/internal/controller/http"
|
||||
)
|
||||
|
||||
// Перечень адресного пространства для проверок журнала. Обработчиков он здесь
|
||||
// не вешает: журналу нужны только границы, а не то, что стоит за ними.
|
||||
func journalMounts() []httpcontroller.Mount {
|
||||
return []httpcontroller.Mount{
|
||||
{Path: httpcontroller.StorageRoot},
|
||||
{Path: httpcontroller.PanelRoot},
|
||||
{Path: httpcontroller.AppRoot},
|
||||
{Path: httpcontroller.AuthRoot},
|
||||
{Path: httpcontroller.HealthPath, Exact: true},
|
||||
{Path: httpcontroller.MetricsPath, Exact: true},
|
||||
}
|
||||
}
|
||||
|
||||
// Имя файла в хранилище в журнал не идёт: оно последняя часть ссылки
|
||||
// `/api/files/...`, и строка журнала вместе с идентификатором записи собрала бы
|
||||
// ссылку целиком. Инвариант проекта, critical.
|
||||
func TestJournalRouteHidesStoredFileName(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
path string
|
||||
want string
|
||||
}{
|
||||
{
|
||||
name: "ссылка на файл теряет имя",
|
||||
path: "/api/files/files/abc123def456ghi/9f1c-3b2a.mp3",
|
||||
want: "/api/files/files/abc123def456ghi/<имя>",
|
||||
},
|
||||
{
|
||||
name: "маршрут остаётся различимым",
|
||||
path: "/api/files/files/abc123def456ghi/запись.ogg",
|
||||
want: "/api/files/files/abc123def456ghi/<имя>",
|
||||
},
|
||||
{
|
||||
name: "прочие пути не трогаются",
|
||||
path: "/app/audiorecords/abc123def456ghi",
|
||||
want: "/app/audiorecords/abc123def456ghi",
|
||||
},
|
||||
{
|
||||
name: "приём не трогается",
|
||||
path: "/app/audiorecords",
|
||||
want: "/app/audiorecords",
|
||||
},
|
||||
{
|
||||
name: "сам префикс без имени не портится",
|
||||
path: "/api/files/",
|
||||
want: "/api/files/",
|
||||
},
|
||||
}
|
||||
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
assert.Equal(t, c.want, journalRoute(c.path, journalMounts()))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Отдельно и прямо: имени в готовой строке нет. Проверка судит результат, а не
|
||||
// устройство — переписанная реализация обязана остаться зелёной.
|
||||
func TestJournalRouteDropsNameEntirely(t *testing.T) {
|
||||
const stored = "0f7b8dd3-d1cc-424c.mp3"
|
||||
|
||||
route := journalRoute("/api/files/files/rec0000000000000/"+stored, journalMounts())
|
||||
|
||||
assert.NotContains(t, route, stored, "имя файла в хранилище не доезжает до журнала")
|
||||
assert.Contains(t, route, "rec0000000000000", "идентификатор записи остаётся: по нему прослеживается путь")
|
||||
}
|
||||
|
||||
// Путь, не принадлежащий сервису, уходит приложению и в журнал дословно не
|
||||
// идёт: множеством его значений распоряжается спрашивающий.
|
||||
func TestJournalRouteHidesWebappPath(t *testing.T) {
|
||||
cases := []string{
|
||||
"/",
|
||||
"/records/abc123def456ghi",
|
||||
"/" + strings.Repeat("a", 1024),
|
||||
}
|
||||
|
||||
for _, path := range cases {
|
||||
route := journalRoute(path, journalMounts())
|
||||
|
||||
assert.Equal(t, webappRoute, route)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,362 @@
|
||||
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"
|
||||
"git.vakhrushev.me/av/transcriber/web"
|
||||
"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)
|
||||
|
||||
// Адресное пространство сервиса объявлено одним перечнем, и он порождает
|
||||
// регистрацию, а не описывает её: корень, заведённый мимо перечня, не
|
||||
// получит обработчика вовсе. Отсюда же уровень журнала для адресов
|
||||
// наблюдения и правило неизвестного пути у раздачи приложения.
|
||||
mounts := httpcontroller.ServiceMounts(appHandler, authHandler, promhttp.Handler())
|
||||
|
||||
dist, appBuilt := web.Dist()
|
||||
webappHandler := httpcontroller.NewWebappHandler(dist, appBuilt, mounts, 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 httpcontroller.IsObservationAddress(mounts, e.Request.URL.Path) {
|
||||
// Опрос здоровья и метрик идёт постоянно и полезного не несёт.
|
||||
level = slog.LevelDebug
|
||||
}
|
||||
|
||||
attrs := []any{
|
||||
"http.method", e.Request.Method,
|
||||
"http.route", journalRoute(e.Request.URL.Path, mounts),
|
||||
"http.status_code", e.Status(),
|
||||
"duration_ms", time.Since(start).Milliseconds(),
|
||||
"transport", "http",
|
||||
}
|
||||
|
||||
// Путь, отданный приложению, в журнал не идёт — вместо него исход
|
||||
// и длина: по ним видно, что происходит, а множеством значений
|
||||
// самого пути распоряжается спрашивающий.
|
||||
if outcome := httpcontroller.WebappOutcome(e); outcome != "" {
|
||||
attrs = append(attrs,
|
||||
"webapp.outcome", outcome,
|
||||
"http.path_length", len(e.Request.URL.Path))
|
||||
}
|
||||
|
||||
logger.Log(e.Request.Context(), level, "Incoming request", attrs...)
|
||||
|
||||
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)
|
||||
}
|
||||
|
||||
httpcontroller.RegisterServiceRoutes(se.Router, mounts)
|
||||
|
||||
// Раздача приложения вешается последней: она занимает корень, и всё,
|
||||
// что не совпало ни с одним адресом сервиса, доходит до неё.
|
||||
webappHandler.Register(se.Router)
|
||||
|
||||
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/"
|
||||
|
||||
// webappRoute — чем в журнале обозначается всякий путь, отданный приложению.
|
||||
const webappRoute = "<приложение>"
|
||||
|
||||
// journalRoute готовит путь запроса к записи в журнал.
|
||||
//
|
||||
// Инвариант проекта запрещает имени файла в хранилище попадать в журнал: имя —
|
||||
// последняя часть ссылки `/api/files/...`, и строка журнала вместе с
|
||||
// идентификатором записи собирала бы ссылку целиком. Слой журнала пишет путь
|
||||
// всякого запроса, поэтому имя срезается здесь — иначе оно уезжало бы в
|
||||
// собранные логи при каждом скачивании записи.
|
||||
//
|
||||
// Срезается только имя: маршрут остаётся различимым, и наблюдаемость от этого не
|
||||
// теряется.
|
||||
//
|
||||
// Путь, не принадлежащий сервису, в журнал не идёт вовсе. До появления раздачи
|
||||
// приложения такой путь ловил отказ маршрутизатора, а теперь получает разметку
|
||||
// с кодом `200`: множеством его значений распоряжается спрашивающий, и
|
||||
// дословная строка сделала бы журнал местом, куда аноним пишет свой текст.
|
||||
func journalRoute(path string, mounts []httpcontroller.Mount) string {
|
||||
if !httpcontroller.IsServiceAddress(mounts, path) {
|
||||
return webappRoute
|
||||
}
|
||||
|
||||
if !strings.HasPrefix(path, filesPathPrefix) {
|
||||
return path
|
||||
}
|
||||
|
||||
cut := strings.LastIndex(path, "/")
|
||||
if cut < len(filesPathPrefix) {
|
||||
return path
|
||||
}
|
||||
|
||||
return path[:cut+1] + "<имя>"
|
||||
}
|
||||
Reference in New Issue
Block a user