use indexmap::IndexMap;
use lex_bytecode::vm::Tracer;
use lex_bytecode::Value;
use serde::{Deserialize, Serialize};
use std::sync::{Arc, Mutex};
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct RunId(pub String);
impl RunId {
pub fn new(seed: &str) -> Self {
use sha2::{Digest, Sha256};
let mut h = Sha256::new();
h.update(seed.as_bytes());
h.update(format!("{:?}", std::time::SystemTime::now()).as_bytes());
let r = h.finalize();
let mut hex = String::with_capacity(64);
for b in r { hex.push_str(&format!("{:02x}", b)); }
RunId(hex)
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum TraceNodeKind { Call, Effect }
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct TraceNode {
pub node_id: String,
pub kind: TraceNodeKind,
pub target: String,
pub input: serde_json::Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
pub started_at: u64,
pub ended_at: u64,
#[serde(default)]
pub children: Vec<TraceNode>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct TraceTree {
pub run_id: String,
pub root_target: String,
pub root_input: serde_json::Value,
pub root_output: Option<serde_json::Value>,
pub root_error: Option<String>,
pub started_at: u64,
pub ended_at: u64,
pub nodes: Vec<TraceNode>,
}
impl TraceTree {
pub fn find(&self, node_id: &str) -> Option<&TraceNode> {
for n in &self.nodes {
if let Some(found) = find_in(n, node_id) { return Some(found); }
}
None
}
}
fn find_in<'a>(n: &'a TraceNode, target: &str) -> Option<&'a TraceNode> {
if n.node_id == target { return Some(n); }
for c in &n.children {
if let Some(f) = find_in(c, target) { return Some(f); }
}
None
}
pub struct Recorder {
state: Arc<Mutex<RecorderState>>,
}
pub(crate) struct RecorderState {
open: Vec<OpenFrame>,
completed: Vec<TraceNode>,
pub(crate) overrides: IndexMap<String, serde_json::Value>,
}
struct OpenFrame {
node: TraceNode,
children: Vec<TraceNode>,
}
impl Recorder {
pub fn new() -> Self {
Self {
state: Arc::new(Mutex::new(RecorderState {
open: Vec::new(),
completed: Vec::new(),
overrides: IndexMap::new(),
})),
}
}
pub fn handle(&self) -> Handle {
Handle { state: Arc::clone(&self.state) }
}
pub fn with_overrides(self, overrides: IndexMap<String, serde_json::Value>) -> Self {
self.state.lock().unwrap().overrides = overrides;
self
}
}
impl Default for Recorder { fn default() -> Self { Self::new() } }
#[derive(Clone)]
pub struct Handle {
state: Arc<Mutex<RecorderState>>,
}
impl Handle {
pub fn finalize(
&self,
root_target: impl Into<String>,
root_input: serde_json::Value,
root_output: Option<serde_json::Value>,
root_error: Option<String>,
started_at: u64,
ended_at: u64,
) -> TraceTree {
let st = self.state.lock().unwrap();
TraceTree {
run_id: RunId::new(&format!("{}-{}", started_at, ended_at)).0,
root_target: root_target.into(),
root_input,
root_output,
root_error,
started_at,
ended_at,
nodes: st.completed.clone(),
}
}
}
fn now_unix() -> u64 {
use std::time::{SystemTime, UNIX_EPOCH};
SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0)
}
fn values_to_json(args: &[Value]) -> serde_json::Value {
serde_json::Value::Array(args.iter().map(value_to_json).collect())
}
fn value_to_json(v: &Value) -> serde_json::Value {
use serde_json::Value as J;
match v {
Value::Int(n) => J::from(*n),
Value::Float(f) => J::from(*f),
Value::Bool(b) => J::Bool(*b),
Value::Str(s) => J::String(s.to_string()),
Value::Bytes(b) => J::String(b.iter().map(|b| format!("{:02x}", b)).collect()),
Value::Unit => J::Null,
Value::List(items) => J::Array(items.iter().map(value_to_json).collect()),
Value::Tuple(items) => J::Array(items.iter().map(value_to_json).collect()),
Value::Record(fields) => {
let mut m = serde_json::Map::new();
for (k, v) in fields { m.insert(k.clone(), value_to_json(v)); }
J::Object(m)
}
Value::Variant { name, args } => {
let mut m = serde_json::Map::new();
m.insert("$variant".into(), J::String(name.clone()));
m.insert("args".into(), J::Array(args.iter().map(value_to_json).collect()));
J::Object(m)
}
Value::Closure { body_hash, .. } => {
let prefix: String = body_hash.iter().take(4)
.map(|b| format!("{b:02x}")).collect();
J::String(format!("<closure {prefix}>"))
}
Value::F64Array { rows, cols, data } => {
let mut m = serde_json::Map::new();
m.insert("$f64_array".into(), J::Bool(true));
m.insert("rows".into(), J::from(*rows));
m.insert("cols".into(), J::from(*cols));
m.insert("data".into(), J::Array(data.iter().map(|f| J::from(*f)).collect()));
J::Object(m)
}
Value::Map(m) => {
let mut o = serde_json::Map::new();
o.insert("$map".into(), J::Bool(true));
o.insert("entries".into(), J::Array(m.iter().map(|(k, v)| {
J::Array(vec![value_to_json(&k.as_value()), value_to_json(v)])
}).collect()));
J::Object(o)
}
Value::Set(s) => {
let mut o = serde_json::Map::new();
o.insert("$set".into(), J::Bool(true));
o.insert("items".into(), J::Array(
s.iter().map(|k| value_to_json(&k.as_value())).collect()));
J::Object(o)
}
Value::Deque(items) => {
let mut o = serde_json::Map::new();
o.insert("$deque".into(), J::Bool(true));
o.insert("items".into(), J::Array(
items.iter().map(value_to_json).collect()));
J::Object(o)
}
Value::Actor(_) => J::String("<actor>".into()),
Value::Ticker(_) => J::String("<ticker>".into()),
Value::ArrowTable(t) => {
let mut o = serde_json::Map::new();
o.insert("$arrow_table".into(), J::Bool(true));
o.insert("nrows".into(), J::from(t.num_rows() as i64));
o.insert("ncols".into(), J::from(t.num_columns() as i64));
J::Object(o)
}
}
}
pub(crate) fn json_to_value(v: &serde_json::Value) -> Value {
use serde_json::Value as J;
match v {
J::Null => Value::Unit,
J::Bool(b) => Value::Bool(*b),
J::Number(n) => {
if let Some(i) = n.as_i64() { Value::Int(i) }
else if let Some(f) = n.as_f64() { Value::Float(f) }
else { Value::Unit }
}
J::String(s) => Value::Str(s.as_str().into()),
J::Array(items) => Value::List(items.iter().map(json_to_value).collect()),
J::Object(map) => {
if let (Some(serde_json::Value::String(name)), Some(serde_json::Value::Array(args))) =
(map.get("$variant"), map.get("args"))
{
return Value::Variant {
name: name.clone(),
args: args.iter().map(json_to_value).collect(),
};
}
let mut out = indexmap::IndexMap::new();
for (k, v) in map { out.insert(k.clone(), json_to_value(v)); }
Value::Record(out)
}
}
}
impl Tracer for Recorder {
fn enter_call(&mut self, node_id: &str, name: &str, args: &[Value]) {
push_call_frame(&self.state, node_id, name, args);
}
fn enter_effect(&mut self, node_id: &str, kind: &str, op: &str, args: &[Value]) {
push_effect_frame(&self.state, node_id, kind, op, args);
}
fn exit_ok(&mut self, value: &Value) { exit_ok_frame(&self.state, value); }
fn exit_err(&mut self, message: &str) { exit_err_frame(&self.state, message); }
fn exit_call_tail(&mut self) { exit_tail_frame(&self.state); }
fn override_effect(&mut self, node_id: &str) -> Option<Value> {
lookup_override(&self.state, node_id)
}
}
impl Tracer for Handle {
fn enter_call(&mut self, node_id: &str, name: &str, args: &[Value]) {
push_call_frame(&self.state, node_id, name, args);
}
fn enter_effect(&mut self, node_id: &str, kind: &str, op: &str, args: &[Value]) {
push_effect_frame(&self.state, node_id, kind, op, args);
}
fn exit_ok(&mut self, value: &Value) { exit_ok_frame(&self.state, value); }
fn exit_err(&mut self, message: &str) { exit_err_frame(&self.state, message); }
fn exit_call_tail(&mut self) { exit_tail_frame(&self.state); }
fn override_effect(&mut self, node_id: &str) -> Option<Value> {
lookup_override(&self.state, node_id)
}
}
fn push_call_frame(state: &Mutex<RecorderState>, node_id: &str, name: &str, args: &[Value]) {
let mut st = state.lock().unwrap();
st.open.push(OpenFrame {
node: TraceNode {
node_id: node_id.to_string(),
kind: TraceNodeKind::Call,
target: name.to_string(),
input: values_to_json(args),
output: None,
error: None,
started_at: now_unix(),
ended_at: 0,
children: Vec::new(),
},
children: Vec::new(),
});
}
fn push_effect_frame(state: &Mutex<RecorderState>, node_id: &str, kind: &str, op: &str, args: &[Value]) {
let mut st = state.lock().unwrap();
st.open.push(OpenFrame {
node: TraceNode {
node_id: node_id.to_string(),
kind: TraceNodeKind::Effect,
target: format!("{kind}.{op}"),
input: values_to_json(args),
output: None,
error: None,
started_at: now_unix(),
ended_at: 0,
children: Vec::new(),
},
children: Vec::new(),
});
}
fn exit_ok_frame(state: &Mutex<RecorderState>, value: &Value) {
let mut st = state.lock().unwrap();
if let Some(mut frame) = st.open.pop() {
frame.node.ended_at = now_unix();
frame.node.output = Some(value_to_json(value));
frame.node.children = frame.children;
attach_completed(&mut st, frame.node);
}
}
fn exit_err_frame(state: &Mutex<RecorderState>, message: &str) {
let mut st = state.lock().unwrap();
if let Some(mut frame) = st.open.pop() {
frame.node.ended_at = now_unix();
frame.node.error = Some(message.to_string());
frame.node.children = frame.children;
attach_completed(&mut st, frame.node);
}
}
fn exit_tail_frame(state: &Mutex<RecorderState>) {
let mut st = state.lock().unwrap();
if let Some(mut frame) = st.open.pop() {
frame.node.ended_at = now_unix();
frame.node.output = Some(serde_json::Value::Null);
frame.node.children = frame.children;
attach_completed(&mut st, frame.node);
}
}
fn lookup_override(state: &Mutex<RecorderState>, node_id: &str) -> Option<Value> {
let st = state.lock().unwrap();
st.overrides.get(node_id).map(json_to_value)
}
fn attach_completed(st: &mut RecorderState, node: TraceNode) {
if let Some(parent) = st.open.last_mut() {
parent.children.push(node);
} else {
st.completed.push(node);
}
}