Skip to content

Commit 22dfab2

Browse files
authored
Merge pull request #5 from compforge/bolt-event-store
feat(go): add durable Bolt event store
2 parents 4e16d31 + ad017a7 commit 22dfab2

7 files changed

Lines changed: 390 additions & 8 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,7 @@
11
name: CI
22

33
on:
4-
push:
5-
pull_request:
4+
workflow_dispatch:
65

76
jobs:
87
python:

‎AGENTS.md‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,11 @@ Timeline、轨迹分析和评测等下游用途消费。
2626
│ └── packages/
2727
│ ├── core/ # TypeScript Core 与 Memory Store
2828
│ └── pi/ # Pi hooks 与 lossless native SessionStorage
29-
└── go/ # Go Core、Memory Store 与公共 Recorder API
30-
└── adapters/
31-
└── agentgo/ # AgentGo hooks 与原生恢复适配
29+
└── go/ # Go Core、Store 与公共 Recorder API
30+
├── adapters/
31+
│ └── agentgo/ # AgentGo hooks 与原生恢复适配
32+
└── stores/
33+
└── bolt/ # 单文件持久化 EventStore
3234
```
3335

3436
## 核心模型与关键约定

‎README.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ not silently replay a side-effecting tool.
4848
| `conformance/` | Cross-language golden vectors and adapter contract tests |
4949
| `python/` | Python core SDK plus memory, Redis, and SQLAlchemy stores |
5050
| `typescript/` | TypeScript core SDK and Pi adapter |
51-
| `go/` | Go core SDK and AgentGo adapter |
51+
| `go/` | Go core SDK, memory/Bolt stores, and AgentGo adapter |
5252

5353
Current framework profiles are integration examples, not definitions of the core session model:
5454

‎go/go.mod‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,9 @@ module github.com/compforge/agent-ledger/go
33
go 1.25.0
44

55
require (
6-
github.com/compforge/agentgo v0.0.1 // indirect
7-
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467 // indirect
6+
github.com/compforge/agentgo v0.0.1
7+
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467
8+
go.etcd.io/bbolt v1.5.0
89
)
10+
11+
require golang.org/x/sys v0.45.0 // indirect

‎go/go.sum‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,3 +2,17 @@ github.com/compforge/agentgo v0.0.1 h1:e3JiF7za1xN9NdGw4M8BFH9vdVwMJuD/K55eNSclK
22
github.com/compforge/agentgo v0.0.1/go.mod h1:5EkjADRpln5pwK2wd1cNwUwndq2VQDbeefybMVNMhpk=
33
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467 h1:uX1JmpONuD549D73r6cgnxyUu18Zb7yHAy5AYU0Pm4Q=
44
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467/go.mod h1:uzvlm1mxhHkdfqitSA92i7Se+S9ksOn3a3qmv/kyOCw=
5+
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
6+
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
7+
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
8+
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
9+
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
10+
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
11+
go.etcd.io/bbolt v1.5.0 h1:S7GAl7Fxv12yohbwFfIbQCGDWbQbtDGPET4P/bD4lxU=
12+
go.etcd.io/bbolt v1.5.0/go.mod h1:mkltfYE5aUHQxUct9N9V+Kp7aSjFqjgrhcXIS70Lrdk=
13+
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
14+
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
15+
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
16+
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
17+
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
18+
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=

‎go/stores/bolt/store.go‎

Lines changed: 280 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,280 @@
1+
package boltstore
2+
3+
import (
4+
"context"
5+
"encoding/binary"
6+
"encoding/json"
7+
"errors"
8+
"fmt"
9+
"iter"
10+
"strconv"
11+
"time"
12+
13+
agentledger "github.com/compforge/agent-ledger/go"
14+
bolt "go.etcd.io/bbolt"
15+
)
16+
17+
var (
18+
streamsBucket = []byte("streams")
19+
sessionsBucket = []byte("sessions")
20+
receiptsBucket = []byte("receipts")
21+
eventIDsBucket = []byte("event_ids")
22+
)
23+
24+
// Store persists Agent Ledger streams in one Bolt database. Bolt serializes
25+
// writers, so the EventStore append contract and its optimistic version check
26+
// are committed in the same transaction.
27+
type Store struct {
28+
db *bolt.DB
29+
}
30+
31+
func Open(path string, timeout time.Duration) (*Store, error) {
32+
if timeout <= 0 {
33+
return nil, errors.New("open bolt event store: timeout must be positive")
34+
}
35+
db, err := bolt.Open(path, 0o600, &bolt.Options{Timeout: timeout})
36+
if err != nil {
37+
return nil, fmt.Errorf("open bolt event store %q: %w", path, err)
38+
}
39+
return &Store{db: db}, nil
40+
}
41+
42+
func (s *Store) Close() error {
43+
if err := s.db.Close(); err != nil {
44+
return fmt.Errorf("close bolt event store: %w", err)
45+
}
46+
return nil
47+
}
48+
49+
func (s *Store) Append(
50+
ctx context.Context,
51+
stream agentledger.EventStream,
52+
expectedVersion int64,
53+
appendID string,
54+
events ...agentledger.ProposedEvent,
55+
) (agentledger.CommitReceipt, error) {
56+
if err := ctx.Err(); err != nil {
57+
return agentledger.CommitReceipt{}, err
58+
}
59+
if len(events) == 0 {
60+
return agentledger.CommitReceipt{}, errors.New("append requires at least one event")
61+
}
62+
batch, err := clone(events)
63+
if err != nil {
64+
return agentledger.CommitReceipt{}, fmt.Errorf("snapshot append batch: %w", err)
65+
}
66+
seen := make(map[string]struct{}, len(batch))
67+
for _, event := range batch {
68+
if event.SessionID != stream.SessionID {
69+
return agentledger.CommitReceipt{}, errors.New("all events must belong to the target stream's session")
70+
}
71+
if _, duplicate := seen[event.EventID]; duplicate {
72+
return agentledger.CommitReceipt{}, fmt.Errorf("%w: %s", agentledger.ErrDuplicateEvent, event.EventID)
73+
}
74+
seen[event.EventID] = struct{}{}
75+
}
76+
digest, err := agentledger.CanonicalAppendDigest(batch)
77+
if err != nil {
78+
return agentledger.CommitReceipt{}, err
79+
}
80+
81+
var receipt agentledger.CommitReceipt
82+
err = s.db.Update(func(tx *bolt.Tx) error {
83+
if err := ctx.Err(); err != nil {
84+
return err
85+
}
86+
streams, err := tx.CreateBucketIfNotExists(streamsBucket)
87+
if err != nil {
88+
return err
89+
}
90+
sessions, err := tx.CreateBucketIfNotExists(sessionsBucket)
91+
if err != nil {
92+
return err
93+
}
94+
receipts, err := tx.CreateBucketIfNotExists(receiptsBucket)
95+
if err != nil {
96+
return err
97+
}
98+
eventIDs, err := tx.CreateBucketIfNotExists(eventIDsBucket)
99+
if err != nil {
100+
return err
101+
}
102+
103+
receiptKey := composite(stream.SessionID, stream.StreamID, appendID)
104+
if encoded := receipts.Get(receiptKey); encoded != nil {
105+
if err := json.Unmarshal(encoded, &receipt); err != nil {
106+
return fmt.Errorf("decode append receipt: %w", err)
107+
}
108+
if receipt.Digest != digest {
109+
return agentledger.ErrIdempotencyViolation
110+
}
111+
return nil
112+
}
113+
114+
streamBucket, err := streams.CreateBucketIfNotExists(composite(stream.SessionID, stream.StreamID))
115+
if err != nil {
116+
return err
117+
}
118+
currentVersion := int64(streamBucket.Sequence()) - 1
119+
if currentVersion != expectedVersion {
120+
return fmt.Errorf("%w: expected %d, actual %d", agentledger.ErrStreamConflict, expectedVersion, currentVersion)
121+
}
122+
for _, event := range batch {
123+
key := composite(stream.SessionID, event.EventID)
124+
if eventIDs.Get(key) != nil {
125+
return fmt.Errorf("%w: %s", agentledger.ErrDuplicateEvent, event.EventID)
126+
}
127+
}
128+
129+
sessionBucket, err := sessions.CreateBucketIfNotExists([]byte(stream.SessionID))
130+
if err != nil {
131+
return err
132+
}
133+
committedAt := time.Now().UTC().Format(time.RFC3339Nano)
134+
stored := make([]agentledger.StoredEvent, 0, len(batch))
135+
for _, event := range batch {
136+
streamSequence, err := streamBucket.NextSequence()
137+
if err != nil {
138+
return err
139+
}
140+
sessionSequence, err := sessionBucket.NextSequence()
141+
if err != nil {
142+
return err
143+
}
144+
item := agentledger.StoredEvent{
145+
ProposedEvent: event,
146+
StreamID: stream.StreamID,
147+
StreamVersion: int64(streamSequence) - 1,
148+
CommitCursor: strconv.FormatUint(sessionSequence-1, 10),
149+
CommittedAt: committedAt,
150+
}
151+
encoded, err := json.Marshal(item)
152+
if err != nil {
153+
return fmt.Errorf("encode stored event: %w", err)
154+
}
155+
if err := streamBucket.Put(sequenceKey(streamSequence-1), encoded); err != nil {
156+
return err
157+
}
158+
if err := sessionBucket.Put(sequenceKey(sessionSequence-1), encoded); err != nil {
159+
return err
160+
}
161+
if err := eventIDs.Put(composite(stream.SessionID, event.EventID), []byte{1}); err != nil {
162+
return err
163+
}
164+
stored = append(stored, item)
165+
}
166+
167+
receipt = agentledger.CommitReceipt{
168+
Stream: stream,
169+
AppendID: appendID,
170+
Digest: digest,
171+
FirstVersion: stored[0].StreamVersion,
172+
LastVersion: stored[len(stored)-1].StreamVersion,
173+
FirstCursor: stored[0].CommitCursor,
174+
LastCursor: stored[len(stored)-1].CommitCursor,
175+
CommittedAt: committedAt,
176+
}
177+
for _, event := range stored {
178+
receipt.EventIDs = append(receipt.EventIDs, event.EventID)
179+
}
180+
encoded, err := json.Marshal(receipt)
181+
if err != nil {
182+
return fmt.Errorf("encode append receipt: %w", err)
183+
}
184+
return receipts.Put(receiptKey, encoded)
185+
})
186+
if err != nil {
187+
return agentledger.CommitReceipt{}, fmt.Errorf("append bolt event batch: %w", err)
188+
}
189+
return clone(receipt)
190+
}
191+
192+
func (s *Store) Load(ctx context.Context, stream agentledger.EventStream, afterVersion int64) iter.Seq2[agentledger.StoredEvent, error] {
193+
return s.read(ctx, streamsBucket, composite(stream.SessionID, stream.StreamID), afterVersion)
194+
}
195+
196+
func (s *Store) ScanSession(ctx context.Context, sessionID, afterCursor string) iter.Seq2[agentledger.StoredEvent, error] {
197+
after := int64(-1)
198+
if afterCursor != "" {
199+
value, err := strconv.ParseInt(afterCursor, 10, 64)
200+
if err != nil || value < 0 {
201+
return errorSequence(fmt.Errorf("invalid cursor %q", afterCursor))
202+
}
203+
after = value
204+
}
205+
return s.read(ctx, sessionsBucket, []byte(sessionID), after)
206+
}
207+
208+
func (s *Store) read(ctx context.Context, rootName, childName []byte, after int64) iter.Seq2[agentledger.StoredEvent, error] {
209+
return func(yield func(agentledger.StoredEvent, error) bool) {
210+
var encodedEvents [][]byte
211+
err := s.db.View(func(tx *bolt.Tx) error {
212+
root := tx.Bucket(rootName)
213+
if root == nil {
214+
return nil
215+
}
216+
child := root.Bucket(childName)
217+
if child == nil {
218+
return nil
219+
}
220+
cursor := child.Cursor()
221+
for key, value := cursor.Seek(sequenceKey(uint64(after + 1))); key != nil; key, value = cursor.Next() {
222+
encodedEvents = append(encodedEvents, append([]byte(nil), value...))
223+
}
224+
return nil
225+
})
226+
if err != nil {
227+
yield(agentledger.StoredEvent{}, fmt.Errorf("read bolt event stream: %w", err))
228+
return
229+
}
230+
for _, encoded := range encodedEvents {
231+
if err := ctx.Err(); err != nil {
232+
yield(agentledger.StoredEvent{}, err)
233+
return
234+
}
235+
var event agentledger.StoredEvent
236+
if err := json.Unmarshal(encoded, &event); err != nil {
237+
yield(agentledger.StoredEvent{}, fmt.Errorf("decode stored event: %w", err))
238+
return
239+
}
240+
if !yield(event, nil) {
241+
return
242+
}
243+
}
244+
}
245+
}
246+
247+
func composite(parts ...string) []byte {
248+
var result []byte
249+
for index, part := range parts {
250+
if index > 0 {
251+
result = append(result, 0)
252+
}
253+
result = append(result, part...)
254+
}
255+
return result
256+
}
257+
258+
func sequenceKey(value uint64) []byte {
259+
key := make([]byte, 8)
260+
binary.BigEndian.PutUint64(key, value)
261+
return key
262+
}
263+
264+
func clone[T any](value T) (T, error) {
265+
var result T
266+
encoded, err := json.Marshal(value)
267+
if err != nil {
268+
return result, err
269+
}
270+
if err := json.Unmarshal(encoded, &result); err != nil {
271+
return result, err
272+
}
273+
return result, nil
274+
}
275+
276+
func errorSequence(err error) iter.Seq2[agentledger.StoredEvent, error] {
277+
return func(yield func(agentledger.StoredEvent, error) bool) {
278+
yield(agentledger.StoredEvent{}, err)
279+
}
280+
}

0 commit comments

Comments
 (0)