contentkit

package module
v0.30.0 Latest Latest
Warning

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

Go to latest
Published: Sep 24, 2026 License: MIT Imports: 20 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
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 ContentResolver port and its Resolution{Ref, Visible, Accessible, PreviewLimit}, shared by content and media
media per-item folders and keys, kind registry, the Store port, manifests with conditional-write edits, direct uploads and their HTTP API, the optional UploadLimiter, sweep, folder deletion and processing 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) processor: WebP variants, public slots, zip downloads
media/token media access tokens, shared by hosts and the access worker
media/video ffmpeg encode jobs: byte-range fMP4 HLS ladder, AAC per audio track, WebVTT per text subtitle, sprite, per-quality MP4 downloads, poster frames and hover previews; Frames for the poster picker; River in schema media_worker (cmd/media-worker)
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 one PostgreSQL baseline and one ClickHouse baseline
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},
})

These are fresh-store baselines, not an in-place upgrade of old 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
})

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

Design: MEDIA-DESIGN.md. One private bucket; each item owns a folder the library keys:

{tenant}/{kind}/{id}/manifest.json | manifests/{version}.json
                    /originals/{sha256-hex | u-uuid | slot | slot.json | i-uuid}   never served
                    /blobs/sha256-{hex}                                            immutable derivatives
                    /public/{slot_width | i-uuid}.webp                             slots, inline images

Host wiring (one tenant; errors elided):

kinds, _ := media.NewRegistry(
	media.Kind{Name: "gallery", Versioned: true, Types: []string{"image/png", "image/jpeg"}, MaxBytes: 10 << 20,
		Specs: map[string]media.Spec{"thumb": {Width: 460, Height: 650, Fit: media.FitCover, Quality: 80}, "high": {Quality: 90}},
		Slots: map[string]media.Slot{"cover": {Aspect: 460.0 / 650, Widths: []int{230, 460, 920}}},
		Zip:   "high"},
	media.Kind{Name: "video", Types: []string{"video/mp4", "video/x-matroska"}, MaxBytes: 20 << 30, Video: true})
store, _ := s3.New(s3.Config{Bucket: "media", Endpoint: rgw, PublicEndpoint: "https://s3.doujins.ai", UsePathStyle: true,
	AccessKeyID: id, SecretAccessKey: secret, Capabilities: caps}) // caps from media.Probe
key, _ := token.ParseKey(os.Getenv("MEDIA_TOKEN_KEY")) // "{kid}:{base64}", shared with media-access
ring, _ := token.NewRing(key, nil)

jobs, _ := media.NewJobs(media.JobsConfig{Store: store, Kinds: kinds, Tenants: []string{"d"}, Limiter: limiter})
manifests, _ := media.NewManifests(store, kinds, media.ManifestOptions{Locker: media.PGLocker(pool), Jobs: jobs})
images, _ := image.New(image.Config{Store: store, Kinds: kinds, Manifests: manifests})
_ = jobs.AddProcessor(images.Process)
_ = video.Migrate(ctx, pool) // River schema media_worker, run by cmd/media-worker
videos, _ := video.NewEnqueuer(pool, kinds)
_ = jobs.AddProcessor(videos.Processor())
client, _ := riverhelpers.New(ctx, pool, &river.Config{Schema: "public"}, runtime.RiverJobs(), jobs.RiverJobs())

uploads, _ := media.NewUploads(media.UploadOptions{Store: store, Kinds: kinds, Manifests: manifests,
	Authorizer: hostUploads, Tickets: &ring, Limiter: limiter, Queue: jobs})
reader, _ := media.NewReader(media.ReaderOptions{Manifests: manifests, Kinds: kinds, Resolver: resolver, Hooks: hooks,
	Delivery: media.Delivery{Mode: media.DeliverCookie, BaseURL: "https://media.doujins.com", CookieDomain: "doujins.com", SigningKey: key}})
mux.Handle("/api/media/upload/", http.StripPrefix("/api/media/upload", media.UploadHandler(uploads, media.UploadHandlerOptions{Tenant: "d", Actor: actorOf,
	PublicBaseURL: "https://media.doujins.com"})))
mux.Handle("/api/media/", http.StripPrefix("/api/media", reader.Handler(media.HandlerOptions{Tenant: "d", Identity: identity})))

_ = jobs.DeleteItemsTx(ctx, tx, media.Deletion{Ref: ref, Owner: owner}) // in the host's delete transaction
_ = jobs.EraseUserTx(ctx, tx, "d", userID, deletions...)                 // the user's items plus user/{id}/

The access worker (cmd/media-access, image ghcr.io/open-rails/contentkit-media-access:{tag}, same tag as the hosts' ContentKit) serves BaseURL. It needs MEDIA_ACCESS_S3_ENDPOINT, _S3_BUCKET, a read-only key (_S3_ACCESS_KEY_ID, _S3_SECRET_ACCESS_KEY) allowed only */blobs/* and */public/*, MEDIA_ACCESS_TOKEN_KEY and _TOKEN_KEY_PREVIOUS (the hosts' {kid}:{base64} ring), and optionally MEDIA_ACCESS_HOSTS and MEDIA_ACCESS_CORS_ORIGINS (the sites, with credentials); secrets may be given as {VAR}_FILE. public/ is served without a token (no-cache); blobs/ needs ?t= or an mt cookie (private, immutable; 403 without a valid token); manifests, originals/ and unknown keys are 404. cmd/media-worker documents its environment.

Edit writes with If-Match (or If-None-Match: *) and retries on conflict; without Capabilities.ConditionalPut it serializes on a Postgres advisory lock instead. Reads are cached in process and revalidated by ETag. Presigned PUTs bind Content-Type, Content-Length and x-amz-checksum-sha256.

Uploads go straight to the bucket (media.Uploads, served by media.UploadHandler): the host's UploadAuthorizer.CanUpload (AuthKit) runs at presign and commit for the folder written (slots and inline images: the work, ref.Content()), and the kind's types and size cap bind every presign. Up to 64 MiB is one PUT to originals/sha256-{hex} signed with its type, length and SHA-256; larger files are multipart to originals/u-{uuid} with 8–16 MiB parts, each signed with its length and SHA-256, resumed through ListParts and completed by the server (a signed ticket carries the S3 UploadId; nothing is stored). Slot originals PUT to originals/{slot}. A kind with Inline takes inline images: presign with inline: true names a new i-{uuid}, whose original PUTs to originals/{id} and is committed with commit-slot; it is re-encoded with the Inline spec to public/{id}.webp. Commit is one conditional manifest edit (insert, replace, move, rename, remove, edit) that HEAD-checks each new original, re-hashes it when the store does not enforce checksums, and enqueues a ProcessJob. Kind.MaxFiles and per-type Kind.TypeLimits ("image", "video": MaxBytes replacing the kind's, MaxFiles) cap a manifest: a commit that ends over a cap and adds to it fails with 409 too_many_files. A kind may mix images (Specs) and videos (Video); each processor handles only its own files.

Edits are non-destructive: File.Edit{Crop{x,y,w,h}, Rotate} crops in the source's pixels (EXIF orientation applied), then rotates clockwise by 0, 90, 180 or 270. The edit op sets or (without edit) clears it; insert and replace may carry one. It is checked against File.Dims, the source's size recorded by processing (before that, by the processor, which reports an out-of-bounds edit to Hooks.Failed). A variant's spec is Spec.For(edit), so changing or clearing an edit re-derives that file's variants (and the zip) from the untouched original. Spec.EditorOnly variants (recorded editor: true) are signed only when the resolver's Resolution.Editor is set; Spec.Unedited (must be EditorOnly) ignores edits: an editor's view of the whole source. Under a folder cookie such a blob is unlisted, not locked: its content-hash name is never sent to other viewers. No master is written. meta.w/h is the edited size; the read API returns edit and dims to editors.

Slots are fixed public images such as avatars and covers, rendered at several widths for high-density screens:

Slots: map[string]media.Slot{
	"avatar": {Aspect: 1, Widths: []int{128, 256, 512}},
	"cover":  {Aspect: 3, Widths: []int{1500, 3000}, MinWidth: 1500},
}

A slot's Edit uses the same crop (original pixels, EXIF-oriented) and rotate; the crop's height follows its width at Aspect (the edited width/height), and no crop means the largest centred one. The original PUTs to originals/{slot} and POST /commit-slot {ref, slot, sha256, edit} commits it; POST /edit-slot {ref, slot, edit} re-edits the kept original without an upload; Uploads.SetSlotFromFile(ctx, actor, SlotFromFile{Ref, Slot, From, File, Edit}) (POST /commit-slot-from-file {ref, slot, from, file, edit}) copies a manifest image (of From, default Ref: another item of the tenant needs CanUpload on both; default edit: the file's own); POST /slot-original returns the original to uploaders for the editor. The edit and the job's result live in originals/{slot}.json, so spec changes re-encode with it. Each width is public/{slot}_{width}.webp; widths wider than the edited image are skipped, never upscaled, and an edit outside the original or narrower than MinWidth (at least the smallest width) is refused (by the job when the original's size is not yet known: Hooks.Failed, keeping the served outputs). Slot routes and the read API's GET /{kind}/{id}/slots/{slot} answer SlotManifest{aspect, edit, dims, version, outputs: [{name, w, h, url}], pending, error}. Output URLs carry ?v={version}, which the access worker serves immutable while current; listings build them without reads from one stored value per slot: Hooks.SlotEncoded reports each encode's SlotStamp ("{version}:{w},{w}…", the version and produced widths) before the slot record shows that version, and Reader.SlotOutputs(ref, slot, stamp) returns exactly that encode's outputs (the zero stamp: widths up to MinWidth, unversioned). SlotManifest.Stamp() backfills a stamp from a read.

Image processing (media/image, CGO over libvips via govips; install libvips-dev to build it). image.New(Config{Store, Kinds, Manifests, Specs, Hooks}) gives Process(ctx, media.ProcessJob); register it with jobs.AddProcessor(proc.Process) and pass jobs as UploadOptions.Queue. It derives WebP variants per the kind's Specs (or a per-file SpecChooser) from each file's master, else original, through its edit, only where a variant is missing or its spec differs, stores them as blobs/sha256-…, and records them in one manifest edit per pass that drops results for sources replaced meanwhile; it repeats until a commit that landed during the run is covered too. A kind with Zip set gets downloads.zip: a stored zip of that variant in file order, rebuilt only when its inputs hash changes; its display name comes from Hooks.DownloadName at read time. Slots and inline images re-encode originals/{slot} into public/ in place (no-cache, ETag), with writes conditional on the output's previous ETag and skipped when they already match. Undecodable sources go to Hooks.Failed and are not retried.

The optional UploadLimiter (media.NewPGLimiter over the baseline's content_media_* tables) rate-limits uploaders (files/hour, bytes/day → 429) and enforces per-owner quota. The rule: a commit that grows the owner's stored originals (the change in the manifest's distinct originals) past its quota fails with 413 quota_exceeded and writes nothing. Presign reservations only refuse early (used + pending + size); they lapse after a day, which never lets a late commit past the quota. Exempt grants are charged but never refused.

The bucket needs CORS allowing PUT from the app origins with the Content-Type and x-amz-checksum-sha256 headers, and the AbortIncompleteMultipartUpload: 1 day rule Store.Configure sets.

Video (media/video, run by cmd/media-worker) encodes each video/* manifest file in one ffmpeg pass: H.264 High (CRF 22, preset fast, keyframes every 4 s) at the kind's ladder (Kind.Video = &media.Video{Ladder: []int{1080, 720, 480}}; default media.DefaultLadder, 2160/1440/1080/720/480), AAC per audio track, WebVTT per text subtitle and a 10×10 sprite whose tiles keep the source aspect (short side 90). A rung N is the output's short side (a 1080 rung of a vertical video is 1080 wide); rungs above the source's short side are dropped (a smaller source gets one rung at its own short side). Every frame is then capped, aspect kept, at 4096 px per side and a 3840×2160 area (common H.264 hardware decode limits), so a 21:9 2160 rung is 4096×1756 and an 8K source is downscaled to 3840×2160; a rung whose capped frame repeats the next one's is dropped. Output is square-pixel (SAR applied), at most 60 fps, and frames above 1080p-class carry the lowest fitting level (5.0/5.1/5.2); smaller ones keep x264's. hls.video[] records rung and the true w/h. Sources whose display aspect is outside Video.MinAspect–MaxAspect (default 1/2.4–2.4, admitting 2560×1080 and 2.39:1 cinema; 0.5% slack) fail permanently: the file's hls becomes {source, spec, error} with no renditions, Hooks.Failed (in the encoder's process) gets video.ErrAspect, editors see failed in the read API, and it is retried only when the source or the kind's bounds change. Each rendition and audio track is one single-file fMP4 blob whose segments are [offset, length, seconds] (the init segment is [0, segments[0].offset)); each rung also gets a muxed MP4 in downloads["{file}-{N}p"] (video, every audio track, subtitles). Blobs are written first; one manifest edit then records hls and downloads only if the file still derives from the encoded original, so a replaced file keeps its previous hls until then. Outputs are byte-identical on retry. ffmpeg reads only local files (-protocol_whitelist file) through container demuxers (mov/mp4, matroska/webm, avi, mpegts, flv, ogg, asf, mpeg): playlists and concat lists are refused. Changing the ladder changes video.Spec(ladder), so files re-encode. Jobs live in River schema media_worker in the host database: hosts run video.Migrate and enqueue through video.NewEnqueuer (insert-only; register enqueuer.Processor() with media.Jobs.AddProcessor); the worker's environment is documented in cmd/media-worker.

After each encode the job grabs the item's poster frame (the poster slot; the image job encodes it) and renders its hover preview (silent MP4 and animated WebP loops) from their selections; see HOST_INTEGRATION "Video posters and hover previews".

Encode progress: with ReaderOptions.Progress: video.NewProgressSource(pool) the read API adds progress to each visible video file still pending (none yet, or a replaced source), and GET /{kind}/{id}/video-images adds the item's current step. The worker parses ffmpeg -progress and writes, at most every Config.ProgressInterval (2 s) plus on phase changes, a per-file map to its own River row (metadata.contentkit_progress, cleared when the job ends; no extra table). Contract (media.EncodeProgress): phase (queued downloading probing encoding muxing uploading publishing, then item-wide images), queue_position (1 = next; waiting jobs only), segments_done/segments_total (HLS segments, ceil(duration/4)), percent (time-based, never decreasing), speed (×realtime, smoothed over ~8 s), eta (seconds: remaining media / speed plus projected uploads / measured throughput), at (unix ms), stalled (a running job silent for a minute). It reveals only timing and queue depth, so every viewer allowed the file gets it; one indexed query, only for items with a pending video.

Playback is served by Reader.Handler next to the read API, generated per request after one Resolve (private, no-store; the folder cookie is set in cookie mode): /{kind}/{id}/hls/{file}/master.m3u8?audio=&subs= (optional id/language filters; RESOLUTION is the rung's true w×h), video/{N}.m3u8, audio/{id}.m3u8, subs/{id}.m3u8, sprite.vtt, and /{kind}/{id}/download/{key} (302 to the signed dl= URL, full access only; name from Hooks.DownloadName). Media playlists are EXT-X-BYTERANGE lines over one blob URL per rendition. A file plays when the grant allows it (full access, inside a preview cut, or a teaser); preview viewers get per-file URL tokens. In the browser, hls.js needs xhrSetup: xhr => { xhr.withCredentials = true } in cookie mode and the worker's Origins must list the site; native Safari/iOS HLS should be checked in cookie mode and switched to URL mode if it does not send the cookie.

Tokens are kid.exp.base64url(HMAC-SHA256(secret, "{scope}|{exp}")): a scope is a folder (…/blobs/, covering the objects directly under it), one key, or {key}#dl={name} for a download name. Expiry is window-aligned (default 4 h); token.Ring verifies the current and previous key.

Media's River jobs (jobs.RiverJobs()) compose into the host client through helpers/river; edits schedule a sweep and commits enqueue processing:

  • Sweep (per folder, 24 h after each edit and in a daily pass over Tenants): deletes blobs/ and hash-named originals/ no manifest in the folder references, only once every manifest and the object itself are older than Grace (24 h; plus 1 day for u- multipart objects, which may be dated at initiation). Slot originals, public/ and manifests are never swept. Invariant: it deletes only objects no manifest references and no in-flight commit can newly reference. Presign reuses an existing original, and a commit accepts one, only while a manifest references it or it is well before the sweep's cutoff (grace/2 for presign; grace/4, at most 1 h, for commit); otherwise the client uploads it again. Set UploadOptions.Grace to the same grace (taken from Queue when it is the *Jobs).
  • Deletion removes the whole folder, manifests first, then again after LateUploadWindow (25 h) for PUTs and multipart completions that land late. With a Limiter, the owner's quota (the manifests' OriginalBytes) is released once.
  • Processing: jobs.Enqueue (the uploads' ProcessQueue) runs one pending job per ref and slot through every AddProcessor processor. An Enqueue (or ScheduleSweep) while an equal job runs queues one follow-up that starts after it, since the running job may have read its inputs before the change; an equal job still waiting absorbs it.
  • Media packages add workers with jobs.Register(func(*river.Config) error) before composition and enqueue with jobs.Insert/InsertTx.
  • Restore: docs/restore.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. On a backend without conditional PUT the tests edit manifests under PGLocker, so they 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 installs all PostgreSQL features in one host-selected schema and the optional ClickHouse signal plane. This baseline initializes fresh stores.

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 CatalogQuery

type CatalogQuery struct {
	Limit int
}

CatalogQuery bounds a Universe read. Limit 0 = host-defined default.

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.

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 {
	Universe(ctx context.Context, tenant string, contentKind string, q CatalogQuery) ([]string, error)
}

ContentCatalog supplies the content "universe" for Unseen: live, non-deleted content ids of a kind, read from the host's own tables. The host owns visibility and gating (premium, region, ...) — ContentKit never interprets them. Order defines Unseen order (recommended: newest first).

type ContentCatalogFunc

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

ContentCatalogFunc adapts a function to the ContentCatalog interface.

func (ContentCatalogFunc) Universe

func (f ContentCatalogFunc) Universe(ctx context.Context, tenant string, contentKind string, q CatalogQuery) ([]string, 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"): host universe MINUS the subject's seen-set. The host catalog applies its own visibility/premium gating against its own tables.

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 the two ContentKit baselines.

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.

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 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 which catalog universe to diff against. Required.
	ContentKind string
	// Limit caps the returned ids (default 50). Order follows the host
	// catalog's Universe order.
	Limit int
	// CatalogLimit is passed through to the host catalog's Universe call
	// (0 = host default).
	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.
cmd
media-access command
Command media-access is the media access worker.
Command media-access is the media access worker.
media-worker command
Command media-worker runs media/video encode jobs from River schema media_worker in the host database.
Command media-worker runs media/video encode jobs from River schema media_worker in the host database.
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.
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.
Package media stores host content files in per-item folders of one private bucket: library-built keys, the generic manifest with conditional-write edits, and the Store port.
Package media stores host content files in per-item folders of one private bucket: library-built keys, the generic manifest with conditional-write edits, and the Store port.
accessworker
Package accessworker is the media access worker's HTTP handler, run by cmd/media-access: it checks the token for a blob path (URL `?t=` or cookie `mt`), serves public/ paths without one (immutable at a current ?v= version), refuses manifests and originals/, and streams the object from the private bucket with its own read-only key.
Package accessworker is the media access worker's HTTP handler, run by cmd/media-access: it checks the token for a blob path (URL `?t=` or cookie `mt`), serves public/ paths without one (immutable at a current ?v= version), refuses manifests and originals/, and streams the object from the private bucket with its own read-only key.
image
Package image derives WebP variants, public slots and zip downloads with libvips (CGO).
Package image derives WebP variants, public slots and zip downloads with libvips (CGO).
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 media.UploadHandler over a fresh MinIO/RGW bucket for the browser SDK's integration tests (sdk/upload/test).
Command uploadtestserver serves media.UploadHandler over a fresh MinIO/RGW bucket for the browser SDK's integration tests (sdk/upload/test).
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 API wire types as TypeScript for the browser SDK (sdk/upload/src/wire.gen.ts).
Package wirets renders the upload 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 worker can classify paths without importing the media runtime:
Package layout defines media object keys, dependency-free so the access worker 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 worker so the format cannot drift:
Package token signs and verifies media access tokens, shared by the host signer and the access worker so the format cannot drift:
video
Package video encodes an item's video files with ffmpeg into a byte-range HLS ladder (one single-file fMP4 blob per rendition and audio track), WebVTT subtitles, a thumbnail sprite and one muxed MP4 download per quality, and records them in the manifest's hls and downloads.
Package video encodes an item's video files with ffmpeg into a byte-range HLS ladder (one single-file fMP4 blob per rendition and audio track), WebVTT subtitles, a thumbnail sprite and one muxed MP4 download per quality, and records them in the manifest's hls and downloads.
Package migrations owns ContentKit's PostgreSQL and ClickHouse baselines.
Package migrations owns ContentKit's PostgreSQL and ClickHouse baselines.
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