From 50753a8c96fd56c68a8390699fc67e40e75abe3b Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Mon, 4 May 2026 01:42:14 +0900 Subject: [PATCH] wip: knotmirror/stream: diskpersist Signed-off-by: Seongmin Lee --- go.mod | 1 + go.sum | 2 + .../stream/persist/diskpersist/diskpersist.go | 104 +++++++++++++++++- 3 files changed, 105 insertions(+), 2 deletions(-) diff --git a/go.mod b/go.mod index 3394bd28..0df21629 100644 --- a/go.mod +++ b/go.mod @@ -204,6 +204,7 @@ require ( github.com/hashicorp/go-secure-stdlib/strutil v0.1.2 // indirect github.com/hashicorp/go-sockaddr v1.0.7 // indirect github.com/hashicorp/golang-lru v1.0.2 // indirect + github.com/hashicorp/golang-lru/arc/v2 v2.0.7 // indirect github.com/hashicorp/hcl v1.0.1-vault-7 // indirect github.com/hexops/gotextdiff v1.0.3 // indirect github.com/ipfs/bbloom v0.0.4 // indirect diff --git a/go.sum b/go.sum index d942d433..4722b2a5 100644 --- a/go.sum +++ b/go.sum @@ -449,6 +449,8 @@ github.com/hashicorp/go-version v1.8.0 h1:KAkNb1HAiZd1ukkxDFGmokVZe1Xy9HG6NUp+bP github.com/hashicorp/go-version v1.8.0/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09ZGVZPK5anwXA= github.com/hashicorp/golang-lru v1.0.2 h1:dV3g9Z/unq5DpblPpw+Oqcv4dU/1omnb4Ok8iPY6p1c= github.com/hashicorp/golang-lru v1.0.2/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4= +github.com/hashicorp/golang-lru/arc/v2 v2.0.7 h1:QxkVTxwColcduO+LP7eJO56r2hFiG8zEbfAAzRv52KQ= +github.com/hashicorp/golang-lru/arc/v2 v2.0.7/go.mod h1:Pe7gBlGdc8clY5LJ0LpJXMt5AmgmWNH1g+oFFVUHOEc= github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/hashicorp/hcl v1.0.1-vault-7 h1:ag5OxFVy3QYTFTJODRzTKVZ6xvdfLLCA1cy/Y6xGI0I= diff --git a/knotmirror/stream/persist/diskpersist/diskpersist.go b/knotmirror/stream/persist/diskpersist/diskpersist.go index 00a2dbf8..86ea4222 100644 --- a/knotmirror/stream/persist/diskpersist/diskpersist.go +++ b/knotmirror/stream/persist/diskpersist/diskpersist.go @@ -1,4 +1,104 @@ package diskpersist -// TODO: implement DiskPersistence -// but without gorm +import ( + "bytes" + "context" + "log/slog" + "os" + "sync" + "time" + + arc "github.com/hashicorp/golang-lru/arc/v2" + + "tangled.org/core/knotmirror/stream" + "tangled.org/core/knotmirror/stream/persist" +) + +type DiskPersistence struct { + primaryDir string + archiveDir string + eventsPerFile int64 + writeBufferSize int + retention time.Duration + + // meta *gorm.DB + + broadcast func(*stream.XRPCStreamEvent) + + logfi *os.File + + eventCounter int64 + curSeq int64 + initialSeq int64 + + uids UidSource + uidCache *arc.ARCCache[uint64, string] + didCache *arc.ARCCache[string, uint64] + + writers *sync.Pool + buffers *sync.Pool + scratch []byte + + outbuf *bytes.Buffer + evtbuf []persistJob + + shutdown chan struct{} + + log *slog.Logger + + lk sync.Mutex +} + +type persistJob struct { + Bytes []byte + Evt *stream.XRPCStreamEvent + Buffer *bytes.Buffer // so we can put it back in the pool when we're done +} + +const ( + EvtFlagTakedown = 1 << iota + EvtFlagRebased +) + +var _ (persist.EventPersistence) = (*DiskPersistence)(nil) + +type UidSource interface { + DidToUid(ctx context.Context, did string) (uint64, error) +} + +// func NewDiskPersistence(primaryDir, archiveDir string, db *gorm.DB, opts *DiskPersistOptions) (*DiskPersistence, error) { +// panic("unimplemented") +// } + +func (dp *DiskPersistence) SetUidSource(uids UidSource) { + dp.uids = uids +} + +func (dp *DiskPersistence) Persist(ctx context.Context, e *stream.XRPCStreamEvent) error { + panic("unimplemented") +} + +func (dp *DiskPersistence) Playback(ctx context.Context, since int64, cb func(*stream.XRPCStreamEvent) error) error { + panic("unimplemented") +} + +func (dp *DiskPersistence) TakeDownRepo(ctx context.Context, usr uint64) error { + panic("unimplemented") +} + +func (dp *DiskPersistence) Flush(ctx context.Context) error { + panic("unimplemented") +} + +func (dp *DiskPersistence) Shutdown(ctx context.Context) error { + close(dp.shutdown) + if err := dp.Flush(ctx); err != nil { + return err + } + + return dp.logfi.Close() +} + +func (dp *DiskPersistence) SetEventBroadcaster(f func(*stream.XRPCStreamEvent)) { + dp.broadcast = f +} -- 2.51.2