data

package
v2.0.0-alpha.57 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: Apache-2.0 Imports: 18 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

Functions

func DecodeAggregateKey

func DecodeAggregateKey(key runtime.Value) (runtime.Value, int, bool)

func NewDataSet

func NewDataSet(distinct bool) runtime.List

func NewKVIterator

func NewKVIterator(source runtime.Iterator) runtime.Iterator

func NewStreamGroupValue

func NewStreamGroupValue(streams []runtime.Stream) runtime.Value

NewStreamGroupValue creates one VM-owned value that closes every subscription.

func NewStreamValue

func NewStreamValue(stream runtime.Stream) runtime.Value

func ReduceAggregate

func ReduceAggregate(ctx context.Context, values runtime.List, kind bytecode.AggregateKind) (runtime.Value, error)

ReduceAggregate finalizes a collected list with the same semantics as native aggregate collectors. The list and its elements remain borrowed. The caller validates the aggregate kind while preparing the program.

func Sleep

func Sleep(ctx context.Context, duration runtime.Duration) error

func Stringify

func Stringify(ctx context.Context, input runtime.Value) (string, error)

Stringify converts a Value to a String. If the input is an Iterable, it concatenates

Types

type AggregateCollector

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

func (*AggregateCollector) Close

func (c *AggregateCollector) Close() error

func (*AggregateCollector) Copy

func (c *AggregateCollector) Copy() runtime.Value

func (*AggregateCollector) Get

func (*AggregateCollector) Hash

func (c *AggregateCollector) Hash() uint64

func (*AggregateCollector) Iterate

func (*AggregateCollector) Length

func (*AggregateCollector) MarshalJSON

func (c *AggregateCollector) MarshalJSON() ([]byte, error)

func (*AggregateCollector) Set

func (c *AggregateCollector) Set(ctx context.Context, key, value runtime.Value) error

func (*AggregateCollector) String

func (c *AggregateCollector) String() string

func (*AggregateCollector) UpdateAggregate

func (c *AggregateCollector) UpdateAggregate(idx int, value runtime.Value) error

type AggregateKey

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

func NewAggregateKey

func NewAggregateKey(groupKey runtime.Value, selectorIdx int) *AggregateKey

func (*AggregateKey) Copy

func (k *AggregateKey) Copy() runtime.Value

func (*AggregateKey) GroupKey

func (k *AggregateKey) GroupKey() runtime.Value

func (*AggregateKey) Hash

func (k *AggregateKey) Hash() uint64

func (*AggregateKey) SelectorIndex

func (k *AggregateKey) SelectorIndex() int

func (*AggregateKey) String

func (k *AggregateKey) String() string

func (*AggregateKey) VMUntracked

func (*AggregateKey) VMUntracked()

type ClosableIterator

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

func NewClosableIterator

func NewClosableIterator(src runtime.Iterator) *ClosableIterator

func (*ClosableIterator) Close

func (it *ClosableIterator) Close() error

func (*ClosableIterator) Copy

func (it *ClosableIterator) Copy() runtime.Value

func (*ClosableIterator) Hash

func (it *ClosableIterator) Hash() uint64

func (*ClosableIterator) Key

func (it *ClosableIterator) Key() runtime.Value

func (*ClosableIterator) MarshalJSON

func (it *ClosableIterator) MarshalJSON() ([]byte, error)

func (*ClosableIterator) Next

func (it *ClosableIterator) Next(ctx context.Context) error

func (*ClosableIterator) String

func (it *ClosableIterator) String() string

func (*ClosableIterator) Value

func (it *ClosableIterator) Value() runtime.Value

type CounterCollector

type CounterCollector struct {
	*runtime.Box[runtime.Int]
}

CounterCollector is a Transformer implementation that tracks and increments a counter of processed values.

func (*CounterCollector) Close

func (c *CounterCollector) Close() error

func (*CounterCollector) Get

func (*CounterCollector) Increment

func (c *CounterCollector) Increment()

func (*CounterCollector) Iterate

func (*CounterCollector) Length

func (c *CounterCollector) Length(_ context.Context) (runtime.Int, error)

func (*CounterCollector) Set

type DataSet

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

func (*DataSet) Append

func (ds *DataSet) Append(ctx context.Context, item runtime.Value) error

func (*DataSet) At

func (ds *DataSet) At(ctx context.Context, idx runtime.Int) (runtime.Value, error)

func (*DataSet) Clear

func (ds *DataSet) Clear(ctx context.Context) error

func (*DataSet) Clone

func (ds *DataSet) Clone(ctx context.Context) (runtime.Cloneable, error)

func (*DataSet) Compare

func (ds *DataSet) Compare(ctx context.Context, other runtime.Value) (runtime.Ordering, error)

func (*DataSet) Concat

func (ds *DataSet) Concat(ctx context.Context, other runtime.List) error

func (*DataSet) Contains

func (ds *DataSet) Contains(ctx context.Context, value runtime.Value) (runtime.Boolean, error)

func (*DataSet) Copy

func (ds *DataSet) Copy() runtime.Value

func (*DataSet) Equal

func (ds *DataSet) Equal(ctx context.Context, other runtime.Value) (bool, error)

func (*DataSet) Filter

func (ds *DataSet) Filter(ctx context.Context, predicate runtime.IndexReadablePredicate) (runtime.List, error)

func (*DataSet) Find

func (*DataSet) First

func (ds *DataSet) First(ctx context.Context) (runtime.Value, error)

func (*DataSet) ForEach

func (ds *DataSet) ForEach(ctx context.Context, predicate runtime.IndexReadablePredicate) error

func (*DataSet) Hash

func (ds *DataSet) Hash() uint64

func (*DataSet) IndexOf

func (ds *DataSet) IndexOf(ctx context.Context, value runtime.Value) (runtime.Int, error)

func (*DataSet) Insert

func (ds *DataSet) Insert(ctx context.Context, idx runtime.Int, value runtime.Value) error

func (*DataSet) Iterate

func (ds *DataSet) Iterate(ctx context.Context) (runtime.Iterator, error)

func (*DataSet) Last

func (ds *DataSet) Last(ctx context.Context) (runtime.Value, error)

func (*DataSet) Length

func (ds *DataSet) Length(ctx context.Context) (runtime.Int, error)

func (*DataSet) LookupAt

func (ds *DataSet) LookupAt(ctx context.Context, index runtime.Int) (runtime.Value, bool, error)

func (*DataSet) New

func (ds *DataSet) New(ctx context.Context) (runtime.List, error)

New creates an independent empty collection with the receiver's configuration.

func (*DataSet) Remove

func (ds *DataSet) Remove(ctx context.Context, value runtime.Value) error

func (*DataSet) RemoveAt

func (ds *DataSet) RemoveAt(ctx context.Context, idx runtime.Int) (runtime.Value, error)

func (*DataSet) SetAt

func (ds *DataSet) SetAt(ctx context.Context, idx runtime.Int, value runtime.Value) error

func (*DataSet) Slice

func (ds *DataSet) Slice(ctx context.Context, start, end runtime.Int) (runtime.List, error)

func (*DataSet) SortAsc

func (ds *DataSet) SortAsc(ctx context.Context) error

func (*DataSet) SortDesc

func (ds *DataSet) SortDesc(ctx context.Context) error

func (*DataSet) String

func (ds *DataSet) String() string

func (*DataSet) Swap

func (ds *DataSet) Swap(ctx context.Context, a, b runtime.Int) error

func (*DataSet) VMUntracked

func (*DataSet) VMUntracked()

type FastObject

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

func NewFastObject

func NewFastObject(cache *ShapeCache, dictThreshold int) *FastObject

func NewFastObjectOf

func NewFastObjectOf(cache *ShapeCache, dictThreshold int, size int) *FastObject

func (*FastObject) Clear

func (t *FastObject) Clear(_ context.Context) error

func (*FastObject) Clone

func (t *FastObject) Clone(ctx context.Context) (runtime.Cloneable, error)

func (*FastObject) Compare

func (t *FastObject) Compare(ctx context.Context, other runtime.Value) (runtime.Ordering, error)

func (*FastObject) Contains

func (t *FastObject) Contains(ctx context.Context, target runtime.Value) (runtime.Boolean, error)

func (*FastObject) ContainsKey

func (t *FastObject) ContainsKey(_ context.Context, key runtime.Value) (runtime.Boolean, error)

func (*FastObject) ContainsValue

func (t *FastObject) ContainsValue(ctx context.Context, target runtime.Value) (runtime.Boolean, error)

func (*FastObject) Copy

func (t *FastObject) Copy() runtime.Value

func (*FastObject) Equal

func (t *FastObject) Equal(ctx context.Context, other runtime.Value) (bool, error)

func (*FastObject) Filter

func (t *FastObject) Filter(ctx context.Context, predicate runtime.KeyReadablePredicate) (runtime.List, error)

func (*FastObject) Find

func (*FastObject) ForEach

func (t *FastObject) ForEach(ctx context.Context, predicate runtime.KeyReadablePredicate) error

func (*FastObject) Get

func (*FastObject) Hash

func (t *FastObject) Hash() uint64

func (*FastObject) IsEmpty

func (t *FastObject) IsEmpty(_ context.Context) (runtime.Boolean, error)

func (*FastObject) Iterate

func (t *FastObject) Iterate(_ context.Context) (runtime.Iterator, error)

func (*FastObject) Keys

func (t *FastObject) Keys(_ context.Context) (runtime.List, error)

func (*FastObject) Length

func (t *FastObject) Length(_ context.Context) (runtime.Int, error)

func (*FastObject) Lookup

func (t *FastObject) Lookup(_ context.Context, key runtime.Value) (runtime.Value, bool, error)

func (*FastObject) LookupSlot

func (t *FastObject) LookupSlot(key string) (int, bool)

func (*FastObject) MarshalJSON

func (t *FastObject) MarshalJSON() ([]byte, error)

func (*FastObject) Merge

func (t *FastObject) Merge(ctx context.Context, other runtime.Map) error

func (*FastObject) New

func (t *FastObject) New(ctx context.Context) (runtime.Map, error)

New creates an independent empty collection with the receiver's configuration.

func (*FastObject) ObjectLike

func (t *FastObject) ObjectLike()

func (*FastObject) Remove

func (t *FastObject) Remove(ctx context.Context, value runtime.Value) error

func (*FastObject) RemoveKey

func (t *FastObject) RemoveKey(_ context.Context, key runtime.Value) error

func (*FastObject) Set

func (t *FastObject) Set(_ context.Context, key runtime.Value, value runtime.Value) error

func (*FastObject) SetSlotWithShape

func (t *FastObject) SetSlotWithShape(next *Shape, slot int, value runtime.Value) bool

func (*FastObject) SetStringCached

func (t *FastObject) SetStringCached(key string, value runtime.Value) (*Shape, *Shape, int, bool)

func (*FastObject) Shape

func (t *FastObject) Shape() *Shape

func (*FastObject) ShapeID

func (t *FastObject) ShapeID() uint64

func (*FastObject) SlotValue

func (t *FastObject) SlotValue(slot int) (runtime.Value, bool)

func (*FastObject) String

func (t *FastObject) String() string

func (*FastObject) Type

func (t *FastObject) Type() runtime.Type

func (*FastObject) VMUntracked

func (*FastObject) VMUntracked()

func (*FastObject) Values

func (t *FastObject) Values(_ context.Context) (runtime.List, error)

type GroupedAggregateCollector

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

func (*GroupedAggregateCollector) Close

func (c *GroupedAggregateCollector) Close() error

func (*GroupedAggregateCollector) Copy

func (*GroupedAggregateCollector) Get

func (*GroupedAggregateCollector) Hash

func (*GroupedAggregateCollector) Iterate

func (*GroupedAggregateCollector) Length

func (*GroupedAggregateCollector) MarshalJSON

func (c *GroupedAggregateCollector) MarshalJSON() ([]byte, error)

func (*GroupedAggregateCollector) Set

func (c *GroupedAggregateCollector) Set(ctx context.Context, key, value runtime.Value) error

func (*GroupedAggregateCollector) String

func (c *GroupedAggregateCollector) String() string

func (*GroupedAggregateCollector) UpdateAggregate

func (c *GroupedAggregateCollector) UpdateAggregate(ctx context.Context, groupKey, value runtime.Value, idx int) error

type Iterator

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

func NewIterator

func NewIterator(src runtime.Iterator) *Iterator

func (*Iterator) Copy

func (it *Iterator) Copy() runtime.Value

func (*Iterator) Hash

func (it *Iterator) Hash() uint64

func (*Iterator) Key

func (it *Iterator) Key() runtime.Value

func (*Iterator) MarshalJSON

func (it *Iterator) MarshalJSON() ([]byte, error)

func (*Iterator) Next

func (it *Iterator) Next(ctx context.Context) error

func (*Iterator) String

func (it *Iterator) String() string

func (*Iterator) VMUntracked

func (*Iterator) VMUntracked()

func (*Iterator) Value

func (it *Iterator) Value() runtime.Value

type IteratorState

type IteratorState interface {
	runtime.Value
	Next(context.Context) error
	Value() runtime.Value
	Key() runtime.Value
}

func WrapIterator

func WrapIterator(src runtime.Iterator) IteratorState

type KV

type KV struct {
	Key   runtime.Value
	Value runtime.Value
}

KV represents a key-value pair where both the key and value are of type runtime.Value.

func NewKV

func NewKV(key, value runtime.Value) *KV

NewKV creates and returns a new KV instance with the provided key and value.

func (*KV) Copy

func (p *KV) Copy() runtime.Value

func (*KV) Hash

func (p *KV) Hash() uint64

func (*KV) String

func (p *KV) String() string

func (*KV) VMUntracked

func (*KV) VMUntracked()

type KVIterator

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

func (*KVIterator) Next

func (iter *KVIterator) Next(ctx context.Context) (runtime.Value, runtime.Value, error)

type KeyCollector

type KeyCollector struct {
	*runtime.Box[runtime.List]
	// contains filtered or unexported fields
}

KeyCollector collects semantically unique keys and sorts them in ascending order when iterated.

func (*KeyCollector) Close

func (c *KeyCollector) Close() error

func (*KeyCollector) Get

func (*KeyCollector) Iterate

func (c *KeyCollector) Iterate(ctx context.Context) (runtime.Iterator, error)

func (*KeyCollector) Length

func (c *KeyCollector) Length(ctx context.Context) (runtime.Int, error)

func (*KeyCollector) Set

func (c *KeyCollector) Set(ctx context.Context, key, _ runtime.Value) error

type KeyCounterCollector

type KeyCounterCollector struct {
	*runtime.Box[runtime.List]
	// contains filtered or unexported fields
}

func (*KeyCounterCollector) Close

func (c *KeyCounterCollector) Close() error

func (*KeyCounterCollector) Get

func (*KeyCounterCollector) Iterate

func (*KeyCounterCollector) Length

func (c *KeyCounterCollector) Length(ctx context.Context) (runtime.Int, error)

func (*KeyCounterCollector) Set

func (c *KeyCounterCollector) Set(ctx context.Context, key, _ runtime.Value) error

type KeyGroupCollector

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

func (*KeyGroupCollector) Close

func (c *KeyGroupCollector) Close() error

func (*KeyGroupCollector) Copy

func (c *KeyGroupCollector) Copy() runtime.Value

func (*KeyGroupCollector) Get

func (*KeyGroupCollector) Hash

func (c *KeyGroupCollector) Hash() uint64

func (*KeyGroupCollector) Iterate

func (*KeyGroupCollector) Length

func (*KeyGroupCollector) Set

func (c *KeyGroupCollector) Set(ctx context.Context, key, value runtime.Value) error

func (*KeyGroupCollector) String

func (c *KeyGroupCollector) String() string

type MultiSorter

type MultiSorter struct {
	*runtime.Box[runtime.List]
	// contains filtered or unexported fields
}

func (*MultiSorter) Close

func (s *MultiSorter) Close() error

func (*MultiSorter) Get

func (*MultiSorter) Iterate

func (s *MultiSorter) Iterate(ctx context.Context) (runtime.Iterator, error)

func (*MultiSorter) Length

func (s *MultiSorter) Length(ctx context.Context) (runtime.Int, error)

func (*MultiSorter) Set

func (s *MultiSorter) Set(ctx context.Context, key, value runtime.Value) error

type Shape

type Shape = fastShape

type ShapeCache

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

func NewShapeCache

func NewShapeCache(limit int) *ShapeCache

func (*ShapeCache) Root

func (c *ShapeCache) Root() *Shape

func (*ShapeCache) Transition

func (c *ShapeCache) Transition(shape *Shape, key string) *Shape

type Sorter

type Sorter struct {
	*runtime.Box[runtime.List]
	// contains filtered or unexported fields
}

func (*Sorter) Close

func (s *Sorter) Close() error

func (*Sorter) Get

func (*Sorter) Iterate

func (s *Sorter) Iterate(ctx context.Context) (runtime.Iterator, error)

func (*Sorter) Length

func (s *Sorter) Length(ctx context.Context) (runtime.Int, error)

func (*Sorter) Set

func (s *Sorter) Set(ctx context.Context, key, value runtime.Value) error

type StreamGroupController

type StreamGroupController interface {
	runtime.Value
	ArmDone(index runtime.Int) error
}

StreamGroupController exposes completion only to stream-group bytecode.

type StreamGroupIterator

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

StreamGroupIterator fans in active arms under one deadline and reports declaration indexes as keys. An exhausted arm is reported once with a None key and its declaration index as the value.

func (*StreamGroupIterator) ArmDone

func (it *StreamGroupIterator) ArmDone(index runtime.Int) error

ArmDone prevents queued messages or errors from a satisfied arm from affecting the group.

func (*StreamGroupIterator) Close

func (it *StreamGroupIterator) Close() error

func (*StreamGroupIterator) Copy

func (it *StreamGroupIterator) Copy() runtime.Value

func (*StreamGroupIterator) Hash

func (it *StreamGroupIterator) Hash() uint64

func (*StreamGroupIterator) Key

func (it *StreamGroupIterator) Key() runtime.Value

func (*StreamGroupIterator) Next

func (it *StreamGroupIterator) Next(ctx context.Context) error

func (*StreamGroupIterator) String

func (it *StreamGroupIterator) String() string

func (*StreamGroupIterator) Value

func (it *StreamGroupIterator) Value() runtime.Value

type StreamGroupValue

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

StreamGroupValue owns the subscriptions established for one grouped wait.

func (*StreamGroupValue) Close

func (v *StreamGroupValue) Close() error

func (*StreamGroupValue) Copy

func (v *StreamGroupValue) Copy() runtime.Value

func (*StreamGroupValue) Hash

func (v *StreamGroupValue) Hash() uint64

func (*StreamGroupValue) Iterate

func (v *StreamGroupValue) Iterate(timeout runtime.Duration) IteratorState

func (*StreamGroupValue) String

func (v *StreamGroupValue) String() string

type StreamIterable

type StreamIterable interface {
	runtime.Value
	Iterate(timeout runtime.Duration) IteratorState
}

StreamIterable is the VM-internal iterator factory shared by singular and grouped streams.

type StreamValue

type StreamValue struct {
	*runtime.Box[runtime.Stream]
}

func (*StreamValue) Close

func (v *StreamValue) Close() error

func (*StreamValue) Iterate

func (v *StreamValue) Iterate(timeout runtime.Duration) IteratorState

type Transformer

func NewAggregateCollector

func NewAggregateCollector(plan bytecode.AggregatePlan) Transformer

func NewCollector

func NewCollector(typ bytecode.CollectorType) Transformer

func NewCollectorSafe

func NewCollectorSafe(typ bytecode.CollectorType) (Transformer, error)

func NewCounterCollector

func NewCounterCollector() Transformer

func NewGroupedAggregateCollector

func NewGroupedAggregateCollector(plan bytecode.AggregatePlan) Transformer

func NewKeyCollector

func NewKeyCollector() Transformer

func NewKeyCounterCollector

func NewKeyCounterCollector() Transformer

func NewKeyGroupCollector

func NewKeyGroupCollector() Transformer

func NewMultiSorter

func NewMultiSorter(directions []runtime.SortDirection) Transformer

func NewNoopCollector

func NewNoopCollector() Transformer

func NewSorter

func NewSorter(direction runtime.SortDirection) Transformer

Jump to

Keyboard shortcuts

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