у записи появился владелец: чужую больше не отдают

- колонка `owner` связью с `users` в обеих коллекциях новым шагом схемы
  `202608140001`; чтение задачи сужено владельцем, и чужая, ничья и
  несуществующая дают один ответ; правило просмотра файлов сужено им же
- приём по HTTP берёт владельца из сессии, а предъявителя без учётной записи
  пользователя отвергает до чтения тела: позже пришлось бы убирать уложенный
  файл, а уборки файлов сервис не умеет. Выборка воркера владельцем не сужается
- удаление учётной записи с записями отвергается стражем, и вешает его сама
  сборка хранилища: сборка, забывшая его позвать, теряла защиту молча
This commit is contained in:
av
2026-08-14 12:18:11 +03:00
parent b7d4660aef
commit 8af8ec2e54
46 changed files with 2520 additions and 101 deletions
+1 -1
View File
@@ -26,7 +26,7 @@ type stubJobRepo struct {
func (r *stubJobRepo) Create(*entity.TranscribeJob) error { return nil }
func (r *stubJobRepo) Save(*entity.TranscribeJob, string) error { return nil }
func (r *stubJobRepo) GetByID(string) (*entity.TranscribeJob, error) {
func (r *stubJobRepo) GetByID(string, string) (*entity.TranscribeJob, error) {
return nil, errors.New("не зовётся этими проверками")
}
+105
View File
@@ -0,0 +1,105 @@
package service
import (
"strings"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"git.vakhrushev.me/av/transcriber/internal/contract"
"git.vakhrushev.me/av/transcriber/internal/entity"
)
// Выборка воркера владельцем не сужается: владелец решает, кому запись
// показывать, а не кому её считать. Сужение остановило бы расшифровку записей
// бота вовсе, а записи остальных поставило бы в зависимость от того, кто первым
// завёл учётную запись.
func TestWorkerTakesJobsOfEveryOwner(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
first, err := env.service.CreateJobFromApi(t.Context(),
strings.NewReader("первая"), "one.mp3", newOwner(t, env.app))
require.NoError(t, err)
second, err := env.service.CreateJobFromApi(t.Context(),
strings.NewReader("вторая"), "two.mp3", newOwner(t, env.app))
require.NoError(t, err)
// Третья пришла ботом, и владельца у неё нет вовсе.
third := newTelegramJob(t, env)
// Срок протухания в прошлом: захваченная задача остаётся за держателем, и
// следующий вызов берёт следующую, а не ту же самую.
taken := map[string]bool{}
for _, holder := range []string{"one", "two", "three"} {
job, err := env.jobRepo.FindAndAcquire(entity.StateCreated, holder, time.Now().Add(-time.Hour))
require.NoError(t, err, "воркер берёт задачи подряд, владельцем не сужаясь")
taken[job.Id] = true
}
assert.True(t, taken[first.Id], "задача первого владельца досталась воркеру")
assert.True(t, taken[second.Id], "и второго")
assert.True(t, taken[third.Id], "и задача без владельца")
_, err = env.jobRepo.FindAndAcquire(entity.StateCreated, "next", time.Now().Add(-time.Hour))
var missing *contract.JobNotFoundError
assert.ErrorAs(t, err, &missing, "больше в этом состоянии никого")
}
// Захват читает владельца: снимок задачи, в котором он всегда пуст, был бы
// ловушкой для первого же шага, начавшего сохранять задачу целиком, — и сегодня
// уже ломал бы файл, который шаг заводит.
func TestAcquireCarriesOwner(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
owner := newOwner(t, env.app)
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", owner)
require.NoError(t, err)
acquired, err := env.jobRepo.FindAndAcquire(entity.StateCreated, "holder", time.Now().Add(-time.Hour))
require.NoError(t, err)
require.Equal(t, job.Id, acquired.Id)
require.NotNil(t, acquired.OwnerID, "владелец приехал из захвата")
assert.Equal(t, owner, *acquired.OwnerID)
}
// Шаг конвейера владельца не затирает: сохранение кладёт только то, чем
// распоряжается конвейер, и владельца среди этого нет.
func TestPipelineStepKeepsOwner(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
owner := newOwner(t, env.app)
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", owner)
require.NoError(t, err)
// Отказ конвертации — приговор записи: шаг переводит задачу в `failed` и
// сохраняет её. Это сохранение владельца тронуть не должно.
require.NoError(t, env.service.FindAndRunConversionJob(t.Context()))
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
require.Equal(t, entity.StateFailed, after.State, "шаг записал свой приговор")
require.NotNil(t, after.OwnerID, "владелец пережил шаг")
assert.Equal(t, owner, *after.OwnerID)
}
// Приём из веба без владельца задачи не заводит. Обязательность держит здесь
// код, а не схема: колонка допускает пустое значение ради записей бота.
func TestCreateJobFromApiRequiresOwner(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "one.mp3", "")
require.ErrorIs(t, err, contract.ErrOwnerRequired)
records, err := env.app.FindAllRecords("transcribe_jobs")
require.NoError(t, err)
assert.Empty(t, records, "задачи не заведено")
files, err := env.app.FindAllRecords("files")
require.NoError(t, err)
assert.Empty(t, files, "и файла тоже: отказ наступает раньше укладки")
}
+80 -9
View File
@@ -11,6 +11,7 @@ import (
"testing"
"time"
"github.com/google/uuid"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/types"
"github.com/stretchr/testify/assert"
@@ -155,7 +156,7 @@ func TestJobDiesAfterAttemptLimit(t *testing.T) {
var noop *contract.NoopJobError
require.ErrorAs(t, err, &noop, "мёртвая задача шагу не отдаётся")
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDead, after.State, "задача видна отбором по состоянию")
assert.Greater(t, after.Attempts, maxAttempts, "число попыток сохранено")
@@ -201,7 +202,7 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
// Ссылку переставляем на запись без содержимого: шаг отказывает на получении
// рабочей копии — то есть отказом, а не приговором записи.
empty, err := env.fileRepo.CreateRemote("object-key", 1)
empty, err := env.fileRepo.CreateRemote("object-key", 1, "")
require.NoError(t, err)
record, err := env.app.FindRecordById(migrations.JobsCollection, job.Id)
@@ -212,7 +213,7 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
// Первый отказ.
require.Error(t, env.service.FindAndRunConversionJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
require.Nil(t, after.AcquisitionID, "захват снят: задача пригодна к повтору")
require.NotNil(t, after.DelayTime, "пауза поставлена")
@@ -223,7 +224,7 @@ func TestFailedStepSchedulesRetryWithGrowingDelay(t *testing.T) {
clearDelay(t, env, job.Id)
require.Error(t, env.service.FindAndRunConversionJob(t.Context()))
after, err = env.jobRepo.GetByID(job.Id)
after, err = readJob(env.app, job.Id)
require.NoError(t, err)
require.NotNil(t, after.DelayTime)
@@ -248,7 +249,7 @@ func TestWorkFileRemovedAfterIntakeFailure(t *testing.T) {
env := newPipelineEnv(t, &failingMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3")
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env.app))
require.Error(t, err, "отказ источника метаданных роняет приём")
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
@@ -264,7 +265,7 @@ func TestWorkFileRemovedAfterSuccessfulIntake(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3")
_, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "sample.mp3", newOwner(t, env.app))
require.NoError(t, err)
leftovers, err := filepath.Glob(filepath.Join(tempDir, "transcriber-*"))
@@ -283,7 +284,7 @@ func TestJobNeverPointsToMissingFile(t *testing.T) {
// исходную запись, а не на несозданный результат.
require.NoError(t, env.service.FindAndRunConversionJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateFailed, after.State)
require.NotNil(t, after.FileID)
@@ -299,7 +300,7 @@ func TestStoredContentSurvivesRoundTrip(t *testing.T) {
content := strings.Repeat("запись ", 1000)
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader(content), "sample.mp3")
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader(content), "sample.mp3", newOwner(t, env.app))
require.NoError(t, err)
require.NotNil(t, job.FileID)
@@ -322,7 +323,7 @@ func TestStoredContentSurvivesRoundTrip(t *testing.T) {
func TestLocalizeGivesReadableCopy(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("содержимое"), "sample.mp3")
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("содержимое"), "sample.mp3", newOwner(t, env.app))
require.NoError(t, err)
require.NotNil(t, job.FileID)
@@ -362,3 +363,73 @@ func TestWorkFilesRemovedAfterConversionFailure(t *testing.T) {
require.NoError(t, err)
assert.Empty(t, leftovers, "ни исходной копии, ни копии под результат не осталось")
}
// readJob читает задачу мимо сужения владельцем.
//
// Читающий метод хранилища отдаёт задачу только её владельцу, а проверки
// конвейера смотрят задачи, принятые ботом: владельца у таких нет вовсе, и по
// правилу разграничения они не достаются никому. Проверке нужен не доступ, а
// состояние записи после шага, поэтому она берёт его прямо из хранилища. Это
// не второй способ читать задачу в сервисе — в сервисе способ по-прежнему один.
func readJob(app core.App, id string) (*entity.TranscribeJob, error) {
record, err := app.FindRecordById(migrations.JobsCollection, id)
if err != nil {
return nil, err
}
job := &entity.TranscribeJob{
Id: record.Id,
State: record.GetString("state"),
Source: record.GetString("source"),
Attempts: record.GetInt("attempts"),
CreatedAt: record.GetDateTime("created").Time(),
UpdatedAt: record.GetDateTime("updated").Time(),
}
for _, field := range []struct {
name string
dst **string
}{
{"owner", &job.OwnerID},
{"file", &job.FileID},
{"error_text", &job.ErrorText},
{"acquisition_id", &job.AcquisitionID},
{"recognition_op_id", &job.RecognitionOpID},
{"transcription_text", &job.TranscriptionText},
} {
if value := record.GetString(field.name); value != "" {
stored := value
*field.dst = &stored
}
}
if delay := record.GetDateTime("delay_time"); !delay.IsZero() {
moment := delay.Time()
job.DelayTime = &moment
}
if acquired := record.GetDateTime("acquire_time"); !acquired.IsZero() {
moment := acquired.Time()
job.AcquireTime = &moment
}
return job, nil
}
// newOwner заводит учётную запись и отдаёт её идентификатор.
//
// Владелец — связь с коллекцией пользователей, и хранилище проверяет, что такая
// запись есть: выдуманный идентификатор задачу завести не даст.
func newOwner(t *testing.T, app core.App) string {
t.Helper()
users, err := app.FindCollectionByNameOrId(migrations.UsersCollection)
require.NoError(t, err)
record := core.NewRecord(users)
record.Set("email", uuid.NewString()+"@example.test")
record.Set("verified", true)
record.SetPassword(uuid.NewString())
require.NoError(t, app.Save(record))
return record.Id
}
+7 -7
View File
@@ -98,7 +98,7 @@ func TestTranscribeJobHandsRecordOverAndMovesOn(t *testing.T) {
assert.Equal(t, 1, rec.recognizeCalls, "содержимое отдано распознавателю")
assert.NotEmpty(t, rec.lastObjectKey, "ключ объекта назван")
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateTranscribe, after.State)
require.NotNil(t, after.RecognitionOpID)
@@ -122,7 +122,7 @@ func TestTranscribeJobKeepsJobRetryableOnRecognizerFailure(t *testing.T) {
require.Error(t, svc.FindAndRunTranscribeJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateConverted, after.State, "задача осталась на своём шаге")
assert.Nil(t, after.AcquisitionID, "захват снят: задача пригодна к повтору")
@@ -153,7 +153,7 @@ func TestCheckJobWaitsWithoutSpendingAttempts(t *testing.T) {
for i := 0; i < 3; i++ {
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateTranscribe, after.State)
assert.Equal(t, 0, after.Attempts, "ожидание операции попытку не тратит")
@@ -177,7 +177,7 @@ func TestCheckJobFailsJobAndTellsSender(t *testing.T) {
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateFailed, after.State)
@@ -199,7 +199,7 @@ func TestCheckJobCompletesAndAnswersOnce(t *testing.T) {
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDone, after.State)
require.NotNil(t, after.TranscriptionText)
@@ -227,7 +227,7 @@ func TestCheckJobCompletesEmptyTextWithExplanation(t *testing.T) {
require.NoError(t, svc.FindAndRunTranscribeCheckJob(t.Context()))
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDone, after.State)
@@ -261,7 +261,7 @@ func TestCheckJobWritesNothingWhenAcquisitionLost(t *testing.T) {
var lost *contract.LostAcquisitionError
require.ErrorAs(t, err, &lost)
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateTranscribe, after.State, "результат не записан")
assert.Empty(t, env.sender.messages, "отправителю ничего не отправлено")
+2 -2
View File
@@ -45,7 +45,7 @@ func TestShutdownDuringConversionKeepsJobRetryable(t *testing.T) {
require.Error(t, err, "шаг обязан сообщить об обрыве наверх")
require.ErrorIs(t, err, context.Canceled, "обрыв узнаётся по смыслу, а не по тексту")
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateCreated, after.State, "задача осталась на повтор, а не похоронена")
@@ -72,7 +72,7 @@ func TestShutdownBeforeStepLeavesJobUntouched(t *testing.T) {
var noop *contract.NoopJobError
require.ErrorAs(t, err, &noop)
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateCreated, after.State)
assert.Equal(t, 0, after.Attempts, "захвата не было — попытке взяться неоткуда")
+31 -6
View File
@@ -87,10 +87,24 @@ func (s *TranscribeService) CreateJobFromTelegram(ctx context.Context, file io.R
return s.createTranscribeJob(ctx, job, file, fileName)
}
func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader, fileName string) (*entity.TranscribeJob, error) {
// CreateJobFromApi заводит задачу от имени вошедшего. Владелец обязателен:
// пустой отвергается здесь, потому что колонка владельца допускает пустое
// значение ради записей бота, и приём по HTTP — то место, где обязательность
// держится.
//
// Отказ этот — последний рубеж, а не первый: предъявителя без учётной записи
// пользователя транспорт отвергает раньше, до чтения тела. Здесь он остаётся на
// случай нового вызывающего, который такой проверки не поставит.
func (s *TranscribeService) CreateJobFromApi(ctx context.Context, file io.Reader, fileName, ownerID string) (*entity.TranscribeJob, error) {
if ownerID == "" {
s.logger.Error("Refusing to create job without owner")
return nil, contract.ErrOwnerRequired
}
job := &entity.TranscribeJob{
State: entity.StateCreated,
Source: entity.SourceApi,
State: entity.StateCreated,
Source: entity.SourceApi,
OwnerID: &ownerID,
}
return s.createTranscribeJob(ctx, job, file, fileName)
@@ -135,7 +149,7 @@ func (s *TranscribeService) createTranscribeJob(ctx context.Context, job *entity
return nil, err
}
fileRecord, err := s.fileRepo.CreateLocal(storageFileName, work)
fileRecord, err := s.fileRepo.CreateLocal(storageFileName, work, ownerOf(job))
if err != nil {
s.logger.Error("Failed to create file record", "error", err, "file_ext", ext)
return nil, err
@@ -284,7 +298,7 @@ func (s *TranscribeService) convertJob(ctx context.Context, job *entity.Transcri
metrics.OutputFileSizeHistogram.WithLabelValues("ogg").Observe(float64(destSize))
destFileName := fmt.Sprintf("%s%s", uuid.NewString(), ".ogg")
destFileRecord, err := s.fileRepo.CreateLocal(destFileName, dest)
destFileRecord, err := s.fileRepo.CreateLocal(destFileName, dest, ownerOf(job))
if err != nil {
s.logger.Error("Failed to create converted file record", "error", err, "job_id", job.Id)
return err
@@ -346,7 +360,7 @@ func (s *TranscribeService) transcribeJob(ctx context.Context, job *entity.Trans
"job_id", job.Id,
"operation_id", operationID)
destFileRecord, err := s.fileRepo.CreateRemote(fileRecord.FileName, fileRecord.Size)
destFileRecord, err := s.fileRepo.CreateRemote(fileRecord.FileName, fileRecord.Size, ownerOf(job))
if err != nil {
s.logger.Error("Failed to create S3 file record", "error", err, "job_id", job.Id)
return err
@@ -635,3 +649,14 @@ func (s *TranscribeService) closeWork(work contract.WorkFile) {
s.logger.Error("Failed to remove work file", "error", err)
}
}
// ownerOf — владелец задачи строкой; пустая значит «владельца нет», и таковы
// записи, принятые ботом. Файл наследует владельца своей задачи: правило
// просмотра коллекции файлов сужено этой колонкой, и файл, заведённый шагом
// конвейера без неё, перестал бы доставаться собственному владельцу.
func ownerOf(job *entity.TranscribeJob) string {
if job.OwnerID == nil {
return ""
}
return *job.OwnerID
}
+3 -3
View File
@@ -71,7 +71,7 @@ func TestUndeliveredOnDownChannelKeepsJobDone(t *testing.T) {
assert.Equal(t, 1, sender.calls, "ответ до отправителя доехал")
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDone, after.State, "задача осталась в достигнутом состоянии")
require.NotNil(t, after.TranscriptionText)
@@ -108,7 +108,7 @@ func TestUndeliveredWithoutChatKeepsJobDone(t *testing.T) {
assert.Equal(t, 0, sender.calls, "до отправителя дело не дошло: адресата нет")
after, err := env.jobRepo.GetByID(job.Id)
after, err := readJob(env.app, job.Id)
require.NoError(t, err)
assert.Equal(t, entity.StateDone, after.State)
assert.Nil(t, after.ErrorText, "отказ задаче не приписан")
@@ -127,7 +127,7 @@ func TestUndeliveredWithoutChatKeepsJobDone(t *testing.T) {
func TestApiJobDoesNotReachSenderAndLogsNothing(t *testing.T) {
env := newPipelineEnv(t, &okMetaViewer{}, &failingConverter{})
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "voice.ogg")
job, err := env.service.CreateJobFromApi(t.Context(), strings.NewReader("запись"), "voice.ogg", newOwner(t, env.app))
require.NoError(t, err)
sender := &downSender{}