From c162b72247e06d15bf652891ad0fb500e73d1739 Mon Sep 17 00:00:00 2001 From: "liqiankun.1111" Date: Wed, 19 Aug 2026 13:39:11 +0800 Subject: [PATCH 1/2] feat(go): add durable Bolt event store Add a single-file EventStore that preserves atomic append, idempotency, optimistic concurrency, session scans, and persistence across restarts. --- AGENTS.md | 8 +- README.md | 2 +- go/go.mod | 7 +- go/go.sum | 14 ++ go/stores/bolt/store.go | 280 +++++++++++++++++++++++++++++++++++ go/stores/bolt/store_test.go | 84 +++++++++++ 6 files changed, 389 insertions(+), 6 deletions(-) create mode 100644 go/stores/bolt/store.go create mode 100644 go/stores/bolt/store_test.go diff --git a/AGENTS.md b/AGENTS.md index bff302b..ff9c1d8 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 ``` ## 核心模型与关键约定 diff --git a/README.md b/README.md index 045794a..23160a1 100644 --- a/README.md +++ b/README.md @@ -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: diff --git a/go/go.mod b/go/go.mod index a1abc91..b2c0b47 100644 --- a/go/go.mod +++ b/go/go.mod @@ -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 diff --git a/go/go.sum b/go/go.sum index 1d61692..6a199a3 100644 --- a/go/go.sum +++ b/go/go.sum @@ -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= diff --git a/go/stores/bolt/store.go b/go/stores/bolt/store.go new file mode 100644 index 0000000..4be75a1 --- /dev/null +++ b/go/stores/bolt/store.go @@ -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) + } +} diff --git a/go/stores/bolt/store_test.go b/go/stores/bolt/store_test.go new file mode 100644 index 0000000..55aee3b --- /dev/null +++ b/go/stores/bolt/store_test.go @@ -0,0 +1,84 @@ +package boltstore + +import ( + "context" + "errors" + "path/filepath" + "testing" + "time" + + agentledger "github.com/compforge/agent-ledger/go" +) + +func TestStorePersistsAndResumesAppendContract(t *testing.T) { + path := filepath.Join(t.TempDir(), "ledger.db") + store := openTestStore(t, path) + stream := agentledger.EventStream{SessionID: "session-1", StreamID: "run-1"} + event := agentledger.NewEvent("run.started", stream.SessionID, "run-1", agentledger.Actor{Type: "agent", ID: "test"}) + + receipt, err := store.Append(context.Background(), stream, -1, "append-1", event) + if err != nil { + t.Fatalf("append: %v", err) + } + if err := store.Close(); err != nil { + t.Fatalf("close: %v", err) + } + + store = openTestStore(t, path) + t.Cleanup(func() { _ = store.Close() }) + replayed, err := store.Append(context.Background(), stream, -1, "append-1", event) + if err != nil { + t.Fatalf("replay idempotent append: %v", err) + } + if replayed.Digest != receipt.Digest || replayed.LastVersion != receipt.LastVersion { + t.Fatalf("replayed receipt = %#v, want %#v", replayed, receipt) + } + + second := agentledger.NewEvent("run.completed", stream.SessionID, "run-1", agentledger.Actor{Type: "agent", ID: "test"}) + if _, err := store.Append(context.Background(), stream, -1, "append-2", second); !errors.Is(err, agentledger.ErrStreamConflict) { + t.Fatalf("stale append error = %v, want ErrStreamConflict", err) + } + if _, err := store.Append(context.Background(), stream, 0, "append-2", second); err != nil { + t.Fatalf("append after reopen: %v", err) + } + + var eventTypes []string + for stored, err := range store.ScanSession(context.Background(), stream.SessionID, "") { + if err != nil { + t.Fatalf("scan: %v", err) + } + eventTypes = append(eventTypes, stored.EventType) + } + if len(eventTypes) != 2 || eventTypes[0] != "run.started" || eventTypes[1] != "run.completed" { + t.Fatalf("event types = %v", eventTypes) + } +} + +func TestStoreRejectsDuplicateEventAcrossStreamsInSession(t *testing.T) { + store := openTestStore(t, filepath.Join(t.TempDir(), "ledger.db")) + t.Cleanup(func() { _ = store.Close() }) + ctx := context.Background() + event := agentledger.NewEvent("test.recorded", "session", "run-1", agentledger.Actor{Type: "agent", ID: "test"}) + if _, err := store.Append(ctx, agentledger.EventStream{SessionID: "session", StreamID: "run-1"}, -1, "one", event); err != nil { + t.Fatalf("append first: %v", err) + } + event.RunID = "run-2" + if _, err := store.Append(ctx, agentledger.EventStream{SessionID: "session", StreamID: "run-2"}, -1, "two", event); !errors.Is(err, agentledger.ErrDuplicateEvent) { + t.Fatalf("duplicate error = %v, want ErrDuplicateEvent", err) + } +} + +func TestOpenRequiresExplicitTimeout(t *testing.T) { + if _, err := Open(filepath.Join(t.TempDir(), "ledger.db"), 0); err == nil { + t.Fatal("Open accepted zero timeout") + } +} + +func openTestStore(t *testing.T, path string) *Store { + t.Helper() + store, err := Open(path, time.Second) + if err != nil { + t.Fatalf("open store: %v", err) + } + return store +} From ad017a7eef97ecfad2cf76cff2be998830190add Mon Sep 17 00:00:00 2001 From: "liqiankun.1111" Date: Wed, 19 Aug 2026 13:41:40 +0800 Subject: [PATCH 2/2] ci: disable automatic workflows --- .github/workflows/ci.yml | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3eaeba5..e0e279b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -1,8 +1,7 @@ name: CI on: - push: - pull_request: + workflow_dispatch: jobs: python: