use std::sync::Arc;
use async_trait::async_trait;
use operon::options::{OperonOptions, PsqlMetaStorageOptions, PsqlStorageOptions};
use operon::{Operon, OperonService, PsqlMetaStorage, define_operon};
use serde::{Deserialize, Serialize};
type Input = String;
#[derive(Debug, Clone, Serialize, Deserialize)]
struct Intermediate(String);
#[derive(Debug, Clone, Serialize, Deserialize)]
struct Output(char);
define_operon! {
splitter = {
Input<input_no> = get_inputs();
Intermediate<word_no> = get_words(Input) for input_no;
Output<char_no> = get_chars(Intermediate) for input_no, word_no;
}
}
#[derive(OperonService)]
#[operon(error = std::convert::Infallible)]
struct MySplitterService;
#[async_trait]
impl SplitterService for MySplitterService {
async fn get_inputs(&self) -> Result<Vec<Input>, Self::Error> {
Ok(vec![
Input::from("Hello World"),
Input::from("Hello Operon"),
Input::from(""),
])
}
async fn get_words(&self, input: Input) -> Result<Vec<Intermediate>, Self::Error> {
Ok(input
.split_whitespace()
.map(|s| Intermediate(s.to_string()))
.collect())
}
async fn get_chars(&self, intermediate: Intermediate) -> Result<Vec<Output>, Self::Error> {
log::info!("Processing intermediate: {}", intermediate.0);
Ok(intermediate.0.chars().map(Output).collect())
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let database_uri = std::env::var("POSTGRES_URI")?;
let operon_options = OperonOptions::new();
let service = Arc::new(MySplitterService);
let storage = Arc::new(
PsqlStorageOptions::new(&database_uri)
.with_schema("ex1_data")
.build::<PsqlSplitterStorage>()?,
);
let meta = PsqlMetaStorageOptions::new(&database_uri)
.with_schema("ex1_meta")
.build()?;
let operon_instance: Operon<MySplitterService, PsqlSplitterStorage, PsqlMetaStorage> =
Operon::new(service, storage.clone(), meta).with_options(operon_options);
operon_instance.run().await?;
println!("{:?}", storage.get_intermediate([0, 0]).await?);
println!("{:?}", storage.get_output([0, 0, 0]).await?);
Ok(())
}