contentkit

package module
v0.58.2 Latest Latest
Warning

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

Go to latest
Published: Sep 26, 2026 License: MIT Imports: 21 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, PreviewLimit, Editor}, 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: byte-range fMP4 HLS ladder, AAC per audio track, WebVTT per text subtitle and subtitle sidecar, audio files, sprite, per-quality MP4 downloads, poster frames; Frames for the poster picker
media/worker the media worker: one process for placement, images and video, built by the host from its media config (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 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            the one manifest (versions, slots, index); never served
                    /originals/sha256-{hex}   uploads, deduped per item; never served
                    /temp/u-{uuid}            staged multipart uploads until placed; never served
                    /temp/e-{hex}             editor views (Kind.Editor); editor token only
                    /private/sha256-{hex}     every rendition (token)
                    /public/sha256-{hex}      copies of the exposed renditions: slots and inline images of an item that is not hidden

originals/, private/ and public/ names are their content's SHA-256, so those objects are immutable: a change writes new names and the manifest-driven sweep deletes what the manifest no longer lists. temp/ is intermediary and discardable: nothing a viewer needs lives there, and the sweep wipes it by age. Clients never build URLs; the API returns them.

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: media.Ratio("46:65"), Widths: []int{230, 460, 920}}},
		Editor: &media.Spec{Width: 1200, Height: 1200, Fit: media.FitInside, Quality: 80},
		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}) // never dials; capabilities come from the first Check
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, Locker: media.PGLocker(pool), Kinds: kinds, Tenants: []string{"d"}, Limiter: limiter, Resolver: resolver})
manifests, _ := media.NewManifests(store, kinds, media.ManifestOptions{Locker: media.PGLocker(pool), Sweeps: jobs})
_ = workqueue.Migrate(ctx, pool, "doujins_media_worker") // this host's worker schema, drained by its media worker
queue, _ := workqueue.New(pool, kinds, "doujins_media_worker")
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: queue, ProcessOnUpload: true})
reader, _ := media.NewReader(media.ReaderOptions{Manifests: manifests, Kinds: kinds, Resolver: resolver, Hooks: hooks,
	Progress: workqueue.NewProgressSource(pool), Queue: queue,
	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,
	Reader: reader})))
mux.Handle("/api/media/", http.StripPrefix("/api/media", reader.Handler(media.HandlerOptions{Tenant: "d", Identity: identity})))

_ = jobs.ExposeTx(ctx, tx, ref)                                           // in any transaction that changes whether anonymous viewers see it
_ = 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 */private/*, */public/* and */temp/e-*, MEDIA_ACCESS_TOKEN_KEY and _TOKEN_KEY_PREVIOUS (the hosts' {kid}:{base64} ring), MEDIA_ACCESS_HOSTS (the media host names; empty serves any Host, warned) and MEDIA_ACCESS_CORS_ORIGINS (the sites' exact origins, with credentials; empty breaks hls.js, warned; wildcards and paths are refused); secrets may be given as {VAR}_FILE. public/ is served without a token (public, immutable); private/ needs ?t= or an mt cookie (private, immutable); a temp/e- editor view needs ?t= with an editor token (token.EditorScope, which no viewer token carries). Everything refused (no or bad token, the manifest, originals/, staged uploads, unknown keys) is one identical no-store 404, so denials look like absence. Every object carries Cross-Origin-Resource-Policy: same-site (MEDIA_ACCESS_RESOURCE_POLICY=cross-origin only when the pages live on another site than the media), so other sites cannot embed it with <img>/<video>. See HOST_INTEGRATION "Production media delivery".

The media worker (media/worker) is the one process that does media work: it hashes and places staged uploads, derives image variants, zips, slot outputs and inline images (libvips) and encodes video and poster frames (ffmpeg), from the host's worker River schema (media/workqueue, worker.Config.Schema, MEDIA_WORKER_SCHEMA) in the host database. The schema is required and per host (e.g. doujins_media_worker, hentai0_media_worker): hosts sharing a database must not share one, or each worker takes the other's jobs. Queue names are fixed within it. The host presigns, commits, publishes and reads, and links only media/workqueue (no libvips, no ffmpeg). The worker must apply the host's exact kinds and policy, so the host builds it from the same code that builds its media.Registry, image.SpecChooser and media.Hooks (Failed, SlotEncoded and ItemReady run in the worker), e.g. as a subcommand of the host binary:

cfg, _ := worker.FromEnv(ctx) // DATABASE_URL, MEDIA_S3_*, MEDIA_WORKER_SCHEMA, MEDIA_HOST_RIVER_SCHEMA, MEDIA_WORKER_* (see worker.FromEnv)
cfg.Kinds, cfg.Specs, cfg.Hooks = kinds, specs, hooks // the host's media config package
w, _ := worker.New(ctx, cfg) // no DDL: the host's migrate step runs workqueue.Migrate(ctx, pool, schema)
_ = w.Run(ctx) // until SIGTERM; running jobs get MEDIA_WORKER_SHUTDOWN_GRACE

worker.New needs no DDL rights, so the worker runs as the host's unprivileged app role; the host applies workqueue.Migrate with its other migrations.

cmd/media-worker (image ghcr.io/open-rails/contentkit-media-worker) is the stock build for hosts whose kinds are plain data: it reads them from MEDIA_KINDS_FILE (a JSON array of media.Kind). The worker hands a video item's poster publish, and folder sweeps after its edits, back to the host's River schema (MEDIA_HOST_RIVER_SCHEMA), where jobs.RiverJobs() runs them with the host's Resolver. It never exits for a missing dependency: it runs no DDL (the host migrates MEDIA_WORKER_SCHEMA) and its build is retried while Postgres is down; Run takes no jobs until the bucket answers, and MEDIA_METRICS_ADDR serves /livez, /readyz (built), /statusz and app_dependency_up with /metrics.

Process on upload (UploadOptions.ProcessOnUpload, default false): the presign reply tells the SDK to commit each file as soon as it is uploaded, {op: "insert", unattached: true}, so the worker processes it while the user is still arranging the upload. Unattached files are charged to the quota, count against the kind's caps and are left out of every read (editors ask for them with ReadOptions.Unattached, POST /files), zips and the automatic poster. {op: "attach", name} makes one part of the item after the attached files (or at index), without reprocessing; remove of an unattached file discards it: the item's worker jobs are cancelled (workqueue.Queue.Cancel, every stage) and re-enqueued for the rest, and its staged or placed original and derivatives no manifest references are deleted at once. The SDK's UploadQueue does all of this: commit() attaches in queue order, remove() discards, and item.processing (dims, hls, failed, progress) is polled until item.processed.

Readiness. Manifests.Readiness(ctx, ref) (Root.Readiness(kind)) is ready when every attached file, set slot and video poster is processed (videos: every stage, no hls.pending; images: variants for the current source and edit), processing while any is not, and failed once nothing is processing and some could not be (Failed names them). After every image or video job that leaves an item settled, the worker calls Hooks.ItemReady(ctx, tx, ref, readiness) in a transaction on the host database; an error retries the job, so it must be idempotent. A host that holds content back until its media is ready publishes it there (and enqueues its Expose with HostQueue.ExposeTx in the same tx). Independently, reads never show a non-editor a file with nothing processed to serve, or a failed one (File.Servable): media added to live content appears once processed.

Edit always runs under the Locker (required: a Postgres advisory lock, PGLocker, 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 bucket is optional at startup. s3.New never dials; register store.Check(ctx, prefix) as the host's optional S3 dependency probe (helpers deps). Its first success probes the backend's capabilities unless Config.Capabilities declares them; a probe records nothing unless every step succeeded or was refused cleanly (412, checksum mismatch, 501), so a throttled or cut-off probe is retried. Until then the store claims none (locked unconditional edits, server-side rehash). An unreachable or 5xx bucket surfaces as media.ErrUnavailable: 503 unavailable from the read and upload handlers. In media jobs it becomes a River snooze (not an attempt) only when the job's own context is live and a fresh bounded Check confirms the bucket is down, capped by MaxOutageSnoozes; a job that outran its timeout or one broken object spends attempts (media.SnoozeUnavailable). 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 temp/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). The manifest names a staged upload u-{uuid} until Manifests.Place moves it to originals/sha256-{hex} with the hash computed while reading it (server-side copy, or none when the folder already holds the hash; every reference renamed; the temp upload deleted; idempotent). Slot and inline originals are hash-named too (originals/sha256-{hex}, deduped). A kind with Inline takes inline images: presign with inline: true names a new i-{uuid}, committed with commit-slot; it is rendered with the Inline spec to private/ and copied to public/, and Reader.InlineURL returns its URL (ErrPending until rendered). 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. No master is written. meta.w/h is the edited size; the read API returns edit and dims to editors.

Editor views (Kind.Editor, a Spec) are what croppers draw on: the whole source, EXIF-oriented, ignoring crop and rotate, for image files and slot originals. They are an input-keyed cache, temp/e-{hex} of (source, spec) (Item.EditorView), never in the manifest: the image job renders them, the sweep deletes them after JobsConfig.EditorTTL, and a missing one is rendered again when an editor asks (ReaderOptions.Queue; the slot routes use UploadOptions.Queue). Editors (Resolution.Editor) get them as the read API's variant=editor and as editor_url in slot manifests, signed with an editor token no viewer token equals; a crop in progress is drawn by the client, so nothing uncommitted is stored.

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: media.Aspect1x1, Widths: []int{128, 512}}, // small, large
	"cover":  {Aspect: media.Ratio("3:1"), Widths: []int{900, 3000}, MinWidth: 600},
}

A slot's Edit uses the same crop (original pixels, EXIF-oriented) and rotate; the crop's height follows its width at Aspect (a media.Aspect ratio in lowest terms, written "W:H" in JSON and config: media.Ratio("9:16"), ParseAspect, constants Aspect1x1, Aspect3x1, Aspect4x5, Aspect16x9, Aspect9x16, Aspect21x9; all maths is integer, heights round half up), and no crop means the largest centred one. The original PUTs to originals/sha256-{hex} and POST /commit-slot {ref, slot, sha256, edit, filename} 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). Originals never leave the server: editors re-crop on the slot's editor_url. The record (original, edit, result) lives in the manifest's slots, so spec changes re-encode with it. Each width is a new private/sha256-{hex}, copied to public/ unless the item is hidden; a change writes new names and swaps the record. Nothing is upscaled: a width wider than the edited image is rendered at the edited width, so every width exists once the slot is set. An edit outside the original or narrower than MinWidth (default 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, outputs: [{w, h, url}], pending, error}: public URLs for viewers, private/ URLs with a token for editors of a hidden item. Listings link a slot without reads: Hooks.SlotEncoded(ctx, ref, slot, listing) hands the host a SlotListing to store, and Reader.ListedSlot(ref, slot, listing) builds its URLs. Slot{Aspect: media.AspectNative} keeps the edited image's own shape: no crop by default, crops of any shape.

Image processing (media/image, CGO over libvips via govips; install libvips-dev to build it; run by the media worker). image.New(Config{Store, Kinds, Manifests, Specs, Hooks}) gives Process(ctx, media.ProcessJob). A staged source is hashed from the bytes read for decoding and placed first. 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 private/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 render from their original only when the record's fingerprint changed. 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 the media worker) encodes each video/* manifest file at the kind's ladder (Kind.Video = &media.Video{Ladder: []int{1080, 480}}; default media.DefaultLadder, 2160/1080/480) in each of the worker's codecs (video.Config.Codecs, MEDIA_WORKER_CODECS; default av1,h264, hevc optional), plus 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, and a source below 1080 that is not a rung gets one at its own short side (720p: 720 + 480). Every frame is then capped, aspect kept, at 4096 px per side and a 3840×2160 area (common 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), 8-bit 4:2:0, at most 60 fps, an IDR every 4 s without scene cuts (one per 4 s segment), closed GOPs. Rates are a capped CRF per rung (live-action H.264: 2160 CRF 23 at most 32 Mbit/s, 1440 23/18M, 1080 23/12M, 720 22/7M, 480 21/3M; VBV buffer 2× the cap; HEVC one CRF lower and AV1 CRF 34–30, both at 0.6× the caps; Video.Profile: media.VideoAnimation tunes for animation at lower CRFs and caps). H.264 is High (frames above 1080p-class carry the lowest fitting level 5.0/5.1/5.2), HEVC Main tagged hvc1 (Safari requires it), AV1 Main. Config.Encoder auto (default) uses NVENC for each codec whose probe encode works and the CPU otherwise (libx264 and libx265 preset fast, Config.Preset up to 1080 and TopPreset above; SVT-AV1 preset 8), cpu never NVENC, nvenc requires it; a file NVENC fails on re-encodes on the CPU. New probe-encodes each codec and fails when its encoder is missing or ignores forced keyframes (libsvtav1 needs ffmpeg ≥ 7 with SVT-AV1 ≥ 2). hls.video[] records rung, codec, CODECS and the true w/h, ordered by codec (as configured), then largest rung first. 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 H.264 MP4 (the first codec when H.264 is not configured) 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 (same encoders, presets and Threads). Progressive stages: one stage per rung, smallest first. A stage decodes the source once, scales it (lanczos) to its rung and encodes it in every codec; the first also makes the tracks and the sprite. Each stage is published at once (hls.pending lists the rungs to come), so a viewer plays 480p while 1080p and 2160p encode; the worker queues each next stage as a follow-up job (same args, River priority 2), behind other uploads' first stages. Every codec advances together because hls.js picks one codec set at start and never switches it for bandwidth. Progress reports stage/stages. Passthrough: when the source already is a compliant top rung (MP4/MOV constant-rate 8-bit 4:2:0 progressive H.264 High/Main or HEVC Main tagged hvc1, ≤ level 5.2, unrotated, at the rung's exact frame, within its bitrate cap in that codec, with an IDR starting each 4 s segment) that rung is stream-copied in its codec, provided its segments match the published rung below; otherwise it is encoded. Playback: the master playlist lists every rung in every codec with its CODECS (avc1…, hvc1…, av01…), codecs in configured order, each starting at its 1080 rung; media playlists are video/{rung}-{codec}.m3u8. hls.js drops variants MediaSource.isTypeSupported refuses and Safari those it cannot decode, so H.264 is the fallback; the SDK player keeps one codec set (see sdk/upload). 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, profile, codecs or recipe (versioned; bumped when the encode defaults change) changes Encoder.Spec(video), so files re-encode once; encoder (CPU/NVENC) and preset choices are not part of it. A staged source is hashed while it downloads for ffmpeg and placed before the encode, so hls.source names the placed original. Jobs are {ref} (workqueue.VideoArgs; the worker takes the kind from its registry), not unique, and a job for a fresh manifest is a no-op; workqueue.Queue.Cancel cancels an item's queued and running jobs of every stage.

Audio (Kind.Audio = &media.Audio{}; the media worker's own media_audio queue and per-manifest lock, so audio never waits behind video encodes; MEDIA_WORKER_AUDIO_CONCURRENCY, default 2) encodes each audio/* file (mp3, m4a/mp4, wav, flac, ogg/opus, aac, mka/webm, aiff, caf, wma): its default (else first) audio stream to AAC-LC 128 kbit/s, 48 kHz stereo (downmixed before any measuring), as a one-track HLS ladder (hls.audio, no video; the master playlist is one audio-only variant) and a faststart M4A remuxed from it, the file's audio variant (?variant=audio, for an <audio> element) and its download {file}-audio. Language and label come from the stream or container tags. Audio.Loudness (LUFS, e.g. -16; default 0 = off) measures EBU R128 loudness and true peak of the stereo mix, then applies one linear gain, min(target − I, −1.5 − TP) dB (volume): one extra decode, no dynamics processing, so a quiet source limited by its peak stays below the target. The SDK MediaGallery plays audio items from their audio variant (request it in the read) with the file's download. An unreadable source fails like video (hls.error, Hooks.Failed). A kind that lists audio/ types must set Audio.

Subtitle sidecars (media.SubtitleTypes: WebVTT, SRT, SSA/ASS; a video kind only) are manifest files beside the video. The worker's video job converts each one to WebVTT, which becomes the file's vtt variant. The parser follows the file's type, never its content:

  • SRT goes through ffmpeg's WebVTT encoder (the one used for a source's own text tracks) after its timings are normalized (01:02,5 → 00:01:02,500);
  • SSA/ASS goes through ffmpeg after vector drawings ({\p1}…{\p0}) and {comment} blocks are stripped;
  • WebVTT is read directly.

Sidecars over 32 MB fail before they are downloaded. The charset comes from meta.charset, else a BOM, BOM-less UTF-16 or valid UTF-8, else the legacy charsets of meta.lang (Shift_JIS, GB18030, Big5, EUC-KR, Windows-125x). Without a hint, a CJK charset needs most high bytes to pair as its common characters (kana, Hangul, frequent Han), else Cyrillic scoring, else Windows-1252.

Every WebVTT output, sidecar or source track, is then cleaned:

  • only b/i/u markup stays (other tags are dropped and their text kept);
  • text is escaped and safe cue settings are kept;
  • NOTE, STYLE and REGION blocks are dropped, along with ASS positioning and colors and ASS vector drawings;
  • cues are sorted by start time.

A source track over 32 MB, or with no cues, is dropped; a sidecar with no cues fails (Failure, Hooks.Failed); the read API's ready marks a converted one. A source's tracks carry hls.subs_spec: a new cleaning recipe re-extracts them from the source without re-encoding the ladder. The master playlist lists the video's own tracks and then its sidecars (meta.for names the video, default the first, and follows a rename; meta.lang, label, forced; track id = file name). Adding, replacing or removing a sidecar re-encodes nothing, because playlists are built per request. MP4 downloads carry only the source's tracks.

After each encode the job grabs the item's poster frame (the poster slot; the image job encodes it) from its selection. There is no preview clip: the SDK previews the HLS itself inline; see HOST_INTEGRATION "Video posters and inline previews".

Encode progress: with ReaderOptions.Progress: workqueue.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; the first variant is the highest rung up to 1080p, where Safari/iOS native HLS starts, then the rest by descending bandwidth), 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 (…/private/, 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.

Reader.Handler limits each viewer (HandlerOptions.Limit, default 2 requests/s, burst 120; keyed by Actor.ID, else Actor.IP, else the peer address) with 429 rate_limited + Retry-After, and logs every signed response (media urls signed: viewer, ref, access, expiry, and a short hash of a folder token) so a leaked URL traces to its viewer.

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 originals/, private/ and public/ objects outside the manifest's index once the manifest and the object are older than Grace (24 h), and temp/ whatever the manifest's age: editor views older than EditorTTL (7 days) and staged uploads no file references older than TempUploadTTL (48 h: above the bucket's 1-day multipart abort rule, since multipart objects may be dated at initiation). A staged upload still being uploaded is not an object yet, and one being processed is referenced. Deleted public/ keys go to Hooks.PublicRemoved (CDN purge). S3 lifecycle rules cannot match */temp/* (filters are prefixes), so the sweep is the mechanism; AbortIncompleteMultipartUpload stays the backstop for uploads never completed. 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; a quarter of its retention, at most 1 h, for commit); otherwise the client uploads it again. Set UploadOptions.Grace and TempUploadTTL to the sweep's (taken from the Manifests' Sweeps when it is the *Jobs).
  • Deletion removes the whole folder, manifests and public/ 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.
  • Expose (jobs.ExposeTx in every transaction that changes whether anonymous viewers see an item: create a draft, publish, hide, delete, restore) resolves the item anonymously. Hidden: the manifest records it, public/ is emptied at once and the keys go to Hooks.PublicRemoved. Visible: the slot and inline outputs are copied back. private/ is never touched; free vs members-only is only whether the host grants a token.
  • Processing: workqueue.Queue (the uploads' ProcessQueue) inserts one pending image job per ref and slot, and a video job for a video or audio kind's manifest, into the worker's schema. 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. 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 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 is the stock media worker (media/worker) for hosts whose kinds are plain data: it reads them from a JSON file.
Command media-worker is the stock media worker (media/worker) for hosts whose kinds are plain data: it reads them from a JSON file.
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 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 private/ path (URL `?t=` or cookie `mt`), serves public/ paths without one and temp/ editor views only under an editor token (URL `?t=`, token.EditorScope, which viewer tokens never carry), refuses the manifest, originals/ and staged uploads, 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 private/ path (URL `?t=` or cookie `mt`), serves public/ paths without one and temp/ editor views only under an editor token (URL `?t=`, token.EditorScope, which viewer tokens never carry), refuses the manifest, originals/ and staged uploads, and streams the object from the private bucket with its own read-only key.
image
Package image derives WebP variants, slot and inline renditions and zip downloads with libvips (CGO).
Package image derives WebP variants, slot and inline renditions 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.
worker
Package worker is the media worker: the one process that does all media work, from the host's worker River schema (Config.Schema) in its database.
Package worker is the media worker: the one process that does all media work, from the host's worker River schema (Config.Schema) in its database.
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 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