Skip to content

Commit 8b3de74

Browse files
vparfonovopenshift-merge-bot[bot]
authored andcommitted
test(buffer): add functional test for disk buffer corruption recovery (LOG-9386)
Verify that Vector does not CrashLoop when the disk buffer contains corrupted protobuf records. The test corrupts buffer record payloads while the collector is running, then kills it so the restart reads the corrupted buffer. Asserts that Vector logs a warning about corrupted records during seek, starts successfully, and continues delivering logs after recovery. Signed-off-by: Vitalii Parfonov <vparfono@redhat.com>
1 parent b948959 commit 8b3de74

3 files changed

Lines changed: 443 additions & 0 deletions

File tree

‎test/framework/functional/framework.go‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -192,6 +192,22 @@ func (f *CollectorFunctionalFramework) AddOutputContainersVisitors() []runtime.P
192192
return visitors
193193
}
194194

195+
const ToolsContainerName = "tools"
196+
197+
func (f *CollectorFunctionalFramework) AddToolsContainerVisitor(volumeName, mountPath string) runtime.PodBuilderVisitor {
198+
return func(b *runtime.PodBuilder) error {
199+
b.AddEmptyDirVolume(volumeName)
200+
b.GetContainer(constants.CollectorName).
201+
AddVolumeMount(volumeName, mountPath, "", false).
202+
Update()
203+
b.AddContainer(ToolsContainerName, "registry.access.redhat.com/ubi9/ubi-minimal:latest").
204+
AddVolumeMount(volumeName, mountPath, "", false).
205+
WithCmd([]string{"sleep", "infinity"}).
206+
End()
207+
return nil
208+
}
209+
}
210+
195211
// Deploy the objects needed to functional Test
196212
func (f *CollectorFunctionalFramework) Deploy() (err error) {
197213
return f.DeployWithVisitors(f.AddOutputContainersVisitors())
Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
1+
package misc
2+
3+
import (
4+
"encoding/binary"
5+
"fmt"
6+
"hash/crc32"
7+
"os"
8+
)
9+
10+
// corruptDiskBufferPayload corrupts protobuf payloads in a Vector disk buffer
11+
// data file while keeping CRC32 checksums valid, causing ReaderError::Decode
12+
// on startup seek. The last record is left intact (writer validates it).
13+
// Returns the number of corrupted records.
14+
func corruptDiskBufferPayload(datFilePath string) (int, error) {
15+
data, err := os.ReadFile(datFilePath)
16+
if err != nil {
17+
return 0, fmt.Errorf("read file: %w", err)
18+
}
19+
20+
type record struct{ start, length int }
21+
var records []record
22+
offset := 0
23+
for offset+8 < len(data) {
24+
recLen := int(binary.BigEndian.Uint64(data[offset : offset+8]))
25+
recStart := offset + 8
26+
if recLen < 32 || recStart+recLen > len(data) {
27+
break
28+
}
29+
records = append(records, record{recStart, recLen})
30+
offset = recStart + recLen
31+
}
32+
33+
if len(records) < 2 {
34+
return 0, fmt.Errorf("need at least 2 records to corrupt (found %d); last record must stay intact for writer validation", len(records))
35+
}
36+
37+
corruptUpTo := len(records) - 1
38+
corrupted := 0
39+
for i := 0; i < corruptUpTo; i++ {
40+
r := records[i]
41+
if err := corruptRecordPayload(data[r.start : r.start+r.length]); err != nil {
42+
return 0, fmt.Errorf("corrupt record %d: %w", i, err)
43+
}
44+
corrupted++
45+
}
46+
47+
if err := os.WriteFile(datFilePath, data, 0o640); err != nil {
48+
return 0, err
49+
}
50+
return corrupted, nil
51+
}
52+
53+
// corruptRecordPayload XOR-flips the payload bytes of a single rkyv-serialized
54+
// record and recalculates the CRC32 checksum so the record passes checksum
55+
// validation but fails protobuf decode.
56+
func corruptRecordPayload(recordData []byte) error {
57+
if len(recordData) < 32 {
58+
return fmt.Errorf("record too small: %d bytes", len(recordData))
59+
}
60+
61+
// ArchivedRecord layout (repr(C), little-endian, 32 bytes at end of record):
62+
// +0: checksum u32, +4: padding, +8: id u64,
63+
// +16: metadata u32, +20: rel_ptr i32, +24: payload_len u32, +28: padding
64+
structOff := len(recordData) - 32
65+
66+
id := binary.LittleEndian.Uint64(recordData[structOff+8 : structOff+16])
67+
metadata := binary.LittleEndian.Uint32(recordData[structOff+16 : structOff+20])
68+
relPtr := int32(binary.LittleEndian.Uint32(recordData[structOff+20 : structOff+24]))
69+
payloadLen := binary.LittleEndian.Uint32(recordData[structOff+24 : structOff+28])
70+
71+
payloadStart := int(structOff+20) + int(relPtr)
72+
if payloadStart < 0 || payloadStart+int(payloadLen) > len(recordData) {
73+
return fmt.Errorf("payload out of bounds: start=%d len=%d record=%d", payloadStart, payloadLen, len(recordData))
74+
}
75+
76+
payload := recordData[payloadStart : payloadStart+int(payloadLen)]
77+
for i := range payload {
78+
payload[i] ^= 0xFF
79+
}
80+
81+
// Recalculate CRC32-IEEE: BE(id) + BE(metadata) + payload
82+
h := crc32.NewIEEE()
83+
idBE := make([]byte, 8)
84+
binary.BigEndian.PutUint64(idBE, id)
85+
_, _ = h.Write(idBE)
86+
metaBE := make([]byte, 4)
87+
binary.BigEndian.PutUint32(metaBE, metadata)
88+
_, _ = h.Write(metaBE)
89+
_, _ = h.Write(payload)
90+
binary.LittleEndian.PutUint32(recordData[structOff:structOff+4], h.Sum32())
91+
92+
return nil
93+
}

0 commit comments

Comments
 (0)