asyncapi_rust_codegen/lib.rs
1//! Procedural macro implementation for asyncapi-rust
2//!
3//! This crate provides the procedural macros that power `asyncapi-rust`, enabling
4//! compile-time generation of AsyncAPI 3.0 specifications from Rust code.
5//!
6//! ## Overview
7//!
8//! Two derive macros are provided:
9//!
10//! ### `#[derive(ToAsyncApiMessage)]`
11//!
12//! Generates message metadata and JSON schemas from Rust types (structs or enums).
13//!
14//! - Works with [`serde`](https://serde.rs) for serialization patterns
15//! - Uses [`schemars`](https://docs.rs/schemars) for JSON Schema generation
16//! - Supports `#[asyncapi(...)]` helper attributes for documentation
17//! - Generates methods: `asyncapi_message_names()`, `asyncapi_messages()`, etc.
18//!
19//! **Example:**
20//! ```rust,ignore
21//! use asyncapi_rust::{ToAsyncApiMessage, schemars::JsonSchema};
22//! use serde::{Deserialize, Serialize};
23//!
24//! #[derive(Serialize, Deserialize, JsonSchema, ToAsyncApiMessage)]
25//! #[serde(tag = "type")]
26//! pub enum ChatMessage {
27//! #[serde(rename = "user.join")]
28//! #[asyncapi(
29//! summary = "User joins",
30//! description = "Sent when a user enters a room"
31//! )]
32//! UserJoin { username: String, room: String },
33//!
34//! #[serde(rename = "chat.message")]
35//! #[asyncapi(summary = "Chat message")]
36//! Chat { username: String, room: String, text: String },
37//! }
38//!
39//! // Generated methods available:
40//! let names = ChatMessage::asyncapi_message_names();
41//! let messages = ChatMessage::asyncapi_messages(); // Requires JsonSchema
42//! ```
43//!
44//! ### `#[derive(AsyncApi)]`
45//!
46//! Generates complete AsyncAPI 3.0 specifications with servers, channels, and operations.
47//!
48//! - Requires `title` and `version` attributes
49//! - Supports optional `description` attribute
50//! - Use `#[asyncapi_server(...)]` to define servers
51//! - Use `#[asyncapi_channel(...)]` to define channels
52//! - Use `#[asyncapi_operation(...)]` to define operations
53//! - Can use multiple of each attribute type
54//!
55//! **Example:**
56//! ```rust,ignore
57//! use asyncapi_rust::AsyncApi;
58//!
59//! #[derive(AsyncApi)]
60//! #[asyncapi(
61//! title = "Chat API",
62//! version = "1.0.0",
63//! description = "Real-time chat application"
64//! )]
65//! #[asyncapi_server(
66//! name = "production",
67//! host = "chat.example.com",
68//! protocol = "wss",
69//! description = "Production WebSocket server"
70//! )]
71//! #[asyncapi_channel(
72//! name = "chat",
73//! address = "/ws/chat"
74//! )]
75//! #[asyncapi_operation(
76//! name = "sendMessage",
77//! action = "send",
78//! channel = "chat"
79//! )]
80//! #[asyncapi_operation(
81//! name = "receiveMessage",
82//! action = "receive",
83//! channel = "chat"
84//! )]
85//! struct ChatApi;
86//!
87//! // Generated method:
88//! let spec = ChatApi::asyncapi_spec();
89//! ```
90//!
91//! ## Supported Attributes
92//!
93//! ### `#[asyncapi(...)]` on message types
94//!
95//! Helper attributes for documenting messages (used with `ToAsyncApiMessage`):
96//!
97//! - `summary = "..."` - Short summary of the message
98//! - `description = "..."` - Detailed description
99//! - `title = "..."` - Human-readable title (defaults to message name)
100//! - `content_type = "..."` - Content type (defaults to "application/json")
101//! - `triggers_binary` - Flag for binary messages (sets content_type to "application/octet-stream")
102//!
103//! ### `#[asyncapi(...)]` on API specs
104//!
105//! Required attributes for complete specifications (used with `AsyncApi`):
106//!
107//! - `title = "..."` - API title (required)
108//! - `version = "..."` - API version (required)
109//! - `description = "..."` - API description (optional)
110//!
111//! ### `#[asyncapi_server(...)]`
112//!
113//! Define server connection information:
114//!
115//! - `name = "..."` - Server identifier (required)
116//! - `host = "..."` - Server host/URL (required)
117//! - `protocol = "..."` - Protocol (e.g., "wss", "ws", "grpc") (required)
118//! - `description = "..."` - Server description (optional)
119//!
120//! ### `#[asyncapi_channel(...)]`
121//!
122//! Define communication channels:
123//!
124//! - `name = "..."` - Channel identifier (required)
125//! - `address = "..."` - Channel path/address (optional)
126//!
127//! ### `#[asyncapi_operation(...)]`
128//!
129//! Define send/receive operations:
130//!
131//! - `name = "..."` - Operation identifier (required)
132//! - `action = "send"|"receive"` - Operation type (required)
133//! - `channel = "..."` - Channel reference (required)
134//!
135//! ## Integration with serde
136//!
137//! The macros respect serde attributes for naming and structure:
138//!
139//! - `#[serde(rename = "...")]` - Use custom name in AsyncAPI spec
140//! - `#[serde(tag = "...")]` - Tagged enum with discriminator field
141//! - `#[serde(skip)]` - Exclude fields from schema
142//! - `#[serde(skip_serializing_if = "...")]` - Optional fields
143//!
144//! ## Integration with schemars
145//!
146//! JSON schemas are generated automatically using schemars:
147//!
148//! - Requires `JsonSchema` derive on message types
149//! - Generates complete JSON Schema from Rust type definitions
150//! - Supports nested types, generics, and references
151//! - Schemas include validation rules from type constraints
152//!
153//! ## Generated Code
154//!
155//! The macros generate implementations with these methods:
156//!
157//! **From `ToAsyncApiMessage`:**
158//! - `asyncapi_message_names() -> Vec<&'static str>` - Get all message names
159//! - `asyncapi_message_count() -> usize` - Number of messages
160//! - `asyncapi_tag_field() -> Option<&'static str>` - Serde tag field if present
161//! - `asyncapi_messages() -> Vec<Message>` - Generate messages with schemas
162//!
163//! **From `AsyncApi`:**
164//! - `asyncapi_spec() -> AsyncApiSpec` - Generate complete specification
165//!
166//! ## Implementation Notes
167//!
168//! - All code generation happens at compile time (proc macros)
169//! - Zero runtime cost - generates plain Rust code
170//! - Compile errors if documentation drifts from code
171//! - Type-safe - uses Rust's type system for validation
172
173#![warn(clippy::all)]
174
175use proc_macro::TokenStream;
176use quote::quote;
177use syn::{Data, DeriveInput, parse_macro_input};
178
179mod asyncapi_attrs;
180mod asyncapi_spec_attrs;
181mod serde_attrs;
182
183use asyncapi_attrs::extract_asyncapi_meta;
184use asyncapi_spec_attrs::extract_asyncapi_spec_meta;
185use serde_attrs::{extract_serde_rename, extract_serde_tag};
186
187/// Derive macro for generating AsyncAPI message metadata
188///
189/// # Example
190///
191/// ```rust,ignore
192/// use asyncapi_rust::ToAsyncApiMessage;
193/// use serde::{Deserialize, Serialize};
194///
195/// #[derive(Serialize, Deserialize, ToAsyncApiMessage)]
196/// #[serde(tag = "type")]
197/// pub enum Message {
198/// #[serde(rename = "chat")]
199/// Chat { room: String, text: String },
200/// Echo { id: i64, text: String },
201/// }
202/// ```
203#[proc_macro_derive(ToAsyncApiMessage, attributes(asyncapi))]
204pub fn derive_to_asyncapi_message(input: TokenStream) -> TokenStream {
205 let input = parse_macro_input!(input as DeriveInput);
206 let name = &input.ident;
207
208 // Extract serde tag attribute from enum
209 let tag_field = extract_serde_tag(&input.attrs);
210
211 // Struct to hold message metadata
212 struct MessageMeta {
213 /// Stable message identity used in components.messages and asyncapi_message_names().
214 /// Defaults to the Rust variant/type identifier; overridable via
215 /// `#[asyncapi(message_name = "...")]`.
216 name: String,
217 /// Wire discriminant value from serde rename (used for payload schema lookup).
218 /// May be an empty string when `#[serde(rename = "")]`; defaults to variant ident.
219 discriminant: String,
220 summary: Option<String>,
221 description: Option<String>,
222 title: Option<String>,
223 content_type: Option<String>,
224 triggers_binary: bool,
225 mqtt: Option<crate::asyncapi_attrs::MqttMessageBindingsMeta>,
226 }
227
228 // Parse enum variants or struct
229 let messages = match &input.data {
230 Data::Enum(data_enum) => {
231 let mut message_metas = Vec::new();
232
233 for variant in &data_enum.variants {
234 let variant_name = &variant.ident;
235 let variant_ident_str = variant_name.to_string();
236
237 // Wire discriminant: serde rename if present (even if empty), else variant ident.
238 let discriminant = extract_serde_rename(&variant.attrs)
239 .unwrap_or_else(|| variant_ident_str.clone());
240
241 // Extract asyncapi metadata
242 let asyncapi_meta = match extract_asyncapi_meta(&variant.attrs) {
243 Ok(m) => m,
244 Err(e) => return e.to_compile_error().into(),
245 };
246
247 // Message identity: explicit message_name override, else variant ident.
248 // We deliberately do NOT use the serde rename here — it may be empty,
249 // non-unique across enums, or unsuitable as a code identifier.
250 let message_name = asyncapi_meta
251 .message_name
252 .clone()
253 .unwrap_or_else(|| variant_ident_str.clone());
254
255 message_metas.push(MessageMeta {
256 name: message_name,
257 discriminant,
258 summary: asyncapi_meta.summary,
259 description: asyncapi_meta.description,
260 title: asyncapi_meta.title,
261 content_type: asyncapi_meta.content_type,
262 triggers_binary: asyncapi_meta.triggers_binary,
263 mqtt: asyncapi_meta.mqtt,
264 });
265 }
266
267 message_metas
268 }
269 Data::Struct(_) => {
270 // For structs, extract metadata from the struct itself
271 let asyncapi_meta = match extract_asyncapi_meta(&input.attrs) {
272 Ok(m) => m,
273 Err(e) => return e.to_compile_error().into(),
274 };
275 let struct_name = name.to_string();
276 let message_name = asyncapi_meta
277 .message_name
278 .clone()
279 .unwrap_or_else(|| struct_name.clone());
280
281 vec![MessageMeta {
282 name: message_name,
283 discriminant: struct_name,
284 summary: asyncapi_meta.summary,
285 description: asyncapi_meta.description,
286 title: asyncapi_meta.title,
287 content_type: asyncapi_meta.content_type,
288 triggers_binary: asyncapi_meta.triggers_binary,
289 mqtt: asyncapi_meta.mqtt,
290 }]
291 }
292 Data::Union(_) => {
293 return syn::Error::new_spanned(name, "ToAsyncApiMessage cannot be derived for unions")
294 .to_compile_error()
295 .into();
296 }
297 };
298
299 let message_count = messages.len();
300 let message_literals = messages.iter().map(|m| m.name.as_str());
301
302 // Prepare metadata for message generation
303 let message_names_for_gen = messages.iter().map(|m| m.name.as_str());
304 // Wire discriminant for each variant — used to look up per-variant schemas at runtime.
305 let message_discriminants = messages.iter().map(|m| m.discriminant.as_str());
306 let message_titles = messages.iter().map(|m| {
307 if let Some(ref title) = m.title {
308 quote! { Some(#title.to_string()) }
309 } else {
310 let name = &m.name;
311 quote! { Some(#name.to_string()) }
312 }
313 });
314 let message_summaries = messages.iter().map(|m| {
315 if let Some(ref summary) = m.summary {
316 quote! { Some(#summary.to_string()) }
317 } else {
318 quote! { None }
319 }
320 });
321 let message_descriptions = messages.iter().map(|m| {
322 if let Some(ref desc) = m.description {
323 quote! { Some(#desc.to_string()) }
324 } else {
325 quote! { None }
326 }
327 });
328 let message_content_types = messages.iter().map(|m| {
329 if let Some(ref ct) = m.content_type {
330 quote! { Some(#ct.to_string()) }
331 } else if m.triggers_binary {
332 quote! { Some("application/octet-stream".to_string()) }
333 } else {
334 quote! { Some("application/json".to_string()) }
335 }
336 });
337
338 let message_mqtt_bindings = messages.iter().map(|m| {
339 if let Some(ref mqtt) = m.mqtt {
340 let payload_format_indicator = if let Some(s) = mqtt.payload_format_indicator {
341 quote! { Some(#s) }
342 } else {
343 quote! { None }
344 };
345 let content_type = if let Some(s) = &mqtt.content_type {
346 quote! { Some(#s.to_string()) }
347 } else {
348 quote! { None }
349 };
350 let binding_version = if let Some(s) = &mqtt.binding_version {
351 quote! { Some(#s.to_string()) }
352 } else {
353 quote! { None }
354 };
355 let response_topic = match &mqtt.response_topic {
356 Some(crate::asyncapi_attrs::ResponseTopicMeta::Uri(t)) => {
357 quote! {
358 Some(asyncapi_rust::MqttResponseTopic::Uri(
359 #t.to_string()
360 ))
361 }
362 }
363 Some(crate::asyncapi_attrs::ResponseTopicMeta::Reference(r)) => {
364 quote! {
365 Some(asyncapi_rust::MqttResponseTopic::Schema({
366 let schema = schemars::schema_for!(#r);
367
368 let schema_json = serde_json::to_value(&schema)
369 .expect("Failed to serialize schema");
370
371 serde_json::from_value(schema_json)
372 .expect("Failed to deserialize schema")
373 }))
374 }
375 }
376 _ => quote! { None },
377 };
378
379 let correlation_data = if let Some(c) = &mqtt.correlation_data {
380 quote! {
381 Some({
382 let schema = schemars::schema_for!(#c);
383
384 let schema_json = serde_json::to_value(&schema)
385 .expect("Failed to serialize schema");
386
387 serde_json::from_value::<asyncapi_rust::Schema>(schema_json)
388 .expect("Failed to deserialize schema")
389 })
390 }
391 } else {
392 quote! { None }
393 };
394
395 quote! {
396 Some(
397 asyncapi_rust::MessageBindings {
398 mqtt: Some(asyncapi_rust::MqttMessageBindings {
399 payload_format_indicator: #payload_format_indicator,
400 content_type: #content_type,
401 binding_version: #binding_version,
402 response_topic: #response_topic,
403 correlation_data: #correlation_data
404 })
405 }
406 )
407 }
408 } else {
409 quote! { None }
410 }
411 });
412
413 let tag_info = if let Some(tag) = tag_field {
414 quote! {
415 Some(#tag)
416 }
417 } else {
418 quote! { None }
419 };
420
421 let expanded = quote! {
422 // const _: () scopes the helper so it doesn't leak into the user's namespace
423 const _: () = {
424 /// Rewrites schemars' `#/$defs/X` refs to `#/components/schemas/X` in-place.
425 fn rewrite_defs_refs(value: &mut serde_json::Value) {
426 match value {
427 serde_json::Value::Object(map) => {
428 if let Some(r) = map.get_mut("$ref") {
429 if let Some(s) = r.as_str() {
430 if let Some(name) = s.strip_prefix("#/$defs/") {
431 *r = serde_json::Value::String(
432 format!("#/components/schemas/{}", name)
433 );
434 }
435 }
436 }
437 for v in map.values_mut() {
438 rewrite_defs_refs(v);
439 }
440 }
441 serde_json::Value::Array(arr) => {
442 for v in arr.iter_mut() {
443 rewrite_defs_refs(v);
444 }
445 }
446 _ => {}
447 }
448 }
449
450 impl #name {
451 /// Get AsyncAPI message names for this type
452 pub fn asyncapi_message_names() -> Vec<&'static str> {
453 vec![#(#message_literals),*]
454 }
455
456 /// Get the number of messages in this type
457 pub fn asyncapi_message_count() -> usize {
458 #message_count
459 }
460
461 /// Get the serde tag field name if this is a tagged enum
462 pub fn asyncapi_tag_field() -> Option<&'static str> {
463 #tag_info
464 }
465
466 /// Return shared schema definitions for this type, keyed by name.
467 ///
468 /// These are the `$defs` that schemars generates for sub-types referenced
469 /// by this type's variants. The `AsyncApi` derive collects them into
470 /// `components.schemas` so message payloads can reference them via
471 /// `#/components/schemas/X` instead of embedding them inline.
472 pub fn asyncapi_schemas() -> asyncapi_rust::indexmap::IndexMap<String, asyncapi_rust::Schema>
473 where
474 Self: schemars::JsonSchema,
475 {
476 use schemars::schema_for;
477 let schema = schema_for!(Self);
478 let schema_json = serde_json::to_value(&schema)
479 .expect("Failed to serialize schema");
480
481 let mut result = asyncapi_rust::indexmap::IndexMap::new();
482 if let Some(defs) = schema_json.get("$defs").and_then(|v| v.as_object()) {
483 for (name, def_schema) in defs {
484 let mut def = def_schema.clone();
485 rewrite_defs_refs(&mut def);
486 if let Ok(s) = serde_json::from_value::<asyncapi_rust::Schema>(def) {
487 result.insert(name.clone(), s);
488 }
489 }
490 }
491 result
492 }
493
494 /// Generate AsyncAPI Message objects with JSON schemas.
495 ///
496 /// For internally-tagged enums each message carries only its own variant
497 /// schema. `$ref`s within payloads point to `#/components/schemas/X`;
498 /// the corresponding definitions are available via `asyncapi_schemas()`.
499 pub fn asyncapi_messages() -> Vec<asyncapi_rust::Message>
500 where
501 Self: schemars::JsonSchema,
502 {
503 use schemars::schema_for;
504
505 let schema = schema_for!(Self);
506 let schema_json = serde_json::to_value(&schema)
507 .expect("Failed to serialize schema");
508
509 // Build a discriminant→schema map using the actual serde tag field name.
510 let tag_field = Self::asyncapi_tag_field();
511 let mut variant_schemas: asyncapi_rust::indexmap::IndexMap<String, serde_json::Value> =
512 asyncapi_rust::indexmap::IndexMap::new();
513 if let Some(tag) = tag_field {
514 if let Some(variants) = schema_json.get("oneOf").and_then(|v| v.as_array()) {
515 for variant in variants {
516 let discriminant = variant
517 .get("properties")
518 .and_then(|props| props.get(tag))
519 .and_then(|tag_prop| {
520 tag_prop.get("const").or_else(|| {
521 tag_prop
522 .get("enum")
523 .and_then(|e| e.as_array())
524 .and_then(|a| a.first())
525 })
526 })
527 .and_then(|v| v.as_str())
528 .map(|s| s.to_string());
529
530 if let Some(name) = discriminant {
531 let mut variant_schema = variant.clone();
532 // Drop $defs — they live in components.schemas, not the payload.
533 if let Some(obj) = variant_schema.as_object_mut() {
534 obj.remove("$defs");
535 }
536 rewrite_defs_refs(&mut variant_schema);
537 variant_schemas.insert(name, variant_schema);
538 }
539 }
540 }
541 }
542
543 // Metadata arrays are baked in at compile time; schemas resolved at runtime.
544 let names: &[&str] = &[#(#message_names_for_gen),*];
545 // Discriminants are the serde rename values — used to look up per-variant
546 // schemas. Separate from names so empty renames and cross-enum collisions
547 // don't affect message identity.
548 let discriminants: &[&str] = &[#(#message_discriminants),*];
549 let titles: &[Option<String>] = &[#(#message_titles),*];
550 let summaries: &[Option<String>] = &[#(#message_summaries),*];
551 let descriptions: &[Option<String>] = &[#(#message_descriptions),*];
552 let content_types: &[Option<String>] = &[#(#message_content_types),*];
553 let bindings: &[Option<asyncapi_rust::MessageBindings>] = &[#(#message_mqtt_bindings),*];
554
555 let mut messages = Vec::with_capacity(names.len());
556 for i in 0..names.len() {
557 let msg_name = names[i];
558 let discriminant = discriminants[i];
559 let payload = if let Some(v) = variant_schemas.get(discriminant) {
560 serde_json::from_value(v.clone()).ok()
561 } else {
562 // Structs, untagged enums, or variants not in the map:
563 // remove $defs and rewrite refs in the full schema.
564 let mut fallback = schema_json.clone();
565 if let Some(obj) = fallback.as_object_mut() {
566 obj.remove("$defs");
567 }
568 rewrite_defs_refs(&mut fallback);
569 serde_json::from_value(fallback).ok()
570 };
571 messages.push(asyncapi_rust::Message {
572 name: Some(msg_name.to_string()),
573 title: titles[i].clone(),
574 summary: summaries[i].clone(),
575 description: descriptions[i].clone(),
576 content_type: content_types[i].clone(),
577 payload,
578 bindings: bindings[i].clone()
579 });
580 }
581 messages
582 }
583 }
584 };
585 };
586
587 TokenStream::from(expanded)
588}
589
590/// Derive macro for generating complete AsyncAPI specification
591///
592/// # Example
593///
594/// ```rust,ignore
595/// use asyncapi_rust::AsyncApi;
596///
597/// #[derive(AsyncApi)]
598/// #[asyncapi(
599/// title = "Chat API",
600/// version = "1.0.0",
601/// description = "A real-time chat API"
602/// )]
603/// struct ChatApi;
604/// ```
605#[proc_macro_derive(
606 AsyncApi,
607 attributes(
608 asyncapi,
609 asyncapi_server,
610 asyncapi_channel,
611 asyncapi_operation,
612 asyncapi_messages
613 )
614)]
615pub fn derive_asyncapi(input: TokenStream) -> TokenStream {
616 let input = parse_macro_input!(input as DeriveInput);
617 let name = &input.ident;
618
619 // Extract asyncapi spec metadata
620 let spec_meta = match extract_asyncapi_spec_meta(&input.attrs) {
621 Ok(m) => m,
622 Err(e) => return e.to_compile_error().into(),
623 };
624
625 // Validate required fields
626 let title = match spec_meta.title {
627 Some(t) => t,
628 None => {
629 return syn::Error::new_spanned(
630 name,
631 "AsyncApi requires a title attribute: #[asyncapi(title = \"...\")]",
632 )
633 .to_compile_error()
634 .into();
635 }
636 };
637
638 let version = match spec_meta.version {
639 Some(v) => v,
640 None => {
641 return syn::Error::new_spanned(
642 name,
643 "AsyncApi requires a version attribute: #[asyncapi(version = \"...\")]",
644 )
645 .to_compile_error()
646 .into();
647 }
648 };
649
650 let description = if let Some(desc) = spec_meta.description {
651 quote! { Some(#desc.to_string()) }
652 } else {
653 quote! { None }
654 };
655
656 // Generate servers
657 let servers_code = if spec_meta.servers.is_empty() {
658 quote! { None }
659 } else {
660 let server_entries = spec_meta.servers.iter().map(|server| {
661 let name = &server.name;
662 let host = &server.host;
663 let protocol = &server.protocol;
664 let pathname = if let Some(p) = &server.pathname {
665 quote! { Some(#p.to_string()) }
666 } else {
667 quote! { None }
668 };
669 let desc = if let Some(d) = &server.description {
670 quote! { Some(#d.to_string()) }
671 } else {
672 quote! { None }
673 };
674
675 // Generate server variables
676 let variables = if server.variables.is_empty() {
677 quote! { None }
678 } else {
679 let var_entries = server.variables.iter().map(|var| {
680 let var_name = &var.name;
681 let var_desc = if let Some(d) = &var.description {
682 quote! { Some(#d.to_string()) }
683 } else {
684 quote! { None }
685 };
686 let var_default = if let Some(d) = &var.default {
687 quote! { Some(#d.to_string()) }
688 } else {
689 quote! { None }
690 };
691 let var_enum = if var.enum_values.is_empty() {
692 quote! { None }
693 } else {
694 let enum_vals = &var.enum_values;
695 quote! { Some(vec![#(#enum_vals.to_string()),*]) }
696 };
697 let var_examples = if var.examples.is_empty() {
698 quote! { None }
699 } else {
700 let examples = &var.examples;
701 quote! { Some(vec![#(#examples.to_string()),*]) }
702 };
703
704 quote! {
705 server_variables.insert(
706 #var_name.to_string(),
707 asyncapi_rust::ServerVariable {
708 description: #var_desc,
709 default: #var_default,
710 enum_values: #var_enum,
711 examples: #var_examples,
712 }
713 );
714 }
715 });
716
717 quote! {
718 {
719 let mut server_variables = asyncapi_rust::indexmap::IndexMap::new();
720 #(#var_entries)*
721 Some(server_variables)
722 }
723 }
724 };
725
726 let bindings = if let Some(mqtt) = &server.mqtt {
727 let client_id = if let Some(s) = &mqtt.client_id {
728 quote! { Some(#s.to_string()) }
729 } else {
730 quote! { None }
731 };
732 let clean_session = if let Some(s) = mqtt.clean_session {
733 quote! { Some(#s) }
734 } else {
735 quote! { None }
736 };
737 let last_will = if let Some(lw) = &mqtt.last_will {
738 let topic = &lw.topic;
739 let qos = lw.qos;
740 let retain = lw.retain;
741 let message = &lw.message;
742 quote! {
743 Some(
744 asyncapi_rust::MqttLastWill {
745 topic: #topic.to_string(),
746 qos: #qos,
747 retain: #retain,
748 message: #message.to_string()
749 }
750 )
751 }
752 } else {
753 quote! { None }
754 };
755 let binding_version = if let Some(s) = &mqtt.binding_version {
756 quote! { Some(#s.to_string()) }
757 } else {
758 quote! { None }
759 };
760
761 let session_expiry_interval = match &mqtt.session_expiry_interval {
762 Some(crate::asyncapi_spec_attrs::MqttBindingNumValueMeta::Value(t)) => {
763 quote! {
764 Some(asyncapi_rust::MqttBindingNumValue::Value(
765 #t
766 ))
767 }
768 }
769 Some(crate::asyncapi_spec_attrs::MqttBindingNumValueMeta::Reference(r)) => {
770 quote! {
771 Some(asyncapi_rust::MqttBindingNumValue::Schema({
772 let schema = schemars::schema_for!(#r);
773
774 let schema_json = serde_json::to_value(&schema)
775 .expect("Failed to serialize schema");
776
777 serde_json::from_value(schema_json)
778 .expect("Failed to deserialize schema")
779 }))
780 }
781 }
782 _ => quote! { None },
783 };
784
785 let keep_alive = if let Some(k) = mqtt.keep_alive {
786 quote! { Some(#k) }
787 } else {
788 quote! { None }
789 };
790
791 let max_packet_size = match &mqtt.maximum_packet_size {
792 Some(crate::asyncapi_spec_attrs::MqttBindingNumValueMeta::Value(t)) => {
793 quote! {
794 Some(asyncapi_rust::MqttBindingNumValue::Value(
795 #t
796 ))
797 }
798 }
799 Some(crate::asyncapi_spec_attrs::MqttBindingNumValueMeta::Reference(r)) => {
800 quote! {
801 Some(asyncapi_rust::MqttBindingNumValue::Schema({
802 let schema = schemars::schema_for!(#r);
803
804 let schema_json = serde_json::to_value(&schema)
805 .expect("Failed to serialize schema");
806
807 serde_json::from_value(schema_json)
808 .expect("Failed to deserialize schema")
809 }))
810 }
811 }
812 _ => quote! { None },
813 };
814
815 quote! {
816 Some(asyncapi_rust::ServerBindings {
817 mqtt: Some(asyncapi_rust::MqttServerBindings {
818 client_id: #client_id,
819 clean_session: #clean_session,
820 session_expiry_interval: #session_expiry_interval,
821 binding_version: #binding_version,
822 last_will: #last_will,
823 keep_alive: #keep_alive,
824 max_packet_size: #max_packet_size
825 })
826 })
827 }
828 } else {
829 quote! { None }
830 };
831
832 quote! {
833 servers.insert(
834 #name.to_string(),
835 asyncapi_rust::Server {
836 host: #host.to_string(),
837 protocol: #protocol.to_string(),
838 pathname: #pathname,
839 description: #desc,
840 variables: #variables,
841 bindings: #bindings
842 }
843 );
844 }
845 });
846
847 quote! {
848 {
849 let mut servers = asyncapi_rust::indexmap::IndexMap::new();
850 #(#server_entries)*
851 Some(servers)
852 }
853 }
854 };
855
856 // Generate channels
857 let channels_code = if spec_meta.channels.is_empty() {
858 quote! { None }
859 } else {
860 let channel_entries = spec_meta.channels.iter().map(|channel| {
861 let name = &channel.name;
862 let address = if let Some(addr) = &channel.address {
863 quote! { Some(#addr.to_string()) }
864 } else {
865 quote! { None }
866 };
867
868 // Generate channel parameters
869 let parameters = if channel.parameters.is_empty() {
870 quote! { None }
871 } else {
872 let param_entries = channel.parameters.iter().map(|param| {
873 let param_name = ¶m.name;
874 let param_desc = if let Some(d) = ¶m.description {
875 quote! { Some(#d.to_string()) }
876 } else {
877 quote! { None }
878 };
879 let param_default = if let Some(d) = ¶m.default {
880 quote! { Some(#d.to_string()) }
881 } else {
882 quote! { None }
883 };
884 let param_enum = if param.enum_values.is_empty() {
885 quote! { None }
886 } else {
887 let vals = ¶m.enum_values;
888 quote! { Some(vec![#(#vals.to_string()),*]) }
889 };
890 let param_examples = if param.examples.is_empty() {
891 quote! { None }
892 } else {
893 let vals = ¶m.examples;
894 quote! { Some(vec![#(#vals.to_string()),*]) }
895 };
896 let param_location = if let Some(l) = ¶m.location {
897 quote! { Some(#l.to_string()) }
898 } else {
899 quote! { None }
900 };
901
902 quote! {
903 channel_parameters.insert(
904 #param_name.to_string(),
905 asyncapi_rust::Parameter {
906 description: #param_desc,
907 default: #param_default,
908 enum_values: #param_enum,
909 examples: #param_examples,
910 location: #param_location,
911 }
912 );
913 }
914 });
915
916 quote! {
917 {
918 let mut channel_parameters = asyncapi_rust::indexmap::IndexMap::new();
919 #(#param_entries)*
920 Some(channel_parameters)
921 }
922 }
923 };
924
925 quote! {
926 channels.insert(
927 #name.to_string(),
928 asyncapi_rust::Channel {
929 address: #address,
930 messages: None,
931 parameters: #parameters,
932 }
933 );
934 }
935 });
936
937 quote! {
938 {
939 let mut channels = asyncapi_rust::indexmap::IndexMap::new();
940 #(#channel_entries)*
941 Some(channels)
942 }
943 }
944 };
945
946 // Generate operations
947 let operations_code = if spec_meta.operations.is_empty() {
948 quote! { None }
949 } else {
950 let operation_entries = spec_meta.operations.iter().map(|operation| {
951 let name = &operation.name;
952 let channel_ref = &operation.channel;
953 let action = &operation.action;
954
955 // Convert action string to OperationAction enum
956 let action_enum = if action == "send" {
957 quote! { asyncapi_rust::OperationAction::Send }
958 } else if action == "receive" {
959 quote! { asyncapi_rust::OperationAction::Receive }
960 } else {
961 return syn::Error::new_spanned(
962 name,
963 format!("Invalid action '{}', must be 'send' or 'receive'", action),
964 )
965 .to_compile_error();
966 };
967
968 let bindings = if let Some(mqtt) = &operation.mqtt {
969 let qos = if let Some(s) = mqtt.qos {
970 quote! { Some(#s) }
971 } else {
972 quote! { None }
973 };
974
975 let retain = if let Some(b) = mqtt.retain {
976 quote! { Some(#b) }
977 } else {
978 quote! { None }
979 };
980
981 let binding_version = if let Some(s) = &mqtt.binding_version {
982 quote! { Some(#s.to_string()) }
983 } else {
984 quote! { None }
985 };
986
987 let message_expiry = match &mqtt.message_expiry_interval {
988 Some(crate::asyncapi_spec_attrs::MqttBindingNumValueMeta::Value(t)) => {
989 quote! {
990 Some(asyncapi_rust::MqttBindingNumValue::Value(
991 #t
992 ))
993 }
994 }
995 Some(crate::asyncapi_spec_attrs::MqttBindingNumValueMeta::Reference(r)) => {
996 quote! {
997 Some(asyncapi_rust::MqttBindingNumValue::Schema({
998 let schema = schemars::schema_for!(#r);
999
1000 let schema_json = serde_json::to_value(&schema)
1001 .expect("Failed to serialize schema");
1002
1003 serde_json::from_value(schema_json)
1004 .expect("Failed to deserialize schema")
1005 }))
1006 }
1007 }
1008 _ => quote! { None },
1009 };
1010
1011 quote! {
1012 Some(asyncapi_rust::OperationBindings {
1013 mqtt: Some(asyncapi_rust::MqttOperationBindings {
1014 qos: #qos,
1015 retain: #retain,
1016 message_expiry_interval: #message_expiry,
1017 binding_version: #binding_version,
1018 })
1019 })
1020 }
1021 } else {
1022 quote! { None }
1023 };
1024
1025 quote! {
1026 operations.insert(
1027 #name.to_string(),
1028 asyncapi_rust::Operation {
1029 action: #action_enum,
1030 channel: asyncapi_rust::ChannelRef {
1031 reference: format!("#/channels/{}", #channel_ref),
1032 },
1033 messages: None,
1034 bindings: #bindings
1035 }
1036 );
1037 }
1038 });
1039
1040 quote! {
1041 {
1042 let mut operations = asyncapi_rust::indexmap::IndexMap::new();
1043 #(#operation_entries)*
1044 Some(operations)
1045 }
1046 }
1047 };
1048
1049 // Generate components with messages and hoisted shared schemas
1050 let components_code = if spec_meta.message_types.is_empty() {
1051 quote! { None }
1052 } else {
1053 let type_calls = spec_meta.message_types.iter().map(|type_name| {
1054 quote! {
1055 for msg in #type_name::asyncapi_messages() {
1056 if let Some(ref name) = msg.name {
1057 if messages.contains_key(name.as_str()) {
1058 panic!(
1059 "asyncapi-rust: message name collision for '{}' from {}. \
1060 Use #[asyncapi(message_name = \"...\")] on one variant to disambiguate.",
1061 name,
1062 stringify!(#type_name)
1063 );
1064 }
1065 messages.insert(name.clone(), msg.clone());
1066 }
1067 }
1068 // Hoist shared $defs into components.schemas (first writer wins on name collision)
1069 for (name, schema) in #type_name::asyncapi_schemas() {
1070 schemas.entry(name).or_insert(schema);
1071 }
1072 }
1073 });
1074
1075 quote! {
1076 {
1077 let mut messages = asyncapi_rust::indexmap::IndexMap::new();
1078 let mut schemas = asyncapi_rust::indexmap::IndexMap::new();
1079 #(#type_calls)*
1080 Some(asyncapi_rust::Components {
1081 messages: if messages.is_empty() { None } else { Some(messages) },
1082 schemas: if schemas.is_empty() { None } else { Some(schemas) },
1083 })
1084 }
1085 }
1086 };
1087
1088 let expanded = quote! {
1089 impl #name {
1090 /// Generate the AsyncAPI specification
1091 ///
1092 /// Returns an AsyncApiSpec with Info, Servers, Channels, and Operations
1093 /// sections populated from attributes.
1094 pub fn asyncapi_spec() -> asyncapi_rust::AsyncApiSpec {
1095 asyncapi_rust::AsyncApiSpec {
1096 asyncapi: "3.0.0".to_string(),
1097 info: asyncapi_rust::Info {
1098 title: #title.to_string(),
1099 version: #version.to_string(),
1100 description: #description,
1101 },
1102 servers: #servers_code,
1103 channels: #channels_code,
1104 operations: #operations_code,
1105 components: #components_code,
1106 }
1107 }
1108 }
1109 };
1110
1111 TokenStream::from(expanded)
1112}
1113
1114#[cfg(test)]
1115mod tests {
1116 #[test]
1117 fn test_placeholder() {
1118 // Macro expansion tests will go here
1119 }
1120}