use std::sync::Arc;
use async_trait::async_trait;
use dataflow_rs::engine::error::DataflowError;
use dataflow_rs::engine::functions::AsyncFunctionHandler;
use dataflow_rs::engine::functions::PublishKafkaConfig;
use dataflow_rs::engine::task_context::TaskContext;
use dataflow_rs::engine::task_outcome::TaskOutcome;
use serde_json::Value;
use super::schema::{FieldKind, FieldSchema};
use crate::connector::ConnectorRegistry;
const NAME: &str = "publish_kafka";
pub struct PublishKafkaHandler {
pub registry: Arc<ConnectorRegistry>,
pub producers: Option<Arc<crate::kafka::producer::KafkaProducerCache>>,
}
#[async_trait]
impl AsyncFunctionHandler for PublishKafkaHandler {
type Input = PublishKafkaConfig;
async fn execute(
&self,
ctx: &mut TaskContext<'_>,
input: &PublishKafkaConfig,
) -> dataflow_rs::Result<TaskOutcome> {
let channel = super::extract_channel(ctx.message()).to_string();
super::connector_helpers::guarded_handler(
NAME,
&self.registry,
&input.connector,
&channel,
async move {
let connector =
super::connector_helpers::resolve_connector(&self.registry, &input.connector)
.await?;
let kafka_config = super::connector_helpers::require_kafka_connector(
connector.as_ref(),
&input.connector,
)?;
super::connector_helpers::require_op(
kafka_config.operations.publish,
"publish",
&input.connector,
)?;
let producers = match &self.producers {
Some(p) => p,
None => {
return Err(DataflowError::FunctionExecution {
context: format!(
"Kafka publishing to topic '{}' is not available. \
Enable Kafka in configuration to use publish_kafka.",
input.topic
),
source: None,
});
}
};
if !kafka_config.brokers.is_empty() {
crate::validation::check_broker_endpoints(
&input.connector,
&kafka_config.brokers,
kafka_config.allow_private_urls,
)
.await
.map_err(crate::errors::connector_detail_error)?;
}
let producer = producers
.for_brokers(&kafka_config.brokers)
.await
.map_err(|e| {
DataflowError::function_execution(
format!(
"Failed to create Kafka producer for connector '{}': {e}",
input.connector
),
None,
)
})?;
let key = input.resolve_key(ctx)?;
let value_json: Value = match input.resolve_value(ctx)? {
Some(value) => value,
None => ctx.data().into(),
};
let value = serde_json::to_string(&value_json).map_err(|e| {
DataflowError::function_execution(
format!("Failed to serialize Kafka message value: {e}"),
None,
)
})?;
producer
.send(&input.topic, key.as_deref(), value.as_bytes())
.await
.map_err(|e| {
DataflowError::function_execution(
format!("Kafka publish to '{}' failed: {e}", input.topic),
None,
)
})?;
tracing::debug!(
topic = %input.topic,
"Published message to Kafka"
);
Ok(TaskOutcome::Success)
},
)
.await
}
}
pub(super) const PUBLISH_KAFKA_FIELDS: &[FieldSchema] = &[
FieldSchema {
name: "connector",
description: "Name of the Kafka connector to publish through.",
kind: FieldKind::String,
required: true,
resolvable: false,
alias: None,
},
FieldSchema {
name: "topic",
description: "Target topic name.",
kind: FieldKind::String,
required: true,
resolvable: false,
alias: None,
},
FieldSchema {
name: "key_logic",
description: "JSONLogic expression to derive the message key.",
kind: FieldKind::Any,
required: false,
resolvable: false,
alias: None,
},
FieldSchema {
name: "value_logic",
description: "JSONLogic expression to derive the message value.",
kind: FieldKind::Any,
required: false,
resolvable: false,
alias: None,
},
];