Skip to content

Commit 69ada72

Browse files
committed
test(pulsar): §5 conformance runner + vendored suite +
drift-guard CI
1 parent fe9130f commit 69ada72

11 files changed

Lines changed: 469 additions & 1 deletion

File tree

‎.github/workflows/ci.yml‎

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,10 +30,23 @@ jobs:
3030
- name: Test
3131
run: mvn -B --no-transfer-progress verify
3232

33+
conformance:
34+
name: Conformance suite in sync
35+
runs-on: ubuntu-latest
36+
steps:
37+
- uses: actions/checkout@v5
38+
- name: Verify vendored conformance matches the canonical suite
39+
run: |
40+
git clone --depth 1 https://github.com/BabelQueue/conformance.git "$RUNNER_TEMP/conformance"
41+
diff -ru "$RUNNER_TEMP/conformance/manifest.json" "src/test/resources/conformance/manifest.json"
42+
diff -ru "$RUNNER_TEMP/conformance/fixtures" "src/test/resources/conformance/fixtures"
43+
diff -ru "$RUNNER_TEMP/conformance/schema" "src/test/resources/conformance/schema"
44+
echo "Vendored conformance is in sync with the canonical suite."
45+
3346
ci-green:
3447
name: CI green
3548
runs-on: ubuntu-latest
36-
needs: [test]
49+
needs: [test, conformance]
3750
if: ${{ always() }}
3851
steps:
3952
- name: Fail if any required job did not pass

‎pom.xml‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,13 @@
8282
<version>${mockito.version}</version>
8383
<scope>test</scope>
8484
</dependency>
85+
<!-- Parses the vendored conformance manifest in PulsarConformanceTest. -->
86+
<dependency>
87+
<groupId>org.json</groupId>
88+
<artifactId>json</artifactId>
89+
<version>20240303</version>
90+
<scope>test</scope>
91+
</dependency>
8592
</dependencies>
8693

8794
<build>
Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,83 @@
1+
package com.babelqueue.pulsar;
2+
3+
import static org.junit.jupiter.api.Assertions.assertEquals;
4+
import static org.mockito.ArgumentMatchers.any;
5+
import static org.mockito.ArgumentMatchers.anyInt;
6+
import static org.mockito.Mockito.mock;
7+
import static org.mockito.Mockito.when;
8+
9+
import com.babelqueue.Envelope;
10+
import com.babelqueue.EnvelopeCodec;
11+
import java.io.InputStream;
12+
import java.nio.charset.StandardCharsets;
13+
import java.util.Map;
14+
import org.apache.pulsar.client.api.Consumer;
15+
import org.apache.pulsar.client.api.Message;
16+
import org.json.JSONArray;
17+
import org.json.JSONObject;
18+
import org.junit.jupiter.api.Test;
19+
20+
/**
21+
* Apache Pulsar binding conformance against the vendored canonical suite's {@code pulsar}
22+
* block: the §5 property projection (bq-* string→string) and the
23+
* {@code attempts = max(body, RedeliveryCount)} reconciliation (no −1; the redelivery count is
24+
* 0-based). The Pulsar consumer/message are mocked with Mockito — no Pulsar, no network.
25+
*/
26+
class PulsarConformanceTest {
27+
28+
private static final String URN = "urn:babel:orders:created";
29+
30+
private static String resource(String path) throws Exception {
31+
try (InputStream in = PulsarConformanceTest.class.getResourceAsStream("/conformance/" + path)) {
32+
if (in == null) {
33+
throw new IllegalStateException("vendored conformance resource missing: " + path);
34+
}
35+
return new String(in.readAllBytes(), StandardCharsets.UTF_8);
36+
}
37+
}
38+
39+
private static JSONObject pulsarBlock() throws Exception {
40+
return new JSONObject(resource("manifest.json")).getJSONObject("pulsar");
41+
}
42+
43+
@Test
44+
void propertyProjectionMatchesGolden() throws Exception {
45+
JSONObject projection = pulsarBlock().getJSONObject("property_projection");
46+
Envelope envelope = EnvelopeCodec.decode(resource(projection.getString("envelope_file")));
47+
Map<String, String> got = PulsarProperties.of(envelope);
48+
49+
JSONObject want = projection.getJSONObject("properties");
50+
assertEquals(want.keySet(), got.keySet());
51+
for (String key : want.keySet()) {
52+
assertEquals(want.getString(key), got.get(key), key);
53+
}
54+
}
55+
56+
@Test
57+
@SuppressWarnings("unchecked")
58+
void attemptsReconciliationMatchesGolden() throws Exception {
59+
JSONArray cases = pulsarBlock().getJSONObject("attempts_reconciliation").getJSONArray("cases");
60+
for (int i = 0; i < cases.length(); i++) {
61+
JSONObject testCase = cases.getJSONObject(i);
62+
Envelope base = EnvelopeCodec.make(URN, Map.of("x", 1), "orders", null);
63+
Envelope bumped = new Envelope(
64+
base.job(), base.traceId(), base.data(), base.meta(),
65+
testCase.getInt("body_attempts"), base.deadLetter());
66+
67+
Message<byte[]> msg = mock(Message.class);
68+
when(msg.getValue()).thenReturn(EnvelopeCodec.encode(bumped).getBytes(StandardCharsets.UTF_8));
69+
when(msg.getRedeliveryCount()).thenReturn(testCase.getInt("redelivery_count"));
70+
71+
Consumer<byte[]> consumer = mock(Consumer.class);
72+
when(consumer.receive(anyInt(), any())).thenReturn(msg);
73+
74+
int[] seen = {-1};
75+
PulsarConsumer.builder(consumer)
76+
.handler(URN, (env, message) -> seen[0] = env.attempts())
77+
.build()
78+
.poll();
79+
80+
assertEquals(testCase.getInt("expected_attempts"), seen[0], testCase.getString("name"));
81+
}
82+
}
83+
}
Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
{
2+
"job": "urn:babel:orders:created",
3+
"trace_id": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b",
4+
"data": {
5+
"order_id": 1042
6+
},
7+
"meta": {
8+
"id": "f1e2d3c4-b5a6-4789-90ab-cdef01234567",
9+
"queue": "orders",
10+
"lang": "php",
11+
"schema_version": 1,
12+
"created_at": 1749132727000
13+
},
14+
"attempts": 3,
15+
"dead_letter": {
16+
"reason": "failed",
17+
"error": "Payment gateway timeout",
18+
"exception": "App\\Exceptions\\GatewayTimeout",
19+
"failed_at": 1749132730000,
20+
"original_queue": "orders",
21+
"attempts": 3,
22+
"lang": "php"
23+
}
24+
}
Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
{
2+
"trace_id": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b",
3+
"data": {
4+
"order_id": 1042
5+
},
6+
"meta": {
7+
"id": "f1e2d3c4-b5a6-4789-90ab-cdef01234567",
8+
"queue": "orders",
9+
"lang": "php",
10+
"schema_version": 1,
11+
"created_at": 1749132727000
12+
},
13+
"attempts": 0
14+
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
{
2+
"job": "urn:babel:orders:created",
3+
"trace_id": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b",
4+
"data": {
5+
"order_id": 1042
6+
},
7+
"meta": {
8+
"id": "f1e2d3c4-b5a6-4789-90ab-cdef01234567",
9+
"queue": "orders",
10+
"lang": "php",
11+
"schema_version": 2,
12+
"created_at": 1749132727000
13+
},
14+
"attempts": 0
15+
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
{
2+
"job": "urn:babel:orders:created",
3+
"trace_id": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b",
4+
"data": {
5+
"order_id": 1042
6+
},
7+
"meta": {
8+
"id": "f1e2d3c4-b5a6-4789-90ab-cdef01234567",
9+
"queue": "orders",
10+
"lang": "php",
11+
"schema_version": 1,
12+
"created_at": 1749132727000
13+
},
14+
"attempts": 0
15+
}
Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,20 @@
1+
{
2+
"job": "urn:babel:catalog:item.indexed",
3+
"trace_id": "3f7a1d2e-9b4c-4a8d-bc1e-0f5a6b7c8d90",
4+
"data": {
5+
"title": "Café — naïve ☕",
6+
"qty": 7,
7+
"price_cents": 1299,
8+
"ratio": 0.5,
9+
"active": true,
10+
"note": null
11+
},
12+
"meta": {
13+
"id": "b2c3d4e5-f607-4890-a1b2-c3d4e5f60718",
14+
"queue": "catalog",
15+
"lang": "python",
16+
"schema_version": 1,
17+
"created_at": 1749132727000
18+
},
19+
"attempts": 2
20+
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
{
2+
"urn": "urn:babel:orders:created",
3+
"trace_id": "9c1e0b44-7a2d-4e6f-8a10-2b3c4d5e6f70",
4+
"data": {
5+
"order_id": 1042
6+
},
7+
"meta": {
8+
"id": "a1b2c3d4-e5f6-4789-90ab-cdef01234567",
9+
"queue": "orders",
10+
"lang": "go",
11+
"schema_version": 1,
12+
"created_at": 1749132727000
13+
},
14+
"attempts": 0
15+
}
Lines changed: 152 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,152 @@
1+
{
2+
"schema_version": 1,
3+
"description": "Cross-SDK conformance cases. Every BabelQueue SDK core must satisfy these against the canonical wire envelope. Per-message fields (meta.id, trace_id, meta.created_at) are intrinsically unique and are NOT asserted by value.",
4+
"cases": [
5+
{
6+
"name": "order-created",
7+
"file": "fixtures/order-created.json",
8+
"valid": true,
9+
"description": "A normal produced envelope.",
10+
"expect": {
11+
"urn": "urn:babel:orders:created",
12+
"data": { "order_id": 1042 },
13+
"attempts": 0,
14+
"lang": "php",
15+
"schema_version": 1
16+
}
17+
},
18+
{
19+
"name": "urn-alias",
20+
"file": "fixtures/urn-alias.json",
21+
"valid": true,
22+
"description": "Consumers MUST accept 'urn' as an inbound alias for 'job'.",
23+
"expect": {
24+
"urn": "urn:babel:orders:created",
25+
"data": { "order_id": 1042 },
26+
"attempts": 0,
27+
"lang": "go",
28+
"schema_version": 1
29+
}
30+
},
31+
{
32+
"name": "dead-lettered",
33+
"file": "fixtures/dead-lettered.json",
34+
"valid": true,
35+
"description": "A dead-lettered message: original preserved + additive dead_letter block.",
36+
"expect": {
37+
"urn": "urn:babel:orders:created",
38+
"data": { "order_id": 1042 },
39+
"attempts": 3,
40+
"lang": "php",
41+
"schema_version": 1,
42+
"dead_letter": { "reason": "failed", "original_queue": "orders" }
43+
}
44+
},
45+
{
46+
"name": "unicode-and-numbers",
47+
"file": "fixtures/unicode-and-numbers.json",
48+
"valid": true,
49+
"description": "UTF-8 strings, integers, an exact float, boolean and null round-trip identically.",
50+
"expect": {
51+
"urn": "urn:babel:catalog:item.indexed",
52+
"data": { "title": "Café — naïve ☕", "qty": 7, "price_cents": 1299, "ratio": 0.5, "active": true, "note": null },
53+
"attempts": 2,
54+
"lang": "python",
55+
"schema_version": 1
56+
}
57+
},
58+
{
59+
"name": "invalid-unknown-schema-version",
60+
"file": "fixtures/invalid-unknown-schema-version.json",
61+
"valid": false,
62+
"reason": "meta.schema_version is not a version this SDK supports"
63+
},
64+
{
65+
"name": "invalid-missing-urn",
66+
"file": "fixtures/invalid-missing-urn.json",
67+
"valid": false,
68+
"reason": "no 'job' or 'urn' — the message has no identity"
69+
}
70+
],
71+
"sqs": {
72+
"description": "Amazon SQS binding conformance (broker-bindings.md §3). Every SDK that ships an SQS transport must satisfy these. The envelope body stays byte-identical (the 'cases' above); these lock the native projection + reconciliation the binding adds. Per-message values reuse fixtures/order-created.json so the expected attributes are deterministic.",
73+
"attribute_projection": {
74+
"description": "On produce, the transport MUST project these native MessageAttributes from the envelope (a redundant, routable view of the body; ids/strings are DataType String, counters are DataType Number). Applies to every SQS-producing SDK.",
75+
"envelope_file": "fixtures/order-created.json",
76+
"message_attributes": {
77+
"bq-job": { "DataType": "String", "StringValue": "urn:babel:orders:created" },
78+
"bq-trace-id": { "DataType": "String", "StringValue": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b" },
79+
"bq-message-id": { "DataType": "String", "StringValue": "f1e2d3c4-b5a6-4789-90ab-cdef01234567" },
80+
"bq-schema-version": { "DataType": "Number", "StringValue": "1" },
81+
"bq-source-lang": { "DataType": "String", "StringValue": "php" },
82+
"bq-created-at": { "DataType": "Number", "StringValue": "1749132727000" }
83+
}
84+
},
85+
"attempts_reconciliation": {
86+
"description": "On consume, attempts = max(body.attempts, ApproximateReceiveCount - 1): a first delivery reads 0, an absent/garbage count is ignored, a runtime-incremented count is never lowered. Applies to SDKs that reconcile the envelope body on consume (the framework-less/runtime transports). A drop-in driver that surfaces the broker's native delivery count instead (e.g. Laravel's SqsJob.attempts() = ApproximateReceiveCount) is exempt — it documents that divergence.",
87+
"cases": [
88+
{ "name": "first-delivery", "body_attempts": 0, "approximate_receive_count": "1", "expected_attempts": 0 },
89+
{ "name": "third-delivery", "body_attempts": 0, "approximate_receive_count": "3", "expected_attempts": 2 },
90+
{ "name": "native-exceeds-body", "body_attempts": 2, "approximate_receive_count": "5", "expected_attempts": 4 },
91+
{ "name": "never-lower-runtime", "body_attempts": 5, "approximate_receive_count": "1", "expected_attempts": 5 },
92+
{ "name": "garbage-count-ignored", "body_attempts": 4, "approximate_receive_count": "not-a-number", "expected_attempts": 4 },
93+
{ "name": "absent-count", "body_attempts": 3, "approximate_receive_count": null, "expected_attempts": 3 }
94+
]
95+
}
96+
},
97+
"asb": {
98+
"description": "Azure Service Bus binding conformance (broker-bindings.md §4). Every SDK that ships an ASB transport must satisfy these. The envelope body stays byte-identical (the 'cases' above); these lock the native projection + reconciliation the binding adds. Per-message values reuse fixtures/order-created.json so the expected projection is deterministic.",
99+
"property_projection": {
100+
"description": "On produce, the transport MUST project these native Service Bus message fields from the envelope: Subject = job (the URN), CorrelationId = trace_id, MessageId = meta.id, ContentType = application/json, plus the bq- ApplicationProperties as native AMQP-typed values (numbers stay numbers, not strings). Applies to every ASB-producing SDK.",
101+
"envelope_file": "fixtures/order-created.json",
102+
"message": {
103+
"subject": "urn:babel:orders:created",
104+
"correlation_id": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b",
105+
"message_id": "f1e2d3c4-b5a6-4789-90ab-cdef01234567",
106+
"content_type": "application/json"
107+
},
108+
"application_properties": {
109+
"bq-schema-version": 1,
110+
"bq-source-lang": "php",
111+
"bq-created-at": 1749132727000
112+
}
113+
},
114+
"attempts_reconciliation": {
115+
"description": "On consume, attempts = max(body.attempts, DeliveryCount - 1): a first delivery (DeliveryCount 1) reads 0, a runtime-incremented body count is never lowered, and DeliveryCount <= 1 leaves the body's own count untouched (the runtime retries by republishing with attempts+1). DeliveryCount is the native 1-based ASB redelivery counter; the rule is identical across the native-consumer SDKs (.NET/Java/Node) and the Transport+App SDKs (Python/Go).",
116+
"cases": [
117+
{ "name": "first-delivery", "body_attempts": 0, "delivery_count": 1, "expected_attempts": 0 },
118+
{ "name": "third-delivery", "body_attempts": 0, "delivery_count": 3, "expected_attempts": 2 },
119+
{ "name": "native-exceeds-body", "body_attempts": 2, "delivery_count": 5, "expected_attempts": 4 },
120+
{ "name": "never-lower-runtime", "body_attempts": 5, "delivery_count": 2, "expected_attempts": 5 },
121+
{ "name": "first-delivery-keeps-body", "body_attempts": 4, "delivery_count": 1, "expected_attempts": 4 },
122+
{ "name": "zero-count-keeps-body", "body_attempts": 3, "delivery_count": 0, "expected_attempts": 3 }
123+
]
124+
}
125+
},
126+
"pulsar": {
127+
"description": "Apache Pulsar binding conformance (broker-bindings.md §5). Every SDK that ships a Pulsar transport must satisfy these. The envelope body stays byte-identical (the 'cases' above); these lock the native projection + reconciliation the binding adds. Per-message values reuse fixtures/order-created.json so the expected projection is deterministic.",
128+
"property_projection": {
129+
"description": "On produce, the transport MUST project these native Pulsar message properties from the envelope, all string->string (Pulsar properties are string-typed, so numbers are stringified): bq-job = job (the URN), bq-trace-id = trace_id, bq-message-id = meta.id, bq-schema-version = str(meta.schema_version), bq-source-lang = meta.lang, bq-attempts = str(attempts). The payload is the byte-identical envelope; the native publishTime mirrors meta.created_at (broker-set, body authoritative). Applies to every Pulsar-producing SDK.",
130+
"envelope_file": "fixtures/order-created.json",
131+
"properties": {
132+
"bq-job": "urn:babel:orders:created",
133+
"bq-trace-id": "7b3f9c2a-e41d-4f88-9b2a-1c0d5e6f7a8b",
134+
"bq-message-id": "f1e2d3c4-b5a6-4789-90ab-cdef01234567",
135+
"bq-schema-version": "1",
136+
"bq-source-lang": "php",
137+
"bq-attempts": "0"
138+
}
139+
},
140+
"attempts_reconciliation": {
141+
"description": "On consume, attempts = max(body.attempts, RedeliveryCount): Pulsar's RedeliveryCount is 0-based (0 on first delivery) so it maps directly with NO -1, a runtime-incremented body count is never lowered, and RedeliveryCount 0 leaves the body's own count untouched (the runtime retries by republishing with attempts+1, which resets the broker's redelivery count to 0). The rule is identical across the native-consumer SDKs (.NET/Java/Node) and the Transport+App SDKs (Python/Go).",
142+
"cases": [
143+
{ "name": "first-delivery", "body_attempts": 0, "redelivery_count": 0, "expected_attempts": 0 },
144+
{ "name": "third-delivery", "body_attempts": 0, "redelivery_count": 2, "expected_attempts": 2 },
145+
{ "name": "native-exceeds-body", "body_attempts": 2, "redelivery_count": 5, "expected_attempts": 5 },
146+
{ "name": "never-lower-runtime", "body_attempts": 5, "redelivery_count": 1, "expected_attempts": 5 },
147+
{ "name": "first-delivery-keeps-body", "body_attempts": 4, "redelivery_count": 0, "expected_attempts": 4 },
148+
{ "name": "native-equals-body", "body_attempts": 3, "redelivery_count": 3, "expected_attempts": 3 }
149+
]
150+
}
151+
}
152+
}

0 commit comments

Comments
 (0)