Голосовой центр L2: прокси транскрибации POST /v1/ai/transcribe
CI / test (push) Successful in 7m45s
CI / deploy-ift (push) Successful in 2m53s
CI / e2e-ift (push) Successful in 1m23s
CI / deploy-stage (push) Successful in 2m9s
CI / e2e-stage (push) Successful in 1m50s

Back-прокси STT для press-to-talk (канон PRODUCT-AI, голос и тап —
два канала одного движка действий):

- handler_ai_transcribe + logic_ai_transcribe (паттерн logic_geo):
  адаптеры провайдеров xai и openai-совместимый; env STT_PROVIDER,
  STT_URL, STT_API_KEY, STT_MODEL, STT_TIMEOUT_MS, STT_PROXY_URL
- multipart read_file_part в handler_upload_utils (лимит 10 MB)
- ETS rate-limit: 10 req/min на пользователя + глобальный троттлинг
  апстрима; без ключа — 503 unavailable, ошибки апстрима — 502/503
- аудио и транскрипт не персистятся (приват-инвариант канона)
- отдельный httpc-профиль: STT_PROXY_URL применяется только к STT
  (эгресс stt-egress для гео-блокированных провайдеров)
- trails в eventhub_trails.erl, маршрут в eventhub_app.erl,
  docker/.env.example дополнен STT_*

Проверка: eunit модуля 8/8; полный прогон eunit зелёный (1211 тестов),
compile без новых warning'ов.

Refs EventHub/EventHubBack#79
Refs EventHub/EventHubFront#76
This commit is contained in:
2026-08-20 23:56:35 +03:00
parent f4d89a262c
commit 1ef9020c40
7 changed files with 547 additions and 20 deletions
+30 -18
View File
@@ -1,28 +1,28 @@
# Скопируйте в .env (файл .env в .gitignore). Для локальной разработки: EVENTHUB_ENV=dev # 小泻芯锌懈褉褍泄褌械 胁 .env (褎邪泄谢 .env .gitignore). 袛谢褟 谢芯泻邪谢褜薪芯泄 褉邪蟹褉邪斜芯褌泻懈: EVENTHUB_ENV=dev
# Режим: dev | stage | prod (stage/prod — обязательные сильные секреты) # 袪械卸懈屑: dev | stage | prod (stage/prod 芯斜褟蟹邪褌械谢褜薪褘械 褋懈谢褜薪褘械 褋械泻褉械褌褘)
EVENTHUB_ENV=dev EVENTHUB_ENV=dev
RELEASE_COOKIE=ваш-очень-длинный-секретный-куки-минимум-32-символа RELEASE_COOKIE=胁邪褕-芯褔械薪褜-写谢懈薪薪褘泄-褋械泻褉械褌薪褘泄-泻褍泻懈-屑懈薪懈屑褍屑-32-褋懈屑胁芯谢邪
GRAFANA_ADMIN_PASSWORD=сложный-уникальный-пароль GRAFANA_ADMIN_PASSWORD=褋谢芯卸薪褘泄-褍薪懈泻邪谢褜薪褘泄-锌邪褉芯谢褜
# JWT (минимум 32 символа, не из denylist — см. infra_secrets) # JWT (屑懈薪懈屑褍屑 32 褋懈屑胁芯谢邪, 薪械 懈蟹 denylist 褋屑. infra_secrets)
JWT_SECRET=замените-на-случайную-строку-минимум-32-символа-для-user-jwt JWT_SECRET=蟹邪屑械薪懈褌械-薪邪-褋谢褍褔邪泄薪褍褞-褋褌褉芯泻褍-屑懈薪懈屑褍屑-32-褋懈屑胁芯谢邪-写谢褟-user-jwt
ADMIN_JWT_SECRET=замените-на-другую-случайную-строку-минимум-32-символа ADMIN_JWT_SECRET=蟹邪屑械薪懈褌械-薪邪-写褉褍谐褍褞-褋谢褍褔邪泄薪褍褞-褋褌褉芯泻褍-屑懈薪懈屑褍屑-32-褋懈屑胁芯谢邪
# Seed админов (обязательны в stage/prod при пустой таблице admin; в dev — опционально) # Seed 邪写屑懈薪芯胁 (芯斜褟蟹邪褌械谢褜薪褘 胁 stage/prod 锌褉懈 锌褍褋褌芯泄 褌邪斜谢懈褑械 admin; dev 芯锌褑懈芯薪邪谢褜薪芯)
ADMIN_SUPER_EMAIL=superadmin@eventhub.local ADMIN_SUPER_EMAIL=superadmin@eventhub.local
ADMIN_SUPER_PASSWORD=замените-сильный-пароль-суперадмина ADMIN_SUPER_PASSWORD=蟹邪屑械薪懈褌械-褋懈谢褜薪褘泄-锌邪褉芯谢褜-褋褍锌械褉邪写屑懈薪邪
ADMIN_EMAIL=admin@eventhub.local ADMIN_EMAIL=admin@eventhub.local
ADMIN_PASSWORD=замените-сильный-пароль-админа ADMIN_PASSWORD=蟹邪屑械薪懈褌械-褋懈谢褜薪褘泄-锌邪褉芯谢褜-邪写屑懈薪邪
ADMIN_MODER_EMAIL=moderator@eventhub.local ADMIN_MODER_EMAIL=moderator@eventhub.local
ADMIN_MODER_PASSWORD=замените-сильный-пароль-модератора ADMIN_MODER_PASSWORD=蟹邪屑械薪懈褌械-褋懈谢褜薪褘泄-锌邪褉芯谢褜-屑芯写械褉邪褌芯褉邪
ADMIN_SUPPORT_EMAIL=support@eventhub.local ADMIN_SUPPORT_EMAIL=support@eventhub.local
ADMIN_SUPPORT_PASSWORD=замените-сильный-пароль-поддержки ADMIN_SUPPORT_PASSWORD=蟹邪屑械薪懈褌械-褋懈谢褜薪褘泄-锌邪褉芯谢褜-锌芯写写械褉卸泻懈
# Email: EMAIL_TRANSPORT=auto|smtp|http_api|log (auto: API key http_api, else SMTP_HOST smtp, else log) # Email: EMAIL_TRANSPORT=auto|smtp|http_api|log (auto: API key http_api, else SMTP_HOST smtp, else log)
EMAIL_TRANSPORT=auto EMAIL_TRANSPORT=auto
# HTTP API (HTTPS when outbound SMTP ports blocked). Canon: Resend (Brevo optional). # HTTP API (HTTPS when outbound SMTP ports blocked). Canon: Resend (Brevo optional).
EMAIL_API_KEY= EMAIL_API_KEY=
EMAIL_API_PROVIDER=resend EMAIL_API_PROVIDER=resend
# EMAIL_API_URL=https://api.resend.com/emails # EMAIL_API_URL=https://api.resend.com/emails
@@ -30,7 +30,7 @@ EMAIL_API_PROVIDER=resend
VAPID_PUBLIC_KEY= VAPID_PUBLIC_KEY=
VAPID_PRIVATE_KEY= VAPID_PRIVATE_KEY=
VAPID_SUBJECT=mailto:admin@calentiq.com VAPID_SUBJECT=mailto:admin@calentiq.com
# SMTP (пусто SMTP_HOST + нет API key = только лог; пароль не коммитить) # SMTP (锌褍褋褌芯 SMTP_HOST + 薪械褌 API key = 褌芯谢褜泻芯 谢芯谐; 锌邪褉芯谢褜 薪械 泻芯屑屑懈褌懈褌褜)
SMTP_HOST= SMTP_HOST=
SMTP_PORT=587 SMTP_PORT=587
SMTP_USER= SMTP_USER=
@@ -38,12 +38,24 @@ SMTP_PASS=
SMTP_FROM=noreply@calentiq.com SMTP_FROM=noreply@calentiq.com
SMTP_TLS=if_available SMTP_TLS=if_available
PUBLIC_APP_URL=https://stage.calentiq.com PUBLIC_APP_URL=https://stage.calentiq.com
# Окно напоминаний о booking (часы до старта; Back#70) # 袨泻薪芯 薪邪锌芯屑懈薪邪薪懈泄 芯 booking (褔邪褋褘 写芯 褋褌邪褉褌邪; Back#70)
REMINDER_LEAD_HOURS=24 REMINDER_LEAD_HOURS=24
# Uploads (avatar/cover) — каталог на volume /app/data (Back#71) # Uploads (avatar/cover) 泻邪褌邪谢芯谐 薪邪 volume /app/data (Back#71)
UPLOAD_DIR=/app/data/uploads UPLOAD_DIR=/app/data/uploads
UPLOAD_MAX_BYTES=2097152 UPLOAD_MAX_BYTES=2097152
# OpenCage (DevOps#16). = /v1/geo -> 503 (degrade). Free: 2500/, 1 rps. # OpenCage (DevOps#16). 象耱铋 觌 = /v1/geo -> 503 (degrade). Free: 2500/耋蜿, 1 rps.
OPENCAGE_API_KEY= OPENCAGE_API_KEY=
OPENCAGE_TIMEOUT_MS=3000 OPENCAGE_TIMEOUT_MS=3000
# STT-прокси голосового центра (Back#79, Front#76). Без STT_API_KEY
# /v1/ai/transcribe отвечает 503 unavailable, клиент деградирует в палитру.
# Провайдеры: xai (POST {STT_URL}, Bearer) | openai (OpenAI-совместимый).
STT_PROVIDER=xai
STT_URL=https://api.x.ai/v1/stt
STT_API_KEY=
# Опционально: модель для openai-провайдера и таймаут апстрима.
STT_MODEL=
STT_TIMEOUT_MS=15000
# Опционально: HTTP-эгресс (сервис stt-egress в стеке, vless-образ);
# гео-блокированные STT-провайдеры вызываются через него.
STT_PROXY_URL=
+1
View File
@@ -112,6 +112,7 @@ start_http() ->
{"/v1/geo/geocode", handler_geo, []}, {"/v1/geo/geocode", handler_geo, []},
{"/v1/geo/reverse", handler_geo, []}, {"/v1/geo/reverse", handler_geo, []},
{"/v1/ai/hint-metrics", handler_ai_metrics, []}, {"/v1/ai/hint-metrics", handler_ai_metrics, []},
{"/v1/ai/transcribe", handler_ai_transcribe, []},
{"/v1/calendars", handler_calendars, []}, {"/v1/calendars", handler_calendars, []},
{"/v1/calendars/:id", handler_calendar_by_id, []}, {"/v1/calendars/:id", handler_calendar_by_id, []},
{"/v1/calendars/:id/cover", handler_calendar_cover, []}, {"/v1/calendars/:id/cover", handler_calendar_cover, []},
+68
View File
@@ -0,0 +1,68 @@
%%%-------------------------------------------------------------------
%%% @doc POST /v1/ai/transcribe — прокси транскрибации голоса
%%% (PRODUCT-AI, Front#76). Bearer required; multipart-поле file.
%%% Аудио и транскрипт не персистятся (приват-инвариант канона).
%%% @end
%%%-------------------------------------------------------------------
-module(handler_ai_transcribe).
-behaviour(cowboy_handler).
-export([init/2, trails/0]).
init(Req, Opts) ->
handle(Req, Opts).
trails() ->
[
#{
path => <<"/v1/ai/transcribe">>,
method => <<"POST">>,
description => <<"Transcribe press-to-talk audio via STT provider proxy. "
"Bearer required; multipart field `file` (audio).">>,
tags => [<<"AI">>],
responses => #{
200 => #{description => <<"{transcript, lang}; empty transcript = unheard">>},
400 => #{description => <<"missing_file / bad_multipart / too_large">>},
401 => #{description => <<"Unauthorized">>},
429 => #{description => <<"rate_limited">>},
502 => #{description => <<"upstream_error">>},
503 => #{description => <<"unavailable">>}
}
}
].
handle(Req, _Opts) ->
case cowboy_req:method(Req) of
<<"POST">> -> post(Req);
_ -> handler_utils:send_error(Req, 405, <<"Method not allowed">>)
end.
post(Req) ->
case handler_utils:auth_user(Req) of
{ok, UserId, Req1} ->
Max = logic_ai_transcribe:max_bytes(),
case handler_upload_utils:read_file_part(Req1, Max) of
{ok, Audio, CType, Req2} ->
case logic_ai_transcribe:transcribe(UserId, Audio, CType) of
{ok, Transcript, Lang} ->
handler_utils:send_json(Req2, 200,
#{<<"transcript">> => Transcript, <<"lang">> => Lang});
{error, rate_limited} ->
handler_utils:send_error(Req2, 429, <<"rate_limited">>);
{error, unavailable} ->
handler_utils:send_error(Req2, 503, <<"unavailable">>);
{error, upstream_error} ->
handler_utils:send_error(Req2, 502, <<"upstream_error">>);
{error, invalid_audio} ->
handler_utils:send_error(Req2, 400, <<"invalid_audio">>)
end;
{error, too_large, Req2} ->
handler_utils:send_error(Req2, 400, <<"too_large">>);
{error, missing_file, Req2} ->
handler_utils:send_error(Req2, 400, <<"missing_file">>);
{error, bad_multipart, Req2} ->
handler_utils:send_error(Req2, 400, <<"bad_multipart">>)
end;
{error, Code, Message, Req1} ->
handler_utils:send_error(Req1, Code, Message)
end.
+34 -2
View File
@@ -1,10 +1,10 @@
%%%------------------------------------------------------------------- %%%-------------------------------------------------------------------
%%% @doc Multipart helpers for image uploads (Back#71). %%% @doc Multipart helpers for file uploads (Back#71, Front#76).
%%% Field name: file %%% Field name: file
%%% @end %%% @end
%%%------------------------------------------------------------------- %%%-------------------------------------------------------------------
-module(handler_upload_utils). -module(handler_upload_utils).
-export([read_image_part/1]). -export([read_image_part/1, read_file_part/2]).
-spec read_image_part(cowboy_req:req()) -> -spec read_image_part(cowboy_req:req()) ->
{ok, binary(), cowboy_req:req()} | {ok, binary(), cowboy_req:req()} |
@@ -39,6 +39,38 @@ read_parts(Req, Max) ->
{error, missing_file, Req1} {error, missing_file, Req1}
end. end.
%% Произвольный файл без проверки содержимого (аудио для STT):
%% возвращает байты и content-type части.
-spec read_file_part(cowboy_req:req(), non_neg_integer()) ->
{ok, binary(), binary(), cowboy_req:req()} |
{error, missing_file | too_large | bad_multipart, cowboy_req:req()}.
read_file_part(Req, Max) ->
try read_file_parts(Req, Max)
catch
_:_ -> {error, bad_multipart, Req}
end.
read_file_parts(Req, Max) ->
case cowboy_req:read_part(Req) of
{ok, Headers, Req1} ->
case cow_multipart:form_data(Headers) of
{file, _Name, _Filename, CType} ->
case read_limited(Req1, Max, <<>>) of
{ok, Bin, Req2} -> {ok, Bin, part_ctype(CType), Req2};
{error, too_large, Req2} -> {error, too_large, Req2}
end;
{data, _Name} ->
{ok, _Body, Req2} = cowboy_req:read_part_body(Req1),
read_file_parts(Req2, Max)
end;
{done, Req1} ->
{error, missing_file, Req1}
end.
part_ctype(undefined) -> <<"application/octet-stream">>;
part_ctype(B) when is_binary(B) -> B;
part_ctype(L) when is_list(L) -> list_to_binary(L).
read_limited(Req, Max, Acc) when byte_size(Acc) > Max -> read_limited(Req, Max, Acc) when byte_size(Acc) > Max ->
_ = drain(Req), _ = drain(Req),
{error, too_large, Req}; {error, too_large, Req};
+284
View File
@@ -0,0 +1,284 @@
%%%-------------------------------------------------------------------
%%% @doc Прокси транскрибации голоса (PRODUCT-AI, Front#76).
%%%
%%% Front шлёт аудио press-to-talk записи; Back проксирует его в
%%% продуктовый STT-провайдер (адаптеры: xai, openai-совместимый).
%%% Ключи провайдера не покидают Back; аудио и транскрипт не
%%% персистятся (приват-инвариант канона). Rate-limit в ETS.
%%% @end
%%%-------------------------------------------------------------------
-module(logic_ai_transcribe).
-export([transcribe/3, max_bytes/0, ensure/0]).
%% Экспортировано для тестов и хэндлера.
-export([parse_provider_response/2, build_multipart/4]).
-define(RL_TABLE, eventhub_stt_rl).
-define(UP, eventhub_stt_upstream).
-define(RL_LIMIT, 10). % запросов на пользователя
-define(RL_WINDOW_MS, 60000).
-define(MAX_AUDIO_BYTES, 10 * 1024 * 1024).
-define(DEFAULT_XAI_URL, <<"https://api.x.ai/v1/stt">>).
-define(DEFAULT_OPENAI_URL, <<"http://localhost:8000/v1/audio/transcriptions">>).
-define(DEFAULT_TIMEOUT_MS, 15000).
-define(HTTP_PROFILE, eventhub_stt).
%%%===================================================================
%%% API
%%%===================================================================
-spec max_bytes() -> non_neg_integer().
max_bytes() -> ?MAX_AUDIO_BYTES.
-spec ensure() -> ok.
ensure() ->
case ets:info(?RL_TABLE) of
undefined ->
ets:new(?RL_TABLE, [named_table, public, set, {write_concurrency, true}]);
_ -> ok
end,
case ets:info(?UP) of
undefined ->
ets:new(?UP, [named_table, public, set]);
_ -> ok
end,
ok.
%% Транскрибация аудио пользователя. Пустой transcript в ответе
%% провайдера возвращаем как {ok, <<>>, Lang} — клиент покажет
%% тост «не расслышал» без retry-цикла (degrade-матрица канона).
-spec transcribe(binary(), binary(), binary()) ->
{ok, binary(), binary()} |
{error, unavailable | rate_limited | upstream_error | invalid_audio}.
transcribe(UserId, Audio, ContentType)
when is_binary(UserId), is_binary(Audio), byte_size(Audio) > 0,
byte_size(Audio) =< ?MAX_AUDIO_BYTES ->
ensure(),
case api_key() of
<<>> -> {error, unavailable};
Key ->
case allow(UserId) of
false -> {error, rate_limited};
true ->
case throttle_allow() of
false -> {error, rate_limited};
true ->
Provider = provider(),
call_provider(Provider, stt_url(Provider), Key, Audio, ContentType)
end
end
end;
transcribe(_UserId, _Audio, _ContentType) ->
{error, invalid_audio}.
%%%===================================================================
%%% Конфигурация (env, паттерн logic_geo)
%%%===================================================================
provider() ->
case env_bin("STT_PROVIDER") of
<<>> -> xai;
<<"openai">> -> openai;
<<"xai">> -> xai;
_ -> xai
end.
stt_url(xai) ->
case env_bin("STT_URL") of
<<>> -> ?DEFAULT_XAI_URL;
Url -> Url
end;
stt_url(openai) ->
case env_bin("STT_URL") of
<<>> -> ?DEFAULT_OPENAI_URL;
Url -> Url
end.
api_key() -> env_bin("STT_API_KEY").
timeout_ms() -> env_int("STT_TIMEOUT_MS", ?DEFAULT_TIMEOUT_MS).
proxy_url() -> env_bin("STT_PROXY_URL").
env_bin(Name) ->
case os:getenv(Name) of
false -> <<>>;
"" -> <<>>;
V -> list_to_binary(string:trim(V))
end.
env_int(Name, Default) ->
case os:getenv(Name) of
false -> Default;
"" -> Default;
S ->
try list_to_integer(S) of
N when N > 0 -> N;
_ -> Default
catch
_:_ -> Default
end
end.
%%%===================================================================
%%% Вызов провайдера
%%%===================================================================
call_provider(xai, Url, Key, Audio, CType) ->
%% xAI: multipart; опциональные поля до поля file, file — последним.
Body = build_multipart([{<<"format">>, <<"true">>}], Audio, <<"voice.dat">>, CType),
post_and_parse(xai, Url, Key, Body);
call_provider(openai, Url, Key, Audio, CType) ->
Model = case env_bin("STT_MODEL") of
<<>> -> <<"whisper-1">>;
M -> M
end,
Body = build_multipart([{<<"model">>, Model}], Audio, <<"voice.dat">>, CType),
post_and_parse(openai, Url, Key, Body).
post_and_parse(Provider, Url, Key, Body) ->
Headers = [{"authorization", binary_to_list(<<"Bearer ", Key/binary>>)}],
CTypeHdr = "multipart/form-data; boundary=" ++ binary_to_list(boundary(Body)),
case http_post(Url, Headers, CTypeHdr, Body) of
{ok, RespBody} -> parse_provider_response(Provider, RespBody);
{error, _} -> {error, upstream_error}
end.
%% Хук для тестов: fun(Url :: binary(), Headers, ContentType, Body) ->
%% {ok, RespBinary} | {error, term()}.
http_post(Url, Headers, ContentType, Body) ->
Fun = case application:get_env(eventhub, stt_http_post) of
{ok, F} when is_function(F, 4) -> F;
_ -> fun default_http_post/4
end,
Fun(Url, Headers, ContentType, Body).
default_http_post(Url, Headers, ContentType, Body) ->
_ = application:ensure_all_started(inets),
ensure_profile(),
UrlS = binary_to_list(Url),
case httpc:request(post, {UrlS, Headers, ContentType, Body},
[{timeout, timeout_ms()}], [{body_format, binary}], ?HTTP_PROFILE) of
{ok, {{_, Code, _}, _, Resp}} when Code >= 200, Code < 300, is_binary(Resp) ->
{ok, Resp};
{ok, {{_, Code, _}, _, Resp}} when Code >= 200, Code < 300 ->
{ok, iolist_to_binary(Resp)};
{ok, _} -> {error, http_status};
{error, Reason} -> {error, Reason}
end.
%% Отдельный профиль httpc: прокси (STT_PROXY_URL) применяется только
%% к STT-апстриму, остальные httpc-вызовы не затрагиваются.
ensure_profile() ->
case inets:start(httpc, [{profile, ?HTTP_PROFILE}]) of
{ok, _Pid} -> apply_proxy();
{error, {already_started, _Pid}} -> apply_proxy();
_ -> ok
end.
apply_proxy() ->
case parse_proxy(proxy_url()) of
{ok, Host, Port} ->
Addr = {Host, Port},
httpc:set_options([{proxy, {Addr, []}}, {https_proxy, {Addr, []}}],
?HTTP_PROFILE);
error -> ok
end.
parse_proxy(<<>>) -> error;
parse_proxy(Url) ->
try uri_string:parse(binary_to_list(Url)) of
#{host := Host, port := Port} -> {ok, Host, Port};
#{host := Host} -> {ok, Host, 8080};
_ -> error
catch
_:_ -> error
end.
%%%===================================================================
%%% Multipart и парсинг ответа
%%%===================================================================
%% Multipart/form-data: Fields (name=value) идут первыми, файл —
%% последним полем (требование xAI для стриминг-загрузки).
-spec build_multipart([{binary(), binary()}], binary(), binary(), binary()) -> binary().
build_multipart(Fields, FileBin, FileName, FileCType) ->
B = boundary_bin(),
FieldParts = [field_part(B, Name, Value) || {Name, Value} <- Fields],
FilePart = [
<<"--", B/binary, "\r\n">>,
<<"Content-Disposition: form-data; name=\"file\"; filename=\"">>,
FileName, <<"\"\r\n">>,
<<"Content-Type: ">>, FileCType, <<"\r\n\r\n">>,
FileBin, <<"\r\n">>
],
iolist_to_binary([FieldParts, FilePart, <<"--", B/binary, "--\r\n">>]).
field_part(B, Name, Value) ->
[
<<"--", B/binary, "\r\n">>,
<<"Content-Disposition: form-data; name=\"">>, Name, <<"\"\r\n\r\n">>,
Value, <<"\r\n">>
].
boundary(Body) ->
{Pos, _} = binary:match(Body, <<"\r\n">>),
<<"--", B/binary>> = binary:part(Body, 0, Pos),
B.
boundary_bin() ->
Hex = [io_lib:format("~2.16.0b", [rand:uniform(256) - 1]) || _ <- lists:seq(1, 12)],
<<"ehstt", (iolist_to_binary(Hex))/binary>>.
%% Маппинг ответа провайдера в {ok, Transcript, Lang}.
%% Оба адаптера читают `text`; язык — из `language` (xAI определяет сам).
-spec parse_provider_response(xai | openai, binary()) ->
{ok, binary(), binary()} | {error, upstream_error}.
parse_provider_response(_Provider, RespBody) when is_binary(RespBody) ->
try jsx:decode(RespBody, [return_maps]) of
#{<<"text">> := Text} = M when is_binary(Text) ->
Lang = case maps:get(<<"language">>, M, null) of
L when is_binary(L) -> L;
_ -> <<>>
end,
{ok, Text, Lang};
_ -> {error, upstream_error}
catch
_:_ -> {error, upstream_error}
end;
parse_provider_response(_Provider, _) ->
{error, upstream_error}.
%%%===================================================================
%%% Rate-limit (паттерн logic_ai_metrics)
%%%===================================================================
allow(UserId) ->
ensure(),
Now = erlang:system_time(millisecond),
Key = {rl, UserId},
Result = case ets:lookup(?RL_TABLE, Key) of
[{_, WindowStart, Count}] when Now - WindowStart < ?RL_WINDOW_MS ->
Count < ?RL_LIMIT;
_ ->
true
end,
case Result of
true ->
ets:update_counter(?RL_TABLE, Key, {3, 1}, {Key, Now, 0}),
true;
false ->
false
end.
%% Глобальный троттлинг апстрима: не более 2 rps на весь узел.
throttle_allow() ->
ensure(),
NowMs = erlang:monotonic_time(millisecond),
MinGap = 500,
case ets:lookup(?UP, last) of
[{last, Prev}] when NowMs - Prev < MinGap -> false;
_ ->
ets:insert(?UP, {last, NowMs}),
true
end.
+1
View File
@@ -97,6 +97,7 @@ user() ->
handler_search, handler_search,
handler_geo, handler_geo,
handler_ai_metrics, handler_ai_metrics,
handler_ai_transcribe,
handler_subscription, handler_subscription,
handler_ticket_by_id, handler_ticket_by_id,
handler_tickets, handler_tickets,
+129
View File
@@ -0,0 +1,129 @@
%%%-------------------------------------------------------------------
%%% EUnit: прокси транскрибации голоса (logic_ai_transcribe, Front#76).
%%% Без реальных HTTP: апстрим подменяется хуком stt_http_post.
%%%-------------------------------------------------------------------
-module(logic_ai_transcribe_tests).
-include_lib("eunit/include/eunit.hrl").
-define(RL_TABLE, eventhub_stt_rl).
-define(UP, eventhub_stt_upstream).
-define(XAI_URL, <<"https://stt.test/xai">>).
setup() ->
logic_ai_transcribe:ensure(),
ets:delete_all_objects(?RL_TABLE),
ets:delete_all_objects(?UP),
os:putenv("STT_API_KEY", "test-key"),
os:putenv("STT_PROVIDER", "xai"),
os:putenv("STT_URL", binary_to_list(?XAI_URL)),
application:set_env(eventhub, stt_http_post, fun stub_post/4),
ok.
cleanup(_) ->
os:unsetenv("STT_API_KEY"),
os:unsetenv("STT_PROVIDER"),
os:unsetenv("STT_URL"),
os:unsetenv("STT_MODEL"),
application:unset_env(eventhub, stt_http_post),
ok.
logic_ai_transcribe_test_() ->
{foreach, fun setup/0, fun cleanup/1, [
{"no key means unavailable", fun test_unavailable/0},
{"xai adapter posts to configured url", fun test_xai_success/0},
{"openai adapter uses model field", fun test_openai_adapter/0},
{"upstream errors mapped", fun test_upstream_errors/0},
{"per-user rate limit", fun test_user_rate_limit/0},
{"global upstream throttle", fun test_global_throttle/0},
{"invalid audio rejected", fun test_invalid_audio/0},
{"multipart layout", fun test_multipart/0}
]}.
%% Стаб апстрима: фиксирует запрос и отвечает успешным JSON.
stub_post(Url, Headers, ContentType, Body) ->
put(stt_call, {Url, Headers, ContentType, Body}),
{ok, <<"{\"text\":\"привет\",\"language\":\"ru\"}">>}.
test_unavailable() ->
os:unsetenv("STT_API_KEY"),
?assertEqual({error, unavailable},
logic_ai_transcribe:transcribe(<<"u1">>, <<"audio">>, <<"audio/webm">>)).
test_xai_success() ->
?assertEqual({ok, <<"привет">>, <<"ru">>},
logic_ai_transcribe:transcribe(<<"u1">>, <<"audio-bytes">>, <<"audio/webm">>)),
{Url, Headers, ContentType, Body} = get(stt_call),
?assertEqual(?XAI_URL, Url),
?assertEqual([{"authorization", "Bearer test-key"}], Headers),
?assertMatch("multipart/form-data; boundary=" ++ _, ContentType),
?assertMatch({_, _}, binary:match(Body, <<"name=\"format\"">>)),
?assertMatch({_, _}, binary:match(Body, <<"name=\"file\"">>)),
?assertMatch({_, _}, binary:match(Body, <<"audio-bytes">>)),
%% Поле file — последнее (требование xAI).
{PosFormat, _} = binary:match(Body, <<"name=\"format\"">>),
{PosFile, _} = binary:match(Body, <<"name=\"file\"">>),
?assert(PosFormat < PosFile).
test_openai_adapter() ->
os:putenv("STT_PROVIDER", "openai"),
os:putenv("STT_URL", "https://stt.test/openai"),
os:putenv("STT_MODEL", "whisper-large-v3"),
?assertEqual({ok, <<"привет">>, <<"ru">>},
logic_ai_transcribe:transcribe(<<"u1">>, <<"audio">>, <<"audio/webm">>)),
{Url, _Headers, _CType, Body} = get(stt_call),
?assertEqual(<<"https://stt.test/openai">>, Url),
?assertMatch({_, _}, binary:match(Body, <<"name=\"model\"">>)),
?assertMatch({_, _}, binary:match(Body, <<"whisper-large-v3">>)).
test_upstream_errors() ->
%% Ошибка транспорта.
application:set_env(eventhub, stt_http_post,
fun(_Url, _H, _C, _B) -> {error, nxdomain} end),
?assertEqual({error, upstream_error},
logic_ai_transcribe:transcribe(<<"u1">>, <<"audio">>, <<"audio/webm">>)),
%% Некорректный JSON от апстрима.
ets:delete_all_objects(?UP),
application:set_env(eventhub, stt_http_post,
fun(_Url, _H, _C, _B) -> {ok, <<"not-json">>} end),
?assertEqual({error, upstream_error},
logic_ai_transcribe:transcribe(<<"u2">>, <<"audio">>, <<"audio/webm">>)),
%% JSON без поля text.
ets:delete_all_objects(?UP),
application:set_env(eventhub, stt_http_post,
fun(_Url, _H, _C, _B) -> {ok, <<"{\"language\":\"ru\"}">>} end),
?assertEqual({error, upstream_error},
logic_ai_transcribe:transcribe(<<"u3">>, <<"audio">>, <<"audio/webm">>)).
test_user_rate_limit() ->
Now = erlang:system_time(millisecond),
ets:insert(?RL_TABLE, {{rl, <<"u-rl">>}, Now, 10}),
?assertEqual({error, rate_limited},
logic_ai_transcribe:transcribe(<<"u-rl">>, <<"audio">>, <<"audio/webm">>)),
%% Другой пользователь не затронут.
?assertEqual({ok, <<"привет">>, <<"ru">>},
logic_ai_transcribe:transcribe(<<"u-ok">>, <<"audio">>, <<"audio/webm">>)).
test_global_throttle() ->
ets:insert(?UP, {last, erlang:monotonic_time(millisecond)}),
?assertEqual({error, rate_limited},
logic_ai_transcribe:transcribe(<<"u1">>, <<"audio">>, <<"audio/webm">>)).
test_invalid_audio() ->
?assertEqual({error, invalid_audio},
logic_ai_transcribe:transcribe(<<"u1">>, <<>>, <<"audio/webm">>)),
?assertEqual({error, invalid_audio},
logic_ai_transcribe:transcribe(<<"u1">>,
binary:copy(<<"x">>, logic_ai_transcribe:max_bytes() + 1),
<<"audio/webm">>)).
test_multipart() ->
Body = logic_ai_transcribe:build_multipart(
[{<<"a">>, <<"1">>}], <<"DATA">>, <<"voice.dat">>, <<"audio/webm">>),
?assertMatch({_, _}, binary:match(Body, <<"name=\"a\"\r\n\r\n1\r\n">>)),
?assertMatch({_, _}, binary:match(
Body, <<"name=\"file\"; filename=\"voice.dat\"">>)),
?assertMatch({_, _}, binary:match(Body, <<"Content-Type: audio/webm">>)),
?assertMatch({_, _}, binary:match(Body, <<"DATA">>)),
?assert(binary:match(Body, <<"--\r\n">>) =/= nomatch),
%% Завершающий boundary в конце тела: --<B>--\r\n.
?assertEqual(<<"--\r\n">>, binary:part(Body, byte_size(Body), -4)).