From 1d243ad2f6ec543d1c840605c16171d216875c1a Mon Sep 17 00:00:00 2001 From: Anton Vakhrushev Date: Thu, 13 Aug 2026 09:01:36 +0300 Subject: [PATCH] =?UTF-8?q?=D0=92=D1=80=D0=B5=D0=BC=D1=8F=20=D1=87=D0=B8?= =?UTF-8?q?=D1=82=D0=B0=D0=B5=D1=82=D1=81=D1=8F=20=D0=B5=D0=B4=D0=B8=D0=BD?= =?UTF-8?q?=D0=BE=D0=B9=20=D1=82=D0=BE=D1=87=D0=BA=D0=BE=D0=B9,=20=D0=B0?= =?UTF-8?q?=20=D0=BA=D0=BE=D0=BD=D0=B5=D1=86=20=D0=BF=D0=BE=D1=82=D0=BE?= =?UTF-8?q?=D0=BA=D0=B0=20=D1=83=D0=B7=D0=BD=D0=B0=D1=91=D1=82=D1=81=D1=8F?= =?UTF-8?q?=20=D0=BD=D0=B5=20=D0=BF=D0=BE=20=D1=82=D0=B5=D0=BA=D1=81=D1=82?= =?UTF-8?q?=D1=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - заведён internal/clock: Now даёт метку в UTC, Start — начало измерения длительности с монотонными часами; девять мест рабочего кода и запрос захвата переведены на него, долг «время time.Now() по месту» закрыт - клиент SpeechKit узнавал конец потока сравнением err.Error() == "EOF": отказ с тем же текстом вернул бы усечённую расшифровку как готовую, теперь errors.Is(err, io.EOF) - починены находки новых линтеров: две опечатки, два slog.DiscardHandler, четыре неэкранированные подстановки в docker/entrypoint.sh --- docker/entrypoint.sh | 8 ++--- .../adapter/recognizer/yandex/speechkit.go | 6 +++- .../repo/pocketbase/transcript_job_repo.go | 12 +++++-- internal/clock/clock.go | 32 +++++++++++++++++++ internal/entity/job.go | 8 +++-- internal/service/find_job_test.go | 3 +- internal/service/pipeline_test.go | 2 +- internal/service/transcribe.go | 14 ++++---- main.go | 4 ++- 9 files changed, 69 insertions(+), 20 deletions(-) create mode 100644 internal/clock/clock.go diff --git a/docker/entrypoint.sh b/docker/entrypoint.sh index ac372e7..34b1ccd 100755 --- a/docker/entrypoint.sh +++ b/docker/entrypoint.sh @@ -12,21 +12,21 @@ if [ "${USER}" != "transcriber" ]; then fi if [ -z "${USER_GID}" ]; then - USER_GID="$(id -g ${USER})" + USER_GID="$(id -g "${USER}")" fi if [ -z "${USER_UID}" ]; then - USER_UID="$(id -u ${USER})" + USER_UID="$(id -u "${USER}")" fi # Change GID for USER? -if [ -n "${USER_GID}" ] && [ "${USER_GID}" != "$(id -g ${USER})" ]; then +if [ -n "${USER_GID}" ] && [ "${USER_GID}" != "$(id -g "${USER}")" ]; then sed -i -e "s/^${USER}:\([^:]*\):[0-9]*/${USER}:\1:${USER_GID}/" /etc/group sed -i -e "s/^${USER}:\([^:]*\):\([0-9]*\):[0-9]*/${USER}:\1:\2:${USER_GID}/" /etc/passwd fi # Change UID for USER? -if [ -n "${USER_UID}" ] && [ "${USER_UID}" != "$(id -u ${USER})" ]; then +if [ -n "${USER_UID}" ] && [ "${USER_UID}" != "$(id -u "${USER}")" ]; then sed -i -e "s/^${USER}:\([^:]*\):[0-9]*:\([0-9]*\)/${USER}:\1:${USER_UID}:\2/" /etc/passwd fi diff --git a/internal/adapter/recognizer/yandex/speechkit.go b/internal/adapter/recognizer/yandex/speechkit.go index 19ad3d6..6a19bb2 100644 --- a/internal/adapter/recognizer/yandex/speechkit.go +++ b/internal/adapter/recognizer/yandex/speechkit.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "io" "strings" "google.golang.org/grpc" @@ -157,7 +158,10 @@ func (s *speechKitService) getRecognitionText(operationID string) (string, error for { resp, err := stream.Recv() if err != nil { - if err.Error() == "EOF" { + // Конец потока библиотека отдаёт ровно `io.EOF`. Прежде он узнавался + // сравнением текста сообщения: так же выглядел бы и настоящий отказ + // с текстом «EOF», и распознавание молча вернуло бы половину текста. + if errors.Is(err, io.EOF) { break } return "", fmt.Errorf("failed to receive recognition response: %w", err) diff --git a/internal/adapter/repo/pocketbase/transcript_job_repo.go b/internal/adapter/repo/pocketbase/transcript_job_repo.go index 9357c5f..b24dece 100644 --- a/internal/adapter/repo/pocketbase/transcript_job_repo.go +++ b/internal/adapter/repo/pocketbase/transcript_job_repo.go @@ -14,6 +14,8 @@ import ( "git.vakhrushev.me/av/transcriber/internal/entity" "git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase/migrations" + + "git.vakhrushev.me/av/transcriber/internal/clock" ) type TranscriptJobRepository struct { @@ -107,7 +109,13 @@ const acquireColumns = `id, state, source, file, error_text, acquisition_id, ` + // разделителем, обратил бы условие срока в постоянную истину или постоянную // ложь — молча. func (repo *TranscriptJobRepository) FindAndAcquire(state, acquisitionId string, rottingTime time.Time) (*entity.TranscribeJob, error) { - now := types.NowDateTime() + // Метка времени берётся единой точкой, а не `types.NowDateTime()`: обёртка + // хранилища читает часы сама, и запрет линтера её не видит — новая метка в + // этом запросе обошла бы единую точку молча. + now, err := types.ParseDateTime(clock.Now()) + if err != nil { + return nil, fmt.Errorf("failed to parse current time: %w", err) + } query := repo.app.DB().NewQuery(` UPDATE {{` + migrations.JobsCollection + `}} @@ -142,7 +150,7 @@ func (repo *TranscriptJobRepository) FindAndAcquire(state, acquisitionId string, if errors.Is(err, sql.ErrNoRows) { return nil, &contract.JobNotFoundError{State: state, Message: "appropriate job not found"} } - return nil, fmt.Errorf("failed to aquire job with state %s: %w", state, err) + return nil, fmt.Errorf("failed to acquire job with state %s: %w", state, err) } return row.toJob(), nil diff --git a/internal/clock/clock.go b/internal/clock/clock.go new file mode 100644 index 0000000..df64944 --- /dev/null +++ b/internal/clock/clock.go @@ -0,0 +1,32 @@ +// Package clock — единая точка чтения времени. +// +// Прежде время брали `time.Now()` по месту вызова, и конвенция +// [docs/conventions/database.md] числила это долгом: метка времени в локальной +// зоне, а хранилище сравнивает времена **строками** побайтово. Долг закрыт +// заведением этого пакета; правило держит `forbidigo` в `.golangci.yml` — +// `time.Now` вне этого пакета запрещён. +// +// Метка времени и измерение длительности читаются по-разному, и потому здесь две +// функции, а не одна. +package clock + +import "time" + +// Now — метка времени: UTC, как её пишет и сравнивает хранилище. +// +// Приведение к UTC снимает монотонные часы, и для метки это верно: её кладут в +// колонку и сравнивают с чужими значениями, а не с собственным прошлым +// показанием. +func Now() time.Time { + return time.Now().UTC() +} + +// Start — начало измерения длительности: время **с монотонными часами**. +// +// Зоны у него нет намеренно: значение не выходит наружу и годится только на +// вход `time.Since`. Монотонные часы здесь и нужны — иначе перевод стрелок или +// поправка ntp посреди конвертации дала бы отрицательную или скачущую +// длительность в журнале и в метрике. +func Start() time.Time { + return time.Now() +} diff --git a/internal/entity/job.go b/internal/entity/job.go index e498947..4c29edf 100644 --- a/internal/entity/job.go +++ b/internal/entity/job.go @@ -2,6 +2,8 @@ package entity import ( "time" + + "git.vakhrushev.me/av/transcriber/internal/clock" ) type TranscribeJob struct { @@ -52,13 +54,13 @@ func (j *TranscribeJob) MoveToState(state string) { // именно отказавшие: иначе задача, прошедшая конвейер целиком, накопила бы // их поштучно и умерла бы здоровой. j.Attempts = 0 - j.UpdatedAt = time.Now() + j.UpdatedAt = clock.Now() } func (j *TranscribeJob) MoveToStateAndDelay(state string, delay *time.Time) { j.MoveToState(state) j.DelayTime = delay - j.UpdatedAt = time.Now() + j.UpdatedAt = clock.Now() } func (j *TranscribeJob) Done(transcriptionText string) { @@ -78,7 +80,7 @@ func (j *TranscribeJob) RetryAfter(delay time.Time) { j.AcquisitionID = nil j.AcquireTime = nil j.DelayTime = &delay - j.UpdatedAt = time.Now() + j.UpdatedAt = clock.Now() } // Die переводит задачу, исчерпавшую попытки, в состояние «мертва». Число diff --git a/internal/service/find_job_test.go b/internal/service/find_job_test.go index d749c0a..c30a93a 100644 --- a/internal/service/find_job_test.go +++ b/internal/service/find_job_test.go @@ -3,7 +3,6 @@ package service import ( "errors" "fmt" - "io" "log/slog" "testing" "time" @@ -36,7 +35,7 @@ func (r *stubJobRepo) FindAndAcquire(string, string, time.Time) (*entity.Transcr } func serviceWithRepo(repo contract.TranscriptJobRepository) *TranscribeService { - logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + logger := slog.New(slog.DiscardHandler) return NewTranscribeService(repo, nil, nil, nil, nil, nil, logger) } diff --git a/internal/service/pipeline_test.go b/internal/service/pipeline_test.go index 4a86493..fe75e2f 100644 --- a/internal/service/pipeline_test.go +++ b/internal/service/pipeline_test.go @@ -90,7 +90,7 @@ func newPipelineEnv(t *testing.T, metaviewer contract.AudioMetaViewer, converter converter, &recognizer.MemoryAudioRecognizer{}, sender, - slog.New(slog.NewTextHandler(io.Discard, nil)), + slog.New(slog.DiscardHandler), ) return &pipelineEnv{app: app, service: svc, jobRepo: jobRepo, fileRepo: fileRepo, sender: sender} diff --git a/internal/service/transcribe.go b/internal/service/transcribe.go index 0e2a6c7..8f177da 100644 --- a/internal/service/transcribe.go +++ b/internal/service/transcribe.go @@ -14,6 +14,8 @@ import ( "git.vakhrushev.me/av/transcriber/internal/entity" "git.vakhrushev.me/av/transcriber/internal/metrics" "github.com/google/uuid" + + "git.vakhrushev.me/av/transcriber/internal/clock" ) const ( @@ -224,7 +226,7 @@ func (s *TranscribeService) convertJob(job *entity.TranscribeJob, holder string) s.logger.Info("Converting file", "job_id", job.Id, "src_format", srcExt) // Измеряем время конвертации - startTime := time.Now() + startTime := clock.Start() err = s.converter.Convert(src.Path(), dest.Path()) conversionDuration := time.Since(startTime) @@ -321,7 +323,7 @@ func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder stri // Обновляем задачу с ID операции распознавания job.FileID = &destFileRecord.Id job.RecognitionOpID = &operationID - delayTime := time.Now().Add(firstCheckDelay) + delayTime := clock.Now().Add(firstCheckDelay) job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime) if err := s.jobRepo.Save(job, holder); err != nil { @@ -336,7 +338,7 @@ func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder stri func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder string) error { if job.RecognitionOpID == nil { s.logger.Error("Recognition operation ID not found", "job_id", job.Id) - return fmt.Errorf("recogniton opId not found for job: %s", job.Id) + return fmt.Errorf("recognition opId not found for job: %s", job.Id) } opId := *job.RecognitionOpID @@ -354,7 +356,7 @@ func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder // здесь своя, числом, а число попыток обнуляется переходом: ожидание // чужой операции попытку не тратит. s.logger.Info("Operation in progress", "job_id", job.Id, "operation_id", opId) - delayTime := time.Now().Add(nextCheckDelay) + delayTime := clock.Now().Add(nextCheckDelay) job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime) if err := s.jobRepo.Save(job, holder); err != nil { s.logger.Error("Failed to save job", "error", err, "job_id", job.Id) @@ -397,7 +399,7 @@ func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder // переводят в «мертва» и сообщают об этом отправителю. func (s *TranscribeService) findJob(state string, expiration time.Duration) (*entity.TranscribeJob, string, error) { acquisitionId := uuid.NewString() - rottingTime := time.Now().Add(-1 * expiration) + rottingTime := clock.Now().Add(-1 * expiration) job, err := s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime) if err != nil { @@ -448,7 +450,7 @@ func (s *TranscribeService) scheduleRetry(job *entity.TranscribeJob, holder stri return } - job.RetryAfter(time.Now().Add(retryDelay(job.Attempts))) + job.RetryAfter(clock.Now().Add(retryDelay(job.Attempts))) if err := s.jobRepo.Save(job, holder); err != nil { var lostOnSave *contract.LostAcquisitionError diff --git a/main.go b/main.go index dc5767d..b9ff958 100644 --- a/main.go +++ b/main.go @@ -27,6 +27,8 @@ import ( "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() { @@ -201,7 +203,7 @@ func main() { // `sloggin`, а хранилище пишет запросы в свою таблицу, которой в // журнале контейнера не видно. Поля — те, что просит конвенция. se.Router.BindFunc(func(e *core.RequestEvent) error { - start := time.Now() + start := clock.Start() err := e.Next() level := slog.LevelInfo