runtime

package
v0.37.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: May 16, 2026 License: MIT Imports: 12 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func Match added in v0.32.0

func Match(msg Msg, pattern Msg) bool

Match compares two messages and return true if they matches and false otherwise. Unlike Equal it compares only some aspects of the messages.

func Run added in v0.26.0

func Run(ctx context.Context, prog Program, registry map[string]FuncCreator) (int, error)

func Terminate added in v0.37.0

func Terminate(ctx context.Context, exitCode int)

Terminate requests graceful runtime termination with the provided process exit code.

func Uint8Index added in v0.34.0

func Uint8Index(idx int) uint8

Uint8Index validates idx and returns it as uint8 or panics.

Types

type ArrayInport added in v0.25.0

type ArrayInport struct {
	// contains filtered or unexported fields
}

func NewArrayInport added in v0.25.0

func NewArrayInport(
	tracer *Tracer,
	chans []<-chan OrderedMsg,
	addr PortAddr,
	interceptor Interceptor,
) *ArrayInport

func (ArrayInport) Len added in v0.25.0

func (a ArrayInport) Len() int

func (*ArrayInport) Receive added in v0.25.0

func (a *ArrayInport) Receive(ctx context.Context, idx int) (OrderedMsg, bool)

Receive receives a message from a specific array slot together with its runtime ordering metadata.

func (*ArrayInport) ReceiveAll added in v0.26.0

func (a *ArrayInport) ReceiveAll(ctx context.Context, f func(idx int, ordered OrderedMsg) bool) bool

ReceiveAll receives messages from all available array inport slots just once. It returns false if context is done or if the provided function returns false. The function is called for each message received. The function should return false if it wants to stop receiving messages. Functions receive full transport envelopes and are called in order of incoming messages, not in order of slots.

func (*ArrayInport) Select added in v0.26.0

func (a *ArrayInport) Select(ctx context.Context) (SelectedMsg, bool)

Select returns oldest available message across all available array inport slots.

type ArrayOutport added in v0.25.0

type ArrayOutport struct {
	// contains filtered or unexported fields
}

func NewArrayOutport added in v0.25.0

func NewArrayOutport(tracer *Tracer, addr PortAddr, interceptor Interceptor, slots []chan<- OrderedMsg) *ArrayOutport

func (*ArrayOutport) Len added in v0.25.0

func (a *ArrayOutport) Len() int

func (*ArrayOutport) Send added in v0.25.0

func (a *ArrayOutport) Send(ctx context.Context, idx uint8, msg Msg, causes ...OrderedMsg) bool

func (*ArrayOutport) SendAll added in v0.25.0

func (a *ArrayOutport) SendAll(ctx context.Context, msg Msg, causes ...OrderedMsg) bool

SendAllV2 sends the same message to all slots of the array outport. It returns false if context is done. It blocks until message is sent to all slots. Slots are not guaranteed to be handled in order, message is sent to first available slot. Each slot is guaranteed to be handled only once. TODO: figure out why this is the only working version of `SendAll`

type BoolMsg

type BoolMsg struct {
	// contains filtered or unexported fields
}

func NewBoolMsg

func NewBoolMsg(b bool) BoolMsg

func (BoolMsg) Bool

func (msg BoolMsg) Bool() bool

func (BoolMsg) Bytes added in v0.35.0

func (BoolMsg) Bytes() []byte

func (BoolMsg) Dict added in v0.26.0

func (BoolMsg) Dict() map[string]Msg

func (BoolMsg) Equal added in v0.26.0

func (msg BoolMsg) Equal(other Msg) bool

func (BoolMsg) Float

func (BoolMsg) Float() float64

func (BoolMsg) Int

func (BoolMsg) Int() int64

func (BoolMsg) List

func (BoolMsg) List() []Msg

func (BoolMsg) MarshalJSON added in v0.20.0

func (msg BoolMsg) MarshalJSON() ([]byte, error)

func (BoolMsg) Str

func (BoolMsg) Str() string

func (BoolMsg) String

func (msg BoolMsg) String() string

func (BoolMsg) Struct added in v0.26.0

func (BoolMsg) Struct() StructMsg

func (BoolMsg) Union added in v0.28.1

func (BoolMsg) Union() UnionMsg

type BytesMsg added in v0.35.0

type BytesMsg struct {
	// contains filtered or unexported fields
}

--- BYTES ---

func NewBytesMsg added in v0.35.0

func NewBytesMsg(v []byte) BytesMsg

func (BytesMsg) Bool added in v0.35.0

func (BytesMsg) Bool() bool

func (BytesMsg) Bytes added in v0.35.0

func (msg BytesMsg) Bytes() []byte

func (BytesMsg) Dict added in v0.35.0

func (BytesMsg) Dict() map[string]Msg

func (BytesMsg) Equal added in v0.35.0

func (msg BytesMsg) Equal(other Msg) bool

func (BytesMsg) Float added in v0.35.0

func (BytesMsg) Float() float64

func (BytesMsg) Int added in v0.35.0

func (BytesMsg) Int() int64

func (BytesMsg) List added in v0.35.0

func (BytesMsg) List() []Msg

func (BytesMsg) MarshalJSON added in v0.35.0

func (msg BytesMsg) MarshalJSON() ([]byte, error)

func (BytesMsg) Str added in v0.35.0

func (BytesMsg) Str() string

func (BytesMsg) String added in v0.35.0

func (msg BytesMsg) String() string

func (BytesMsg) Struct added in v0.35.0

func (BytesMsg) Struct() StructMsg

func (BytesMsg) Union added in v0.35.0

func (BytesMsg) Union() UnionMsg

type DictMsg added in v0.26.0

type DictMsg struct {
	// contains filtered or unexported fields
}

--- DICT ---

func NewDictMsg added in v0.26.0

func NewDictMsg(d map[string]Msg) DictMsg

func (DictMsg) Bool added in v0.26.0

func (DictMsg) Bool() bool

func (DictMsg) Bytes added in v0.35.0

func (DictMsg) Bytes() []byte

func (DictMsg) Dict added in v0.26.0

func (msg DictMsg) Dict() map[string]Msg

func (DictMsg) Equal added in v0.26.0

func (msg DictMsg) Equal(other Msg) bool

func (DictMsg) Float added in v0.26.0

func (DictMsg) Float() float64

func (DictMsg) Int added in v0.26.0

func (DictMsg) Int() int64

func (DictMsg) List added in v0.26.0

func (DictMsg) List() []Msg

func (DictMsg) MarshalJSON added in v0.26.0

func (msg DictMsg) MarshalJSON() ([]byte, error)

func (DictMsg) Str added in v0.26.0

func (DictMsg) Str() string

func (DictMsg) String added in v0.26.0

func (msg DictMsg) String() string

func (DictMsg) Struct added in v0.26.0

func (DictMsg) Struct() StructMsg

func (DictMsg) Union added in v0.28.1

func (DictMsg) Union() UnionMsg

type FloatMsg

type FloatMsg struct {
	// contains filtered or unexported fields
}

func NewFloatMsg

func NewFloatMsg(n float64) FloatMsg

func (FloatMsg) Bool

func (FloatMsg) Bool() bool

func (FloatMsg) Bytes added in v0.35.0

func (FloatMsg) Bytes() []byte

func (FloatMsg) Dict added in v0.26.0

func (FloatMsg) Dict() map[string]Msg

func (FloatMsg) Equal added in v0.26.0

func (msg FloatMsg) Equal(other Msg) bool

func (FloatMsg) Float

func (msg FloatMsg) Float() float64

func (FloatMsg) Int

func (FloatMsg) Int() int64

func (FloatMsg) List

func (FloatMsg) List() []Msg

func (FloatMsg) MarshalJSON added in v0.20.0

func (msg FloatMsg) MarshalJSON() ([]byte, error)

func (FloatMsg) Str

func (FloatMsg) Str() string

func (FloatMsg) String

func (msg FloatMsg) String() string

func (FloatMsg) Struct added in v0.26.0

func (FloatMsg) Struct() StructMsg

func (FloatMsg) Union added in v0.28.1

func (FloatMsg) Union() UnionMsg

type FuncCall

type FuncCall struct {
	Config Msg
	IO     IO
	Ref    string
}

type FuncCreator

type FuncCreator interface {
	Create(IO, Msg) (func(context.Context), error)
}

type IO added in v0.26.0

type IO struct {
	In  Inports
	Out Outports
}

type Inport added in v0.26.0

type Inport struct {
	// contains filtered or unexported fields
}

func NewInport added in v0.26.0

func NewInport(
	array *ArrayInport,
	single *SingleInport,
) Inport

func (Inport) Array added in v0.26.0

func (f Inport) Array() *ArrayInport

func (Inport) Single added in v0.26.0

func (f Inport) Single() *SingleInport

type Inports added in v0.26.0

type Inports struct {
	// contains filtered or unexported fields
}

func NewInports added in v0.26.0

func NewInports(ports map[string]Inport) Inports

func (Inports) Array added in v0.26.0

func (f Inports) Array(name string) (ArrayInport, error)

func (Inports) Ports added in v0.26.0

func (f Inports) Ports() map[string]Inport

func (Inports) Single added in v0.26.0

func (f Inports) Single(name string) (SingleInport, error)

type IntMsg

type IntMsg struct {
	// contains filtered or unexported fields
}

func NewIntMsg

func NewIntMsg(n int64) IntMsg

func (IntMsg) Bool

func (IntMsg) Bool() bool

func (IntMsg) Bytes added in v0.35.0

func (IntMsg) Bytes() []byte

func (IntMsg) Dict added in v0.26.0

func (IntMsg) Dict() map[string]Msg

func (IntMsg) Equal added in v0.26.0

func (msg IntMsg) Equal(other Msg) bool

func (IntMsg) Float

func (IntMsg) Float() float64

func (IntMsg) Int

func (msg IntMsg) Int() int64

func (IntMsg) List

func (IntMsg) List() []Msg

func (IntMsg) MarshalJSON added in v0.20.0

func (msg IntMsg) MarshalJSON() ([]byte, error)

func (IntMsg) Str

func (IntMsg) Str() string

func (IntMsg) String

func (msg IntMsg) String() string

func (IntMsg) Struct added in v0.26.0

func (IntMsg) Struct() StructMsg

func (IntMsg) Union added in v0.28.1

func (IntMsg) Union() UnionMsg

type Interceptor added in v0.25.0

type Interceptor interface {
	Sent(context.Context, PortSlotAddr, OrderedMsg, TraceHop)
	Received(context.Context, PortSlotAddr, OrderedMsg) OrderedMsg
}

func NewInterceptor added in v0.37.0

func NewInterceptor(tracePath, comment string) (Interceptor, func() error, error)

NewInterceptor always enables in-memory tracing. When tracePath is empty, it skips JSONL emission and returns the production interceptor.

type JSONLTraceFileWriter added in v0.37.0

type JSONLTraceFileWriter struct {
	// contains filtered or unexported fields
}

JSONLTraceFileWriter is an interceptor that writes each tracing event into a file in a JSONL format.

func NewDebugInterceptor added in v0.26.0

func NewDebugInterceptor(comment string) *JSONLTraceFileWriter

func (*JSONLTraceFileWriter) Open added in v0.37.0

func (d *JSONLTraceFileWriter) Open(filepath string) (func() error, error)

Open safely opens the file for writing and returns its close func.

func (*JSONLTraceFileWriter) Received added in v0.37.0

func (d *JSONLTraceFileWriter) Received(_ context.Context, receiver PortSlotAddr, ordered OrderedMsg) OrderedMsg

Received implements Interceptor interface. The only thing this interceptor really does - is writes JSONL line to a file.

func (*JSONLTraceFileWriter) Sent added in v0.37.0

func (d *JSONLTraceFileWriter) Sent(
	_ context.Context,
	sender PortSlotAddr,
	ordered OrderedMsg,
	hop TraceHop,
)

Sent implements Interceptor interface. The only thing this interceptor really does - is writes JSONL line to a file.

type ListMsg

type ListMsg struct {
	// contains filtered or unexported fields
}

--- LIST ---

func NewListMsg

func NewListMsg(v []Msg) ListMsg

func (ListMsg) Bool

func (ListMsg) Bool() bool

func (ListMsg) Bytes added in v0.35.0

func (ListMsg) Bytes() []byte

func (ListMsg) Dict added in v0.26.0

func (ListMsg) Dict() map[string]Msg

func (ListMsg) Equal added in v0.26.0

func (msg ListMsg) Equal(other Msg) bool

func (ListMsg) Float

func (ListMsg) Float() float64

func (ListMsg) Int

func (ListMsg) Int() int64

func (ListMsg) List

func (msg ListMsg) List() []Msg

func (ListMsg) MarshalJSON added in v0.20.0

func (msg ListMsg) MarshalJSON() ([]byte, error)

func (ListMsg) Str

func (ListMsg) Str() string

func (ListMsg) String

func (msg ListMsg) String() string

func (ListMsg) Struct added in v0.26.0

func (ListMsg) Struct() StructMsg

func (ListMsg) Union added in v0.28.1

func (ListMsg) Union() UnionMsg

type Msg

type Msg interface {
	Bool() bool
	Int() int64
	Float() float64
	Str() string
	Bytes() []byte
	List() []Msg
	Dict() map[string]Msg
	Struct() StructMsg
	Union() UnionMsg

	Equal(Msg) bool
}

func Call added in v0.33.0

func Call(ctx context.Context, prog Program, registry map[string]FuncCreator, input Msg) (Msg, int, error)

Call runs a single request-response round-trip using program Start/Stop. It sends the provided input to Start, waits for one message on Stop, then cancels and waits for all handlers to finish.

type NoEffectInterceptor added in v0.37.0

type NoEffectInterceptor struct{}

NoEffectInterceptor exist to be used as default interceptor when no actual interception is needed. Just to satisfy compiler.

func (NoEffectInterceptor) Prepare added in v0.37.0

func (NoEffectInterceptor) Prepare() error

func (NoEffectInterceptor) Received added in v0.37.0

func (NoEffectInterceptor) Sent added in v0.37.0

func (p NoEffectInterceptor) Sent(
	_ context.Context,
	_ PortSlotAddr,
	ordered OrderedMsg,
	_ TraceHop,
)

type OrderedMsg added in v0.26.0

type OrderedMsg struct {
	Msg
	// contains filtered or unexported fields
}

OrderedMsg is a transport envelope with payload and runtime ordering metadata.

func (OrderedMsg) String added in v0.26.0

func (o OrderedMsg) String() string

String is just a simple stringer that ignores index while formatting.

type Outport added in v0.26.0

type Outport struct {
	// contains filtered or unexported fields
}

func NewOutport added in v0.26.0

func NewOutport(
	single *SingleOutport,
	array *ArrayOutport,
) Outport

type Outports added in v0.26.0

type Outports struct {
	// contains filtered or unexported fields
}

func NewOutports added in v0.26.0

func NewOutports(ports map[string]Outport) Outports

func (Outports) Array added in v0.26.0

func (f Outports) Array(name string) (ArrayOutport, error)

func (Outports) Single added in v0.26.0

func (f Outports) Single(name string) (SingleOutport, error)

type PortAddr

type PortAddr struct {
	Path string `json:"Path"`
	Port string `json:"Port"`
}

type PortSlotAddr added in v0.26.0

type PortSlotAddr struct {
	Index *uint8 `json:",omitempty"` // nil means single port
	PortAddr
}

type Program

type Program struct {
	// for programmer start is inport and stop is outport, but for runtime it's inverted
	Start     *SingleOutport // Start must be inport of the first function
	Stop      *SingleInport  // Stop must be outport of the (one of the) terminator function(s)
	FuncCalls []FuncCall
}

type SelectedMsg added in v0.26.0

type SelectedMsg struct {
	OrderedMsg OrderedMsg
	SlotIdx    uint8
}

SelectedMsg is a message selected from available messages on all array inport slots.

func (SelectedMsg) String added in v0.26.0

func (s SelectedMsg) String() string

type SingleInport added in v0.25.0

type SingleInport struct {
	// contains filtered or unexported fields
}

func NewSingleInport added in v0.25.0

func NewSingleInport(
	tracer *Tracer,
	ch <-chan OrderedMsg,
	addr PortAddr,
	interceptor Interceptor,
) *SingleInport

func (SingleInport) Receive added in v0.25.0

func (s SingleInport) Receive(ctx context.Context) (OrderedMsg, bool)

Receive returns the next incoming transport envelope with its runtime ordering metadata.

type SingleOutport added in v0.25.0

type SingleOutport struct {
	// contains filtered or unexported fields
}

func NewSingleOutport added in v0.25.0

func NewSingleOutport(
	tracer *Tracer,
	addr PortAddr,
	interceptor Interceptor,
	outCh chan<- OrderedMsg,
) *SingleOutport

func (SingleOutport) Send added in v0.25.0

func (s SingleOutport) Send(ctx context.Context, msg Msg, causes ...OrderedMsg) bool

type StringMsg added in v0.26.0

type StringMsg struct {
	// contains filtered or unexported fields
}

--- STRING ---

func NewStringMsg added in v0.26.0

func NewStringMsg(s string) StringMsg

func (StringMsg) Bool added in v0.26.0

func (StringMsg) Bool() bool

func (StringMsg) Bytes added in v0.35.0

func (StringMsg) Bytes() []byte

func (StringMsg) Dict added in v0.26.0

func (StringMsg) Dict() map[string]Msg

func (StringMsg) Equal added in v0.26.0

func (msg StringMsg) Equal(other Msg) bool

func (StringMsg) Float added in v0.26.0

func (StringMsg) Float() float64

func (StringMsg) Int added in v0.26.0

func (StringMsg) Int() int64

func (StringMsg) List added in v0.26.0

func (StringMsg) List() []Msg

func (StringMsg) MarshalJSON added in v0.26.0

func (msg StringMsg) MarshalJSON() ([]byte, error)

func (StringMsg) Str added in v0.26.0

func (msg StringMsg) Str() string

func (StringMsg) String added in v0.26.0

func (msg StringMsg) String() string

func (StringMsg) Struct added in v0.26.0

func (StringMsg) Struct() StructMsg

func (StringMsg) Union added in v0.28.1

func (StringMsg) Union() UnionMsg

type StructField added in v0.33.0

type StructField struct {
	// contains filtered or unexported fields
}

structfield is a helper to construct structs via runtime.newstruct api without exposing fields.

func NewStructField added in v0.33.0

func NewStructField(name string, value Msg) StructField

newstructfield constructs a structfield with provided name and value.

type StructMsg added in v0.26.0

type StructMsg struct {
	// contains filtered or unexported fields
}

--- STRUCT ---

func NewStructMsg added in v0.26.0

func NewStructMsg(fields []StructField) StructMsg

newstruct builds a struct message from a slice of structfield. underlying struct representation remains unchanged for now.

func (StructMsg) Bool added in v0.26.0

func (StructMsg) Bool() bool

func (StructMsg) Bytes added in v0.35.0

func (StructMsg) Bytes() []byte

func (StructMsg) Dict added in v0.26.0

func (StructMsg) Dict() map[string]Msg

func (StructMsg) Equal added in v0.26.0

func (msg StructMsg) Equal(other Msg) bool

Equal implements strict equality for StructMsg messages. It returns false if the lengths of the names and fields are different. It returns false if any of the fields are not equal.

func (StructMsg) Float added in v0.26.0

func (StructMsg) Float() float64

func (StructMsg) Get added in v0.26.0

func (msg StructMsg) Get(name string) Msg

get returns the value of a field by name. it panics if the field is not found. it uses linear scan to find the field.

func (StructMsg) Int added in v0.26.0

func (StructMsg) Int() int64

func (StructMsg) List added in v0.26.0

func (StructMsg) List() []Msg

func (StructMsg) MarshalJSON added in v0.26.0

func (msg StructMsg) MarshalJSON() ([]byte, error)

func (StructMsg) Str added in v0.26.0

func (StructMsg) Str() string

func (StructMsg) String added in v0.26.0

func (msg StructMsg) String() string

func (StructMsg) Struct added in v0.26.0

func (msg StructMsg) Struct() StructMsg

func (StructMsg) Union added in v0.28.1

func (StructMsg) Union() UnionMsg

type TraceHop added in v0.37.0

type TraceHop struct {
	Sender   *PortSlotAddr
	Receiver *PortSlotAddr
	Message  string
	// CauseIndexes is the single source of truth for causal edges in traceStore.
	// Tree views are reconstructed from these indexes on demand.
	CauseIndexes []uint64
	// Index is materialized on read from traceStore map key.
	// We do not persist it in storage separately to avoid duplicated state.
	Index uint64
}

type Tracer added in v0.37.0

type Tracer struct {
	// contains filtered or unexported fields
}

func NewTracer added in v0.37.0

func NewTracer() *Tracer

func TracerFromIO added in v0.37.0

func TracerFromIO(runtimeIO IO) *Tracer

TracerFromIO returns the runtime tracer bound to this IO wiring.

This is a pragmatic bridge used by runtime funcs (for example runtime.Panic) to read the current dataflow trace from runtime state. It is intentionally wiring-based in the current implementation and may be redesigned later.

func (*Tracer) HopByOrderedMsg added in v0.37.0

func (t *Tracer) HopByOrderedMsg(ordered OrderedMsg) (TraceHop, bool)

HopByOrderedMsg returns a normalized hop for a concrete ordered message.

func (*Tracer) HopsByCauseIndexes added in v0.37.0

func (t *Tracer) HopsByCauseIndexes(causeIndexes []uint64) []TraceHop

HopsByCauseIndexes resolves parent hops by stored cause indexes. Missing indexes are ignored to keep trace-read path resilient.

type UnionMsg added in v0.28.1

type UnionMsg struct {
	// contains filtered or unexported fields
}

--- UNION ---

func AsUnion added in v0.37.0

func AsUnion(msg Msg) (UnionMsg, bool)

func NewUnionMsg added in v0.28.1

func NewUnionMsg(tag string, data Msg) UnionMsg

func (UnionMsg) Bool added in v0.28.1

func (UnionMsg) Bool() bool

func (UnionMsg) Bytes added in v0.35.0

func (UnionMsg) Bytes() []byte

func (UnionMsg) Data added in v0.32.0

func (msg UnionMsg) Data() Msg

func (UnionMsg) Dict added in v0.28.1

func (UnionMsg) Dict() map[string]Msg

func (UnionMsg) Equal added in v0.28.1

func (msg UnionMsg) Equal(other Msg) bool

Equal implements strict equality for UnionMsg messages. If one union has data and another doesn't, it returns false. It returns false if tags are different. It returns false if data is different. Tags are compared as Go strings and data is compared recursevely using Equal method.

func (UnionMsg) Float added in v0.28.1

func (UnionMsg) Float() float64

func (UnionMsg) Int added in v0.28.1

func (UnionMsg) Int() int64

func (UnionMsg) List added in v0.28.1

func (UnionMsg) List() []Msg

func (UnionMsg) MarshalJSON added in v0.35.0

func (msg UnionMsg) MarshalJSON() ([]byte, error)

func (UnionMsg) Str added in v0.28.1

func (UnionMsg) Str() string

func (UnionMsg) String added in v0.28.1

func (msg UnionMsg) String() string

func (UnionMsg) Struct added in v0.28.1

func (UnionMsg) Struct() StructMsg

func (UnionMsg) Tag added in v0.28.1

func (msg UnionMsg) Tag() string

func (UnionMsg) Union added in v0.28.1

func (msg UnionMsg) Union() UnionMsg

Directories

Path Synopsis
Package funcs implements low-level flows (runtime functions).
Package funcs implements low-level flows (runtime functions).

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL