Skip to content
Merged
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
10 changes: 10 additions & 0 deletions server/errors.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": ""
}
]
75 changes: 75 additions & 0 deletions server/jetstream_batching_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"encoding/json"
"errors"
"fmt"
"math"
"math/big"
"strconv"
"strings"
Expand Down Expand Up @@ -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)
}
14 changes: 14 additions & 0 deletions server/jetstream_errors_generated.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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}"},
Expand Down Expand Up @@ -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)
Expand Down
32 changes: 31 additions & 1 deletion server/stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down
Loading