fee71e7740
Feature/postprocessing * tests * passtest * fixshorttests * mosttests * improvingbasedockerfile * testspeeds * testing * host * canparallel * clean * passfullsuite * singlepagemax * test * findfeatures * findstables * tbls * tablestoo * tablestoo * lateraltests * tableloc * cleanup * inlinetable * childids * cleanup * tests
119 lines
2.6 KiB
Go
119 lines
2.6 KiB
Go
package documenttext
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"database/sql"
|
|
"errors"
|
|
"io"
|
|
"time"
|
|
|
|
"queryorchestration/internal/database/repository"
|
|
"queryorchestration/internal/serviceconfig/build"
|
|
"queryorchestration/internal/serviceconfig/objectstore"
|
|
|
|
"github.com/aws/aws-sdk-go-v2/service/s3"
|
|
"github.com/google/uuid"
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
)
|
|
|
|
func (s *Service) storeProcess(ctx context.Context, text io.Reader, clean *repository.Currentcleanentry) error {
|
|
version := build.GetVersionUnixTimestamp()
|
|
|
|
hashReader, storeReader, err := duplicateReader(text)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
hash, err := s.cfg.CalculateETag(ctx, hashReader)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return s.cfg.ExecuteDBTransaction(ctx, func(ctx context.Context, q *repository.Queries) error {
|
|
existingId, err := q.GetDocumentTextExtractionByHash(ctx, &repository.GetDocumentTextExtractionByHashParams{
|
|
Cleanentryid: clean.ID,
|
|
Hash: hash,
|
|
})
|
|
if err != nil && !errors.Is(err, sql.ErrNoRows) {
|
|
return err
|
|
}
|
|
|
|
var textId uuid.UUID
|
|
if existingId == uuid.Nil || errors.Is(err, sql.ErrNoRows) {
|
|
createdAt := time.Now().UTC()
|
|
key := objectstore.BucketKey{
|
|
ClientID: clean.Clientid,
|
|
Location: objectstore.Text,
|
|
CreatedAt: createdAt,
|
|
EntityID: clean.ID,
|
|
}
|
|
|
|
currentPart, err := q.GetTextOutCurrentPart(ctx, &repository.GetTextOutCurrentPartParams{
|
|
Clientid: &key.ClientID,
|
|
Querydate: pgtype.Timestamptz{
|
|
Time: key.CreatedAt,
|
|
Valid: true,
|
|
},
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
part := s.cfg.GetDirectoryPart(currentPart.Part, currentPart.Count)
|
|
key.Part = &part
|
|
|
|
keyStr := key.String()
|
|
|
|
textId, err = q.AddDocumentText(ctx, &repository.AddDocumentTextParams{
|
|
Cleanid: clean.ID,
|
|
Bucket: *clean.Bucket,
|
|
Key: keyStr,
|
|
Hash: hash,
|
|
Part: *key.Part,
|
|
Createdat: pgtype.Timestamp{
|
|
Time: key.CreatedAt,
|
|
Valid: true,
|
|
},
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
_, err = s.cfg.GetStoreClient().PutObject(ctx, &s3.PutObjectInput{
|
|
Bucket: clean.Bucket,
|
|
Key: &keyStr,
|
|
Body: storeReader,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
} else {
|
|
textId = existingId
|
|
}
|
|
|
|
err = q.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{
|
|
Textid: textId,
|
|
Version: version,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
})
|
|
}
|
|
|
|
func duplicateReader(reader io.Reader) (io.Reader, io.Reader, error) {
|
|
content, err := io.ReadAll(reader)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
|
|
reader1 := bytes.NewReader(content)
|
|
reader2 := bytes.NewReader(content)
|
|
|
|
return reader1, reader2, nil
|
|
}
|