Skip to content

Commit 9fd93dc

Browse files
committed
feat: Add retention size limit for rotated files
This adds a new flag to the observer command to limit the size of the rotated files. When the rotated files exceed this size, the oldest files are deleted. Signed-off-by: Jakub Sztandera <oss@kubuxu.com>
1 parent e7f84e7 commit 9fd93dc

3 files changed

Lines changed: 54 additions & 7 deletions

File tree

cmd/f3/observer.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,11 @@ var observerCmd = cli.Command{
6666
Usage: "The maximum length of time to keep the rotated files.",
6767
Value: 2 * 7 * 24 * time.Hour,
6868
},
69+
&cli.Int64Flag{
70+
Name: "retentionSize",
71+
Usage: "The maximum size of the rotated files in megabytes. If not set, no limit is applied.",
72+
Value: 0,
73+
},
6974
&cli.StringFlag{
7075
Name: "dataSourceName",
7176
Usage: "The observer database DSN",
@@ -144,7 +149,9 @@ var observerCmd = cli.Command{
144149
observer.WithMaxBatchSize(cctx.Int("maxBatchSize")),
145150
observer.WithMaxBatchDelay(cctx.Duration("maxBatchDelay")),
146151
observer.WithChainExchangeMaxMessageAge(cctx.Duration("chainExchangeMaxMessageAge")),
152+
observer.WithMaxRetentionSize(cctx.Int64("retentionSize") * 1024 * 1024),
147153
}
154+
148155
var identity crypto.PrivKey
149156
if cctx.IsSet("identity") {
150157
marshaledKey, err := os.ReadFile(cctx.String("identity"))

observer/observer.go

Lines changed: 30 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,10 +8,12 @@ import (
88
"encoding/json"
99
"errors"
1010
"fmt"
11+
"io/fs"
1112
"net/http"
1213
"os"
1314
"path"
1415
"path/filepath"
16+
"sort"
1517
"sync"
1618
"time"
1719

@@ -527,7 +529,8 @@ func (o *Observer) rotateMessages(ctx context.Context) error {
527529
if err != nil {
528530
return err
529531
}
530-
var foundAtLeastOneParquet bool
532+
retainedSize := int64(0)
533+
var retained []fs.FileInfo
531534
for _, entry := range dir {
532535
if !entry.IsDir() && filepath.Ext(entry.Name()) == ".parquet" {
533536
info, err := entry.Info()
@@ -541,12 +544,36 @@ func (o *Observer) rotateMessages(ctx context.Context) error {
541544
logger.Infow("Removed old file", "olderThan", o.retention, "file", entry.Name())
542545
}
543546
} else {
544-
foundAtLeastOneParquet = true
547+
retainedSize += info.Size()
548+
retained = append(retained, info)
545549
}
546550
}
547551
}
548552

549-
return o.createOrReplaceMessagesView(ctx, foundAtLeastOneParquet)
553+
logger.Infow("Retention size", "retainedSize", retainedSize, "maxRetentionSize", o.maxRetentionSize)
554+
555+
if retainedSize > o.maxRetentionSize {
556+
logger.Infow("Retention size exceeded, deleting oldest files", "retainedSize", retainedSize, "maxRetentionSize", o.maxRetentionSize)
557+
// sort retained by modification time, oldest last
558+
sort.Slice(retained, func(i, j int) bool {
559+
return retained[i].ModTime().After(retained[j].ModTime())
560+
})
561+
// iterate in reverse order to delete oldest first
562+
for i := len(retained) - 1; i >= 0; i-- {
563+
if retainedSize < o.maxRetentionSize {
564+
break
565+
}
566+
if err := os.Remove(filepath.Join(o.rotatePath, retained[i].Name())); err != nil {
567+
logger.Errorw("Failed to remove retention policy for file", "file", retained[i].Name(), "err", err)
568+
} else {
569+
logger.Infow("Removed old file", "olderThan", o.retention, "file", retained[i].Name())
570+
retainedSize -= retained[i].Size()
571+
retained = retained[:i]
572+
}
573+
}
574+
}
575+
576+
return o.createOrReplaceMessagesView(ctx, len(retained) > 0)
550577
}
551578

552579
func (o *Observer) listenAndServeQueries() error {

observer/options.go

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ import (
1313
"github.com/libp2p/go-libp2p/core/host"
1414
"github.com/libp2p/go-libp2p/core/peer"
1515
"github.com/multiformats/go-multiaddr"
16-
"github.com/multiformats/go-multiaddr-dns"
16+
madns "github.com/multiformats/go-multiaddr-dns"
1717
)
1818

1919
type Option func(*options) error
@@ -35,9 +35,10 @@ type options struct {
3535
queryServerListenAddress string
3636
queryServerReadTimeout time.Duration
3737

38-
rotatePath string
39-
rotateInterval time.Duration
40-
retention time.Duration
38+
rotatePath string
39+
rotateInterval time.Duration
40+
retention time.Duration
41+
maxRetentionSize int64
4142

4243
pubSub *pubsub.PubSub
4344
pubSubValidatorDisabled bool
@@ -81,6 +82,7 @@ func newOptions(opts ...Option) (*options, error) {
8182
finalityCertsMaxPollInterval: 2 * time.Minute,
8283
chainExchangeBufferSize: 1000,
8384
chainExchangeMaxMessageAge: 3 * time.Minute,
85+
maxRetentionSize: 0,
8486
}
8587
for _, apply := range opts {
8688
if err := apply(&opt); err != nil {
@@ -270,6 +272,17 @@ func WithRetention(retention time.Duration) Option {
270272
}
271273
}
272274

275+
// WithMaxRetentionSize sets the maximum size of the retention directory.
276+
// This is weakly enforced, and the directory may grow larger than this
277+
// size. If the directory grows larger than this size, the oldest files
278+
// will be deleted until the directory size is below this size.
279+
func WithMaxRetentionSize(size int64) Option {
280+
return func(o *options) error {
281+
o.maxRetentionSize = size
282+
return nil
283+
}
284+
}
285+
273286
func WithDataSourceName(dataSourceName string) Option {
274287
return func(o *options) error {
275288
o.dataSourceName = dataSourceName

0 commit comments

Comments
 (0)