Skip to content

Commit

Permalink
[NUC-233] Troubeshooting (#152)
Browse files Browse the repository at this point in the history
  • Loading branch information
alxtkr77 authored Aug 20, 2024
1 parent 919d54d commit 8d1f15a
Show file tree
Hide file tree
Showing 2 changed files with 30 additions and 0 deletions.
4 changes: 4 additions & 0 deletions pkg/dataplane/streamconsumergroup/sequencenumberhandler.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,10 @@ func (snh *sequenceNumberHandler) stop() error {
}

func (snh *sequenceNumberHandler) markShardSequenceNumber(shardID int, sequenceNumber uint64) error {
err := snh.member.streamConsumerGroup.checkShardExists(shardID, sequenceNumber)
if err != nil {
return errors.Wrapf(err, "Failed checking shard exists. Current sequenceNumber %v", snh.markedShardSequenceNumbers[shardID])
}
snh.markedShardSequenceNumbers[shardID] = sequenceNumber

return nil
Expand Down
26 changes: 26 additions & 0 deletions pkg/dataplane/streamconsumergroup/streamconsumergroup.go
Original file line number Diff line number Diff line change
Expand Up @@ -367,10 +367,36 @@ func (scg *streamConsumerGroup) setShardSequenceNumberInPersistency(shardID int,
Attributes: map[string]interface{}{
scg.getShardCommittedSequenceNumberAttributeName(): sequenceNumber,
},
Condition: "__obj_type == 3",
})
return err
}

func (scg *streamConsumerGroup) checkShardExists(shardID int, sequenceNumber uint64) error {
shardPath, err := scg.getShardPath(shardID)
if err != nil {
return errors.Wrapf(err, "Failed getting shard path: %v", shardID)
}

response, err := scg.container.GetItemSync(&v3io.GetItemInput{
Path: shardPath,
AttributeNames: []string{
"__obj_type",
},
})
defer response.Release()
if err != nil {
return errors.Wrapf(err, "Updating shard counter %v to %v failed.Failed to fetch shard object type", shardPath, sequenceNumber)
}
getItemOutput := response.Output.(*v3io.GetItemOutput)

objType, err2 := getItemOutput.Item.GetFieldInt("__obj_type")
if err2 != nil || objType != 3 {
return errors.Wrapf(err2, "Updating shard counter %v to %v failed. Shard is a regular file (%v)", shardPath, sequenceNumber, objType)
}
return nil
}

// returns true if the states are equal, ignoring heartbeat times
func (scg *streamConsumerGroup) statesEqual(state0 *State, state1 *State) bool {
if state0.SchemasVersion != state1.SchemasVersion {
Expand Down

0 comments on commit 8d1f15a

Please sign in to comment.