From 08fe6ec6ec6d86b12538340208bfaca744725be1 Mon Sep 17 00:00:00 2001 From: Ankit Khullar Date: Sun, 4 Oct 2026 14:53:05 +0530 Subject: [PATCH] [feat][client] Expose message schema version Add Message.getSchemaVersion(), which returns the version of the schema the producer used to serialize the message, or null when the message carries no schema version. It wraps the C++ client's hasSchemaVersion() and getLongSchemaVersion(). Consumers of Avro or JSON topics need the writer schema version to decode messages produced with an older, compatible schema. Master Issue: #362 Co-Authored-By: Claude Opus 5.5 (1M context) --- index.d.ts | 5 ++++ src/Message.cc | 14 +++++++++ src/Message.h | 1 + tests/end_to_end.test.js | 64 ++++++++++++++++++++++++++++++++++++++++ tstest.ts | 1 + 5 files changed, 85 insertions(+) diff --git a/index.d.ts b/index.d.ts index 160fd52..c6311a2 100644 --- a/index.d.ts +++ b/index.d.ts @@ -194,6 +194,11 @@ export class Message { getProducerName(): string; isReplicated(): boolean; getReplicatedFrom(): string; + /** + * The version of the schema the producer used to serialize this message, + * or `null` if the message carries no schema version. + */ + getSchemaVersion(): number | null; getEncryptionContext(): EncryptionContext | null; } diff --git a/src/Message.cc b/src/Message.cc index 55fa7e3..07dfe04 100644 --- a/src/Message.cc +++ b/src/Message.cc @@ -65,6 +65,7 @@ Napi::Object Message::Init(Napi::Env env, Napi::Object exports) { InstanceMethod("getProducerName", &Message::GetProducerName), InstanceMethod("isReplicated", &Message::IsReplicated), InstanceMethod("getReplicatedFrom", &Message::GetReplicatedFrom), + InstanceMethod("getSchemaVersion", &Message::GetSchemaVersion), InstanceMethod("getEncryptionContext", &Message::GetEncryptionContext)}); constructor = Napi::Persistent(func); @@ -197,6 +198,19 @@ Napi::Value Message::GetReplicatedFrom(const Napi::CallbackInfo &info) { return Napi::String::New(env, replicatedFrom); } +Napi::Value Message::GetSchemaVersion(const Napi::CallbackInfo &info) { + Napi::Env env = info.Env(); + if (!ValidateCMessage(env)) { + return env.Null(); + } + + const pulsar::Message &message = this->cMessage.get()->message; + if (!message.hasSchemaVersion()) { + return env.Null(); + } + return Napi::Number::New(env, message.getLongSchemaVersion()); +} + Napi::Value Message::GetEncryptionContext(const Napi::CallbackInfo &info) { Napi::Env env = info.Env(); if (!ValidateCMessage(env)) { diff --git a/src/Message.h b/src/Message.h index c84126c..fadba41 100644 --- a/src/Message.h +++ b/src/Message.h @@ -48,6 +48,7 @@ class Message : public Napi::ObjectWrap { Napi::Value GetProducerName(const Napi::CallbackInfo &info); Napi::Value IsReplicated(const Napi::CallbackInfo &info); Napi::Value GetReplicatedFrom(const Napi::CallbackInfo &info); + Napi::Value GetSchemaVersion(const Napi::CallbackInfo &info); Napi::Value GetRedeliveryCount(const Napi::CallbackInfo &info); Napi::Value GetEncryptionContext(const Napi::CallbackInfo &info); bool ValidateCMessage(Napi::Env env); diff --git a/tests/end_to_end.test.js b/tests/end_to_end.test.js index d60be82..565d4fd 100644 --- a/tests/end_to_end.test.js +++ b/tests/end_to_end.test.js @@ -178,6 +178,70 @@ const Pulsar = require('../index'); await client.close(); }); + test('getSchemaVersion', async () => { + const client = new Pulsar.Client({ + serviceUrl: 'pulsar://localhost:6650', + operationTimeoutSeconds: 30, + }); + + const schemaV1 = { + schemaType: 'Json', + schema: JSON.stringify({ + type: 'record', + name: 'User', + fields: [{ name: 'name', type: 'string' }], + }), + }; + // Compatible evolution of schemaV1: adds a field with a default. + const schemaV2 = { + schemaType: 'Json', + schema: JSON.stringify({ + type: 'record', + name: 'User', + fields: [ + { name: 'name', type: 'string' }, + { name: 'age', type: 'int', default: 0 }, + ], + }), + }; + + const topic = `persistent://public/default/schema-version-${Date.now()}`; + const consumer = await client.subscribe({ + topic, + subscription: 'sub1', + schema: schemaV1, + }); + + const producerV1 = await client.createProducer({ topic, schema: schemaV1 }); + await producerV1.send({ data: Buffer.from(JSON.stringify({ name: 'alice' })) }); + const msgV1 = await consumer.receive(); + consumer.acknowledge(msgV1); + + const producerV2 = await client.createProducer({ topic, schema: schemaV2 }); + await producerV2.send({ data: Buffer.from(JSON.stringify({ name: 'bob', age: 30 })) }); + const msgV2 = await consumer.receive(); + consumer.acknowledge(msgV2); + + expect(msgV1.getSchemaVersion()).toBe(0); + expect(msgV2.getSchemaVersion()).toBe(1); + + const bytesTopic = `persistent://public/default/schema-version-bytes-${Date.now()}`; + const bytesConsumer = await client.subscribe({ topic: bytesTopic, subscription: 'sub1' }); + const bytesProducer = await client.createProducer({ topic: bytesTopic }); + await bytesProducer.send({ data: Buffer.from('no-schema') }); + const bytesMsg = await bytesConsumer.receive(); + bytesConsumer.acknowledge(bytesMsg); + + expect(bytesMsg.getSchemaVersion()).toBeNull(); + + await producerV1.close(); + await producerV2.close(); + await consumer.close(); + await bytesProducer.close(); + await bytesConsumer.close(); + await client.close(); + }); + test('Produce/Consume Listener', async () => { const client = new Pulsar.Client({ serviceUrl: 'pulsar://localhost:6650', diff --git a/tstest.ts b/tstest.ts index c58f660..e8acda5 100644 --- a/tstest.ts +++ b/tstest.ts @@ -309,6 +309,7 @@ import Pulsar = require('./index'); const redeliveryCount: number = message1.getRedeliveryCount(); const isReplicated: boolean = message1.isReplicated(); const replicatedFrom: string = message1.getReplicatedFrom(); + const schemaVersion: number | null = message1.getSchemaVersion(); const partitionKey: string = message1.getPartitionKey(); const message3: Pulsar.Message = await reader1.readNext();