Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
18 kB · 525 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526package viewlog
import ( "bufio" "compress/gzip" "context" "encoding/json" "errors" "fmt" "io" "path" "sort" "strings" "time"
"stream.place/streamplace/pkg/comatproto"
"stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/vod")
// 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 // 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 // the boundary) is still read. Per-event filtering by ev.Ts keeps // the count correct; this just decides which files to crack open. // Defaults to 1 hour, which is generously larger than the writer's // flush interval. ReadMargin time.Duration // FetchMetafile resolves a CID to its parsed metafile. Defaults to // reading blobs/<cid>.json from the same store the logs live in. // A nil return + nil error means "metafile not found" and the // 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// sids per place.stream.video over the window plus per-track totals// of bytes + duration transferred.type VideoCount struct { VideoURI string Count int64 Tracks []TrackUsage}
// TrackUsage is the objective half of the aggregator's output: how// many bytes + how much playback duration were served from one// place.stream.media.track record inside the window. The strongRef// points at the track record whose bytes were transferred.type TrackUsage struct { Track comatproto.RepoStrongRef Bytes int64 DurationMS int64}
// AggregateResult is the return shape of AggregateWindow. The caller// turns each VideoCount into a place.stream.media.viewCount record.type AggregateResult struct { Window AggregateInput VideoCounts []VideoCount // FilesRead and EventsRead surface aggregation effort for ops // observability; not load-bearing for the records themselves. FilesRead int EventsRead int MetafilesLoaded int}
// AggregateWindow lists every view-log file under view-logs/ whose// filename timestamp falls in [WindowStart - ReadMargin, WindowEnd),// reads each as gzipped JSONL, and produces per-video totals://// - Count: distinct sids that hit ≥ ThresholdSegments segment_requests// - Tracks: per-track bytes + duration credited by intersecting each// segment_request's HTTP Range with the metafile's byte layout//// Manifest events provide the sid → video URI mapping segment events// inherit from; segment_requests for a sid we never saw a manifest// for are unattributed and dropped (no video to credit).func AggregateWindow(ctx context.Context, store blob.Store, in AggregateInput) (*AggregateResult, error) { if in.WindowEnd.Before(in.WindowStart) || in.WindowEnd.Equal(in.WindowStart) { return nil, fmt.Errorf("viewlog: AggregateWindow needs a positive window, got start=%s end=%s", in.WindowStart, in.WindowEnd) } threshold := in.ThresholdSegments if threshold <= 0 { threshold = 1 } readMargin := in.ReadMargin if readMargin <= 0 { readMargin = time.Hour } fetchMetafile := in.FetchMetafile if fetchMetafile == nil { fetchMetafile = func(ctx context.Context, cid string) (*vod.Metafile, error) { return fetchMetafileFromStore(ctx, store, cid) } }
keys, err := store.List(ctx, viewLogsPrefix) if err != nil { return nil, fmt.Errorf("viewlog: list logs: %w", err) }
// Pick files whose openedAt timestamp falls in // [WindowStart - ReadMargin, WindowEnd). Events inside each file // are filtered to the exact window below — the read margin just // makes sure we don't skip a file whose contents straddle the // boundary. readStart := in.WindowStart.Add(-readMargin) type pickedFile struct { key string ts time.Time } var picked []pickedFile for _, k := range keys { ts, ok := parseViewLogKeyTime(k) if !ok { continue } if ts.Before(readStart) || !ts.Before(in.WindowEnd) { continue } picked = append(picked, pickedFile{key: k, ts: ts}) } // Time-order the read so per-sid manifest_request events arrive // before their segment_requests, regardless of which node wrote // them. Pure-lex sort on the key would interleave by node-DID // before timestamp; we want timestamp first. sort.Slice(picked, func(i, j int) bool { if picked[i].ts.Equal(picked[j].ts) { return picked[i].key < picked[j].key } return picked[i].ts.Before(picked[j].ts) })
// sidVideo holds the most recent video URI a sid asked about. A // segment_request inherits this; if the sid is unknown, the event // is dropped (unattributed). sidVideo := make(map[string]string) // segments counts segment_requests per (sid, video) pair. type pair struct { sid string video string } segments := make(map[pair]int) // trackTotals accumulates objective bytes + duration per // (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 trackURI string } trackTotals := make(map[trackKey]*TrackUsage) // 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 for _, pf := range picked { if err := readJSONLGz(ctx, store, pf.key, func(ev *Event) { // Manifest events flow into the sid→video map even when // they're outside the window — a session might have // started before WindowStart but kept fetching segments // inside it. Only segment_request events are clipped to // [WindowStart, WindowEnd). switch ev.Type { case EventTypeManifestRequest: if ev.SID != "" && ev.VideoURI != "" { sidVideo[ev.SID] = ev.VideoURI } eventsRead++ case EventTypeSegmentRequest: if ev.Ts.Before(in.WindowStart) || !ev.Ts.Before(in.WindowEnd) { return } if ev.SID == "" { return } video := sidVideo[ev.SID] if video == "" { return } segments[pair{sid: ev.SID, video: video}]++ eventsRead++
// Per-track byte/duration accounting needs (a) the // metafile that maps blob offsets to per-track // 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 } meta, cached := metafileCache[ev.CID] if !cached { m, err := fetchMetafile(ctx, ev.CID) if err != nil { log.Debug(ctx, "viewlog: fetch metafile", "cid", ev.CID, "error", err) } meta = m metafileCache[ev.CID] = meta if meta != nil { metafilesLoaded++ } } 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 } // CDN-ingested events carry bytes-sent but no Range; // prorate across tracks by byte share. Node-logged // events carry a Range over the blob, which includes // the flat-MP4 header the metafile's fragment-relative // offsets sit behind, so shift it before intersecting. prorated := ev.BytesSent > 0 && ev.RangeStart == 0 && ev.RangeEnd == 0 totalBody := metafileBodyBytes(meta) for tid, track := range meta.Tracks { var bytes, durMS int64 if prorated { bytes, durMS = prorateInTrack(ev.BytesSent, totalBody, track) } else { bytes, durMS = rangeOverlapInTrack(ev.RangeStart-meta.FlatHeaderSize, ev.RangeEnd-meta.FlatHeaderSize, track) } if bytes == 0 { continue } ref := refs[tid] if ref.Uri == "" { // 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{Track: ref} trackTotals[k] = tu } tu.Bytes += bytes tu.DurationMS += durMS } } }); err != nil { // Per-file read failures are logged + skipped rather than // aborting the entire window. One corrupt file shouldn't // drop every video's count. log.Error(ctx, "viewlog: read log file", "key", pf.key, "error", err) continue } }
// Tally distinct qualifying sids per video. perVideo := make(map[string]int64) for p, n := range segments { if n >= threshold { perVideo[p.video]++ } }
// Group tracks under each video. A video can show up here with // zero qualifying sids (bytes were transferred but every sid fell // below threshold) — emit it anyway so the objective totals // aren't lost. videoSet := make(map[string]struct{}) for v := range perVideo { videoSet[v] = struct{}{} } tracksByVideo := make(map[string][]TrackUsage) for k, tu := range trackTotals { tracksByVideo[k.video] = append(tracksByVideo[k.video], *tu) videoSet[k.video] = struct{}{} }
out := make([]VideoCount, 0, len(videoSet)) for v := range videoSet { tracks := tracksByVideo[v] // Stable per-video track order so re-runs produce byte- // 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], Tracks: tracks, }) } // Sort for deterministic output (mostly for tests + record-key // stability when callers iterate in slice order). sort.Slice(out, func(i, j int) bool { return out[i].VideoURI < out[j].VideoURI })
return &AggregateResult{ Window: AggregateInput{WindowStart: in.WindowStart, WindowEnd: in.WindowEnd, ThresholdSegments: threshold, ReadMargin: readMargin}, VideoCounts: out, FilesRead: len(picked), EventsRead: eventsRead, MetafilesLoaded: metafilesLoaded, }, nil}
// rangeOverlapInTrack returns the (bytes, durationMS) credit one// HTTP Range request earns against a single muxlTrack's byte layout.//// HLS playback typically issues one byte-range request per// EXT-X-BYTERANGE entry, so the request usually matches one// MetafileSegment exactly. Partial overlaps (chunked range fetches,// resumed downloads) credit a proportional share of duration so the// numbers stay roughly honest under either pattern.//// rangeStart and rangeEnd are inclusive (RFC 7233). track.Segments'// Offset is the absolute byte offset within the blob; Size is the// segment's byte length; DurationTicks is the segment's playback// length in track timescale units.func rangeOverlapInTrack(rangeStart, rangeEnd int64, track vod.MetafileTrack) (bytes, durationMS int64) { if rangeEnd < rangeStart || track.Timescale == 0 { return 0, 0 } for _, seg := range track.Segments { if seg.Size <= 0 { continue } segLast := seg.Offset + seg.Size - 1 overlapStart := rangeStart if seg.Offset > overlapStart { overlapStart = seg.Offset } overlapEnd := rangeEnd if segLast < overlapEnd { overlapEnd = segLast } if overlapEnd < overlapStart { continue } overlap := overlapEnd - overlapStart + 1 bytes += overlap // Byte-proportional duration: (overlap / size) * segment-ms. // Avoid float math so deterministic runs produce identical // records. segDurMS := int64(seg.DurationTicks) * 1000 / int64(track.Timescale) durationMS += segDurMS * overlap / seg.Size } return bytes, durationMS}
// metafileBodyBytes is the byte length of every track's segments// summed: the denominator for prorating a bytes-sent figure.func metafileBodyBytes(meta *vod.Metafile) int64 { var total int64 for _, t := range meta.Tracks { for _, seg := range t.Segments { total += seg.Size } } return total}
// prorateInTrack credits one track with its byte-share of a request// whose size we know but whose Range we don't (CDN access logs). The// track gets bytesSent * (trackBytes / totalBody) bytes and the same// fraction of its playback duration. bytesSent is capped at the blob// body so a whole-blob download that also carried the flat header// doesn't over-credit.func prorateInTrack(bytesSent, totalBody int64, track vod.MetafileTrack) (bytes, durationMS int64) { if bytesSent <= 0 || totalBody <= 0 || track.Timescale == 0 { return 0, 0 } if bytesSent > totalBody { bytesSent = totalBody } var trackBytes int64 var trackTicks uint64 for _, seg := range track.Segments { trackBytes += seg.Size trackTicks += seg.DurationTicks } if trackBytes == 0 { return 0, 0 } // Integer math throughout so re-runs are byte-identical. bytes = bytesSent * trackBytes / totalBody trackDurMS := int64(trackTicks) * 1000 / int64(track.Timescale) durationMS = trackDurMS * bytesSent / totalBody return bytes, durationMS}
// viewLogsPrefix is the top of the per-node log tree, shared by every// writer that targets the same store.const viewLogsPrefix = "view-logs/"
// keyTimeFormat matches the suffix the Writer attaches to each rotated// blob. RFC3339-ish but with ':' replaced by '-' so the key stays a// valid filename on every platform. Nanosecond precision keeps// closely-spaced flushes (size-triggered storms during a popular video's// hot moment) from sharing a key and overwriting each other.const keyTimeFormat = "2006-01-02T15-04-05.000000000Z"
// parseViewLogKeyTime extracts the rotated-at timestamp from a key// shaped like `view-logs/<node-did>/<window>.jsonl.gz`. Returns// (zero, false) for anything that doesn't fit the shape.func parseViewLogKeyTime(key string) (time.Time, bool) { base := path.Base(key) const suffix = ".jsonl.gz" if !strings.HasSuffix(base, suffix) { return time.Time{}, false } stamp := strings.TrimSuffix(base, suffix) t, err := time.Parse(keyTimeFormat, stamp) if err != nil { return time.Time{}, false } return t.UTC(), true}
// fetchMetafileFromStore is the default MetafileFetcher: read// blobs/<cid>.json from the same blob.Store the logs live in.// Returns (nil, nil) when the metafile is missing so callers can// distinguish "no such metafile" from "store error".func fetchMetafileFromStore(ctx context.Context, store blob.Store, cid string) (*vod.Metafile, error) { key := vod.BlobsPrefix + cid + ".json" r, err := store.Open(ctx, key) if err != nil { if errors.Is(err, blob.ErrNotFound) { return nil, nil } return nil, fmt.Errorf("open metafile %s: %w", key, err) } defer r.Close() body := make([]byte, r.Size()) if _, err := r.ReadAt(body, 0); err != nil && !errors.Is(err, io.EOF) { return nil, fmt.Errorf("read metafile %s: %w", key, err) } var meta vod.Metafile if err := json.Unmarshal(body, &meta); err != nil { return nil, fmt.Errorf("parse metafile %s: %w", key, err) } return &meta, nil}
// readJSONLGz opens key, gunzips, and calls fn for each decoded Event.// JSON parse failures on a single line are logged + skipped rather// than aborting the whole file — a malformed line shouldn't lose// every event that follows it.func readJSONLGz(ctx context.Context, store blob.Store, key string, fn func(*Event)) error { r, err := store.Open(ctx, key) if err != nil { return fmt.Errorf("open %s: %w", key, err) } defer r.Close() gz, err := gzip.NewReader(io.NewSectionReader(r, 0, r.Size())) if err != nil { return fmt.Errorf("gunzip %s: %w", key, err) } defer gz.Close() sc := bufio.NewScanner(gz) // JSONL lines can be larger than the default 64KB scanner buffer. sc.Buffer(make([]byte, 0, 64*1024), 1<<20) for sc.Scan() { var ev Event if err := json.Unmarshal(sc.Bytes(), &ev); err != nil { log.Debug(ctx, "viewlog: skipping malformed line", "key", key, "error", err) continue } fn(&ev) } if err := sc.Err(); err != nil && !errors.Is(err, io.EOF) { return fmt.Errorf("scan %s: %w", key, err) } return nil}