
Pipeline.pudu
Pudu150 lines6.8 KB
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 }444546export fn defaults() -> Options {47 Options{name: None, instance: None, clock: Clock.system(), randomizer: Randomizer.system(), listeners: [], severityOf: None}48}495051export fn build[T, E](strategies: Array[Strategy.Strategy[T, E]]) -> Result[Pipeline[T, E], Invalid] {52 buildWith(&defaults(), strategies)53}54555657export 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}848586export 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}929394export fn execute[T, E](pipeline: &Pipeline[T, E], callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {95 executeWith(pipeline, Context.create(), callback)96}979899100export 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}112113114115export 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}118119120export 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}123124125export 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}138139140export 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}147148149export fn explain(invalid: &Invalid) -> Str { Template.fill(Messages.PIPELINE_INVALID, &[invalid.problems.join(Messages.PROBLEM_SEPARATOR)]) }150