Skip to content

Commit 2d52882

Browse files
committed
[raft/rid] Add rid Raftstore subscriptions
1 parent 2e2ef9b commit 2d52882

5 files changed

Lines changed: 461 additions & 32 deletions

File tree

pkg/rid/application/subscription.go

Lines changed: 37 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -8,15 +8,11 @@ import (
88
dssmodels "github.com/interuss/dss/pkg/models"
99
ridmodels "github.com/interuss/dss/pkg/rid/models"
1010
"github.com/interuss/dss/pkg/rid/repos"
11+
ridraftstore "github.com/interuss/dss/pkg/rid/store/raftstore"
1112
"github.com/interuss/stacktrace"
1213
"go.uber.org/zap"
1314
)
1415

15-
const (
16-
// Defined in requirement DSS0030.
17-
maxSubscriptionsPerArea = 10
18-
)
19-
2016
// SubscriptionApp provides the interface to the application logic for Subscription entities
2117
// AppInterface provides the interface to the application logic for ISA entities
2218
// Note that there is no need for the applciation layer to have the same API as
@@ -60,7 +56,7 @@ func (a *app) InsertSubscription(ctx context.Context, s *ridmodels.Subscription)
6056
return nil, stacktrace.Propagate(err, "Unable to adjust time range")
6157
}
6258
var sub *ridmodels.Subscription
63-
_, err := a.store.Transact(ctx, "", nil, func(ctx context.Context, repo repos.Repository) error {
59+
raftResult, err := a.store.Transact(ctx, ridraftstore.InsertSubscriptionTransaction, s, func(ctx context.Context, repo repos.Repository) error {
6460

6561
// ensure it doesn't exist yet
6662
old, err := repo.GetSubscription(ctx, s.ID)
@@ -78,7 +74,7 @@ func (a *app) InsertSubscription(ctx context.Context, s *ridmodels.Subscription)
7874
return stacktrace.Propagate(err,
7975
"Failed to fetch subscription count, rejecting request")
8076
}
81-
if count >= maxSubscriptionsPerArea {
77+
if count >= ridmodels.MaxSubscriptionsPerArea {
8278
return stacktrace.Propagate(
8379
stacktrace.NewErrorWithCode(dsserr.Exhausted, "Too many existing subscriptions in this area already"),
8480
"%s had %d subscriptions in the area", s.Owner, count)
@@ -91,14 +87,23 @@ func (a *app) InsertSubscription(ctx context.Context, s *ridmodels.Subscription)
9187

9288
return nil
9389
})
90+
91+
if raftResult != nil {
92+
var ok bool
93+
sub, ok = raftResult.(*ridmodels.Subscription)
94+
if !ok {
95+
return nil, stacktrace.NewError("invalid result type: %T", raftResult)
96+
}
97+
}
98+
9499
return sub, err
95100
}
96101

97102
// InsertSubscription implements the App InsertSubscription method
98103
func (a *app) UpdateSubscription(ctx context.Context, s *ridmodels.Subscription) (*ridmodels.Subscription, error) {
99104
var sub *ridmodels.Subscription
100105

101-
_, err := a.store.Transact(ctx, "", nil, func(ctx context.Context, repo repos.Repository) error {
106+
raftResult, err := a.store.Transact(ctx, ridraftstore.UpdateSubscriptionTransaction, s, func(ctx context.Context, repo repos.Repository) error {
102107
old, err := repo.GetSubscription(ctx, s.ID)
103108
switch {
104109
case err != nil:
@@ -128,7 +133,7 @@ func (a *app) UpdateSubscription(ctx context.Context, s *ridmodels.Subscription)
128133
return stacktrace.Propagate(err,
129134
"Failed to fetch subscription count, rejecting request")
130135
}
131-
if count >= maxSubscriptionsPerArea {
136+
if count >= ridmodels.MaxSubscriptionsPerArea {
132137
return stacktrace.Propagate(
133138
stacktrace.NewErrorWithCode(dsserr.Exhausted, "Too many existing subscriptions in this area already"),
134139
"%s had %d subscriptions in the area", s.Owner, count)
@@ -139,13 +144,26 @@ func (a *app) UpdateSubscription(ctx context.Context, s *ridmodels.Subscription)
139144
}
140145
return nil
141146
})
147+
148+
if raftResult != nil {
149+
var ok bool
150+
sub, ok = raftResult.(*ridmodels.Subscription)
151+
if !ok {
152+
return nil, stacktrace.NewError("invalid result type: %T", raftResult)
153+
}
154+
}
155+
142156
return sub, err
143157
}
144158

145159
// DeleteSubscription deletes the Subscription identified by "id" and owned by "owner".
146160
func (a *app) DeleteSubscription(ctx context.Context, id dssmodels.ID, owner dssmodels.Owner, version *dssmodels.Version) (*ridmodels.Subscription, error) {
147161
var ret *ridmodels.Subscription
148-
_, err := a.store.Transact(ctx, "", nil, func(ctx context.Context, repo repos.Repository) error {
162+
raftResult, err := a.store.Transact(ctx, ridraftstore.DeleteSubscriptionTransaction, ridraftstore.DeleteSubscriptionPayload{
163+
ID: id,
164+
Owner: owner,
165+
Version: version,
166+
}, func(ctx context.Context, repo repos.Repository) error {
149167
var err error
150168
old, err := repo.GetSubscription(ctx, id)
151169
switch {
@@ -169,5 +187,14 @@ func (a *app) DeleteSubscription(ctx context.Context, id dssmodels.ID, owner dss
169187
}
170188
return nil
171189
})
190+
191+
if raftResult != nil {
192+
var ok bool
193+
ret, ok = raftResult.(*ridmodels.Subscription)
194+
if !ok {
195+
return nil, stacktrace.NewError("invalid result type: %T", raftResult)
196+
}
197+
}
198+
172199
return ret, err
173200
}

pkg/rid/models/models.go

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,14 @@
11
package models
22

33
import (
4-
"github.com/interuss/stacktrace"
54
"net/url"
5+
6+
"github.com/interuss/stacktrace"
7+
)
8+
9+
const (
10+
// Defined in requirement DSS0030.
11+
MaxSubscriptionsPerArea = 10
612
)
713

814
// ValidateURL ensures https

pkg/rid/store/raftstore/store.go

Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,9 @@ import (
55
"encoding/json"
66
"slices"
77

8+
"github.com/golang/geo/s2"
89
"github.com/interuss/dss/pkg/memstore"
10+
dssmodels "github.com/interuss/dss/pkg/models"
911
"github.com/interuss/dss/pkg/raftstore"
1012
"github.com/interuss/dss/pkg/raftstore/consensus"
1113
ridmodels "github.com/interuss/dss/pkg/rid/models"
@@ -30,13 +32,35 @@ const (
3032
DeleteISATransaction raftstore.RequestType = "deleteISATransaction"
3133
InsertISATransaction raftstore.RequestType = "insertISATransaction"
3234
UpdateISATransaction raftstore.RequestType = "updateISATransaction"
35+
36+
getSubscription raftstore.RequestType = "getSubscription"
37+
deleteSubscription raftstore.RequestType = "deleteSubscription"
38+
insertSubscription raftstore.RequestType = "insertSubscription"
39+
updateSubscription raftstore.RequestType = "updateSubscription"
40+
searchSubscriptions raftstore.RequestType = "searchSubscriptions"
41+
searchSubscriptionsByOwner raftstore.RequestType = "searchSubscriptionsByOwner"
42+
updateNotificationIdxsInCells raftstore.RequestType = "updateNotificationIdxsInCells"
43+
maxSubscriptionCountInCellsByOwner raftstore.RequestType = "maxSubscriptionCountInCellsByOwner"
44+
listExpiredSubscriptions raftstore.RequestType = "listExpiredSubscriptions"
45+
countSubscriptions raftstore.RequestType = "countSubscriptions"
46+
47+
DeleteSubscriptionTransaction raftstore.RequestType = "deleteSubscriptionTransaction"
48+
InsertSubscriptionTransaction raftstore.RequestType = "insertSubscriptionTransaction"
49+
UpdateSubscriptionTransaction raftstore.RequestType = "updateSubscriptionTransaction"
3350
)
3451

3552
var readOnlyRequests = []raftstore.RequestType{
3653
getISA,
3754
searchISAs,
3855
listExpiredISAs,
3956
countISAs,
57+
58+
getSubscription,
59+
searchSubscriptions,
60+
searchSubscriptionsByOwner,
61+
maxSubscriptionCountInCellsByOwner,
62+
listExpiredSubscriptions,
63+
countSubscriptions,
4064
}
4165

4266
// repo is a full implementation of rid.repos.Repository for Raft-based storage.
@@ -144,6 +168,83 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err
144168
case UpdateISATransaction:
145169
return r.updateISATransactionApplier(ctx, proposal, mem)
146170

171+
// Subscriptions
172+
173+
case getSubscription:
174+
var id dssmodels.ID
175+
if err := json.Unmarshal(proposal.Value, &id); err != nil {
176+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", getSubscription)
177+
}
178+
return mem.GetSubscription(ctx, id)
179+
180+
case deleteSubscription:
181+
var sub ridmodels.Subscription
182+
if err := json.Unmarshal(proposal.Value, &sub); err != nil {
183+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", deleteSubscription)
184+
}
185+
return mem.DeleteSubscription(ctx, &sub)
186+
187+
case insertSubscription:
188+
var sub ridmodels.Subscription
189+
if err := json.Unmarshal(proposal.Value, &sub); err != nil {
190+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", insertSubscription)
191+
}
192+
return mem.InsertSubscription(ctx, &sub)
193+
194+
case updateSubscription:
195+
var sub ridmodels.Subscription
196+
if err := json.Unmarshal(proposal.Value, &sub); err != nil {
197+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", updateSubscription)
198+
}
199+
return mem.UpdateSubscription(ctx, &sub)
200+
201+
case searchSubscriptions:
202+
var cells s2.CellUnion
203+
if err := json.Unmarshal(proposal.Value, &cells); err != nil {
204+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", searchSubscriptions)
205+
}
206+
return mem.SearchSubscriptions(ctx, cells)
207+
208+
case searchSubscriptionsByOwner:
209+
var p searchSubscriptionsByOwnerPayload
210+
if err := json.Unmarshal(proposal.Value, &p); err != nil {
211+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", searchSubscriptionsByOwner)
212+
}
213+
return mem.SearchSubscriptionsByOwner(ctx, p.Cells, p.Owner)
214+
215+
case updateNotificationIdxsInCells:
216+
var cells s2.CellUnion
217+
if err := json.Unmarshal(proposal.Value, &cells); err != nil {
218+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", updateNotificationIdxsInCells)
219+
}
220+
return mem.UpdateNotificationIdxsInCells(ctx, cells)
221+
222+
case maxSubscriptionCountInCellsByOwner:
223+
var p maxSubscriptionCountInCellsByOwnerPayload
224+
if err := json.Unmarshal(proposal.Value, &p); err != nil {
225+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", maxSubscriptionCountInCellsByOwner)
226+
}
227+
return mem.MaxSubscriptionCountInCellsByOwner(ctx, p.Cells, p.Owner)
228+
229+
case listExpiredSubscriptions:
230+
var p listExpiredSubscriptionsPayload
231+
if err := json.Unmarshal(proposal.Value, &p); err != nil {
232+
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", listExpiredSubscriptions)
233+
}
234+
return mem.ListExpiredSubscriptions(ctx, p.Writer, p.Threshold)
235+
236+
case countSubscriptions:
237+
return mem.CountSubscriptions(ctx)
238+
239+
case DeleteSubscriptionTransaction:
240+
return r.deleteSubscriptionTransactionApplier(ctx, proposal, mem)
241+
242+
case InsertSubscriptionTransaction:
243+
return r.insertSubscriptionTransactionApplier(ctx, proposal, mem)
244+
245+
case UpdateSubscriptionTransaction:
246+
return r.updateSubscriptionTransactionApplier(ctx, proposal, mem)
247+
147248
default:
148249
return nil, stacktrace.NewError("unknown request type: %q", proposal.RequestType)
149250
}

0 commit comments

Comments
 (0)