Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
@@ -1,8 +1,7 @@
name: CI

on:
push:
pull_request:
workflow_dispatch:

jobs:
python:
Expand Down
8 changes: 5 additions & 3 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,11 @@ Timeline、轨迹分析和评测等下游用途消费。
│ └── packages/
│ ├── core/ # TypeScript Core 与 Memory Store
│ └── pi/ # Pi hooks 与 lossless native SessionStorage
└── go/ # Go Core、Memory Store 与公共 Recorder API
└── adapters/
└── agentgo/ # AgentGo hooks 与原生恢复适配
└── go/ # Go Core、Store 与公共 Recorder API
├── adapters/
│ └── agentgo/ # AgentGo hooks 与原生恢复适配
└── stores/
└── bolt/ # 单文件持久化 EventStore
```

## 核心模型与关键约定
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ not silently replay a side-effecting tool.
| `conformance/` | Cross-language golden vectors and adapter contract tests |
| `python/` | Python core SDK plus memory, Redis, and SQLAlchemy stores |
| `typescript/` | TypeScript core SDK and Pi adapter |
| `go/` | Go core SDK and AgentGo adapter |
| `go/` | Go core SDK, memory/Bolt stores, and AgentGo adapter |

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

Expand Down
7 changes: 5 additions & 2 deletions go/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,9 @@ module github.com/compforge/agent-ledger/go
go 1.25.0

require (
github.com/compforge/agentgo v0.0.1 // indirect
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467 // indirect
github.com/compforge/agentgo v0.0.1
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467
go.etcd.io/bbolt v1.5.0
)

require golang.org/x/sys v0.45.0 // indirect
14 changes: 14 additions & 0 deletions go/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -2,3 +2,17 @@ github.com/compforge/agentgo v0.0.1 h1:e3JiF7za1xN9NdGw4M8BFH9vdVwMJuD/K55eNSclK
github.com/compforge/agentgo v0.0.1/go.mod h1:5EkjADRpln5pwK2wd1cNwUwndq2VQDbeefybMVNMhpk=
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467 h1:uX1JmpONuD549D73r6cgnxyUu18Zb7yHAy5AYU0Pm4Q=
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467/go.mod h1:uzvlm1mxhHkdfqitSA92i7Se+S9ksOn3a3qmv/kyOCw=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
go.etcd.io/bbolt v1.5.0 h1:S7GAl7Fxv12yohbwFfIbQCGDWbQbtDGPET4P/bD4lxU=
go.etcd.io/bbolt v1.5.0/go.mod h1:mkltfYE5aUHQxUct9N9V+Kp7aSjFqjgrhcXIS70Lrdk=
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
280 changes: 280 additions & 0 deletions go/stores/bolt/store.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,280 @@
package boltstore

import (
"context"
"encoding/binary"
"encoding/json"
"errors"
"fmt"
"iter"
"strconv"
"time"

agentledger "github.com/compforge/agent-ledger/go"
bolt "go.etcd.io/bbolt"
)

var (
streamsBucket = []byte("streams")
sessionsBucket = []byte("sessions")
receiptsBucket = []byte("receipts")
eventIDsBucket = []byte("event_ids")
)

// Store persists Agent Ledger streams in one Bolt database. Bolt serializes
// writers, so the EventStore append contract and its optimistic version check
// are committed in the same transaction.
type Store struct {
db *bolt.DB
}

func Open(path string, timeout time.Duration) (*Store, error) {
if timeout <= 0 {
return nil, errors.New("open bolt event store: timeout must be positive")
}
db, err := bolt.Open(path, 0o600, &bolt.Options{Timeout: timeout})
if err != nil {
return nil, fmt.Errorf("open bolt event store %q: %w", path, err)
}
return &Store{db: db}, nil
}

func (s *Store) Close() error {
if err := s.db.Close(); err != nil {
return fmt.Errorf("close bolt event store: %w", err)
}
return nil
}

func (s *Store) Append(
ctx context.Context,
stream agentledger.EventStream,
expectedVersion int64,
appendID string,
events ...agentledger.ProposedEvent,
) (agentledger.CommitReceipt, error) {
if err := ctx.Err(); err != nil {
return agentledger.CommitReceipt{}, err
}
if len(events) == 0 {
return agentledger.CommitReceipt{}, errors.New("append requires at least one event")
}
batch, err := clone(events)
if err != nil {
return agentledger.CommitReceipt{}, fmt.Errorf("snapshot append batch: %w", err)
}
seen := make(map[string]struct{}, len(batch))
for _, event := range batch {
if event.SessionID != stream.SessionID {
return agentledger.CommitReceipt{}, errors.New("all events must belong to the target stream's session")
}
if _, duplicate := seen[event.EventID]; duplicate {
return agentledger.CommitReceipt{}, fmt.Errorf("%w: %s", agentledger.ErrDuplicateEvent, event.EventID)
}
seen[event.EventID] = struct{}{}
}
digest, err := agentledger.CanonicalAppendDigest(batch)
if err != nil {
return agentledger.CommitReceipt{}, err
}

var receipt agentledger.CommitReceipt
err = s.db.Update(func(tx *bolt.Tx) error {
if err := ctx.Err(); err != nil {
return err
}
streams, err := tx.CreateBucketIfNotExists(streamsBucket)
if err != nil {
return err
}
sessions, err := tx.CreateBucketIfNotExists(sessionsBucket)
if err != nil {
return err
}
receipts, err := tx.CreateBucketIfNotExists(receiptsBucket)
if err != nil {
return err
}
eventIDs, err := tx.CreateBucketIfNotExists(eventIDsBucket)
if err != nil {
return err
}

receiptKey := composite(stream.SessionID, stream.StreamID, appendID)
if encoded := receipts.Get(receiptKey); encoded != nil {
if err := json.Unmarshal(encoded, &receipt); err != nil {
return fmt.Errorf("decode append receipt: %w", err)
}
if receipt.Digest != digest {
return agentledger.ErrIdempotencyViolation
}
return nil
}

streamBucket, err := streams.CreateBucketIfNotExists(composite(stream.SessionID, stream.StreamID))
if err != nil {
return err
}
currentVersion := int64(streamBucket.Sequence()) - 1
if currentVersion != expectedVersion {
return fmt.Errorf("%w: expected %d, actual %d", agentledger.ErrStreamConflict, expectedVersion, currentVersion)
}
for _, event := range batch {
key := composite(stream.SessionID, event.EventID)
if eventIDs.Get(key) != nil {
return fmt.Errorf("%w: %s", agentledger.ErrDuplicateEvent, event.EventID)
}
}

sessionBucket, err := sessions.CreateBucketIfNotExists([]byte(stream.SessionID))
if err != nil {
return err
}
committedAt := time.Now().UTC().Format(time.RFC3339Nano)
stored := make([]agentledger.StoredEvent, 0, len(batch))
for _, event := range batch {
streamSequence, err := streamBucket.NextSequence()
if err != nil {
return err
}
sessionSequence, err := sessionBucket.NextSequence()
if err != nil {
return err
}
item := agentledger.StoredEvent{
ProposedEvent: event,
StreamID: stream.StreamID,
StreamVersion: int64(streamSequence) - 1,
CommitCursor: strconv.FormatUint(sessionSequence-1, 10),
CommittedAt: committedAt,
}
encoded, err := json.Marshal(item)
if err != nil {
return fmt.Errorf("encode stored event: %w", err)
}
if err := streamBucket.Put(sequenceKey(streamSequence-1), encoded); err != nil {
return err
}
if err := sessionBucket.Put(sequenceKey(sessionSequence-1), encoded); err != nil {
return err
}
if err := eventIDs.Put(composite(stream.SessionID, event.EventID), []byte{1}); err != nil {
return err
}
stored = append(stored, item)
}

receipt = agentledger.CommitReceipt{
Stream: stream,
AppendID: appendID,
Digest: digest,
FirstVersion: stored[0].StreamVersion,
LastVersion: stored[len(stored)-1].StreamVersion,
FirstCursor: stored[0].CommitCursor,
LastCursor: stored[len(stored)-1].CommitCursor,
CommittedAt: committedAt,
}
for _, event := range stored {
receipt.EventIDs = append(receipt.EventIDs, event.EventID)
}
encoded, err := json.Marshal(receipt)
if err != nil {
return fmt.Errorf("encode append receipt: %w", err)
}
return receipts.Put(receiptKey, encoded)
})
if err != nil {
return agentledger.CommitReceipt{}, fmt.Errorf("append bolt event batch: %w", err)
}
return clone(receipt)
}

func (s *Store) Load(ctx context.Context, stream agentledger.EventStream, afterVersion int64) iter.Seq2[agentledger.StoredEvent, error] {
return s.read(ctx, streamsBucket, composite(stream.SessionID, stream.StreamID), afterVersion)
}

func (s *Store) ScanSession(ctx context.Context, sessionID, afterCursor string) iter.Seq2[agentledger.StoredEvent, error] {
after := int64(-1)
if afterCursor != "" {
value, err := strconv.ParseInt(afterCursor, 10, 64)
if err != nil || value < 0 {
return errorSequence(fmt.Errorf("invalid cursor %q", afterCursor))
}
after = value
}
return s.read(ctx, sessionsBucket, []byte(sessionID), after)
}

func (s *Store) read(ctx context.Context, rootName, childName []byte, after int64) iter.Seq2[agentledger.StoredEvent, error] {
return func(yield func(agentledger.StoredEvent, error) bool) {
var encodedEvents [][]byte
err := s.db.View(func(tx *bolt.Tx) error {
root := tx.Bucket(rootName)
if root == nil {
return nil
}
child := root.Bucket(childName)
if child == nil {
return nil
}
cursor := child.Cursor()
for key, value := cursor.Seek(sequenceKey(uint64(after + 1))); key != nil; key, value = cursor.Next() {
encodedEvents = append(encodedEvents, append([]byte(nil), value...))
}
return nil
})
if err != nil {
yield(agentledger.StoredEvent{}, fmt.Errorf("read bolt event stream: %w", err))
return
}
for _, encoded := range encodedEvents {
if err := ctx.Err(); err != nil {
yield(agentledger.StoredEvent{}, err)
return
}
var event agentledger.StoredEvent
if err := json.Unmarshal(encoded, &event); err != nil {
yield(agentledger.StoredEvent{}, fmt.Errorf("decode stored event: %w", err))
return
}
if !yield(event, nil) {
return
}
}
}
}

func composite(parts ...string) []byte {
var result []byte
for index, part := range parts {
if index > 0 {
result = append(result, 0)
}
result = append(result, part...)
}
return result
}

func sequenceKey(value uint64) []byte {
key := make([]byte, 8)
binary.BigEndian.PutUint64(key, value)
return key
}

func clone[T any](value T) (T, error) {
var result T
encoded, err := json.Marshal(value)
if err != nil {
return result, err
}
if err := json.Unmarshal(encoded, &result); err != nil {
return result, err
}
return result, nil
}

func errorSequence(err error) iter.Seq2[agentledger.StoredEvent, error] {
return func(yield func(agentledger.StoredEvent, error) bool) {
yield(agentledger.StoredEvent{}, err)
}
}
Loading