From 09044c543655d1c7010d9302020adeaac4afa9b1 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sun, 17 May 2026 17:19:21 -0700 Subject: [PATCH] =?UTF-8?q?viewCount:=20drop=20methodology/threshold,=20tr?= =?UTF-8?q?ack=20id=20=E2=86=92=20strongRef?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two lexicon changes: - methodology / thresholdSegments come off the main record. The contract is now \"the reporting node did its best with what it had\"; revisit when we actually have hot-swappable algorithms to advertise. - trackUsage.trackId (stringified u32 within the MUXL container) becomes a com.atproto.repo.strongRef pointing at the place.stream.media.track record. The in-container id is a local detail; the strongRef is a globally stable identity, so user-contributed tracks (transcripts, transcodes published by accounts other than the streamer) attribute to their own records when those start flowing. Aggregator gains a FetchTrackRefs callback that resolves a blob CID to a {trackIdInContainer → strongRef} map; one fetch per CID, cached alongside the metafile. Rows whose trackId can't be resolved (no track record yet, or the record was deleted) get dropped — a TrackUsage without a stable strongRef wouldn't help any consumer. Bootstrap supplies the production resolver via model.GetMediaTracksByBlob + ToRecord. Tests pass an in-memory map. Co-Authored-By: Claude Opus 4.7 --- .../media/place-stream-media-viewcount.md | 69 ++++----- lexicons/place/stream/media/viewCount.json | 37 ++--- pkg/cmd/streamplace.go | 46 +++++- pkg/streamplace/cbor_gen.go | 146 +++--------------- pkg/streamplace/mediaviewCount.go | 17 +- pkg/viewlog/aggregate.go | 97 ++++++++---- pkg/viewlog/aggregate_test.go | 88 ++++++++++- pkg/viewlog/publish.go | 41 +++-- 8 files changed, 278 insertions(+), 263 deletions(-) diff --git a/js/docs/src/content/docs/lex-reference/media/place-stream-media-viewcount.md b/js/docs/src/content/docs/lex-reference/media/place-stream-media-viewcount.md index 4c76b10dc..f964312ac 100644 --- a/js/docs/src/content/docs/lex-reference/media/place-stream-media-viewcount.md +++ b/js/docs/src/content/docs/lex-reference/media/place-stream-media-viewcount.md @@ -13,22 +13,20 @@ description: Reference for the place.stream.media.viewCount lexicon **Type:** `record` -A streamplace node's report of view counts for one place.stream.video over a closed time window. Published in the reporting node's server repo (not the streamer's), so a video served by multiple nodes accumulates multiple records — consumers are expected to sum across trusted reporters. The rkey is conventionally `-` so re-running the aggregator over the same window is idempotent. The methodology field documents the floor used to count a session (e.g. "any-segment" = any sid that fetched ≥ threshold segment_requests); the embedded `tracks` array carries objective per-track bytes + duration totals that don't depend on the methodology. +A streamplace node's report of view counts for one place.stream.video over a closed time window. Published in the reporting node's server repo (not the streamer's), so a video served by multiple nodes accumulates multiple records — consumers are expected to sum across trusted reporters. The rkey is conventionally `-` so re-running the aggregator over the same window is idempotent. Counts represent the reporting node's best effort given the data it has; the `tracks` array carries the objective byte / duration totals it observed. **Record Key:** `any` **Record Properties:** -| Name | Type | Req'd | Description | Constraints | -| ------------------- | ------------------------------------- | ----- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------ | -| `video` | `string` | ✅ | AT-URI of the place.stream.video this count is for. | Format: `at-uri` | -| `count` | `integer` | ✅ | Number of distinct sessions that qualified as a view over [windowStart, windowEnd) by the stated methodology. | Min: 0 | -| `windowStart` | `string` | ✅ | Inclusive lower bound of the aggregation window. | Format: `datetime` | -| `windowEnd` | `string` | ✅ | Exclusive upper bound of the aggregation window. | Format: `datetime` | -| `methodology` | `string` | ✅ | Identifier for the counting algorithm. "any-segment" means: a distinct sid that fetched at least `thresholdSegments` segments for this video over the window. Future methodologies (e.g. "ms-from-metafile") will use new tags so older records remain interpretable. | | -| `thresholdSegments` | `integer` | ❌ | Floor on segment_request count per (sid, video) for the "any-segment" methodology. Defaults to 1. | Min: 0 | -| `tracks` | Array of [`#trackUsage`](#trackusage) | ❌ | Per-track totals of bytes + playback duration actually transferred over the window. Computed by intersecting each segment_request's HTTP Range with the metafile's per-track byte layout, so the numbers are objective regardless of the chosen view-counting methodology. | | -| `indexedAt` | `string` | ✅ | When the reporting node ran this aggregation. Useful for ordering successive reports. | Format: `datetime` | +| Name | Type | Req'd | Description | Constraints | +| ------------- | ------------------------------------- | ----- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------ | +| `video` | `string` | ✅ | AT-URI of the place.stream.video this count is for. | Format: `at-uri` | +| `count` | `integer` | ✅ | Number of distinct sessions the reporting node observed as a view over [windowStart, windowEnd). | Min: 0 | +| `windowStart` | `string` | ✅ | Inclusive lower bound of the aggregation window. | Format: `datetime` | +| `windowEnd` | `string` | ✅ | Exclusive upper bound of the aggregation window. | Format: `datetime` | +| `tracks` | Array of [`#trackUsage`](#trackusage) | ❌ | Per-track totals of bytes + playback duration the reporting node served over the window. Each entry references the place.stream.media.track record whose bytes were transferred, so user-contributed tracks (transcripts, transcodes published by other accounts) attribute naturally to their own records. | | +| `indexedAt` | `string` | ✅ | When the reporting node ran this aggregation. Useful for ordering successive reports. | Format: `datetime` | --- @@ -38,15 +36,15 @@ A streamplace node's report of view counts for one place.stream.video over a clo **Type:** `object` -One row of the tracks array: bytes + duration transferred for a single muxlTrack inside the video's MUXL container over the window. trackId matches the muxlTrack record's `trackId` field. +One row of the tracks array: bytes + duration transferred for a single place.stream.media.track record over the window. **Properties:** -| Name | Type | Req'd | Description | Constraints | -| ------------ | --------- | ----- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------- | -| `trackId` | `string` | ✅ | Stringified u32 matching the MUXL container's per-track id. | | -| `bytes` | `integer` | ✅ | Total bytes served from this track's byte ranges over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets. | Min: 0 | -| `durationMs` | `integer` | ✅ | Total playback duration served from this track, in milliseconds. Per HLS segment in the range: (overlap bytes / segment bytes) \* segment duration, so partial-segment fetches credit a proportional share of duration. | Min: 0 | +| Name | Type | Req'd | Description | Constraints | +| ------------ | -------------------------------------------------------------------------------------------------------------------------------------- | ----- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------- | +| `track` | [`com.atproto.repo.strongRef`](https://github.com/bluesky-social/atproto/tree/main/lexicons/com/atproto/repo/strongref.json#undefined) | ✅ | Strong reference to the place.stream.media.track record whose bytes were transferred. | | +| `bytes` | `integer` | ✅ | Total bytes served from this track over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets in the metafile. | Min: 0 | +| `durationMs` | `integer` | ✅ | Total playback duration served from this track, in milliseconds. Per HLS segment in the range: (overlap bytes / segment bytes) \* segment duration, so partial-segment fetches credit a proportional share of duration. | Min: 0 | --- @@ -59,18 +57,11 @@ One row of the tracks array: bytes + duration transferred for a single muxlTrack "defs": { "main": { "type": "record", - "description": "A streamplace node's report of view counts for one place.stream.video over a closed time window. Published in the reporting node's server repo (not the streamer's), so a video served by multiple nodes accumulates multiple records — consumers are expected to sum across trusted reporters. The rkey is conventionally `-` so re-running the aggregator over the same window is idempotent. The methodology field documents the floor used to count a session (e.g. \"any-segment\" = any sid that fetched ≥ threshold segment_requests); the embedded `tracks` array carries objective per-track bytes + duration totals that don't depend on the methodology.", + "description": "A streamplace node's report of view counts for one place.stream.video over a closed time window. Published in the reporting node's server repo (not the streamer's), so a video served by multiple nodes accumulates multiple records — consumers are expected to sum across trusted reporters. The rkey is conventionally `-` so re-running the aggregator over the same window is idempotent. Counts represent the reporting node's best effort given the data it has; the `tracks` array carries the objective byte / duration totals it observed.", "key": "any", "record": { "type": "object", - "required": [ - "video", - "count", - "windowStart", - "windowEnd", - "methodology", - "indexedAt" - ], + "required": ["video", "count", "windowStart", "windowEnd", "indexedAt"], "properties": { "video": { "type": "string", @@ -80,7 +71,7 @@ One row of the tracks array: bytes + duration transferred for a single muxlTrack "count": { "type": "integer", "minimum": 0, - "description": "Number of distinct sessions that qualified as a view over [windowStart, windowEnd) by the stated methodology." + "description": "Number of distinct sessions the reporting node observed as a view over [windowStart, windowEnd)." }, "windowStart": { "type": "string", @@ -92,18 +83,9 @@ One row of the tracks array: bytes + duration transferred for a single muxlTrack "format": "datetime", "description": "Exclusive upper bound of the aggregation window." }, - "methodology": { - "type": "string", - "description": "Identifier for the counting algorithm. \"any-segment\" means: a distinct sid that fetched at least `thresholdSegments` segments for this video over the window. Future methodologies (e.g. \"ms-from-metafile\") will use new tags so older records remain interpretable." - }, - "thresholdSegments": { - "type": "integer", - "minimum": 0, - "description": "Floor on segment_request count per (sid, video) for the \"any-segment\" methodology. Defaults to 1." - }, "tracks": { "type": "array", - "description": "Per-track totals of bytes + playback duration actually transferred over the window. Computed by intersecting each segment_request's HTTP Range with the metafile's per-track byte layout, so the numbers are objective regardless of the chosen view-counting methodology.", + "description": "Per-track totals of bytes + playback duration the reporting node served over the window. Each entry references the place.stream.media.track record whose bytes were transferred, so user-contributed tracks (transcripts, transcodes published by other accounts) attribute naturally to their own records.", "items": { "type": "ref", "ref": "#trackUsage" @@ -119,17 +101,18 @@ One row of the tracks array: bytes + duration transferred for a single muxlTrack }, "trackUsage": { "type": "object", - "description": "One row of the tracks array: bytes + duration transferred for a single muxlTrack inside the video's MUXL container over the window. trackId matches the muxlTrack record's `trackId` field.", - "required": ["trackId", "bytes", "durationMs"], + "description": "One row of the tracks array: bytes + duration transferred for a single place.stream.media.track record over the window.", + "required": ["track", "bytes", "durationMs"], "properties": { - "trackId": { - "type": "string", - "description": "Stringified u32 matching the MUXL container's per-track id." + "track": { + "type": "ref", + "ref": "com.atproto.repo.strongRef", + "description": "Strong reference to the place.stream.media.track record whose bytes were transferred." }, "bytes": { "type": "integer", "minimum": 0, - "description": "Total bytes served from this track's byte ranges over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets." + "description": "Total bytes served from this track over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets in the metafile." }, "durationMs": { "type": "integer", diff --git a/lexicons/place/stream/media/viewCount.json b/lexicons/place/stream/media/viewCount.json index b95cf0d69..d0c4b361b 100644 --- a/lexicons/place/stream/media/viewCount.json +++ b/lexicons/place/stream/media/viewCount.json @@ -4,18 +4,11 @@ "defs": { "main": { "type": "record", - "description": "A streamplace node's report of view counts for one place.stream.video over a closed time window. Published in the reporting node's server repo (not the streamer's), so a video served by multiple nodes accumulates multiple records — consumers are expected to sum across trusted reporters. The rkey is conventionally `-` so re-running the aggregator over the same window is idempotent. The methodology field documents the floor used to count a session (e.g. \"any-segment\" = any sid that fetched ≥ threshold segment_requests); the embedded `tracks` array carries objective per-track bytes + duration totals that don't depend on the methodology.", + "description": "A streamplace node's report of view counts for one place.stream.video over a closed time window. Published in the reporting node's server repo (not the streamer's), so a video served by multiple nodes accumulates multiple records — consumers are expected to sum across trusted reporters. The rkey is conventionally `-` so re-running the aggregator over the same window is idempotent. Counts represent the reporting node's best effort given the data it has; the `tracks` array carries the objective byte / duration totals it observed.", "key": "any", "record": { "type": "object", - "required": [ - "video", - "count", - "windowStart", - "windowEnd", - "methodology", - "indexedAt" - ], + "required": ["video", "count", "windowStart", "windowEnd", "indexedAt"], "properties": { "video": { "type": "string", @@ -25,7 +18,7 @@ "count": { "type": "integer", "minimum": 0, - "description": "Number of distinct sessions that qualified as a view over [windowStart, windowEnd) by the stated methodology." + "description": "Number of distinct sessions the reporting node observed as a view over [windowStart, windowEnd)." }, "windowStart": { "type": "string", @@ -37,18 +30,9 @@ "format": "datetime", "description": "Exclusive upper bound of the aggregation window." }, - "methodology": { - "type": "string", - "description": "Identifier for the counting algorithm. \"any-segment\" means: a distinct sid that fetched at least `thresholdSegments` segments for this video over the window. Future methodologies (e.g. \"ms-from-metafile\") will use new tags so older records remain interpretable." - }, - "thresholdSegments": { - "type": "integer", - "minimum": 0, - "description": "Floor on segment_request count per (sid, video) for the \"any-segment\" methodology. Defaults to 1." - }, "tracks": { "type": "array", - "description": "Per-track totals of bytes + playback duration actually transferred over the window. Computed by intersecting each segment_request's HTTP Range with the metafile's per-track byte layout, so the numbers are objective regardless of the chosen view-counting methodology.", + "description": "Per-track totals of bytes + playback duration the reporting node served over the window. Each entry references the place.stream.media.track record whose bytes were transferred, so user-contributed tracks (transcripts, transcodes published by other accounts) attribute naturally to their own records.", "items": { "type": "ref", "ref": "#trackUsage" @@ -64,17 +48,18 @@ }, "trackUsage": { "type": "object", - "description": "One row of the tracks array: bytes + duration transferred for a single muxlTrack inside the video's MUXL container over the window. trackId matches the muxlTrack record's `trackId` field.", - "required": ["trackId", "bytes", "durationMs"], + "description": "One row of the tracks array: bytes + duration transferred for a single place.stream.media.track record over the window.", + "required": ["track", "bytes", "durationMs"], "properties": { - "trackId": { - "type": "string", - "description": "Stringified u32 matching the MUXL container's per-track id." + "track": { + "type": "ref", + "ref": "com.atproto.repo.strongRef", + "description": "Strong reference to the place.stream.media.track record whose bytes were transferred." }, "bytes": { "type": "integer", "minimum": 0, - "description": "Total bytes served from this track's byte ranges over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets." + "description": "Total bytes served from this track over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets in the metafile." }, "durationMs": { "type": "integer", diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 18a887e58..f23b1a671 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -18,6 +18,7 @@ import ( "syscall" "time" + comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/carstore" "github.com/ethereum/go-ethereum/common/hexutil" "github.com/livepeer/go-livepeer/cmd/livepeer/starter" @@ -402,13 +403,48 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu // goroutine fires per --view-count-aggregate-interval; statedb's // unique TaskKey makes the cross-node race a no-op for losers. if vodStore != nil && cli.ViewCountAggregateInterval > 0 { + // Resolver: for a blob CID, return strongRefs of every + // place.stream.media.track record whose muxlTrack lives in + // that blob, keyed by in-container trackId. Used by the + // aggregator to attribute bytes/duration to the right track + // record (the streamer's original track records or, later, + // user-contributed transcript/transcode tracks). + fetchTrackRefs := func(ctx context.Context, cid string) (map[string]*comatproto.RepoStrongRef, error) { + rows, err := mod.GetMediaTracksByBlob(ctx, cid) + if err != nil { + return nil, err + } + out := make(map[string]*comatproto.RepoStrongRef, len(rows)) + for _, row := range rows { + rec, err := row.ToRecord() + if err != nil { + log.Warn(ctx, "viewlog refs: decode track record", + "uri", row.URI, "error", err) + continue + } + if rec.Track == nil || rec.Track.MediaDefs_MuxlTrack == nil { + continue + } + tid := rec.Track.MediaDefs_MuxlTrack.TrackId + if tid == "" { + continue + } + out[tid] = &comatproto.RepoStrongRef{ + LexiconTypeID: "com.atproto.repo.strongRef", + Uri: row.URI, + Cid: row.CID, + } + } + return out, nil + } state.SetViewCountAggregator(func(ctx context.Context, t statedb.ViewCountAggregateTask) error { return viewlog.RunAggregation(ctx, viewlog.RunAggregationInput{ - Store: vodStore, - CLI: cli, - WindowStart: t.WindowStart, - WindowEnd: t.WindowEnd, - ReadMargin: 2 * cli.ViewLogFlushInterval, + Store: vodStore, + CLI: cli, + WindowStart: t.WindowStart, + WindowEnd: t.WindowEnd, + ReadMargin: 2 * cli.ViewLogFlushInterval, + FetchTrackRefs: fetchTrackRefs, }) }) group.Go(func() error { diff --git a/pkg/streamplace/cbor_gen.go b/pkg/streamplace/cbor_gen.go index 9e91a309b..f7208a209 100644 --- a/pkg/streamplace/cbor_gen.go +++ b/pkg/streamplace/cbor_gen.go @@ -10562,11 +10562,7 @@ func (t *MediaViewCount) MarshalCBOR(w io.Writer) error { } cw := cbg.NewCborWriter(w) - fieldCount := 9 - - if t.ThresholdSegments == nil { - fieldCount-- - } + fieldCount := 7 if t.Tracks == nil { fieldCount-- @@ -10715,29 +10711,6 @@ func (t *MediaViewCount) MarshalCBOR(w io.Writer) error { return err } - // t.Methodology (string) (string) - if len("methodology") > 1000000 { - return xerrors.Errorf("Value in field \"methodology\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("methodology"))); err != nil { - return err - } - if _, err := cw.WriteString(string("methodology")); err != nil { - return err - } - - if len(t.Methodology) > 1000000 { - return xerrors.Errorf("Value in field t.Methodology was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Methodology))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Methodology)); err != nil { - return err - } - // t.WindowStart (string) (string) if len("windowStart") > 1000000 { return xerrors.Errorf("Value in field \"windowStart\" was too long") @@ -10760,38 +10733,6 @@ func (t *MediaViewCount) MarshalCBOR(w io.Writer) error { if _, err := cw.WriteString(string(t.WindowStart)); err != nil { return err } - - // t.ThresholdSegments (int64) (int64) - if t.ThresholdSegments != nil { - - if len("thresholdSegments") > 1000000 { - return xerrors.Errorf("Value in field \"thresholdSegments\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("thresholdSegments"))); err != nil { - return err - } - if _, err := cw.WriteString(string("thresholdSegments")); err != nil { - return err - } - - if t.ThresholdSegments == nil { - if _, err := cw.Write(cbg.CborNull); err != nil { - return err - } - } else { - if *t.ThresholdSegments >= 0 { - if err := cw.WriteMajorTypeHeader(cbg.MajUnsignedInt, uint64(*t.ThresholdSegments)); err != nil { - return err - } - } else { - if err := cw.WriteMajorTypeHeader(cbg.MajNegativeInt, uint64(-*t.ThresholdSegments-1)); err != nil { - return err - } - } - } - - } return nil } @@ -10820,7 +10761,7 @@ func (t *MediaViewCount) UnmarshalCBOR(r io.Reader) (err error) { n := extra - nameBuf := make([]byte, 17) + nameBuf := make([]byte, 11) for i := uint64(0); i < n; i++ { nameLen, ok, err := cbg.ReadFullStringIntoBuf(cr, nameBuf, 1000000) if err != nil { @@ -10955,17 +10896,6 @@ func (t *MediaViewCount) UnmarshalCBOR(r io.Reader) (err error) { t.WindowEnd = string(sval) } - // t.Methodology (string) (string) - case "methodology": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Methodology = string(sval) - } // t.WindowStart (string) (string) case "windowStart": @@ -10977,42 +10907,6 @@ func (t *MediaViewCount) UnmarshalCBOR(r io.Reader) (err error) { t.WindowStart = string(sval) } - // t.ThresholdSegments (int64) (int64) - case "thresholdSegments": - { - - b, err := cr.ReadByte() - if err != nil { - return err - } - if b != cbg.CborNull[0] { - if err := cr.UnreadByte(); err != nil { - return err - } - maj, extra, err := cr.ReadHeader() - if err != nil { - return err - } - var extraI int64 - switch maj { - case cbg.MajUnsignedInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 positive overflow") - } - case cbg.MajNegativeInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 negative overflow") - } - extraI = -1 - extraI - default: - return fmt.Errorf("wrong type for int64 field: %d", maj) - } - - t.ThresholdSegments = (*int64)(&extraI) - } - } default: // Field doesn't exist on this type, so ignore it @@ -11058,26 +10952,19 @@ func (t *MediaViewCount_TrackUsage) MarshalCBOR(w io.Writer) error { } } - // t.TrackId (string) (string) - if len("trackId") > 1000000 { - return xerrors.Errorf("Value in field \"trackId\" was too long") + // t.Track (atproto.RepoStrongRef) (struct) + if len("track") > 1000000 { + return xerrors.Errorf("Value in field \"track\" was too long") } - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("trackId"))); err != nil { + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("track"))); err != nil { return err } - if _, err := cw.WriteString(string("trackId")); err != nil { + if _, err := cw.WriteString(string("track")); err != nil { return err } - if len(t.TrackId) > 1000000 { - return xerrors.Errorf("Value in field t.TrackId was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.TrackId))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.TrackId)); err != nil { + if err := t.Track.MarshalCBOR(cw); err != nil { return err } @@ -11173,16 +11060,25 @@ func (t *MediaViewCount_TrackUsage) UnmarshalCBOR(r io.Reader) (err error) { t.Bytes = int64(extraI) } - // t.TrackId (string) (string) - case "trackId": + // t.Track (atproto.RepoStrongRef) (struct) + case "track": { - sval, err := cbg.ReadStringWithMax(cr, 1000000) + + b, err := cr.ReadByte() if err != nil { return err } + if b != cbg.CborNull[0] { + if err := cr.UnreadByte(); err != nil { + return err + } + t.Track = new(atproto.RepoStrongRef) + if err := t.Track.UnmarshalCBOR(cr); err != nil { + return xerrors.Errorf("unmarshaling t.Track pointer: %w", err) + } + } - t.TrackId = string(sval) } // t.DurationMs (int64) (int64) case "durationMs": diff --git a/pkg/streamplace/mediaviewCount.go b/pkg/streamplace/mediaviewCount.go index 87d950608..1d811951a 100644 --- a/pkg/streamplace/mediaviewCount.go +++ b/pkg/streamplace/mediaviewCount.go @@ -5,6 +5,7 @@ package streamplace import ( + comatproto "github.com/bluesky-social/indigo/api/atproto" lexutil "github.com/bluesky-social/indigo/lex/util" ) @@ -14,15 +15,11 @@ func init() { type MediaViewCount struct { LexiconTypeID string `json:"$type" cborgen:"$type,const=place.stream.media.viewCount"` - // count: Number of distinct sessions that qualified as a view over [windowStart, windowEnd) by the stated methodology. + // count: Number of distinct sessions the reporting node observed as a view over [windowStart, windowEnd). Count int64 `json:"count" cborgen:"count"` // indexedAt: When the reporting node ran this aggregation. Useful for ordering successive reports. IndexedAt string `json:"indexedAt" cborgen:"indexedAt"` - // methodology: Identifier for the counting algorithm. "any-segment" means: a distinct sid that fetched at least `thresholdSegments` segments for this video over the window. Future methodologies (e.g. "ms-from-metafile") will use new tags so older records remain interpretable. - Methodology string `json:"methodology" cborgen:"methodology"` - // thresholdSegments: Floor on segment_request count per (sid, video) for the "any-segment" methodology. Defaults to 1. - ThresholdSegments *int64 `json:"thresholdSegments,omitempty" cborgen:"thresholdSegments,omitempty"` - // tracks: Per-track totals of bytes + playback duration actually transferred over the window. Computed by intersecting each segment_request's HTTP Range with the metafile's per-track byte layout, so the numbers are objective regardless of the chosen view-counting methodology. + // tracks: Per-track totals of bytes + playback duration the reporting node served over the window. Each entry references the place.stream.media.track record whose bytes were transferred, so user-contributed tracks (transcripts, transcodes published by other accounts) attribute naturally to their own records. Tracks []*MediaViewCount_TrackUsage `json:"tracks,omitempty" cborgen:"tracks,omitempty"` // video: AT-URI of the place.stream.video this count is for. Video string `json:"video" cborgen:"video"` @@ -34,12 +31,12 @@ type MediaViewCount struct { // MediaViewCount_TrackUsage is a "trackUsage" in the place.stream.media.viewCount schema. // -// One row of the tracks array: bytes + duration transferred for a single muxlTrack inside the video's MUXL container over the window. trackId matches the muxlTrack record's `trackId` field. +// One row of the tracks array: bytes + duration transferred for a single place.stream.media.track record over the window. type MediaViewCount_TrackUsage struct { - // bytes: Total bytes served from this track's byte ranges over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets. + // bytes: Total bytes served from this track over the window. Sum across attributed segment_requests' Range intersections with the track's segment offsets in the metafile. Bytes int64 `json:"bytes" cborgen:"bytes"` // durationMs: Total playback duration served from this track, in milliseconds. Per HLS segment in the range: (overlap bytes / segment bytes) * segment duration, so partial-segment fetches credit a proportional share of duration. DurationMs int64 `json:"durationMs" cborgen:"durationMs"` - // trackId: Stringified u32 matching the MUXL container's per-track id. - TrackId string `json:"trackId" cborgen:"trackId"` + // track: Strong reference to the place.stream.media.track record whose bytes were transferred. + Track *comatproto.RepoStrongRef `json:"track" cborgen:"track"` } diff --git a/pkg/viewlog/aggregate.go b/pkg/viewlog/aggregate.go index 9a2653bec..0eeee4e5b 100644 --- a/pkg/viewlog/aggregate.go +++ b/pkg/viewlog/aggregate.go @@ -13,30 +13,39 @@ import ( "strings" "time" + comatproto "github.com/bluesky-social/indigo/api/atproto" + "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/vod" ) -// MethodologyAnySegment is the initial counting heuristic: a distinct -// sid that fetched at least `ThresholdSegments` segments for a video -// inside the aggregation window. Recorded on every published view- -// count record so consumers can interpret the number against the -// algorithm that produced it; later methodologies (ms-from-metafile, -// client-reported, …) ship under their own tags. -const MethodologyAnySegment = "any-segment" - // MetafileFetcher resolves a blob CID to its parsed metafile. The // aggregator calls it once per unique CID seen in the window's // segment_request events; tests can substitute a fake. type MetafileFetcher func(ctx context.Context, cid string) (*vod.Metafile, error) +// TrackRefFetcher resolves a blob CID to the strongRefs of the +// place.stream.media.track records whose muxlTrack lives in that +// blob, keyed by the muxlTrack's in-container trackId. The aggregator +// looks each in-container tid up to attribute bytes/duration to the +// owning track record — user-contributed tracks (transcripts, +// transcodes) attribute to their own record, the original tracks +// attribute to the streamer's record. +// +// Returning a nil map (no error) is fine; the corresponding bytes +// just don't get a per-track row in the output. +type TrackRefFetcher func(ctx context.Context, cid string) (map[string]*comatproto.RepoStrongRef, error) + // AggregateInput bundles the window + tunables for one aggregation // pass. WindowStart is inclusive, WindowEnd is exclusive — matches the // AT-URI record's wire shape. type AggregateInput struct { - WindowStart time.Time - WindowEnd time.Time + WindowStart time.Time + WindowEnd time.Time + // ThresholdSegments is the floor on segment_requests per (sid, + // video) before that sid counts as a view. Defaults to 1. + // Operator-level knob, not surfaced on the published record. ThresholdSegments int // ReadMargin extends the file-listing window backwards so a file // opened slightly before WindowStart (whose events may straddle @@ -51,6 +60,13 @@ type AggregateInput struct { // segment_request is counted toward the sid view but contributes // zero bytes / duration. FetchMetafile MetafileFetcher + // FetchTrackRefs resolves a CID to the strongRefs of the + // place.stream.media.track records whose muxlTrack lives in that + // blob, keyed by in-container trackId. Required for per-track + // usage rows; when nil or returning no entry for a given tid the + // bytes/duration credit for that track is dropped (the sid view + // still counts). + FetchTrackRefs TrackRefFetcher } // VideoCount is one row in the AggregateResult: distinct qualifying @@ -64,9 +80,10 @@ type VideoCount struct { // TrackUsage is the objective half of the aggregator's output: how // many bytes + how much playback duration were served from one -// muxlTrack inside the window. Methodology-independent. +// place.stream.media.track record inside the window. The strongRef +// points at the track record whose bytes were transferred. type TrackUsage struct { - TrackID string + Track *comatproto.RepoStrongRef Bytes int64 DurationMS int64 } @@ -162,17 +179,19 @@ func AggregateWindow(ctx context.Context, store blob.Store, in AggregateInput) ( } segments := make(map[pair]int) // trackTotals accumulates objective bytes + duration per - // (video, track). Methodology-independent; the published record - // carries both alongside the heuristic Count. + // (video, track-record). Keyed on (video, trackURI) since the + // strongRef's URI is globally unique — same CID across two + // videos (parent + clip) still attributes to the right tracks. type trackKey struct { - video string - trackID string + video string + trackURI string } trackTotals := make(map[trackKey]*TrackUsage) - // metafileCache: one fetch per CID across the whole window. The - // nil value is cached too, so a missing metafile doesn't trigger - // repeated fetches. + // metafileCache + trackRefCache: one fetch per CID across the + // whole window. The nil values are cached too, so a missing + // metafile / refs lookup doesn't trigger repeated fetches. metafileCache := make(map[string]*vod.Metafile) + trackRefCache := make(map[string]map[string]*comatproto.RepoStrongRef) metafilesLoaded := 0 var eventsRead int @@ -203,10 +222,12 @@ func AggregateWindow(ctx context.Context, store blob.Store, in AggregateInput) ( segments[pair{sid: ev.SID, video: video}]++ eventsRead++ - // Per-track byte/duration accounting needs the + // Per-track byte/duration accounting needs (a) the // metafile that maps blob offsets to per-track - // segments. Cache once per CID; nil cache entries - // mean "tried + missing" so we don't refetch. + // segments and (b) the strongRef for each track's + // place.stream.media.track record. Both cached per + // CID; nil cache entries mean "tried + missing" so + // we don't refetch. if ev.CID == "" || ev.RangeEnd < ev.RangeStart { return } @@ -226,15 +247,35 @@ func AggregateWindow(ctx context.Context, store blob.Store, in AggregateInput) ( if meta == nil { return } + refs, refsCached := trackRefCache[ev.CID] + if !refsCached { + if in.FetchTrackRefs != nil { + r, err := in.FetchTrackRefs(ctx, ev.CID) + if err != nil { + log.Debug(ctx, "viewlog: fetch track refs", + "cid", ev.CID, "error", err) + } + refs = r + } + trackRefCache[ev.CID] = refs + } for tid, track := range meta.Tracks { bytes, durMS := rangeOverlapInTrack(ev.RangeStart, ev.RangeEnd, track) if bytes == 0 { continue } - k := trackKey{video: video, trackID: tid} + ref := refs[tid] + if ref == nil { + // No track record found for this in-container + // tid. Drop the credit — a TrackUsage row + // without a stable strongRef wouldn't be + // useful to consumers. + continue + } + k := trackKey{video: video, trackURI: ref.Uri} tu, ok := trackTotals[k] if !ok { - tu = &TrackUsage{TrackID: tid} + tu = &TrackUsage{Track: ref} trackTotals[k] = tu } tu.Bytes += bytes @@ -276,8 +317,12 @@ func AggregateWindow(ctx context.Context, store blob.Store, in AggregateInput) ( for v := range videoSet { tracks := tracksByVideo[v] // Stable per-video track order so re-runs produce byte- - // identical records. - sort.Slice(tracks, func(i, j int) bool { return tracks[i].TrackID < tracks[j].TrackID }) + // identical records. Sort on the strongRef URI (globally + // unique) since that's the only stable identity now that + // in-container tids aren't surfaced. + sort.Slice(tracks, func(i, j int) bool { + return tracks[i].Track.Uri < tracks[j].Track.Uri + }) out = append(out, VideoCount{ VideoURI: v, Count: perVideo[v], diff --git a/pkg/viewlog/aggregate_test.go b/pkg/viewlog/aggregate_test.go index 42ef0497f..84c534144 100644 --- a/pkg/viewlog/aggregate_test.go +++ b/pkg/viewlog/aggregate_test.go @@ -7,6 +7,7 @@ import ( "testing" "time" + comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/stretchr/testify/require" @@ -14,6 +15,24 @@ import ( "stream.place/streamplace/pkg/vod" ) +// fixtureTrackRefs returns the strongRefs that pair with the fixture +// metafile, keyed by in-container trackId. Same shape the real +// resolver in cmd/streamplace.go produces from MediaTrack rows. +func fixtureTrackRefs() map[string]*comatproto.RepoStrongRef { + return map[string]*comatproto.RepoStrongRef{ + "1": { + LexiconTypeID: "com.atproto.repo.strongRef", + Uri: "at://did:plc:alice/place.stream.media.track/track-video-1", + Cid: "bafytrack1", + }, + "2": { + LexiconTypeID: "com.atproto.repo.strongRef", + Uri: "at://did:plc:alice/place.stream.media.track/track-audio-2", + Cid: "bafytrack2", + }, + } +} + // loadStore returns an in-memory-ish FileStore + a writer pointed at // it, sharing the same NodeDID across test cases so the file naming // stays predictable. @@ -323,8 +342,8 @@ func fixtureMetafile() *vod.Metafile { Type: "audio", Timescale: 1000, Segments: []vod.MetafileSegment{ - {Offset: 400, Size: 100, DurationTicks: 2000}, // [400..499] = 2s - {Offset: 900, Size: 100, DurationTicks: 2000}, // [900..999] = 2s + {Offset: 400, Size: 100, DurationTicks: 2000}, // [400..499] = 2s + {Offset: 900, Size: 100, DurationTicks: 2000}, // [900..999] = 2s {Offset: 1400, Size: 100, DurationTicks: 2000}, }, }, @@ -407,6 +426,7 @@ func TestAggregateWindowTrackUsage(t *testing.T) { store, err := blob.NewFileStore(root) require.NoError(t, err) + refs := fixtureTrackRefs() res, err := AggregateWindow(context.Background(), store, AggregateInput{ WindowStart: now.Add(-time.Minute), WindowEnd: now.Add(time.Minute), @@ -414,18 +434,26 @@ func TestAggregateWindowTrackUsage(t *testing.T) { require.Equal(t, cidA, cid) return fixtureMetafile(), nil }, + FetchTrackRefs: func(ctx context.Context, cid string) (map[string]*comatproto.RepoStrongRef, error) { + require.Equal(t, cidA, cid) + return refs, nil + }, }) require.NoError(t, err) require.Len(t, res.VideoCounts, 1) require.Equal(t, int64(1), res.VideoCounts[0].Count) require.Equal(t, 1, res.MetafilesLoaded, "metafile cache: one fetch even though three segment_requests share the CID") - byTID := map[string]TrackUsage{} + byURI := map[string]TrackUsage{} for _, t := range res.VideoCounts[0].Tracks { - byTID[t.TrackID] = t + byURI[t.Track.Uri] = t } - require.Equal(t, TrackUsage{TrackID: "1", Bytes: 800, DurationMS: 4000}, byTID["1"], "two video segments fully fetched") - require.Equal(t, TrackUsage{TrackID: "2", Bytes: 100, DurationMS: 2000}, byTID["2"], "one audio segment fully fetched") + require.Equal(t, + TrackUsage{Track: refs["1"], Bytes: 800, DurationMS: 4000}, + byURI[refs["1"].Uri], "two video segments fully fetched") + require.Equal(t, + TrackUsage{Track: refs["2"], Bytes: 100, DurationMS: 2000}, + byURI[refs["2"].Uri], "one audio segment fully fetched") } func TestAggregateWindowVideoWithUsageButNoCount(t *testing.T) { @@ -449,6 +477,7 @@ func TestAggregateWindowVideoWithUsageButNoCount(t *testing.T) { store, err := blob.NewFileStore(root) require.NoError(t, err) + refs := fixtureTrackRefs() res, err := AggregateWindow(context.Background(), store, AggregateInput{ WindowStart: now.Add(-time.Minute), WindowEnd: now.Add(time.Minute), @@ -456,15 +485,60 @@ func TestAggregateWindowVideoWithUsageButNoCount(t *testing.T) { FetchMetafile: func(ctx context.Context, cid string) (*vod.Metafile, error) { return fixtureMetafile(), nil }, + FetchTrackRefs: func(ctx context.Context, cid string) (map[string]*comatproto.RepoStrongRef, error) { + return refs, nil + }, }) require.NoError(t, err) require.Len(t, res.VideoCounts, 1) require.Equal(t, int64(0), res.VideoCounts[0].Count, "below threshold ⇒ no view") require.Len(t, res.VideoCounts[0].Tracks, 1, "but the bytes are still reported") - require.Equal(t, "1", res.VideoCounts[0].Tracks[0].TrackID) + require.Equal(t, refs["1"].Uri, res.VideoCounts[0].Tracks[0].Track.Uri) require.Equal(t, int64(400), res.VideoCounts[0].Tracks[0].Bytes) } +func TestAggregateWindowTidWithoutStrongRefIsDropped(t *testing.T) { + // A track that has a metafile entry but no place.stream.media.track + // record (e.g. mid-publish, or the record was deleted) doesn't get + // a usage row — we can't reference it. The sid view still counts. + root := t.TempDir() + now := time.Date(2026, 5, 17, 12, 0, 0, 0, time.UTC) + w := newAggTestWriter(t, root, "did:web:node1", now) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go w.Run(ctx) + + const ( + videoA = "at://did:plc:alice/place.stream.video/v1" + cidA = "bafyfixture" + ) + w.Log(ctx, Event{Ts: now, Type: EventTypeManifestRequest, VideoURI: videoA, SID: "sid"}) + // Range hits both video (track "1") and audio (track "2"), but the + // resolver only knows about track "1". + w.Log(ctx, Event{Ts: now.Add(time.Second), Type: EventTypeSegmentRequest, CID: cidA, SID: "sid", RangeStart: 0, RangeEnd: 499}) + require.NoError(t, w.Close()) + + store, err := blob.NewFileStore(root) + require.NoError(t, err) + refs := fixtureTrackRefs() + partial := map[string]*comatproto.RepoStrongRef{"1": refs["1"]} + res, err := AggregateWindow(context.Background(), store, AggregateInput{ + WindowStart: now.Add(-time.Minute), + WindowEnd: now.Add(time.Minute), + FetchMetafile: func(ctx context.Context, cid string) (*vod.Metafile, error) { + return fixtureMetafile(), nil + }, + FetchTrackRefs: func(ctx context.Context, cid string) (map[string]*comatproto.RepoStrongRef, error) { + return partial, nil + }, + }) + require.NoError(t, err) + require.Len(t, res.VideoCounts, 1) + require.Equal(t, int64(1), res.VideoCounts[0].Count) + require.Len(t, res.VideoCounts[0].Tracks, 1, "only the resolved track gets a row") + require.Equal(t, refs["1"].Uri, res.VideoCounts[0].Tracks[0].Track.Uri) +} + func TestAggregateWindowMissingMetafileDoesNotPoison(t *testing.T) { // A segment_request whose CID has no metafile still counts toward // the sid view (the threshold cares about request count, not byte diff --git a/pkg/viewlog/publish.go b/pkg/viewlog/publish.go index c7c5c5365..fedbb9756 100644 --- a/pkg/viewlog/publish.go +++ b/pkg/viewlog/publish.go @@ -24,12 +24,13 @@ var aggregateTracer = otel.Tracer("viewlog") // RunAggregationInput plumbs everything the bootstrap-installed // aggregator function needs to do one window pass end-to-end. type RunAggregationInput struct { - Store blob.Store - CLI *config.CLI - WindowStart time.Time - WindowEnd time.Time - ReadMargin time.Duration - FetchMetafile MetafileFetcher + Store blob.Store + CLI *config.CLI + WindowStart time.Time + WindowEnd time.Time + ReadMargin time.Duration + FetchMetafile MetafileFetcher + FetchTrackRefs TrackRefFetcher } // RunAggregation reads logs for the window, computes per-video view @@ -46,10 +47,11 @@ func RunAggregation(ctx context.Context, in RunAggregationInput) error { defer span.End() result, err := AggregateWindow(ctx, in.Store, AggregateInput{ - WindowStart: in.WindowStart, - WindowEnd: in.WindowEnd, - ReadMargin: in.ReadMargin, - FetchMetafile: in.FetchMetafile, + WindowStart: in.WindowStart, + WindowEnd: in.WindowEnd, + ReadMargin: in.ReadMargin, + FetchMetafile: in.FetchMetafile, + FetchTrackRefs: in.FetchTrackRefs, }) if err != nil { return fmt.Errorf("aggregate window: %w", err) @@ -68,26 +70,23 @@ func RunAggregation(ctx context.Context, in RunAggregationInput) error { ) indexedAt := aqtime.FromTime(time.Now().UTC()).String() - threshold := int64(result.Window.ThresholdSegments) for _, vc := range result.VideoCounts { tracks := make([]*streamplace.MediaViewCount_TrackUsage, 0, len(vc.Tracks)) for _, t := range vc.Tracks { tracks = append(tracks, &streamplace.MediaViewCount_TrackUsage{ - TrackId: t.TrackID, + Track: t.Track, Bytes: t.Bytes, DurationMs: t.DurationMS, }) } rec := &streamplace.MediaViewCount{ - LexiconTypeID: constants.PLACE_STREAM_MEDIA_VIEW_COUNT, - Video: vc.VideoURI, - Count: vc.Count, - WindowStart: in.WindowStart.UTC().Format(time.RFC3339), - WindowEnd: in.WindowEnd.UTC().Format(time.RFC3339), - Methodology: MethodologyAnySegment, - ThresholdSegments: &threshold, - Tracks: tracks, - IndexedAt: indexedAt, + LexiconTypeID: constants.PLACE_STREAM_MEDIA_VIEW_COUNT, + Video: vc.VideoURI, + Count: vc.Count, + WindowStart: in.WindowStart.UTC().Format(time.RFC3339), + WindowEnd: in.WindowEnd.UTC().Format(time.RFC3339), + Tracks: tracks, + IndexedAt: indexedAt, } rkey, err := viewCountRkey(vc.VideoURI, in.WindowStart) if err != nil { -- 2.51.2