2025-04-02 18:50:03 +00:00
|
|
|
package documenttext
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"bytes"
|
|
|
|
|
"context"
|
|
|
|
|
"database/sql"
|
|
|
|
|
"errors"
|
|
|
|
|
"io"
|
|
|
|
|
"strings"
|
2025-04-03 12:13:16 +00:00
|
|
|
"time"
|
2025-04-02 18:50:03 +00:00
|
|
|
|
|
|
|
|
"queryorchestration/internal/database/repository"
|
|
|
|
|
"queryorchestration/internal/serviceconfig/build"
|
2025-04-03 12:13:16 +00:00
|
|
|
"queryorchestration/internal/serviceconfig/objectstore"
|
2025-04-02 18:50:03 +00:00
|
|
|
|
|
|
|
|
"github.com/aws/aws-sdk-go-v2/service/s3"
|
|
|
|
|
"github.com/google/uuid"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
type ProcessParams struct {
|
2025-04-03 19:17:24 +00:00
|
|
|
Bucket string
|
|
|
|
|
Key objectstore.BucketKey
|
|
|
|
|
Hash string
|
2025-04-02 18:50:03 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type Page string
|
|
|
|
|
|
|
|
|
|
func (s *Service) getPages(baseTextract io.Reader) ([]*Page, error) {
|
|
|
|
|
// For each page
|
|
|
|
|
// Look for table, form, signature
|
|
|
|
|
// If one/many found - run appropriate tier and merge in
|
|
|
|
|
// tiers?? - analyze (for tables?)
|
|
|
|
|
|
|
|
|
|
// cleanPages
|
|
|
|
|
// linkPages
|
|
|
|
|
|
|
|
|
|
return []*Page{}, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *Service) mergePages(pages []*Page) (io.Reader, error) {
|
|
|
|
|
// merge document - pages to text
|
|
|
|
|
// e.g. common tables
|
|
|
|
|
return strings.NewReader("hello"), nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *Service) processTrigger(ctx context.Context, trigger *repository.GetDocumentTextTriggerRow, params *ProcessParams) error {
|
2025-04-03 12:13:16 +00:00
|
|
|
keyStr := params.Key.String()
|
2025-04-02 18:50:03 +00:00
|
|
|
response, err := s.cfg.GetStoreClient().GetObject(ctx, &s3.GetObjectInput{
|
|
|
|
|
Bucket: ¶ms.Bucket,
|
2025-04-03 12:13:16 +00:00
|
|
|
Key: &keyStr,
|
2025-04-02 18:50:03 +00:00
|
|
|
IfMatch: ¶ms.Hash,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pages, err := s.getPages(response.Body)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
text, err := s.mergePages(pages)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err = s.storeProcess(ctx, text, trigger, params)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *Service) storeProcess(ctx context.Context, text io.Reader, trigger *repository.GetDocumentTextTriggerRow, params *ProcessParams) 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: trigger.Cleanid,
|
|
|
|
|
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) {
|
2025-04-03 19:17:24 +00:00
|
|
|
createdAt := time.Now().UTC()
|
2025-04-03 12:13:16 +00:00
|
|
|
key := objectstore.BucketKey{
|
2025-04-03 19:17:24 +00:00
|
|
|
ClientID: params.Key.ClientID,
|
2025-04-03 12:13:16 +00:00
|
|
|
Location: objectstore.TextOut,
|
2025-04-03 19:17:24 +00:00
|
|
|
CreatedAt: createdAt,
|
|
|
|
|
EntityID: trigger.ID,
|
2025-04-03 12:13:16 +00:00
|
|
|
}
|
|
|
|
|
keyStr := key.String()
|
2025-04-02 18:50:03 +00:00
|
|
|
|
|
|
|
|
textId, err = q.AddDocumentText(ctx, &repository.AddDocumentTextParams{
|
|
|
|
|
Bucket: params.Bucket,
|
2025-04-03 12:13:16 +00:00
|
|
|
Key: keyStr,
|
2025-04-02 18:50:03 +00:00
|
|
|
Hash: hash,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
_, err = s.cfg.GetStoreClient().PutObject(ctx, &s3.PutObjectInput{
|
|
|
|
|
Bucket: ¶ms.Bucket,
|
2025-04-03 12:13:16 +00:00
|
|
|
Key: &keyStr,
|
2025-04-02 18:50:03 +00:00
|
|
|
Body: storeReader,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
} else {
|
|
|
|
|
textId = existingId
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err = q.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{
|
|
|
|
|
Textid: textId,
|
|
|
|
|
Triggerid: trigger.ID,
|
|
|
|
|
Version: version,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *Service) Process(ctx context.Context, params *ProcessParams) error {
|
2025-04-03 19:17:24 +00:00
|
|
|
trigger, err := s.cfg.GetDBQueries().GetDocumentTextTrigger(ctx, params.Key.EntityID)
|
2025-04-02 18:50:03 +00:00
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err = s.processTrigger(ctx, trigger, params)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err = s.informExtraction(ctx, trigger.Documentid)
|
|
|
|
|
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
|
|
|
|
|
}
|