#[cfg(test)]
mod tests;
pub mod http_validator;
pub mod sql_classifier;
pub mod storage;
use std::{collections::HashSet, sync::Arc};
use fraiseql_core::security::SecurityContext;
use fraiseql_error::Result;
use crate::{
HostContext,
types::{EventPayload, LogEntry, LogLevel},
};
#[derive(Debug, Clone)]
pub struct HostContextConfig {
pub allowed_domains: Vec<String>,
pub allowed_env_vars: HashSet<String>,
pub max_http_response_bytes: usize,
pub http_connect_timeout_ms: u64,
pub http_read_timeout_ms: u64,
pub max_storage_upload_bytes: usize,
}
impl Default for HostContextConfig {
fn default() -> Self {
Self {
allowed_domains: vec![],
allowed_env_vars: HashSet::new(),
max_http_response_bytes: 10 * 1024 * 1024, http_connect_timeout_ms: 5000,
http_read_timeout_ms: 30000,
max_storage_upload_bytes: 100 * 1024 * 1024, }
}
}
pub trait QueryExecutor: Send + Sync {
fn execute_query(
&self,
query: &str,
variables: Option<&serde_json::Value>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<serde_json::Value>> + Send + '_>>;
}
pub struct LiveHostContext {
event_payload: EventPayload,
config: HostContextConfig,
logs: Arc<std::sync::Mutex<Vec<LogEntry>>>,
query_executor: Option<Arc<dyn QueryExecutor>>,
http_client: Option<Arc<reqwest::Client>>,
pub storage_backend: Option<Arc<dyn storage::StorageBackend>>,
security_context: Option<SecurityContext>,
sender_resolver: Option<Arc<dyn crate::outbound::SenderIdentityResolver>>,
email_transport: Option<Arc<dyn crate::outbound::EmailTransport>>,
idempotency_token: Option<String>,
source_cursor: Option<SourceCursorBinding>,
}
struct SourceCursorBinding {
source_name: String,
store: fraiseql_observers::PostgresSourceCursorStore,
}
impl LiveHostContext {
#[must_use]
pub fn new(event_payload: EventPayload, config: HostContextConfig) -> Self {
Self {
event_payload,
config,
logs: Arc::new(std::sync::Mutex::new(Vec::new())),
query_executor: None,
http_client: None,
storage_backend: None,
security_context: None,
sender_resolver: None,
email_transport: None,
idempotency_token: None,
source_cursor: None,
}
}
#[must_use]
pub fn with_executor(mut self, executor: Arc<dyn QueryExecutor>) -> Self {
self.query_executor = Some(executor);
self
}
#[must_use]
pub fn with_http_client(
event_payload: EventPayload,
config: HostContextConfig,
http_client: Arc<reqwest::Client>,
) -> Self {
Self {
event_payload,
config,
logs: Arc::new(std::sync::Mutex::new(Vec::new())),
query_executor: None,
http_client: Some(http_client),
storage_backend: None,
security_context: None,
sender_resolver: None,
email_transport: None,
idempotency_token: None,
source_cursor: None,
}
}
#[must_use]
pub fn with_security_context(mut self, security_context: SecurityContext) -> Self {
self.security_context = Some(security_context);
self
}
#[must_use]
pub fn captured_logs(&self) -> Vec<LogEntry> {
self.logs.lock().expect("log mutex poisoned").clone()
}
#[must_use]
pub fn with_email(
mut self,
sender_resolver: Arc<dyn crate::outbound::SenderIdentityResolver>,
email_transport: Arc<dyn crate::outbound::EmailTransport>,
) -> Self {
self.sender_resolver = Some(sender_resolver);
self.email_transport = Some(email_transport);
self
}
#[must_use]
pub fn with_idempotency_token(mut self, token: impl Into<String>) -> Self {
self.idempotency_token = Some(token.into());
self
}
#[must_use]
pub fn with_source_cursor(
mut self,
source_name: impl Into<String>,
store: fraiseql_observers::PostgresSourceCursorStore,
) -> Self {
self.source_cursor = Some(SourceCursorBinding {
source_name: source_name.into(),
store,
});
self
}
}
#[allow(unknown_lints, clippy::unused_async_trait_impl)]
impl HostContext for LiveHostContext {
async fn query(
&self,
graphql: &str,
variables: serde_json::Value,
) -> Result<serde_json::Value> {
let executor = self.query_executor.as_ref().ok_or_else(|| {
fraiseql_error::FraiseQLError::Unsupported {
message: "query executor not configured".to_string(),
}
})?;
executor.execute_query(graphql, Some(&variables)).await
}
async fn sql_query(
&self,
sql: &str,
_params: &[serde_json::Value],
) -> Result<Vec<serde_json::Value>> {
let classification = sql_classifier::classify_sql(sql)?;
match classification {
sql_classifier::SqlClassification::ReadOnly => {
Err(fraiseql_error::FraiseQLError::Unsupported {
message: "sql_query host function is not implemented: the statement \
was accepted as read-only but no execution backend is wired"
.to_string(),
})
},
sql_classifier::SqlClassification::Rejected(reason) => {
Err(fraiseql_error::FraiseQLError::Authorization {
message: format!("SQL query not allowed: {}", reason),
action: Some("execute_sql_query".to_string()),
resource: None,
})
},
}
}
async fn http_request(
&self,
method: &str,
url: &str,
headers: &[(String, String)],
body: Option<&[u8]>,
) -> Result<crate::host::HttpResponse> {
crate::host::outbound_http::perform(
&self.config,
self.http_client.as_ref(),
method,
url,
headers,
body,
)
.await
}
async fn storage_get(&self, bucket: &str, key: &str) -> Result<Vec<u8>> {
let backend = self.storage_backend.as_ref().ok_or_else(|| {
fraiseql_error::FraiseQLError::Unsupported {
message: "storage backend not configured".to_string(),
}
})?;
backend.get(bucket, key).await
}
async fn storage_put(
&self,
bucket: &str,
key: &str,
body: &[u8],
content_type: &str,
) -> Result<()> {
if body.len() > self.config.max_storage_upload_bytes {
return Err(fraiseql_error::FraiseQLError::Validation {
message: format!(
"upload size {} exceeds limit {}",
body.len(),
self.config.max_storage_upload_bytes
),
path: None,
});
}
let backend = self.storage_backend.as_ref().ok_or_else(|| {
fraiseql_error::FraiseQLError::Unsupported {
message: "storage backend not configured".to_string(),
}
})?;
backend.put(bucket, key, body, content_type).await
}
async fn send_email(
&self,
request: &crate::outbound::SendEmailRequest,
) -> Result<crate::outbound::SendEmailResponse> {
let resolver = self.sender_resolver.as_ref().ok_or_else(|| {
fraiseql_error::FraiseQLError::Unsupported {
message: "send_email is not configured: no sender-identity resolver is wired \
(configure a mailbox with an SMTP send half)"
.to_string(),
}
})?;
let transport = self.email_transport.as_ref().ok_or_else(|| {
fraiseql_error::FraiseQLError::Unsupported {
message: "send_email is not configured: no email transport is wired (configure a \
mailbox with an SMTP send half)"
.to_string(),
}
})?;
let auth = self.auth_context()?;
let sender = resolver.resolve_sender(&auth).await.map_err(|error| {
if error.retryable {
fraiseql_error::FraiseQLError::ServiceUnavailable {
message: error.message,
retry_after: None,
}
} else {
fraiseql_error::FraiseQLError::Authorization {
message: error.message,
action: Some("send_email".to_string()),
resource: None,
}
}
})?;
let context = crate::outbound::SendContext {
send_id: self.idempotency_token.as_deref(),
tenant: self
.security_context
.as_ref()
.and_then(|ctx| ctx.tenant_id.as_ref())
.map(|tenant| tenant.as_str()),
};
transport.send(&sender, request, context).await
}
fn auth_context(&self) -> Result<serde_json::Value> {
let context = self.security_context.as_ref().ok_or_else(|| {
fraiseql_error::FraiseQLError::Unsupported {
message: "no security context is wired on this host: the dispatch path must \
inject the caller's authenticated context or the function's run_as \
identity via with_security_context"
.to_string(),
}
})?;
Ok(crate::host::auth_context_json(context))
}
fn env_var(&self, name: &str) -> Result<Option<String>> {
if self.config.allowed_env_vars.contains(name) {
Ok(std::env::var(name).ok())
} else {
Err(fraiseql_error::FraiseQLError::Authorization {
message: format!(
"environment variable `{name}` is not on this function's allowlist \
(grant it via FRAISEQL_FUNCTIONS_ALLOWED_ENV_VARS or [sources] \
allowed_env_vars)"
),
action: Some("env_var".to_string()),
resource: Some(name.to_string()),
})
}
}
fn event_payload(&self) -> &EventPayload {
&self.event_payload
}
fn log(&self, level: LogLevel, message: &str) {
let entry = LogEntry {
level,
message: message.to_string(),
timestamp: chrono::Utc::now(),
};
self.logs.lock().expect("log mutex poisoned").push(entry);
}
fn idempotency_token(&self) -> Option<String> {
self.idempotency_token.clone()
}
async fn cursor(&self) -> Result<Option<serde_json::Value>> {
use fraiseql_observers::SourceCursorStore;
match &self.source_cursor {
Some(binding) => {
let snapshot =
binding.store.load(&binding.source_name).await.map_err(|error| {
fraiseql_error::FraiseQLError::database(error.to_string())
})?;
Ok(snapshot.value)
},
None => Ok(None),
}
}
async fn advance_cursor(&self, value: serde_json::Value) -> Result<()> {
use fraiseql_observers::SourceCursorStore;
let Some(binding) = &self.source_cursor else {
return Err(fraiseql_error::FraiseQLError::validation(
"advance_cursor: this function is not a scheduled source (no cursor binding)",
));
};
let snapshot = binding
.store
.load(&binding.source_name)
.await
.map_err(|error| fraiseql_error::FraiseQLError::database(error.to_string()))?;
let applied = binding
.store
.advance(&binding.source_name, &snapshot, value)
.await
.map_err(|error| fraiseql_error::FraiseQLError::database(error.to_string()))?;
if applied {
tracing::debug!(source = %binding.source_name, "source cursor advanced");
Ok(())
} else {
Err(fraiseql_error::FraiseQLError::database(
"advance_cursor: the cursor moved concurrently (lost the compare-and-swap)",
))
}
}
}