diff --git a/js/docs/src/content/docs/lex-reference/multistream/place-stream-multistream-defs.md b/js/docs/src/content/docs/lex-reference/multistream/place-stream-multistream-defs.md index 582f21ab0..f2255f626 100644 --- a/js/docs/src/content/docs/lex-reference/multistream/place-stream-multistream-defs.md +++ b/js/docs/src/content/docs/lex-reference/multistream/place-stream-multistream-defs.md @@ -15,11 +15,28 @@ description: Reference for the place.stream.multistream.defs lexicon **Properties:** -| Name | Type | Req'd | Description | Constraints | -| -------- | --------- | ----- | ----------- | ---------------- | -| `uri` | `string` | ✅ | | Format: `at-uri` | -| `cid` | `string` | ✅ | | Format: `cid` | -| `record` | `unknown` | ✅ | | | +| Name | Type | Req'd | Description | Constraints | +| ------------- | ------------------------------------------------------------------------------------------- | ----- | ----------- | ---------------- | +| `uri` | `string` | ✅ | | Format: `at-uri` | +| `cid` | `string` | ✅ | | Format: `cid` | +| `record` | `unknown` | ✅ | | | +| `latestEvent` | [`place.stream.multistream.defs#event`](/lex-reference/place-stream-multistream-defs#event) | ❌ | | | + +--- + + + +### `event` + +**Type:** `object` + +**Properties:** + +| Name | Type | Req'd | Description | Constraints | +| ----------- | -------- | ----- | ----------- | ---------------------------------------------- | +| `message` | `string` | ✅ | | | +| `status` | `string` | ✅ | | Enum: `inactive`, `pending`, `active`, `error` | +| `createdAt` | `string` | ✅ | | Format: `datetime` | --- @@ -44,6 +61,27 @@ description: Reference for the place.stream.multistream.defs lexicon }, "record": { "type": "unknown" + }, + "latestEvent": { + "type": "ref", + "ref": "place.stream.multistream.defs#event" + } + } + }, + "event": { + "type": "object", + "required": ["message", "status", "createdAt"], + "properties": { + "message": { + "type": "string" + }, + "status": { + "type": "string", + "enum": ["inactive", "pending", "active", "error"] + }, + "createdAt": { + "type": "string", + "format": "datetime" } } } diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index e6dae185e..1894c5d49 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -1889,10 +1889,30 @@ "type": "string", "format": "cid" }, - "record": {} + "record": {}, + "latestEvent": { + "$ref": "#/components/schemas/place.stream.multistream.defs_event" + } }, "required": ["uri", "cid", "record"] }, + "place.stream.multistream.defs_event": { + "type": "object", + "properties": { + "message": { + "type": "string" + }, + "status": { + "type": "string", + "enum": ["inactive", "pending", "active", "error"] + }, + "createdAt": { + "type": "string", + "format": "date-time" + } + }, + "required": ["message", "status", "createdAt"] + }, "place.stream.livestream_livestreamView": { "type": "object", "properties": { diff --git a/lexicons/place/stream/multistream/defs.json b/lexicons/place/stream/multistream/defs.json index f06ca553d..e2dec518a 100644 --- a/lexicons/place/stream/multistream/defs.json +++ b/lexicons/place/stream/multistream/defs.json @@ -8,7 +8,23 @@ "properties": { "uri": { "type": "string", "format": "at-uri" }, "cid": { "type": "string", "format": "cid" }, - "record": { "type": "unknown" } + "record": { "type": "unknown" }, + "latestEvent": { + "type": "ref", + "ref": "place.stream.multistream.defs#event" + } + } + }, + "event": { + "type": "object", + "required": ["message", "status", "createdAt"], + "properties": { + "message": { "type": "string" }, + "status": { + "type": "string", + "enum": ["inactive", "pending", "active", "error"] + }, + "createdAt": { "type": "string", "format": "datetime" } } } } diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index b953d1e6b..ebe9848b1 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -127,10 +127,6 @@ type runningMultistream struct { uri string } -func (rm *runningMultistream) Cancel() { - rm.cancel() -} - // we're making an attempt here not to log (sensitive) stream keys, so we're // referencing by atproto URI func (ss *StreamSession) HandleMultistreamTargets(ctx context.Context) error { @@ -145,18 +141,22 @@ func (ss *StreamSession) HandleMultistreamTargets(ctx context.Context) error { return fmt.Errorf("failed to list multistream targets: %w", err) } currentRunning := map[string]bool{} - for _, target := range targets { - rec, ok := target.Record.Val.(*streamplace.MultistreamTarget) + for _, targetView := range targets { + rec, ok := targetView.Record.Val.(*streamplace.MultistreamTarget) if !ok { - log.Error(ctx, "failed to convert multistream target to streamplace multistream target", "uri", target.Uri) + log.Error(ctx, "failed to convert multistream target to streamplace multistream target", "uri", targetView.Uri) continue } - key := fmt.Sprintf("%s:%s", target.Uri, rec.Url) + key := fmt.Sprintf("%s:%s", targetView.Uri, rec.Url) if running[key] == nil { childCtx, childCancel := context.WithCancel(ctx) ss.Go(ctx, func() error { - log.Log(ctx, "starting multistream target", "uri", target.Uri) - return ss.StartMultistreamTarget(childCtx, target.Record.Val.(*streamplace.MultistreamTarget)) + log.Log(ctx, "starting multistream target", "uri", targetView.Uri) + err := ss.statefulDB.CreateMultistreamEvent(targetView.Uri, "starting multistream target", "pending") + if err != nil { + log.Error(ctx, "failed to create multistream event", "error", err) + } + return ss.StartMultistreamTarget(childCtx, targetView) }) running[key] = &runningMultistream{ cancel: childCancel, @@ -168,7 +168,7 @@ func (ss *StreamSession) HandleMultistreamTargets(ctx context.Context) error { for key := range running { if !currentRunning[key] { log.Log(ctx, "stopping multistream target", "uri", running[key].uri) - running[key].Cancel() + running[key].cancel() delete(running, key) } } @@ -181,11 +181,15 @@ func (ss *StreamSession) HandleMultistreamTargets(ctx context.Context) error { } } -func (ss *StreamSession) StartMultistreamTarget(ctx context.Context, target *streamplace.MultistreamTarget) error { +func (ss *StreamSession) StartMultistreamTarget(ctx context.Context, targetView *streamplace.MultistreamDefs_TargetView) error { for { - err := ss.mm.RTMPPush(ctx, ss.repoDID, "source", target.Url) + err := ss.mm.RTMPPush(ctx, ss.repoDID, "source", targetView) if err != nil { log.Error(ctx, "failed to push to RTMP server", "error", err) + err := ss.statefulDB.CreateMultistreamEvent(targetView.Uri, err.Error(), "error") + if err != nil { + log.Error(ctx, "failed to create multistream event", "error", err) + } } select { case <-ctx.Done(): diff --git a/pkg/media/rtmp_push.go b/pkg/media/rtmp_push.go index da61d8aab..5cf578061 100644 --- a/pkg/media/rtmp_push.go +++ b/pkg/media/rtmp_push.go @@ -7,24 +7,32 @@ import ( "io" "net" "net/url" + "reflect" "strings" + "time" "github.com/go-gst/go-gst/gst" "github.com/google/uuid" "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/streamplace" ) // This function remains in scope for the duration of a single users' playback -func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition string, targetURL string) error { +func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition string, targetView *streamplace.MultistreamDefs_TargetView) error { uu, err := uuid.NewV7() if err != nil { return err } ctx, cancel := context.WithCancel(ctx) defer cancel() - ctx = log.WithLogValues(ctx, "webrtcID", uu.String()) + ctx = log.WithLogValues(ctx, "pushID", uu.String()) ctx = log.WithLogValues(ctx, "mediafunc", "RTMPPush") + rec, ok := targetView.Record.Val.(*streamplace.MultistreamTarget) + if !ok { + return fmt.Errorf("failed to convert target view to multistream target") + } + targetURL := rec.Url pipelineSlice := []string{ "flvmux name=muxer ! rtmp2sink name=rtmp2sink", @@ -66,6 +74,52 @@ func (mm *MediaManager) RTMPPush(ctx context.Context, user string, rendition str return fmt.Errorf("invalid target URL scheme: %s", u.Scheme) } + go func() { + pollFreq := time.Second * 1 + for { + select { + case <-ctx.Done(): + return + case <-time.After(pollFreq): + prop, err := rtmp2sink.GetProperty("stats") + if err != nil { + log.Error(ctx, "error getting rtmp2sink peak-kbps", "error", err) + continue + } + if prop == nil { + log.Error(ctx, "failed to get rtmp2sink peak-kbps", "prop", prop) + continue + } + log.Warn(ctx, "rtmp2sink peak-kbps", "prop", reflect.TypeOf(prop)) + propVal, ok := prop.(*gst.Structure) + if !ok { + log.Error(ctx, "failed to convert rtmp2sink peak-kbps", "prop", prop) + continue + } + outBytesAcked, err := propVal.GetValue("out-bytes-acked") + if err != nil { + log.Error(ctx, "failed to get rtmp2sink out-bytes-acked", "error", err) + continue + } + outBytesAckedVal, ok := outBytesAcked.(uint64) + if !ok { + log.Error(ctx, "failed to convert rtmp2sink out-bytes-acked", "prop", prop) + continue + } + if outBytesAckedVal > 0 { + err = mm.atsync.StatefulDB.CreateMultistreamEvent(targetView.Uri, fmt.Sprintf("wrote %d bytes", outBytesAckedVal), "active") + if err != nil { + log.Error(ctx, "failed to create multistream event", "error", err) + } + // increase pollFreq, once it's working we don't need to spam the database + pollFreq = time.Second * 15 + } + log.Debug(ctx, "rtmp2sink out-bytes-acked", "outBytesAckedVal", outBytesAckedVal) + } + + } + }() + segBuffer := make(chan *bus.Seg, 1024) go func() { segChan := mm.bus.SubscribeSegment(ctx, user, rendition) diff --git a/pkg/statedb/multistream_event.go b/pkg/statedb/multistream_event.go new file mode 100644 index 000000000..323df638e --- /dev/null +++ b/pkg/statedb/multistream_event.go @@ -0,0 +1,34 @@ +package statedb + +import ( + "time" + + "github.com/google/uuid" +) + +type MultistreamEvent struct { + ID string `gorm:"column:id;primarykey"` + TargetURI string `gorm:"column:target_uri;primarykey;index:idx_target_created,priority:1"` + Message string `gorm:"column:message"` + Status string `gorm:"column:status"` + CreatedAt time.Time `gorm:"column:created_at;index:idx_target_created,priority:2"` +} + +func (m *MultistreamEvent) TableName() string { + return "multistream_events" +} + +func (state *StatefulDB) CreateMultistreamEvent(targetURI, message, status string) error { + uu, err := uuid.NewV7() + if err != nil { + return err + } + event := &MultistreamEvent{ + ID: uu.String(), + TargetURI: targetURI, + Message: message, + Status: status, + CreatedAt: time.Now().UTC(), + } + return state.DB.Create(event).Error +} diff --git a/pkg/statedb/multistream_target.go b/pkg/statedb/multistream_target.go index 669a04c1c..5b60eca11 100644 --- a/pkg/statedb/multistream_target.go +++ b/pkg/statedb/multistream_target.go @@ -3,8 +3,10 @@ package statedb import ( "bytes" "fmt" + "time" - "github.com/bluesky-social/indigo/lex/util" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/util" "stream.place/streamplace/pkg/spid" "stream.place/streamplace/pkg/streamplace" ) @@ -76,7 +78,7 @@ func (state *StatefulDB) CreateMultistreamTarget(input *streamplace.MultistreamC return &streamplace.MultistreamDefs_TargetView{ Uri: uri, Cid: cid.String(), - Record: &util.LexiconTypeDecoder{Val: input.MultistreamTarget}, + Record: &lexutil.LexiconTypeDecoder{Val: input.MultistreamTarget}, }, nil } @@ -84,9 +86,22 @@ func (state *StatefulDB) GetMultistreamTarget(uri string) (*streamplace.Multistr return nil, nil } +type TargetWithEvent struct { + MultistreamTarget + LatestEventID *string `gorm:"column:latest_event_id"` + LatestEventStatus *string `gorm:"column:latest_event_status"` + LatestEventMessage *string `gorm:"column:latest_event_message"` + LatestEventCreatedAt *time.Time `gorm:"column:latest_event_created_at"` +} + func (state *StatefulDB) ListMultistreamTargets(repoDID string, limit int, offset int, active *bool) ([]*streamplace.MultistreamDefs_TargetView, error) { - var targets []MultistreamTarget - query := state.DB.Where("repo_did = ?", repoDID) + + var targets []TargetWithEvent + query := state.DB.Table("multistream_targets"). + Select("multistream_targets.*, me.id as latest_event_id, me.status as latest_event_status, me.message as latest_event_message, me.created_at as latest_event_created_at"). + Joins(`LEFT JOIN multistream_events me ON multistream_targets.uri = me.target_uri + AND me.created_at = (SELECT MAX(created_at) FROM multistream_events WHERE target_uri = multistream_targets.uri)`). + Where("repo_did = ?", repoDID) if active != nil { query = query.Where("active = ?", *active) @@ -103,7 +118,7 @@ func (state *StatefulDB) ListMultistreamTargets(repoDID string, limit int, offse result := make([]*streamplace.MultistreamDefs_TargetView, len(targets)) for i, target := range targets { var multistreamTarget streamplace.MultistreamTarget - err = multistreamTarget.UnmarshalCBOR(bytes.NewReader(target.MultistreamTarget)) + err = multistreamTarget.UnmarshalCBOR(bytes.NewReader(target.MultistreamTarget.MultistreamTarget)) if err != nil { return nil, fmt.Errorf("failed to unmarshal multistream target: %w", err) } @@ -112,11 +127,23 @@ func (state *StatefulDB) ListMultistreamTargets(repoDID string, limit int, offse return nil, fmt.Errorf("failed to get CID: %w", err) } - result[i] = &streamplace.MultistreamDefs_TargetView{ + targetView := &streamplace.MultistreamDefs_TargetView{ Uri: target.URI, Cid: cid.String(), - Record: &util.LexiconTypeDecoder{Val: &multistreamTarget}, + Record: &lexutil.LexiconTypeDecoder{Val: &multistreamTarget}, } + + // Add the latest event if it exists + if target.LatestEventID != nil { + event := &streamplace.MultistreamDefs_Event{ + Status: *target.LatestEventStatus, + Message: *target.LatestEventMessage, + CreatedAt: target.LatestEventCreatedAt.Format(util.ISO8601), + } + targetView.LatestEvent = event + } + + result[i] = targetView } return result, nil @@ -177,7 +204,7 @@ func (state *StatefulDB) UpdateMultistreamTarget(uri string, input *streamplace. return &streamplace.MultistreamDefs_TargetView{ Uri: uri, Cid: cid.String(), - Record: &util.LexiconTypeDecoder{Val: input.MultistreamTarget}, + Record: &lexutil.LexiconTypeDecoder{Val: input.MultistreamTarget}, }, nil } diff --git a/pkg/statedb/statedb.go b/pkg/statedb/statedb.go index 631fbd72f..1bf108896 100644 --- a/pkg/statedb/statedb.go +++ b/pkg/statedb/statedb.go @@ -48,6 +48,7 @@ var StatefulDBModels = []any{ Repo{}, Webhook{}, MultistreamTarget{}, + MultistreamEvent{}, } var NoPostgresDatabaseCode = "3D000" diff --git a/pkg/streamplace/multistreamdefs.go b/pkg/streamplace/multistreamdefs.go index 95605e4cf..f6b3c2a67 100644 --- a/pkg/streamplace/multistreamdefs.go +++ b/pkg/streamplace/multistreamdefs.go @@ -8,9 +8,17 @@ import ( "github.com/bluesky-social/indigo/lex/util" ) +// MultistreamDefs_Event is a "event" in the place.stream.multistream.defs schema. +type MultistreamDefs_Event struct { + CreatedAt string `json:"createdAt" cborgen:"createdAt"` + Message string `json:"message" cborgen:"message"` + Status string `json:"status" cborgen:"status"` +} + // MultistreamDefs_TargetView is a "targetView" in the place.stream.multistream.defs schema. type MultistreamDefs_TargetView struct { - Cid string `json:"cid" cborgen:"cid"` - Record *util.LexiconTypeDecoder `json:"record" cborgen:"record"` - Uri string `json:"uri" cborgen:"uri"` + Cid string `json:"cid" cborgen:"cid"` + LatestEvent *MultistreamDefs_Event `json:"latestEvent,omitempty" cborgen:"latestEvent,omitempty"` + Record *util.LexiconTypeDecoder `json:"record" cborgen:"record"` + Uri string `json:"uri" cborgen:"uri"` }