contentkit

package module
v0.62.0 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 22 Imported by: 0

README

contentkit

contentkit is the deterministic content library for the Doujins, Hentai0 and marketplace hosts: tenant-scoped interactions (posts, comments, reactions, favorites, polls), keyword search over host content, the ClickHouse signal plane (consumption, feedback, exposures, popularity, erasure) and the discovery reads over both. It needs no model provider, API key or vector extension. Semantic search belongs entirely to the separate, deferred User Intelligence library; ContentKit has no semantic search configuration or runtime hook.

Design: open-rails-tracker/contentkit/DESIGN.md. Host contract: HOST_INTEGRATION.md. Migrations: docs/migration.md, docs/restore.md.

Release policy

ContentKit is pre-stable and releases on v0.x; its API may change between minor versions. The historical v1.0.0-rc.1 through v1.1.1 tags were premature and are retracted. v1.1.2 is a self-retracted, withdrawal-only marker carrying Go module metadata, not a supported stable API release.

The marker and v0.12.2 identify the same source commit. Go reads retractions from the highest release before applying them, so the marker makes fresh @latest requests select the latest unretracted v0.x release, including future v0.x releases. Existing tags remain immutable; retraction preserves explicit pins and does not automatically downgrade existing consumers. See Go module retractions.

Vocabulary

Name Meaning
tenant_id one site: doujins, hentai0, the marketplace
content_id the host-owned work: a gallery, a video, a listing; a canonical UUIDv7, never reused (Content ids)
content_version_id one selectable version of that work (optional)
taxonomy_id a generic ContentKit record: tag, artist, series, creator, character, voice actor
ContentRef {TenantID, ContentKind, ContentID, ContentVersionID *string} — the typed reference every API, key, index and cursor carries

A language is metadata on a document or version, never part of a reference. Every read is scoped to the tenant pinned at construction; a reference of another tenant is an error, never remapped.

Packages

Package Owns
contentref ContentRef, ContentKey, TaxonomyID
access Actor, the batch ContentResolver port (Resolve(ctx, refs, actor) → map[ContentKey]Resolution; an omitted ref denies) and its Resolution{Ref, Visible, Accessible, Editor}, shared by content and media
media the registry (Config, kinds, upload paths, private and public presets), ordered manifests with provenance and conditional-write edits, the Store port, direct uploads and commit ops with their HTTP API, reads and HLS playlists, the optional UploadLimiter, and the sweep, deletion, Expose and relays as River jobs
media/s3 Store over aws-sdk-go-v2 (Ceph RGW in production, MinIO in tests), bucket policy and point-in-time Restore
media/image libvips (CGO) producer: Image presets, public presets, zips, editor views, PublishDefaults
media/token media access tokens, shared by hosts and the access agent
media/layout object keys and the access agent's host and default rules, dependency-free
media/agent the access agent's handler (cmd/media-access)
media/video ffmpeg producers: byte-range fMP4 HLS ladders with audio, subtitle and sprite tracks, MP4 per rung, audio, subtitles, frame grabs; Frames for the frame picker
media/worker the media worker: one process for every producer, built from the host's registry (cmd/media-worker is the stock build)
media/workqueue the host's side of the worker: its per-host River schema, insert-only Queue (enqueue, cancel), encode progress
media/tiered optional public/members/ppv/members_ppv/premium policy over an entitlement Checker (hosts adapt OpenRails CheckEntitlements)
content posts, comments, reactions, favorites, polls (multiple-choice and free-text) and their counts over ContentRef, in the host schema's content_* interaction tables; the Identity/Authorizer/UserEnricher/ContentProcessor ports, post and poll images through Media, the optional ContentModerator (held/review queue) and AnswerClassifier ports, and the HTTP routes
search PGroonga keyword search (exact/alias/prefix/typo, EN/ZH/JA/KO), documents and dirty queue, RRF, the DocumentSink port
worker one tenant's document maintenance: dirty queue, bounded backfill, sink delivery
taxonomy generic catalog: nodes (tags, artists, creators, characters, series, seasons, voice actors), localized names/aliases, edges, content assignments, effective tags, per-language counts, typeahead documents, admin routes
signal ClickHouse signal plane: canonical signals, compact subject state, daily rollups, windows, erasure fence, exposures/attribution, repair
popularity named ranking policy (PolicyV1) over the window metrics: ClickHouse RankExpr and Go Score in agreement, literal windows, session scorer, taxonomy popularity through the host Catalog port
discovery SimilarTo/Recommend: the Candidates port, the default co-engagement source (Engagement), Fallback, and the shared exclusion/fill policy (Recommender)
eval lexical golden-case evaluation, reports, baselines
migrations PostgreSQL migration chain and ClickHouse baseline
adapters/authkit its own module (opt-in): account avatars and content authors from AuthKit; see HOST_INTEGRATION "Account avatars"
root Runtime (one constructor: hub + content + HTTP mount), Migrate (all PostgreSQL features and optional ClickHouse signals), Client (keyword search + typeahead), EmbeddedHub (signal + discovery)

Install

One call installs all PostgreSQL features in a host-selected schema (which may also hold application tables) and the signal plane in ClickHouse. PGroonga and pg_trgm live in public; no vector extension is required:

_ = signal.CreateDatabase(ctx, adminCH, "hub", cluster)
_ = contentkit.Migrate(ctx, contentkit.MigrateConfig{
	DB: sqlDB, Schema: "doujins",
	ClickHouse: &chmigrate.Config{ClientAddr: addr, Database: "hub", App: "contentkit_signal", Cluster: cluster},
})

The baselines initialize fresh stores; newer PostgreSQL migrations upgrade the current lineage, not retired migration chains. See docs/migration.md.

Runtime

rt, _ := contentkit.NewRuntime(ctx, contentkit.RuntimeConfig{
	EmbeddedConfig: contentkit.EmbeddedConfig{PG: pool, PGSchema: "doujins", Tenant: "doujins", CH: ch, CHDatabase: "hub"},
	Content: content.Options{Schema: "doujins", Identity: identity, Authz: authz, Resolver: resolver, ContentKinds: []string{"gallery", "post"}},
})
mux.Handle("/api/social/", http.StripPrefix("/api/social", rt.Handler()))
counts, _ := rt.Content.Counts(ctx, []contentkit.ContentRef{rt.Content.Ref("gallery", "42")})
_ = worker.SyncOnce(ctx, rt.WorkerOptions(hostWorkerOptions)) // host documents + posts

rt is the Hub (search, typeahead, signals, discovery) plus rt.Content (interactions). See HOST_INTEGRATION.md.

Documents are one per content reference and language: a gallery with an English original, an English colored edition and a Spanish original is three documents keyed by version. Hosts mark changes in the transaction that changes content and run one worker tick per schedule:

_ = search.MarkDirty(ctx, tx, schema, []search.DirtyMark{{DocumentKey: search.DocumentKey{ContentRef: ref, Language: "en"}}})

_ = worker.SyncOnce(ctx, worker.Options{
	Pool: pool, Schema: schema, Tenant: "doujins",
	SupportedLanguages: []string{"en", "es"}, ContentKinds: []string{"gallery"},
	ListContent:           listGalleries,          // bounded pages of refs, for backfill
	BuildKeywordDocuments: buildGalleryDocuments,  // refs -> KeywordDocument{Title, Aliases, Keywords}
	Sink:                  nil,                    // optional DocumentSink
})

The worker records rejected host documents in content_search_invalid with their dirty-queue revision, validation error and failure time. It continues with other documents and backfill, but does not acknowledge the rejected row. Inspect that table for repair; marking the document dirty after correcting its source advances the revision and retries it. A deletion clears the invalid record. Run the PostgreSQL migrations before starting a worker built against this schema.

Querying groups documents per work before paging; the host's eligibility join decides, per document, whether that one row is visible and which is preferred:

client, _ := contentkit.NewClient(contentkit.ClientConfig{Pool: pool, Schema: schema, Tenant: "doujins"})
page, _ := client.Search(ctx, query, contentkit.SearchOptions{
	Language: "es", ContentKinds: []string{"gallery"}, Limit: 20,
	Eligibility: &contentkit.Eligibility{SQL: eligibilitySQL, Args: args},
})
for _, hit := range page.Hits { /* hit.ContentID (work), hit.Version(), hit.Language, hit.Score */ }

SearchHit.Score is the keyword match tier (exact title 1, alias 0.9, token/prefix 0.75–0.77, typo 0.5–0.52). SearchWithTrace returns retrieval provenance for evaluation. Query limits: 256 characters, 16 tokens, a two-second ceiling per request.

Ports

Port Called when Contract
DocumentSink the worker publishes or deletes a document Upsert(PublishedDocument), Delete(DocumentKey, Version); at-least-once, atomic newer-version wins across both operations, with a retained deletion tombstone; a failing sink keeps the row queued and never blocks the keyword index

DocumentSink is a neutral document change feed for external indexes, caches or audit consumers. ContentKit ships no sink implementation or AI-specific configuration.

Errors

Both mounted handlers (content.Runtime.Handler, taxonomy.Handler) answer failures with one flat body. Branch on code; error is a human message and may change.

{"error":"not found","code":"not_found"}
Status Code Meaning
400 invalid_request malformed or semantically invalid input
401 unauthorized no identity
403 forbidden identity present, not permitted
404 not_found absent, unpublished or soft-deleted (existence is hidden)
409 conflict state or revision conflict
422 moderation_rejected a ContentModerator refused the write; error is the author-facing reason
501 not_configured the host never wired the port this route needs (Media, AnswerClassifier)
500 tenant_mismatch a host port answered with another tenant's data
500 internal_error anything else

5xx bodies carry no cause: it goes to Options.Logger (slog.Default() when unset) with the request method, path, status and duration. Postgres constraint names, driver text and stack traces are logged, never served.

Media

Media is an app-defined, self-describing file system in one private bucket; no database table records what media exists. The app declares its kinds in one registry (media.Config); HOST_INTEGRATION "Media" covers the registry, commit ops, reads and exposure.

{namespace}/{kind}/{id}/manifest.json         gzip JSON: the ordered file list with provenance; never served
                       /private/sha256-{hex}  every blob: uploads, derived files, editor views (token)
                       /public/{name}         app-declared names, e.g. cover-460.webp (anyone)
                       /temp/{name}           staged uploads (u-{uuid}) and in-flight writes; never served
{namespace}/{kind}/_default/public/{name}     a public preset's default image

An item is the host's version (gallery/456 English, gallery/789 Korean). Private blobs are content-addressed and immutable; public names are fixed and overwritten in place (ETags and a CDN purge keep them fresh). Every URL is https://media.<site>/v1/{namespace}/{kind}/{id}/{public|private}/{name}.

reg, _ := media.NewRegistry(media.Config{Namespace: "doujins", BaseURL: "https://media.doujins.ai",
	Kinds: []media.Kind{Gallery, accountmedia.User},
	Hooks: media.Hooks{Resolver: resolver, CanUpload: authorizer, PurgePublic: purge, ItemReady: ready}})
store, _ := s3.New(s3.Config{Bucket: "media", Endpoint: rgw, PublicEndpoint: "https://s3.doujins.ai", UsePathStyle: true,
	AccessKeyID: id, SecretAccessKey: secret})
key, _ := token.ParseKey(os.Getenv("MEDIA_TOKEN_KEY")) // "{kid}:{base64}", shared with media-access
_ = workqueue.Migrate(ctx, pool, "doujins_media_worker") // this host's worker schema
queue, _ := workqueue.New(pool, reg, "doujins_media_worker")
jobs, _ := media.NewJobs(media.JobsConfig{Store: store, Registry: reg, Locker: media.PGLocker(pool), Pool: pool,
	Processes: queue, Limiter: limiter})
uploads, _ := media.NewUploads(media.UploadOptions{Store: store, Manifests: jobs.Manifests(), Tickets: &ring,
	Limiter: limiter, Queue: queue, Frames: frames})
reader, _ := media.NewReader(media.ReaderOptions{Manifests: jobs.Manifests(), Queue: queue, Progress: progress,
	Delivery: media.Delivery{Mode: media.DeliverCookie, CookieDomain: "doujins.ai", SigningKey: key}})
mux.Handle("/api/media/upload/", http.StripPrefix("/api/media/upload", media.UploadHandler(uploads, media.UploadHandlerOptions{Actor: actorOf})))
mux.Handle("/api/media/", http.StripPrefix("/api/media", reader.Handler(media.HandlerOptions{Identity: identity})))
// composed into the host's River client: jobs.RiverJobs()
_ = jobs.ExposeTx(ctx, tx, ref)                                          // whenever anonymous visibility changes
_ = jobs.DeleteItemsTx(ctx, tx, media.Deletion{Ref: ref, Owner: owner}) // in the host's delete transaction

Manifests. Every edit runs under the Locker (required: PGLocker, a Postgres advisory lock shared by every process on the bucket) and writes with If-Match (or If-None-Match: *) once the store reports conditional PUT, retrying on conflict. The edit is normalized (canonical order, download names, dropped originals) and validated; an unchanged manifest is not written. Reads go through an in-process cache bounded by bytes and revalidated by ETag, so a read is never stale. A 2,000-page gallery is about 1.5 MB of JSON and 400 KB stored; a 2-hour video about 5 KB, its segment tables living in index blobs. Every read and write stops at MaxManifestBytes (8 MiB of JSON; the default 128 MiB cache holds five). Edits stop 4 KiB short of it and only when they grow the manifest, so a full item can still shrink and the hidden and full flags always fit. A commit is refused (413 too_large) when the item, once processed, would pass that: it projects every output its presets have yet to write. When outputs still overrun it, the worker marks the item full instead of recording them: producers then write no private output for it (its public files still render), and editor reads report full (state full), until a commit frees the bytes the refused record was short of (deficit, counting what removed uploads would still have added). An item holds at most MaxUploads (10,000) uploads, an upload's meta at most 4 KiB of JSON and the item's 16 KiB; names are capped at 200 bytes of JSON, written unescaped (< is one byte). The manifest object records the quota its uploads were charged, so deleting an item never needs it to decode, and hiding deletes public/ before touching it.

The bucket is optional at startup. s3.New never dials; register store.Check(ctx, prefix) as the host's optional S3 dependency probe. Its first success probes the backend's capabilities unless Config.Capabilities declares them; a throttled or cut-off probe records nothing. An unreachable or 5xx bucket is media.ErrUnavailable: 503 unavailable over HTTP, and in media jobs a River snooze (not an attempt) while a fresh Check confirms the outage, capped by MaxOutageSnoozes (media.SnoozeUnavailable).

Uploads go straight to the bucket, to a staged name: up to 64 MiB is one PUT to temp/u-{uuid} signed with its type, length and SHA-256, larger files are multipart (8–16 MiB parts, each signed with its length and SHA-256, resumed through ListParts, completed by the server from a signed ticket; nothing is stored). The browser's hash only finds an identical blob already in the folder (no upload). After the commit, the worker's place job hashes each staged upload and writes it to private/sha256-{hex} of the bytes it read (Manifests.Place); an existing blob is never overwritten, so a blob always holds the bytes of its name. The optional UploadLimiter (media.NewPGLimiter) rate-limits uploaders and charges each upload's size (at least 64 KiB) to the grant's Owner; growth past the quota fails with 413 quota_exceeded, and deleting an item releases it.

Access agent (cmd/media-access, media/agent, image ghcr.io/open-rails/contentkit-media-access:{tag}): public/ to anyone (public, max-age=300, stale-while-revalidate=86400, falling back to the kind's _default for declared names), private/sha256-{hex} with the item's token in ?t= or the mt cookie (private, immutable). A token opens every private file of its item or none; ?dl={name} serves the file as a download under that name when it has the file's type. Everything else and every denial is one no-store 404. It keeps no state: rate limit the media host at the ingress; the read API limits the items a viewer opens per hour (ReaderOptions.Issuance). It needs MEDIA_ACCESS_S3_ENDPOINT, _S3_BUCKET, a key reading only */private/* and */public/*, MEDIA_ACCESS_TOKEN_KEY and _TOKEN_KEY_PREVIOUS, MEDIA_ACCESS_HOSTS (media.doujins.ai=doujins,accounts; media.hanime.media=hentai0,accounts), MEDIA_ACCESS_CORS_ORIGINS and MEDIA_ACCESS_DEFAULTS (layout.FormatDefaults(media.AgentConfig(reg).Defaults)).

The media worker (media/worker) runs every producer from the host's worker River schema (media/workqueue, MEDIA_WORKER_SCHEMA, per host: hosts sharing a database never share one): images, zips, public presets and editor views (media/image, libvips) and HLS, MP4, audio, subtitles and frames (media/video, ffmpeg). The host presigns, commits, exposes and reads, and links only media/workqueue. The worker hands readiness (Hooks.ItemReady), purges (Hooks.PurgePublic) and sweeps back to the host's River schema (MEDIA_HOST_RIVER_SCHEMA), so the stock build (cmd/media-worker, image ghcr.io/open-rails/contentkit-media-worker) needs only the registry as JSON (MEDIA_KINDS_FILE). Hosts with a Private.Choose build their own:

cfg, _ := worker.FromEnv(ctx) // DATABASE_URL, MEDIA_S3_*, MEDIA_WORKER_SCHEMA, MEDIA_HOST_RIVER_SCHEMA, MEDIA_WORKER_*
cfg.Kinds = reg                // the host's registry
w, _ := worker.New(ctx, cfg)   // no DDL: the host's migrate step runs workqueue.Migrate
_ = w.Run(ctx)                 // until SIGTERM; running jobs get MEDIA_WORKER_SHUTDOWN_GRACE

It never exits for a missing dependency: Run takes no jobs until the bucket answers, and MEDIA_METRICS_ADDR serves /livez, /readyz, /statusz and app_dependency_up with /metrics.

Images (media/image, CGO over libvips): WebP at each Image spec (inside or cover box, quality, blur), through the upload's edit (crop in EXIF-oriented source pixels, then a clockwise quarter rotation); nothing is upscaled. GIF and WebP animations keep every frame, delay and loop count; MaxPixels (100 MP over all frames), MaxFrames (1000) and MaxAnimationSeconds (60) bound sources, and a stored source over its Upload's MaxBytes is refused unread. Uploads feeding image presets take only media.ImageTypes (JPEG, PNG, WebP, GIF, AVIF, HEIF, TIFF), the declared type binds the decoder, and every other libvips loader is blocked (libvips 8.13 or later). Public presets render every width through the same spec (a width past the edited image at its width) and carry from and fp as object metadata.

Video (media/video): an HLS preset encodes the ladder (rung = short side, default 2160/1080/480, none above the source) in each of the worker's codecs (MEDIA_WORKER_CODECS, default av1,h264), plus AAC per audio track, WebVTT per text subtitle and a seek sprite. Each rendition is one byte-range fMP4 blob whose segment table is its own index blob; stages are published rung by rung (the upload stays pending until the last), a compliant top rung is stream-copied, and sources outside the aspect bounds fail. An MP4 preset muxes H.264 at one rung with the default audio; Audio gives an HLS track and an M4A (optional EBU R128 loudness); Subtitles converts SRT and SSA/ASS to clean WebVTT. Playlists are built per request by the read API. Encode progress comes from workqueue.NewProgressSource (ReaderOptions.Progress).

Jobs (media.Jobs, composed into the host's River client): the sweep collects garbage by manifest reference (unreferenced private blobs a grace period old, unexpected public names at once, temp/ by age); folder deletion (with a late-upload second pass and quota release); Expose; Regenerate; SweepOrphans; and the worker's relays.

Tiered access (media/tiered) maps public/members/ppv/ members_ppv/premium to Resolutions over an entitlement Checker.

Browser SDK (sdk/upload, @openrails/contentkit-upload, attached to each release): hashing, uploads, commit ops, reads, waitFor, the frame picker, publicURL/srcSet, React hooks and UI. See sdk/upload/README.md.

Taxonomy

Nodes, names, edges and assignments are tenant-scoped; effective tags are the work's assignments ∪ the selected version's; RequireAll makes a multi-node filter hold on one eligible version inside the same join as search. The PostgreSQL baseline always installs the taxonomy tables; see docs/taxonomy-migration.md.

store, _ := taxonomy.New(taxonomy.Options{Pool: pool, Schema: schema, Tenant: "doujins",
	Kinds: []string{"tag", "artist", "character", "series", "voice_actor"}, Languages: []string{"en", "es"},
	CountEligibility: &search.Eligibility{SQL: releasedVersionSQL}})
_ = store.WithTx(tx).Assign(ctx, []taxonomy.Assignment{{ContentRef: g1.WithVersion(v2), TaxonomyID: "colored"}}, taxonomy.AssignOptions{})
tags, _ := store.EffectiveTags(ctx, []contentkit.ContentRef{g1.WithVersion(v2)})
filter, args, _ := taxonomy.RequireAll(schema, []taxonomy.TaxonomyID{"colored"})
page, _ := client.Search(ctx, q, contentkit.SearchOptions{Language: "es", ContentKinds: []string{"gallery"}, FilterSQL: filter, FilterArgs: args, Eligibility: elig})
mux.Handle("/admin/taxonomy/", http.StripPrefix("/admin/taxonomy", taxonomy.Handler(store)))

ListNodes backs a catalog index page directly: the display name in the request language (falling back to the store's configured order, or pinned with LanguageMode), the per-language content count, an A-Z index, a name+alias search, hide-empty, five orders and offset paging with a total. The zero ListOptions keeps the keyset contract a full admin sync wants.

page, _ := store.ListNodes(ctx, taxonomy.ListOptions{Kind: "artist", Language: "es",
	ContentKind: "gallery", MinCount: 1, Sort: taxonomy.SortCount, Offset: 40, Limit: 20})
// page.Total, and per row: Name, NameLanguage, Count.
letter, _ := store.ListNodes(ctx, taxonomy.ListOptions{Kind: "artist", Language: "es", NamePrefix: "a"})
cast, _ := store.ListNodes(ctx, taxonomy.ListOptions{Kind: "character", Related: "s-fate", Relation: taxonomy.RelationMemberOf})

Worker: ContentKinds: append(hostKinds, store.Kinds()...), ListContent: store.Lister(listGalleries), BuildKeywordDocuments: store.Builder(buildGalleryDocuments).

Signal plane and discovery

hub, _ := contentkit.NewEmbedded(contentkit.EmbeddedConfig{
	PG: pool, PGSchema: schema, CH: ch, CHDatabase: "hub", Tenant: "doujins",
	Scorers:  map[string]signal.Scorer{"gallery": galleryScorer},
	Catalogs: map[string]contentkit.ContentCatalog{"gallery": galleryCatalog},
})
g1 := hub.Content("gallery", "g1")
_ = hub.RecordSignals(ctx, []signal.Signal{{ContentRef: g1, Subject: user, Type: signal.TypeView, EventID: sessionID, Revision: checkpoint, OccurredAt: start, Progress: 95, ProgressMax: 100}})
states, _ := hub.States(ctx, user, []contentkit.ContentRef{g1, g1.WithVersion("v2")})
top, _ := hub.Popular(ctx, "gallery", signal.PopularOptions{Window: signal.LastDays(30, time.Now())})
  • Identity is (tenant, content ref, subject, type, event id); retries, reordering and revisions converge, nothing is incremented. Record the work view and, separately, the selected version's view (g1.WithVersion(v)).
  • Work reads (History, Popular, SeenIDs, co-engagement) count the work once across versions; States and Metrics read exactly the references given, work or version.
  • Windows are whole UTC days, 7/30/90/365/all, no decay.
  • Rank by a named policy, not the default rank: popularity.ByName("v1"), popularity.New(popularity.Config{Source: hub, Policy: policy}), then ranker.Popular / ranker.Scores / ranker.Taxonomy (docs/popularity-policy.md).
  • SimilarTo and Recommend draw candidates from EmbeddedConfig.Candidates (default: co-engagement); see discovery candidates.
  • EraseSubjects is account erasure with a quorum-written fence; see HOST_INTEGRATION.md and, for interaction data, interaction erasure.
  • Reactions and favorites reach the signal plane from their own revisioned rows: schedule rt.SyncPreferences (watermark with a commit overlap) and rt.ResyncPreferences (full re-send); see HOST_INTEGRATION.md.

Testing

CONTENTKIT_TEST_URL=postgres://...  CONTENTKIT_PROFILE_URL=postgres://... \
CONTENTKIT_TEST_CH_ADDR=localhost:9000 CONTENTKIT_TEST_CH_USER=... CONTENTKIT_TEST_CH_PASSWORD=... CONTENTKIT_TEST_CH_CLUSTER=... \
go test ./... -race -count=1 -p 2

Media tests also need an S3 backend (CONTENTKIT_TEST_S3_ENDPOINT, _ACCESS_KEY, _SECRET_KEY, optional _REGION, _BUCKET and _REQUIRE; see media/internal/s3test). CI runs them on MinIO; media/image runs in its own CI job with libvips, and the other jobs exclude it. To record a Ceph RGW release's capabilities, point the same variables at an RGW bucket and run go test ./media/... -v -count=1; the log prints the probed capabilities. Manifest edits run under PGLocker, so the tests also need CONTENTKIT_TEST_URL (they skip without it). media/video tests also need ffmpeg and ffprobe on PATH (they skip without them unless CONTENTKIT_TEST_FFMPEG=1).

Backend Conditional PUT SHA-256 enforced Notes
MinIO RELEASE.2025-09-07 yes yes drops AbortIncompleteMultipartUpload (expires uploads itself)
Ceph 19.2 / 20.2 standalone dbstore RGW no no wrong Range bytes; not representative of RADOS-backed RGW
production Ceph RGW (RADOS, 2026-09-24) no (If-Match yes, If-None-Match: * ignored) no Range, response-content-disposition, versioning, lifecycle (incl. AbortIncompleteMultipartUpload) and multipart work; hosts wire PGLocker

Tests run against real PGroonga Postgres and ClickHouse+Keeper and skip without the variables; CONTENTKIT_PROFILE_URL needs CREATEDB. Regenerate the eval baseline with CONTENTKIT_EVAL_UPDATE=1.

Documentation

Overview

Package contentkit is the deterministic content library: tenant-scoped keyword search over host content, the ClickHouse signal plane, and the discovery reads over both. The DocumentSink port publishes document changes to optional external consumers.

Index

Constants

View Source
const PreferenceEventID = "current"

PreferenceEventID is the stable signal identity of a subject's current preference on one axis: (tenant, canonical reference, subject, axis, "current"). A newer revision supersedes the previous one; a re-send carries the same revision and converges.

Variables

View Source
var ErrSignalPlaneDisabled = errors.New("contentkit: signal plane disabled (no ClickHouse configured)")

ErrSignalPlaneDisabled is returned by signal/discovery methods when the hub was constructed without a ClickHouse connection.

Functions

func Migrate

func Migrate(ctx context.Context, cfg MigrateConfig) error

Migrate applies PostgreSQL migrations in one host-selected schema and the optional ClickHouse signal baseline.

func NewEvalRunner

func NewEvalRunner(client *Client, base SearchOptions) eval.CaseRunner

NewEvalRunner adapts a Client to eval.CaseRunner so a golden suite can be executed against real search. The base options carry cross-case settings (LanguageMode and filters); each case overrides Language, ContentKinds, and Limit from its own definition.

This adapter is the single seam where the client meets the dependency-free eval package.

Types

type CandidateTrace

type CandidateTrace struct {
	Key   TraceKey `json:"key"`
	Rank  int      `json:"rank"`
	Score float32  `json:"score"`
}

CandidateTrace records one source candidate at its raw source rank.

type CatalogPage added in v0.58.8

type CatalogPage struct {
	IDs        []string
	NextCursor string
}

CatalogPage contains ids in catalog order and the cursor for the next page. An empty NextCursor means the catalog is exhausted.

type CatalogQuery

type CatalogQuery struct {
	Limit  int
	Cursor string
}

CatalogQuery requests one bounded page. An empty Cursor starts at the first page; subsequent cursors are opaque to ContentKit and owned by the host.

type Client

type Client struct {
	// contains filtered or unexported fields
}

Client answers keyword search and typeahead for one tenant.

func NewClient

func NewClient(cfg ClientConfig) (*Client, error)

func (*Client) Search

func (c *Client) Search(ctx context.Context, userText string, opts SearchOptions) (SearchResult, error)

Search returns one page of content items. Documents from every searched language are grouped per work before Offset and Limit apply; each hit carries the matched document's reference and language.

func (*Client) SearchWithTrace

func (c *Client) SearchWithTrace(ctx context.Context, userText string, opts SearchOptions) (SearchResult, SearchTrace, error)

SearchWithTrace executes Search and returns opt-in retrieval provenance. On failure, the returned trace contains all work completed before the error.

func (*Client) Tenant

func (c *Client) Tenant() string

Tenant returns the tenant this client is scoped to.

func (*Client) Typeahead

func (c *Client) Typeahead(ctx context.Context, userText string, opts TypeaheadOptions) ([]TypeaheadHit, error)

Typeahead returns suggestions while a user is typing (typos/substring matching), one per content item, grouped before Limit.

func (*Client) WalkSearchMatches added in v0.59.0

func (c *Client) WalkSearchMatches(ctx context.Context, tx pgx.Tx, text string, opts SearchMatchOptions, visit func(context.Context, []SearchHit) error) error

WalkSearchMatches visits batches of up to 256 matched works in identity order, each carrying its best eligible matched edition. It does not paginate or rank the works. A successful return means the complete match set was visited; errors must not be interpreted as a complete count or page. The visitor receives the request's bounded context and must propagate it to downstream work. The caller owns tx, using the client's database, and must roll back on error.

type ClientConfig

type ClientConfig struct {
	Pool   *pgxpool.Pool
	Schema string
	// Tenant scopes every document, query and result. Required.
	Tenant string

	// Defaults.
	DefaultLanguage string
	DefaultLimit    int
}

ClientConfig configures the keyword search client of one tenant.

type ContentCatalog

type ContentCatalog interface {
	Page(ctx context.Context, tenant string, contentKind string, q CatalogQuery) (CatalogPage, error)
}

ContentCatalog pages live, visible content ids of a kind from the host's tables. The host owns visibility and ordering; a cursor must continue that same stable order without repeating or skipping ids.

type ContentCatalogFunc

type ContentCatalogFunc func(ctx context.Context, tenant string, contentKind string, q CatalogQuery) (CatalogPage, error)

ContentCatalogFunc adapts a function to the ContentCatalog interface.

func (ContentCatalogFunc) Page added in v0.58.8

func (f ContentCatalogFunc) Page(ctx context.Context, tenant string, contentKind string, q CatalogQuery) (CatalogPage, error)

type ContentKey

type ContentKey = contentref.ContentKey

ContentKey is the comparable form of a ContentRef.

type ContentRef

type ContentRef = contentref.ContentRef

ContentRef is the tenant-scoped reference to host-owned content (the work or one of its versions). See contentref.

type ContributionTrace

type ContributionTrace struct {
	SourceIndex  int     `json:"source_index"`
	SourceRank   int     `json:"source_rank"`
	Weight       float32 `json:"weight"`
	Contribution float32 `json:"contribution"`
}

ContributionTrace records one exact source contribution to a result score.

type DocumentKey

type DocumentKey = search.DocumentKey

DocumentKey identifies one keyword document: a ContentRef in one language.

type DocumentSink

type DocumentSink = search.DocumentSink

DocumentSink is the optional document port (see search.DocumentSink): the worker delivers every published keyword document to it at least once.

type Eligibility

type Eligibility = search.Eligibility

Eligibility is the host's per-document eligibility join; see search.Eligibility.

type EmbeddedConfig

type EmbeddedConfig struct {
	// Content plane (Postgres). PG + PGSchema are required. PGSchema is the
	// host-selected schema containing all ContentKit tables; it may also hold
	// application tables.
	PG       *pgxpool.Pool
	PGSchema string

	// Content-plane defaults (as in ClientConfig).
	DefaultLanguage string
	DefaultLimit    int
	// DefaultRRFK controls deterministic popularity/discovery fusion (default 60).
	DefaultRRFK int

	// Signal plane (ClickHouse). Optional: omit CH to run content-only
	// (signal/discovery methods return ErrSignalPlaneDisabled). CHDatabase
	// is the hub's dedicated ClickHouse database; apply
	// migrations.ClickHouse and gate startup on signal.CheckSchema.
	CH         signal.Conn
	CHDatabase string

	// Tenant is the single tenant of this embedded hub. Required: every
	// document, signal, cursor and result carries it.
	Tenant string

	// Scorers maps content kind → host Scorer. When a signal arrives for a
	// registered kind, the scorer's result overwrites Score / Progress /
	// ProgressMax / Completed before recording.
	Scorers map[string]signal.Scorer

	// Catalogs maps content kind → host ContentCatalog (the Unseen universe).
	Catalogs map[string]ContentCatalog

	// Candidates sources SimilarTo and Recommend (requires CH). Nil =
	// discovery.Engagement (co-engagement) over this hub's signal store.
	Candidates discovery.Candidates
}

EmbeddedConfig configures an in-process hub against the shared DB.

type EmbeddedHub

type EmbeddedHub struct {
	// contains filtered or unexported fields
}

EmbeddedHub implements Hub in-process. Construct with NewEmbedded.

func NewEmbedded

func NewEmbedded(cfg EmbeddedConfig) (*EmbeddedHub, error)

NewEmbedded builds the embedded hub: in-process, shared DB, one tenant.

func (*EmbeddedHub) Attribution

Attribution exports renders at one stage with their clicks joined (see signal.Store.Attribution); the evaluation dataset source.

func (*EmbeddedHub) Client

func (h *EmbeddedHub) Client() *Client

Client returns the underlying content-plane client (advanced use).

func (*EmbeddedHub) Content

func (h *EmbeddedHub) Content(contentKind, contentID string) ContentRef

Content returns a reference to a work of this hub's tenant.

func (*EmbeddedHub) EnforceErasures

func (h *EmbeddedHub) EnforceErasures(ctx context.Context) (signal.ErasureReport, error)

EnforceErasures physically removes residue of every recorded erasure of this tenant (see signal.Store.EnforceErasures): schedule it and run it after every restore.

func (*EmbeddedHub) EraseSubjects

func (h *EmbeddedHub) EraseSubjects(ctx context.Context, subjects []signal.Subject) (signal.ErasureReport, error)

EraseSubjects permanently erases subjects from this tenant's signal plane: see signal.Store.EraseSubjects for the completion contract. Shared accounts exist in several tenants: each host erases its own tenant.

func (*EmbeddedHub) Forget

func (h *EmbeddedHub) Forget(ctx context.Context, subject signal.Subject, contentKind, contentID string) error

Forget erases the subject's signals for one work and its versions (contentID set) or a whole content kind (contentID empty) — host "clear my history" support.

func (*EmbeddedHub) ForgetExposures

func (h *EmbeddedHub) ForgetExposures(ctx context.Context, subject signal.Subject) error

ForgetExposures clears a subject's result-list exposures (search history).

func (*EmbeddedHub) ForgetExposuresBefore added in v0.15.0

func (h *EmbeddedHub) ForgetExposuresBefore(ctx context.Context, subject signal.Subject, before time.Time) error

ForgetExposuresBefore clears only exposures from before the host's clear request.

func (*EmbeddedHub) History

func (h *EmbeddedHub) History(ctx context.Context, subject signal.Subject, opts signal.HistoryOptions) ([]signal.StateRow, error)

func (*EmbeddedHub) HistoryCount

func (h *EmbeddedHub) HistoryCount(ctx context.Context, subject signal.Subject, opts signal.HistoryOptions) (int64, error)

HistoryCount returns the total row count History would paginate over.

func (*EmbeddedHub) Inventory

func (h *EmbeddedHub) Inventory(ctx context.Context) ([]signal.InventoryRow, error)

Inventory reports canonical event volume per content kind and signal type.

func (*EmbeddedHub) Metrics

func (h *EmbeddedHub) Metrics(ctx context.Context, refs []ContentRef, window signal.Window) (map[ContentKey]signal.ContentMetrics, error)

Metrics returns named window metrics (viewers, views, completions, feedback, ...) for the references (works or versions).

func (*EmbeddedHub) Popular

func (h *EmbeddedHub) Popular(ctx context.Context, contentKind string, opts signal.PopularOptions) ([]signal.PopularHit, error)

func (*EmbeddedHub) PopularityFor

func (h *EmbeddedHub) PopularityFor(ctx context.Context, contentKind string, ids []string, window signal.Window) (map[string]float64, error)

PopularityFor scores a fixed candidate set (work ids of one kind) by the popularity ranking, returning content_id -> score. Use to rank a host- filtered universe (e.g. "galleries of artist X by popularity").

func (*EmbeddedHub) PurgeContentKinds

func (h *EmbeddedHub) PurgeContentKinds(ctx context.Context, contentKinds []string) error

PurgeContentKinds irreversibly deletes whole content kinds from this tenant's signal plane (see signal.Store.PurgeContentKinds).

func (*EmbeddedHub) Recommend

func (h *EmbeddedHub) Recommend(ctx context.Context, subject signal.Subject, opts RecommendOptions) ([]RecHit, error)

Recommend returns "for you" works for a subject from the configured Candidates source, excluding seen (unless IncludeSeen) and disliked works, with a popularity fill for cold start. The host hydrates the references.

func (*EmbeddedHub) RecordExposures

func (h *EmbeddedHub) RecordExposures(ctx context.Context, exposures []signal.Exposure) error

RecordExposures logs one row per result list and stage (served, rendered, visible) so clicks can be attributed to what was actually exposed. Hosts call it once per list per stage, never per item.

func (*EmbeddedHub) RecordSignals

func (h *EmbeddedHub) RecordSignals(ctx context.Context, signals []signal.Signal) error

RecordSignals applies each content kind's registered Scorer, then records the batch (see signal.Store.RecordSignals). A scorer error records nothing.

func (*EmbeddedHub) RefreshCoEngagement

func (h *EmbeddedHub) RefreshCoEngagement(ctx context.Context, opts signal.RefreshCoEngagementOptions) error

RefreshCoEngagement (re)materializes the content_pairs co-engagement rollup for this tenant (see signal.Store.RefreshCoEngagement). Run periodically.

func (*EmbeddedHub) RepairProjections

func (h *EmbeddedHub) RepairProjections(ctx context.Context, opts signal.RepairOptions) (signal.RepairResult, error)

RepairProjections is the bounded, host-scheduled projection repair (see signal.Store.RepairProjections): run it periodically with IngestedSince for crash repair, and with Rebuild over a window after projection changes.

func (*EmbeddedHub) Search

func (h *EmbeddedHub) Search(ctx context.Context, userText string, opts HubSearchOptions) (SearchResult, error)

func (*EmbeddedHub) SeenIDs

func (h *EmbeddedHub) SeenIDs(ctx context.Context, subject signal.Subject, contentKind string) (map[string]struct{}, error)

SeenIDs returns the subject's seen-set for one content kind (the signal-plane half of the unseen anti-join). Use when the host wants to run its own diff against a custom-filtered universe instead of Unseen's registered catalog.

func (*EmbeddedHub) SimilarTo

func (h *EmbeddedHub) SimilarTo(ctx context.Context, ref ContentRef, opts SimilarOptions) ([]RecHit, error)

SimilarTo returns works like the anchor ("more like this") from the configured Candidates source (default: co-engagement).

func (*EmbeddedHub) States

func (h *EmbeddedHub) States(ctx context.Context, subject signal.Subject, refs []ContentRef) (map[ContentKey]signal.State, error)

func (*EmbeddedHub) Tenant

func (h *EmbeddedHub) Tenant() string

Tenant returns the pinned tenant value.

func (*EmbeddedHub) Typeahead

func (h *EmbeddedHub) Typeahead(ctx context.Context, userText string, opts TypeaheadOptions) ([]TypeaheadHit, error)

func (*EmbeddedHub) Unseen

func (h *EmbeddedHub) Unseen(ctx context.Context, subject signal.Subject, opts UnseenOptions) ([]string, error)

Unseen returns catalog ids the subject has not seen (max_progress > 0 defines "seen"). The host catalog applies its own visibility decisions; ContentKit checks the subject's state only for each bounded candidate page.

type EmptyReason

type EmptyReason string

EmptyReason explains a successful empty response.

const (
	EmptyReasonNormalizedQuery EmptyReason = "normalized_query_empty"
	EmptyReasonNoCandidates    EmptyReason = "no_candidates"
)

type Hub

type Hub interface {
	// Tenant returns the tenant this hub instance is scoped to.
	Tenant() string

	// Content plane.
	Search(ctx context.Context, userText string, opts HubSearchOptions) (SearchResult, error)
	Typeahead(ctx context.Context, userText string, opts TypeaheadOptions) ([]TypeaheadHit, error)
	SimilarTo(ctx context.Context, ref ContentRef, opts SimilarOptions) ([]RecHit, error)

	// Signal plane.
	RecordSignals(ctx context.Context, signals []signal.Signal) error
	RecordExposures(ctx context.Context, exposures []signal.Exposure) error
	ForgetExposures(ctx context.Context, subject signal.Subject) error
	ForgetExposuresBefore(ctx context.Context, subject signal.Subject, before time.Time) error
	Attribution(ctx context.Context, opts signal.AttributionOptions) (signal.AttributionPage, error)
	Forget(ctx context.Context, subject signal.Subject, contentKind, contentID string) error
	EraseSubjects(ctx context.Context, subjects []signal.Subject) (signal.ErasureReport, error)
	EnforceErasures(ctx context.Context) (signal.ErasureReport, error)

	// Discovery plane.
	History(ctx context.Context, subject signal.Subject, opts signal.HistoryOptions) ([]signal.StateRow, error)
	HistoryCount(ctx context.Context, subject signal.Subject, opts signal.HistoryOptions) (int64, error)
	SeenIDs(ctx context.Context, subject signal.Subject, contentKind string) (map[string]struct{}, error)
	Unseen(ctx context.Context, subject signal.Subject, opts UnseenOptions) ([]string, error)
	States(ctx context.Context, subject signal.Subject, refs []ContentRef) (map[ContentKey]signal.State, error)
	Metrics(ctx context.Context, refs []ContentRef, window signal.Window) (map[ContentKey]signal.ContentMetrics, error)
	Popular(ctx context.Context, contentKind string, opts signal.PopularOptions) ([]signal.PopularHit, error)
	PopularityFor(ctx context.Context, contentKind string, ids []string, window signal.Window) (map[string]float64, error)
	Recommend(ctx context.Context, subject signal.Subject, opts RecommendOptions) ([]RecHit, error)

	// Maintenance.
	RefreshCoEngagement(ctx context.Context, opts signal.RefreshCoEngagementOptions) error
	RepairProjections(ctx context.Context, opts signal.RepairOptions) (signal.RepairResult, error)
	Inventory(ctx context.Context) ([]signal.InventoryRow, error)
	PurgeContentKinds(ctx context.Context, contentKinds []string) error
}

Hub is the single surface host apps program against: content-plane queries (search/typeahead), the signal plane (RecordSignals), and the discovery plane (reads over content × signals). All methods return ranked content references (+ per-subject State); the host hydrates them into cards from its own DB. Every method is scoped to the tenant pinned at construction; a reference of another tenant is an error.

type HubSearchOptions

type HubSearchOptions struct {
	SearchOptions
	Personalize *Personalization
}

HubSearchOptions extends content SearchOptions with optional signal-aware personalization.

type KeywordDocument

type KeywordDocument = search.KeywordDocument

KeywordDocument is the host's canonical search input for one document.

type LanguageMode

type LanguageMode string
const (
	// LanguageModeExact uses only the requested language.
	LanguageModeExact LanguageMode = "exact"
	// LanguageModeFallbackEnglish uses requested language first, then English.
	LanguageModeFallbackEnglish LanguageMode = "fallback_en"
)

type MigrateConfig

type MigrateConfig struct {
	// DB holds PostgreSQL DDL credentials. Required.
	DB *sql.DB
	// Schema receives every PostgreSQL table; it may be the application's schema.
	// Required. Identifiers contain letters, numbers or underscores.
	Schema string
	// ClickHouse applies the signal baseline when set; PostgresDB defaults to DB.
	ClickHouse *chmigrate.Config
}

MigrateConfig selects the host-owned stores for ContentKit migrations.

type Personalization

type Personalization struct {
	Subject signal.Subject

	// PopularityWeight is the RRF weight of the candidate-set popularity
	// list blended with the content ranking. Defaults to 0.25.
	PopularityWeight float32
	// PopularityWindow bounds candidate popularity (zero = all time).
	PopularityWindow signal.Window

	// AffinityWeight boosts works the subject already engaged with by
	// (1 + AffinityWeight·last_score/100). 0 = off.
	AffinityWeight float32

	// DemoteSeen demotes already-seen / completed works.
	DemoteSeen bool
	// SeenPenalty multiplies seen-but-not-completed scores (default 0.85).
	SeenPenalty float32
	// CompletedPenalty multiplies completed scores (default 0.6).
	CompletedPenalty float32
	// DislikePenalty multiplies works the subject has net-negative explicit
	// feedback for (default 0.3). Always applied when view context is loaded
	// (i.e. AffinityWeight > 0 or DemoteSeen).
	DislikePenalty float32
}

Personalization fuses signal aggregates into search ranking. Recall is unchanged — this is a ranking-only layer over the candidate set. A per-request toggle the host flips (e.g. only for logged-in users).

type PublishedDocument

type PublishedDocument = search.PublishedDocument

PublishedDocument is one keyword document as delivered to a DocumentSink.

type RecHit

type RecHit = discovery.Hit

RecHit is one ranked work of SimilarTo or Recommend.

type RecommendOptions

type RecommendOptions = discovery.RecommendOptions

RecommendOptions controls Recommend (see discovery.RecommendOptions).

type ResultTrace

type ResultTrace struct {
	Key           TraceKey            `json:"key"`
	Rank          int                 `json:"rank"`
	Score         float32             `json:"score"`
	ScoreKind     ScoreKind           `json:"score_kind"`
	Contributions []ContributionTrace `json:"contributions"`
}

ResultTrace records one returned item and its best document's contributions.

type RetrievalBackend

type RetrievalBackend string

RetrievalBackend identifies one candidate source.

const (
	BackendKeyword RetrievalBackend = "keyword"
)

type Runtime

type Runtime struct {
	*EmbeddedHub
	Content *content.Runtime
}

Runtime is the one surface a host wires: the Hub (search, typeahead, signals, discovery) plus the content module and its HTTP routes.

func NewRuntime

func NewRuntime(ctx context.Context, cfg RuntimeConfig) (*Runtime, error)

NewRuntime builds the hub and the content module over the host pool.

func (*Runtime) EraseSubjects

func (r *Runtime) EraseSubjects(ctx context.Context, subjects []signal.Subject) (signal.ErasureReport, error)

EraseSubjects erases every configured runtime plane: signals, interaction data and data retained by moderation/classifier providers. Current approved authored content remains under host retention policy (content.EraseSubjects). ContentKit-owned reactions, favorites, poll votes and unpublished submissions are removed atomically behind the source fence. EmbeddedHub.EraseSubjects is the explicit analytics-only lower-level API.

AuthKit ACK means durable acceptance by the host's deletion ledger, not this downstream completion. Retry while error != nil or !report.Complete(). A disabled signal plane is intentionally absent. Remaining may include pending plane markers content_plane or signal_plane when a plane could not complete; these are not estimates of retained provider rows.

func (*Runtime) Handler

func (r *Runtime) Handler() http.Handler

Handler returns the content routes (comments, reactions, favorites, polls, posts). Mount it under a prefix after the host's auth middleware.

func (*Runtime) ResyncPreferences added in v0.15.0

func (r *Runtime) ResyncPreferences(ctx context.Context) (content.PreferenceSyncReport, error)

ResyncPreferences re-sends every exportable preference: the periodic repair for sink loss. Newer revisions still win.

func (*Runtime) SyncPreferences added in v0.15.0

func (r *Runtime) SyncPreferences(ctx context.Context) (content.PreferenceSyncReport, error)

SyncPreferences sends reactions and favorites changed since the last sync into the signal plane (see content.Runtime.SyncPreferences). Schedule it from the host's worker; with the signal plane disabled it returns ErrSignalPlaneDisabled.

func (*Runtime) WorkerOptions

func (r *Runtime) WorkerOptions(host worker.Options) worker.Options

WorkerOptions returns the host's keyword worker options extended with ContentKit's own documents: posts (content.KindPost) are listed and built by the content module, every other kind by the host callbacks. Missing host callbacks for configured kinds fail backfill.

type RuntimeConfig

type RuntimeConfig struct {
	EmbeddedConfig
	Content content.Options
}

RuntimeConfig configures one tenant's full ContentKit: the search and signal planes (EmbeddedConfig) and the interaction module (Content). Pool, tenant and schema are shared: Content.Pool, Content.Tenant and Content.Schema are filled from the hub configuration when empty. Posts join the keyword queue.

type ScoreKind

type ScoreKind string

ScoreKind identifies the numeric domain of a candidate or result score.

const (
	ScoreKeywordMatch ScoreKind = "keyword_match"
)

type SearchHit

type SearchHit struct {
	// ContentRef names the matched document's content: the work, or the
	// version when the document was indexed per version.
	ContentRef
	// Language is the matched document's language.
	Language string
	// Score ranks the item by its best matching document in any searched
	// language, using the keyword match tier.
	Score float32
}

SearchHit is one content item, represented by its matched document.

type SearchMatchOptions added in v0.59.0

type SearchMatchOptions struct {
	Language     string
	ContentKinds []string
	Eligibility  *Eligibility
	FilterSQL    string
	FilterArgs   map[string]any
}

SearchMatchOptions selects all matches in one language, without a candidate window. Hosts use this for catalog-wide non-relevance ordering and counts.

type SearchOptions

type SearchOptions struct {
	Language string
	// Defaults to LanguageModeExact when omitted.
	LanguageMode LanguageMode

	// ContentKinds selects the kinds searched. Required.
	ContentKinds []string

	// Limit is the page size in content items; Offset skips items. Documents
	// are grouped per item before either applies.
	Limit  int
	Offset int
	// CandidateLimit is the document window requested from each retrieval
	// source (per language) before grouping. It defaults to twice Offset+Limit
	// (at least 100) and is clamped to at least Offset+Limit. Pass the same
	// value on every page when a truncated window must stay identical.
	CandidateLimit int

	// Eligibility maps each document to the host's access, publication and
	// version-trait rules on that one document. Without it every document of
	// a work is eligible.
	Eligibility *Eligibility

	FilterSQL  string
	FilterArgs map[string]any
}

type SearchResult

type SearchResult struct {
	Hits []SearchHit
	// HasMore is true when items follow this page in the grouped retrieval, or
	// when Truncated: documents beyond the window were never ranked, so the
	// next page may still be non-empty.
	HasMore bool
	// Truncated reports that a candidate window filled. Raise CandidateLimit
	// for complete deep pagination.
	Truncated bool
}

SearchResult is one page of items.

type SearchTrace

type SearchTrace struct {
	NormalizedQuery         string        `json:"normalized_query"`
	RequestedLanguage       string        `json:"requested_language,omitempty"`
	RequestedLanguageMode   LanguageMode  `json:"requested_language_mode,omitempty"`
	Languages               []string      `json:"languages,omitempty"`
	RequestedResultLimit    int           `json:"requested_result_limit"`
	ResultLimit             int           `json:"result_limit"`
	RequestedCandidateLimit int           `json:"requested_candidate_limit"`
	CandidateLimit          int           `json:"candidate_limit"`
	Sources                 []SourceTrace `json:"sources,omitempty"`
	Results                 []ResultTrace `json:"results,omitempty"`
	EmptyReason             EmptyReason   `json:"empty_reason,omitempty"`
	ErrorCategory           string        `json:"error_category,omitempty"`
}

SearchTrace contains opt-in effective configuration and retrieval provenance.

type SimilarOptions

type SimilarOptions = discovery.SimilarOptions

SimilarOptions controls SimilarTo (see discovery.SimilarOptions).

type SourceStatus

type SourceStatus string

SourceStatus records whether an attempted retrieval source succeeded.

const (
	SourceStatusSucceeded SourceStatus = "succeeded"
	SourceStatusFailed    SourceStatus = "failed"
)

type SourceTrace

type SourceTrace struct {
	Backend       RetrievalBackend `json:"backend"`
	Language      string           `json:"language"`
	ScoreKind     ScoreKind        `json:"score_kind"`
	Limit         int              `json:"limit"`
	Status        SourceStatus     `json:"status"`
	ErrorCategory string           `json:"error_category,omitempty"`
	Candidates    []CandidateTrace `json:"candidates,omitempty"`
}

SourceTrace records one language-specific keyword retrieval and its candidates.

type TaxonomyID

type TaxonomyID = contentref.TaxonomyID

TaxonomyID identifies a generic ContentKit catalog record (tag, artist, series, creator, character, voice actor).

type TraceKey

type TraceKey struct {
	ContentRef
	Language string `json:"language"`
}

TraceKey identifies a document in retrieval provenance.

type TypeaheadHit

type TypeaheadHit struct {
	ContentRef
	Language string
	Score    float32
}

TypeaheadHit is one suggested content item and its matched document.

type TypeaheadOptions

type TypeaheadOptions struct {
	Language string
	// Defaults to LanguageModeExact when omitted.
	LanguageMode  LanguageMode
	ContentKinds  []string
	Limit         int
	MinSimilarity float32
	FilterSQL     string
	FilterArgs    map[string]any
	// Eligibility groups suggestions per content item; see SearchOptions.
	Eligibility *Eligibility
}

type UnseenOptions

type UnseenOptions struct {
	// ContentKind selects the catalog to read. Required.
	ContentKind string
	// Limit caps the returned ids (default 50). Order follows the host catalog.
	Limit int
	// CatalogLimit is the per-page candidate count (default and max 1,000).
	CatalogLimit int
}

UnseenOptions controls Unseen reads.

Directories

Path Synopsis
Package access is the host's gating vocabulary shared by content interactions and media: the authenticated Actor, the ContentResolver port and its Resolution.
Package access is the host's gating vocabulary shared by content interactions and media: the authenticated Actor, the ContentResolver port and its Resolution.
adapters
authkit module
cmd
media-access command
Command media-access is the media access agent (media/agent).
Command media-access is the media access agent (media/agent).
media-worker command
Command media-worker is the stock media worker (media/worker) for hosts whose registry is plain data: it reads it from a JSON file (media.Registry's JSON).
Command media-worker is the stock media worker (media/worker) for hosts whose registry is plain data: it reads it from a JSON file (media.Registry's JSON).
Package content is ContentKit's interaction module: posts, comments, reactions, favorites and polls over tenant-scoped content references, stored in the host schema's content_* interaction tables.
Package content is ContentKit's interaction module: posts, comments, reactions, favorites and polls over tenant-scoped content references, stored in the host schema's content_* interaction tables.
Package contentref defines the reference vocabulary every ContentKit package and port shares: a tenant-scoped reference to host-owned content and the identity of a generic taxonomy record.
Package contentref defines the reference vocabulary every ContentKit package and port shares: a tenant-scoped reference to host-owned content and the identity of a generic taxonomy record.
Package discovery produces "similar to this work" and "for you" lists.
Package discovery produces "similar to this work" and "for you" lists.
Package eval provides dependency-free golden-query evaluation, aggregate search-quality metrics, versioned reports, baseline comparison, and score-domain-safe threshold sweeps.
Package eval provides dependency-free golden-query evaluation, aggregate search-quality metrics, versioned reports, baseline comparison, and score-domain-safe threshold sweeps.
internal
boundaries
Package boundaries holds a test enforcing ContentKit's package layering from the module's dependency graph (go list; no compilation or services).
Package boundaries holds a test enforcing ContentKit's package layering from the module's dependency graph (go list; no compilation or services).
normalize
Package normalize cleans user query text before retrieval.
Package normalize cleans user query text before retrieval.
pglock
Package pglock holds Postgres session advisory locks on dedicated connections, so the holder's work may use every connection of its pool.
Package pglock holds Postgres session advisory locks on dedicated connections, so the holder's work may use every connection of its pool.
pgtest
Package pgtest provisions disposable keyword-profile schemas on the CONTENTKIT_TEST_URL Postgres for integration tests.
Package pgtest provisions disposable keyword-profile schemas on the CONTENTKIT_TEST_URL Postgres for integration tests.
signaltest
Package signaltest provisions disposable signal-plane ClickHouse databases for integration tests by applying the real migration lineage.
Package signaltest provisions disposable signal-plane ClickHouse databases for integration tests by applying the real migration lineage.
tcpproxy
Package tcpproxy is a test TCP forwarder that a test takes down (connections refused) and brings back on the same address, to stand in for a provider outage.
Package tcpproxy is a test TCP forwarder that a test takes down (connections refused) and brings back on the same address, to stand in for a provider outage.
Package media stores app media as a self-describing file system in one private bucket.
Package media stores app media as a self-describing file system in one private bucket.
agent
Package agent is the media access agent run by cmd/media-access.
Package agent is the media access agent run by cmd/media-access.
image
Package image is the media worker's image producer (libvips, CGO): the Image presets' private WebP files, the Zip presets, the public presets' fixed names, editor views, and the kinds' default public images.
Package image is the media worker's image producer (libvips, CGO): the Image presets' private WebP files, the Zip presets, the public presets' fixed names, editor views, and the kinds' default public images.
internal/s3test
Package s3test opens the test bucket from the environment:
Package s3test opens the test bucket from the environment:
internal/uploadtestserver command
Command uploadtestserver serves the media upload and read APIs over an isolated MinIO/RGW namespace for the browser SDK's integration tests (sdk/upload/test) and e2e specs.
Command uploadtestserver serves the media upload and read APIs over an isolated MinIO/RGW namespace for the browser SDK's integration tests (sdk/upload/test) and e2e specs.
internal/videotest
Package videotest builds synthetic videos whose frames identify their time and orientation, and classifies decoded pixels, for poster and preview tests.
Package videotest builds synthetic videos whose frames identify their time and orientation, and classifies decoded pixels, for poster and preview tests.
internal/wirets
Package wirets renders the upload and read API wire types as TypeScript for the browser SDK (sdk/upload/src/wire.gen.ts).
Package wirets renders the upload and read API wire types as TypeScript for the browser SDK (sdk/upload/src/wire.gen.ts).
layout
Package layout defines media object keys, dependency-free so the access agent can classify paths without importing the media runtime:
Package layout defines media object keys, dependency-free so the access agent can classify paths without importing the media runtime:
s3
Package s3 implements media.Store over aws-sdk-go-v2 for Ceph RGW (production) and MinIO (tests).
Package s3 implements media.Store over aws-sdk-go-v2 for Ceph RGW (production) and MinIO (tests).
tiered
Package tiered is an optional visibility policy: it maps an item's level to entitlement keys and asks a Checker which ones the actor holds.
Package tiered is an optional visibility policy: it maps an item's level to entitlement keys and asks a Checker which ones the actor holds.
token
Package token signs and verifies media access tokens, shared by the host signer and the access agent so the format cannot drift:
Package token signs and verifies media access tokens, shared by the host signer and the access agent so the format cannot drift:
video
Package video runs an item's video-family presets with ffmpeg: HLS (a byte-range fMP4 ladder with the source's audio and text tracks and a seek sprite), MP4 (a muxed H.264 file at one rung), Audio (an HLS track and an M4A) and Subtitles (clean WebVTT), and grabs the frames of uploads that declare Upload.Frames.
Package video runs an item's video-family presets with ffmpeg: HLS (a byte-range fMP4 ladder with the source's audio and text tracks and a seek sprite), MP4 (a muxed H.264 file at one rung), Audio (an HLS track and an M4A) and Subtitles (clean WebVTT), and grabs the frames of uploads that declare Upload.Frames.
worker
Package worker is the media worker: the one process that places uploads and runs every producer, from the host's worker River schema (Config.Schema) in its database: staged uploads hashed and placed at their blobs (media.Manifests.Place), images, zips, public presets and editor views (media/image, libvips), and HLS, MP4, audio, subtitles and frames (media/video, ffmpeg).
Package worker is the media worker: the one process that places uploads and runs every producer, from the host's worker River schema (Config.Schema) in its database: staged uploads hashed and placed at their blobs (media.Manifests.Place), images, zips, public presets and editor views (media/image, libvips), and HLS, MP4, audio, subtitles and frames (media/video, ffmpeg).
workqueue
Package workqueue is the host's side of the media worker (media/worker): the River schema it drains in the host database, insert-only enqueueing and cancelling of its jobs, and processing progress for the read API.
Package workqueue is the host's side of the media worker (media/worker): the River schema it drains in the host database, insert-only enqueueing and cancelling of its jobs, and processing progress for the read API.
workqueue/metrics
Package metrics exports the host's media worker queue health.
Package metrics exports the host's media worker queue health.
Package migrations owns ContentKit's PostgreSQL and ClickHouse migrations.
Package migrations owns ContentKit's PostgreSQL and ClickHouse migrations.
Package popularity ranks host content by a named, bounded policy over the signal plane's canonical window metrics (docs/popularity-policy.md).
Package popularity ranks host content by a named, bounded policy over the signal plane's canonical window metrics (docs/popularity-policy.md).
Package search is ContentKit's keyword retrieval over content_search_documents: exact names and aliases, native-script prefixes and bounded typos for every language, with the host's eligibility join applied inside every route.
Package search is ContentKit's keyword retrieval over content_search_documents: exact names and aliases, native-script prefixes and bounded typos for every language, with the host's eligibility join applied inside every route.
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.
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.
Package taxonomy is ContentKit's generic catalog: tenant-scoped nodes (tags, artists, creators, characters, series, seasons, voice actors, ...) with localized names and aliases, typed node relationships and typed assignments of host content (a work or one of its versions) to nodes.
Package taxonomy is ContentKit's generic catalog: tenant-scoped nodes (tags, artists, creators, characters, series, seasons, voice actors, ...) with localized names and aliases, typed node relationships and typed assignments of host content (a work or one of its versions) to nodes.
Package worker maintains one tenant's keyword documents: it drains the dirty queue, runs a bounded cursor backfill and delivers every published document to the optional DocumentSink.
Package worker maintains one tenant's keyword documents: it drains the dirty queue, runs a bounded cursor backfill and delivers every published document to the optional DocumentSink.

Jump to

Keyboard shortcuts

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