From e6f30f22cbd3df24ff03bf076ee110aa8ddd9a06 Mon Sep 17 00:00:00 2001 From: Maurice van Veen Date: Mon, 29 Sep 2025 10:44:59 +0200 Subject: [PATCH] (2.12) Atomic batch: reject unsupported commit Signed-off-by: Maurice van Veen --- server/errors.json | 10 ++++ server/jetstream_batching_test.go | 75 ++++++++++++++++++++++++++++ server/jetstream_errors_generated.go | 14 ++++++ server/stream.go | 32 +++++++++++- 4 files changed, 130 insertions(+), 1 deletion(-) diff --git a/server/errors.json b/server/errors.json index 9293361dba5..7f04b23469a 100644 --- a/server/errors.json +++ b/server/errors.json @@ -1978,5 +1978,15 @@ "help": "", "url": "", "deprecates": "" + }, + { + "constant": "JSAtomicPublishInvalidBatchCommitErr", + "code": 400, + "error_code": 10200, + "description": "atomic publish batch commit is invalid", + "comment": "", + "help": "", + "url": "", + "deprecates": "" } ] diff --git a/server/jetstream_batching_test.go b/server/jetstream_batching_test.go index 82e16ba4ff8..c00d3c5fdf6 100644 --- a/server/jetstream_batching_test.go +++ b/server/jetstream_batching_test.go @@ -20,6 +20,7 @@ import ( "encoding/json" "errors" "fmt" + "math" "math/big" "strconv" "strings" @@ -2615,3 +2616,77 @@ func TestJetStreamAtomicBatchPublishExpectedLastSubjectSequence(t *testing.T) { require_Equal(t, resp.PubAck.BatchId, "uuid") require_Equal(t, resp.PubAck.BatchSize, 2) } + +func TestJetStreamAtomicBatchPublishCommitUnsupported(t *testing.T) { + s := RunBasicJetStreamServer(t) + defer s.Shutdown() + + nc := clientConnectToServer(t, s) + defer nc.Close() + + cfg := &StreamConfig{ + Name: "TEST", + Subjects: []string{"foo"}, + Storage: MemoryStorage, + Replicas: 1, + AllowAtomicPublish: true, + } + _, err := jsStreamCreate(t, nc, cfg) + require_NoError(t, err) + + mset, err := s.globalAccount().lookupStream("TEST") + require_NoError(t, err) + + var resp JSPubAckResponse + for _, unsupportedCommit := range []string{"", "unsupported", "0"} { + m := nats.NewMsg("foo") + m.Header.Set("Nats-Batch-Id", "uuid") + m.Header.Set("Nats-Batch-Sequence", "1") + _, err = nc.RequestMsg(m, time.Second) + require_NoError(t, err) + + m.Header.Set("Nats-Batch-Sequence", "2") + m.Header.Set("Nats-Batch-Commit", unsupportedCommit) + msg, err := nc.RequestMsg(m, time.Second) + require_NoError(t, err) + resp = JSPubAckResponse{} + require_NoError(t, json.Unmarshal(msg.Data, &resp)) + require_True(t, resp.Error != nil) + require_Error(t, resp.Error, NewJSAtomicPublishInvalidBatchCommitError()) + + // Confirm no batches are left. + mset.mu.RLock() + batches := mset.batches + mset.mu.RUnlock() + batches.mu.Lock() + groups := len(batches.group) + batches.mu.Unlock() + require_Len(t, groups, 0) + } + + // The required API level should allow the batch to be rejected. + m := nats.NewMsg("foo") + m.Header.Set("Nats-Batch-Id", "uuid") + m.Header.Set("Nats-Batch-Sequence", "1") + m.Header.Set("Nats-Batch-Commit", "1") + m.Header.Set("Nats-Required-Api-Level", strconv.Itoa(math.MaxInt)) + msg, err := nc.RequestMsg(m, time.Second) + require_NoError(t, err) + resp = JSPubAckResponse{} + require_NoError(t, json.Unmarshal(msg.Data, &resp)) + require_True(t, resp.Error != nil) + require_Error(t, resp.Error, NewJSRequiredApiLevelError()) + + // If required API level check passes, the header should be stripped. + m.Header.Set("Nats-Required-Api-Level", "0") + msg, err = nc.RequestMsg(m, time.Second) + require_NoError(t, err) + resp = JSPubAckResponse{} + require_NoError(t, json.Unmarshal(msg.Data, &resp)) + require_True(t, resp.Error == nil) + require_Equal(t, resp.PubAck.Sequence, 1) + + sm, err := mset.getMsg(1) + require_NoError(t, err) + require_Len(t, len(sliceHeader(JSRequiredApiLevel, sm.Header)), 0) +} diff --git a/server/jetstream_errors_generated.go b/server/jetstream_errors_generated.go index a1ca08cfd09..dbc7f0499d7 100644 --- a/server/jetstream_errors_generated.go +++ b/server/jetstream_errors_generated.go @@ -14,6 +14,9 @@ const ( // JSAtomicPublishIncompleteBatchErr atomic publish batch is incomplete JSAtomicPublishIncompleteBatchErr ErrorIdentifier = 10176 + // JSAtomicPublishInvalidBatchCommitErr atomic publish batch commit is invalid + JSAtomicPublishInvalidBatchCommitErr ErrorIdentifier = 10200 + // JSAtomicPublishInvalidBatchIDErr atomic publish batch ID is invalid JSAtomicPublishInvalidBatchIDErr ErrorIdentifier = 10179 @@ -605,6 +608,7 @@ var ( JSAccountResourcesExceededErr: {Code: 400, ErrCode: 10002, Description: "resource limits exceeded for account"}, JSAtomicPublishDisabledErr: {Code: 400, ErrCode: 10174, Description: "atomic publish is disabled"}, JSAtomicPublishIncompleteBatchErr: {Code: 400, ErrCode: 10176, Description: "atomic publish batch is incomplete"}, + JSAtomicPublishInvalidBatchCommitErr: {Code: 400, ErrCode: 10200, Description: "atomic publish batch commit is invalid"}, JSAtomicPublishInvalidBatchIDErr: {Code: 400, ErrCode: 10179, Description: "atomic publish batch ID is invalid"}, JSAtomicPublishMissingSeqErr: {Code: 400, ErrCode: 10175, Description: "atomic publish sequence is missing"}, JSAtomicPublishTooLargeBatchErrF: {Code: 400, ErrCode: 10199, Description: "atomic publish batch is too large: {size}"}, @@ -855,6 +859,16 @@ func NewJSAtomicPublishIncompleteBatchError(opts ...ErrorOption) *ApiError { return ApiErrors[JSAtomicPublishIncompleteBatchErr] } +// NewJSAtomicPublishInvalidBatchCommitError creates a new JSAtomicPublishInvalidBatchCommitErr error: "atomic publish batch commit is invalid" +func NewJSAtomicPublishInvalidBatchCommitError(opts ...ErrorOption) *ApiError { + eopts := parseOpts(opts) + if ae, ok := eopts.err.(*ApiError); ok { + return ae + } + + return ApiErrors[JSAtomicPublishInvalidBatchCommitErr] +} + // NewJSAtomicPublishInvalidBatchIDError creates a new JSAtomicPublishInvalidBatchIDErr error: "atomic publish batch ID is invalid" func NewJSAtomicPublishInvalidBatchIDError(opts ...ErrorOption) *ApiError { eopts := parseOpts(opts) diff --git a/server/stream.go b/server/stream.go index 1a080e6cd44..1bf64daba42 100644 --- a/server/stream.go +++ b/server/stream.go @@ -6264,7 +6264,6 @@ func (mset *stream) processJetStreamBatchMsg(batchId, subject, reply string, hdr } return err } - commit := len(sliceHeader(JSBatchCommit, hdr)) != 0 mset.mu.Lock() if mset.batches == nil { @@ -6339,6 +6338,37 @@ func (mset *stream) processJetStreamBatchMsg(batchId, subject, reply string, hdr batches.group[batchId] = b } + var commit bool + if c := sliceHeader(JSBatchCommit, hdr); c != nil { + // Reject the batch if the commit is not recognized. + if !bytes.Equal(c, []byte("1")) { + b.cleanupLocked(batchId, batches) + batches.mu.Unlock() + err := NewJSAtomicPublishInvalidBatchCommitError() + if canRespond { + b, _ := json.Marshal(&JSPubAckResponse{PubAck: &PubAck{Stream: name}, Error: err}) + outq.send(newJSPubMsg(reply, _EMPTY_, _EMPTY_, nil, b, nil, 0)) + } + return err + } + commit = true + } + + // The required API level can have the batch be rejected. But the header is always removed. + if len(sliceHeader(JSRequiredApiLevel, hdr)) != 0 { + if errorOnRequiredApiLevel(hdr) { + b.cleanupLocked(batchId, batches) + batches.mu.Unlock() + err := NewJSRequiredApiLevelError() + if canRespond { + b, _ := json.Marshal(&JSPubAckResponse{PubAck: &PubAck{Stream: name}, Error: err}) + outq.send(newJSPubMsg(reply, _EMPTY_, _EMPTY_, nil, b, nil, 0)) + } + return err + } + hdr = removeHeaderIfPresent(hdr, JSRequiredApiLevel) + } + // Detect gaps. b.lseq++ if b.lseq != batchSeq {