From 1acc598b54269fe6fd035e3714025905b5e6a979 Mon Sep 17 00:00:00 2001
From: "Natalie B." <22222885+espeon@users.noreply.github.com>
Date: Thu, 18 Dec 2025 22:30:28 -0600
Subject: [PATCH] Added lexicon + backend
---
.../server/place-stream-server-defs.md | 34 +++++
.../place-stream-server-deletestorage.md | 58 ++++++++
.../server/place-stream-server-getstorage.md | 58 ++++++++
.../place-stream-server-upsertstorage.md | 114 ++++++++++++++++
lexicons/place/stream/server/defs.json | 17 +++
.../place/stream/server/deleteStorage.json | 22 +++
lexicons/place/stream/server/getStorage.json | 22 +++
.../place/stream/server/upsertStorage.json | 57 ++++++++
pkg/spxrpc/storage.go | 127 ++++++++++++++++++
pkg/spxrpc/stubs.go | 47 +++++++
pkg/statedb/statedb.go | 1 +
pkg/statedb/storage.go | 86 ++++++++++++
pkg/streamplace/serverdefs.go | 10 ++
pkg/streamplace/serverdeleteStorage.go | 26 ++++
pkg/streamplace/servergetStorage.go | 26 ++++
pkg/streamplace/serverupsertStorage.go | 34 +++++
16 files changed, 739 insertions(+)
create mode 100644 js/docs/src/content/docs/lex-reference/server/place-stream-server-deletestorage.md
create mode 100644 js/docs/src/content/docs/lex-reference/server/place-stream-server-getstorage.md
create mode 100644 js/docs/src/content/docs/lex-reference/server/place-stream-server-upsertstorage.md
create mode 100644 lexicons/place/stream/server/deleteStorage.json
create mode 100644 lexicons/place/stream/server/getStorage.json
create mode 100644 lexicons/place/stream/server/upsertStorage.json
create mode 100644 pkg/spxrpc/storage.go
create mode 100644 pkg/statedb/storage.go
create mode 100644 pkg/streamplace/serverdeleteStorage.go
create mode 100644 pkg/streamplace/servergetStorage.go
create mode 100644 pkg/streamplace/serverupsertStorage.go
diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md
index a5eeb4765..e06974271 100644
--- a/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md
+++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-defs.md
@@ -51,6 +51,23 @@ A webhook configuration for receiving Streamplace events.
---
+
+
+### `storage`
+
+**Type:** `object`
+
+S3 storage configuration for backups.
+
+**Properties:**
+
+| Name | Type | Req'd | Description | Constraints |
+| ---------------------------- | --------- | ----- | --------------------------------------------------------------------------------------------- | ------------------ |
+| `url` | `string` | ✅ | S3 storage URL with masked secret key in format: s3+https://ACCESS_KEY:\*\*\*@endpoint/bucket | |
+| `requestedSecondsPerSegment` | `integer` | ✅ | Requested duration for each HLS segment in seconds. | Min: 1
Max: 60 |
+
+---
+
## Lexicon Source
```json
@@ -157,6 +174,23 @@ A webhook configuration for receiving Streamplace events.
"description": "Text to replace with."
}
}
+ },
+ "storage": {
+ "type": "object",
+ "description": "S3 storage configuration for backups.",
+ "required": ["url", "requestedSecondsPerSegment"],
+ "properties": {
+ "url": {
+ "type": "string",
+ "description": "S3 storage URL with masked secret key in format: s3+https://ACCESS_KEY:***@endpoint/bucket"
+ },
+ "requestedSecondsPerSegment": {
+ "type": "integer",
+ "minimum": 1,
+ "maximum": 60,
+ "description": "Requested duration for each HLS segment in seconds."
+ }
+ }
}
}
}
diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-deletestorage.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-deletestorage.md
new file mode 100644
index 000000000..46675d832
--- /dev/null
+++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-deletestorage.md
@@ -0,0 +1,58 @@
+---
+title: place.stream.server.deleteStorage
+description: Reference for the place.stream.server.deleteStorage lexicon
+---
+
+**Lexicon Version:** 1
+
+## Definitions
+
+
+
+### `main`
+
+**Type:** `procedure`
+
+Delete S3 storage configuration.
+
+**Parameters:** _(None defined)_
+
+**Output:**
+
+- **Encoding:** `application/json`
+- **Schema:**
+
+**Schema Type:** `object`
+
+| Name | Type | Req'd | Description | Constraints |
+| --------- | --------- | ----- | ----------- | ----------- |
+| `success` | `boolean` | ✅ | | |
+
+---
+
+## Lexicon Source
+
+```json
+{
+ "lexicon": 1,
+ "id": "place.stream.server.deleteStorage",
+ "defs": {
+ "main": {
+ "type": "procedure",
+ "description": "Delete S3 storage configuration.",
+ "output": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "required": ["success"],
+ "properties": {
+ "success": {
+ "type": "boolean"
+ }
+ }
+ }
+ }
+ }
+ }
+}
+```
diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-getstorage.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-getstorage.md
new file mode 100644
index 000000000..9b52c33a0
--- /dev/null
+++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-getstorage.md
@@ -0,0 +1,58 @@
+---
+title: place.stream.server.getStorage
+description: Reference for the place.stream.server.getStorage lexicon
+---
+
+**Lexicon Version:** 1
+
+## Definitions
+
+
+
+### `main`
+
+**Type:** `query`
+
+Get S3 storage configuration (with masked secret key).
+
+**Parameters:** _(None defined)_
+
+**Output:**
+
+- **Encoding:** `application/json`
+- **Schema:**
+
+**Schema Type:** `object`
+
+| Name | Type | Req'd | Description | Constraints |
+| --------- | ------------------------------------------------------------------------------------- | ----- | ----------- | ----------- |
+| `storage` | [`place.stream.server.defs#storage`](/lex-reference/place-stream-server-defs#storage) | ❌ | | |
+
+---
+
+## Lexicon Source
+
+```json
+{
+ "lexicon": 1,
+ "id": "place.stream.server.getStorage",
+ "defs": {
+ "main": {
+ "type": "query",
+ "description": "Get S3 storage configuration (with masked secret key).",
+ "output": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "properties": {
+ "storage": {
+ "type": "ref",
+ "ref": "place.stream.server.defs#storage"
+ }
+ }
+ }
+ }
+ }
+ }
+}
+```
diff --git a/js/docs/src/content/docs/lex-reference/server/place-stream-server-upsertstorage.md b/js/docs/src/content/docs/lex-reference/server/place-stream-server-upsertstorage.md
new file mode 100644
index 000000000..df9be97bb
--- /dev/null
+++ b/js/docs/src/content/docs/lex-reference/server/place-stream-server-upsertstorage.md
@@ -0,0 +1,114 @@
+---
+title: place.stream.server.upsertStorage
+description: Reference for the place.stream.server.upsertStorage lexicon
+---
+
+**Lexicon Version:** 1
+
+## Definitions
+
+
+
+### `main`
+
+**Type:** `procedure`
+
+Create or update S3 storage configuration for backups.
+
+**Parameters:** _(None defined)_
+
+**Input:**
+
+- **Encoding:** `application/json`
+- **Schema:**
+
+**Schema Type:** `object`
+
+| Name | Type | Req'd | Description | Constraints |
+| ---------------------------- | --------- | ----- | -------------------------------------------------------------------------- | ----------------------------------- |
+| `url` | `string` | ❌ | S3 storage URL in format: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket | |
+| `requestedSecondsPerSegment` | `integer` | ❌ | Requested duration for each HLS segment in seconds. | Min: 1
Max: 60
Default: `6` |
+
+**Output:**
+
+- **Encoding:** `application/json`
+- **Schema:**
+
+**Schema Type:** `object`
+
+| Name | Type | Req'd | Description | Constraints |
+| --------- | ------------------------------------------------------------------------------------- | ----- | ----------- | ----------- |
+| `storage` | [`place.stream.server.defs#storage`](/lex-reference/place-stream-server-defs#storage) | ✅ | | |
+
+**Possible Errors:**
+
+- `InvalidUrl`: The provided S3 URL is invalid or malformed.
+- `ConnectionFailed`: Could not connect to the S3 endpoint with the provided
+ credentials.
+- `MaskedCredentialsModified`: Cannot modify URL while keeping masked
+ credentials. Provide full credentials or omit URL to keep existing
+ configuration.
+
+---
+
+## Lexicon Source
+
+```json
+{
+ "lexicon": 1,
+ "id": "place.stream.server.upsertStorage",
+ "defs": {
+ "main": {
+ "type": "procedure",
+ "description": "Create or update S3 storage configuration for backups.",
+ "input": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "required": [],
+ "properties": {
+ "url": {
+ "type": "string",
+ "description": "S3 storage URL in format: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket"
+ },
+ "requestedSecondsPerSegment": {
+ "type": "integer",
+ "minimum": 1,
+ "maximum": 60,
+ "default": 6,
+ "description": "Requested duration for each HLS segment in seconds."
+ }
+ }
+ }
+ },
+ "output": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "required": ["storage"],
+ "properties": {
+ "storage": {
+ "type": "ref",
+ "ref": "place.stream.server.defs#storage"
+ }
+ }
+ }
+ },
+ "errors": [
+ {
+ "name": "InvalidUrl",
+ "description": "The provided S3 URL is invalid or malformed."
+ },
+ {
+ "name": "ConnectionFailed",
+ "description": "Could not connect to the S3 endpoint with the provided credentials."
+ },
+ {
+ "name": "MaskedCredentialsModified",
+ "description": "Cannot modify URL while keeping masked credentials. Provide full credentials or omit URL to keep existing configuration."
+ }
+ ]
+ }
+ }
+}
+```
diff --git a/lexicons/place/stream/server/defs.json b/lexicons/place/stream/server/defs.json
index 459309144..6c4020b0e 100644
--- a/lexicons/place/stream/server/defs.json
+++ b/lexicons/place/stream/server/defs.json
@@ -98,6 +98,23 @@
"description": "Text to replace with."
}
}
+ },
+ "storage": {
+ "type": "object",
+ "description": "S3 storage configuration for backups.",
+ "required": ["url", "requestedSecondsPerSegment"],
+ "properties": {
+ "url": {
+ "type": "string",
+ "description": "S3 storage URL with masked secret key in format: s3+https://ACCESS_KEY:***@endpoint/bucket"
+ },
+ "requestedSecondsPerSegment": {
+ "type": "integer",
+ "minimum": 1,
+ "maximum": 60,
+ "description": "Requested duration for each HLS segment in seconds."
+ }
+ }
}
}
}
diff --git a/lexicons/place/stream/server/deleteStorage.json b/lexicons/place/stream/server/deleteStorage.json
new file mode 100644
index 000000000..ca581d796
--- /dev/null
+++ b/lexicons/place/stream/server/deleteStorage.json
@@ -0,0 +1,22 @@
+{
+ "lexicon": 1,
+ "id": "place.stream.server.deleteStorage",
+ "defs": {
+ "main": {
+ "type": "procedure",
+ "description": "Delete S3 storage configuration.",
+ "output": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "required": ["success"],
+ "properties": {
+ "success": {
+ "type": "boolean"
+ }
+ }
+ }
+ }
+ }
+ }
+}
diff --git a/lexicons/place/stream/server/getStorage.json b/lexicons/place/stream/server/getStorage.json
new file mode 100644
index 000000000..aeec23d51
--- /dev/null
+++ b/lexicons/place/stream/server/getStorage.json
@@ -0,0 +1,22 @@
+{
+ "lexicon": 1,
+ "id": "place.stream.server.getStorage",
+ "defs": {
+ "main": {
+ "type": "query",
+ "description": "Get S3 storage configuration (with masked secret key).",
+ "output": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "properties": {
+ "storage": {
+ "type": "ref",
+ "ref": "place.stream.server.defs#storage"
+ }
+ }
+ }
+ }
+ }
+ }
+}
diff --git a/lexicons/place/stream/server/upsertStorage.json b/lexicons/place/stream/server/upsertStorage.json
new file mode 100644
index 000000000..9b06f66c3
--- /dev/null
+++ b/lexicons/place/stream/server/upsertStorage.json
@@ -0,0 +1,57 @@
+{
+ "lexicon": 1,
+ "id": "place.stream.server.upsertStorage",
+ "defs": {
+ "main": {
+ "type": "procedure",
+ "description": "Create or update S3 storage configuration for backups.",
+ "input": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "required": [],
+ "properties": {
+ "url": {
+ "type": "string",
+ "description": "S3 storage URL in format: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket"
+ },
+ "requestedSecondsPerSegment": {
+ "type": "integer",
+ "minimum": 1,
+ "maximum": 60,
+ "default": 6,
+ "description": "Requested duration for each HLS segment in seconds."
+ }
+ }
+ }
+ },
+ "output": {
+ "encoding": "application/json",
+ "schema": {
+ "type": "object",
+ "required": ["storage"],
+ "properties": {
+ "storage": {
+ "type": "ref",
+ "ref": "place.stream.server.defs#storage"
+ }
+ }
+ }
+ },
+ "errors": [
+ {
+ "name": "InvalidUrl",
+ "description": "The provided S3 URL is invalid or malformed."
+ },
+ {
+ "name": "ConnectionFailed",
+ "description": "Could not connect to the S3 endpoint with the provided credentials."
+ },
+ {
+ "name": "MaskedCredentialsModified",
+ "description": "Cannot modify URL while keeping masked credentials. Provide full credentials or omit URL to keep existing configuration."
+ }
+ ]
+ }
+ }
+}
diff --git a/pkg/spxrpc/storage.go b/pkg/spxrpc/storage.go
new file mode 100644
index 000000000..14540626f
--- /dev/null
+++ b/pkg/spxrpc/storage.go
@@ -0,0 +1,127 @@
+package spxrpc
+
+import (
+ "context"
+ "net/http"
+ "regexp"
+ "strings"
+
+ "github.com/bluesky-social/indigo/xrpc"
+ "github.com/labstack/echo/v4"
+ "github.com/streamplace/oatproxy/pkg/oatproxy"
+ "stream.place/streamplace/pkg/log"
+ "stream.place/streamplace/pkg/statedb"
+ placestreamtypes "stream.place/streamplace/pkg/streamplace"
+)
+
+func (s *Server) handlePlaceStreamServerUpsertStorage(ctx context.Context, input *placestreamtypes.ServerUpsertStorage_Input) (*placestreamtypes.ServerUpsertStorage_Output, error) {
+ // Get authenticated user
+ session, _ := oatproxy.GetOAuthSession(ctx)
+ if session == nil {
+ return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found")
+ }
+
+ // Get existing storage if any
+ existing, _ := s.statefulDB.GetStorage(session.DID)
+
+ var url string
+ if input.Url != nil {
+ url = *input.Url
+ // If the client sent back the masked version, use the existing URL
+ if existing != nil && url != "" {
+ // Check if the URL matches the masked format from ToLexicon
+ maskedExisting := existing.ToLexicon().Url
+ if url == maskedExisting {
+ url = existing.URL
+ }
+ }
+
+ // Check if URL contains masked credentials but was otherwise modified
+ if strings.Contains(url, ":***@") {
+ return nil, &xrpc.Error{
+ StatusCode: http.StatusBadRequest,
+ Wrapped: &xrpc.XRPCError{ErrStr: "MaskedCredentialsModified", Message: "Cannot modify URL while keeping masked credentials. Provide full credentials or omit URL to keep existing configuration."},
+ }
+ }
+
+ // Validate S3 URL format if we have a new one
+ if url != "" && !isValidS3URL(url) {
+ return nil, echo.NewHTTPError(http.StatusBadRequest, "Invalid S3 URL format. Expected: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket")
+ }
+ } else if existing != nil {
+ url = existing.URL
+ }
+
+ if url == "" {
+ return nil, echo.NewHTTPError(http.StatusBadRequest, "S3 URL is required")
+ }
+
+ // Convert input to database model
+ storage := statedb.StorageFromLexiconInput(input, session.DID)
+ storage.URL = url
+
+ // Upsert storage
+ err := s.statefulDB.UpsertStorage(storage)
+ if err != nil {
+ log.Error(ctx, "failed to upsert storage", "err", err)
+ return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to save storage configuration")
+ }
+
+ // Get the saved storage to return
+ savedStorage, err := s.statefulDB.GetStorage(session.DID)
+ if err != nil {
+ log.Error(ctx, "failed to get storage after upsert", "err", err)
+ return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to retrieve storage configuration")
+ }
+
+ return &placestreamtypes.ServerUpsertStorage_Output{
+ Storage: savedStorage.ToLexicon(),
+ }, nil
+}
+
+func (s *Server) handlePlaceStreamServerGetStorage(ctx context.Context) (*placestreamtypes.ServerGetStorage_Output, error) {
+ // Get authenticated user
+ session, _ := oatproxy.GetOAuthSession(ctx)
+ if session == nil {
+ return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found")
+ }
+
+ // Get storage
+ storage, err := s.statefulDB.GetStorage(session.DID)
+ if err != nil {
+ // Return empty response if no storage configured
+ return &placestreamtypes.ServerGetStorage_Output{
+ Storage: nil,
+ }, nil
+ }
+
+ return &placestreamtypes.ServerGetStorage_Output{
+ Storage: storage.ToLexicon(),
+ }, nil
+}
+
+func (s *Server) handlePlaceStreamServerDeleteStorage(ctx context.Context) (*placestreamtypes.ServerDeleteStorage_Output, error) {
+ // Get authenticated user
+ session, _ := oatproxy.GetOAuthSession(ctx)
+ if session == nil {
+ return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session not found")
+ }
+
+ // Delete storage
+ err := s.statefulDB.DeleteStorage(session.DID)
+ if err != nil {
+ log.Error(ctx, "failed to delete storage", "err", err)
+ return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to delete storage configuration")
+ }
+
+ return &placestreamtypes.ServerDeleteStorage_Output{
+ Success: true,
+ }, nil
+}
+
+// isValidS3URL validates the S3 URL format
+func isValidS3URL(url string) bool {
+ // Format: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket
+ re := regexp.MustCompile(`^s3\+https?://[^:]+:[^@]+@[^/]+/.+$`)
+ return re.MatchString(url)
+}
diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go
index 7367ebc78..e36468f2e 100644
--- a/pkg/spxrpc/stubs.go
+++ b/pkg/spxrpc/stubs.go
@@ -305,11 +305,14 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error {
e.POST("/xrpc/place.stream.multistream.putTarget", s.HandlePlaceStreamMultistreamPutTarget)
e.POST("/xrpc/place.stream.playback.whep", s.HandlePlaceStreamPlaybackWhep)
e.POST("/xrpc/place.stream.server.createWebhook", s.HandlePlaceStreamServerCreateWebhook)
+ e.POST("/xrpc/place.stream.server.deleteStorage", s.HandlePlaceStreamServerDeleteStorage)
e.POST("/xrpc/place.stream.server.deleteWebhook", s.HandlePlaceStreamServerDeleteWebhook)
e.GET("/xrpc/place.stream.server.getServerTime", s.HandlePlaceStreamServerGetServerTime)
+ e.GET("/xrpc/place.stream.server.getStorage", s.HandlePlaceStreamServerGetStorage)
e.GET("/xrpc/place.stream.server.getWebhook", s.HandlePlaceStreamServerGetWebhook)
e.GET("/xrpc/place.stream.server.listWebhooks", s.HandlePlaceStreamServerListWebhooks)
e.POST("/xrpc/place.stream.server.updateWebhook", s.HandlePlaceStreamServerUpdateWebhook)
+ e.POST("/xrpc/place.stream.server.upsertStorage", s.HandlePlaceStreamServerUpsertStorage)
return nil
}
@@ -795,6 +798,19 @@ func (s *Server) HandlePlaceStreamServerCreateWebhook(c echo.Context) error {
return c.JSON(200, out)
}
+func (s *Server) HandlePlaceStreamServerDeleteStorage(c echo.Context) error {
+ ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerDeleteStorage")
+ defer span.End()
+ var out *placestream.ServerDeleteStorage_Output
+ var handleErr error
+ // func (s *Server) handlePlaceStreamServerDeleteStorage(ctx context.Context) (*placestream.ServerDeleteStorage_Output, error)
+ out, handleErr = s.handlePlaceStreamServerDeleteStorage(ctx)
+ if handleErr != nil {
+ return handleErr
+ }
+ return c.JSON(200, out)
+}
+
func (s *Server) HandlePlaceStreamServerDeleteWebhook(c echo.Context) error {
ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerDeleteWebhook")
defer span.End()
@@ -826,6 +842,19 @@ func (s *Server) HandlePlaceStreamServerGetServerTime(c echo.Context) error {
return c.JSON(200, out)
}
+func (s *Server) HandlePlaceStreamServerGetStorage(c echo.Context) error {
+ ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerGetStorage")
+ defer span.End()
+ var out *placestream.ServerGetStorage_Output
+ var handleErr error
+ // func (s *Server) handlePlaceStreamServerGetStorage(ctx context.Context) (*placestream.ServerGetStorage_Output, error)
+ out, handleErr = s.handlePlaceStreamServerGetStorage(ctx)
+ if handleErr != nil {
+ return handleErr
+ }
+ return c.JSON(200, out)
+}
+
func (s *Server) HandlePlaceStreamServerGetWebhook(c echo.Context) error {
ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerGetWebhook")
defer span.End()
@@ -893,6 +922,24 @@ func (s *Server) HandlePlaceStreamServerUpdateWebhook(c echo.Context) error {
return c.JSON(200, out)
}
+func (s *Server) HandlePlaceStreamServerUpsertStorage(c echo.Context) error {
+ ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamServerUpsertStorage")
+ defer span.End()
+
+ var body placestream.ServerUpsertStorage_Input
+ if err := c.Bind(&body); err != nil {
+ return err
+ }
+ var out *placestream.ServerUpsertStorage_Output
+ var handleErr error
+ // func (s *Server) handlePlaceStreamServerUpsertStorage(ctx context.Context,body *placestream.ServerUpsertStorage_Input) (*placestream.ServerUpsertStorage_Output, error)
+ out, handleErr = s.handlePlaceStreamServerUpsertStorage(ctx, &body)
+ if handleErr != nil {
+ return handleErr
+ }
+ return c.JSON(200, out)
+}
+
func (s *Server) RegisterHandlersToolsOzone(e *echo.Echo) error {
return nil
}
diff --git a/pkg/statedb/statedb.go b/pkg/statedb/statedb.go
index 604fbf08d..14fbcba11 100644
--- a/pkg/statedb/statedb.go
+++ b/pkg/statedb/statedb.go
@@ -53,6 +53,7 @@ var StatefulDBModels = []any{
MultistreamEvent{},
BrandingBlob{},
ModerationAuditLog{},
+ Storage{},
BroadcastOrigin{},
}
diff --git a/pkg/statedb/storage.go b/pkg/statedb/storage.go
new file mode 100644
index 000000000..190bada5d
--- /dev/null
+++ b/pkg/statedb/storage.go
@@ -0,0 +1,86 @@
+package statedb
+
+import (
+ "regexp"
+ "time"
+
+ "github.com/google/uuid"
+ "gorm.io/gorm"
+ streamplace "stream.place/streamplace/pkg/streamplace"
+)
+
+// Storage represents S3 storage configuration for a user
+type Storage struct {
+ ID string `gorm:"column:id;primarykey"`
+ UserDID string `gorm:"column:user_did;not null;unique"`
+ URL string `gorm:"column:url;not null"`
+ RequestedSecondsPerSegment int `gorm:"column:requested_seconds_per_segment;default:6"`
+ CreatedAt time.Time `gorm:"column:created_at"`
+ UpdatedAt time.Time `gorm:"column:updated_at"`
+}
+
+func (s *Storage) TableName() string {
+ return "storage"
+}
+
+func maskSecretKey(url string) string {
+ // Format: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket
+ re := regexp.MustCompile(`(s3\+https?://[^:]+:)([^@]+)(@.+)`)
+ return re.ReplaceAllString(url, "${1}***${3}")
+}
+
+func (state *StatefulDB) UpsertStorage(storage *Storage) error {
+ if storage.ID == "" {
+ storage.ID = uuid.New().String()
+ }
+
+ var existing Storage
+ err := state.DB.Where("user_did = ?", storage.UserDID).First(&existing).Error
+ if err == nil {
+ storage.ID = existing.ID
+ storage.CreatedAt = existing.CreatedAt
+ return state.DB.Save(storage).Error
+ } else if err == gorm.ErrRecordNotFound {
+ return state.DB.Create(storage).Error
+ }
+
+ return err
+}
+
+func (state *StatefulDB) GetStorage(userDID string) (*Storage, error) {
+ var storage Storage
+ err := state.DB.Where("user_did = ?", userDID).First(&storage).Error
+ if err != nil {
+ return nil, err
+ }
+ return &storage, nil
+}
+
+func (state *StatefulDB) DeleteStorage(userDID string) error {
+ return state.DB.Where("user_did = ?", userDID).Delete(&Storage{}).Error
+}
+
+func (s *Storage) ToLexicon() *streamplace.ServerDefs_Storage {
+ return &streamplace.ServerDefs_Storage{
+ Url: maskSecretKey(s.URL),
+ RequestedSecondsPerSegment: int64(s.RequestedSecondsPerSegment),
+ }
+}
+
+func StorageFromLexiconInput(input *streamplace.ServerUpsertStorage_Input, userDID string) *Storage {
+ requestedSeconds := 20
+ if input.RequestedSecondsPerSegment != nil {
+ requestedSeconds = int(*input.RequestedSecondsPerSegment)
+ }
+
+ var url string
+ if input.Url != nil {
+ url = *input.Url
+ }
+
+ return &Storage{
+ UserDID: userDID,
+ URL: url,
+ RequestedSecondsPerSegment: requestedSeconds,
+ }
+}
diff --git a/pkg/streamplace/serverdefs.go b/pkg/streamplace/serverdefs.go
index 18eff55e8..31519d751 100644
--- a/pkg/streamplace/serverdefs.go
+++ b/pkg/streamplace/serverdefs.go
@@ -12,6 +12,16 @@ type ServerDefs_RewriteRule struct {
To string `json:"to" cborgen:"to"`
}
+// ServerDefs_Storage is a "storage" in the place.stream.server.defs schema.
+//
+// S3 storage configuration for backups.
+type ServerDefs_Storage struct {
+ // requestedSecondsPerSegment: Requested duration for each HLS segment in seconds.
+ RequestedSecondsPerSegment int64 `json:"requestedSecondsPerSegment" cborgen:"requestedSecondsPerSegment"`
+ // url: S3 storage URL with masked secret key in format: s3+https://ACCESS_KEY:***@endpoint/bucket
+ Url string `json:"url" cborgen:"url"`
+}
+
// ServerDefs_Webhook is a "webhook" in the place.stream.server.defs schema.
//
// A webhook configuration for receiving Streamplace events.
diff --git a/pkg/streamplace/serverdeleteStorage.go b/pkg/streamplace/serverdeleteStorage.go
new file mode 100644
index 000000000..b63d067cf
--- /dev/null
+++ b/pkg/streamplace/serverdeleteStorage.go
@@ -0,0 +1,26 @@
+// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT.
+
+// Lexicon schema: place.stream.server.deleteStorage
+
+package streamplace
+
+import (
+ "context"
+
+ lexutil "github.com/bluesky-social/indigo/lex/util"
+)
+
+// ServerDeleteStorage_Output is the output of a place.stream.server.deleteStorage call.
+type ServerDeleteStorage_Output struct {
+ Success bool `json:"success" cborgen:"success"`
+}
+
+// ServerDeleteStorage calls the XRPC method "place.stream.server.deleteStorage".
+func ServerDeleteStorage(ctx context.Context, c lexutil.LexClient) (*ServerDeleteStorage_Output, error) {
+ var out ServerDeleteStorage_Output
+ if err := c.LexDo(ctx, lexutil.Procedure, "", "place.stream.server.deleteStorage", nil, nil, &out); err != nil {
+ return nil, err
+ }
+
+ return &out, nil
+}
diff --git a/pkg/streamplace/servergetStorage.go b/pkg/streamplace/servergetStorage.go
new file mode 100644
index 000000000..2a02da55d
--- /dev/null
+++ b/pkg/streamplace/servergetStorage.go
@@ -0,0 +1,26 @@
+// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT.
+
+// Lexicon schema: place.stream.server.getStorage
+
+package streamplace
+
+import (
+ "context"
+
+ lexutil "github.com/bluesky-social/indigo/lex/util"
+)
+
+// ServerGetStorage_Output is the output of a place.stream.server.getStorage call.
+type ServerGetStorage_Output struct {
+ Storage *ServerDefs_Storage `json:"storage,omitempty" cborgen:"storage,omitempty"`
+}
+
+// ServerGetStorage calls the XRPC method "place.stream.server.getStorage".
+func ServerGetStorage(ctx context.Context, c lexutil.LexClient) (*ServerGetStorage_Output, error) {
+ var out ServerGetStorage_Output
+ if err := c.LexDo(ctx, lexutil.Query, "", "place.stream.server.getStorage", nil, nil, &out); err != nil {
+ return nil, err
+ }
+
+ return &out, nil
+}
diff --git a/pkg/streamplace/serverupsertStorage.go b/pkg/streamplace/serverupsertStorage.go
new file mode 100644
index 000000000..7449fa286
--- /dev/null
+++ b/pkg/streamplace/serverupsertStorage.go
@@ -0,0 +1,34 @@
+// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT.
+
+// Lexicon schema: place.stream.server.upsertStorage
+
+package streamplace
+
+import (
+ "context"
+
+ lexutil "github.com/bluesky-social/indigo/lex/util"
+)
+
+// ServerUpsertStorage_Input is the input argument to a place.stream.server.upsertStorage call.
+type ServerUpsertStorage_Input struct {
+ // requestedSecondsPerSegment: Requested duration for each HLS segment in seconds.
+ RequestedSecondsPerSegment *int64 `json:"requestedSecondsPerSegment,omitempty" cborgen:"requestedSecondsPerSegment,omitempty"`
+ // url: S3 storage URL in format: s3+https://ACCESS_KEY:SECRET_KEY@endpoint/bucket
+ Url *string `json:"url,omitempty" cborgen:"url,omitempty"`
+}
+
+// ServerUpsertStorage_Output is the output of a place.stream.server.upsertStorage call.
+type ServerUpsertStorage_Output struct {
+ Storage *ServerDefs_Storage `json:"storage" cborgen:"storage"`
+}
+
+// ServerUpsertStorage calls the XRPC method "place.stream.server.upsertStorage".
+func ServerUpsertStorage(ctx context.Context, c lexutil.LexClient, input *ServerUpsertStorage_Input) (*ServerUpsertStorage_Output, error) {
+ var out ServerUpsertStorage_Output
+ if err := c.LexDo(ctx, lexutil.Procedure, "application/json", "place.stream.server.upsertStorage", nil, input, &out); err != nil {
+ return nil, err
+ }
+
+ return &out, nil
+}
--
2.51.2