signal

package
v1.1.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 20, 2026 License: MIT Imports: 11 Imported by: 0

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

View Source
const (
	DefaultAttributionRenders = 500
	DefaultAttributionClicks  = 5000
)

Default page bounds of Attribution.

View Source
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.

View Source
const (
	SubjectKindUser = "user"
	SubjectKindAnon = "anon"
)

SubjectKind values stored in the subject_kind column.

View Source
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.

View Source
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.

View Source
const MaxErasureSubjects = 100

MaxErasureSubjects bounds one EraseSubjects call.

View Source
const TypeClick = "click"

TypeClick is the selection signal type Attribution joins to exposures.

View Source
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

func CheckSchema(ctx context.Context, conn Conn, database string) error

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

func CreateDatabase(ctx context.Context, conn Exec, database, cluster string) error

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

type Exec interface {
	Exec(ctx context.Context, query string, args ...any) error
}

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

type LimitError struct {
	Field string
	Limit int
	Got   int
}

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

type SchemaMismatchError struct {
	Database string
	Problems []string
}

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 Scored

type Scored struct {
	Score       int16
	Progress    uint32
	ProgressMax uint32
	Completed   bool
}

Scored is the result of a content kind's Scorer.

type Scorer

type Scorer interface {
	Score(ctx context.Context, s Signal) (Scored, error)
}

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

type ScorerFunc func(ctx context.Context, s Signal) (Scored, error)

ScorerFunc adapts a function to the Scorer interface.

func (ScorerFunc) Score

func (f ScorerFunc) Score(ctx context.Context, s Signal) (Scored, error)

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

func NewStore(conn Conn, database string) (*Store, error)

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

func (st *Store) EnforceErasures(ctx context.Context, tenant string) (ErasureReport, error)

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(),

  1. 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;
  2. 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;
  3. 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

func (st *Store) ForgetExposures(ctx context.Context, tenant string, subject Subject) error

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

func (st *Store) Inventory(ctx context.Context, tenant string) ([]InventoryRow, error)

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

func (st *Store) PurgeContentKinds(ctx context.Context, tenant string, contentKinds []string) error

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

func (st *Store) RecordExposures(ctx context.Context, tenant string, exposures []Exposure) error

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

func (st *Store) RecordSignals(ctx context.Context, tenant string, signals []Signal) error

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.

func (*Store) TopStates

func (st *Store) TopStates(ctx context.Context, tenant string, subject Subject, opts TopStatesOptions) ([]StateRow, error)

TopStates returns the subject's highest-signal works (recommendation seeds): ordered by last_score DESC, then recency.

type Subject

type Subject struct {
	UserID  string
	AnonKey string
}

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.

func (Subject) Key

func (s Subject) Key() string

Key returns the stored subject identifier.

func (Subject) Kind

func (s Subject) Kind() string

Kind returns "user" or "anon".

func (Subject) Validate

func (s Subject) Validate() error

Validate checks that exactly one of UserID / AnonKey is set.

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 AllTime

func AllTime() Window

AllTime returns the unbounded window.

func Between

func Between(from, to time.Time) Window

Between returns the whole UTC days [from, to); both must be UTC midnights.

func LastDays

func LastDays(n int, now time.Time) Window

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).

func (Window) String

func (w Window) String() string

String renders the window as "[from,to)" dates ("*" when unbounded).

func (Window) Validate

func (w Window) Validate() error

Validate reports a window that is not a non-empty range of whole UTC days.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL