Skip to content

Commit e5da1b2

Browse files
authored
alerts: create metas only if alert was not discarded (#4633)
1 parent 83f688a commit e5da1b2

4 files changed

Lines changed: 165 additions & 8 deletions

File tree

pkg/database/alerts.go

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -653,11 +653,6 @@ func (c *Client) createAlertBatch(ctx context.Context, machineID string, owner *
653653
return nil, rollbackOnError(tx, err, fmt.Sprintf("building events for alert %s", alertItem.UUID))
654654
}
655655

656-
metas, err := buildMetaCreates(ctx, c.Log, txEnt, alertItem)
657-
if err != nil {
658-
c.Log.Warningf("error creating alert meta: %s", err)
659-
}
660-
661656
decisions, discardCount, err := c.buildDecisions(ctx, c.Log, txEnt, alertItem, stopAtTime)
662657
if err != nil {
663658
return nil, rollbackOnError(tx, err, fmt.Sprintf("building decisions for alert %s", alertItem.UUID))
@@ -669,6 +664,13 @@ func (c *Client) createAlertBatch(ctx context.Context, machineID string, owner *
669664
continue
670665
}
671666

667+
// after the discard check: metas are linked to the alert only when it is built,
668+
// so inserting them earlier leaves orphan rows behind when the alert is dropped
669+
metas, err := buildMetaCreates(ctx, c.Log, txEnt, alertItem)
670+
if err != nil {
671+
c.Log.Warningf("error creating alert meta: %s", err)
672+
}
673+
672674
builder := txEnt.Alert.
673675
Create().
674676
SetScenario(*alertItem.Scenario).

pkg/database/alerts_pg_race_test.go

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -147,11 +147,17 @@ func TestAlertCreateVsFlushOrphansRace(t *testing.T) {
147147
require.Zero(t, edgeErrCount.Load(), "got %d edge-constraint errors (events deleted out from under alert insert)", edgeErrCount.Load())
148148
require.Zero(t, otherErrCount.Load(), "got %d unexpected errors", otherErrCount.Load())
149149
require.Positive(t, successCount.Load(), "no alerts were inserted at all — test setup is wrong")
150+
151+
// every meta of a committed alert must still be there: the orphan flush must not
152+
// see rows that are still in flight
153+
metaCount, err := dbClient.Ent.Meta.Query().Count(ctx)
154+
require.NoError(t, err)
155+
require.Equal(t, int(successCount.Load())*2, metaCount, "orphan flush deleted metas of in-flight alerts")
150156
}
151157

152-
// makeRaceAlerts builds a batch of distinct alerts with several events and a
153-
// decision each, so the create path actually exercises Event.CreateBulk plus
154-
// the alert↔events O2M edge update plus decision attachment.
158+
// makeRaceAlerts builds a batch of distinct alerts with several events, metas and
159+
// a decision each, so the create path actually exercises Event.CreateBulk and
160+
// Meta.CreateBulk plus the alert↔events O2M edge update plus decision attachment.
155161
func makeRaceAlerts(workerID, batchID, eventsPerAlert int) []*models.Alert {
156162
now := time.Now().UTC().Format(time.RFC3339)
157163

@@ -201,6 +207,10 @@ func makeRaceAlerts(workerID, batchID, eventsPerAlert int) []*models.Alert {
201207
IP: value,
202208
},
203209
Events: events,
210+
Meta: models.Meta{
211+
{Key: "datasource_path", Value: "/var/log/test.log"},
212+
{Key: "datasource_type", Value: "file"},
213+
},
204214
Decisions: []*models.Decision{
205215
{
206216
Duration: &duration,

pkg/database/flush.go

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ import (
1717
"github.com/crowdsecurity/crowdsec/pkg/database/ent/decision"
1818
"github.com/crowdsecurity/crowdsec/pkg/database/ent/event"
1919
"github.com/crowdsecurity/crowdsec/pkg/database/ent/machine"
20+
"github.com/crowdsecurity/crowdsec/pkg/database/ent/meta"
2021
"github.com/crowdsecurity/crowdsec/pkg/database/ent/metric"
2122
"github.com/crowdsecurity/crowdsec/pkg/database/ent/predicate"
2223
"github.com/crowdsecurity/crowdsec/pkg/logging"
@@ -27,6 +28,11 @@ const (
2728
// how long to keep metrics in the local database
2829
defaultMetricsMaxAge = 7 * 24 * time.Hour
2930
flushInterval = 1 * time.Minute
31+
// orphan metas are deleted in bounded batches: an install upgrading with a
32+
// backlog can carry hundreds of millions of them, and a single unbounded
33+
// DELETE would hold the table for minutes.
34+
orphanMetaBatchSize = 5_000
35+
orphanMetaMaxPerRun = 500_000
3036
)
3137

3238
func (c *Client) StartFlushScheduler(ctx context.Context, config *csconfig.FlushDBCfg) (gocron.Scheduler, error) {
@@ -183,6 +189,54 @@ func (c *Client) FlushOrphans(ctx context.Context) {
183189
if eventsCount > 0 {
184190
c.Log.Infof("%d deleted orphan decisions", eventsCount)
185191
}
192+
193+
metasCount, err := c.flushOrphanMetas(ctx)
194+
if err != nil {
195+
c.Log.Warningf("error while deleting orphan metas: %s", err)
196+
return
197+
}
198+
199+
if metasCount > 0 {
200+
c.Log.Infof("%d deleted orphan metas", metasCount)
201+
}
202+
}
203+
204+
// flushOrphanMetas deletes meta rows that were never linked to an alert, up to
205+
// orphanMetaMaxPerRun per call. Versions before 1.8.0 inserted metas outside the
206+
// alert transaction, so an upgraded database can hold a large backlog: it is
207+
// drained over several runs instead of in one statement.
208+
func (c *Client) flushOrphanMetas(ctx context.Context) (int, error) {
209+
deleted := 0
210+
211+
for deleted < orphanMetaMaxPerRun {
212+
ids, err := c.Ent.Meta.Query().Where(meta.Not(meta.HasOwner())).Limit(orphanMetaBatchSize).IDs(ctx)
213+
if err != nil {
214+
return deleted, err
215+
}
216+
217+
if len(ids) == 0 {
218+
break
219+
}
220+
221+
count, err := c.Ent.Meta.Delete().Where(meta.IDIn(ids...)).Exec(ctx)
222+
if err != nil {
223+
return deleted, err
224+
}
225+
226+
// nothing deleted for a full batch means someone else got there first: stop
227+
// rather than spin on the same rows
228+
if count == 0 {
229+
break
230+
}
231+
232+
deleted += count
233+
234+
if len(ids) < orphanMetaBatchSize {
235+
break
236+
}
237+
}
238+
239+
return deleted, nil
186240
}
187241

188242
func (c *Client) flushBouncers(ctx context.Context, authType string, duration *time.Duration) {

pkg/database/flush_test.go

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,15 @@ package database
33
import (
44
"context"
55
"fmt"
6+
"slices"
67
"testing"
78
"time"
89

910
"github.com/go-openapi/strfmt"
1011
"github.com/stretchr/testify/require"
1112

13+
"github.com/crowdsecurity/crowdsec/pkg/database/ent"
14+
"github.com/crowdsecurity/crowdsec/pkg/database/ent/meta"
1215
"github.com/crowdsecurity/crowdsec/pkg/models"
1316
"github.com/crowdsecurity/crowdsec/pkg/types"
1417
)
@@ -156,3 +159,91 @@ func TestFlushAlerts_MaxItemsKeepsActiveDecisions(t *testing.T) {
156159
require.NoError(t, err)
157160
require.Equal(t, 1, decCount, "active decision must not be cascade-deleted")
158161
}
162+
163+
func countMetas(t *testing.T, ctx context.Context, c *Client) (total int, orphans int) {
164+
t.Helper()
165+
166+
total, err := c.Ent.Meta.Query().Count(ctx)
167+
require.NoError(t, err)
168+
169+
orphans, err = c.Ent.Meta.Query().Where(meta.Not(meta.HasOwner())).Count(ctx)
170+
require.NoError(t, err)
171+
172+
return total, orphans
173+
}
174+
175+
// A dropped alert must not leave its metas behind: they are inserted before the
176+
// alert row, and only the alert insert links them.
177+
func TestCreateAlert_DroppedAlertLeavesNoMetas(t *testing.T) {
178+
ctx := t.Context()
179+
c := getDBClient(t, ctx)
180+
181+
machineID := "test-dropped-alert-metas"
182+
registerFlushTestMachine(t, ctx, c, machineID)
183+
184+
// the decision value is not a valid address, so the decision is discarded and
185+
// the alert with it
186+
dropped := makeFlushAlert("not-an-ip", true)
187+
dropped.Meta = models.Meta{{Key: "k1", Value: "v1"}, {Key: "k2", Value: "v2"}}
188+
189+
kept := makeFlushAlert("1.2.3.4", true)
190+
kept.Meta = models.Meta{{Key: "k1", Value: "v1"}}
191+
192+
_, err := c.CreateAlert(ctx, machineID, []*models.Alert{dropped, kept})
193+
require.NoError(t, err)
194+
195+
alerts, err := c.Ent.Alert.Query().Count(ctx)
196+
require.NoError(t, err)
197+
require.Equal(t, 1, alerts)
198+
199+
total, orphans := countMetas(t, ctx, c)
200+
require.Equal(t, 1, total)
201+
require.Zero(t, orphans)
202+
}
203+
204+
// Metas orphaned by an older version are reaped, linked ones are not.
205+
func TestFlushOrphans_DeletesOrphanMetas(t *testing.T) {
206+
tests := []struct {
207+
name string
208+
orphans int
209+
}{
210+
{name: "single batch", orphans: 3},
211+
{name: "several batches", orphans: orphanMetaBatchSize + 1},
212+
}
213+
214+
for _, tc := range tests {
215+
t.Run(tc.name, func(t *testing.T) {
216+
ctx := t.Context()
217+
c := getDBClient(t, ctx)
218+
219+
machineID := "test-orphan-metas"
220+
registerFlushTestMachine(t, ctx, c, machineID)
221+
222+
alert := makeFlushAlert("1.2.3.4", true)
223+
alert.Meta = models.Meta{{Key: "k1", Value: "v1"}}
224+
225+
_, err := c.CreateAlert(ctx, machineID, []*models.Alert{alert})
226+
require.NoError(t, err)
227+
228+
for chunk := range slices.Chunk(make([]struct{}, tc.orphans), 500) {
229+
builders := make([]*ent.MetaCreate, len(chunk))
230+
for i := range chunk {
231+
builders[i] = c.Ent.Meta.Create().SetKey("orphan").SetValue("v")
232+
}
233+
234+
_, err := c.Ent.Meta.CreateBulk(builders...).Save(ctx)
235+
require.NoError(t, err)
236+
}
237+
238+
total, orphans := countMetas(t, ctx, c)
239+
require.Equal(t, tc.orphans+1, total)
240+
require.Equal(t, tc.orphans, orphans)
241+
242+
c.FlushOrphans(ctx)
243+
244+
total, orphans = countMetas(t, ctx, c)
245+
require.Equal(t, 1, total)
246+
require.Zero(t, orphans)
247+
})
248+
}
249+
}

0 commit comments

Comments
 (0)