use crate::contracts::availability::{SchemaAvailability, SchemaFreshness};
use crate::contracts::knowledge_validity::item_validity_for;
use crate::contracts::view::ContractSchemaState;
use saya_store::{KnowledgeItem, KnowledgeItemStore, SqliteStateStore};
use saya_types::{
ClaimId, ClaimPayload, ClaimStatus, DatabaseObjectRef, KnowledgeState, ProfileIdentity,
};
use crate::contracts::ContractOpError;
pub(crate) const QUEUE_LIMIT_CAP: usize = 200;
pub(crate) const QUEUE_DEFAULT_LIMIT: usize = 50;
#[derive(Debug)]
pub(crate) struct QueuedClaim {
pub payload: Option<ClaimPayload>,
pub id: ClaimId,
pub status: ClaimStatus,
pub object: DatabaseObjectRef,
}
impl QueuedClaim {
pub(crate) fn from_item(item: &KnowledgeItem) -> Option<Self> {
let id = ClaimId::parse(&item.id).ok()?;
Some(Self {
payload: Some(item.value.clone()),
id,
status: status_for_queue(item.state),
object: item.object.clone(),
})
}
}
fn status_for_queue(state: KnowledgeState) -> ClaimStatus {
match state {
KnowledgeState::Active => ClaimStatus::Confirmed,
KnowledgeState::Pending => ClaimStatus::Candidate,
KnowledgeState::Dismissed => ClaimStatus::Rejected,
_ => ClaimStatus::Rejected,
}
}
#[derive(Debug)]
pub(crate) struct QueuedCandidate {
pub claim: QueuedClaim,
pub schema_state: ContractSchemaState,
pub evidence_count: usize,
}
pub(crate) async fn review_queue(
store: &SqliteStateStore,
profiles: &[ProfileIdentity],
schemas: &[(ProfileIdentity, SchemaAvailability)],
limit: usize,
) -> Result<Vec<QueuedCandidate>, ContractOpError> {
let limit = limit.min(QUEUE_LIMIT_CAP);
let mut entries: Vec<(KnowledgeItem, ContractSchemaState)> = Vec::new();
for profile in profiles {
for item in store.knowledge_for_profile(profile).await? {
if item.state != KnowledgeState::Pending {
continue;
}
let live = live_schema(schemas, item.object.profile());
let schema_state = item_validity_for(&item, live, SchemaFreshness::Unbounded).into();
entries.push((item, schema_state));
}
}
entries.sort_by(|(a, _), (b, _)| {
a.created_unix_ms
.cmp(&b.created_unix_ms)
.then_with(|| a.slot.as_str().cmp(&b.slot.as_str()))
.then_with(|| a.id.cmp(&b.id))
});
entries.truncate(limit);
let queued = entries
.into_iter()
.filter_map(|(item, schema_state)| {
let claim = QueuedClaim::from_item(&item)?;
Some(QueuedCandidate {
claim,
schema_state,
evidence_count: 0,
})
})
.collect();
Ok(queued)
}
fn live_schema<'s>(
schemas: &'s [(ProfileIdentity, SchemaAvailability)],
profile: &ProfileIdentity,
) -> &'s SchemaAvailability {
schemas
.iter()
.find(|(p, _)| p == profile)
.map(|(_, avail)| avail)
.unwrap_or(&SchemaAvailability::Missing)
}