use std::{
cmp::Ordering,
collections::HashMap,
path::PathBuf,
sync::Arc,
time::{Duration, Instant},
};
use anyhow::Result;
use bnf_sampler::{grammar::Grammar, vocabulary::Vocabulary};
use derivative::Derivative;
use flume::{Receiver, Sender};
use itertools::Itertools;
use qp_trie::Trie;
use salvo::oapi::ToSchema;
use serde::{Deserialize, Serialize};
use tokio::sync::{Mutex, RwLock};
use web_rwkv::{
context::Context,
runtime::{
infer::{InferChunk, InferInfo, InferInput, InferInputBatch, InferOption, InferOutput},
model::{ModelInfo, ModelRuntime, State},
softmax::softmax,
Job, JobBuilder, JobRuntime,
},
tensor::{TensorCpu, TensorInit},
tokenizer::Tokenizer,
};
use crate::{
sampler::{bnf::BnfSampler, Transformer},
Environment, FinishReason, GenerateRequest, ReloadRequest, Token, TokenCounter,
};
const END_OF_LINE_TOKEN: u16 = 261;
const PROMPT_CACHE_TOKENS: usize = 32;
const MAX_CACHE_ITEMS: usize = 256;
const SAMPLER_ARENA_CAPACITY: usize = 1048576;
const GRAMMAR_ARENA_CAPACITY: usize = 1024;
#[derive(Debug)]
pub enum SlotResult {
Success(usize),
Fault(usize),
Failure(Box<GenerateContext>),
Error(String),
}
#[derive(Debug)]
enum SlotState {
Idle(Tokens, Instant),
Wait(Box<GenerateContext>),
Busy,
}
impl Default for SlotState {
fn default() -> Self {
Self::Idle(Default::default(), Instant::now())
}
}
#[derive(Debug, PartialEq, Eq)]
enum SlotChoice {
Continue(usize, usize),
Back(usize),
Empty(usize),
}
impl std::cmp::Ord for SlotChoice {
fn cmp(&self, other: &Self) -> Ordering {
use SlotChoice::{Back, Continue, Empty};
match (self, other) {
(Continue(_, x), Continue(_, y)) => x.cmp(y),
(Continue(_, _), _) => Ordering::Greater,
(_, Continue(_, _)) => Ordering::Less,
(Empty(_), Empty(_)) => Ordering::Equal,
(Empty(_), Back(_)) => Ordering::Greater,
(Back(_), Empty(_)) => Ordering::Less,
(Back(_), Back(_)) => Ordering::Equal,
}
}
}
impl std::cmp::PartialOrd for SlotChoice {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
#[derive(Debug, Clone, Default)]
pub enum Payload {
#[default]
Empty,
Busy(GenerateContext),
Done(GenerateContext),
}
impl Payload {
pub fn take(&mut self) -> Option<GenerateContext> {
match std::mem::take(self) {
Payload::Done(context) => Some(context),
payload => {
*self = payload;
None
}
}
}
pub fn finalize(&mut self) {
*self = match std::mem::take(self) {
Payload::Busy(context) => Payload::Done(context),
payload => payload,
}
}
#[must_use]
pub fn is_empty(&self) -> bool {
matches!(self, Self::Empty)
}
}
#[repr(transparent)]
#[derive(Debug, Default, Clone)]
pub struct Tokens(pub Vec<u16>);
impl std::ops::Deref for Tokens {
type Target = TokenSlice;
fn deref(&self) -> &Self::Target {
self.0.as_token_slice()
}
}
impl std::borrow::Borrow<[u8]> for Tokens {
fn borrow(&self) -> &[u8] {
bytemuck::cast_slice(&self.0)
}
}
impl std::borrow::Borrow<[u16]> for Tokens {
fn borrow(&self) -> &[u16] {
&self.0
}
}
impl std::borrow::Borrow<TokenSlice> for Tokens {
fn borrow(&self) -> &TokenSlice {
self.0[..].as_token_slice()
}
}
impl qp_trie::Break for Tokens {
type Split = TokenSlice;
fn empty<'a>() -> &'a Self::Split {
Default::default()
}
fn find_break(&self, loc: usize) -> &Self::Split {
self.0[..loc >> 1].as_token_slice()
}
}
#[repr(transparent)]
pub struct TokenSlice([u16]);
impl std::ops::Deref for TokenSlice {
type Target = [u16];
fn deref(&self) -> &Self::Target {
&self.0
}
}
impl std::borrow::Borrow<[u8]> for TokenSlice {
fn borrow(&self) -> &[u8] {
bytemuck::cast_slice(&self.0)
}
}
impl Default for &TokenSlice {
fn default() -> Self {
<&[u16]>::default().as_token_slice()
}
}
pub trait AsTokenSlice {
fn as_token_slice(&self) -> &TokenSlice;
}
impl AsTokenSlice for [u16] {
fn as_token_slice(&self) -> &TokenSlice {
let ptr = self as *const [u16] as *const TokenSlice;
unsafe { &*ptr }
}
}
#[derive(Derivative, Clone)]
#[derivative(Debug)]
pub struct GenerateContext {
pub prompt_tokens: Vec<u16>,
pub prompt_cached: bool,
pub prefix: Tokens,
pub suffix: Tokens,
pub model_text: Vec<u8>,
pub buffer: Vec<u8>,
pub model_tokens: Vec<u16>,
#[derivative(Debug = "ignore")]
pub transformers: Vec<Arc<RwLock<dyn Transformer + Send + Sync>>>,
pub instant: Option<Instant>,
pub request: GenerateRequest,
pub sender: Sender<Token>,
}
#[derive(Debug)]
struct CachedItem<T> {
item: Arc<T>,
instant: Instant,
}
impl<T> CachedItem<T> {
pub fn new(backed: T) -> Self {
Self {
item: Arc::new(backed),
instant: Instant::now(),
}
}
pub fn update(cached: CachedItem<T>) -> Self {
Self {
item: cached.item,
instant: Instant::now(),
}
}
}
impl<T> Clone for CachedItem<T> {
fn clone(&self) -> Self {
Self {
item: self.item.clone(),
instant: self.instant,
}
}
}
#[derive(Debug, Default)]
struct Cache {
state: Option<InitState>,
cache: Trie<Tokens, CachedItem<TensorCpu<f32>>>,
}
impl Cache {
fn maintain(&mut self) {
let cache = &mut self.cache;
if cache.count() <= MAX_CACHE_ITEMS {
return;
}
let mut remove = vec![];
for (tokens, _) in cache
.iter()
.sorted_unstable_by_key(|(_, item)| item.instant.elapsed())
.skip(MAX_CACHE_ITEMS)
{
remove.push(tokens.to_owned());
}
for tokens in remove.into_iter() {
cache.remove(&tokens);
}
}
}
#[derive(Debug, Default)]
struct CacheHub {
backed: HashMap<StateId, Cache>,
default: Cache,
}
impl CacheHub {
fn fetch(&mut self, id: StateId) -> &mut Cache {
match self.backed.get_mut(&id) {
Some(item) => item,
None => &mut self.default,
}
}
}
#[derive(
Derivative, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, ToSchema,
)]
#[derivative(Debug = "transparent")]
#[serde(transparent)]
pub struct StateId(uuid::Uuid);
impl StateId {
pub fn new() -> Self {
Self(uuid::Uuid::new_v4())
}
}
#[derive(Derivative, Clone)]
#[derivative(Debug)]
pub struct InitState {
pub name: String,
pub id: StateId,
pub default: bool,
#[derivative(Debug = "ignore")]
pub data: TensorCpu<f32>,
}
struct Model<M>(M);
trait ModelSerialize {
fn serialize(&self, file: std::fs::File) -> Result<()>;
}
impl<M: Serialize> ModelSerialize for Model<M> {
fn serialize(&self, file: std::fs::File) -> Result<()> {
use cbor4ii::{core::enc::Write, serde::Serializer};
use std::{fs::File, io::Write as _};
struct FileWriter(File);
impl Write for FileWriter {
type Error = std::io::Error;
fn push(&mut self, input: &[u8]) -> Result<(), Self::Error> {
self.0.write_all(input)
}
}
let file = FileWriter(file);
let mut serializer = Serializer::new(file);
self.0.serialize(&mut serializer)?;
Ok(())
}
}
impl Environment {
pub async fn enqueue(&self, context: GenerateContext) -> Vec<GenerateContext> {
let mut queue = vec![];
match self {
Environment::Loaded(runtime) => {
match runtime.queue(context).await.expect("queue task error") {
SlotResult::Success(batch) => log::info!("queued task at slot {batch}"),
SlotResult::Fault(batch) => log::info!("swapped task at slot {batch}"),
SlotResult::Failure(context) => queue.push(*context),
SlotResult::Error(reason) => log::warn!("queue task failed: {}", reason),
}
}
Environment::None => queue.push(context),
};
queue
}
}
pub struct Runtime {
context: Context,
reload: ReloadRequest,
info: ModelInfo,
state: Arc<dyn State + Send + Sync>,
model: Arc<dyn ModelSerialize + Send + Sync>,
runtime: JobRuntime<InferInput, InferOutput>,
tokenizer: Arc<Tokenizer>,
vocab: Arc<Vocabulary>,
slots: Mutex<Vec<SlotState>>,
caches: Mutex<CacheHub>,
}
impl Runtime {
pub async fn new<J, B>(
context: Context,
builder: B,
reload: ReloadRequest,
states: Vec<InitState>,
tokenizer: Tokenizer,
vocab: Vocabulary,
) -> Self
where
J: Job<Info = InferInfo, Input = InferChunk, Output = InferOutput>,
B: JobBuilder<J, Info = InferInfo> + ModelRuntime,
{
let slots = (0..reload.max_batch)
.map(|_| SlotState::default())
.collect();
let mut caches = CacheHub::default();
if let Some(state) = states.iter().find(|state| state.default) {
caches.default.state = Some(state.clone());
}
for state in states {
let id = state.id;
let item = Cache {
state: Some(state),
cache: Trie::new(),
};
caches.backed.insert(id, item);
}
let info = builder.info();
let state = Arc::new(builder.state());
let model = Arc::new(Model(builder.model()));
let runtime = JobRuntime::new(builder).await;
Self {
context,
reload,
info,
state,
model,
runtime,
tokenizer: Arc::new(tokenizer),
vocab: Arc::new(vocab),
slots: Mutex::new(slots),
caches: Mutex::new(caches),
}
}
#[inline]
pub fn context(&self) -> &Context {
&self.context
}
#[inline]
pub fn reload(&self) -> &ReloadRequest {
&self.reload
}
#[inline]
pub fn info(&self) -> &ModelInfo {
&self.info
}
#[inline]
pub fn num_batch(&self) -> usize {
self.state.num_batch()
}
#[inline]
pub fn tokenizer(&self) -> Arc<Tokenizer> {
self.tokenizer.clone()
}
pub async fn states(&self) -> Vec<(StateId, InitState)> {
let caches = self.caches.lock().await;
let mut states = vec![];
if let Some(state) = &caches.default.state {
states.push((state.id, state.clone()));
}
for item in caches.backed.values() {
if let Some(state) = &item.state {
states.push((state.id, state.clone()));
}
}
states
}
pub async fn load_init_state(&self, state: InitState) {
let mut caches = self.caches.lock().await;
caches.backed.insert(
state.id,
Cache {
state: Some(state),
cache: Trie::new(),
},
);
}
pub async fn unload_init_state(&self, id: StateId) {
let mut caches = self.caches.lock().await;
caches.backed.remove(&id);
}
pub async fn serialize_model(&self, path: PathBuf) -> Result<()> {
let model = self.model.clone();
let handle = tokio::task::spawn_blocking(move || {
let file = std::fs::File::create(path)?;
model.serialize(file)
});
handle.await?
}
async fn checkout(
&self,
id: StateId,
tokens: &[u16],
batch: usize,
) -> (Vec<u16>, Arc<TensorCpu<f32>>) {
let mut caches = self.caches.lock().await;
let Cache { state, cache } = caches.fetch(id);
let prefix = cache.longest_common_prefix(tokens.as_token_slice());
let len = (1..=prefix.len())
.rev()
.find(|len| cache.contains_key(prefix[0..*len].as_token_slice()))
.unwrap_or_default();
log::info!("slot {} checks out backed cache of length {}", batch, len);
let prefix = prefix[0..len].to_vec();
let state = state.clone().map(|state| state.data);
let reload = match cache.remove(prefix[..].as_token_slice()) {
Some(reload) => CachedItem::update(reload),
None => CachedItem::new(state.unwrap_or_else(|| self.state.init())),
};
if len > 0 {
let key = Tokens(prefix.clone());
cache.insert(key, reload.clone());
}
(prefix, reload.item)
}
async fn compile_bnf_schema(&self, schema: String) -> Result<BnfSampler> {
let grammar = Grammar::new(&schema, self.vocab.clone(), GRAMMAR_ARENA_CAPACITY)?;
let start_nonterminal = self.reload.bnf.start_nonterminal.clone();
let sampler = bnf_sampler::sampler::Sampler::new(
grammar,
start_nonterminal,
self.vocab.clone(),
SAMPLER_ARENA_CAPACITY,
self.reload.bnf.enable_bytes_cache,
)?;
Ok(BnfSampler::new(sampler))
}
pub async fn queue(&self, context: GenerateContext) -> Result<SlotResult> {
let mut slots = self.slots.lock().await;
let (last, tokens) = match [context.prefix, context.suffix].concat().split_last() {
Some((last, tokens)) => (*last, tokens.to_vec()),
None => return Ok(SlotResult::Error("empty task is not queued".into())),
};
let mut transformers = Vec::<Arc<RwLock<dyn Transformer + Send + Sync>>>::new();
if let Some(schema) = context.request.bnf_schema.clone() {
match self.compile_bnf_schema(schema).await {
Ok(bnf) => transformers.push(Arc::new(RwLock::new(bnf))),
Err(err) => return Ok(SlotResult::Error(err.to_string())),
}
}
let choice = slots
.iter()
.enumerate()
.filter_map(|(batch, slot)| match slot {
SlotState::Idle(content, time) => {
let delta = time.elapsed().as_millis();
match (content.is_empty(), tokens.starts_with(content)) {
(true, _) => Some((SlotChoice::Empty(batch), delta)),
(false, true) => Some((SlotChoice::Continue(batch, content.len()), delta)),
(false, false) => Some((SlotChoice::Back(batch), delta)),
}
}
_ => None,
})
.max_by(|lhs, rhs| lhs.0.cmp(&rhs.0).then(lhs.1.cmp(&rhs.1)))
.map(|(x, _)| x);
match choice {
None => Ok(SlotResult::Failure(
GenerateContext {
prefix: Default::default(),
suffix: Tokens([tokens, vec![last]].concat()),
transformers,
..context
}
.into(),
)),
Some(SlotChoice::Back(batch)) => {
log::info!("start at non-empty slot {}", batch);
let (prefix, reload) = self.checkout(context.request.state, &tokens, batch).await;
self.state.load(batch, reload.as_ref().clone())?;
let tokens = [tokens, vec![last]].concat();
let len = prefix.len();
let mut state = SlotState::Wait(
GenerateContext {
prefix: Tokens(tokens[..len].to_vec()),
suffix: Tokens(tokens[len..].to_vec()),
transformers,
..context
}
.into(),
);
std::mem::swap(&mut state, &mut slots[batch]);
Ok(SlotResult::Fault(batch))
}
Some(SlotChoice::Empty(batch)) => {
log::info!("start at empty slot {}", batch);
let (prefix, reload) = self.checkout(context.request.state, &tokens, batch).await;
self.state.load(batch, reload.as_ref().clone())?;
let tokens = [tokens, vec![last]].concat();
let len = prefix.len();
let state = SlotState::Wait(
GenerateContext {
prefix: Tokens(tokens[..len].to_vec()),
suffix: Tokens(tokens[len..].to_vec()),
transformers,
..context
}
.into(),
);
slots[batch] = state;
Ok(SlotResult::Fault(batch))
}
Some(SlotChoice::Continue(batch, len)) => {
log::info!("continue at slot {}", batch);
let tokens = [tokens, vec![last]].concat();
let state = SlotState::Wait(
GenerateContext {
prefix: Tokens(tokens[..len].to_vec()),
suffix: Tokens(tokens[len..].to_vec()),
transformers,
..context
}
.into(),
);
slots[batch] = state;
Ok(SlotResult::Success(batch))
}
}
}
async fn prepare(&self, payloads: &mut [Payload]) -> Result<()> {
let mut slots = self.slots.lock().await;
for (slot, payload) in slots.iter().zip_eq(payloads.iter_mut()) {
if !(payload.is_empty() || matches!(slot, SlotState::Busy)) {
log::warn!("payload should either be empty or slot should be busy");
*payload = Payload::Empty;
}
}
for (batch, payload) in payloads.iter_mut().enumerate() {
let Some(context) = payload.take() else {
continue;
};
let backed = self.state.back(batch).await?;
if context.request.embed {
let layer = context
.request
.embed_layer
.clamp(0, self.info.num_layer - 1);
let backed = backed.clone();
let embed = self.state.embed(layer, backed)?.to_vec();
let _ = context.sender.send(Token::Embed(embed));
}
let mut caches = self.caches.lock().await;
let cache = &mut caches.fetch(context.request.state).cache;
cache.insert(context.prefix.clone(), CachedItem::new(backed));
log::info!(
"backed completed slot {} of length {}",
batch,
context.prefix.len()
);
assert!(matches!(slots[batch], SlotState::Busy));
slots[batch] = SlotState::Idle(context.prefix, Instant::now());
}
let occupancy = payloads
.iter()
.filter(|x| matches!(x, Payload::Busy(_)))
.count();
let remain = self.reload.max_batch - self.reload.max_batch.min(occupancy);
let batches = slots
.iter()
.enumerate()
.filter(|(_, slot)| matches!(slot, SlotState::Wait(_)))
.take(remain)
.map(|(batch, _)| batch)
.collect_vec();
for batch in batches {
let mut slot = SlotState::Busy;
std::mem::swap(&mut slots[batch], &mut slot);
match slot {
SlotState::Wait(context) => {
let _ = context.sender.send(Token::Start);
assert!(matches!(payloads[batch], Payload::Empty));
payloads[batch] = Payload::Busy(*context);
}
_ => unreachable!(),
};
}
Ok(())
}
async fn process(&self, payloads: &mut [Payload]) -> Result<()> {
self.prepare(payloads).await?;
let batches = payloads
.iter()
.map(|payload| match payload {
Payload::Busy(context) => context.suffix.0.clone(),
_ => vec![],
})
.map(|tokens| InferInputBatch {
tokens,
option: InferOption::Last,
})
.collect();
let inference = InferInput::new(batches, self.reload.token_chunk_size);
if inference.num_token() == 0 {
return Ok(());
}
let mut inference = Some(inference);
let outputs = loop {
let input = inference.take().unwrap();
let (input, output) = self.runtime.infer(input).await;
inference = Some(input);
if output.iter().any(|batch| batch.size() > 0) {
break output;
}
};
let mut set = tokio::task::JoinSet::new();
for (batch, (payload, output)) in payloads.iter().zip_eq(outputs.iter()).enumerate() {
match (payload, output) {
(Payload::Busy(context), output) if output.size() > 0 => {
let num_vocab = self.info.num_vocab;
let output = output.0.clone();
let transformers = context.transformers.clone();
let sampler = context.request.sampler.clone();
let bias = context.request.bias.clone();
set.spawn(async move {
let mut data = output.to_vec();
assert_eq!(data.len(), num_vocab);
sampler.read().await.transform(&mut data);
for (token, bias) in bias.iter() {
data[*token as usize] += *bias;
}
for transformer in transformers {
transformer.read().await.transform(&mut data);
}
(batch, data)
});
}
_ => {}
}
}
let mut outputs = HashMap::new();
while let Some(Ok((batch, data))) = set.join_next().await {
outputs.insert(batch, data);
}
let outputs = (0..payloads.len())
.map(|batch| outputs.remove(&batch))
.map(|data| match data {
Some(data) => TensorCpu::from_data([self.info.num_vocab, 1, 1, 1], data),
None => TensorCpu::from_data([self.info.num_vocab, 0, 1, 1], vec![]),
})
.try_collect()?;
let outputs = softmax(&self.context, outputs).await?;
let mut set = tokio::task::JoinSet::new();
for (batch, (payload, output)) in
payloads.iter_mut().zip_eq(outputs.into_iter()).enumerate()
{
match (payload, output) {
(Payload::Busy(context), output) if output.size() > 0 => {
let num_vocab = self.info.num_vocab;
let sampler = context.request.sampler.clone();
set.spawn(async move {
let data = output.to_vec();
assert_eq!(data.len(), num_vocab);
let token = sampler.write().await.sample(&data);
(batch, token)
});
}
_ => {}
}
}
let mut tokens = HashMap::new();
while let Some(Ok((batch, token))) = set.join_next().await {
tokens.insert(batch, token);
}
let inference = inference.unwrap();
for (batch, (payload, input)) in
itertools::multizip((payloads.iter_mut(), inference.batches.into_iter())).enumerate()
{
let Payload::Busy(context) = payload else {
continue;
};
let instant = context.instant.get_or_insert(Instant::now());
let prefix = std::mem::take(&mut context.prefix);
let suffix = std::mem::take(&mut context.suffix);
let model_tokens = [prefix.0, suffix.0].concat();
assert!(model_tokens.len() >= input.tokens.len());
let len = model_tokens.len() - input.tokens.len();
context.prefix = Tokens(model_tokens[..len].to_vec());
context.suffix = Tokens(model_tokens[len..].to_vec());
let Some(&token) = tokens.get(&batch) else {
continue;
};
if !context.prompt_cached && context.prompt_tokens.len() > PROMPT_CACHE_TOKENS {
let mut caches = self.caches.lock().await;
let cache = &mut caches.fetch(context.request.state).cache;
let backed = self.state.back(batch).await?;
cache.insert(context.prefix.clone(), CachedItem::new(backed));
context.prompt_cached = true;
log::info!(
"backed prompt of slot {} of length {}",
batch,
context.prefix.len()
);
}
let token = match token {
0 => END_OF_LINE_TOKEN,
_ => token,
};
assert_eq!(context.suffix.len(), 0);
context.suffix.0.push(token);
let mut word = self.tokenizer.decode(&[token])?;
context.model_text.append(&mut word.clone());
context.buffer.append(&mut word);
context.model_tokens.push(token);
let mut done = false;
let mut finish = |reason| {
let counter = {
let prompt = context.prompt_tokens.len();
let completion = context.model_tokens.len();
let total = prompt + completion;
let duration = instant.elapsed();
TokenCounter {
prompt,
completion,
total,
duration,
}
};
let _ = context.sender.send(Token::Stop(reason, counter));
let _ = context.sender.send(Token::Done);
done = true;
};
let mut exhausted = false;
for transformer in context.transformers.iter() {
let mut transformer = transformer.write().await;
exhausted |= transformer.update(token);
}
let ((head, tail), stop_matched) = context
.request
.stop
.iter()
.map(|stop| {
let stop = stop.as_bytes();
let mut index_safe = 0;
let mut index_unsafe = 0;
while index_unsafe < context.buffer.len() {
let index_stop = index_unsafe - index_safe;
if index_stop >= stop.len() {
return (index_safe, true);
}
let output = context.buffer[index_unsafe];
let stop = stop[index_stop];
index_unsafe += 1;
if output != stop {
index_safe = index_unsafe;
}
}
(index_safe, index_unsafe - index_safe >= stop.len())
})
.min_by(|x, y| match (x.1, y.1) {
(true, false) => Ordering::Less,
(false, true) => Ordering::Greater,
_ => x.0.cmp(&y.0),
})
.map(|(mid, matched)| (context.buffer.split_at(mid), matched))
.unwrap_or(((&context.buffer[..], &[]), false));
if context.sender.is_disconnected() {
done = true;
} else if exhausted || stop_matched {
let output = String::from_utf8_lossy(head);
let _ = context.sender.send(Token::Content(output.into()));
finish(FinishReason::Stop);
} else if context.model_tokens.len() >= context.request.max_tokens {
finish(FinishReason::Length);
} else if let Ok(word) = String::from_utf8(head.to_vec()) {
let _ = context.sender.send(Token::Content(word));
context.buffer = tail.to_vec();
}
done.then(|| payload.finalize());
}
Ok(())
}
async fn maintain_cache(&self) {
let mut caches = self.caches.lock().await;
caches.default.maintain();
caches.backed.iter_mut().for_each(|(_, x)| x.maintain());
}
}
pub async fn run(receiver: Receiver<()>, env: Arc<RwLock<Environment>>) {
{
let env = env.clone();
tokio::spawn(async move {
loop {
if let Environment::Loaded(runtime) = &*env.read().await {
runtime.maintain_cache().await;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
});
}
while let Ok(()) = receiver.recv_async().await {
if let Environment::Loaded(runtime) = &*env.read().await {
let mut payloads = vec![Payload::default(); runtime.num_batch()];
'run: loop {
if let Err(err) = runtime.process(&mut payloads).await {
log::error!("{}", err);
break 'run;
}
if payloads.iter().all(Payload::is_empty) {
break 'run;
}
}
}
}
}