Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
d17c86d
[raft/scd] Extract update and create opintent
MariemBaccari Aug 21, 2026
0c04bf8
[scd] Return conflict without error
MariemBaccari Sep 9, 2026
8fe495a
[raftstore] Rename actions packages to operations
MariemBaccari Aug 28, 2026
49b4f3a
[raftstore] Unexport operations
MariemBaccari Aug 28, 2026
e724a41
[raftstore] Rename context methods
MariemBaccari Aug 28, 2026
cf95e9c
[raftstore] Embed memstore and use checkpoint
MariemBaccari Aug 19, 2026
fb30ae4
[raft/scd] Implement constraints repo methods
MariemBaccari Aug 28, 2026
0bdb05b
[raft/scd] Implement subscriptions repo methods
MariemBaccari Aug 28, 2026
c7ec73d
[raft/scd] Implement operational intents repo methods
MariemBaccari Aug 28, 2026
d12d152
[raft/scd] Implement availability repo methods
MariemBaccari Aug 28, 2026
69ba1e7
[raft/rid] Implement ISA repo methods
MariemBaccari Aug 28, 2026
1abf559
[raft/rid] Implement subscriptions repo methods
MariemBaccari Aug 28, 2026
af323f4
[raft] Fix start action missing required context
the-glu Sep 11, 2026
1bbe0d9
[raft/rid] Extract InsertSubscription and route Create through the store
MariemBaccari Aug 31, 2026
30a2ff6
[raft/rid] Extract UpdateSubscription and route Update through the store
MariemBaccari Aug 31, 2026
6814f24
[raft/rid] Extract DeleteISA and route Delete through the store
MariemBaccari Aug 31, 2026
17a68f3
[raft/rid] Extract InsertISA and route Create through the store
MariemBaccari Aug 31, 2026
57cc09f
[raft/rid] Extract UpdateISA and route Update through the store
MariemBaccari Aug 31, 2026
16d9599
[raft/rid] Extract Get/Search and remove application layer
MariemBaccari Sep 15, 2026
af4ce63
[raft] Make raft runnable-locally
the-glu Jun 16, 2026
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
1 change: 0 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,6 @@ test-go-units-crdb: cleanup-test-go-units-crdb
go run ./cmds/db-manager/main.go migrate --schemas_dir ./build/db_schemas/scd --db_version latest --datastore_host localhost
go run ./cmds/db-manager/main.go migrate --schemas_dir ./build/db_schemas/aux_ --db_version latest --datastore_host localhost
go test -cover -count=1 -v ./pkg/rid/store/sqlstore --datastore_host localhost --datastore_port 26257 --datastore_ssl_mode disable --datastore_user root -test.gocoverdir=$(COVERDATA_DIR)
go test -cover -count=1 -v ./pkg/rid/application --datastore_host localhost --datastore_port 26257 --datastore_ssl_mode disable --datastore_user root -test.gocoverdir=$(COVERDATA_DIR)
go test -cover -count=1 -v ./pkg/scd/store/sqlstore --datastore_host localhost --datastore_port 26257 --datastore_ssl_mode disable --datastore_user root -test.gocoverdir=$(COVERDATA_DIR)
go test -cover -count=1 -v ./pkg/aux_/store/sqlstore --datastore_host localhost --datastore_port 26257 --datastore_ssl_mode disable --datastore_user root -test.gocoverdir=$(COVERDATA_DIR)
@docker stop dss-crdb-for-testing > /dev/null
Expand Down
4 changes: 3 additions & 1 deletion build/dev/docker-compose_dss.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,7 @@ services:
- $PWD/../test-certs:/var/test-certs:ro
- $PWD/startup/core_service.sh:/startup/core_service.sh:ro
- $PWD/startup/coverdata:/startup/coverdata:rw # we will save coverage info here
- raftdata:/raftdata
environment:
COMPOSE_PROFILES: ${COMPOSE_PROFILES}
# Note: requires the Dockerfile to have been built with "-cover" in the EXTRA_GO_INSTALL_FLAGS var
Expand Down Expand Up @@ -147,7 +148,7 @@ services:
interval: 3m
start_period: 30s
start_interval: 5s
profiles: ["", "with-yugabyte", "with-monitoring"]
profiles: ["", "with-yugabyte", "with-monitoring", "with-raft"]

local-dss-dummy-oauth:
build:
Expand Down Expand Up @@ -216,3 +217,4 @@ networks:

volumes:
local-dss-data:
raftdata:
9 changes: 7 additions & 2 deletions build/dev/startup/core_service.sh
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,13 @@ if [ "${COMPOSE_PROFILES#*"with-yugabyte"}" != "${COMPOSE_PROFILES}" ]; then
echo "Using Yugabyte"
DATASTORE_CONNECTION="-datastore_host local-dss-ybdb -datastore_user yugabyte --datastore_port 5433"
else
echo "Using CockroachDB"
DATASTORE_CONNECTION="-datastore_host local-dss-crdb"
if [ "${COMPOSE_PROFILES#*"with-raft"}" != "${COMPOSE_PROFILES}" ]; then
echo "Using raft"
DATASTORE_CONNECTION="-store_type raft -raft_node_id=1 -rid_raft_peers=1=http://127.0.0.1:9011 -scd_raft_peers=1=http://127.0.0.1:9021 -aux_raft_peers=1=http://127.0.0.1:9031 -raft_datadir /raftdata"
else
echo "Using CockroachDB"
DATASTORE_CONNECTION="-datastore_host local-dss-crdb"
fi
fi

if [ "${COMPOSE_PROFILES#*"with-monitoring"}" != "${COMPOSE_PROFILES}" ]; then
Expand Down
17 changes: 11 additions & 6 deletions cmds/core-service/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@ import (
requestlocality "github.com/interuss/dss/pkg/locality"
"github.com/interuss/dss/pkg/logging"
"github.com/interuss/dss/pkg/random"
"github.com/interuss/dss/pkg/rid/application"
rid_v1 "github.com/interuss/dss/pkg/rid/server/v1"
rid_v2 "github.com/interuss/dss/pkg/rid/server/v2"
rids "github.com/interuss/dss/pkg/rid/store"
Expand Down Expand Up @@ -112,6 +111,15 @@ func createAuxServer(ctx context.Context, locality string, publicEndpoint string
return nil, stacktrace.Propagate(err, "Unable to interact with store")
}

ctx = timestamp.NewContext(ctx, time.Now())

seed, err := random.NewSeed()
if err != nil {
return nil, stacktrace.Propagate(err, "Unable to generate seed")
}

ctx = random.NewContext(ctx, seed)

err = repo.SaveOwnMetadata(ctx, locality, publicEndpoint)

if err != nil {
Expand Down Expand Up @@ -141,15 +149,12 @@ func createRIDServers(ctx context.Context, locality string, logger *zap.Logger)
}
}

app := application.NewFromTransactor(ridStore, logger)
return &rid_v1.Server{
Store: ridStore,
App: app,
Locality: locality,
AllowHTTPBaseUrls: *allowHTTPBaseUrls,
}, &rid_v2.Server{
Store: ridStore,
App: app,
Locality: locality,
AllowHTTPBaseUrls: *allowHTTPBaseUrls,
}, nil
Expand Down Expand Up @@ -368,9 +373,9 @@ func RunHTTPServer(ctx context.Context, ctxCanceler func(), address, locality st
handler = authorizer.TokenMiddleware(handler)
handler = http.TimeoutHandler(handler, *timeout, "request timeout")
handler = logging.HTTPMiddleware(logger, *dumpRequests, handler)
handler = timestamp.RequestTimestampMiddleware(handler)
handler = timestamp.Middleware(handler)
handler = random.Middleware(handler)
handler = requestlocality.LocalityMiddleware(locality)(handler)
handler = requestlocality.Middleware(locality)(handler)

if *enableMetrics || *enableTracing {
// We use the default settings; the APIRouter handler will override the span value accordingly, as it has more information.
Expand Down
2 changes: 1 addition & 1 deletion pkg/aux_/pool_participants.go
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ func (a *Server) PutDSSInstancesHeartbeat(ctx context.Context, req *restapi.PutD
}
heartbeat.Timestamp = &ts
} else {
now := timestamp.MustGetRequestTimestamp(ctx)
now := timestamp.MustFromContext(ctx)
heartbeat.Timestamp = &now
}

Expand Down
2 changes: 1 addition & 1 deletion pkg/aux_/store/memstore/dss.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import (
)

func (r *repo) SaveOwnMetadata(ctx context.Context, loc string, publicEndpoint string) error {
now := timestamp.MustGetRequestTimestamp(ctx)
now := timestamp.MustFromContext(ctx)

r.state.Participants[locality(loc)] = &participant{
PublicEndpoint: publicEndpoint,
Expand Down
8 changes: 4 additions & 4 deletions pkg/aux_/store/memstore/dss_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ var fakeClock = clockwork.NewFakeClock()

func TestSaveOwnMetadataRoundTrip(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
r := newRepo()

require.NoError(t, r.SaveOwnMetadata(ctx, "dss-1", "https://example.com"))
Expand All @@ -35,7 +35,7 @@ func TestSaveOwnMetadataRoundTrip(t *testing.T) {

func TestSaveOwnMetadataUpsert(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
r := newRepo()

require.NoError(t, r.SaveOwnMetadata(ctx, "dss-1", "https://old.example.com"))
Expand All @@ -50,7 +50,7 @@ func TestSaveOwnMetadataUpsert(t *testing.T) {

func TestGetDSSMetadataPicksLatestHeartbeat(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
r := newRepo()

require.NoError(t, r.SaveOwnMetadata(ctx, "dss-1", "https://example.com"))
Expand All @@ -71,7 +71,7 @@ func TestGetDSSMetadataPicksLatestHeartbeat(t *testing.T) {

func TestGetDSSMetadataUpdatesHeartbeatPerSource(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
r := newRepo()

require.NoError(t, r.SaveOwnMetadata(ctx, "dss-1", "https://example.com"))
Expand Down
4 changes: 2 additions & 2 deletions pkg/aux_/store/memstore/snapshot_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ import (

func TestSnapshotRoundTrip(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
src := newRepo()
require.NoError(t, src.SaveOwnMetadata(ctx, "dss-1", "https://example.com"))
ts := time.Now().UTC()
Expand All @@ -40,7 +40,7 @@ func TestSnapshotRoundTrip(t *testing.T) {

func TestRestoreFromSnapshotReplacesState(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
src := newRepo()
require.NoError(t, src.SaveOwnMetadata(ctx, "dss-1", "https://example.com"))
data, err := src.GetSnapshot()
Expand Down
4 changes: 2 additions & 2 deletions pkg/aux_/store/memstore/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ import (

func TestCheckpointRestore(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())

r := newRepo()

Expand All @@ -34,7 +34,7 @@ func TestCheckpointRestore(t *testing.T) {

func TestCheckpointIsolatesUpsert(t *testing.T) {
ctx := context.Background()
ctx = timestamp.WithRequestTimestamp(ctx, fakeClock.Now())
ctx = timestamp.NewContext(ctx, fakeClock.Now())
r := newRepo()

require.NoError(t, r.SaveOwnMetadata(ctx, "dss-1", "https://old.example.com"))
Expand Down
19 changes: 5 additions & 14 deletions pkg/aux_/store/raftstore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,7 @@ const (
// repo is a full implementation of aux_.repos.Repository for Raft-based storage.
type repo struct {
consensus *consensus.Consensus
memStore *memstore.Store[repos.Repository]
memRepo repos.Repository
*memstore.Store[repos.Repository]
}

func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.Store[repos.Repository], error) {
Expand All @@ -39,7 +38,7 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.
return nil, stacktrace.Propagate(err, "failed to initialize aux memstore")
}

r := &repo{memStore: memStore, memRepo: memStore.GetRepo()}
r := &repo{Store: memStore}
store, err := raftstore.Init(ctx, logger.With(zap.String("service", "aux_")), locality, params, r, nil)
if err != nil {
return nil, stacktrace.Propagate(err, "failed to initialize aux raftstore")
Expand All @@ -52,14 +51,6 @@ func Init(ctx context.Context, logger *zap.Logger, locality string) (*raftstore.

func (r *repo) GetRepo() repos.Repository { return r }

func (r *repo) GetSnapshot() ([]byte, error) {
return r.memStore.GetSnapshot()
}

func (r *repo) RestoreFromSnapshot(data []byte) error {
return r.memStore.RestoreFromSnapshot(data)
}

func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, error) {
switch proposal.RequestType {
case saveOwnMetadata:
Expand All @@ -68,18 +59,18 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", saveOwnMetadata)
}

return nil, r.memRepo.SaveOwnMetadata(ctx, payload.Locality, payload.PublicEndpoint)
return nil, r.Store.GetRepo().SaveOwnMetadata(ctx, payload.Locality, payload.PublicEndpoint)

case getDSSMetadata:
return r.memRepo.GetDSSMetadata(ctx)
return r.Store.GetRepo().GetDSSMetadata(ctx)

case recordHeartbeat:
var heartbeat auxmodels.Heartbeat
if err := json.Unmarshal(proposal.Value, &heartbeat); err != nil {
return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", recordHeartbeat)
}

return nil, r.memRepo.RecordHeartbeat(ctx, heartbeat)
return nil, r.Store.GetRepo().RecordHeartbeat(ctx, heartbeat)

default:
return nil, stacktrace.NewError("unknown request type: %q", proposal.RequestType)
Expand Down
4 changes: 0 additions & 4 deletions pkg/errors/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,6 @@ const (
// larger than the max area allowed. See geo/s2.go.
AreaTooLarge = stacktrace.ErrorCode(iota)

// MissingOVNs is the error to signal that an AirspaceConflictResponse should
// be returned rather than the standard error response.
MissingOVNs

// AlreadyExists is used when attempting to create a resource that already
// exists.
AlreadyExists
Expand Down
20 changes: 10 additions & 10 deletions pkg/locality/locality.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,30 +7,30 @@ import (
"github.com/interuss/stacktrace"
)

type localityKey struct{}
type key struct{}

// MustGetRequestLocality returns the request locality from the context and panics if it is not
// MustFromContext returns the request locality from the context and panics if it is not
// present, which is a programming error.
func MustGetRequestLocality(ctx context.Context) string {
locality, ok := ctx.Value(localityKey{}).(string)
func MustFromContext(ctx context.Context) string {
locality, ok := ctx.Value(key{}).(string)
if !ok {
panic(stacktrace.NewError("request locality not present in context"))
}

return locality
}

// WithRequestLocality returns a new context with the given locality.
func WithRequestLocality(ctx context.Context, locality string) context.Context {
return context.WithValue(ctx, localityKey{}, locality)
// NewContext returns a new context with the given locality.
func NewContext(ctx context.Context, locality string) context.Context {
return context.WithValue(ctx, key{}, locality)
}

// LocalityMiddleware is an HTTP middleware that stamps each incoming request with this
// Middleware is an HTTP middleware that stamps each incoming request with this
// DSS instance's locality so that locality-dependent operations execute deterministically across nodes.
func LocalityMiddleware(locality string) func(http.Handler) http.Handler {
func Middleware(locality string) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
next.ServeHTTP(w, r.WithContext(WithRequestLocality(r.Context(), locality)))
next.ServeHTTP(w, r.WithContext(NewContext(r.Context(), locality)))
})
}
}
78 changes: 78 additions & 0 deletions pkg/models/geo.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package models

import (
"encoding/json"
"time"

"github.com/golang/geo/s2"
Expand Down Expand Up @@ -46,6 +47,83 @@ type Volume3D struct {
Footprint Geometry
}

type Volume3DJSON struct {
AltitudeHi *float32 `json:"altitude_hi,omitempty"`
AltitudeLo *float32 `json:"altitude_lo,omitempty"`
Footprint *geometryJSON `json:"footprint,omitempty"`
}

type geometryType string

const (
circle geometryType = "circle"
polygon geometryType = "polygon"
cells geometryType = "cells"
)

// geometryJSON is a helper struct for marshaling and unmarshaling Geometry types to/from JSON.
type geometryJSON struct {
Type geometryType `json:"type"`
Polygon *GeoPolygon `json:"polygon,omitempty"`
Circle *GeoCircle `json:"circle,omitempty"`
Cells []s2.CellID `json:"cells,omitempty"`
}

func (v Volume3D) MarshalJSON() ([]byte, error) {
w := Volume3DJSON{AltitudeHi: v.AltitudeHi, AltitudeLo: v.AltitudeLo}
if v.Footprint != nil {
switch f := v.Footprint.(type) {
case *GeoPolygon:
w.Footprint = &geometryJSON{Type: polygon, Polygon: f}

case *GeoCircle:
w.Footprint = &geometryJSON{Type: circle, Circle: f}

case precomputedCellGeometry:
cellsResult := make([]s2.CellID, 0, len(f))
for id := range f {
cellsResult = append(cellsResult, id)
}
w.Footprint = &geometryJSON{Type: cells, Cells: cellsResult}

default:
return nil, stacktrace.NewError("Volume3D: unsupported Footprint type %T for JSON marshaling", v.Footprint)
}
}

return json.Marshal(w)
}

func (v *Volume3D) UnmarshalJSON(data []byte) error {
var w Volume3DJSON
if err := json.Unmarshal(data, &w); err != nil {
return err
}
v.AltitudeHi = w.AltitudeHi
v.AltitudeLo = w.AltitudeLo
if w.Footprint != nil {
switch w.Footprint.Type {
case polygon:
v.Footprint = w.Footprint.Polygon

case circle:
v.Footprint = w.Footprint.Circle

case cells:
pcg := make(precomputedCellGeometry, len(w.Footprint.Cells))
for _, id := range w.Footprint.Cells {
pcg[id] = struct{}{}
}

v.Footprint = pcg
default:
return stacktrace.NewError("Volume3D: unknown geometry type %q", w.Footprint.Type)
}
}

return nil
}

// Geometry models a geometry.
type Geometry interface {
// CalculateCovering returns an s2 cell covering for a geometry.
Expand Down
Loading
Loading