use std::cell::Cell;
use std::fs::{self, File};
use std::io::{BufWriter, Write};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{channel, sync_channel, Receiver, Sender, SyncSender, TrySendError};
use std::sync::{Condvar, Mutex, OnceLock};
const DEFAULT_QUEUE_BOUND: usize = 4096;
const DEFAULT_TOTAL_MAX_BYTES: u64 = 256 * 1024 * 1024;
enum WriterMsg {
Event(String),
Memo { key: String, json: String },
Sidecar { name: String, json: String },
Flush(Sender<()>),
Stop,
}
pub struct JournalWriter {
tx: SyncSender<WriterMsg>,
dropped: Cell<u64>,
}
impl JournalWriter {
pub fn spawn(dir: PathBuf, file: File) -> Self {
let bound = env_usize("SEMA_WORKFLOW_JOURNAL_QUEUE", DEFAULT_QUEUE_BOUND).max(1);
let max_bytes = env_u64("SEMA_WORKFLOW_JOURNAL_MAX_BYTES", DEFAULT_TOTAL_MAX_BYTES);
let (tx, rx) = sync_channel::<WriterMsg>(bound);
let _ = std::thread::Builder::new()
.name("sema-wf-journal".to_string())
.spawn(move || writer_loop(dir, file, max_bytes, rx));
Self {
tx,
dropped: Cell::new(0),
}
}
pub fn enqueue_event(&self, line: String) {
let dropped = self.dropped.get();
if dropped > 0 {
let marker = format!(
r#"{{"event":"journal.overflow","reason":"queue-full","dropped":{dropped}}}"#
);
if self.tx.try_send(WriterMsg::Event(marker)).is_ok() {
self.dropped.set(0);
}
}
if let Err(TrySendError::Full(_)) = self.tx.try_send(WriterMsg::Event(line)) {
self.dropped.set(self.dropped.get() + 1);
}
}
pub fn enqueue_memo(&self, key: String, json: String) {
let _ = self.tx.try_send(WriterMsg::Memo { key, json });
}
pub fn enqueue_sidecar(&self, name: String, json: String) {
let _ = self.tx.try_send(WriterMsg::Sidecar { name, json });
}
pub fn request_flush(&self) -> Receiver<()> {
let (ack_tx, ack_rx) = channel();
let _ = self.tx.try_send(WriterMsg::Flush(ack_tx));
ack_rx
}
}
impl Drop for JournalWriter {
fn drop(&mut self) {
let _ = self.tx.try_send(WriterMsg::Stop);
}
}
fn writer_loop(dir: PathBuf, file: File, max_bytes: u64, rx: Receiver<WriterMsg>) {
let mut out = BufWriter::new(file);
let mut written: u64 = 0;
let mut size_overflowed = false;
while let Ok(msg) = rx.recv() {
wait_if_stalled();
match msg {
WriterMsg::Event(mut line) => {
line.push('\n');
let cost = line.len() as u64;
if size_overflowed {
continue;
}
if written + cost > max_bytes {
size_overflowed = true;
let marker = format!(
"{{\"event\":\"journal.overflow\",\"reason\":\"journal-size-cap\",\"bytes\":{written}}}\n"
);
let _ = out.write_all(marker.as_bytes());
let _ = out.flush();
continue;
}
let _ = out.write_all(line.as_bytes());
let _ = out.flush();
written += cost;
}
WriterMsg::Memo { key, json } => {
let memo_dir = dir.join("memo");
if fs::create_dir_all(&memo_dir).is_ok() {
let _ = fs::write(memo_dir.join(format!("{key}.json")), json);
}
}
WriterMsg::Sidecar { name, json } => {
let _ = fs::write(dir.join(name), json);
}
WriterMsg::Flush(ack) => {
let _ = out.flush();
let _ = ack.send(());
}
WriterMsg::Stop => break,
}
}
let _ = out.flush();
}
fn env_usize(key: &str, default: usize) -> usize {
std::env::var(key)
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(default)
}
fn env_u64(key: &str, default: u64) -> u64 {
std::env::var(key)
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(default)
}
static STALL_ARMED: AtomicBool = AtomicBool::new(false);
static STALL_STATE: OnceLock<(Mutex<bool>, Condvar)> = OnceLock::new();
fn stall_state() -> &'static (Mutex<bool>, Condvar) {
STALL_STATE.get_or_init(|| (Mutex::new(false), Condvar::new()))
}
fn wait_if_stalled() {
if !STALL_ARMED.load(Ordering::Relaxed) {
return;
}
let (lock, cvar) = stall_state();
let mut stalled = lock.lock().unwrap_or_else(|e| e.into_inner());
while STALL_ARMED.load(Ordering::Relaxed) && *stalled {
stalled = cvar.wait(stalled).unwrap_or_else(|e| e.into_inner());
}
}
#[doc(hidden)]
pub fn __test_set_stall(stalled: bool) {
let (lock, cvar) = stall_state();
*lock.lock().unwrap_or_else(|e| e.into_inner()) = stalled;
STALL_ARMED.store(true, Ordering::SeqCst);
cvar.notify_all();
}
#[doc(hidden)]
pub fn __test_disarm_stall() {
let (lock, cvar) = stall_state();
*lock.lock().unwrap_or_else(|e| e.into_inner()) = false;
STALL_ARMED.store(false, Ordering::SeqCst);
cvar.notify_all();
}