Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 57 additions & 15 deletions pkg/scd/store/raftstore/operational_intents.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
Comment on lines +25 to +28

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not for this PR, but with the look of it, looks like we could be pushing the marshalling of the value into HandleClientRequest and change the type of value to any. Except if some operations have more complex things going on?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With the gate 1 work (#1709), I will look into payload size optimizations and that should include investigating custom encoding per payload type. I kept the marshalling outside so that consensus can receive an opaque payload like it would with custom encoding.
I referenced your suggestion in the contributions spreadsheet to keep it in mind in case we don't end up going with custom encoding for some reason. Is it okay if I revisit later ?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sounds good 👍


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)
}
53 changes: 53 additions & 0 deletions pkg/scd/store/raftstore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down
Loading