From 94c3f6dcb4ce078aea959a157daa806f635311f9 Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 09:18:37 +0000 Subject: [PATCH] fix: fence command events by dispatch generation --- internal/server/session/agent_server.go | 6 ++-- internal/server/session/agent_server_test.go | 4 +-- internal/server/store/event.go | 31 ++++++++++++-------- 3 files changed, 23 insertions(+), 18 deletions(-) diff --git a/internal/server/session/agent_server.go b/internal/server/session/agent_server.go index b0b78f1..13d9ac1 100644 --- a/internal/server/session/agent_server.go +++ b/internal/server/session/agent_server.go @@ -244,7 +244,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w continue } if event := envelope.GetCommandEvent(); event != nil { - appendEvent, eventErr := eventAppendFromWire(event, hello.GetClientId(), server.now()) + appendEvent, eventErr := eventAppendFromWire(event, hello.GetClientId(), registration.Generation, server.now()) if eventErr != nil { server.close(connection, websocket.StatusPolicyViolation, "invalid command event") return @@ -263,7 +263,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w } } -func eventAppendFromWire(event *rvboxv1.CommandEvent, clientID string, receipt time.Time) (store.EventAppend, error) { +func eventAppendFromWire(event *rvboxv1.CommandEvent, clientID string, generation uint64, receipt time.Time) (store.EventAppend, error) { issue, err := domain.ParseUUIDv7(event.GetIssueUuid()) if err != nil { return store.EventAppend{}, err @@ -274,7 +274,7 @@ func eventAppendFromWire(event *rvboxv1.CommandEvent, clientID string, receipt t } var immutable [32]byte copy(immutable[:], event.GetImmutableEventSha256()) - result := store.EventAppend{IssueUUID: [16]byte(issue), ClientID: clientID, EventSeq: event.GetEventSeq(), ObservedUnixNano: event.GetObservedAt().AsTime().UnixNano(), ReceiptUnixNano: receipt.UnixNano(), EventType: eventType(event), Compression: 1, RawLength: uint64(len(payload)), Payload: payload, ImmutableSHA256: immutable} + result := store.EventAppend{IssueUUID: [16]byte(issue), ClientID: clientID, SessionGeneration: generation, EventSeq: event.GetEventSeq(), ObservedUnixNano: event.GetObservedAt().AsTime().UnixNano(), ReceiptUnixNano: receipt.UnixNano(), EventType: eventType(event), Compression: 1, RawLength: uint64(len(payload)), Payload: payload, ImmutableSHA256: immutable} if output := event.GetOutput(); output != nil { result.Stream = uint16(output.GetStream()) result.Output = true diff --git a/internal/server/session/agent_server_test.go b/internal/server/session/agent_server_test.go index 050a170..6564489 100644 --- a/internal/server/session/agent_server_test.go +++ b/internal/server/session/agent_server_test.go @@ -116,8 +116,8 @@ func TestWireEventAppendCarriesClientBinding_HP_EVENT_01(t *testing.T) { t.Fatal(err) } event.ImmutableEventSha256 = digest[:] - appendEvent, err := eventAppendFromWire(event, "client-a", time.Now()) - if err != nil || appendEvent.ClientID != "client-a" || appendEvent.EventSeq != 1 || appendEvent.EventType != 4 || appendEvent.ImmutableSHA256 != digest { + appendEvent, err := eventAppendFromWire(event, "client-a", 7, time.Now()) + if err != nil || appendEvent.ClientID != "client-a" || appendEvent.SessionGeneration != 7 || appendEvent.EventSeq != 1 || appendEvent.EventType != 4 || appendEvent.ImmutableSHA256 != digest { t.Fatalf("wire event append = %#v, %v", appendEvent, err) } } diff --git a/internal/server/store/event.go b/internal/server/store/event.go index 738a3b3..d5b2d89 100644 --- a/internal/server/store/event.go +++ b/internal/server/store/event.go @@ -32,19 +32,20 @@ type FaultInjectorFunc func(name string) error func (function FaultInjectorFunc) Checkpoint(name string) error { return function(name) } type EventAppend struct { - IssueUUID [16]byte - ClientID string - EventSeq uint64 - ObservedUnixNano int64 - ReceiptUnixNano int64 - EventType uint16 - Stream uint16 - Compression uint16 - RawLength uint64 - Payload []byte - ImmutableSHA256 [32]byte - Output bool - UseCloseout bool + IssueUUID [16]byte + ClientID string + SessionGeneration uint64 + EventSeq uint64 + ObservedUnixNano int64 + ReceiptUnixNano int64 + EventType uint16 + Stream uint16 + Compression uint16 + RawLength uint64 + Payload []byte + ImmutableSHA256 [32]byte + Output bool + UseCloseout bool } type EventAppendResult struct { @@ -91,6 +92,10 @@ JOIN storage_counters ON storage_counters.singleton = 1 WHERE commands.issue_uui query += ` AND commands.client_id = ?` arguments = append(arguments, event.ClientID) } + if event.SessionGeneration != 0 { + query += ` AND commands.target_session_generation = ?` + arguments = append(arguments, event.SessionGeneration) + } err = database.QueryRowContext(ctx, query, arguments...).Scan( &lastSequence, &commandOutputCharged, &commandCharged, &closeoutRemaining, &clientID, &clientCharged, &serverCharged) if errors.Is(err, sql.ErrNoRows) {