Documentation
¶
Overview ¶
Package signal implements ContentKit's signal plane: an append-only stream of host-defined interaction signals plus a durable per-(subject, content) current-state projection, both stored in ClickHouse.
The signal plane records *what each subject did and thought about each content reference* and projects that into fast reads for history, unseen, engagement, popularity, and (above this package) personalized search and recommendations.
Mechanism vs meaning: this package owns storage, aggregation, and queries. Signal types, content kinds, scoring weights, and completion rules are host-defined data — no business noun appears in the schema. Every key, index and cursor carries the tenant.
Index ¶
- Constants
- func CheckSchema(ctx context.Context, conn Conn, database string) error
- func CreateDatabase(ctx context.Context, conn Exec, database, cluster string) error
- type AttributedClick
- type AttributedRender
- type Attribution
- type AttributionOptions
- type AttributionPage
- type CoEngagedHit
- type CoEngagedOptions
- type Conn
- type ContentKey
- type ContentMetrics
- type ContentRef
- type ErasureReport
- type Exec
- type Exposure
- type ExposureStage
- type HistoryOptions
- type HistoryStatus
- type InventoryRow
- type LimitError
- type Placement
- type PopularHit
- type PopularOptions
- type ProjectionKey
- type RankWeights
- type RefreshCoEngagementOptions
- type RepairOptions
- type RepairResult
- type SchemaMismatchError
- type Scored
- type Scorer
- type ScorerFunc
- type Signal
- type State
- type StateRow
- type Store
- func (st *Store) Attribution(ctx context.Context, tenant string, opts AttributionOptions) (AttributionPage, error)
- func (st *Store) CoEngaged(ctx context.Context, tenant string, ref ContentRef, opts CoEngagedOptions) ([]CoEngagedHit, error)
- func (st *Store) EnforceErasures(ctx context.Context, tenant string) (ErasureReport, error)
- func (st *Store) EraseSubjects(ctx context.Context, tenants []string, subjects []Subject) (ErasureReport, error)
- func (st *Store) ErasedSubjects(ctx context.Context, tenant string, subjects []Subject) (map[Subject]bool, error)
- func (st *Store) Forget(ctx context.Context, tenant string, subject Subject, contentKind string, ...) error
- func (st *Store) ForgetExposures(ctx context.Context, tenant string, subject Subject) error
- func (st *Store) History(ctx context.Context, tenant string, subject Subject, opts HistoryOptions) ([]StateRow, error)
- func (st *Store) HistoryCount(ctx context.Context, tenant string, subject Subject, opts HistoryOptions) (int64, error)
- func (st *Store) Inventory(ctx context.Context, tenant string) ([]InventoryRow, error)
- func (st *Store) Metrics(ctx context.Context, tenant string, refs []ContentRef, window Window) (map[ContentKey]ContentMetrics, error)
- func (st *Store) NegativeIDs(ctx context.Context, tenant string, subject Subject, contentKinds []string) (map[ContentKey]struct{}, error)
- func (st *Store) Popular(ctx context.Context, tenant string, contentKind string, opts PopularOptions) ([]PopularHit, error)
- func (st *Store) PopularityFor(ctx context.Context, tenant string, contentKind string, ids []string, ...) (map[string]float64, error)
- func (st *Store) PurgeContentKinds(ctx context.Context, tenant string, contentKinds []string) error
- func (st *Store) RecordExposures(ctx context.Context, tenant string, exposures []Exposure) error
- func (st *Store) RecordSignals(ctx context.Context, tenant string, signals []Signal) error
- func (st *Store) RefreshCoEngagement(ctx context.Context, tenant string, opts RefreshCoEngagementOptions) error
- func (st *Store) RepairProjections(ctx context.Context, tenant string, opts RepairOptions) (RepairResult, error)
- func (st *Store) SeenIDs(ctx context.Context, tenant string, subject Subject, contentKind string) (map[string]struct{}, error)
- func (st *Store) States(ctx context.Context, tenant string, subject Subject, refs []ContentRef) (map[ContentKey]State, error)
- func (st *Store) TopStates(ctx context.Context, tenant string, subject Subject, opts TopStatesOptions) ([]StateRow, error)
- type Subject
- type TopStatesOptions
- type Window
Constants ¶
const ( DefaultAttributionRenders = 500 DefaultAttributionClicks = 5000 )
Default page bounds of Attribution.
const ( MaxSignalsPerBatch = 500 MaxExposuresPerBatch = 200 MaxShownPerExposure = 200 // MaxIdentifierBytes bounds content kinds/ids/versions, subject keys, // signal types, event ids, query ids, surfaces and languages. MaxIdentifierBytes = 256 MaxResumeBytes = 256 MaxQueryBytes = 256 // MaxPayloadBytes bounds the JSON-encoded payload. MaxPayloadBytes = 2048 )
Ingestion bounds. Oversize input is rejected, never truncated, so hosts see it as backpressure and can count it instead of silently losing detail.
const ( SubjectKindUser = "user" SubjectKindAnon = "anon" )
SubjectKind values stored in the subject_kind column.
const ( SurfaceSearch = "search" SurfaceForYou = "foryou" SurfaceSimilar = "similar" SurfacePopular = "popular" SurfaceOrganic = "organic" )
Render surfaces. Hosts may use other values; these are the common ones evaluation distinguishes.
const ( PayloadKeyRenderID = "render_id" PayloadKeySurface = "surface" PayloadKeyPosition = "position" )
Standardized click-attribution payload keys, set via WithAttribution so evaluation can join clicks to exposures on the render id.
const MaxErasureSubjects = 100
MaxErasureSubjects bounds one EraseSubjects call.
const TypeClick = "click"
TypeClick is the selection signal type Attribution joins to exposures.
const TypeView = "view"
TypeView is the consumption signal type. Viewer counts, popularity and progress/completion/resume state read only view events; clicks and feedback stay separate signals.
Variables ¶
This section is empty.
Functions ¶
func CheckSchema ¶
CheckSchema verifies, read-only, that the signal database has exactly the tables, columns, engines and keys this library version uses. Hosts call it at startup with runtime credentials and must not serve analytics on error: table existence alone is not compatibility. Requires SELECT on system.tables and system.columns for the database.
func CreateDatabase ¶
CreateDatabase creates the dedicated signal database (ON CLUSTER when cluster is set). The schema itself is owned by the versioned migrations in migrations.SignalClickHouse, applied with migratekit chmigrate; migratekit connects to the database, so it must exist first.
Types ¶
type AttributedClick ¶
type AttributedClick struct {
ContentRef
Subject Subject // who clicked; may differ from the render's subject
Position uint32 // position claimed by the click
OccurredAt time.Time
EventID string
// Exposed reports whether the clicked content is in the render's list at
// the requested stage. False means "not exposed at this stage", never negative.
Exposed bool
}
AttributedClick is a canonical click carrying a render id.
type AttributedRender ¶
type AttributedRender struct {
RenderID string
QueryID string
Surface string
Ranker string
Language string
Subject Subject
OccurredAt time.Time
Shown []Placement
Clicks []AttributedClick
// Continued is true when earlier clicks of this render were on the previous
// page: the header repeats, the clicks do not. Merge by RenderID.
Continued bool
}
AttributedRender is one render at one stage with its clicks joined.
type Attribution ¶
type Attribution struct {
RenderID string // the Exposure.RenderID the click came from
Surface string
Position uint32 // 1-based position of the clicked item within that render
}
Attribution links a click/engagement signal to the render that produced it.
type AttributionOptions ¶
type AttributionOptions struct {
// Window bounds exposure and click days (zero = all time).
Window Window
// Stage is the exposure stage clicks are attributed against (required).
Stage ExposureStage
// Surface restricts renders to one surface (empty = all).
Surface string
// After is the opaque cursor from the previous page's Next (empty = start).
After string
// Limit caps renders per page (default DefaultAttributionRenders).
Limit int
// ClickLimit caps click rows per page, attributed and unattributed
// together (default DefaultAttributionClicks). A render with more clicks
// than fit continues on the next page.
ClickLimit int
}
AttributionOptions selects an evaluation export. Stage is required. A paged export passes the previous page's Next as After with the same Window, Stage and Surface; Limit and ClickLimit may change between pages.
type AttributionPage ¶
type AttributionPage struct {
Renders []AttributedRender
// Unattributed are clicks in the window whose render has no exposure at
// the requested stage (prefetched, cancelled, hidden or never recorded):
// they cannot label anything.
Unattributed []AttributedClick
// Next is the cursor for the following page; empty when done.
Next string
}
AttributionPage is one page of the evaluation export.
type CoEngagedHit ¶
type CoEngagedHit struct {
ContentRef
Strength int64
}
CoEngagedHit is one co-engaged work. Strength is the NET co-engaged subject count: subjects with non-negative engagement count for, subjects with negative explicit feedback count against.
type CoEngagedOptions ¶
type CoEngagedOptions struct {
// ContentKinds limits result kinds. Empty = all.
ContentKinds []string
// Window bounds the event scan. Zero = all time.
Window Window
// MaxSubjects caps the engaged-subject set read from the anchor
// (defaults to 10000).
MaxSubjects int
// SkipRollup forces the query-time event scan even when the content_pairs
// rollup has rows (e.g. for freshness-critical reads).
SkipRollup bool
Limit int // default 50
}
CoEngagedOptions controls co-engagement queries ("subjects who engaged with X also engaged with Y"). Co-engagement is work-level.
type Conn ¶
type Conn interface {
Exec(ctx context.Context, query string, args ...any) error
Query(ctx context.Context, query string, args ...any) (driver.Rows, error)
}
Conn is the minimal ClickHouse surface the Store needs; clickhouse-go's driver.Conn satisfies it.
type ContentKey ¶
type ContentKey = contentref.ContentKey
ContentKey is the comparable form of a ContentRef (map keys).
type ContentMetrics ¶
type ContentMetrics struct {
Viewers uint64 // subjects with at least one view
UserViewers uint64
AnonViewers uint64
Views uint64 // view events (sessions)
Completions uint64 // completed views
Completers uint64 // subjects with a completed view
ActiveS uint64 // summed view DurationS
ScoreSum int64 // summed view scores; ScoreSum/Views is the mean
// ViewerEngagementSum sums each viewer's mean session score, normalized
// and clamped to [0,1]. Every viewer contributes at most one unit.
ViewerEngagementSum float64
ReturningViewers uint64 // subjects with more than one canonical view in the window
Events uint64 // canonical events of every type
ValueSum float64
// PositiveSubjects / NegativeSubjects: subjects whose summed feedback
// Value in the window is > 0 / < 0.
PositiveSubjects uint64
NegativeSubjects uint64
SignalCounts map[string]uint64 // canonical events per type
}
ContentMetrics are named, separately defined statistics for one content reference over a window, computed from canonical events. Each subject counts once per metric that says "subjects"; sessions and feedback are never relabeled as views.
type ContentRef ¶
type ContentRef = contentref.ContentRef
ContentRef is the tenant-scoped content reference every signal carries.
type ErasureReport ¶
type ErasureReport struct {
// Remaining counts rows still attributable to the subjects per table after
// the pass; every value is zero when erasure is complete.
Remaining map[string]uint64
// PairsRemoved counts content_pairs rows invalidated because an erased
// subject contributed to a work in them; RefreshCoEngagement rebuilds them.
PairsRemoved uint64
}
ErasureReport describes one erasure or enforcement pass.
func (ErasureReport) Complete ¶
func (r ErasureReport) Complete() bool
Complete reports whether no attributable rows remain.
type Exec ¶
Exec is the minimal ClickHouse execution surface CreateDatabase needs; clickhouse-go's driver.Conn satisfies it.
type Exposure ¶
type Exposure struct {
RenderID string // required, stable per render, shared with its clicks
Stage ExposureStage
Revision uint64
QueryID string
Surface string
Ranker string // ranking configuration identity for offline comparison
Language string
Subject Subject // optional for anonymous renders
Shown []Placement
OccurredAt time.Time // required
}
Exposure is one result list at one stage: the cumulative set of placements the stage reached, as one row. Identity is (tenant, RenderID, Stage); re-sending replaces, a higher Revision supersedes (a visible list that grows as the subject scrolls). Never one row per item or per scroll tick. No query text is stored; QueryID groups the pages/renders of one query.
type ExposureStage ¶
type ExposureStage string
ExposureStage says how far a result list got: served by the API, rendered by the client, or actually visible to the subject. A click is attributed against the stage the evaluation asks for; an item absent from that stage's list was not exposed there and is never a negative example.
const ( StageServed ExposureStage = "served" StageRendered ExposureStage = "rendered" StageVisible ExposureStage = "visible" )
type HistoryOptions ¶
type HistoryOptions struct {
// ContentKind limits results to one kind. Empty = all kinds.
ContentKind string
Status HistoryStatus
// Since drops rows whose relevant activity is older (e.g. host "clear
// history before X" features). Seen statuses use the last view; HistoryAny
// uses the last signal. Zero = no lower bound.
Since time.Time
Limit int // default 50
Offset int
}
HistoryOptions controls History reads. History is work-level: version rows are read by States with explicit references.
type HistoryStatus ¶
type HistoryStatus string
HistoryStatus filters History results.
const ( // HistoryAny returns every content item the subject has any signal for. HistoryAny HistoryStatus = "" // HistorySeen returns items with MaxProgress > 0. HistorySeen HistoryStatus = "seen" // HistoryInProgress returns seen-but-not-completed items. HistoryInProgress HistoryStatus = "in_progress" // HistoryCompleted returns completed items. HistoryCompleted HistoryStatus = "completed" )
type InventoryRow ¶
type InventoryRow struct {
ContentKind string
SignalType string
Events uint64 // canonical events
RawRows uint64 // stored rows incl. unmerged superseded revisions and retries
Subjects uint64
// ContentItems counts distinct content ids (works and their versions
// count once per id).
ContentItems uint64
FirstAt time.Time
LastAt time.Time
}
InventoryRow sizes one content kind × signal type of a tenant's canonical events: the evidence for deciding which signals to keep collecting.
type LimitError ¶
LimitError reports input that exceeds an ingestion bound.
func (*LimitError) Error ¶
func (e *LimitError) Error() string
type Placement ¶
type Placement struct {
ContentRef
Position uint32
}
Placement is one shown content reference and its absolute 1-based position in the render.
type PopularHit ¶
type PopularHit struct {
ContentRef
ContentMetrics
Score float64
}
PopularHit is one ranked work from Popular. Only works with at least one view in the window rank; version rows never rank.
type PopularOptions ¶
type PopularOptions struct {
Window Window
Limit int // default 20
// Weights tunes the default ranking formula.
Weights RankWeights
// RankExpr, when set, REPLACES the default ranking with a host-supplied
// ClickHouse expression (trusted SQL). It may reference the window metric
// columns: viewers, user_viewers, anon_viewers, views, completions,
// completers, active_s, score_sum, viewer_engagement_sum, returning_viewers,
// events, value_sum, positive_subjects,
// negative_subjects.
RankExpr string
}
PopularOptions controls Popular reads.
type ProjectionKey ¶
type ProjectionKey struct {
ContentKey
Subject Subject
}
ProjectionKey identifies one subject × content reference projection.
type RankWeights ¶
type RankWeights struct {
// PriorWeight is the strength of the Bayesian prior in pseudo-views.
// Defaults to 10.
PriorWeight float64
// PriorScore is the prior mean engagement score. Defaults to 0.
PriorScore float64
// QualityFloor keeps pure-volume ranking meaningful when hosts record no
// scores. Defaults to 1.
QualityFloor float64
}
RankWeights tunes the default popularity ranking:
volume = log10(1 + viewers) quality = (score_sum + PriorScore·PriorWeight) / (views + PriorWeight) rank = volume × max(QualityFloor, quality)
Qualifying views have equal time weight inside the selected window. A score of zero is a valid observation and participates in the mean.
type RefreshCoEngagementOptions ¶
type RefreshCoEngagementOptions struct {
// Window bounds which events feed the rollup (zero = all time).
Window Window
// MaxContentPerSubject caps each subject's contribution to pair
// generation (defaults to 100), bounding the cross-product.
MaxContentPerSubject int
}
RefreshCoEngagementOptions controls RefreshCoEngagement.
type RepairOptions ¶
type RepairOptions struct {
// Window limits candidate keys to those with an event on these days
// (zero = all time).
Window Window
// IngestedSince limits candidates to keys with a row ingested at or after
// it: the crash-repair scan. Zero = no bound.
IngestedSince time.Time
// After resumes strictly after this key in key order (nil = start).
After *ProjectionKey
// Limit caps candidate keys examined per call (default 1000).
Limit int
// Rebuild reprojects every candidate instead of only stale ones: after a
// projection semantics change, a restore or a legacy import.
Rebuild bool
}
RepairOptions bounds one RepairProjections call.
type RepairResult ¶
type RepairResult struct {
Examined int
Repaired int
// Next is the cursor for the following call; nil when the range is done.
Next *ProjectionKey
}
RepairResult reports one bounded repair step.
type SchemaMismatchError ¶
SchemaMismatchError lists every difference between the live signal database and the schema this library version requires.
func (*SchemaMismatchError) Error ¶
func (e *SchemaMismatchError) Error() string
type Scorer ¶
Scorer maps a raw signal (the "session") to a normalized engagement score, generic progress, and whether it counts as "completed". Host-provided per content kind; this is where kinds differ while the hub stays generic.
Examples: gallery — progress = max page reached / page count, completed at ≥90%; blog post — read-time + scroll depth; video — watch %.
type ScorerFunc ¶
ScorerFunc adapts a function to the Scorer interface.
type Signal ¶
type Signal struct {
ContentRef
Subject Subject
// Type is host-defined except TypeView.
Type string
// EventID is the required stable source identity. Retries must reuse it;
// never mint a new id (or time) per delivery attempt.
EventID string
// Revision orders cumulative snapshots of one EventID (for example a
// checkpointed consumption session, or a subject's current preference):
// the highest revision is the event, lower ones are superseded. Leave 0 for
// immutable events. Conflicting content at an equal revision resolves by a
// deterministic content hash, not by arrival order.
Revision uint64
// OccurredAt is the required immutable source time (a session's start). It
// selects the UTC day the event counts in.
OccurredAt time.Time
// Consumption measurements, cumulative for the revision.
DurationS uint32 // active time
Progress uint32 // numerator (pages / scroll % / watched s)
ProgressMax uint32 // denominator (page count / 100 / duration)
// Value is explicit feedback (+1 like, -1 dislike, rating). State and
// windows sum the canonical values.
Value float64
// Score is the engagement score; a registered Scorer fills it together with
// Progress/ProgressMax/Completed.
Score int16
// Completed per the content kind's completion rule.
Completed bool
// Resume is an opaque host pointer for "pick up where you left off".
Resume string
// Payload holds bounded context, JSON-encoded at rest.
Payload map[string]any
}
Signal is one logical source event: a consumption session, reaction, click, rating. Identity is (tenant, content ref, subject, Type, EventID): re-delivering the same identity never adds another event, whatever its arrival order, batch or merge state. Never emit one event per scroll/frame tick. A version-scoped reference records that version's consumption; the work-level reference records the work once across versions.
func (Signal) Attribution ¶
func (s Signal) Attribution() Attribution
Attribution reads attribution back from the signal's payload, tolerant of the numeric type a JSON round-trip produces. Missing keys yield zero values.
func (Signal) WithAttribution ¶
func (s Signal) WithAttribution(a Attribution) Signal
WithAttribution returns a copy of the signal with attribution written into a fresh Payload under the standardized keys (existing payload entries are preserved). Zero-valued fields are omitted.
type State ¶
type State struct {
Seen bool // MaxProgress > 0
FirstSeenAt time.Time
LastSignalAt time.Time
LastViewAt time.Time
TotalEvents uint32 // canonical events of every type
Views uint32 // canonical view events (sessions)
Completions uint32 // completed views
ActiveS uint64 // summed view DurationS
MaxProgress uint32 // progress bar = MaxProgress / ProgressMax
ProgressMax uint32
Completed bool
Resume string // latest non-empty resume pointer
LastScore int16 // score of the latest view
// NetValue sums canonical feedback Values. Negative = current sentiment is
// negative; recommendations exclude such content.
NetValue float64
Feedback uint32 // canonical events with a non-zero Value
}
State is one subject's compact, indefinitely retained standing with one content reference, derived from its canonical events.
type StateRow ¶
type StateRow struct {
ContentRef
State
}
StateRow is a State with its content reference, as returned by History.
type Store ¶
type Store struct {
// contains filtered or unexported fields
}
Store reads and writes the signal plane for one ClickHouse database. Every method is tenant-scoped via its tenant argument; every content reference it receives must belong to that tenant.
func NewStore ¶
NewStore returns a Store over an existing ClickHouse connection. The schema comes from migrations.SignalClickHouse; verify it with CheckSchema.
func (*Store) Attribution ¶
func (st *Store) Attribution(ctx context.Context, tenant string, opts AttributionOptions) (AttributionPage, error)
Attribution joins canonical clicks to canonical exposures at one stage, so evaluation sees what was actually exposed, what was clicked, and which clicks have no exposure to learn from.
The export is one deterministic sequence: renders in render_id order, each with its clicks in (occurred_at, content, subject, event_id) order, followed by the unattributed clicks in (render_id, occurred_at, content, subject, event_id) order. Every page is bounded by Limit renders and ClickLimit click rows; Next resumes exactly after the last row emitted, inside a render when its clicks did not fit. Concatenating all pages at any size yields the same rows once each. Erased subjects are excluded at read time.
func (*Store) CoEngaged ¶
func (st *Store) CoEngaged(ctx context.Context, tenant string, ref ContentRef, opts CoEngagedOptions) ([]CoEngagedHit, error)
CoEngaged returns works co-engaged with the anchor work: "subjects who engaged with X also engaged with Y". Strength nets subjects whose summed feedback for the candidate is non-negative against those with negative.
func (*Store) EnforceErasures ¶
EnforceErasures re-applies every recorded erasure of a tenant to every row: it deletes residue left by writers or projections that observed the state before the fence, or rows restored from a backup. Such rows were already unreadable; this removes them physically. There is no cursor, so a delayed write for an old erasure is never skipped. Idempotent; schedule it, and run it after every restore (after re-erasing subjects deleted since the backup, which the host's own deletion ledger knows).
func (*Store) EraseSubjects ¶
func (st *Store) EraseSubjects(ctx context.Context, tenants []string, subjects []Subject) (ErasureReport, error)
EraseSubjects permanently erases subjects from every listed tenant.
Completion contract: when the returned report is Complete(),
- the erasure is recorded in the ledger on every replica (quorum insert), so from that moment every write, read, projection and export on any replica drops or hides the subjects' rows (the barrier): the subject key is dead in the tenant forever;
- every row of the subjects that existed when the pass ran is deleted from signals, compact state, daily contributions, exposures and legacy raw tables on every replica (mutations_sync = 2) and re-counted as zero;
- co-engagement pairs touching works the subjects contributed to are removed (RefreshCoEngagement rebuilds them and re-verifies the ledger around its build).
A writer or projection that observed the pre-erasure state can still leave residue rows; they are unreadable by (1) and EnforceErasures removes them physically. Idempotent. Returns an error without Complete() when a replica is unavailable (the fence would not be durable everywhere) or rows remain after maxErasurePasses. After restoring a backup, replay every erasure newer than the backup and run EnforceErasures.
func (*Store) ErasedSubjects ¶
func (st *Store) ErasedSubjects(ctx context.Context, tenant string, subjects []Subject) (map[Subject]bool, error)
ErasedSubjects reports which of the given subjects a recorded erasure of tenant covers, so a producer can retire their obligations instead of retrying deliveries the ledger will drop.
func (*Store) Forget ¶
func (st *Store) Forget(ctx context.Context, tenant string, subject Subject, contentKind string, contentID string) error
Forget erases a subject's events and projections for one content item (the work and all its versions) or an entire content kind when contentID is empty. Windows and popularity read the per-subject projections, so they stop counting the subject at once. This is NOT account erasure: exposures and the content_pairs rollup remain, and a write racing the deletion can re-project.
func (*Store) ForgetExposures ¶
ForgetExposures removes all result-list exposures of one tenant and subject (host "clear my search history"). Not erasure: see EraseSubjects.
func (*Store) History ¶
func (st *Store) History(ctx context.Context, tenant string, subject Subject, opts HistoryOptions) ([]StateRow, error)
History returns the subject's work-level state rows, most recent relevant activity first. Seen statuses use view recency; HistoryAny uses signal recency.
func (*Store) HistoryCount ¶
func (st *Store) HistoryCount(ctx context.Context, tenant string, subject Subject, opts HistoryOptions) (int64, error)
HistoryCount returns the total row count History would paginate over.
func (*Store) Inventory ¶
Inventory reports canonical and raw event volume per content kind and signal type.
func (*Store) Metrics ¶
func (st *Store) Metrics(ctx context.Context, tenant string, refs []ContentRef, window Window) (map[ContentKey]ContentMetrics, error)
Metrics returns named window metrics for the references (works or versions; zero window = all time). References with no canonical events in the window are absent. A work's metrics never include its version rows.
func (*Store) NegativeIDs ¶
func (st *Store) NegativeIDs(ctx context.Context, tenant string, subject Subject, contentKinds []string) (map[ContentKey]struct{}, error)
NegativeIDs returns the works the subject has net-negative canonical feedback for: the exclusion set for recommendations.
func (*Store) Popular ¶
func (st *Store) Popular(ctx context.Context, tenant string, contentKind string, opts PopularOptions) ([]PopularHit, error)
Popular ranks works of one kind with at least one view in a literal window. Every view in the window has equal time weight.
func (*Store) PopularityFor ¶
func (st *Store) PopularityFor(ctx context.Context, tenant string, contentKind string, ids []string, window Window) (map[string]float64, error)
PopularityFor scores a fixed candidate set of works by the popularity ranking and returns content_id -> score. Candidates without views in the window are absent.
func (*Store) PurgeContentKinds ¶
PurgeContentKinds deletes a tenant's events, projections and co-engagement pairs for whole content kinds (for example retired fan-out rows) and waits for the mutations on every replica. Irreversible; inventory first.
func (*Store) RecordExposures ¶
RecordExposures appends one row per result list and stage, batched into a single INSERT. Exposures of erased subjects are dropped. Empty input is a no-op.
func (*Store) RecordSignals ¶
RecordSignals appends source events, then rebuilds the compact state and daily projections of every touched subject × content reference from its canonical events. Replays, reordering and revisions converge; nothing is incremented. Signals of erased subjects (EraseSubjects) are dropped. If the process stops between the insert and the projections, the events are durable and RepairProjections restores the projections.
func (*Store) RefreshCoEngagement ¶
func (st *Store) RefreshCoEngagement(ctx context.Context, tenant string, opts RefreshCoEngagementOptions) error
RefreshCoEngagement (re)materializes the content_pairs rollup for one tenant: per subject, the distinct works with non-negative summed feedback (capped at MaxContentPerSubject) are cross-joined into pairs; strength = co-engaged subjects minus subjects with negative feedback on the candidate. Pairs carry no subject, so a build that started before an erasure could publish that subject's contribution after EraseSubjects removed its pairs: the build is repeated while the erasure ledger changed during it.
func (*Store) RepairProjections ¶
func (st *Store) RepairProjections(ctx context.Context, tenant string, opts RepairOptions) (RepairResult, error)
RepairProjections is the owned projection repair. It examines up to Limit candidate keys in key order and rebuilds those whose state or daily projection is missing or older than their newest raw event (all of them with Rebuild). Idempotent; call again with After = Next until Next is nil.
func (*Store) SeenIDs ¶
func (st *Store) SeenIDs(ctx context.Context, tenant string, subject Subject, contentKind string) (map[string]struct{}, error)
SeenIDs returns the subject's seen-set for one content kind: work ids with max_progress > 0. This is the signal-plane half of the unseen anti-join.
func (*Store) States ¶
func (st *Store) States(ctx context.Context, tenant string, subject Subject, refs []ContentRef) (map[ContentKey]State, error)
States is the bulk "annotate this list with view context" read: for each requested reference (work or version), the subject's standing (seen?, progress bar, completed?, resume pointer). References the subject has no signals for are absent.
type Subject ¶
Subject is who acted: a resolved user id (e.g. from authkit) or, for anonymous traffic, a session-key hash. Exactly one of the two must be set.
History/unseen are meaningful only for logged-in subjects; anonymous signals still feed popularity/engagement aggregates.
type TopStatesOptions ¶
type TopStatesOptions struct {
ContentKinds []string
// ExcludeNegative drops content the subject has net-negative explicit
// feedback for (a disliked item must not seed recommendations).
ExcludeNegative bool
Limit int // default 10
}
TopStatesOptions controls TopStates (recommendation seeds).
type Window ¶
type Window struct {
From time.Time // inclusive UTC midnight
To time.Time // exclusive UTC midnight
}
Window is a literal range of whole UTC calendar days, [From, To). A zero From or To leaves that side unbounded; the zero Window is all time. Every qualifying event inside the window counts with equal weight: there is no age decay, and an event's contribution does not depend on where in the window it falls.
func LastDays ¶
LastDays returns the n most recent UTC calendar days including the day that contains now: [midnight(now) - (n-1) days, midnight(now) + 1 day). The window moves only at UTC midnight, so String is a stable cache key within a day. n must be positive (Validate reports otherwise).