aboutsummaryrefslogtreecommitdiff
path: root/state/memory_backend.go
diff options
context:
space:
mode:
Diffstat (limited to 'state/memory_backend.go')
-rw-r--r--state/memory_backend.go41
1 files changed, 41 insertions, 0 deletions
diff --git a/state/memory_backend.go b/state/memory_backend.go
new file mode 100644
index 0000000..58e86e0
--- /dev/null
+++ b/state/memory_backend.go
@@ -0,0 +1,41 @@
+package state
+
+import (
+ "context"
+ "goflink/core"
+)
+
+// TODO: currently using Golang Map for go through the DAG pipline
+// This should be use index map index map[uint64]uint32 Hash(KeyGroup + StateName + Key) -> Offset
+// with 1GB of []byte allocation
+// This should be optimize with ZERO GC Scanning, Memory Alignment
+
+type MemoryStateBackend struct {
+ maxParallelism int
+ currentKey string
+ currentGroupKey int
+ // 3D structure StateName -> KeyGroup -> Key -> Value
+ states map[string]map[int]map[string]any
+}
+
+func NewMemoryStateBackend(maxParallelism int) *MemoryStateBackend {
+ return &MemoryStateBackend{
+ maxParallelism: maxParallelism,
+ states: make(map[string]map[int]map[string]any),
+ }
+}
+
+func (ms *MemoryStateBackend) GetCurrentKey() string {
+ return ms.currentKey
+}
+
+func (ms *MemoryStateBackend) SetCurrentKey(key string) error {
+ ms.currentKey = key
+ ms.currentGroupKey = core.AssignToKeyGroup(key, ms.maxParallelism)
+ return nil
+}
+
+func (ms *MemoryStateBackend) GetValueState(ctx context.Context, key string) (core.ValueState[any], error) {
+ //TODO implement me
+ panic("implement me")
+}