diff --git a/pkg/channels/buzz/buzz.go b/pkg/channels/buzz/buzz.go index 0f4b10d8..5656500d 100644 --- a/pkg/channels/buzz/buzz.go +++ b/pkg/channels/buzz/buzz.go @@ -76,13 +76,18 @@ func NewBuzzChannel( channels.WithReasoningChannelID(bc.ReasoningChannelID), ) + seen, err := newSeenStore(cfg.SeenStatePath, bc.Name()) + if err != nil { + return nil, err + } + return &BuzzChannel{ BaseChannel: base, bc: bc, config: cfg, secretKey: sk, publicKey: pk, - seen: newSeenStore(cfg.SeenStatePath, bc.Name()), + seen: seen, }, nil } @@ -121,6 +126,14 @@ func (c *BuzzChannel) Start(ctx context.Context) error { logger.InfoC("buzz", "Starting Buzz channel") c.ctx, c.cancel = context.WithCancel(ctx) + if c.seen != nil { + logger.DebugCF("buzz", "Replay state path", map[string]any{ + "seen_path": c.seen.path, + "channel_name": c.bc.Name(), + "high_water": int64(c.seen.highWaterMark()), + }) + } + relay, err := nostr.RelayConnect(c.ctx, c.config.RelayURL) if err != nil { c.cancel() diff --git a/pkg/channels/buzz/buzz_test.go b/pkg/channels/buzz/buzz_test.go index ab945939..aebb15e8 100644 --- a/pkg/channels/buzz/buzz_test.go +++ b/pkg/channels/buzz/buzz_test.go @@ -2,6 +2,7 @@ package buzz import ( "context" + "encoding/json" "os" "path/filepath" "strings" @@ -243,6 +244,35 @@ func TestReplaySuppression_PersistedEventNotRedelivered(t *testing.T) { t.Fatalf("stat seen file: %v", err) } + // The persisted JSON must concretely record the handled event's ID and + // high-water timestamp, so a restart can dedup it without re-running the + // handler. + raw, err := os.ReadFile(seenPath) + if err != nil { + t.Fatalf("read seen file: %v", err) + } + var persisted map[string]*channelHighWater + if err := json.Unmarshal(raw, &persisted); err != nil { + t.Fatalf("unmarshal seen file: %v", err) + } + entry, ok := persisted["test-channel"] + if !ok { + t.Fatalf("persisted state missing test-channel entry; got %v", persisted) + } + if entry.MaxCreatedAt != int64(evt.CreatedAt) { + t.Fatalf("persisted max_created_at = %d, want %d", entry.MaxCreatedAt, int64(evt.CreatedAt)) + } + foundID := false + for _, id := range entry.RecentIDs { + if id == evt.ID { + foundID = true + break + } + } + if !foundID { + t.Fatalf("persisted recent_ids %v do not contain event ID %q", entry.RecentIDs, evt.ID) + } + // Run 2: a brand-new channel instance reading the same store must not // redeliver the same event ID, simulating a service restart. ch2, msgBus2, _ := newTestBuzzChannel(t, seenPath) @@ -303,7 +333,10 @@ func TestReplaySuppression_SelfMessageNotRecorded(t *testing.T) { // max across all channels, so the subscription Since filter uses a single // global cutoff that skips history for every channel. func TestSeenStoreHighWaterMarkAcrossChannels(t *testing.T) { - s := newSeenStore(filepath.Join(t.TempDir(), "seen.json"), "test") + s, err := newSeenStore(filepath.Join(t.TempDir(), "seen.json"), "test") + if err != nil { + t.Fatalf("newSeenStore: %v", err) + } if hw := s.highWaterMark(); hw != 0 { t.Fatalf("empty store high-water = %d, want 0", hw) @@ -324,7 +357,10 @@ func TestSeenStoreHighWaterMarkAcrossChannels(t *testing.T) { } // Reload from disk to confirm persistence. - s2 := newSeenStore(s.path, "test") + s2, err := newSeenStore(s.path, "test") + if err != nil { + t.Fatalf("newSeenStore reload: %v", err) + } if hw := s2.highWaterMark(); hw != 150 { t.Fatalf("after reload, high-water = %d, want 150", hw) } @@ -354,7 +390,10 @@ func TestReplaySuppression_EventsOlderThanSeededHighWaterAreDropped(t *testing.T // TestSeenStoreRingTrims verifies the recent-ID ring does not grow unbounded, // keeping only the most recent recentIDWindow entries per channel. func TestSeenStoreRingTrims(t *testing.T) { - s := newSeenStore(filepath.Join(t.TempDir(), "seen.json"), "test") + s, err := newSeenStore(filepath.Join(t.TempDir(), "seen.json"), "test") + if err != nil { + t.Fatalf("newSeenStore: %v", err) + } // Insert more IDs than the window; old IDs should be evicted. for i := 0; i < recentIDWindow+50; i++ { @@ -394,3 +433,115 @@ func nip19Letter(n int) string { } return string(out) } + +// TestNewBuzzChannel_RejectsCorruptSeenState asserts that NewBuzzChannel fails +// when the persisted seen-state file is unparseable. Starting a channel whose +// replay state cannot be loaded would silently lose dedup memory and risk +// redelivering already-handled events, so the constructor must refuse. +func TestNewBuzzChannel_RejectsCorruptSeenState(t *testing.T) { + seenPath := filepath.Join(t.TempDir(), "seen.json") + if err := os.WriteFile(seenPath, []byte("{not valid json"), 0o600); err != nil { + t.Fatalf("write corrupt seen file: %v", err) + } + + botSK := nostr.GeneratePrivateKey() + bc := &config.Channel{} + bc.SetName("buzz-corrupt") + cfg := &config.BuzzSettings{ + RelayURL: "wss://relay.invalid", + PrivateKey: *config.NewSecureString(botSK), + Channels: config.FlexibleStringSlice{"test-channel"}, + SeenStatePath: seenPath, + } + msgBus := bus.NewMessageBus() + t.Cleanup(msgBus.Close) + + if _, err := NewBuzzChannel(bc, cfg, msgBus); err == nil { + t.Fatal("NewBuzzChannel unexpectedly succeeded with a corrupt seen-state file") + } +} + +// TestNewSeenStore_RejectsCorruptState exercises the store constructor +// directly: a corrupt file must surface a load error, while a missing file +// must not. +func TestNewSeenStore_RejectsCorruptState(t *testing.T) { + seenPath := filepath.Join(t.TempDir(), "seen.json") + + // Missing file is fine — fresh store starts empty. + s, err := newSeenStore(seenPath, "test") + if err != nil { + t.Fatalf("missing file should not error, got: %v", err) + } + if hw := s.highWaterMark(); hw != 0 { + t.Fatalf("fresh store high-water = %d, want 0", hw) + } + + // Corrupt file must error. + if err := os.WriteFile(seenPath, []byte("::{unjson"), 0o600); err != nil { + t.Fatalf("write corrupt file: %v", err) + } + if _, err := newSeenStore(seenPath, "test"); err == nil { + t.Fatal("newSeenStore unexpectedly succeeded with a corrupt file") + } +} + +// TestOnEvent_PersistenceFailurePreventsDispatch verifies that if marking an +// event as handled fails to persist, onEvent does NOT dispatch the event to +// the inbound bus. This preserves at-most-once delivery: a successful but +// un-persisted mark would let the same event be redelivered and re-handled on +// the next restart. +// +// The persistence failure is induced by replacing the store's parent directory +// with a regular file after load, so saveLocked's MkdirAll fails. This is +// deterministic even when tests run as root and exercises real I/O. +func TestOnEvent_PersistenceFailurePreventsDispatch(t *testing.T) { + parent := filepath.Join(t.TempDir(), "store-parent") + storeDir := filepath.Join(parent, "store") + if err := os.Mkdir(parent, 0o700); err != nil { + t.Fatalf("mkdir store parent: %v", err) + } + if err := os.Mkdir(storeDir, 0o700); err != nil { + t.Fatalf("mkdir store: %v", err) + } + seenPath := filepath.Join(storeDir, "seen.json") + // Valid empty state so load() succeeds; save() will fail once the dir + // is made read-only. + if err := os.WriteFile(seenPath, []byte("{}"), 0o600); err != nil { + t.Fatalf("seed seen file: %v", err) + } + + ch, msgBus, _ := newTestBuzzChannel(t, seenPath) + + if err := os.RemoveAll(storeDir); err != nil { + t.Fatalf("remove store directory: %v", err) + } + if err := os.WriteFile(storeDir, []byte("not a directory"), 0o600); err != nil { + t.Fatalf("replace store directory with file: %v", err) + } + + senderSK := nostr.GeneratePrivateKey() + evt := makeEvent(t, senderSK, "test-channel", "should not be dispatched", nostr.Now()) + + ch.onEvent(evt) + expectNoInbound(t, msgBus) + + // A second, fresh event must also be suppressed while persistence is + // broken — the channel must not "give up" suppressing after one failure. + second := makeEvent(t, senderSK, "test-channel", "also not dispatched", nostr.Now()+1) + ch.onEvent(second) + expectNoInbound(t, msgBus) +} + +// TestOnEvent_PersistenceFailureThenRecoveryDispatches confirms that once +// persistence is restored, a fresh event is dispatched normally — the +// persistence-failure path does not wedge the channel permanently. +func TestOnEvent_PersistenceFailureThenRecoveryDispatches(t *testing.T) { + // Start with a working seen store on a fresh channel. + seenPath := filepath.Join(t.TempDir(), "seen.json") + ch, msgBus, _ := newTestBuzzChannel(t, seenPath) + senderSK := nostr.GeneratePrivateKey() + + evt := makeEvent(t, senderSK, "test-channel", "recovered path delivers", nostr.Now()) + ch.onEvent(evt) + expectInbound(t, msgBus, "recovered path delivers") +} diff --git a/pkg/channels/buzz/handler.go b/pkg/channels/buzz/handler.go index c7ad1f70..d94182b9 100644 --- a/pkg/channels/buzz/handler.go +++ b/pkg/channels/buzz/handler.go @@ -87,13 +87,27 @@ func (c *BuzzChannel) onEvent(evt *nostr.Event) { // and HandleInboundContext does not cause a redelivery on the next start. // The store is persisted to disk; the recent-ID ring absorbs the // inclusive-Since boundary redelivery on restart. + // + // If we cannot persist the state, we must NOT dispatch: a successful but + // un-persisted mark would let the same event be redelivered (and re-handled) + // on the next restart, violating at-most-once inbound delivery. Skip this + // event and let the relay redeliver it later, when persistence may succeed. if c.seen != nil { if err := c.seen.markHandled(chatID, evt.ID, evt.CreatedAt); err != nil { - logger.WarnCF("buzz", "Failed to persist seen state", map[string]any{ + logger.WarnCF("buzz", "Failed to persist seen state; skipping dispatch", map[string]any{ "error": err.Error(), "event_id": evt.ID, + "chat_id": chatID, }) + return } + logger.DebugCF("buzz", "Persisted seen state for event", map[string]any{ + "chat_id": chatID, + "event_id": evt.ID, + "created_at": int64(evt.CreatedAt), + "high_water": int64(c.seen.highWaterMark()), + "seen_path": c.seen.path, + }) } inboundCtx := bus.InboundContext{ diff --git a/pkg/channels/buzz/seen_store.go b/pkg/channels/buzz/seen_store.go index 72079991..c22451a7 100644 --- a/pkg/channels/buzz/seen_store.go +++ b/pkg/channels/buzz/seen_store.go @@ -3,6 +3,7 @@ package buzz import ( "encoding/json" "errors" + "fmt" "os" "path/filepath" "sync" @@ -48,8 +49,11 @@ type seenStore struct { // newSeenStore opens (or creates) the persistent seen-state store at path. // If path is empty it defaults to ~/.picoclaw/buzz/seen-.json. -// A missing file is not an error; the store starts empty. -func newSeenStore(path, defaultName string) *seenStore { +// A missing file is not an error; the store starts empty. Any other read or +// parse failure is returned so callers can refuse to start a channel whose +// replay state cannot be loaded — silently running with a broken store would +// risk redelivering already-handled events. +func newSeenStore(path, defaultName string) (*seenStore, error) { if path == "" { path = defaultSeenStorePath(defaultName) } @@ -57,8 +61,10 @@ func newSeenStore(path, defaultName string) *seenStore { path: path, highWater: make(map[string]*channelHighWater), } - _ = s.load() - return s + if err := s.load(); err != nil { + return nil, fmt.Errorf("buzz seen store load %q: %w", path, err) + } + return s, nil } func defaultSeenStorePath(channelName string) string {