- записи, метаданные и файлы съехались под один каталог данных; появилась панель владельца, а gin, goqu, goose и требование CGO ушли - захват задачи стал одним запросом с RETURNING; заведены число попыток, состояние dead и нарастающая пауза вместо признака is_error - имя файла в хранилище задаёт сервис и в журнал не идёт: вместе с идентификатором записи оно собирало бы ссылку на скачивание
204 lines
6.9 KiB
Go
204 lines
6.9 KiB
Go
package pocketbase
|
|
|
|
import (
|
|
"database/sql"
|
|
"time"
|
|
|
|
"github.com/pocketbase/pocketbase/core"
|
|
"github.com/pocketbase/pocketbase/tools/types"
|
|
|
|
"git.vakhrushev.me/av/transcriber/internal/entity"
|
|
)
|
|
|
|
// Отображение задачи в запись коллекции и обратно живёт одним местом. Прежде
|
|
// список колонок был переписан четырежды — в каждом запросе своего слоя, — и
|
|
// расхождение проявлялось как потерянное при сохранении поле.
|
|
|
|
// applyOwnedByPipeline кладёт в запись только те поля, которыми распоряжается
|
|
// конвейер. Поля, которые он не меняет никогда — куда отвечать отправителю и
|
|
// каким входом пришла запись, — не трогаются вовсе.
|
|
//
|
|
// Разрез нужен потому, что шаг держит задачу снимком с момента захвата и до
|
|
// своего сохранения, а это до восьми часов. Всё, что владелец правил в панели за
|
|
// это время, безусловная запись снимка стёрла бы молча: ни строки в журнале, ни
|
|
// отказа в панели — владелец видел бы успешное сохранение и был бы уверен, что
|
|
// правка на месте.
|
|
func applyOwnedByPipeline(record *core.Record, job *entity.TranscribeJob) {
|
|
record.Set("state", job.State)
|
|
record.Set("file", derefString(job.FileID))
|
|
record.Set("error_text", derefString(job.ErrorText))
|
|
record.Set("acquisition_id", derefString(job.AcquisitionID))
|
|
record.Set("acquire_time", dateOrEmpty(job.AcquireTime))
|
|
record.Set("delay_time", dateOrEmpty(job.DelayTime))
|
|
record.Set("attempts", job.Attempts)
|
|
record.Set("recognition_op_id", derefString(job.RecognitionOpID))
|
|
record.Set("transcription_text", derefString(job.TranscriptionText))
|
|
}
|
|
|
|
// applyToRecord кладёт задачу в запись целиком — это заведение, и спорить за
|
|
// поля здесь не с кем.
|
|
func applyToRecord(record *core.Record, job *entity.TranscribeJob) {
|
|
applyOwnedByPipeline(record, job)
|
|
record.Set("source", job.Source)
|
|
record.Set("tg_chat_id", derefInt64(job.TgChatId))
|
|
record.Set("tg_reply_message_id", derefInt(job.TgReplyMessageId))
|
|
}
|
|
|
|
func recordToJob(record *core.Record) *entity.TranscribeJob {
|
|
return &entity.TranscribeJob{
|
|
Id: record.Id,
|
|
State: record.GetString("state"),
|
|
Source: record.GetString("source"),
|
|
FileID: nilIfEmpty(record.GetString("file")),
|
|
ErrorText: nilIfEmpty(record.GetString("error_text")),
|
|
AcquisitionID: nilIfEmpty(record.GetString("acquisition_id")),
|
|
AcquireTime: timeOrNil(record.GetDateTime("acquire_time")),
|
|
DelayTime: timeOrNil(record.GetDateTime("delay_time")),
|
|
Attempts: record.GetInt("attempts"),
|
|
RecognitionOpID: nilIfEmpty(record.GetString("recognition_op_id")),
|
|
TranscriptionText: nilIfEmpty(record.GetString("transcription_text")),
|
|
TgChatId: nilIfZero64(int64(record.GetInt("tg_chat_id"))),
|
|
TgReplyMessageId: nilIfZeroInt(record.GetInt("tg_reply_message_id")),
|
|
CreatedAt: record.GetDateTime("created").Time(),
|
|
UpdatedAt: record.GetDateTime("updated").Time(),
|
|
}
|
|
}
|
|
|
|
// acquiredRow — задача, прочитанная сырым запросом захвата. Колонки читаются
|
|
// именно так, потому что запрос идёт мимо записей коллекции; связь с их
|
|
// перечнем держит константа acquireColumns и тест захвата, читающий задачу
|
|
// целиком.
|
|
type acquiredRow struct {
|
|
Id string `db:"id"`
|
|
State string `db:"state"`
|
|
Source string `db:"source"`
|
|
FileID sql.NullString `db:"file"`
|
|
ErrorText sql.NullString `db:"error_text"`
|
|
AcquisitionID sql.NullString `db:"acquisition_id"`
|
|
AcquireTime sql.NullString `db:"acquire_time"`
|
|
DelayTime sql.NullString `db:"delay_time"`
|
|
Attempts int `db:"attempts"`
|
|
RecognitionOpID sql.NullString `db:"recognition_op_id"`
|
|
TranscriptionText sql.NullString `db:"transcription_text"`
|
|
TgChatId sql.NullInt64 `db:"tg_chat_id"`
|
|
TgReplyMessageId sql.NullInt64 `db:"tg_reply_message_id"`
|
|
Created sql.NullString `db:"created"`
|
|
Updated sql.NullString `db:"updated"`
|
|
}
|
|
|
|
func (r *acquiredRow) toJob() *entity.TranscribeJob {
|
|
job := &entity.TranscribeJob{
|
|
Id: r.Id,
|
|
State: r.State,
|
|
Source: r.Source,
|
|
FileID: nullToPtr(r.FileID),
|
|
ErrorText: nullToPtr(r.ErrorText),
|
|
AcquisitionID: nullToPtr(r.AcquisitionID),
|
|
AcquireTime: parseTimeOrNil(r.AcquireTime),
|
|
DelayTime: parseTimeOrNil(r.DelayTime),
|
|
Attempts: r.Attempts,
|
|
RecognitionOpID: nullToPtr(r.RecognitionOpID),
|
|
TranscriptionText: nullToPtr(r.TranscriptionText),
|
|
}
|
|
|
|
if r.TgChatId.Valid && r.TgChatId.Int64 != 0 {
|
|
chatId := r.TgChatId.Int64
|
|
job.TgChatId = &chatId
|
|
}
|
|
if r.TgReplyMessageId.Valid && r.TgReplyMessageId.Int64 != 0 {
|
|
msgId := int(r.TgReplyMessageId.Int64)
|
|
job.TgReplyMessageId = &msgId
|
|
}
|
|
if created := parseTimeOrNil(r.Created); created != nil {
|
|
job.CreatedAt = *created
|
|
}
|
|
if updated := parseTimeOrNil(r.Updated); updated != nil {
|
|
job.UpdatedAt = *updated
|
|
}
|
|
|
|
return job
|
|
}
|
|
|
|
func derefString(v *string) string {
|
|
if v == nil {
|
|
return ""
|
|
}
|
|
return *v
|
|
}
|
|
|
|
func derefInt64(v *int64) int64 {
|
|
if v == nil {
|
|
return 0
|
|
}
|
|
return *v
|
|
}
|
|
|
|
func derefInt(v *int) int {
|
|
if v == nil {
|
|
return 0
|
|
}
|
|
return *v
|
|
}
|
|
|
|
// dateOrEmpty отдаёт пустое значение вместо нулевой даты: пустая колонка даты в
|
|
// хранилище это пустая строка, и она же значит «времени нет».
|
|
func dateOrEmpty(v *time.Time) any {
|
|
if v == nil {
|
|
return ""
|
|
}
|
|
date, err := types.ParseDateTime(*v)
|
|
if err != nil {
|
|
return ""
|
|
}
|
|
return date
|
|
}
|
|
|
|
func nilIfEmpty(v string) *string {
|
|
if v == "" {
|
|
return nil
|
|
}
|
|
return &v
|
|
}
|
|
|
|
func timeOrNil(v types.DateTime) *time.Time {
|
|
if v.IsZero() {
|
|
return nil
|
|
}
|
|
t := v.Time()
|
|
return &t
|
|
}
|
|
|
|
func nilIfZero64(v int64) *int64 {
|
|
if v == 0 {
|
|
return nil
|
|
}
|
|
return &v
|
|
}
|
|
|
|
func nilIfZeroInt(v int) *int {
|
|
if v == 0 {
|
|
return nil
|
|
}
|
|
return &v
|
|
}
|
|
|
|
func nullToPtr(v sql.NullString) *string {
|
|
if !v.Valid || v.String == "" {
|
|
return nil
|
|
}
|
|
s := v.String
|
|
return &s
|
|
}
|
|
|
|
func parseTimeOrNil(v sql.NullString) *time.Time {
|
|
if !v.Valid || v.String == "" {
|
|
return nil
|
|
}
|
|
date, err := types.ParseDateTime(v.String)
|
|
if err != nil || date.IsZero() {
|
|
return nil
|
|
}
|
|
t := date.Time()
|
|
return &t
|
|
}
|