mirror of
https://github.com/zjdyzww/harness.git
synced 2026-09-28 14:24:00 +08:00
* 04d35c fix: CODE-5680: apply goimports formatting to wire_gen.go * 6c006f fix: CODE-5680: regenerate gitness wire_gen.go * 3b074c fix: CODE-5680: update create_registry_test NewSystem call with noop collector * 74d32e fix: CODE-5680: update usage_test.go ProvideSystem call with noop collector * ec0a94 fix: CODE-5680: move Prometheus impl out of gitness, keep only interface and noop * 6a19fd fix: CODE-5680: pass Collector as constructor param to NewSystem, remove from WireSet * 59275f fix: CODE-5680: refactor events Collector onto System, no changes to consumer wires * fe3723 fix: CODE-5680: refactor events stream metrics to use injected Collector pattern * 55cf54 fix: CODE-5680: use DiscardedMessageError sentinel type for discard detection * 1d752b fix: CODE-5680: revert redis_consumer.go changes, keep string matching for discard detection * 1c398c fix: CODE-5680: address PR review comments on events stream metrics * e1bf47 feat: CODE-5680: add Prometheus metrics for event
40 lines
1.4 KiB
Go
40 lines
1.4 KiB
Go
// Copyright 2023 Harness, Inc.
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
package stream
|
|
|
|
import "fmt"
|
|
|
|
// DiscardedMessageError is returned (via pushError) when a message is force-acknowledged
|
|
// after exceeding the maximum retry count, regardless of whether the XAck call itself succeeded.
|
|
type DiscardedMessageError struct {
|
|
MessageID string
|
|
StreamID string
|
|
RetryCount int64
|
|
AckErr error // non-nil if the XAck call also failed
|
|
}
|
|
|
|
func (e *DiscardedMessageError) Error() string {
|
|
if e.AckErr != nil {
|
|
return fmt.Sprintf(
|
|
"failed to force acknowledge (discard) message '%s' (Retries: %d) in stream '%s': %s",
|
|
e.MessageID, e.RetryCount, e.StreamID, e.AckErr)
|
|
}
|
|
return fmt.Sprintf(
|
|
"force acknowledged (discarded) message '%s' (Retries: %d) in stream '%s'",
|
|
e.MessageID, e.RetryCount, e.StreamID)
|
|
}
|
|
|
|
func (e *DiscardedMessageError) Unwrap() error { return e.AckErr }
|