Skip to content
Open
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
9 changes: 0 additions & 9 deletions internal/pkg/infrastructure/postgres/deviceinfo.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,12 +168,3 @@ func deleteDeviceInfoById(ctx context.Context, tx pgx.Tx, id int) errors.EdgeX {
}
return nil
}

// updateDeviceInfosDeletableByDeviceName updates deviceInfos by deviceName as deletable
func (c *Client) updateDeviceInfosDeletableByDeviceName(deviceName string) errors.EdgeX {
_, err := c.ConnPool.Exec(context.Background(), fmt.Sprintf("UPDATE %s SET %s = true WHERE %s=$1", deviceInfoTableName, markDeletedCol, deviceNameCol), deviceName)
if err != nil {
return pgClient.WrapDBError(fmt.Sprintf("update %s to deletable by devicename '%s'", deviceInfoTableName, deviceName), err)
}
return nil
}
100 changes: 50 additions & 50 deletions internal/pkg/infrastructure/postgres/event.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (

"github.com/edgexfoundry/edgex-go/internal/core/data/container"
pgClient "github.com/edgexfoundry/edgex-go/internal/pkg/db/postgres"
dbModels "github.com/edgexfoundry/edgex-go/internal/pkg/infrastructure/models"

"github.com/edgexfoundry/go-mod-core-contracts/v4/errors"
model "github.com/edgexfoundry/go-mod-core-contracts/v4/models"
Expand Down Expand Up @@ -200,14 +201,10 @@ func (c *Client) DeleteEventById(id string) errors.EdgeX {
return nil
}

// DeleteEventsByDeviceName deletes specific device's events and corresponding readings
// This function is implemented to starts up two goroutines to delete readings and events in the background to achieve better performance
// DeleteEventsByDeviceName deletes a device's events and corresponding readings.
// The async behavior is handled at the application layer so this database call can
// return real failures instead of hiding data behind mark_deleted on partial errors.
func (c *Client) DeleteEventsByDeviceName(deviceName string) errors.EdgeX {
// update deviceInfo as deletable, then event and reading will not return when the user query event or reading by the specific device
if err := c.updateDeviceInfosDeletableByDeviceName(deviceName); err != nil {
return errors.NewCommonEdgeX(errors.Kind(err), "delete events by deviceName", err)
}
// delete events, readings, deviceInfos
if err := c.deleteEventsByConditions([]string{deviceNameCol}, pgx.NamedArgs{deviceNameCol: deviceName}); err != nil {
return errors.NewCommonEdgeX(errors.Kind(err), "delete events by deviceName", err)
}
Expand All @@ -219,59 +216,62 @@ func (c *Client) DeleteEventsByDeviceNameAndSourceName(deviceName, sourceName st
return c.deleteEventsByConditions([]string{deviceNameCol, sourceNameCol}, pgx.NamedArgs{deviceNameCol: deviceName, sourceNameCol: sourceName})
}

// deleteEventsByConditions deletes specific device's events and corresponding readings
// deleteEventsByConditions deletes specific events, readings, and deviceInfos in one
// transaction so callers only observe success after the delete has actually succeeded.
func (c *Client) deleteEventsByConditions(cols []string, values pgx.NamedArgs) errors.EdgeX {
ctx := context.Background()

sqlStatement := sqlDeleteEventsByColumn(cols...)

go func() {
// deviceInfos are used to remove data by id and remove the id value from the cache after finishing the transaction
deviceInfos, err := c.deviceInfosByConds(cols, values)
if err != nil {
return
}
// delete events and readings in a transaction
pgxErr := pgx.BeginFunc(ctx, c.ConnPool, func(tx pgx.Tx) error {
// select the event-ids of the specified device name from event table as the sub-query of deleting readings
subSqlStatement := sqlQueryEventIdFieldsByCol(cols...)
if err = deleteReadingsBySubQuery(ctx, tx, subSqlStatement, values); err != nil {
if errors.Kind(err) == errors.KindEntityDoesNotExist {
c.loggingClient.Debugf("no readings found for deletion: %s", err.Error())
} else {
c.loggingClient.Errorf("failed delete readings with conditions '%v' '%v': %v", cols, values, err)
return err
}
}
// deviceInfos are removed from the cache only after the transaction commits.
deviceInfos, err := c.deviceInfosByConds(cols, values)
if err != nil {
c.loggingClient.Errorf("failed querying deviceInfos with conditions '%v' '%v': %v", cols, values, err)
return errors.NewCommonEdgeXWrapper(err)
}

err = deleteEvents(ctx, tx, sqlStatement, values)
if err != nil {
if errors.Kind(err) == errors.KindEntityDoesNotExist {
c.loggingClient.Debugf("no events found for deletion: %s", err.Error())
} else {
c.loggingClient.Errorf("failed delete event with conditions '%v' '%v': %v", cols, values, err)
return err
}
}
pgxErr := pgx.BeginFunc(ctx, c.ConnPool, func(tx pgx.Tx) error {
return c.deleteEventsByConditionsInTx(ctx, tx, cols, values, sqlStatement, deviceInfos)
})
if pgxErr != nil {
c.loggingClient.Errorf("failed delete events with conditions '%v' '%v': %v", cols, values, pgxErr)
return errors.NewCommonEdgeXWrapper(pgxErr)
}

for _, deviceInfo := range deviceInfos {
err = deleteDeviceInfoById(ctx, tx, deviceInfo.Id)
if err != nil {
return errors.NewCommonEdgeXWrapper(err)
}
}
return nil
})
if pgxErr != nil {
c.loggingClient.Errorf("failed delete events with conditions '%v' '%v': %v", cols, values, pgxErr)
return
deviceInfoCache := container.DeviceInfoCacheFrom(c.dic.Get)
for _, deviceInfo := range deviceInfos {
deviceInfoCache.Remove(deviceInfo)
}

return nil
}

func (c *Client) deleteEventsByConditionsInTx(ctx context.Context, tx pgx.Tx, cols []string, values pgx.NamedArgs, sqlStatement string, deviceInfos []dbModels.DeviceInfo) errors.EdgeX {
// select the event ids from the event table as the sub-query of deleting readings
subSqlStatement := sqlQueryEventIdFieldsByCol(cols...)
if err := deleteReadingsBySubQuery(ctx, tx, subSqlStatement, values); err != nil {
if errors.Kind(err) == errors.KindEntityDoesNotExist {
c.loggingClient.Debugf("no readings found for deletion: %s", err.Error())
} else {
c.loggingClient.Errorf("failed delete readings with conditions '%v' '%v': %v", cols, values, err)
return err
}
}

if err := deleteEvents(ctx, tx, sqlStatement, values); err != nil {
if errors.Kind(err) == errors.KindEntityDoesNotExist {
c.loggingClient.Debugf("no events found for deletion: %s", err.Error())
} else {
c.loggingClient.Errorf("failed delete event with conditions '%v' '%v': %v", cols, values, err)
return err
}
}

deviceInfoCache := container.DeviceInfoCacheFrom(c.dic.Get)
for _, deviceInfo := range deviceInfos {
deviceInfoCache.Remove(deviceInfo)
for _, deviceInfo := range deviceInfos {
if err := deleteDeviceInfoById(ctx, tx, deviceInfo.Id); err != nil {
return errors.NewCommonEdgeXWrapper(err)
}
}()
}

return nil
}
Expand Down
70 changes: 70 additions & 0 deletions internal/pkg/infrastructure/postgres/event_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,9 @@ import (
"strings"
"testing"

dbModels "github.com/edgexfoundry/edgex-go/internal/pkg/infrastructure/models"
"github.com/edgexfoundry/edgex-go/internal/pkg/infrastructure/postgres/mocks"
"github.com/edgexfoundry/go-mod-core-contracts/v4/clients/logger"
"github.com/edgexfoundry/go-mod-core-contracts/v4/errors"

"github.com/jackc/pgx/v5"
Expand Down Expand Up @@ -147,3 +149,71 @@ func TestDeleteReadingsBySubQuery(t *testing.T) {
})
}
}

func TestDeleteEventsByConditionsInTx(t *testing.T) {
ctx := context.Background()
cols := []string{deviceNameCol}
values := pgx.NamedArgs{deviceNameCol: "test-device"}
sqlStatement := sqlDeleteEventsByColumn(cols...)
deviceInfos := []dbModels.DeviceInfo{{Id: 7}}

tests := []struct {
name string
readingRows int64
readingErr error
eventRows int64
eventErr error
deviceInfoErr error
expectError bool
expectedErrKind errors.ErrKind
}{
{
name: "success",
readingRows: 1,
eventRows: 1,
expectError: false,
expectedErrKind: "",
},
{
name: "reading delete failure returns error",
readingRows: 1,
readingErr: fmt.Errorf("delete readings failed"),
expectError: true,
expectedErrKind: errors.KindDatabaseError,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
tx := new(mocks.Tx)
client := &Client{loggingClient: logger.NewMockClient()}

subQuerySQL := sqlQueryEventIdFieldsByCol(cols...)
readingTag := pgconn.NewCommandTag(fmt.Sprintf("DELETE %d", tt.readingRows))
tx.On("Exec", ctx, mock.MatchedBy(func(sql string) bool {
return strings.Contains(sql, subQuerySQL)
}), values).Return(readingTag, tt.readingErr)

if tt.readingErr == nil {
eventTag := pgconn.NewCommandTag(fmt.Sprintf("DELETE %d", tt.eventRows))
tx.On("Exec", ctx, sqlStatement, values).Return(eventTag, tt.eventErr)

if tt.eventErr == nil && tt.eventRows > 0 {
deviceInfoSQL := fmt.Sprintf("DELETE FROM %s WHERE %s = @%s", deviceInfoTableName, idCol, idCol)
tx.On("Exec", ctx, deviceInfoSQL, pgx.NamedArgs{idCol: deviceInfos[0].Id}).Return(pgconn.NewCommandTag("DELETE 1"), tt.deviceInfoErr)
}
}

err := client.deleteEventsByConditionsInTx(ctx, tx, cols, values, sqlStatement, deviceInfos)

if tt.expectError {
assert.Error(t, err)
assert.Equal(t, tt.expectedErrKind, errors.Kind(err))
} else {
assert.NoError(t, err)
}

tx.AssertExpectations(t)
})
}
}
Loading