{"record":{"id":"fe86a97e2cb67f64","repo":"nats-io/nats-server","slug":"invalid-retained-message-flags","errorCode":null,"errorMessage":"invalid retained message flags","messagePattern":"invalid retained message flags","errorType":"validation","errorClass":null,"httpStatus":null,"severity":"error","filePath":"server/mqtt.go","lineNumber":245,"sourceCode":"\terrMQTTMalformedVarInt            = errors.New(\"malformed variable int\")\n\terrMQTTSecondConnectPacket        = errors.New(\"received a second CONNECT packet\")\n\terrMQTTServerNameMustBeSet        = errors.New(\"mqtt requires server name to be explicitly set\")\n\terrMQTTUserMixWithUsersNKeys      = errors.New(\"mqtt authentication username not compatible with presence of users/nkeys\")\n\terrMQTTTokenMixWIthUsersNKeys     = errors.New(\"mqtt authentication token not compatible with presence of users/nkeys\")\n\terrMQTTAckWaitMustBePositive      = errors.New(\"ack wait must be a positive value\")\n\terrMQTTJSAPITimeoutMustBePositive = errors.New(\"JS API timeout must be a positive value\")\n\terrMQTTStandaloneNeedsJetStream   = errors.New(\"mqtt requires JetStream to be enabled if running in standalone mode\")\n\terrMQTTConnFlagReserved           = errors.New(\"connect flags reserved bit not set to 0\")\n\terrMQTTWillAndRetainFlag          = errors.New(\"if Will flag is set to 0, Will Retain flag must be 0 too\")\n\terrMQTTPasswordFlagAndNoUser      = errors.New(\"password flag set but username flag is not\")\n\terrMQTTCIDEmptyNeedsCleanFlag     = errors.New(\"when client ID is empty, clean session flag must be set to 1\")\n\terrMQTTEmptyWillTopic             = errors.New(\"empty Will topic not allowed\")\n\terrMQTTEmptyUsername              = errors.New(\"empty user name not allowed\")\n\terrMQTTTopicIsEmpty               = errors.New(\"topic cannot be empty\")\n\terrMQTTPacketIdentifierIsZero     = errors.New(\"packet identifier cannot be 0\")\n\terrMQTTUnsupportedCharacters      = errors.New(\"character not supported for MQTT topics\")\n\terrMQTTInvalidSession             = errors.New(\"invalid MQTT session\")\n\terrMQTTInvalidRetainFlags         = errors.New(\"invalid retained message flags\")\n\terrMQTTInvalidRetainedMessage     = errors.New(\"invalid retained message\")\n\terrMQTTSessionCollision           = errors.New(\"stored session does not match client ID\")\n\terrMQTTInvalidPublishLength       = errors.New(\"invalid publish message, variable header exceeds remaining length\")\n\terrMQTTAckPipelineStopped         = errors.New(\"QoS1 PUBACK pipeline has shut down while admitting a message, \" +\n\t\t\"abandoning the wait for its JetStream ack; failing the connection, \" +\n\t\t\"the client will re-send unacknowledged PUBLISH packets on reconnect\")\n)\n\ntype srvMQTT struct {\n\tlistener     net.Listener\n\tlistenerErr  error\n\tauthOverride bool\n\tsessmgr      mqttSessionManager\n}\n\ntype mqttSessionManager struct {\n\tmu       sync.RWMutex\n\tsessions map[string]*mqttAccountSessionManager // key is account name","sourceCodeStart":227,"sourceCodeEnd":263,"githubUrl":"https://github.com/nats-io/nats-server/blob/3a66a489d262bf89b71a71c955c94920394532f3/server/mqtt.go#L227-L263","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["Fix the writer so the retained-message flags byte only contains valid bits (retain + QoS bits, QoS <= 2), keeping the value below 0b1111.","Purge or re-store the corrupted retained message in the retained-message stream.","If upgrading, verify retained-message storage encoding compatibility between server versions and re-publish retained messages after migration."],"exampleFix":"// before: flags include invalid high bits\nflags := byte(0x0f)\n// after: retain bit + QoS only\nqos := byte(1)\nflags := byte(0x01) | (qos << 1) // valid, < mqttPacketFlagMask","handlingStrategy":"validation","validationCode":"func validRetainFlags(flags byte) bool {\n    return flags < 0x0f && (flags>>1)&0x03 <= 2\n}","typeGuard":"func isUsableRetainedMsg(rm *mqttRetainedMsg) bool {\n    return rm != nil && rm.Flags < mqttPacketFlagMask && mqttGetQoS(rm.Flags) <= 2\n}","tryCatchPattern":"rm, err := decodeRetainedMsg(b)\nif err != nil {\n    if strings.Contains(err.Error(), \"invalid retained message flags\") {\n        log.Warn(\"dropping corrupted retained message\")\n        return nil\n    }\n    return err\n}","preventionTips":["Mask the flags byte to retain + QoS bits when writing retained messages","Keep QoS <= 2 when constructing message flags","Validate retained-message storage entries after data migrations or version upgrades"],"tags":["mqtt","retained-messages","data-corruption"],"backgroundTag":"invalid-retained-message-flags","analyzedSha":"3a66a489d262bf89b71a71c955c94920394532f3","analyzedAt":"2026-09-02T04:41:54.247Z","contentChangedAt":null,"schemaVersion":2},"datasetVersion":"2026-09-08T10:18:20.063Z"}