feat: apply durable event acknowledgements
This commit is contained in:
@@ -0,0 +1,29 @@
|
|||||||
|
package agent
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
|
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
|
||||||
|
"github.com/rvbox/rvbox/internal/domain"
|
||||||
|
)
|
||||||
|
|
||||||
|
var ErrInvalidEventAck = errors.New("invalid event acknowledgement")
|
||||||
|
|
||||||
|
type EventAckStore interface {
|
||||||
|
Ack(context.Context, domain.UUID, uint64) error
|
||||||
|
}
|
||||||
|
|
||||||
|
// ApplyEventAck releases only the cumulative, server-confirmed event prefix.
|
||||||
|
// The spool remains the source of truth if the process stops before this call
|
||||||
|
// returns; a repeated acknowledgement is intentionally harmless.
|
||||||
|
func ApplyEventAck(ctx context.Context, store EventAckStore, ack *rvboxv1.EventAck) error {
|
||||||
|
if store == nil || ack == nil || ack.GetThroughEventSeq() == 0 {
|
||||||
|
return ErrInvalidEventAck
|
||||||
|
}
|
||||||
|
issue, err := domain.ParseUUIDv7(ack.GetIssueUuid())
|
||||||
|
if err != nil {
|
||||||
|
return ErrInvalidEventAck
|
||||||
|
}
|
||||||
|
return store.Ack(ctx, issue, ack.GetThroughEventSeq())
|
||||||
|
}
|
||||||
@@ -118,6 +118,17 @@ func TestSendCommandEventCanonicalDigest_HP_EVENT_02(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestApplyEventAckValidatesAndReleasesPrefix_HP_EVENT_03(t *testing.T) {
|
||||||
|
store := &ackStore{}
|
||||||
|
ack := &rvboxv1.EventAck{IssueUuid: "019c46f1-1d02-7000-8000-000000000066", ThroughEventSeq: 4}
|
||||||
|
if err := ApplyEventAck(context.Background(), store, ack); err != nil || store.issue.String() != ack.IssueUuid || store.sequence != 4 {
|
||||||
|
t.Fatalf("ApplyEventAck = issue=%s sequence=%d err=%v", store.issue, store.sequence, err)
|
||||||
|
}
|
||||||
|
if err := ApplyEventAck(context.Background(), store, &rvboxv1.EventAck{IssueUuid: "bad", ThroughEventSeq: 1}); !errors.Is(err, ErrInvalidEventAck) {
|
||||||
|
t.Fatalf("invalid ack error = %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
type fakeTransport struct {
|
type fakeTransport struct {
|
||||||
written []byte
|
written []byte
|
||||||
read []byte
|
read []byte
|
||||||
@@ -130,6 +141,16 @@ func (store *reconcileStore) DiscardTerminal(_ context.Context, issue domain.UUI
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type ackStore struct {
|
||||||
|
issue domain.UUID
|
||||||
|
sequence uint64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (store *ackStore) Ack(_ context.Context, issue domain.UUID, sequence uint64) error {
|
||||||
|
store.issue, store.sequence = issue, sequence
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func (transport *fakeTransport) Write(_ context.Context, value []byte) error {
|
func (transport *fakeTransport) Write(_ context.Context, value []byte) error {
|
||||||
transport.written = append([]byte(nil), value...)
|
transport.written = append([]byte(nil), value...)
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -230,6 +230,12 @@ layer = "unit"
|
|||||||
status = "implemented"
|
status = "implemented"
|
||||||
tests = ["internal/client/agent/handshake_test.go:TestSendCommandEventCanonicalDigest_HP_EVENT_02"]
|
tests = ["internal/client/agent/handshake_test.go:TestSendCommandEventCanonicalDigest_HP_EVENT_02"]
|
||||||
|
|
||||||
|
[[requirements]]
|
||||||
|
id = "HP-EVENT-03"
|
||||||
|
layer = "unit"
|
||||||
|
status = "implemented"
|
||||||
|
tests = ["internal/client/agent/handshake_test.go:TestApplyEventAckValidatesAndReleasesPrefix_HP_EVENT_03"]
|
||||||
|
|
||||||
[[requirements]]
|
[[requirements]]
|
||||||
id = "HP-SES-05"
|
id = "HP-SES-05"
|
||||||
layer = "integration"
|
layer = "integration"
|
||||||
|
|||||||
Reference in New Issue
Block a user