Skip to content

Add Variant, Decimal, and Timestamp CEL functions - #2332

Merged
Robert Yokota (rayokota) merged 76 commits into
masterfrom
add-cel-logical-types-2
Sep 24, 2026
Merged

Robert Yokota (rayokota) merged 76 commits into
masterfrom
add-cel-logical-types-2

Conversation

@rayokota

@rayokota Robert Yokota (rayokota) commented Aug 24, 2026 •

Copy link
Copy Markdown
Member

What

Summary

Adds three families of CEL functions to data contract rules, bringing this client to parity
with the JVM reference implementation:

  • variant(...) / variants.* — read and navigate a Spark/Parquet Variant
  • decimal(...) / decimals.* — exact decimal arithmetic and comparison
  • timestamp(...) extensions — construct from an epoch value at a given precision, and
    accept the temporal shapes the Avro and Protobuf decoders produce

Before this, a rule could not work with any of these three types: a confluent.type.Decimal
or google.protobuf.Timestamp field reached CEL as an opaque message, an Avro decimal reached
it as raw unscaled bytes, and a Variant was unreachable entirely.

What's added

Variant — variant(dyn) and variant(value, metadata) constructors; variants.parseJson
(strict) and variants.tryParseJson (CEL null on a malformed document); variants.type;
navigation via variants.field, variants.index and variants.path (a JSONPath subset:
$, $.field, $[i], $["quoted key"]); typed extraction via variants.as / variants.tryAs;
plus variants.isNull and variants.toJson.

Decimal — decimal(...) from a string, int, uint, double or unscaled-bytes-plus-scale;
arithmetic (add, sub, mul, div, mod); rounding (round, trunc, floor, ceil);
abs and sign; comparisons; and string(...) / double(...) extended to accept a Decimal.

Timestamp — timestamp(value, precision) where precision is one of {0, 3, 6, 9}
(seconds, millis, micros, nanos), and a timestamp(dyn) overload accepting the temporal
representations a decoder hands back. string(...) renders a timestamp with its sub-second
component.

Marshalling boundary — the schema-side value is converted to its CEL type on the way in
and back to the schema's representation on the way out, for both field-level (CEL_FIELD) and
message-level (CEL) rules, across Avro, Protobuf and JSON Schema. A decimal keeps its scale,
a timestamp keeps its unit, and a Variant round-trips as a Variant. A Variant is
converted at the boundary too, but only a message-level rule reaches it.

Semantics

The JVM client is the contract; behaviour here is matched against it rather than against this
language's native conventions. In particular:

  • Exact arithmetic. add, sub, mul and mod are exact, as java.math.BigDecimal is.
    Division is capped at 38 significant digits with HALF_UP, matching the JVM's DIV_MC.
  • Scale is part of the value. 12.34 and 12.340 are the same number in two encodings and
    are rendered differently; round/trunc produce exactly the requested scale, including a
    negative one (round(1234, -2) is 1200). Scale arguments are int32-bounded, as
    BigDecimal's are.
  • Wire form. The unscaled value is minimal big-endian two's complement, byte-identical to
    BigInteger.toByteArray(), and precision is the unscaled value's digit count as
    BigDecimal.precision() reports it.
  • Range and type checks are errors, not coercions. A non-finite double, a timestamp outside
    0001-01-01T00:00:00Z .. 9999-12-31T23:59:59.999999999Z, an out-of-range scale, or a
    wrong-typed argument is a rule error rather than something silently narrowed — the JVM's
    typed overloads reject the same inputs.

Known limitations

These are deliberate and shared across the non-JVM clients:

  • float/double JSON rendering stays native to this language. Byte-identical rendering
    across all clients was designed and implemented, then backed out: the precision walk it
    requires costs 14–23× a native format call and about 75% of the serialization path, and no
    cross-client bug had been reported against it. Values are equal; their shortest-form text may
    differ.
  • precision is informational on read. The JVM applies it as a MathContext when decoding;
    this client returns the value unrounded. Since every client now writes precision as the
    value's own digit count, the two agree for anything these clients produce.
  • Avro local-timestamp-* is not converted. It carries no zone, so the JVM refuses to turn
    it into an instant; conversion support here is tracked separately.

Checklist

  • Contains customer facing changes? Including API/behavior changes
  • Did you add sufficient unit test and/or integration test coverage for this PR?
    • If not, please explain why it is not required

References

JIRA:

Test & Review

Open questions / Follow-ups

Copilot AI lite review requested due to automatic review settings August 24, 2026 23:57
@confluent-cla-assistant

Copy link
Copy Markdown

🎉 All Contributor License Agreements have been signed. Ready to merge.
Please push an empty commit if you would like to re-run the checks to verify CLA status for all contributors.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Adds CEL support for Variant, Decimal, and Timestamp values, with Avro/Protobuf integration and serializer and validator tests.

Changes:

  • Adds Variant codecs, builders, JSON conversion, and path navigation.
  • Adds Decimal conversions/operators and Timestamp overloads.
  • Extends CEL dispatch and serialization integrations.
  • Adds unit and synchronous/asynchronous integration tests.

Reviewed changes

Copilot reviewed 19 out of 20 changed files in this pull request and generated 14 comments.

Show a summary per file
File Summary
tests/schema_registry/test_variant_utils.py Variant codec and timestamp tests.
tests/schema_registry/test_cel_validator.py CEL behavior and integration tests.
tests/schema_registry/_sync/test_proto_serdes.py Synchronous Protobuf integration tests.
tests/schema_registry/_sync/test_avro_serdes.py Synchronous Avro integration tests.
tests/schema_registry/_async/test_proto_serdes.py Asynchronous Protobuf integration tests.
tests/schema_registry/_async/test_avro_serdes.py Asynchronous Avro integration tests.
src/confluent_kafka/schema_registry/rules/cel/variant_path.py Variant path parsing. Nit (2 votes): identifier checks accept Unicode instead of the documented ASCII grammar.
src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py Variant CEL functions. Moderate (2 votes): tryParseJson accepts non-string inputs. Moderate (3 votes): index conversion truncates doubles and accepts booleans.
src/confluent_kafka/schema_registry/rules/cel/timestamp_funcs.py Timestamp CEL overloads. Moderate (4 votes): two-argument overflow escapes as a raw exception. Moderate (3 votes): naive timestamps can bypass rejection.
src/confluent_kafka/schema_registry/rules/cel/extra_func.py Registers extended CEL functions.
src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py Decimal CEL functions. Moderate (2 votes): BoolType is treated as an integer. Moderate (2 votes): nested Decimal message wrappers are not converted correctly.
src/confluent_kafka/schema_registry/rules/cel/cel_validator.py CEL validation and Decimal boundary conversion.
src/confluent_kafka/schema_registry/rules/cel/cel_field_presence.py Namespaced CEL dispatch.
src/confluent_kafka/schema_registry/rules/cel/cel_executor.py CEL value conversion and lazy now binding.
src/confluent_kafka/schema_registry/confluent/types/variant.proto Variant Protobuf schema.
src/confluent_kafka/schema_registry/confluent/types/variant_utils.py Variant codec and builder. Moderate (3 votes): truncated decimals raise IndexError. Moderate (2 votes): integer capacity checks reserve too much space. Moderate (4 votes): negative zero loses its sign in JSON. Moderate (2 votes): decimal capacity checks reject values that fit their selected width.
src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py Generated Variant Protobuf bindings.
src/confluent_kafka/schema_registry/confluent/types/decimal_utils.py Decimal Protobuf conversions. Moderate (4 votes): ambient precision can round large values. Moderate (4 votes): negative boundary values produce non-canonical bytes.
src/confluent_kafka/schema_registry/common/protobuf.py Variant Protobuf integration.
src/confluent_kafka/schema_registry/common/avro.py Avro Variant logical-type integration. Critical (1 vote): logical handlers are registered through incorrect fastavro objects, preventing Variant round-tripping.
Files not reviewed (1)
  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
Suppressed comments (12)

src/confluent_kafka/schema_registry/confluent/types/variant_utils.py:1100

  • The builder accepts a caller-supplied size_limit, but _integer_size() still returns a three-byte width for values above 0xFFFFFF. Any container or metadata larger than that then fails in to_bytes(3) with a raw OverflowError even though the configured limit permits it. Return a four-byte width after the 24-bit range (the header already supports four widths).
def _integer_size(value: int) -> int:
    if value <= U8_MAX:
        return 1
    if value <= U16_MAX:
        return 2
    return U24_SIZE

src/confluent_kafka/schema_registry/confluent/types/variant_utils.py:274

  • Java Float.toString(-0.0f) also preserves the sign, but this branch formats it with int(f) as 0.0. The resulting Variant JSON differs from the documented Java contract; exclude zero from the integer branch so the existing repr() path retains -0.0.
    if f == int(f) and abs(f) < 1e16:
        return "%d.0" % int(f)

src/confluent_kafka/schema_registry/confluent/types/variant_utils.py:278

  • The float formatter claims to match Java Float.toString, but repr(float(s)) uses Python's exponent formatting. For example, a stored float32 value around 1e-7 renders as 1e-07, whereas Java renders 1.0E-7; exact to_json() comparisons therefore diverge for scientific-notation values. Use a formatter with Java's exponent thresholds/casing and required mantissa digit instead of returning Python repr() directly.
    for p in range(1, 10):
        s = "%.*g" % (p, f)
        if struct.unpack("<f", struct.pack("<f", float(s)))[0] == f:
            return repr(float(s))

src/confluent_kafka/schema_registry/confluent/types/variant_utils.py:186

  • Metadata field names are part of the Variant UTF-8 contract, but a malformed byte sequence raises UnicodeDecodeError directly here rather than VariantError. For a raw/protobuf Variant this escapes the CEL function boundary as an unhandled Python exception; normalize invalid UTF-8 to the codec's documented malformed-input error.
    return metadata[string_start + offset:string_start + next_offset].decode("utf-8")

src/confluent_kafka/schema_registry/confluent/types/variant_utils.py:465

  • A malformed UTF-8 string payload also raises UnicodeDecodeError directly from get_string(), despite VariantError being the reader's documented malformed-input exception. This is especially visible through variants.as(..., 'string'), where the raw exception bypasses CEL error handling; catch the decode error and raise VariantError.
        return self.value[start:start + length].decode("utf-8")

src/confluent_kafka/schema_registry/confluent/types/variant_utils.py:538

  • Negative field indexes are not validated here, so Python's negative indexing returns the last object field instead of rejecting the index. get_element_at_index explicitly rejects negative indexes, and JSONPath declares the same non-negative rule; validate the field index before indexing the encoded tables.
        key_id, value_pos = self._field_id_and_offset(idx)

src/confluent_kafka/schema_registry/rules/cel/cel_field_presence.py:139

  • The namespace override calls func(*args) directly, bypassing celpy's normal conversion of function exceptions into CELEvalError. The new variants.* functions can raise VariantError/IndexError from malformed wire data (for example, a proto Variant with invalid metadata), so CelValidator.execute then leaks the raw exception instead of raising its documented RuleError; preserve existing CELEvalError and normalize other runtime exceptions at this dispatch boundary.
                                return func(*args)

src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py:382

  • As with _string, this arm only handles a raw Decimal. A selected protobuf decimal field is a celpy MessageType wrapper, so double(this.decimal_field) falls through to DoubleType with a mapping and raises instead of performing the documented decimal-to-double conversion. Reuse decimal_boundary_value() before delegating.
    if isinstance(v, Decimal):
        return celtypes.DoubleType(float(v))
    return _STDLIB_DOUBLE(v)

src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py:62

  • The (bytes, scale) overload is declared with an integer scale, but int(scale) silently truncates doubles and accepts CEL booleans (2.9 becomes scale 2, true becomes 1). This can produce a valid but unintended decimal instead of reporting an invalid overload argument; validate the CEL integer type before conversion.
def _from_bytes_scale(value: typing.Any, scale: typing.Any) -> Decimal:
    """Construct a Decimal from raw two's-complement big-endian bytes + scale."""
    raw = _coerce_bytes(value)
    s = int(scale)
    if len(raw) == 0:
        return Decimal(0).scaleb(-s, context=_EXACT_CONTEXT)
    return Decimal(int.from_bytes(raw, "big", signed=True)).scaleb(-s, context=_EXACT_CONTEXT)

src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py:296

  • The target scale is documented as an integer, but int(args[1]) silently truncates a CEL double (for example, decimals.round(d, 1.9) rounds at scale 1) and accepts booleans. Validate the CEL integer type rather than coercing arbitrary values; the same validation should be shared with the other scale-taking overloads.
def _decimals_round(*args: typing.Any) -> Decimal:
    """Round to the given scale (HALF_UP). One-arg form rounds to integer."""
    if len(args) == 1:
        return _d(args[0]).quantize(
            Decimal(1), rounding=decimal.ROUND_HALF_UP, context=_EXACT_CONTEXT)
    if len(args) == 2:
        scale = int(args[1])
        return _d(args[0]).quantize(
            Decimal(1).scaleb(-scale), rounding=decimal.ROUND_HALF_UP,
            context=_EXACT_CONTEXT)

src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py:324

  • As in decimals.round, int(args[1]) silently truncates a non-integer CEL value and accepts booleans even though this overload requires an integer target scale. This can truncate at a scale different from the caller's value; reuse the shared integer-scale validation before conversion.
    if len(args) == 2:
        d = _d(args[0])
        scale = int(args[1])
        if scale >= -d.as_tuple().exponent:
            return d
        return d.quantize(
            Decimal(1).scaleb(-scale), rounding=decimal.ROUND_DOWN,
            context=_EXACT_CONTEXT)

src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:230

  • variants.field is documented with a string key, but str(key) accepts arbitrary CEL values. On an object containing a numeric-looking key, variants.field(v, 1) can silently access "1" instead of reporting a bad argument type, unlike the strict parseJson overload. Validate str/StringType before coercing the key.
    return v.get_field_by_key(str(key))

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread src/confluent_kafka/schema_registry/common/avro.py
Comment thread src/confluent_kafka/schema_registry/confluent/types/decimal_utils.py Outdated
Comment thread src/confluent_kafka/schema_registry/confluent/types/decimal_utils.py Outdated
Comment thread src/confluent_kafka/schema_registry/confluent/type/variant_utils.py
Comment thread src/confluent_kafka/schema_registry/rules/cel/timestamp_funcs.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/timestamp_funcs.py
Comment thread src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py
Comment thread src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py
Comment thread src/confluent_kafka/schema_registry/rules/cel/variant_path.py Outdated
@sonarqube-confluent

Copy link
Copy Markdown

Quality Gate failed Quality Gate failed

Failed conditions
76.4% Coverage on New Code (required ≥ 80%)

See analysis details on SonarQube

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Decimal serialization, timestamp formatting, protobuf compatibility, and CEL argument handling have unresolved correctness issues.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Files not reviewed (3)

  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_type_rules_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_types_pb2.py: Generated file

Suppressed comments (2)

src/confluent_kafka/schema_registry/common/protobuf.py:313

  • This length formula is not minimal for negative signed-byte boundaries: -128 is emitted as ff80 rather than Java BigInteger.toByteArray()'s 80, and -32768 similarly gains a leading ff. That contradicts the stated cross-client encoding and makes serialized decimal bytes differ for these values.
    length = (unscaled.bit_length() + 8) // 8
    return unscaled.to_bytes(length, byteorder="big", signed=True)

src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py:207

  • The signed-byte sizing adds a redundant sign byte at every negative boundary (-128 becomes ff80 instead of 80). Since message-level transforms are intended to match the other clients' minimal two's-complement representation, calculate negative magnitude from ~unscaled.
    length = (unscaled.bit_length() + 8) // 8
    return unscaled.to_bytes(length, byteorder="big", signed=True)
  • Files reviewed: 29/32 changed files
  • Comments generated: 8
  • Review effort level: Balanced

Comment thread tests/schema_registry/data/proto/value_type_rules_pb2.py Outdated
Comment thread tests/schema_registry/data/proto/value_types_pb2.py Outdated
Comment thread src/confluent_kafka/schema_registry/common/protobuf.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/timestamp_funcs.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Variant absence handling, Decimal validation, and protobuf write-back have unresolved correctness gaps.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Files not reviewed (3)

  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_type_rules_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_types_pb2.py: Generated file

Suppressed comments (3)

Previously missed (3) — in code that hasn't changed since the last review.

src/confluent_kafka/schema_registry/common/avro.py:36

  • The absent-Avro tests bypass this logical reader by passing a raw map directly. In an actual fastavro decode, an empty {metadata, value} record reaches this hook and Variant(...) raises before CEL can convert it to null. Preserve the empty record shape so the CEL coercion path can recognize absence.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:110
  • Absence is represented by both buffers being empty, but this treats any empty metadata as absent. A corrupt value such as value=b'\x00', metadata=b'' is therefore silently converted to CEL null instead of reporting malformed data.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:359
  • variants.tryAs declares a string type argument, but str(type_str) accepts every CEL value. Consequently variants.tryAs(v, 1) returns null as though extraction failed instead of reporting an invalid call, unlike the typed overload used by the other clients. Validate the argument before invoking the soft extraction path.
  • Files reviewed: 29/32 changed files
  • Comments generated: 3
  • Review effort level: Balanced

Comment thread src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Decimal compatibility and Protobuf write-back paths contain correctness and silent data-loss issues.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Files not reviewed (3)

  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_type_rules_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_types_pb2.py: Generated file

Suppressed comments (3)

Previously missed (2) — in code that hasn't changed since the last review.

src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:110

  • Only the fully default value (value == metadata == b'') represents an absent protobuf/Avro field. Treating any empty metadata as absent also turns a corrupted value with non-empty payload into CEL null, silently hiding malformed data. Reject that partial state instead.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:359
  • variants.tryAs is soft only for extraction/type mismatches, not for an invalid argument type. Coercing the second argument with str() makes variants.tryAs(v, 1) silently return CEL null, whereas the declared (dyn, string) overload should reject it. Validate the type before dispatch, as tryParseJson, field, and index do.

src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py:105

  • fields_by_camelcase_name does not cover an explicitly configured protobuf json_name; it only indexes the computed camel-case name. A transform map using that valid JSON name is silently treated as an unknown field and dropped, despite this function's contract. Resolve against each field's json_name instead.
    return desc.fields_by_camelcase_name.get(name)
  • Files reviewed: 31/34 changed files
  • Comments generated: 5
  • Review effort level: Balanced

Comment thread src/confluent_kafka/schema_registry/confluent/types/decimal_utils.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Decimal parity and protobuf write-back contain unresolved correctness issues that can alter serialized values.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Files not reviewed (3)

  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_type_rules_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_types_pb2.py: Generated file

Suppressed comments (3)

Previously missed (2) — in code that hasn't changed since the last review.

src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py:319

  • Python's Decimal square-root operation always rounds with HALF_EVEN, so _DIV_CONTEXT's HALF_UP setting is not honored here. Precision-38 tie cases therefore diverge from the Java MathContext(38, HALF_UP) contract even though ordinary cases such as sqrt(144) pass. Compute with guard precision and explicitly apply HALF_UP rounding to the final 38-digit result.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:359
  • variants.tryAs stringifies its type argument, so an invalid non-string argument is treated as a soft conversion miss and returns CEL null. The declared (dyn, string) overload should reject argument-type errors; only a valid string naming an incompatible target type should be soft.

src/confluent_kafka/schema_registry/confluent/types/decimal_utils.py:44

  • This conversion ignores the proto's precision, so CEL sees a different number than the protobuf serde for precision-limited messages. For example, value=125, scale=0, precision=2 becomes 125 here, while protobuf_to_decimal (and Java MathContext) produces 1.3E+2. Apply the same HALF_UP precision context before exposing the value to CEL.
    scale = int(msg.scale)
    if not msg.value:
        return Decimal(0).scaleb(-scale, context=_EXACT_CONTEXT)
    unscaled = int.from_bytes(msg.value, "big", signed=True)
    return Decimal(unscaled).scaleb(-scale, context=_EXACT_CONTEXT)
  • Files reviewed: 31/34 changed files
  • Comments generated: 4
  • Review effort level: Balanced

Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated
Comment thread src/confluent_kafka/schema_registry/confluent/types/decimal_utils.py Outdated
Comment thread src/confluent_kafka/schema_registry/rules/cel/decimal_funcs.py
Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Decimal precision, repeated conditions, malformed Variants, and protobuf scalar write-back can currently produce incorrect results.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Files not reviewed (3)

  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_type_rules_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_types_pb2.py: Generated file

Suppressed comments (2)

Previously missed (2) — in code that hasn't changed since the last review.

src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py:252

  • Integer wrapper fields bypass _integral and call int(value) directly, so a transform that returns 1.9 for an Int32Value silently writes 1; range and boolean checks are bypassed too. Route wrapper values through the same scalar conversion used for ordinary protobuf fields.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:214
  • Catching every exception makes a recognized but malformed Variant look like a valid non-null value. For example, invalid metadata raises VariantError during coercion, but variants.isNull returns false while every other accessor reports the corrupt payload. Only suppress the expected type-mismatch CEL error; let malformed Variant errors propagate.
  • Files reviewed: 31/34 changed files
  • Comments generated: 3
  • Review effort level: Balanced

Comment thread src/confluent_kafka/schema_registry/common/protobuf.py
Comment thread src/confluent_kafka/schema_registry/confluent/type/decimal_utils.py
Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Several conversion paths can silently alter data or mishandle malformed and out-of-range values.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Review details

Files not reviewed (3)

  • src/confluent_kafka/schema_registry/confluent/types/variant_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_type_rules_pb2.py: Generated file
  • tests/schema_registry/data/proto/value_types_pb2.py: Generated file

Suppressed comments (10)

Previously missed (8) — in code that hasn't changed since the last review.

src/confluent_kafka/schema_registry/common/avro.py:36

  • The logical reader constructs Variant before the CEL absent-value handling can run. A record whose metadata and value are both empty therefore raises VariantError during actual fastavro decoding, even though the new CEL path defines that representation as absent. Return None for the all-empty representation so end-to-end Avro decoding matches that contract.
    src/confluent_kafka/schema_registry/common/protobuf.py:856
  • The new rescaling path still emits precision = 0 below, while the other three Decimal writers now emit the unscaled value's digit count. This leaves values produced through decimal_to_protobuf with different wire metadata and disables the reader's precision semantics. Set precision from the final rescaled integer (with zero having precision 1).
    src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py:138
  • Silently skipping a null map value removes that entry from the rebuilt message. Protobuf maps cannot contain null values, and the protobuf JSON path this writer is intended to match rejects them rather than changing the map's contents. Raise a rule error instead, as the repeated-field path already does for null elements.
    src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py:252
  • Wrapper fields bypass the validated scalar path and coerce arbitrary values. In particular, int(1.9) silently writes 1 to an integer wrapper, despite the same writer correctly rejecting 1.9 for a plain integer field. Route the wrapper's value through _scalar so wrappers enforce the same type and range rules.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:110
  • This treats any empty metadata buffer as absence, including a malformed value with non-empty payload bytes. That silently converts corrupted protobuf/Avro data to CEL null. Only the all-empty representation should be absent; reject a one-sided representation.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:298
  • A Variant timestamp is an int64, so valid encoded values can fall outside Python datetime's year range. The addition then raises raw OverflowError, which bypasses callers that wrap CELEvalError and leaks an implementation exception from CEL evaluation. Normalize this to a CEL range error, as the regular timestamp(...) overload already does.
    src/confluent_kafka/schema_registry/rules/cel/variant_funcs.py:359
  • Coercing type_str with str() makes a wrong-typed call such as variants.tryAs(v, 123) return CEL null as though it were a normal conversion mismatch. The declared overload requires a string, so an invalid argument type must remain a rule error rather than becoming a soft failure.
    src/confluent_kafka/schema_registry/rules/cel/variant_path.py:73
  • The new parser's bracket indexes, quoted keys/escapes, Unicode identifiers, integer bound, and malformed-path branches are untested; the added suite only exercises $.nested.x. Add focused tests for successful bracket/quoted navigation and each rejection boundary so this public path grammar cannot regress unnoticed.

src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py:105

  • fields_by_camelcase_name only covers protobuf's computed camelCase alias; it does not reliably resolve an explicit [json_name = "..."] option. A transform using that declared JSON name is therefore treated as an unknown key and silently dropped, contrary to this function's contract. Match against each field's json_name property.
    fd = desc.fields_by_name.get(name)
    if fd is not None:
        return fd
    return desc.fields_by_camelcase_name.get(name)

src/confluent_kafka/schema_registry/common/protobuf.py:429

  • Making repeated Decimal/Timestamp fields CEL leaves also makes CONDITION rules run once per element, but transform(...) returns a list of booleans and the condition branch above only checks whether the list object itself is False. Thus [True, False] passes silently. Aggregate repeated results and fail when any element is false before entering the transform write-back branch.
            if fd.type == FieldDescriptor.TYPE_MESSAGE and is_cel_leaf_message(fd.message_type):
                # The rule saw this field as a single value, so it hands back a decimal or a
                # datetime rather than the message; encode it before writing.
                #
                # A repeated leaf field needs the same treatment per element. The walk applies
  • Files reviewed: 31/34 changed files
  • Comments generated: 1
  • Review effort level: Balanced

Comment thread src/confluent_kafka/schema_registry/rules/cel/protobuf_result_writer.py Outdated
@sonarqube-confluent

Copy link
Copy Markdown

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants