use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use libtmux::{Error, Pane};
use crate::identity::{InstanceId, InstanceIdentity};
mod registry;
mod ring;
use registry::{RetainedTail, SnapshotError, Tail, TailTable};
use ring::resume_at;
#[cfg(test)]
use ring::{RING_BYTES, Ring};
fn hold<T>(lock: &Mutex<T>) -> MutexGuard<'_, T> {
lock.lock().unwrap_or_else(PoisonError::into_inner)
}
const MAX_TAILS: usize = 8;
const MAX_TAIL_OPENERS: usize = 1;
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Cursor {
pane: String,
owner: InstanceId,
epoch: u64,
offset: u64,
}
impl Cursor {
#[must_use]
pub fn encode(&self) -> String {
format!(
"{}:{}:{}:{}",
self.pane, self.owner, self.epoch, self.offset
)
}
pub fn decode(text: &str) -> Result<Self, &str> {
let mut fields = text.rsplitn(4, ':');
let offset = fields.next().and_then(|field| field.parse().ok());
let epoch = fields.next().and_then(|field| field.parse().ok());
let owner = fields.next().and_then(InstanceId::decode);
let pane = fields.next();
match (pane, owner, epoch, offset) {
(Some(pane), Some(owner), Some(epoch), Some(offset)) if pane.starts_with('%') => {
Ok(Self {
pane: pane.to_owned(),
owner,
epoch,
offset,
})
}
_ => Err(text),
}
}
#[must_use]
pub fn pane(&self) -> &str {
&self.pane
}
}
#[derive(Debug)]
pub(crate) struct Since {
pub text: String,
pub cursor: Cursor,
pub missed: bool,
pub closed: bool,
}
#[derive(Debug)]
pub(crate) struct Tails {
identity: Arc<InstanceIdentity>,
inner: Mutex<TailTable>,
opening: tokio::sync::Semaphore,
next_epoch: AtomicU64,
}
#[derive(Debug)]
pub(crate) enum TailError {
Tmux(Error),
Snapshot { error: Error, opened: bool },
SnapshotBusy { opened: bool, limit: usize },
ReaderStopped { opened: bool },
OwnerUnavailable,
OpeningAtCapacity { limit: usize },
}
impl From<Error> for TailError {
fn from(error: Error) -> Self {
Self::Tmux(error)
}
}
impl Tails {
#[must_use]
pub(crate) fn new(identity: Arc<InstanceIdentity>) -> Self {
Self {
identity,
inner: Mutex::new(TailTable::new(MAX_TAILS)),
opening: tokio::sync::Semaphore::new(MAX_TAIL_OPENERS),
next_epoch: AtomicU64::new(0),
}
}
#[cfg(test)]
fn with_owner(owner: u128) -> Self {
Self::new(Arc::new(InstanceIdentity::fixed(owner)))
}
fn owner(&self) -> Result<InstanceId, TailError> {
self.identity.get().map_err(|_| TailError::OwnerUnavailable)
}
pub(crate) async fn read(
&self,
pane: &Pane,
cursor: Option<&Cursor>,
) -> Result<Since, TailError> {
let owner = self.owner()?;
let id = pane.id().to_string();
let (tail, opened) = self.ensure(pane, &id).await?;
let epoch = tail.epoch;
if cursor.is_none() {
let snapshot = match tail.snapshot().await {
Ok(snapshot) => snapshot,
Err(SnapshotError::Tmux(error)) => {
return Err(TailError::Snapshot { error, opened });
}
Err(SnapshotError::Busy { limit }) => {
return Err(TailError::SnapshotBusy { opened, limit });
}
Err(SnapshotError::Stopped) => {
hold(&self.inner).remove_if_epoch(&id, epoch);
return Err(TailError::ReaderStopped { opened });
}
};
let closed = hold(&tail.ring).closed;
let text = snapshot
.visible
.iter()
.map(|line| line.to_string_lossy().into_owned())
.collect::<Vec<_>>()
.join("\n");
return Ok(Since {
text,
cursor: Cursor {
pane: id,
owner,
epoch,
offset: snapshot.offset,
},
missed: false,
closed,
});
}
let (read, missed, closed, end) = {
let ring = hold(&tail.ring);
let (from, stale) = resume_at(&ring, cursor, owner, epoch);
let read = ring.snapshot_from(from);
let missed = read.missed || stale;
(read, missed, ring.closed, ring.end())
};
Ok(Since {
text: read.text(),
cursor: Cursor {
pane: id,
owner,
epoch,
offset: end,
},
missed,
closed,
})
}
async fn ensure(&self, pane: &Pane, id: &str) -> Result<(RetainedTail, bool), TailError> {
if let Some(found) = self.touch(id) {
return Ok((found, false));
}
let _opening = self
.opening
.try_acquire()
.map_err(|_| TailError::OpeningAtCapacity {
limit: MAX_TAIL_OPENERS,
})?;
if let Some(found) = self.touch(id) {
return Ok((found, false));
}
let epoch = self.next_epoch();
let tail = Tail::attach(pane, epoch).await?;
let retained = tail.retained();
hold(&self.inner).insert(id.to_owned(), tail);
Ok((retained, true))
}
fn touch(&self, id: &str) -> Option<RetainedTail> {
hold(&self.inner).touch(id)
}
fn next_epoch(&self) -> u64 {
self.next_epoch
.fetch_add(1, Ordering::Relaxed)
.wrapping_add(1)
}
}
#[cfg(test)]
mod tests;