use std::sync::Arc;
use async_trait::async_trait;
use dataflow_rs::engine::functions::AsyncFunctionHandler;
use dataflow_rs::engine::task_context::TaskContext;
use dataflow_rs::engine::task_outcome::TaskOutcome;
use serde_json::Value;
use super::connector_helpers::{
ConnectorCall, apply_output, require_cache_connector, require_op, resolve_required_str,
to_connect_error, to_exec_error,
};
use super::schema::{FieldKind, FieldSchema};
use crate::connector::ConnectorRegistry;
use crate::connector::cache_backend::{CachePool, CachePurpose};
const NAME: &str = "cache_read";
pub struct CacheReadHandler {
pub cache_pool: Arc<CachePool>,
pub registry: Arc<ConnectorRegistry>,
}
#[async_trait]
impl AsyncFunctionHandler for CacheReadHandler {
type Input = Value;
async fn execute(
&self,
ctx: &mut TaskContext<'_>,
input: &Value,
) -> dataflow_rs::Result<TaskOutcome> {
let call = ConnectorCall::begin(NAME, input, ctx)?;
let key = resolve_required_str(input, "key", call.name, ctx)?;
call.run(&self.registry, async {
let connector_config = call.resolve(&self.registry, None).await?;
let cache_config = require_cache_connector(&connector_config, call.connector)?;
require_op(cache_config.operations.read, "read", call.connector)?;
let backend = self
.cache_pool
.get_backend(CachePurpose::Workflow, call.connector, cache_config)
.await
.map_err(to_connect_error)?;
let value = backend.get(&key).await.map_err(to_exec_error)?;
let result = match value {
Some(v) => serde_json::from_str::<Value>(&v).unwrap_or(Value::String(v)),
None => Value::Null,
};
apply_output(ctx, call.output, result);
Ok(TaskOutcome::Success)
})
.await
}
}
pub(super) const CACHE_READ_FIELDS: &[FieldSchema] = &[
FieldSchema {
name: "connector",
description: "Name of the cache connector to read from.",
kind: FieldKind::String,
required: true,
resolvable: false,
alias: None,
},
FieldSchema {
name: "key",
description: "Cache key to look up. Accepts {\"var\": \"path\"} to read the value from the message.",
kind: FieldKind::String,
required: true,
resolvable: true,
alias: None,
},
FieldSchema {
name: "output",
description: "Dotted path in the message where the result is stored. Defaults to \"data\".",
kind: FieldKind::String,
required: false,
resolvable: false,
alias: None,
},
];