Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-resilience

Resilience pipelines for Pudu: retry, circuit breaker, timeout, fallback, hedging, rate limiting, and chaos injection

0.1.0Apache-2.01

InstallClose

Pipeline.pudu

Pudu150 lines6.8 KB

GitHub ↗
1/** @Resilience.Pipeline.Backbone — composes strategies into one executable policy */2module PuduLangResilience.Pipeline34import Std.List as List5import Std.Option as Option6import PuduLangResilience.Clock as Clock7import PuduLangResilience.Constants.Events as Events8import PuduLangResilience.Constants.Messages as Messages9import PuduLangResilience.Context as Context10import PuduLangResilience.Randomizer as Randomizer11import PuduLangResilience as Resilience12import PuduLangResilience.Strategy as Strategy13import PuduLangResilience.Telemetry as Telemetry14import PuduLangResilience.Utils.Template as Template1516/** @Resilience.Pipeline.Options — identity, effects, and telemetry of a pipeline */17export type Options = {18  name: Option[Str],19  instance: Option[Str],20  clock: Clock.Clock,21  randomizer: Randomizer.Randomizer,22  listeners: Array[Telemetry.Listener],23  severityOf: Option[fn(Telemetry.Event) -> Telemetry.Severity]24}2526/** @Resilience.Pipeline.Pipeline — strategies composed outermost first */27export type Pipeline[T, E] = {28  name: Option[Str],29  instance: Option[Str],30  strategies: Array[Strategy.Strategy[T, E]],31  clock: Clock.Clock,32  telemetry: Telemetry.Reporter,33  run: fn(Context.Context, Strategy.Callback[T, E]) -> Resilience.Outcome[T, E]34}3536/** @Resilience.Pipeline.Invalid — every reason a pipeline could not be built */37export type Invalid = { problems: Array[Str] }3839/** @Resilience.Pipeline.Descriptor — what a built pipeline is made of */40export type Descriptor = { name: Option[Str], instance: Option[Str], strategies: Array[Described] }4142/** @Resilience.Pipeline.Described — one strategy of a built pipeline */43export type Described = { name: Option[Str], kind: Str, summary: Str }4445/// No name, the system clock and randomizer, and no listeners.46export fn defaults() -> Options {47  Options{name: None, instance: None, clock: Clock.system(), randomizer: Randomizer.system(), listeners: [], severityOf: None}48}4950/// A pipeline of `strategies` with the default options. The first strategy is the outermost.51export fn build[T, E](strategies: Array[Strategy.Strategy[T, E]]) -> Result[Pipeline[T, E], Invalid] {52  buildWith(&defaults(), strategies)53}5455/// A pipeline of `strategies` with the given options, or every problem found in the options of56/// its strategies and every strategy name used twice.57export fn buildWith[T, E](options: &Options, strategies: Array[Strategy.Strategy[T, E]]) -> Result[Pipeline[T, E], Invalid] {58  var problems: Array[Str] = []59  var seen: Array[Str] = []60  for layer in strategies {61    let label = Option.unwrapOr(layer.name, layer.kind)62    for problem in layer.problems { problems = problems.push(label + ": " + problem) }63    if let Some(name) = layer.name {64      if List.contains(&seen, name) { problems = problems.push("the strategy name '" + name + "' is used more than once") }65      seen = seen.push(name)66    }67  }68  if !problems.isEmpty() { return Err(Invalid{problems: problems}) }69  let base = Telemetry.reporter(Telemetry.Source{pipelineName: options.name, pipelineInstance: options.instance, strategyName: None}, &options.listeners, &options.severityOf)70  var composed = fn(context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] { callback(context) }71  var index = strategies.length() - 172  while index >= 0 {73    let layer = strategies[index]74    let runtime = Strategy.Runtime{clock: options.clock, randomizer: options.randomizer, telemetry: Telemetry.forStrategy(&base, layer.name)}75    (layer.attach)(runtime)76    let inner = composed77    composed = fn(context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {78      (layer.execute)(runtime, context, fn(passed: Context.Context) -> Resilience.Outcome[T, E] { inner(passed, callback) })79    }80    index = index - 181  }82  Ok(Pipeline{name: options.name, instance: options.instance, strategies: strategies, clock: options.clock, telemetry: base, run: composed})83}8485/// A pipeline with no strategies: it runs every callback once, as given.86export fn empty[T, E]() -> Pipeline[T, E] {87  match build([]) {88    case Ok(pipeline) => pipeline89    case Err(_) => panic("a pipeline without strategies has nothing to refuse")90  }91}9293/// Runs a callback through the pipeline with a fresh context.94export fn execute[T, E](pipeline: &Pipeline[T, E], callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {95  executeWith(pipeline, Context.create(), callback)96}9798/// Runs a callback through the pipeline with the given context, reporting the execution's start99/// and, with its duration and outcome, its end.100export fn executeWith[T, E](pipeline: &Pipeline[T, E], context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {101  Telemetry.report(&pipeline.telemetry, &context, Events.PIPELINE_EXECUTING, Telemetry.Debug, None, Telemetry.Executing)102  let started = Clock.now(&pipeline.clock)103  let outcome = (pipeline.run)(context, callback)104  let severity = match outcome {105    case Ok(_) => Telemetry.Information106    case Err(_) => Telemetry.Warning107  }108  let elapsed = Clock.now(&pipeline.clock) - started109  Telemetry.report(&pipeline.telemetry, &context, Events.PIPELINE_EXECUTED, severity, Some(Resilience.summarize(&outcome)), Telemetry.Executed(elapsed))110  outcome111}112113/// Runs a callback that answers a plain `Result` through the pipeline with a fresh context; its114/// error becomes `Raised`.115export fn run[T, E](pipeline: &Pipeline[T, E], action: fn(Context.Context) -> Result[T, E]) -> Resilience.Outcome[T, E] {116  runWith(pipeline, Context.create(), action)117}118119/// `run` with the given context.120export fn runWith[T, E](pipeline: &Pipeline[T, E], context: Context.Context, action: fn(Context.Context) -> Result[T, E]) -> Resilience.Outcome[T, E] {121  executeWith(pipeline, context, fn(passed: Context.Context) -> Resilience.Outcome[T, E] { Resilience.lift(action(passed)) })122}123124/// The pipeline as one strategy of another. It keeps its own clock and telemetry.125export fn asStrategy[T, E](pipeline: &Pipeline[T, E]) -> Strategy.Strategy[T, E] {126  let nested = *pipeline127  Strategy.Strategy {128    name: nested.name,129    kind: "Pipeline",130    summary: show(nested.strategies.length()) + " strategies",131    problems: [],132    attach: Strategy.detached,133    execute: fn(_runtime: Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {134      executeWith(&nested, context, callback)135    }136  }137}138139/// The name, instance, and strategies of a pipeline, outermost first.140export fn describe[T, E](pipeline: &Pipeline[T, E]) -> Descriptor {141  Descriptor {142    name: pipeline.name,143    instance: pipeline.instance,144    strategies: pipeline.strategies.map(fn(layer: Strategy.Strategy[T, E]) -> Described { Described{name: layer.name, kind: layer.kind, summary: layer.summary} })145  }146}147148/// One line per problem.149export fn explain(invalid: &Invalid) -> Str { Template.fill(Messages.PIPELINE_INVALID, &[invalid.problems.join(Messages.PROBLEM_SEPARATOR)]) }150