use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum WorkerJobKind {
SparqlSelect {
query: String,
ntriples: String,
},
ParseTurtle {
turtle: String,
},
CountSubjects {
ntriples: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkerJob {
pub id: u64,
pub kind: WorkerJobKind,
}
impl WorkerJob {
pub fn sparql_select(id: u64, query: impl Into<String>, ntriples: impl Into<String>) -> Self {
Self {
id,
kind: WorkerJobKind::SparqlSelect {
query: query.into(),
ntriples: ntriples.into(),
},
}
}
pub fn parse_turtle(id: u64, turtle: impl Into<String>) -> Self {
Self {
id,
kind: WorkerJobKind::ParseTurtle {
turtle: turtle.into(),
},
}
}
pub fn count_subjects(id: u64, ntriples: impl Into<String>) -> Self {
Self {
id,
kind: WorkerJobKind::CountSubjects {
ntriples: ntriples.into(),
},
}
}
pub fn to_json(&self) -> Result<String, String> {
serde_json::to_string(self).map_err(|e| format!("Serialize error: {e}"))
}
pub fn from_json(json: &str) -> Result<Self, String> {
serde_json::from_str(json).map_err(|e| format!("Deserialize error: {e}"))
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum WorkerResultKind {
SparqlRows { rows_json: String },
ParsedTripleCount { count: usize },
SubjectCount { count: usize },
Error { message: String },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkerResult {
pub job_id: u64,
pub result: WorkerResultKind,
pub elapsed_ms: Option<f64>,
}
impl WorkerResult {
pub fn sparql_ok(job_id: u64, rows_json: impl Into<String>, elapsed_ms: Option<f64>) -> Self {
Self {
job_id,
result: WorkerResultKind::SparqlRows {
rows_json: rows_json.into(),
},
elapsed_ms,
}
}
pub fn parsed_ok(job_id: u64, count: usize, elapsed_ms: Option<f64>) -> Self {
Self {
job_id,
result: WorkerResultKind::ParsedTripleCount { count },
elapsed_ms,
}
}
pub fn error(job_id: u64, message: impl Into<String>) -> Self {
Self {
job_id,
result: WorkerResultKind::Error {
message: message.into(),
},
elapsed_ms: None,
}
}
pub fn to_json(&self) -> Result<String, String> {
serde_json::to_string(self).map_err(|e| format!("Serialize error: {e}"))
}
pub fn from_json(json: &str) -> Result<Self, String> {
serde_json::from_str(json).map_err(|e| format!("Deserialize error: {e}"))
}
pub fn is_error(&self) -> bool {
matches!(self.result, WorkerResultKind::Error { .. })
}
}
pub struct WorkerPool {
capacity: usize,
completed: std::sync::Mutex<Vec<WorkerResult>>,
}
impl WorkerPool {
pub fn new(capacity: usize) -> Self {
Self {
capacity: capacity.max(1),
completed: std::sync::Mutex::new(Vec::new()),
}
}
pub fn capacity(&self) -> usize {
self.capacity
}
pub fn execute(&self, job: WorkerJob) -> WorkerResult {
let result = Self::run_job(&job);
let mut guard = self.completed.lock().unwrap_or_else(|p| p.into_inner());
guard.push(result.clone());
result
}
fn run_job(job: &WorkerJob) -> WorkerResult {
match &job.kind {
WorkerJobKind::CountSubjects { ntriples } => {
let count = ntriples
.lines()
.filter(|l| !l.trim().is_empty() && !l.starts_with('#'))
.count();
WorkerResult::parsed_ok(job.id, count, None)
}
WorkerJobKind::ParseTurtle { turtle } => {
let count = turtle
.lines()
.filter(|l| !l.trim().is_empty() && !l.trim().starts_with('#'))
.count();
WorkerResult::parsed_ok(job.id, count, None)
}
WorkerJobKind::SparqlSelect { query, ntriples: _ } => {
WorkerResult::sparql_ok(
job.id,
format!(
r#"[{{"query": {}}}]"#,
serde_json::to_string(query).unwrap_or_default()
),
None,
)
}
}
}
pub fn completed_results(&self) -> Vec<WorkerResult> {
self.completed
.lock()
.unwrap_or_else(|p| p.into_inner())
.clone()
}
pub fn clear_completed(&self) {
self.completed
.lock()
.unwrap_or_else(|p| p.into_inner())
.clear();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_worker_job_sparql_serialization() {
let job = WorkerJob::sparql_select(1, "SELECT * WHERE {?s ?p ?o}", "");
let json = job.to_json().unwrap();
let parsed = WorkerJob::from_json(&json).unwrap();
assert_eq!(parsed.id, 1);
matches!(parsed.kind, WorkerJobKind::SparqlSelect { .. });
}
#[test]
fn test_worker_job_parse_turtle_serialization() {
let job = WorkerJob::parse_turtle(2, "@prefix ex: <http://example.org/> .");
let json = job.to_json().unwrap();
let back = WorkerJob::from_json(&json).unwrap();
assert_eq!(back.id, 2);
}
#[test]
fn test_worker_result_error_detection() {
let r = WorkerResult::error(42, "something went wrong");
assert!(r.is_error());
assert_eq!(r.job_id, 42);
}
#[test]
fn test_worker_result_serialization() {
let r = WorkerResult::parsed_ok(5, 10, Some(3.15));
let json = r.to_json().unwrap();
let back = WorkerResult::from_json(&json).unwrap();
assert_eq!(back.job_id, 5);
assert!(!back.is_error());
}
#[test]
fn test_pool_count_subjects() {
let pool = WorkerPool::new(2);
let job = WorkerJob::count_subjects(
1,
"<http://a> <http://b> <http://c> .\n<http://d> <http://e> <http://f> .\n",
);
let result = pool.execute(job);
assert!(!result.is_error());
assert_eq!(result.job_id, 1);
if let WorkerResultKind::ParsedTripleCount { count } = result.result {
assert_eq!(count, 2);
} else {
panic!("expected ParsedTripleCount");
}
}
#[test]
fn test_pool_capacity() {
let pool = WorkerPool::new(4);
assert_eq!(pool.capacity(), 4);
}
#[test]
fn test_pool_zero_capacity_clamped() {
let pool = WorkerPool::new(0);
assert_eq!(pool.capacity(), 1, "capacity must be at least 1");
}
#[test]
fn test_pool_multiple_jobs() {
let pool = WorkerPool::new(2);
for i in 0..3 {
pool.execute(WorkerJob::count_subjects(i, ""));
}
assert_eq!(pool.completed_results().len(), 3);
pool.clear_completed();
assert_eq!(pool.completed_results().len(), 0);
}
}