use chrono::{DateTime, Utc};
use crate::ir::{
ComparisonOp, ConflictStrategy, LogicalFilter, LogicalPagination, LogicalProjection,
LogicalRead, LogicalRecord, LogicalValue,
};
use crate::runtime::native_catalog::{NativeModel, native_model};
use super::config::{LOCK_MSG, STATUS_HELD};
pub(crate) fn logical_string(value: impl Into<String>) -> LogicalValue {
LogicalValue::String(value.into())
}
pub(crate) fn lock_filter(
tenant_id: &str,
lock_name: Option<&str>,
status: Option<&str>,
) -> LogicalFilter {
let mut filters = vec![LogicalFilter::Comparison {
field: "tenant_id".to_string(),
op: ComparisonOp::Eq,
value: logical_string(tenant_id),
}];
if let Some(lock_name) = lock_name {
filters.push(LogicalFilter::Comparison {
field: "lock_name".to_string(),
op: ComparisonOp::Eq,
value: logical_string(lock_name),
});
}
if let Some(status) = status {
filters.push(LogicalFilter::Comparison {
field: "status".to_string(),
op: ComparisonOp::Eq,
value: logical_string(status),
});
}
LogicalFilter::And(filters)
}
pub(crate) fn lock_read_by_name(tenant_id: &str, lock_name: &str) -> LogicalRead {
LogicalRead {
message_type: LOCK_MSG.to_string(),
filter: Some(lock_filter(tenant_id, Some(lock_name), None)),
projection: Some(LogicalProjection::fields([
"lock_id".to_string(),
"owner_id".to_string(),
"fencing_token".to_string(),
"status".to_string(),
])),
sort: Vec::new(),
include: Vec::new(),
pagination: Some(LogicalPagination::limit(1)),
}
}
pub(crate) fn held_locks_read(tenant_id: &str, limit: u32, now: DateTime<Utc>) -> LogicalRead {
LogicalRead {
message_type: LOCK_MSG.to_string(),
filter: Some(LogicalFilter::And(vec![
lock_filter(tenant_id, None, Some(STATUS_HELD)),
LogicalFilter::Comparison {
field: "expires_at".to_string(),
op: ComparisonOp::Gt,
value: LogicalValue::Timestamp(now),
},
])),
projection: Some(LogicalProjection::fields(["lock_id".to_string()])),
sort: Vec::new(),
include: Vec::new(),
pagination: Some(LogicalPagination::limit(limit)),
}
}
pub(crate) fn lock_inventory_read(
tenant_id: &str,
lock_name: Option<&str>,
status: Option<&str>,
offset: u64,
limit: u32,
) -> LogicalRead {
let mut pagination = LogicalPagination::limit(limit);
pagination.offset = (offset > 0).then_some(offset);
LogicalRead {
message_type: LOCK_MSG.to_string(),
filter: Some(lock_filter(tenant_id, lock_name, status)),
projection: Some(LogicalProjection::fields([
"lock_id".to_string(),
"tenant_id".to_string(),
"lock_name".to_string(),
"owner_id".to_string(),
"fencing_token".to_string(),
"lease_ttl_seconds".to_string(),
"status".to_string(),
"acquired_at".to_string(),
"expires_at".to_string(),
"metadata_json".to_string(),
])),
sort: vec![crate::ir::LogicalSort {
field: "acquired_at".to_string(),
direction: crate::ir::SortDirection::Desc,
nulls: crate::ir::NullOrder::Default,
}],
include: Vec::new(),
pagination: Some(pagination),
}
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn lock_record(
lock_id: &str,
tenant_id: &str,
lock_name: &str,
owner_id: &str,
fencing_token: i64,
ttl_seconds: i64,
status: &str,
acquired_at: DateTime<Utc>,
expires_at: DateTime<Utc>,
metadata_json: &str,
) -> LogicalRecord {
let mut record = LogicalRecord::new();
record.insert("lock_id".to_string(), logical_string(lock_id));
record.insert("tenant_id".to_string(), logical_string(tenant_id));
record.insert("lock_name".to_string(), logical_string(lock_name));
record.insert("owner_id".to_string(), logical_string(owner_id));
record.insert(
"fencing_token".to_string(),
LogicalValue::Int(fencing_token),
);
record.insert(
"lease_ttl_seconds".to_string(),
LogicalValue::Int(ttl_seconds),
);
record.insert("status".to_string(), logical_string(status));
record.insert(
"acquired_at".to_string(),
LogicalValue::Timestamp(acquired_at),
);
record.insert(
"expires_at".to_string(),
LogicalValue::Timestamp(expires_at),
);
record.insert("metadata_json".to_string(), logical_string(metadata_json));
record
}
pub(crate) fn lock_conflict() -> ConflictStrategy {
ConflictStrategy::update(vec![
"owner_id".to_string(),
"fencing_token".to_string(),
"lease_ttl_seconds".to_string(),
"status".to_string(),
"acquired_at".to_string(),
"expires_at".to_string(),
"metadata_json".to_string(),
])
}
pub(crate) fn lock_model() -> NativeModel {
native_model(
LOCK_MSG,
&[
"lock_id",
"tenant_id",
"lock_name",
"owner_id",
"fencing_token",
"status",
"expires_at",
],
)
}