Skip to content
129 changes: 129 additions & 0 deletions avro/src/test/scala/dev/constructive/eo/avro/AvroBytesSpec.scala
Original file line number Diff line number Diff line change
Expand Up @@ -382,4 +382,133 @@ class AvroBytesSpec extends Specification with ScalaCheck:
.and(AvroBinaryCursor.zigZagInt(Int.MinValue).length === 5)
}

// ---- byte-cursor refusal identities -------------------------------------
//
// Every row below asserts WHICH `AvroFailure` comes back, never `isLeft`. That is the whole
// point: each of these guards, when removed, lets the cursor run on into the Avro runtime, which
// throws, which `locateFrom` catches and reports as `BinaryParseFailed`. A `Left` either way —
// so an `isLeft` assertion pins nothing, and the structured diagnostic the `.record` face
// publishes silently degrades to "parse failed".
//
// Called at the `AvroBinaryCursor` seam (`private[avro]`, same package) rather than through a
// prism: several rows need a path the public drilling macros refuse to build, and the strictness
// flag is not reachable from the surface at all.
//
// covers: AvroBinaryCursor.scala:558 Field step against a non-record schema,
// AvroBinaryCursor.scala:568 UnionBranch step against a non-union schema,
// AvroBinaryCursor.scala:684 declaredBranches' non-union guard,
// AvroBinaryCursor.scala:675 branchOrdinalOf's end-of-alternatives guard,
// AvroBinaryCursor.scala:572 the `requested < 0` refusal under the LENIENT policy,
// AvroBinaryCursor.scala:126 locateElements' non-array terminal,
// AvroBinaryCursor.scala:125 locateElements' strict-union prefix policy
"AvroBinaryCursor.locate: each refusal reports its OWN AvroFailure, not BinaryParseFailed" >> {
val txBytes = toBinary(transactionRecord(Transaction("t-1", Some(42L))), transactionSchema)
val txNullBytes = toBinary(transactionRecord(Transaction("t-2", None)), transactionSchema)
val personBytes = toBinary(personRecord(Person("Alice", 30)), personSchema)
val basketBytes =
toBinary(basketRecord(Basket("ann", List(Order("tea", 2.5, 1)))), basketSchema)

def field(n: String) = PathStep.Field(n)
def branch(n: String) = PathStep.UnionBranch(n)
def at(bytes: Array[Byte], schema: Schema, strict: Boolean, steps: PathStep*) =
AvroBinaryCursor.locate(bytes, schema, steps.toArray, strictTerminalUnion = strict)

// A Field step whose parent schema is a union, not a record.
val notARecord = at(txBytes, transactionSchema, true, field("amount"), field("x")) ===
Left(AvroFailure.NotARecord(field("x")))

// A UnionBranch step whose parent schema is a plain string — and the diagnostic's branch list
// must degrade to Nil rather than interrogating the non-union for alternatives.
val notAUnion = at(personBytes, personSchema, true, field("name"), branch("string")) ===
Left(AvroFailure.UnionResolutionFailed(Nil, branch("string")))

// A branch name that is not declared: the ordinal scan must run off the end and report -1
// rather than indexing past the alternative list. Both policies refuse, and they must refuse
// for THIS reason — the lenient row is the one that separates the `requested < 0` guard from
// the runtime-vs-requested comparison downstream of it.
val unknownBranch = branch("eo.avro.test.NotDeclared")
val declared = Left(AvroFailure.UnionResolutionFailed(List("null", "long"), unknownBranch))
val unknownStrict =
at(txBytes, transactionSchema, true, field("amount"), unknownBranch) === declared
val unknownLenient =
at(txBytes, transactionSchema, false, field("amount"), unknownBranch) === declared

// locateElements: a prefix that resolves to a string, not an array.
val notAnArray = AvroBinaryCursor.locateElements(
basketBytes,
basketSchema,
Array(PathStep.Field("owner")),
Array.empty[PathStep],
) === Left(AvroFailure.NotAnArray(PathStep.Field("owner")))

// locateElements resolves its PREFIX strictly: a union prefix whose runtime branch is `null`
// must fail as a branch mismatch, not be tolerated into a "terminal is not an array".
val elementsPrefixStrict = AvroBinaryCursor.locateElements(
txNullBytes,
transactionSchema,
Array(PathStep.Field("amount"), PathStep.UnionBranch("long")),
Array.empty[PathStep],
) === Left(
AvroFailure.UnionResolutionFailed(List("null", "long"), PathStep.UnionBranch("long"))
)

notARecord
.and(notAUnion)
.and(unknownStrict)
.and(unknownLenient)
.and(notAnArray)
.and(elementsPrefixStrict)
}

// The union-branch POLICY, at both levels of the byte face. Both halves are invisible to every
// other assertion in this module: the payloads still decode, and the decoded values are right.
//
// covers: AvroBinaryCursor.scala:589 the terminal-vs-interior union split,
// AvroPrism.scala:157 the byte-face read's strict terminal-union resolution
"byte-face union policy: an interior step does not anchor the span, a mismatch does not decode" >> {
// An interior union step must NOT anchor the returned span — the span belongs to the step the
// path ends on. Every observable downstream of a mis-anchored span (`getOption`, `.modify`)
// still round-trips, because the span still opens on a decodable value; only `valueSchema`
// gives it away.
val envelope = WireEnvelope("e-1", 7L, Cash(100L), "note")
val envBytes = toBinary(envelopeRecord(envelope), envelopeSchema)
val cashName = summon[AvroCodec[Cash]].schema.getFullName

val interiorOk = AvroBinaryCursor
.locate(
envBytes,
envelopeSchema,
Array(PathStep.Field("payment"), PathStep.UnionBranch(cashName), PathStep.Field("amount")),
strictTerminalUnion = true,
)
.map(_.valueSchema.getType) === Right(Schema.Type.LONG)

// ... and the READ resolves its terminal union strictly. `TwinA` and `TwinB` encode
// identically, so a lenient read does not fail: it DECODES the other branch and hands back a
// well-formed value of the wrong type. The lenient policy is correct for graft/write only.
val twinBytes = toBinaryValue(summon[AvroCodec[TwinA]].encode(TwinA(100L)), twinSchema)
val sameBranch =
AvroPrism.codecPrism[Twin].union[TwinA].getOption(twinBytes) === Some(TwinA(100L))
val crossBranch = AvroPrism.codecPrism[Twin].union[TwinB].getOption(twinBytes) === None

interiorOk.and(sameBranch).and(crossBranch)
}

// covers: AvroTraversal.scala:193 — `.each`'s per-element drilling resolves its prefix to an
// ARRAY or refuses loudly. The public `.each` macro rejects a non-array focus at COMPILE time,
// so this guard is only reachable by constructing the traversal at the internal seam — which
// is exactly what a future drilling entry point would do.
"AvroTraversal: drilling through a non-array prefix throws IllegalArgumentException" >> {
val leaf = new AvroFocus.Leaf[String](Array.empty[PathStep], summon[AvroCodec[String]])
val bogus = new AvroTraversal[String](Array(PathStep.Field("name")), leaf, personSchema)

val thrown =
try
bogus.widenSuffixNamed[String]("whatever")
"no throw"
catch case e: IllegalArgumentException => e.getMessage

thrown must contain("prefix does not point at an array")
}

end AvroBytesSpec
Original file line number Diff line number Diff line change
@@ -1,12 +1,14 @@
package dev.constructive.eo.avro

import scala.language.implicitConversions

import java.io.ByteArrayInputStream
import java.util.concurrent.ConcurrentLinkedQueue
import org.apache.avro.Schema
import org.apache.avro.generic.{GenericData, GenericDatumReader, GenericRecord, IndexedRecord}
import org.apache.avro.io.DecoderFactory
import org.scalacheck.Gen
import org.scalacheck.Prop.forAll
import org.scalacheck.Prop.forAllNoShrink
import org.specs2.ScalaCheck
import org.specs2.mutable.Specification

Expand All @@ -18,7 +20,9 @@ import org.specs2.mutable.Specification
* false` opt-out — must produce a record equal to the OLD fresh-allocation reference decode
* (reproduced verbatim in [[freshDecode]]), across writer→reader evolution (field dropped,
* added-with-default, reordered, promoted) and union-typed payloads. This is the load-bearing
* correctness guard: a wrong decode in a serde library corrupts every consumer.
* correctness guard: a wrong decode in a serde library corrupts every consumer. The evolution
* shapes are DATA ([[reuseScenarios]]) rather than one example each, so every shape is put
* through both entry points.
* 2. '''Thread-safe + non-aliased.''' Concurrent decodes on many threads stay correct, and a
* record decoded earlier on a thread is never mutated by a later decode on that thread (fresh
* datum, no `Utf8`/bytes aliasing).
Expand Down Expand Up @@ -93,81 +97,96 @@ class AvroCodecDecoderReuseSpec extends Specification with ScalaCheck:
yield AvroSpecFixtures.Transaction(id, amount)

// ---- 1. Byte-identical decode vs the fresh-allocation reference ----
//
// Seven near-identical properties collapsed into one: they differed ONLY in the (generator,
// writer schema, reader schema) triple and in which entry point they called, so the triple is
// now data and both entry points — the thread-local cache and the `threadLocalStorage = false`
// opt-out — are checked on every scenario instead of one each.

/** One writer→reader shape. `extra` is a scenario-specific sanity check on the reference decode,
* so a scenario can assert that the evolution it names really happened.
*/
final private case class Reuse(
name: String,
bytes: Gen[Array[Byte]],
writer: Schema,
reader: Schema,
extra: IndexedRecord => Boolean,
)

"non-resolved decodeRecord matches the fresh reference decode (WriterEvent)" >> forAll(
genWriterEvent
) { e =>
val bytes = binaryOf(e)
AvroCodec
.decodeRecord(bytes, writerSchema)
.exists(_ == freshDecode(bytes, writerSchema, writerSchema))
}

"non-resolved decodeRecord matches the fresh reference decode (union payload)" >> forAll(
genTransaction
) { t =>
val bytes = binaryOf(t)
AvroCodec
.decodeRecord(bytes, transactionSchema)
.exists(_ == freshDecode(bytes, transactionSchema, transactionSchema))
}

"resolved decode matches the fresh reference (WriterEvent→ReaderEvent, field dropped)" >> forAll(
genWriterEvent
) { e =>
val bytes = binaryOf(e)
AvroCodec
.decodeResolvedRecord(bytes, writerSchema, readerEventSchema)
.exists(_ == freshDecode(bytes, writerSchema, readerEventSchema))
}

"resolved decode matches the fresh reference (fields reordered, resolved by name)" >> forAll(
genReorderWriter
) { w =>
val bytes = binaryOf(w)
AvroCodec
.decodeResolvedRecord(bytes, reorderWriterSchema, reorderReaderSchema)
.exists(_ == freshDecode(bytes, reorderWriterSchema, reorderReaderSchema))
}

"resolved decode matches the fresh reference (int→long promotion)" >> forAll(genPromoteWriter) {
w =>
val bytes = binaryOf(w)
AvroCodec
.decodeResolvedRecord(bytes, promoteWriterSchema, promoteReaderSchema)
.exists(_ == freshDecode(bytes, promoteWriterSchema, promoteReaderSchema))
}
private def always: IndexedRecord => Boolean = _ => true

"resolved decode matches the fresh reference (field added with default)" >> forAll(genId) { id =>
private val addWriterBytes: Gen[Array[Byte]] = genId.map { id =>
val rec = new GenericData.Record(addWriterSchema)
rec.put("id", id)
val bytes = AvroSpecFixtures.toBinary(rec, addWriterSchema)
val cached = AvroCodec.decodeResolvedRecord(bytes, addWriterSchema, addReaderSchema)
val reference = freshDecode(bytes, addWriterSchema, addReaderSchema)
// Sanity: the reader default really did materialise, so this is a live added-field case.
val defaultApplied = reference.get(addReaderSchema.getField("note").pos).toString == "n/a"
cached.exists(_ == reference) && defaultApplied
AvroSpecFixtures.toBinary(rec, addWriterSchema)
}

// ---- 2. threadLocalStorage = false opt-out (fresh allocation per call) ----

"resolved decode with threadLocalStorage = false matches the fresh reference decode" >> forAll(
genWriterEvent
) { e =>
val bytes = binaryOf(e)
AvroBinaryCursor
.records
.read(
bytes,
0,
bytes.length,
writerSchema,
readerEventSchema,
threadLocalStorage = false,
) == freshDecode(bytes, writerSchema, readerEventSchema)
}
private val reuseScenarios: List[Reuse] = List(
Reuse("non-resolved", genWriterEvent.map(binaryOf(_)), writerSchema, writerSchema, always),
Reuse(
"non-resolved, union payload",
genTransaction.map(binaryOf(_)),
transactionSchema,
transactionSchema,
always,
),
Reuse(
"field dropped",
genWriterEvent.map(binaryOf(_)),
writerSchema,
readerEventSchema,
always,
),
Reuse(
"fields reordered, resolved by name",
genReorderWriter.map(binaryOf(_)),
reorderWriterSchema,
reorderReaderSchema,
always,
),
Reuse(
"int→long promotion",
genPromoteWriter.map(binaryOf(_)),
promoteWriterSchema,
promoteReaderSchema,
always,
),
// The reader default must really materialise, else this is not a live added-field case.
Reuse(
"field added with default",
addWriterBytes,
addWriterSchema,
addReaderSchema,
r => r.get(addReaderSchema.getField("note").pos).toString == "n/a",
),
)

private val genScenario: Gen[(Reuse, Array[Byte])] =
Gen.oneOf(reuseScenarios).flatMap(s => s.bytes.map(b => (s, b)))

"every decode entry point reproduces the fresh-allocation reference, on every evolution shape" >>
forAllNoShrink(genScenario) { (scenario, bytes) =>
val reference = freshDecode(bytes, scenario.writer, scenario.reader)
val cached =
if scenario.writer == scenario.reader then AvroCodec.decodeRecord(bytes, scenario.writer)
else AvroCodec.decodeResolvedRecord(bytes, scenario.writer, scenario.reader)
val uncached = AvroBinaryCursor
.records
.read(
bytes,
0,
bytes.length,
scenario.writer,
scenario.reader,
threadLocalStorage = false,
)
(cached.exists(_ == reference) && (uncached == reference) && scenario.extra(
reference
)) :| scenario.name
}

// ---- 3. Concurrent + non-aliased ----
// ---- 2. Concurrent + non-aliased ----

"concurrent decodes across threads are correct and never alias an earlier record" >> {
val events = (0 until 8).map(i => WriterEvent(s"id-$i", i)).toVector
Expand Down
Loading