Время читается единой точкой, а конец потока узнаётся не по тексту
- заведён internal/clock: Now даёт метку в UTC, Start — начало измерения длительности с монотонными часами; девять мест рабочего кода и запрос захвата переведены на него, долг «время time.Now() по месту» закрыт - клиент SpeechKit узнавал конец потока сравнением err.Error() == "EOF": отказ с тем же текстом вернул бы усечённую расшифровку как готовую, теперь errors.Is(err, io.EOF) - починены находки новых линтеров: две опечатки, два slog.DiscardHandler, четыре неэкранированные подстановки в docker/entrypoint.sh
This commit is contained in:
@@ -12,21 +12,21 @@ if [ "${USER}" != "transcriber" ]; then
|
|||||||
fi
|
fi
|
||||||
|
|
||||||
if [ -z "${USER_GID}" ]; then
|
if [ -z "${USER_GID}" ]; then
|
||||||
USER_GID="$(id -g ${USER})"
|
USER_GID="$(id -g "${USER}")"
|
||||||
fi
|
fi
|
||||||
|
|
||||||
if [ -z "${USER_UID}" ]; then
|
if [ -z "${USER_UID}" ]; then
|
||||||
USER_UID="$(id -u ${USER})"
|
USER_UID="$(id -u "${USER}")"
|
||||||
fi
|
fi
|
||||||
|
|
||||||
# Change GID for USER?
|
# 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]*/${USER}:\1:${USER_GID}/" /etc/group
|
||||||
sed -i -e "s/^${USER}:\([^:]*\):\([0-9]*\):[0-9]*/${USER}:\1:\2:${USER_GID}/" /etc/passwd
|
sed -i -e "s/^${USER}:\([^:]*\):\([0-9]*\):[0-9]*/${USER}:\1:\2:${USER_GID}/" /etc/passwd
|
||||||
fi
|
fi
|
||||||
|
|
||||||
# Change UID for USER?
|
# 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
|
sed -i -e "s/^${USER}:\([^:]*\):[0-9]*:\([0-9]*\)/${USER}:\1:${USER_UID}:\2/" /etc/passwd
|
||||||
fi
|
fi
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"google.golang.org/grpc"
|
"google.golang.org/grpc"
|
||||||
@@ -157,7 +158,10 @@ func (s *speechKitService) getRecognitionText(operationID string) (string, error
|
|||||||
for {
|
for {
|
||||||
resp, err := stream.Recv()
|
resp, err := stream.Recv()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if err.Error() == "EOF" {
|
// Конец потока библиотека отдаёт ровно `io.EOF`. Прежде он узнавался
|
||||||
|
// сравнением текста сообщения: так же выглядел бы и настоящий отказ
|
||||||
|
// с текстом «EOF», и распознавание молча вернуло бы половину текста.
|
||||||
|
if errors.Is(err, io.EOF) {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
return "", fmt.Errorf("failed to receive recognition response: %w", err)
|
return "", fmt.Errorf("failed to receive recognition response: %w", err)
|
||||||
|
|||||||
@@ -14,6 +14,8 @@ import (
|
|||||||
"git.vakhrushev.me/av/transcriber/internal/entity"
|
"git.vakhrushev.me/av/transcriber/internal/entity"
|
||||||
|
|
||||||
"git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase/migrations"
|
"git.vakhrushev.me/av/transcriber/internal/adapter/repo/pocketbase/migrations"
|
||||||
|
|
||||||
|
"git.vakhrushev.me/av/transcriber/internal/clock"
|
||||||
)
|
)
|
||||||
|
|
||||||
type TranscriptJobRepository struct {
|
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) {
|
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(`
|
query := repo.app.DB().NewQuery(`
|
||||||
UPDATE {{` + migrations.JobsCollection + `}}
|
UPDATE {{` + migrations.JobsCollection + `}}
|
||||||
@@ -142,7 +150,7 @@ func (repo *TranscriptJobRepository) FindAndAcquire(state, acquisitionId string,
|
|||||||
if errors.Is(err, sql.ErrNoRows) {
|
if errors.Is(err, sql.ErrNoRows) {
|
||||||
return nil, &contract.JobNotFoundError{State: state, Message: "appropriate job not found"}
|
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
|
return row.toJob(), nil
|
||||||
|
|||||||
@@ -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()
|
||||||
|
}
|
||||||
@@ -2,6 +2,8 @@ package entity
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"git.vakhrushev.me/av/transcriber/internal/clock"
|
||||||
)
|
)
|
||||||
|
|
||||||
type TranscribeJob struct {
|
type TranscribeJob struct {
|
||||||
@@ -52,13 +54,13 @@ func (j *TranscribeJob) MoveToState(state string) {
|
|||||||
// именно отказавшие: иначе задача, прошедшая конвейер целиком, накопила бы
|
// именно отказавшие: иначе задача, прошедшая конвейер целиком, накопила бы
|
||||||
// их поштучно и умерла бы здоровой.
|
// их поштучно и умерла бы здоровой.
|
||||||
j.Attempts = 0
|
j.Attempts = 0
|
||||||
j.UpdatedAt = time.Now()
|
j.UpdatedAt = clock.Now()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (j *TranscribeJob) MoveToStateAndDelay(state string, delay *time.Time) {
|
func (j *TranscribeJob) MoveToStateAndDelay(state string, delay *time.Time) {
|
||||||
j.MoveToState(state)
|
j.MoveToState(state)
|
||||||
j.DelayTime = delay
|
j.DelayTime = delay
|
||||||
j.UpdatedAt = time.Now()
|
j.UpdatedAt = clock.Now()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (j *TranscribeJob) Done(transcriptionText string) {
|
func (j *TranscribeJob) Done(transcriptionText string) {
|
||||||
@@ -78,7 +80,7 @@ func (j *TranscribeJob) RetryAfter(delay time.Time) {
|
|||||||
j.AcquisitionID = nil
|
j.AcquisitionID = nil
|
||||||
j.AcquireTime = nil
|
j.AcquireTime = nil
|
||||||
j.DelayTime = &delay
|
j.DelayTime = &delay
|
||||||
j.UpdatedAt = time.Now()
|
j.UpdatedAt = clock.Now()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Die переводит задачу, исчерпавшую попытки, в состояние «мертва». Число
|
// Die переводит задачу, исчерпавшую попытки, в состояние «мертва». Число
|
||||||
|
|||||||
@@ -3,7 +3,6 @@ package service
|
|||||||
import (
|
import (
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -36,7 +35,7 @@ func (r *stubJobRepo) FindAndAcquire(string, string, time.Time) (*entity.Transcr
|
|||||||
}
|
}
|
||||||
|
|
||||||
func serviceWithRepo(repo contract.TranscriptJobRepository) *TranscribeService {
|
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)
|
return NewTranscribeService(repo, nil, nil, nil, nil, nil, logger)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -90,7 +90,7 @@ func newPipelineEnv(t *testing.T, metaviewer contract.AudioMetaViewer, converter
|
|||||||
converter,
|
converter,
|
||||||
&recognizer.MemoryAudioRecognizer{},
|
&recognizer.MemoryAudioRecognizer{},
|
||||||
sender,
|
sender,
|
||||||
slog.New(slog.NewTextHandler(io.Discard, nil)),
|
slog.New(slog.DiscardHandler),
|
||||||
)
|
)
|
||||||
|
|
||||||
return &pipelineEnv{app: app, service: svc, jobRepo: jobRepo, fileRepo: fileRepo, sender: sender}
|
return &pipelineEnv{app: app, service: svc, jobRepo: jobRepo, fileRepo: fileRepo, sender: sender}
|
||||||
|
|||||||
@@ -14,6 +14,8 @@ import (
|
|||||||
"git.vakhrushev.me/av/transcriber/internal/entity"
|
"git.vakhrushev.me/av/transcriber/internal/entity"
|
||||||
"git.vakhrushev.me/av/transcriber/internal/metrics"
|
"git.vakhrushev.me/av/transcriber/internal/metrics"
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
|
|
||||||
|
"git.vakhrushev.me/av/transcriber/internal/clock"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
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)
|
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())
|
err = s.converter.Convert(src.Path(), dest.Path())
|
||||||
conversionDuration := time.Since(startTime)
|
conversionDuration := time.Since(startTime)
|
||||||
|
|
||||||
@@ -321,7 +323,7 @@ func (s *TranscribeService) transcribeJob(job *entity.TranscribeJob, holder stri
|
|||||||
// Обновляем задачу с ID операции распознавания
|
// Обновляем задачу с ID операции распознавания
|
||||||
job.FileID = &destFileRecord.Id
|
job.FileID = &destFileRecord.Id
|
||||||
job.RecognitionOpID = &operationID
|
job.RecognitionOpID = &operationID
|
||||||
delayTime := time.Now().Add(firstCheckDelay)
|
delayTime := clock.Now().Add(firstCheckDelay)
|
||||||
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
|
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
|
||||||
|
|
||||||
if err := s.jobRepo.Save(job, holder); err != nil {
|
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 {
|
func (s *TranscribeService) checkTranscribeJob(job *entity.TranscribeJob, holder string) error {
|
||||||
if job.RecognitionOpID == nil {
|
if job.RecognitionOpID == nil {
|
||||||
s.logger.Error("Recognition operation ID not found", "job_id", job.Id)
|
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
|
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)
|
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)
|
job.MoveToStateAndDelay(entity.StateTranscribe, &delayTime)
|
||||||
if err := s.jobRepo.Save(job, holder); err != nil {
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
||||||
s.logger.Error("Failed to save job", "error", err, "job_id", job.Id)
|
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) {
|
func (s *TranscribeService) findJob(state string, expiration time.Duration) (*entity.TranscribeJob, string, error) {
|
||||||
acquisitionId := uuid.NewString()
|
acquisitionId := uuid.NewString()
|
||||||
rottingTime := time.Now().Add(-1 * expiration)
|
rottingTime := clock.Now().Add(-1 * expiration)
|
||||||
|
|
||||||
job, err := s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
|
job, err := s.jobRepo.FindAndAcquire(state, acquisitionId, rottingTime)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -448,7 +450,7 @@ func (s *TranscribeService) scheduleRetry(job *entity.TranscribeJob, holder stri
|
|||||||
return
|
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 {
|
if err := s.jobRepo.Save(job, holder); err != nil {
|
||||||
var lostOnSave *contract.LostAcquisitionError
|
var lostOnSave *contract.LostAcquisitionError
|
||||||
|
|||||||
@@ -27,6 +27,8 @@ import (
|
|||||||
"github.com/pocketbase/pocketbase/apis"
|
"github.com/pocketbase/pocketbase/apis"
|
||||||
"github.com/pocketbase/pocketbase/core"
|
"github.com/pocketbase/pocketbase/core"
|
||||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||||
|
|
||||||
|
"git.vakhrushev.me/av/transcriber/internal/clock"
|
||||||
)
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
@@ -201,7 +203,7 @@ func main() {
|
|||||||
// `sloggin`, а хранилище пишет запросы в свою таблицу, которой в
|
// `sloggin`, а хранилище пишет запросы в свою таблицу, которой в
|
||||||
// журнале контейнера не видно. Поля — те, что просит конвенция.
|
// журнале контейнера не видно. Поля — те, что просит конвенция.
|
||||||
se.Router.BindFunc(func(e *core.RequestEvent) error {
|
se.Router.BindFunc(func(e *core.RequestEvent) error {
|
||||||
start := time.Now()
|
start := clock.Start()
|
||||||
err := e.Next()
|
err := e.Next()
|
||||||
|
|
||||||
level := slog.LevelInfo
|
level := slog.LevelInfo
|
||||||
|
|||||||
Reference in New Issue
Block a user