///|
/// Structured observability for production de-identification runs. Events are
/// append-only and contain checksums instead of source payloads, which makes
/// them suitable for logs, traces, and operational dashboards.
pub(all) enum PipelineEventKind {
PipelineRunStarted
PipelineStageStarted
PipelineStageFinished
PipelineRunFinished
PipelineRunFailed
PipelinePolicyLoaded
PipelineReviewOpened
PipelineCheckpointSaved
} derive(Debug, Eq)
///|
pub(all) enum PipelineStageStatus {
PipelineStagePending
PipelineStageActive
PipelineStageSucceeded
PipelineStageFailed
PipelineStageSkipped
} derive(Debug, Eq)
///|
pub(all) struct PipelineEvent {
event_id : String
run_id : String
document_id : String
stage : String
kind : PipelineEventKind
status : PipelineStageStatus
timestamp : String
duration_ms : Int
input_count : Int
output_count : Int
error_code : String
detail : String
payload_checksum : String
} derive(Debug, Eq)
///|
pub(all) struct PipelineStageMetric {
stage : String
status : PipelineStageStatus
started_at : String
finished_at : String
duration_ms : Int
input_count : Int
output_count : Int
error_count : Int
checksum : String
} derive(Debug, Eq)
///|
pub(all) struct ObservableTrace {
run_id : String
document_id : String
mut events : Array[PipelineEvent]
mut stages : Array[PipelineStageMetric]
started_at : String
mut finished_at : String
mut checksum : String
} derive(Debug)
///|
pub fn pipeline_event_kind_name(kind : PipelineEventKind) -> String {
match kind {
PipelineRunStarted => "run_started"
PipelineStageStarted => "stage_started"
PipelineStageFinished => "stage_finished"
PipelineRunFinished => "run_finished"
PipelineRunFailed => "run_failed"
PipelinePolicyLoaded => "policy_loaded"
PipelineReviewOpened => "review_opened"
PipelineCheckpointSaved => "checkpoint_saved"
}
}
///|
pub fn pipeline_stage_status_name(status : PipelineStageStatus) -> String {
match status {
PipelineStagePending => "pending"
PipelineStageActive => "active"
PipelineStageSucceeded => "succeeded"
PipelineStageFailed => "failed"
PipelineStageSkipped => "skipped"
}
}
///|
pub fn pipeline_event(
run_id : String,
document_id : String,
stage : String,
kind : PipelineEventKind,
status : PipelineStageStatus,
timestamp : String,
) -> PipelineEvent {
{
event_id: stable_hash(
run_id + ":" + document_id + ":" + stage + ":" + timestamp,
),
run_id,
document_id,
stage,
kind,
status,
timestamp,
duration_ms: 0,
input_count: 0,
output_count: 0,
error_code: "",
detail: "",
payload_checksum: "",
}
}
///|
pub fn PipelineEvent::with_counts(
event : PipelineEvent,
input_count : Int,
output_count : Int,
) -> PipelineEvent {
{
..event,
input_count: if input_count < 0 {
0
} else {
input_count
},
output_count: if output_count < 0 {
0
} else {
output_count
},
}
}
///|
pub fn PipelineEvent::with_duration(
event : PipelineEvent,
duration_ms : Int,
) -> PipelineEvent {
{ ..event, duration_ms: if duration_ms < 0 { 0 } else { duration_ms } }
}
///|
pub fn PipelineEvent::with_detail(
event : PipelineEvent,
detail : String,
) -> PipelineEvent {
{ ..event, detail, }
}
///|
pub fn PipelineEvent::with_error(
event : PipelineEvent,
error_code : String,
detail : String,
) -> PipelineEvent {
{ ..event, error_code, detail, status: PipelineStageFailed }
}
///|
pub fn PipelineEvent::with_payload_checksum(
event : PipelineEvent,
checksum : String,
) -> PipelineEvent {
{ ..event, payload_checksum: checksum }
}
///|
pub fn pipeline_stage_metric(
stage : String,
started_at : String,
) -> PipelineStageMetric {
{
stage,
status: PipelineStagePending,
started_at,
finished_at: "",
duration_ms: 0,
input_count: 0,
output_count: 0,
error_count: 0,
checksum: "",
}
}
///|
pub fn PipelineStageMetric::start(
self : PipelineStageMetric,
) -> PipelineStageMetric {
{ ..self, status: PipelineStageActive }
}
///|
pub fn PipelineStageMetric::finish(
self : PipelineStageMetric,
status : PipelineStageStatus,
finished_at : String,
duration_ms : Int,
input_count : Int,
output_count : Int,
) -> PipelineStageMetric {
{
..self,
status,
finished_at,
duration_ms: if duration_ms < 0 {
0
} else {
duration_ms
},
input_count: if input_count < 0 {
0
} else {
input_count
},
output_count: if output_count < 0 {
0
} else {
output_count
},
checksum: stable_hash(self.stage + ":" + finished_at),
}
}
///|
pub fn PipelineStageMetric::fail(
self : PipelineStageMetric,
finished_at : String,
duration_ms : Int,
) -> PipelineStageMetric {
{
..self.finish(
PipelineStageFailed,
finished_at,
duration_ms,
self.input_count,
self.output_count,
),
error_count: 1,
}
}
///|
pub fn pipeline_trace(
run_id : String,
document_id : String,
started_at : String,
) -> ObservableTrace {
{
run_id,
document_id,
events: [],
stages: [],
started_at,
finished_at: "",
checksum: stable_hash(run_id + ":" + document_id + ":" + started_at),
}
}
///|
pub fn ObservableTrace::record(
self : ObservableTrace,
event : PipelineEvent,
) -> String {
self.events.push(event)
self.checksum = stable_hash(
self.events.map(fn(item) { item.event_id }).join("\n"),
)
self.checksum
}
///|
pub fn ObservableTrace::add_stage(
self : ObservableTrace,
stage : PipelineStageMetric,
) -> Bool {
if stage.stage.is_empty() {
false
} else {
self.stages.push(stage)
true
}
}
///|
pub fn ObservableTrace::finish(
self : ObservableTrace,
finished_at : String,
) -> ObservableTrace {
self.finished_at = finished_at
self
}
///|
pub fn ObservableTrace::stage(
self : ObservableTrace,
name : String,
) -> PipelineStageMetric? {
let mut result : PipelineStageMetric? = None
for stage in self.stages {
if stage.stage == name {
result = Some(stage)
}
}
result
}
///|
pub fn ObservableTrace::has_failures(self : ObservableTrace) -> Bool {
self.events.any(fn(event) { event.status == PipelineStageFailed }) ||
self.stages.any(fn(stage) { stage.status == PipelineStageFailed })
}
///|
pub fn ObservableTrace::duration(self : ObservableTrace) -> Int {
self.stages.fold(init=0, (sum, stage) => sum + stage.duration_ms)
}
///|
pub fn ObservableTrace::event_count(self : ObservableTrace) -> Int {
self.events.length()
}
///|
pub fn ObservableTrace::stage_count(self : ObservableTrace) -> Int {
self.stages.length()
}
///|
pub fn ObservableTrace::success_rate(self : ObservableTrace) -> Int {
let completed = self.stages
.filter(fn(stage) {
stage.status == PipelineStageSucceeded ||
stage.status == PipelineStageSkipped
})
.length()
if self.stages.is_empty() {
100
} else {
completed * 100 / self.stages.length()
}
}
///|
pub fn ObservableTrace::is_complete(self : ObservableTrace) -> Bool {
!self.finished_at.is_empty() &&
self.stages.all(fn(stage) {
stage.status == PipelineStageSucceeded ||
stage.status == PipelineStageFailed ||
stage.status == PipelineStageSkipped
})
}
///|
pub fn ObservableTrace::summary(self : ObservableTrace) -> String {
[
"run_id=" + self.run_id,
"document_id=" + self.document_id,
"events=" + self.events.length().to_string(),
"stages=" + self.stages.length().to_string(),
"duration_ms=" + self.duration().to_string(),
"success_rate=" + self.success_rate().to_string(),
"failures=" + self.has_failures().to_string(),
"complete=" + self.is_complete().to_string(),
"checksum=" + self.checksum,
].join("\n")
}
///|
pub fn ObservableTrace::to_json(self : ObservableTrace) -> String {
"{" +
"\"run_id\":\"" +
json_escape(self.run_id) +
"\"," +
"\"document_id\":\"" +
json_escape(self.document_id) +
"\"," +
"\"started_at\":\"" +
json_escape(self.started_at) +
"\"," +
"\"finished_at\":\"" +
json_escape(self.finished_at) +
"\"," +
"\"events\":" +
self.events.length().to_string() +
"," +
"\"stages\":" +
self.stages.length().to_string() +
"," +
"\"duration_ms\":" +
self.duration().to_string() +
"," +
"\"success_rate\":" +
self.success_rate().to_string() +
"," +
"\"checksum\":\"" +
json_escape(self.checksum) +
"\"}"
}
///|
pub fn pipeline_trace_event_counts(trace : ObservableTrace) -> Map[String, Int] {
let counts : Map[String, Int] = Map([])
for event in trace.events {
let key = pipeline_event_kind_name(event.kind)
counts[key] = counts.get_or_default(key, 0) + 1
}
counts
}
///|
pub fn pipeline_trace_stage_counts(trace : ObservableTrace) -> Map[String, Int] {
let counts : Map[String, Int] = Map([])
for stage in trace.stages {
let key = pipeline_stage_status_name(stage.status)
counts[key] = counts.get_or_default(key, 0) + 1
}
counts
}
///|
pub fn pipeline_trace_payload_safe(trace : ObservableTrace) -> Bool {
trace.events.all(fn(event) {
event.detail.length() == 0 ||
(!event.detail.contains("original=") && !event.detail.contains("raw_text="))
})
}
///|
pub fn pipeline_trace_is_deterministic(
left : ObservableTrace,
right : ObservableTrace,
) -> Bool {
left.run_id == right.run_id &&
left.document_id == right.document_id &&
left.events.map(fn(item) { item.event_id }).join("\n") ==
right.events.map(fn(item) { item.event_id }).join("\n")
}
///|
pub fn pipeline_trace_export_lines(trace : ObservableTrace) -> Array[String] {
trace.events.map(fn(event) {
[
event.event_id,
event.run_id,
event.document_id,
event.stage,
pipeline_event_kind_name(event.kind),
pipeline_stage_status_name(event.status),
event.timestamp,
event.duration_ms.to_string(),
event.input_count.to_string(),
event.output_count.to_string(),
event.error_code,
event.payload_checksum,
]
.map(json_escape)
.join("\t")
})
}
///|
pub fn pipeline_trace_export_checksum(trace : ObservableTrace) -> String {
stable_hash(pipeline_trace_export_lines(trace).join("\n"))
}
///|
pub fn pipeline_trace_document_count(trace : ObservableTrace) -> Int {
if trace.document_id.is_empty() {
0
} else {
1
}
}
///|
pub fn pipeline_stage_throughput(stage : PipelineStageMetric) -> Int {
if stage.duration_ms <= 0 {
stage.output_count
} else {
stage.output_count * 1000 / stage.duration_ms
}
}
///|
pub fn pipeline_stage_loss(stage : PipelineStageMetric) -> Int {
if stage.input_count <= stage.output_count {
0
} else {
stage.input_count - stage.output_count
}
}
///|
pub fn pipeline_trace_stage_loss(trace : ObservableTrace) -> Int {
trace.stages.fold(init=0, (sum, stage) => sum + pipeline_stage_loss(stage))
}