Skip to content

Commit e94fec4

Browse files
committed
Check server and account limits on stream restore
Signed-off-by: Neil Twigg <neil@nats.io>
1 parent 1d60dce commit e94fec4

2 files changed

Lines changed: 29 additions & 2 deletions

File tree

server/jetstream.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2315,7 +2315,7 @@ func (js *jetStream) checkBytesLimits(selectedLimits *JetStreamAccountLimits, ad
23152315
return NewJSMemoryResourcesExceededError()
23162316
}
23172317
// Check if this server can handle request.
2318-
if checkServer && js.memReserved+addBytes > js.config.MaxMemory {
2318+
if checkServer && js.memReserved+totalBytes > js.config.MaxMemory {
23192319
return NewJSMemoryResourcesExceededError()
23202320
}
23212321
case FileStorage:
@@ -2324,7 +2324,7 @@ func (js *jetStream) checkBytesLimits(selectedLimits *JetStreamAccountLimits, ad
23242324
return NewJSStorageResourcesExceededError()
23252325
}
23262326
// Check if this server can handle request.
2327-
if checkServer && js.storeReserved+addBytes > js.config.MaxStore {
2327+
if checkServer && js.storeReserved+totalBytes > js.config.MaxStore {
23282328
return NewJSStorageResourcesExceededError()
23292329
}
23302330
}

server/stream.go

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6303,6 +6303,10 @@ func (a *Account) RestoreStream(ncfg *StreamConfig, r io.Reader) (*stream, error
63036303
if err != nil {
63046304
return nil, err
63056305
}
6306+
js := jsa.js
6307+
if js == nil {
6308+
return nil, NewJSNotEnabledForAccountError()
6309+
}
63066310

63076311
cfg, apiErr := s.checkStreamCfg(ncfg, a, false)
63086312
if apiErr != nil {
@@ -6337,6 +6341,22 @@ func (a *Account) RestoreStream(ncfg *StreamConfig, r io.Reader) (*stream, error
63376341
}
63386342
sdirCheck := filepath.Clean(sdir) + string(os.PathSeparator)
63396343

6344+
_, isClustered := jsa.jetStreamAndClustered()
6345+
jsa.usageMu.RLock()
6346+
selected, tier, hasTier := jsa.selectLimits(cfg.Replicas)
6347+
jsa.usageMu.RUnlock()
6348+
reserved := int64(0)
6349+
if hasTier {
6350+
if isClustered {
6351+
js.mu.RLock()
6352+
_, reserved = tieredStreamAndReservationCount(js.cluster.streams[a.Name], tier, &cfg)
6353+
js.mu.RUnlock()
6354+
} else {
6355+
reserved = jsa.tieredReservation(tier, &cfg)
6356+
}
6357+
}
6358+
6359+
var bc int64
63406360
tr := tar.NewReader(s2.NewReader(r))
63416361
for {
63426362
hdr, err := tr.Next()
@@ -6349,6 +6369,13 @@ func (a *Account) RestoreStream(ncfg *StreamConfig, r io.Reader) (*stream, error
63496369
if hdr.Typeflag != tar.TypeReg {
63506370
return nil, logAndReturnError()
63516371
}
6372+
bc += hdr.Size
6373+
js.mu.RLock()
6374+
err = js.checkAllLimits(&selected, &cfg, reserved, bc)
6375+
js.mu.RUnlock()
6376+
if err != nil {
6377+
return nil, err
6378+
}
63526379
fpath := filepath.Join(sdir, filepath.Clean(hdr.Name))
63536380
if !strings.HasPrefix(fpath, sdirCheck) {
63546381
return nil, logAndReturnError()

0 commit comments

Comments
 (0)