Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 44 additions & 9 deletions broker/api/api-handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,7 @@ func (a *ApiHandler) Get(w http.ResponseWriter, r *http.Request) {
index.Links.PeersLink = Link(r, Path(PEERS_PATH), nil)
index.Links.BorrowingRequestsLink = Link(r, Path(PATRON_REQUESTS_PATH), Query("side", "borrowing"))
index.Links.LendingRequestsLink = Link(r, Path(PATRON_REQUESTS_PATH), Query("side", "lending"))
index.Links.BatchActionsLink = Link(r, Path("batch_actions"), nil)
WriteJsonResponse(w, index)
}

Expand All @@ -109,29 +110,63 @@ func (a *ApiHandler) GetEvents(w http.ResponseWriter, r *http.Request, params oa
ctx := common.CreateExtCtxWithArgs(r.Context(), &common.LoggerArgs{
Other: logParams,
})
if params.IllTransactionId != nil && events.IsSyntheticID(*params.IllTransactionId) {
AddBadRequestError(ctx, w, errors.New("synthetic IDs are not allowed for event lookup"))
return
}
tran, err := a.getIllTranFromParams(ctx, w, r, params.RequesterSymbol,
params.RequesterReqId, params.IllTransactionId)
if err != nil {
return
}
var resp oapi.Events
resp.Items = make([]oapi.Event, 0)
if tran == nil {
WriteJsonResponse(w, resp)
WriteJsonResponse(w, oapi.Events{Items: make([]oapi.Event, 0)})
return
}
var fullCount int64
var eventList []events.Event
eventList, fullCount, err = a.eventRepo.GetIllTransactionEvents(ctx, tran.ID)
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
resp, err := a.illTransactionEventsResponse(ctx, tran.ID)
if err != nil {
AddInternalError(ctx, w, err)
return
}
WriteJsonResponse(w, resp)
}

func (a *ApiHandler) GetIllTransactionsIdEvents(w http.ResponseWriter, r *http.Request, id string, params oapi.GetIllTransactionsIdEventsParams) {
ctx := common.CreateExtCtxWithArgs(r.Context(), &common.LoggerArgs{
Other: map[string]string{"method": "GetIllTransactionsIdEvents", "id": id},
})
if events.IsSyntheticID(id) {
AddBadRequestError(ctx, w, errors.New("synthetic IDs are not allowed for event lookup"))
return
}

tran, err := a.getIllTranFromParams(ctx, w, r, params.RequesterSymbol, nil, &id)
if err != nil {
return
}
if tran == nil {
AddNotFoundError(w)
return
}
resp, err := a.illTransactionEventsResponse(ctx, tran.ID)
if err != nil {
AddInternalError(ctx, w, err)
return
}
WriteJsonResponse(w, resp)
}

func (a *ApiHandler) illTransactionEventsResponse(ctx common.ExtendedContext, id string) (oapi.Events, error) {
resp := oapi.Events{Items: make([]oapi.Event, 0)}
eventList, fullCount, err := a.eventRepo.GetIllTransactionEvents(ctx, id)
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
return resp, err
}
resp.About.Count = fullCount
for _, event := range eventList {
resp.Items = append(resp.Items, ToApiEvent(event, event.IllTransactionID, nil))
}
WriteJsonResponse(w, resp)
return resp, nil
}

func (a *ApiHandler) GetIllTransactions(w http.ResponseWriter, r *http.Request, params oapi.GetIllTransactionsParams) {
Expand Down Expand Up @@ -652,7 +687,7 @@ func toApiIllTransaction(r *http.Request, trans ill_db.IllTransaction) oapi.IllT
api.SupplierRequestID = getString(trans.SupplierRequestID)
api.LastSupplierStatus = getString(trans.LastSupplierStatus)
api.PrevSupplierStatus = getString(trans.PrevSupplierStatus)
api.EventsLink = Link(r, Path(EVENTS_PATH), Query("ill_transaction_id", trans.ID))
api.EventsLink = Link(r, Path("ill_transactions", trans.ID, "events"), nil)
api.LocatedSuppliersLink = Link(r, Path(LOCATED_SUPPLIERS_PATH), Query("ill_transaction_id", trans.ID))
if trans.RequesterID.Valid {
api.RequesterPeerLink = Link(r, Path(PEERS_PATH, trans.RequesterID.String), nil)
Expand Down
35 changes: 35 additions & 0 deletions broker/api/api_handler_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
package api

import (
"net/http"
"net/http/httptest"
"testing"

"github.com/indexdata/crosslink/broker/events"
"github.com/indexdata/crosslink/broker/oapi"
"github.com/stretchr/testify/assert"
)

func TestGetEventsRejectsSyntheticIDs(t *testing.T) {
for _, id := range []string{events.DEFAULT_ILL_TRANSACTION_ID, events.DEFAULT_PATRON_REQUEST_ID} {
t.Run(id, func(t *testing.T) {
h := ApiHandler{}
req := httptest.NewRequest(http.MethodGet, "/events", nil)
rr := httptest.NewRecorder()
h.GetEvents(rr, req, oapi.GetEventsParams{IllTransactionId: &id})
assert.Equal(t, http.StatusBadRequest, rr.Code)
})
}
}

func TestGetIllTransactionsIdEventsRejectsSyntheticIDs(t *testing.T) {
for _, id := range []string{events.DEFAULT_ILL_TRANSACTION_ID, events.DEFAULT_PATRON_REQUEST_ID} {
t.Run(id, func(t *testing.T) {
h := ApiHandler{}
req := httptest.NewRequest(http.MethodGet, "/ill_transactions/"+id+"/events", nil)
rr := httptest.NewRecorder()
h.GetIllTransactionsIdEvents(rr, req, id, oapi.GetIllTransactionsIdEventsParams{})
assert.Equal(t, http.StatusBadRequest, rr.Code)
})
}
}
2 changes: 1 addition & 1 deletion broker/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -219,7 +219,7 @@ func Init(ctx context.Context) (Context, error) {
}

schedRepoRepo := sched_db.CreateSchedRepo(pool)
schedApiHandler := schedapi.NewSchedulerApiHandler(API_PAGE_SIZE, schedRepoRepo, tenantResolver)
schedApiHandler := schedapi.NewSchedulerApiHandler(API_PAGE_SIZE, schedRepoRepo, eventRepo, tenantResolver)
if err = StartScheduler(ctx, schedRepoRepo, eventBus); err != nil {
return Context{}, err
}
Expand Down
34 changes: 33 additions & 1 deletion broker/descriptors/ModuleDescriptor-template.json
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,15 @@
"broker.ill_transactions.item.get"
]
},
{
"methods": [
"GET"
],
"pathPattern": "/broker/ill_transactions/{id}/events",
"permissionsRequired": [
"broker.ill_transactions.item.events.get"
]
},
{
"methods": [
"GET"
Expand Down Expand Up @@ -267,6 +276,15 @@
"broker.batch_actions.item.get"
]
},
{
"methods": [
"GET"
],
"pathPattern": "/broker/batch_actions/{id}/events",
"permissionsRequired": [
"broker.batch_actions.item.events.get"
]
},
{
"methods": [
"PUT"
Expand Down Expand Up @@ -360,6 +378,7 @@
"visible": true,
"subPermissions": [
"broker.ill_transactions.item.get",
"broker.ill_transactions.item.events.get",
"broker.ill_transactions.get",
"broker.located_suppliers.get",
"broker.events.get",
Expand All @@ -382,6 +401,11 @@
"displayName": "Broker - read ILL transactions",
"permissionName": "broker.ill_transactions.get"
},
{
"description": "Read events for an ILL transaction",
"displayName": "Broker - read ILL transaction events",
"permissionName": "broker.ill_transactions.item.events.get"
},
{
"description": "Read located suppliers",
"displayName": "Broker - read located suppliers",
Expand Down Expand Up @@ -497,6 +521,11 @@
"displayName": "Broker - read batch action",
"permissionName": "broker.batch_actions.item.get"
},
{
"description": "Read events for a batch action",
"displayName": "Broker - read batch action events",
"permissionName": "broker.batch_actions.item.events.get"
},
{
"description": "Update a batch action",
"displayName": "Broker - update batch action",
Expand Down Expand Up @@ -580,7 +609,8 @@
"visible": true,
"subPermissions": [
"broker.batch_actions.get",
"broker.batch_actions.item.get"
"broker.batch_actions.item.get",
"broker.batch_actions.item.events.get"
]
},
{
Expand Down Expand Up @@ -673,6 +703,7 @@
"visible": true,
"subPermissions": [
"broker.ill_transactions.item.get",
"broker.ill_transactions.item.events.get",
"broker.ill_transactions.get",
"broker.located_suppliers.get",
"broker.events.get",
Expand Down Expand Up @@ -701,6 +732,7 @@
"broker.batch_actions.get",
"broker.batch_actions.post",
"broker.batch_actions.item.get",
"broker.batch_actions.item.events.get",
"broker.batch_actions.item.put",
"broker.batch_actions.item.delete",
"broker.batch_actions.item.enable.post",
Expand Down
4 changes: 4 additions & 0 deletions broker/events/eventbus_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,10 @@ func (r *exclusiveCheckErrorRepo) GetIllTransactionEvents(ctx common.ExtendedCon
return nil, 0, nil
}

func (r *exclusiveCheckErrorRepo) GetBatchActionEvents(ctx common.ExtendedContext, taskID string) ([]Event, error) {
return nil, nil
}

func (r *exclusiveCheckErrorRepo) DeleteEventsByIllTransaction(ctx common.ExtendedContext, illTransId string) error {
return nil
}
Expand Down
4 changes: 4 additions & 0 deletions broker/events/eventmodels.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@ const (
EventDomainScheduler EventDomain = "SCHEDULER"
)

func IsSyntheticID(id string) bool {
return id == DEFAULT_ILL_TRANSACTION_ID || id == DEFAULT_PATRON_REQUEST_ID
}

type EventName string

const (
Expand Down
12 changes: 12 additions & 0 deletions broker/events/eventrepo.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ type EventRepo interface {
ClaimEventForSignal(ctx common.ExtendedContext, id string, signal Signal) (Event, error)
Notify(ctx common.ExtendedContext, eventId string, signal Signal, target SignalTarget) error
GetIllTransactionEvents(ctx common.ExtendedContext, id string) ([]Event, int64, error)
GetBatchActionEvents(ctx common.ExtendedContext, taskID string) ([]Event, error)
DeleteEventsByIllTransaction(ctx common.ExtendedContext, illTransId string) error
GetLatestRequestEventByAction(ctx common.ExtendedContext, illTransId string, action string) (Event, error)
GetPatronRequestEvents(ctx common.ExtendedContext, id string) ([]Event, error)
Expand Down Expand Up @@ -117,6 +118,17 @@ func (r *PgEventRepo) GetPatronRequestEvents(ctx common.ExtendedContext, id stri
return events, err
}

func (r *PgEventRepo) GetBatchActionEvents(ctx common.ExtendedContext, taskID string) ([]Event, error) {
rows, err := r.queries.GetBatchActionEvents(ctx, r.GetConnOrTx(), taskID)
var eventList []Event
if err == nil {
for _, row := range rows {
eventList = append(eventList, row.Event)
}
}
return eventList, err
}

func (r *PgEventRepo) DeleteEventsByIllTransaction(ctx common.ExtendedContext, illTransId string) error {
return r.queries.DeleteEventsByIllTransaction(ctx, r.GetConnOrTx(), illTransId)
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
DROP INDEX IF EXISTS idx_event_batch_action_task_timestamp;
3 changes: 3 additions & 0 deletions broker/migrations/051_add_batch_action_event_index.up.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
CREATE INDEX idx_event_batch_action_task_timestamp
ON event ((event_data -> 'batchActionData' ->> 'taskId'), timestamp DESC)
WHERE event_name = 'invoke-batch-action';
Loading
Loading