Skip to content

Commit 74d19f8

Browse files
CROSSLINK-266 Use CQL to search duplicates
1 parent e44d921 commit 74d19f8

8 files changed

Lines changed: 144 additions & 84 deletions

File tree

broker/handler/iso18626-handler.go

Lines changed: 36 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -11,9 +11,10 @@ import (
1111
"sync"
1212
"time"
1313

14+
"github.com/indexdata/cql-go/cqlbuilder"
1415
"github.com/indexdata/crosslink/broker/catalog"
15-
"github.com/indexdata/crosslink/broker/service"
1616
"github.com/indexdata/crosslink/broker/shim"
17+
"github.com/indexdata/crosslink/directory"
1718

1819
"github.com/indexdata/crosslink/broker/adapter"
1920

@@ -34,6 +35,13 @@ var brokerSymbol = utils.GetEnv("BROKER_SYMBOL", "ISIL:BROKER")
3435
const HANDLER_COMP = "iso18626_handler"
3536
const ORIGINAL_INCOMING_MESSAGE = "originalIncomingMessage"
3637

38+
var lookupQueryBuilder = utils.Must(catalog.NewQueryBuilderGen(&directory.QueryConfig{
39+
Identifier: new("supplier_unique_record_id = {term}"),
40+
Type: new(directory.Cql),
41+
}))
42+
43+
var queryTimeFormat = "2006-01-02 15:04:05"
44+
3745
type ErrorValue string
3846

3947
const (
@@ -201,23 +209,37 @@ func checkDuplicateRequest(ctx common.ExtendedContext, request *iso18626.Request
201209
return nil
202210
}
203211

204-
_, err := repo.FindDuplicateIllTransaction(ctx, ill_db.FindDuplicateIllTransactionParams{
205-
RequesterSymbol: createPgText(requesterSymbol),
206-
Hours: service.ToInt32(*windowHours),
207-
PatronID: patronId,
208-
Identifier: lookupParams.Identifier,
209-
Title: lookupParams.Title,
210-
ServiceType: lookupParams.ServiceType,
211-
Isbn: lookupParams.Isbn,
212-
Issn: lookupParams.Issn,
213-
})
212+
cqlList, _, err := lookupQueryBuilder.Build(lookupParams)
213+
if err != nil {
214+
ctx.Logger().Warn("failed build lookup query", "error", err)
215+
return nil
216+
}
217+
218+
lookupCql := strings.Join(cqlList, " or ")
219+
qb, err := cqlbuilder.NewQueryFromString("(" + lookupCql + ")")
220+
if err != nil {
221+
ctx.Logger().Warn("failed to build duplicate check query", "error", err)
222+
return nil
223+
}
224+
formattedTime := time.Now().Add(-time.Duration(*windowHours) * time.Hour).Format(queryTimeFormat)
225+
query, err := qb.And().Search("requester_symbol").Term(requesterSymbol).
226+
And().Search("patron_id").Term(patronId).
227+
And().Search("timestamp").Rel(">=").Term(formattedTime).
228+
And().Search("service_type").Term(lookupParams.ServiceType).
229+
Build()
230+
if err != nil {
231+
ctx.Logger().Warn("failed to build duplicate check query", "error", err)
232+
return nil
233+
}
234+
cql := query.String()
235+
trans, _, err := repo.ListIllTransactions(ctx, ill_db.ListIllTransactionsParams{Limit: 1, Offset: 0}, &cql, []string{requesterSymbol})
214236
if err != nil {
215-
if errors.Is(err, pgx.ErrNoRows) {
216-
return nil // no duplicate found
217-
}
218237
ctx.Logger().Warn("failed to check for duplicate requests, proceeding", "error", err)
219238
return nil // fail open
220239
}
240+
if len(trans) == 0 {
241+
return nil
242+
}
221243
return ErrDuplicateRequest
222244
}
223245

broker/handler/iso18626-handler_test.go

Lines changed: 51 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package handler
33
import (
44
"context"
55
"errors"
6+
"strings"
67
"testing"
78

89
"github.com/indexdata/crosslink/broker/common"
@@ -160,16 +161,21 @@ func (r *MockIllRepositoryNoSelectedSupplier) GetSelectedSupplierForIllTransacti
160161
// configurable results for duplicate-check testing.
161162
type mockDuplicateCheckRepo struct {
162163
mocks.MockIllRepositorySuccess
163-
duplicateId string
164-
err error
165-
called bool
166-
params ill_db.FindDuplicateIllTransactionParams
164+
duplicate bool
165+
err error
166+
called bool
167+
cql string
167168
}
168169

169-
func (r *mockDuplicateCheckRepo) FindDuplicateIllTransaction(ctx common.ExtendedContext, params ill_db.FindDuplicateIllTransactionParams) (string, error) {
170+
func (r *mockDuplicateCheckRepo) ListIllTransactions(ctx common.ExtendedContext, params ill_db.ListIllTransactionsParams, cql *string, symbols []string) ([]ill_db.IllTransaction, int64, error) {
171+
if cql != nil {
172+
r.cql = *cql
173+
}
170174
r.called = true
171-
r.params = params
172-
return r.duplicateId, r.err
175+
if r.duplicate {
176+
return []ill_db.IllTransaction{{ID: "duplicate-id"}}, 1, nil
177+
}
178+
return []ill_db.IllTransaction{}, 0, r.err
173179
}
174180

175181
func TestCheckDuplicateRequest(t *testing.T) {
@@ -207,12 +213,11 @@ func TestCheckDuplicateRequest(t *testing.T) {
207213
name string
208214
request *iso18626.Request
209215
peer ill_db.Peer
210-
duplicateId string
216+
duplicate bool
211217
repoErr error
212218
wantErr error
213219
wantRepoCalled bool
214220
wantPatronId string
215-
wantWindowHrs int32
216221
wantIdentifier string
217222
wantIsbn string
218223
wantIssn string
@@ -248,7 +253,6 @@ func TestCheckDuplicateRequest(t *testing.T) {
248253
wantErr: nil,
249254
wantRepoCalled: true,
250255
wantPatronId: "patron-1",
251-
wantWindowHrs: 1,
252256
wantIdentifier: "rec-1",
253257
wantTitle: "Test Title",
254258
wantSvcType: "Loan",
@@ -261,7 +265,6 @@ func TestCheckDuplicateRequest(t *testing.T) {
261265
wantErr: nil,
262266
wantRepoCalled: true,
263267
wantPatronId: "patron-1",
264-
wantWindowHrs: 1,
265268
wantIdentifier: "rec-1",
266269
wantTitle: "Test Title",
267270
wantSvcType: "Loan",
@@ -270,11 +273,10 @@ func TestCheckDuplicateRequest(t *testing.T) {
270273
name: "duplicate found - returns ErrDuplicateRequest",
271274
request: baseRequest,
272275
peer: ill_db.Peer{CustomData: directory.Entry{DuplicateCheckWindowHours: &window1}},
273-
duplicateId: "existing-tx-id",
276+
duplicate: true,
274277
wantErr: ErrDuplicateRequest,
275278
wantRepoCalled: true,
276279
wantPatronId: "patron-1",
277-
wantWindowHrs: 1,
278280
wantIdentifier: "rec-1",
279281
wantTitle: "Test Title",
280282
wantSvcType: "Loan",
@@ -290,6 +292,22 @@ func TestCheckDuplicateRequest(t *testing.T) {
290292
wantErr: nil,
291293
wantRepoCalled: false,
292294
},
295+
{
296+
name: "no service level - skips duplicate check (can't verify same patron)",
297+
request: &iso18626.Request{
298+
BibliographicInfo: iso18626.BibliographicInfo{
299+
SupplierUniqueRecordId: "rec-1",
300+
Title: "Test Title",
301+
},
302+
PatronInfo: &iso18626.PatronInfo{
303+
PatronId: "patron-1",
304+
},
305+
},
306+
peer: ill_db.Peer{CustomData: directory.Entry{DuplicateCheckWindowHours: &window1}},
307+
repoErr: pgx.ErrNoRows,
308+
wantErr: nil,
309+
wantRepoCalled: false,
310+
},
293311
{
294312
name: "isbn passed as parameter to DB query",
295313
request: isbnRequest,
@@ -298,19 +316,28 @@ func TestCheckDuplicateRequest(t *testing.T) {
298316
wantErr: nil,
299317
wantRepoCalled: true,
300318
wantPatronId: "patron-2",
301-
wantWindowHrs: 1,
302319
wantIsbn: "978-1234",
303320
wantSvcType: "Copy",
304321
},
305322
{
306323
name: "duplicate found via isbn - returns ErrDuplicateRequest",
307324
request: isbnRequest,
308325
peer: ill_db.Peer{CustomData: directory.Entry{DuplicateCheckWindowHours: &window1}},
309-
duplicateId: "existing-tx-isbn",
326+
duplicate: true,
310327
wantErr: ErrDuplicateRequest,
311328
wantRepoCalled: true,
312329
wantPatronId: "patron-2",
313-
wantWindowHrs: 1,
330+
wantIsbn: "978-1234",
331+
wantSvcType: "Copy",
332+
},
333+
{
334+
name: "no duplicate - returns nil",
335+
request: isbnRequest,
336+
peer: ill_db.Peer{CustomData: directory.Entry{DuplicateCheckWindowHours: &window1}},
337+
duplicate: false,
338+
wantErr: nil,
339+
wantRepoCalled: true,
340+
wantPatronId: "patron-2",
314341
wantIsbn: "978-1234",
315342
wantSvcType: "Copy",
316343
},
@@ -320,21 +347,19 @@ func TestCheckDuplicateRequest(t *testing.T) {
320347
t.Run(tt.name, func(t *testing.T) {
321348
appCtx := common.CreateExtCtxWithArgs(context.Background(), nil)
322349
mockRepo := &mockDuplicateCheckRepo{
323-
duplicateId: tt.duplicateId,
324-
err: tt.repoErr,
350+
duplicate: tt.duplicate,
351+
err: tt.repoErr,
325352
}
326353
err := checkDuplicateRequest(appCtx, tt.request, mockRepo, "ISIL:REQ1", tt.peer)
327354
assert.Equal(t, tt.wantErr, err)
328355
assert.Equal(t, tt.wantRepoCalled, mockRepo.called)
329356
if tt.wantRepoCalled {
330-
assert.Equal(t, "ISIL:REQ1", mockRepo.params.RequesterSymbol.String)
331-
assert.Equal(t, tt.wantPatronId, mockRepo.params.PatronID)
332-
assert.Equal(t, tt.wantWindowHrs, mockRepo.params.Hours)
333-
assert.Equal(t, tt.wantIdentifier, mockRepo.params.Identifier)
334-
assert.Equal(t, tt.wantIsbn, mockRepo.params.Isbn)
335-
assert.Equal(t, tt.wantIssn, mockRepo.params.Issn)
336-
assert.Equal(t, tt.wantTitle, mockRepo.params.Title)
337-
assert.Equal(t, tt.wantSvcType, mockRepo.params.ServiceType)
357+
assert.True(t, tt.wantPatronId == "" || strings.Contains(mockRepo.cql, tt.wantPatronId))
358+
assert.True(t, tt.wantIdentifier == "" || strings.Contains(mockRepo.cql, tt.wantIdentifier))
359+
assert.True(t, tt.wantIsbn == "" || strings.Contains(mockRepo.cql, tt.wantIsbn))
360+
assert.True(t, tt.wantIssn == "" || strings.Contains(mockRepo.cql, tt.wantIssn))
361+
assert.True(t, tt.wantTitle == "" || strings.Contains(mockRepo.cql, tt.wantTitle))
362+
assert.True(t, tt.wantSvcType == "" || strings.Contains(mockRepo.cql, tt.wantSvcType))
338363
}
339364
})
340365
}

broker/ill_db/ill_cql.go

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,12 @@ import (
88
"github.com/indexdata/cql-go/cql"
99
"github.com/indexdata/cql-go/cqlbuilder"
1010
"github.com/indexdata/cql-go/pgcql"
11+
pr_db "github.com/indexdata/crosslink/broker/patron_request/db"
12+
"github.com/indexdata/go-utils/utils"
1113
)
1214

15+
var LANGUAGE = utils.GetEnv("LANGUAGE", "english")
16+
1317
func handleIllTransactionsQuery(cqlString string, noBaseArgs int) (pgcql.Query, error) {
1418
def := pgcql.NewPgDefinition()
1519

@@ -28,6 +32,24 @@ func handleIllTransactionsQuery(cqlString string, noBaseArgs int) (pgcql.Query,
2832
f = pgcql.NewFieldString().WithExact()
2933
def.AddField("last_requester_action", f)
3034

35+
def.AddField("isbn", pr_db.NewFieldTextArrayContains("bibliographic_item_identifiers(ill_transaction_data, 'ISBN')").WithFunction("norm_isxn"))
36+
def.AddField("issn", pr_db.NewFieldTextArrayContains("bibliographic_item_identifiers(ill_transaction_data, 'ISSN')").WithFunction("norm_isxn"))
37+
38+
f = pgcql.NewFieldString().WithFullText(LANGUAGE).WithColumn("ill_transaction_data->'bibliographicInfo'->>'title'")
39+
def.AddField("title", f)
40+
41+
nf := pgcql.NewFieldDate()
42+
def.AddField("timestamp", nf)
43+
44+
f = pgcql.NewFieldString().WithLikeOps().WithColumn("ill_transaction_data->'patronInfo'->>'patronId'")
45+
def.AddField("patron_id", f)
46+
47+
f = pgcql.NewFieldString().WithExact().WithColumn("ill_transaction_data->'bibliographicInfo'->>'supplierUniqueRecordId'")
48+
def.AddField("supplier_unique_record_id", f)
49+
50+
f = pgcql.NewFieldString().WithExact().WithColumn("ill_transaction_data->'serviceInfo'->>'serviceType'")
51+
def.AddField("service_type", f)
52+
3153
var parser cql.Parser
3254
query, err := parser.Parse(cqlString)
3355
if err != nil {

broker/ill_db/illrepo.go

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,6 @@ type IllRepo interface {
5050
GetExclusiveBranchSymbolsByPeerId(ctx common.ExtendedContext, peerId string) ([]BranchSymbol, error)
5151
DeleteBranchSymbolByPeerId(ctx common.ExtendedContext, peerId string) error
5252
CallArchiveIllTransactionByDateAndStatus(ctx common.ExtendedContext, toDate time.Time, statuses []string) error
53-
FindDuplicateIllTransaction(ctx common.ExtendedContext, params FindDuplicateIllTransactionParams) (string, error)
5453
}
5554

5655
type PgIllRepo struct {
@@ -478,10 +477,6 @@ func (r *PgIllRepo) CallArchiveIllTransactionByDateAndStatus(ctx common.Extended
478477
return err
479478
}
480479

481-
func (r *PgIllRepo) FindDuplicateIllTransaction(ctx common.ExtendedContext, params FindDuplicateIllTransactionParams) (string, error) {
482-
return r.queries.FindDuplicateIllTransaction(ctx, r.GetConnOrTx(), params)
483-
}
484-
485480
func (r *PgIllRepo) GetExclusiveBranchSymbolsByPeerId(ctx common.ExtendedContext, peerId string) ([]BranchSymbol, error) {
486481
rows, err := r.queries.GetExclusiveBranchSymbolsByPeerId(ctx, r.GetConnOrTx(), peerId)
487482
var symbols []BranchSymbol
Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1,13 @@
1-
DROP INDEX IF EXISTS idx_ill_transaction_requester_timestamp;
1+
DROP INDEX IF EXISTS idx_ill_transaction_requester_timestamp;
2+
3+
DROP INDEX IF EXISTS idx_ill_transaction_isbn;
4+
5+
DROP INDEX IF EXISTS idx_ill_transaction_issn;
6+
7+
DROP INDEX idx_ill_transaction_title_tsv;
8+
9+
DROP INDEX IF EXISTS idx_ill_transaction_supplier_unique_record_id;
10+
11+
DROP INDEX IF EXISTS idx_ill_transaction_patron_id;
12+
13+
DROP INDEX IF EXISTS idx_ill_transaction_service_type;
Lines changed: 22 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,23 @@
11
CREATE INDEX IF NOT EXISTS idx_ill_transaction_requester_timestamp
2-
ON ill_transaction (requester_symbol, timestamp DESC)
3-
INCLUDE (id);
2+
ON ill_transaction (requester_symbol, timestamp)
3+
INCLUDE (id);
4+
5+
CREATE INDEX IF NOT EXISTS idx_ill_transaction_isbn
6+
ON ill_transaction
7+
USING gin (bibliographic_item_identifiers(ill_transaction_data, 'ISBN'));
8+
9+
CREATE INDEX IF NOT EXISTS idx_ill_transaction_issn
10+
ON ill_transaction
11+
USING gin (bibliographic_item_identifiers(ill_transaction_data, 'ISSN'));
12+
13+
CREATE INDEX idx_ill_transaction_title_tsv
14+
ON ill_transaction USING GIN (to_tsvector('simple', ill_transaction_data->'bibliographicInfo'->>'title'));
15+
16+
CREATE INDEX IF NOT EXISTS idx_ill_transaction_supplier_unique_record_id
17+
ON ill_transaction ((ill_transaction_data->'bibliographicInfo'->>'supplierUniqueRecordId'));
18+
19+
CREATE INDEX IF NOT EXISTS idx_ill_transaction_patron_id
20+
ON ill_transaction ((ill_transaction_data->'patronInfo'->>'patronId'));
21+
22+
CREATE INDEX IF NOT EXISTS idx_ill_transaction_service_type
23+
ON ill_transaction ((ill_transaction_data->'serviceInfo'->>'serviceType'));

broker/sqlc/ill_query.sql

Lines changed: 0 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -234,30 +234,3 @@ SELECT archive_ill_transaction_by_date_and_status($1, $2);
234234
SELECT sqlc.embed(branch_symbol)
235235
FROM branch_symbol b
236236
WHERE b.peer_id = $1 AND b.symbol_value not in (SELECT s.symbol_value FROM symbol s);
237-
238-
-- name: FindDuplicateIllTransaction :one
239-
SELECT ill_transaction.id
240-
FROM ill_transaction
241-
WHERE requester_symbol = $1
242-
AND timestamp > NOW() - make_interval(hours => sqlc.arg(hours)::int4)
243-
AND COALESCE(ill_transaction_data->'patronInfo'->>'patronId', '') = sqlc.arg(patron_id)::text
244-
AND COALESCE(ill_transaction_data->'bibliographicInfo'->>'supplierUniqueRecordId', '') = sqlc.arg(identifier)::text
245-
AND COALESCE(ill_transaction_data->'bibliographicInfo'->>'title', '') = sqlc.arg(title)::text
246-
AND COALESCE(ill_transaction_data->'serviceInfo'->>'serviceType', '') = sqlc.arg(service_type)::text
247-
AND COALESCE((
248-
SELECT elem->>'bibliographicItemIdentifier'
249-
FROM jsonb_array_elements(
250-
COALESCE(ill_transaction_data->'bibliographicInfo'->'bibliographicItemId', '[]'::jsonb)
251-
) AS elem
252-
WHERE elem->'bibliographicItemIdentifierCode'->>'#text' = 'ISBN'
253-
LIMIT 1
254-
), '') = sqlc.arg(isbn)::text
255-
AND COALESCE((
256-
SELECT elem->>'bibliographicItemIdentifier'
257-
FROM jsonb_array_elements(
258-
COALESCE(ill_transaction_data->'bibliographicInfo'->'bibliographicItemId', '[]'::jsonb)
259-
) AS elem
260-
WHERE elem->'bibliographicItemIdentifierCode'->>'#text' = 'ISSN'
261-
LIMIT 1
262-
), '') = sqlc.arg(issn)::text
263-
LIMIT 1;

0 commit comments

Comments
 (0)