From f6496f883a9450793b51e5d7e0495ad538abbcb2 Mon Sep 17 00:00:00 2001 From: Mariem Baccari Date: Mon, 21 Sep 2026 22:49:40 +0200 Subject: [PATCH] [raft/scd] Implement operational intents repo methods --- .../store/raftstore/operational_intents.go | 72 +++++++++++++++---- pkg/scd/store/raftstore/store.go | 53 ++++++++++++++ 2 files changed, 110 insertions(+), 15 deletions(-) diff --git a/pkg/scd/store/raftstore/operational_intents.go b/pkg/scd/store/raftstore/operational_intents.go index 110ffd008..7817fb703 100644 --- a/pkg/scd/store/raftstore/operational_intents.go +++ b/pkg/scd/store/raftstore/operational_intents.go @@ -2,38 +2,80 @@ package raftstore import ( "context" + "encoding/json" "time" - dsserr "github.com/interuss/dss/pkg/errors" dssmodels "github.com/interuss/dss/pkg/models" + "github.com/interuss/dss/pkg/raftstore/consensus" scdmodels "github.com/interuss/dss/pkg/scd/models" "github.com/interuss/stacktrace" ) -func (r *repo) GetOperationalIntent(_ context.Context, id dssmodels.ID) (*scdmodels.OperationalIntent, error) { - return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "GetOperationalIntent not implemented for raftstore") +const ( + getOperationalIntent consensus.RequestType[*scdmodels.OperationalIntent] = "getOperationalIntent" + deleteOperationalIntent consensus.RequestType[any] = "deleteOperationalIntent" + upsertOperationalIntent consensus.RequestType[*scdmodels.OperationalIntent] = "upsertOperationalIntent" + searchOperationalIntents consensus.RequestType[[]*scdmodels.OperationalIntent] = "searchOperationalIntents" + getDependentOperationalIntents consensus.RequestType[[]dssmodels.ID] = "getDependentOperationalIntents" + listExpiredOperationalIntents consensus.RequestType[[]*scdmodels.OperationalIntent] = "listExpiredOperationalIntents" + countOperationalIntents consensus.RequestType[int64] = "countOperationalIntents" +) + +func (r *repo) GetOperationalIntent(ctx context.Context, id dssmodels.ID) (*scdmodels.OperationalIntent, error) { + buf, err := json.Marshal(id) + if err != nil { + return nil, stacktrace.Propagate(err, "failed to marshal payload") + } + + return r.consensus.HandleClientRequest(ctx, getOperationalIntent, buf, true) } -func (r *repo) DeleteOperationalIntent(_ context.Context, id dssmodels.ID) error { - return stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteOperationalIntent not implemented for raftstore") +func (r *repo) DeleteOperationalIntent(ctx context.Context, id dssmodels.ID) error { + buf, err := json.Marshal(id) + if err != nil { + return stacktrace.Propagate(err, "failed to marshal payload") + } + + _, err = r.consensus.HandleClientRequest(ctx, deleteOperationalIntent, buf, false) + return err } -func (r *repo) UpsertOperationalIntent(_ context.Context, operation *scdmodels.OperationalIntent) (*scdmodels.OperationalIntent, error) { - return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "UpsertOperationalIntent not implemented for raftstore") +func (r *repo) UpsertOperationalIntent(ctx context.Context, operation *scdmodels.OperationalIntent) (*scdmodels.OperationalIntent, error) { + buf, err := json.Marshal(operation) + if err != nil { + return nil, stacktrace.Propagate(err, "failed to marshal payload") + } + + return r.consensus.HandleClientRequest(ctx, upsertOperationalIntent, buf, false) } -func (r *repo) SearchOperationalIntents(_ context.Context, cellsVolume *dssmodels.CellsVolume4D) ([]*scdmodels.OperationalIntent, error) { - return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "SearchOperationalIntents not implemented for raftstore") +func (r *repo) SearchOperationalIntents(ctx context.Context, cellsVolume *dssmodels.CellsVolume4D) ([]*scdmodels.OperationalIntent, error) { + buf, err := json.Marshal(cellsVolume) + if err != nil { + return nil, stacktrace.Propagate(err, "failed to marshal payload") + } + + return r.consensus.HandleClientRequest(ctx, searchOperationalIntents, buf, true) } -func (r *repo) GetDependentOperationalIntents(_ context.Context, subscriptionID dssmodels.ID) ([]dssmodels.ID, error) { - return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "GetDependentOperationalIntents not implemented for raftstore") +func (r *repo) GetDependentOperationalIntents(ctx context.Context, subscriptionID dssmodels.ID) ([]dssmodels.ID, error) { + buf, err := json.Marshal(subscriptionID) + if err != nil { + return nil, stacktrace.Propagate(err, "failed to marshal payload") + } + + return r.consensus.HandleClientRequest(ctx, getDependentOperationalIntents, buf, true) } -func (r *repo) ListExpiredOperationalIntents(_ context.Context, threshold time.Time) ([]*scdmodels.OperationalIntent, error) { - return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "ListExpiredOperationalIntents not implemented for raftstore") +func (r *repo) ListExpiredOperationalIntents(ctx context.Context, threshold time.Time) ([]*scdmodels.OperationalIntent, error) { + buf, err := json.Marshal(threshold) + if err != nil { + return nil, stacktrace.Propagate(err, "failed to marshal payload") + } + + return r.consensus.HandleClientRequest(ctx, listExpiredOperationalIntents, buf, true) } -func (r *repo) CountOperationalIntents(_ context.Context) (int64, error) { - return 0, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "CountOperationalIntents not implemented for raftstore") +func (r *repo) CountOperationalIntents(ctx context.Context) (int64, error) { + return r.consensus.HandleClientRequest(ctx, countOperationalIntents, nil, true) } diff --git a/pkg/scd/store/raftstore/store.go b/pkg/scd/store/raftstore/store.go index e5b2211c7..8686edd36 100644 --- a/pkg/scd/store/raftstore/store.go +++ b/pkg/scd/store/raftstore/store.go @@ -50,6 +50,9 @@ func (r *repo) GetRepo() repos.Repository { return r } func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, error) { switch proposal.RequestType { + + // Constraints + case string(searchConstraints): var cellsVolume dssmodels.CellsVolume4D if err := json.Unmarshal(proposal.Value, &cellsVolume); err != nil { @@ -81,6 +84,8 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err case string(countConstraints): return r.Store.GetRepo().CountConstraints(ctx) + // Subscriptions + case string(searchSubscriptions): var cellsVolume dssmodels.CellsVolume4D if err := json.Unmarshal(proposal.Value, &cellsVolume); err != nil { @@ -133,6 +138,54 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err case string(countSubscriptions): return r.Store.GetRepo().CountSubscriptions(ctx) + // Operational Intents + + case string(getOperationalIntent): + var id dssmodels.ID + if err := json.Unmarshal(proposal.Value, &id); err != nil { + return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", getOperationalIntent) + } + return r.Store.GetRepo().GetOperationalIntent(ctx, id) + + case string(deleteOperationalIntent): + var id dssmodels.ID + if err := json.Unmarshal(proposal.Value, &id); err != nil { + return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", deleteOperationalIntent) + } + return nil, r.Store.GetRepo().DeleteOperationalIntent(ctx, id) + + case string(upsertOperationalIntent): + var operation scdmodels.OperationalIntent + if err := json.Unmarshal(proposal.Value, &operation); err != nil { + return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", upsertOperationalIntent) + } + return r.Store.GetRepo().UpsertOperationalIntent(ctx, &operation) + + case string(searchOperationalIntents): + var cellsVolume dssmodels.CellsVolume4D + if err := json.Unmarshal(proposal.Value, &cellsVolume); err != nil { + return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", searchOperationalIntents) + } + return r.Store.GetRepo().SearchOperationalIntents(ctx, &cellsVolume) + + case string(getDependentOperationalIntents): + var subscriptionID dssmodels.ID + if err := json.Unmarshal(proposal.Value, &subscriptionID); err != nil { + return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", getDependentOperationalIntents) + } + return r.Store.GetRepo().GetDependentOperationalIntents(ctx, subscriptionID) + + case string(listExpiredOperationalIntents): + var threshold time.Time + if err := json.Unmarshal(proposal.Value, &threshold); err != nil { + return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", listExpiredOperationalIntents) + } + return r.Store.GetRepo().ListExpiredOperationalIntents(ctx, threshold) + + case string(countOperationalIntents): + return r.Store.GetRepo().CountOperationalIntents(ctx) + + // Operations registry (transactions) default: handler, ok := operations.Registry[proposal.RequestType] if !ok {