use std::sync::{Arc, Mutex, OnceLock, Weak};
use std::time::{Duration, Instant};
use anyhow::Result;
use tokio::task::JoinHandle;
use wasmtime::Engine;
use wasmtime::component::Linker;
use crate::bindings::AppPre;
use crate::bindings::terminal as t;
use crate::loader::{self, Fetch, Loaded, UpgradeRequired};
use crate::runner::{compile, link, requires_abi, truncate};
use crate::state::HostState;
use crate::terminal::{EventQueue, Interrupt, Interrupter, PhaseHook};
use crate::{Limits, Phase, ReloadPolicy, Source, UpdateInfo};
pub const VERSION_HEADER: &str = "rattery-app-version";
pub const WATCH_INTERVAL: Duration = Duration::from_millis(750);
pub const GUEST_CHECK_INTERVAL: Duration = Duration::from_secs(1);
pub struct Running {
pub bytes: Arc<Vec<u8>>,
pub instance: AppPre<HostState>,
pub version: Option<String>,
}
struct Pending {
bytes: Arc<Vec<u8>>,
instance: AppPre<HostState>,
version: Option<String>,
shown: Option<String>,
}
struct Deadline {
id: u64,
after: Instant,
by: Instant,
task: JoinHandle<()>,
}
#[derive(Clone, Default, PartialEq, Eq)]
struct Validators {
etag: Option<String>,
last_modified: Option<String>,
}
impl Validators {
fn of(loaded: &Loaded) -> Self {
Self {
etag: loaded.etag.clone(),
last_modified: loaded.last_modified.clone(),
}
}
fn version(&self) -> Option<String> {
self.etag.clone().or_else(|| self.last_modified.clone())
}
}
struct Rejected {
validators: Validators,
version: Option<String>,
bytes: Arc<Vec<u8>>,
}
struct Inner {
running: Running,
pending: Option<Pending>,
newest: Validators,
rejected: Option<Rejected>,
deadline: Option<Deadline>,
next_deadline_id: u64,
last_rejection: Option<String>,
hint_task: Option<JoinHandle<()>>,
last_guest_check: Option<Instant>,
closed: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Offer {
Pending(UpdateInfo),
Unchanged(UpdateInfo),
Current,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Rejection {
pub reason: String,
pub requires_abi: Option<String>,
}
impl std::fmt::Display for Rejection {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.reason)
}
}
impl std::error::Error for Rejection {}
pub struct UpdatesConfig {
pub engine: Engine,
pub linker: Linker<HostState>,
pub policy: ReloadPolicy,
pub watched: bool,
pub source: Source,
pub limits: Limits,
pub queue: Arc<EventQueue>,
pub interrupter: Arc<Interrupter>,
pub on_phase: Option<PhaseHook>,
pub running: Running,
pub etag: Option<String>,
pub last_modified: Option<String>,
}
pub struct Updates {
engine: Engine,
linker: Linker<HostState>,
policy: ReloadPolicy,
watched: bool,
source: Source,
limits: Limits,
queue: Arc<EventQueue>,
interrupter: Arc<Interrupter>,
on_phase: Option<PhaseHook>,
client: OnceLock<reqwest::Client>,
checking: tokio::sync::Mutex<()>,
inner: Mutex<Inner>,
}
impl Updates {
pub fn new(config: UpdatesConfig) -> Self {
Self {
engine: config.engine,
linker: config.linker,
policy: config.policy,
watched: config.watched,
source: config.source,
limits: config.limits,
queue: config.queue,
interrupter: config.interrupter,
on_phase: config.on_phase,
client: OnceLock::new(),
checking: tokio::sync::Mutex::new(()),
inner: Mutex::new(Inner {
running: config.running,
pending: None,
newest: Validators {
etag: config.etag,
last_modified: config.last_modified,
},
rejected: None,
deadline: None,
next_deadline_id: 1,
last_rejection: None,
hint_task: None,
last_guest_check: None,
closed: false,
}),
}
}
fn phase(&self, phase: Phase) {
if let Some(hook) = &self.on_phase {
hook(phase);
}
}
fn closed_rejection() -> Rejection {
Rejection {
reason: "the app has exited".into(),
requires_abi: None,
}
}
pub fn close(&self) {
let mut inner = self.inner.lock().unwrap();
inner.closed = true;
inner.pending = None;
if let Some(deadline) = inner.deadline.take() {
deadline.task.abort();
}
if let Some(task) = inner.hint_task.take() {
task.abort();
}
}
pub fn is_closed(&self) -> bool {
self.inner.lock().unwrap().closed
}
pub fn availability(&self) -> t::Availability {
match (&self.source, self.watched) {
(Source::Bytes(_) | Source::Precompiled(_), _) => t::Availability::Unavailable,
(_, true) => t::Availability::Watched,
(_, false) => t::Availability::OnRequest,
}
}
pub async fn check(self: &Arc<Self>) -> Result<Option<UpdateInfo>> {
let _one_at_a_time = self.checking.lock().await;
self.check_locked().await
}
pub async fn check_throttled(self: &Arc<Self>) -> Result<Option<UpdateInfo>> {
{
let mut inner = self.inner.lock().unwrap();
if inner.closed {
return Err(anyhow::Error::new(Self::closed_rejection()));
}
let now = Instant::now();
if inner
.last_guest_check
.is_some_and(|t| now.duration_since(t) < GUEST_CHECK_INTERVAL)
{
return Ok(info_of(&inner));
}
inner.last_guest_check = Some(now);
}
self.check().await
}
async fn check_locked(self: &Arc<Self>) -> Result<Option<UpdateInfo>> {
let (validators, previous, rejected_standing) = {
let inner = self.inner.lock().unwrap();
if inner.closed {
return Err(anyhow::Error::new(Self::closed_rejection()));
}
match &inner.rejected {
Some(rejected) => (rejected.validators.clone(), rejected.bytes.clone(), true),
None => {
let newest_bytes = inner
.pending
.as_ref()
.map_or_else(|| inner.running.bytes.clone(), |p| p.bytes.clone());
(inner.newest.clone(), newest_bytes, false)
}
}
};
let fetched = match &self.source {
Source::Url(url) => {
let client = match self.client.get() {
Some(client) => client,
None => {
let _ = self.client.set(loader::client(&self.limits)?);
self.client.get().expect("just set")
}
};
loader::fetch_if_changed(
client,
url,
validators.etag.as_deref(),
validators.last_modified.as_deref(),
Some(&previous),
&self.limits,
)
.await
}
Source::Path(path) => loader::read_if_changed(path, &previous, &self.limits)
.await
.map(|loaded| loaded.map_or(Fetch::NotModified, Fetch::New)),
Source::Resolver(resolver) => loader::resolve_if_changed(
resolver,
validators.etag.as_deref(),
&previous,
&self.limits,
)
.await
.map(|loaded| loaded.map_or(Fetch::NotModified, Fetch::New)),
Source::Bytes(_) | Source::Precompiled(_) => Ok(Fetch::NotModified),
};
match fetched {
Ok(Fetch::New(loaded)) => match self.offer_locked(loaded, true).await {
Ok(Offer::Pending(info) | Offer::Unchanged(info)) => Ok(Some(info)),
Ok(Offer::Current) => Ok(None),
Err(rejection) => Err(anyhow::Error::new(rejection)),
},
Ok(Fetch::Same {
etag,
last_modified,
}) => {
let validators = Validators {
etag,
last_modified,
};
let version = validators.version();
let mut inner = self.inner.lock().unwrap();
if rejected_standing {
if let Some(rejected) = &mut inner.rejected {
rejected.validators = validators;
rejected.version = version;
}
} else {
inner.newest = validators;
if version.is_some() {
match &mut inner.pending {
Some(pending) => pending.version = version,
None => inner.running.version = version,
}
}
}
Ok(info_of(&inner))
}
Ok(Fetch::NotModified) => Ok(self.pending()),
Err(err) => {
if let Some(upgrade) = err.downcast_ref::<UpgradeRequired>() {
self.report_rejection(upgrade.to_string(), upgrade.required.clone());
}
Err(err)
}
}
}
pub async fn offer(self: &Arc<Self>, loaded: Loaded) -> Result<Offer, Rejection> {
let _one_at_a_time = self.checking.lock().await;
self.offer_locked(loaded, false).await
}
async fn offer_locked(
self: &Arc<Self>,
loaded: Loaded,
from_source: bool,
) -> Result<Offer, Rejection> {
let validators = from_source.then(|| Validators::of(&loaded));
if let Some(outcome) = self.settle_known(&loaded.bytes, validators.as_ref()) {
return outcome;
}
let version = Validators::of(&loaded).version();
let shown = version
.as_deref()
.map(|v| truncate(crate::sanitize::text(v), self.limits.message_bytes.min(256)));
let updates = self.clone();
let validated = tokio::task::spawn_blocking(move || {
let component = match compile(&updates.engine, &loaded) {
Ok(component) => component,
Err(err) => return Err((err, None, loaded.bytes)),
};
let requires = requires_abi(&updates.engine, &component);
match link(
&updates.engine,
&updates.linker,
&component,
&loaded.description,
) {
Ok(instance) => Ok((loaded.bytes, instance)),
Err(err) => Err((err, requires, loaded.bytes)),
}
})
.await;
let (err, requires_abi, bytes) = match validated {
Ok(Ok((bytes, instance))) => {
let bytes = Arc::new(bytes);
if let Some(outcome) = self.settle_known(&bytes, validators.as_ref()) {
return outcome;
}
let info = self.commit_pending(
Pending {
bytes,
instance,
version,
shown,
},
validators,
);
return Ok(Offer::Pending(info));
}
Ok(Err((err, requires, bytes))) => (format!("{err:#}"), requires, bytes),
Err(err) => (format!("validation failed: {err}"), None, Vec::new()),
};
let reason = truncate(crate::sanitize::text(&err), self.limits.message_bytes);
if let Some(validators) = validators {
let mut inner = self.inner.lock().unwrap();
if !inner.closed {
inner.rejected = Some(Rejected {
validators,
version,
bytes: Arc::new(bytes),
});
}
}
self.report_rejection(reason.clone(), requires_abi.clone());
Err(Rejection {
reason,
requires_abi,
})
}
fn settle_known(
&self,
bytes: &[u8],
validators: Option<&Validators>,
) -> Option<Result<Offer, Rejection>> {
let mut inner = self.inner.lock().unwrap();
if inner.closed {
return Some(Err(Self::closed_rejection()));
}
let version = validators.and_then(Validators::version);
if *inner.running.bytes == *bytes {
if let Some(validators) = validators {
inner.newest = validators.clone();
}
if version.is_some() {
inner.running.version = version;
}
inner.rejected = None;
inner.last_rejection = None;
let withdrawn = inner.pending.take().is_some();
if let Some(deadline) = inner.deadline.take() {
deadline.task.abort();
}
drop(inner);
if withdrawn {
self.phase(Phase::UpdateWithdrawn);
self.queue.push_update_changed();
}
return Some(Ok(Offer::Current));
}
if inner.pending.as_ref().is_some_and(|p| *p.bytes == *bytes) {
if let Some(validators) = validators {
inner.newest = validators.clone();
}
if version.is_some()
&& let Some(pending) = &mut inner.pending
{
pending.version = version;
}
inner.rejected = None;
inner.last_rejection = None;
return Some(Ok(Offer::Unchanged(info_of(&inner).expect("pending"))));
}
None
}
pub fn offer_same(self: &Arc<Self>, shown: Option<String>) {
let pending = {
let inner = self.inner.lock().unwrap();
if inner.closed {
return;
}
Pending {
bytes: inner.running.bytes.clone(),
instance: inner.running.instance.clone(),
version: None,
shown,
}
};
self.commit_pending(pending, None);
}
fn commit_pending(
self: &Arc<Self>,
pending: Pending,
validators: Option<Validators>,
) -> UpdateInfo {
let shown = pending.shown.clone();
let mut inner = self.inner.lock().unwrap();
inner.pending = Some(pending);
if let Some(validators) = validators {
inner.newest = validators;
}
inner.rejected = None;
inner.last_rejection = None;
if let ReloadPolicy::Deferred {
grace,
idle,
hard_limit,
} = self.policy
{
let stale = inner.deadline.as_ref().is_none_or(|d| d.task.is_finished());
if stale {
let id = inner.next_deadline_id;
inner.next_deadline_id += 1;
let now = Instant::now();
let after = now + grace;
let by = now + hard_limit.max(grace);
let task = tokio::spawn(deadline_task(Arc::downgrade(self), id, after, by, idle));
inner.deadline = Some(Deadline {
id,
after,
by,
task,
});
}
}
let info = info_of(&inner).expect("pending was just set");
drop(inner);
self.phase(Phase::UpdateAvailable { version: shown });
if self.policy == ReloadPolicy::Immediate {
self.interrupter.fire(Interrupt::Reload);
} else {
self.queue.push_update_changed();
}
info
}
pub fn pending(&self) -> Option<UpdateInfo> {
info_of(&self.inner.lock().unwrap())
}
pub fn take(&self) -> Option<AppPre<HostState>> {
let mut inner = self.inner.lock().unwrap();
if let Some(deadline) = inner.deadline.take() {
deadline.task.abort();
}
let pending = inner.pending.take()?;
inner.running = Running {
bytes: pending.bytes,
instance: pending.instance.clone(),
version: pending.version,
};
Some(pending.instance)
}
pub fn restore_running(
&self,
instance: AppPre<HostState>,
bytes: Arc<Vec<u8>>,
version: Option<String>,
) {
let mut inner = self.inner.lock().unwrap();
inner.running = Running {
bytes,
instance,
version,
};
}
pub fn running(&self) -> (Arc<Vec<u8>>, Option<String>) {
let inner = self.inner.lock().unwrap();
(inner.running.bytes.clone(), inner.running.version.clone())
}
fn knows(inner: &Inner, version: &str) -> bool {
inner.running.version.as_deref() == Some(version)
|| inner
.pending
.as_ref()
.is_some_and(|p| p.version.as_deref() == Some(version))
|| inner
.rejected
.as_ref()
.is_some_and(|r| r.version.as_deref() == Some(version))
}
pub fn hint(self: &Arc<Self>, version: &str) {
if matches!(self.source, Source::Bytes(_) | Source::Precompiled(_)) {
return;
}
let mut inner = self.inner.lock().unwrap();
if inner.closed
|| Self::knows(&inner, version)
|| inner.hint_task.as_ref().is_some_and(|t| !t.is_finished())
{
return;
}
let weak = Arc::downgrade(self);
let version = version.to_owned();
inner.hint_task = Some(tokio::spawn(async move {
let Some(updates) = weak.upgrade() else {
return;
};
let _one_at_a_time = updates.checking.lock().await;
if Self::knows(&updates.inner.lock().unwrap(), &version) {
return;
}
let _ = updates.check_locked().await;
}));
}
fn report_rejection(&self, reason: String, requires_abi: Option<String>) {
let reason = truncate(crate::sanitize::text(&reason), self.limits.message_bytes);
let mut inner = self.inner.lock().unwrap();
if inner.closed || inner.last_rejection.as_deref() == Some(reason.as_str()) {
return;
}
inner.last_rejection = Some(reason.clone());
drop(inner);
self.phase(Phase::UpdateRejected {
reason,
requires_abi,
});
}
}
async fn deadline_task(
updates: Weak<Updates>,
id: u64,
after: Instant,
by: Instant,
idle: Duration,
) {
tokio::time::sleep_until(after.into()).await;
loop {
let Some(updates) = updates.upgrade() else {
return;
};
let now = Instant::now();
let wait = {
let inner = updates.inner.lock().unwrap();
let mine = inner.deadline.as_ref().is_some_and(|d| d.id == id);
if inner.closed || !mine || inner.pending.is_none() {
return;
}
let idle_for = updates
.queue
.last_input()
.map_or(Duration::MAX, |t| now.saturating_duration_since(t));
if now >= by || idle_for >= idle {
updates.interrupter.fire(Interrupt::Reload);
return;
}
(idle - idle_for)
.min(by - now)
.clamp(Duration::from_millis(50), Duration::from_secs(1))
};
drop(updates);
tokio::time::sleep(wait).await;
}
}
fn info_of(inner: &Inner) -> Option<UpdateInfo> {
let pending = inner.pending.as_ref()?;
let now = Instant::now();
let (reload_after, reload_by) = match &inner.deadline {
Some(d) if !d.task.is_finished() => (
Some(d.after.saturating_duration_since(now)),
Some(d.by.saturating_duration_since(now)),
),
_ => (None, None),
};
Some(UpdateInfo {
version: pending.shown.clone(),
reload_after,
reload_by,
})
}
impl Drop for Updates {
fn drop(&mut self) {
let inner = self.inner.get_mut().unwrap();
if let Some(deadline) = inner.deadline.take() {
deadline.task.abort();
}
if let Some(task) = inner.hint_task.take() {
task.abort();
}
}
}
pub fn to_wit(info: UpdateInfo) -> t::Update {
t::Update {
version: info.version,
reload_after_ms: info.reload_after.map(|d| d.as_millis() as u64),
reload_by_ms: info.reload_by.map(|d| d.as_millis() as u64),
}
}