nats-io/nats-server · error

invalid retained message flags

Error message

invalid retained message flags

What it means

errMQTTInvalidRetainFlags is returned when a retained MQTT message has invalid packet flags. Retained messages are reconstructed from stored flags; the code (server/mqtt.go:3110 and 3130) rejects a stored/encoded error condition by substituting this error, and rejects any flags value at or above the full mask (1111), since valid retained-message flags cannot use all four bits.

Source

Thrown at server/mqtt.go:245

	errMQTTMalformedVarInt            = errors.New("malformed variable int")
	errMQTTSecondConnectPacket        = errors.New("received a second CONNECT packet")
	errMQTTServerNameMustBeSet        = errors.New("mqtt requires server name to be explicitly set")
	errMQTTUserMixWithUsersNKeys      = errors.New("mqtt authentication username not compatible with presence of users/nkeys")
	errMQTTTokenMixWIthUsersNKeys     = errors.New("mqtt authentication token not compatible with presence of users/nkeys")
	errMQTTAckWaitMustBePositive      = errors.New("ack wait must be a positive value")
	errMQTTJSAPITimeoutMustBePositive = errors.New("JS API timeout must be a positive value")
	errMQTTStandaloneNeedsJetStream   = errors.New("mqtt requires JetStream to be enabled if running in standalone mode")
	errMQTTConnFlagReserved           = errors.New("connect flags reserved bit not set to 0")
	errMQTTWillAndRetainFlag          = errors.New("if Will flag is set to 0, Will Retain flag must be 0 too")
	errMQTTPasswordFlagAndNoUser      = errors.New("password flag set but username flag is not")
	errMQTTCIDEmptyNeedsCleanFlag     = errors.New("when client ID is empty, clean session flag must be set to 1")
	errMQTTEmptyWillTopic             = errors.New("empty Will topic not allowed")
	errMQTTEmptyUsername              = errors.New("empty user name not allowed")
	errMQTTTopicIsEmpty               = errors.New("topic cannot be empty")
	errMQTTPacketIdentifierIsZero     = errors.New("packet identifier cannot be 0")
	errMQTTUnsupportedCharacters      = errors.New("character not supported for MQTT topics")
	errMQTTInvalidSession             = errors.New("invalid MQTT session")
	errMQTTInvalidRetainFlags         = errors.New("invalid retained message flags")
	errMQTTInvalidRetainedMessage     = errors.New("invalid retained message")
	errMQTTSessionCollision           = errors.New("stored session does not match client ID")
	errMQTTInvalidPublishLength       = errors.New("invalid publish message, variable header exceeds remaining length")
	errMQTTAckPipelineStopped         = errors.New("QoS1 PUBACK pipeline has shut down while admitting a message, " +
		"abandoning the wait for its JetStream ack; failing the connection, " +
		"the client will re-send unacknowledged PUBLISH packets on reconnect")
)

type srvMQTT struct {
	listener     net.Listener
	listenerErr  error
	authOverride bool
	sessmgr      mqttSessionManager
}

type mqttSessionManager struct {
	mu       sync.RWMutex
	sessions map[string]*mqttAccountSessionManager // key is account name

View on GitHub (pinned to 3a66a489d2)

Solutions

  1. Fix the writer so the retained-message flags byte only contains valid bits (retain + QoS bits, QoS <= 2), keeping the value below 0b1111.
  2. Purge or re-store the corrupted retained message in the retained-message stream.
  3. If upgrading, verify retained-message storage encoding compatibility between server versions and re-publish retained messages after migration.

Example fix

// before: flags include invalid high bits
flags := byte(0x0f)
// after: retain bit + QoS only
qos := byte(1)
flags := byte(0x01) | (qos << 1) // valid, < mqttPacketFlagMask
Defensive patterns

Strategy: validation

Validate before calling

func validRetainFlags(flags byte) bool {
    return flags < 0x0f && (flags>>1)&0x03 <= 2
}

Type guard

func isUsableRetainedMsg(rm *mqttRetainedMsg) bool {
    return rm != nil && rm.Flags < mqttPacketFlagMask && mqttGetQoS(rm.Flags) <= 2
}

Try / catch

rm, err := decodeRetainedMsg(b)
if err != nil {
    if strings.Contains(err.Error(), "invalid retained message flags") {
        log.Warn("dropping corrupted retained message")
        return nil
    }
    return err
}

Prevention

When it happens

Trigger: Decoding a retained message whose Flags byte is >= mqttPacketFlagMask (0b1111); or the upstream decode of the retained message failed and the code deliberately replaces that error with errMQTTInvalidRetainFlags.

Common situations: Corrupted or hand-edited retained-message storage/stream entries; a server version change or data migration producing flags encodings the current decoder rejects; writing raw flags bytes in tests or tooling without masking out reserved bits.

Related errors


AI-assisted analysis of nats-io/nats-server@3a66a489d2 (2026-09-02). Data as JSON: /api/errors/fe86a97e2cb67f64. Report an issue: GitHub.