diff --git a/migration/db/migrations/20260827100500000_add_signage_ai_tables.sql b/migration/db/migrations/20260827100500000_add_signage_ai_tables.sql new file mode 100644 index 00000000..2d0d6019 --- /dev/null +++ b/migration/db/migrations/20260827100500000_add_signage_ai_tables.sql @@ -0,0 +1,151 @@ +-- +micrate Up +-- SQL in section 'Up' is executed when this migration is applied + +-- Provider rows hold the vendor credentials used to generate signage artwork. +-- A row with a NULL authority_id is the shared fallback, mirroring "storages". +CREATE TABLE IF NOT EXISTS "signage_ai_providers" ( + id UUID PRIMARY KEY DEFAULT uuidv7(), + authority_id TEXT, + name TEXT NOT NULL, + provider TEXT NOT NULL, + + -- encrypted JSON object: api key, or Azure resource details, or a service account + credentials TEXT NOT NULL, + endpoint TEXT, + location TEXT, + default_model TEXT, + allowed_models TEXT[] NOT NULL DEFAULT '{}'::TEXT[], + + enabled BOOLEAN NOT NULL DEFAULT TRUE, + is_default BOOLEAN NOT NULL DEFAULT FALSE, + + -- {"user_per_day": 60, "domain_per_month": 2000} + quotas JSONB NOT NULL DEFAULT '{}'::jsonb, + + created_at TIMESTAMPTZ NOT NULL, + updated_at TIMESTAMPTZ NOT NULL, + + CONSTRAINT signage_ai_providers_provider_check + CHECK (provider IN ('OPENAI', 'AZURE_OPENAI', 'GOOGLE_VERTEX')), + CONSTRAINT signage_ai_providers_quotas_check + CHECK (jsonb_typeof(quotas) = 'object'), + FOREIGN KEY (authority_id) REFERENCES authority(id) ON DELETE CASCADE +); + +CREATE INDEX IF NOT EXISTS signage_ai_providers_authority_id_index + ON "signage_ai_providers" USING BTREE (authority_id); + +CREATE UNIQUE INDEX IF NOT EXISTS signage_ai_providers_authority_name_idx + ON "signage_ai_providers" (authority_id, name) + NULLS NOT DISTINCT; + +CREATE UNIQUE INDEX IF NOT EXISTS signage_ai_providers_authority_default_idx + ON "signage_ai_providers" (authority_id) + NULLS NOT DISTINCT + WHERE is_default = TRUE; + +-- +micrate StatementBegin +-- Only one default provider per authority; flip the others back to FALSE. +CREATE OR REPLACE FUNCTION signage_ai_providers_ensure_single_default() + RETURNS TRIGGER AS $$ +BEGIN + IF NEW.is_default THEN + UPDATE signage_ai_providers + SET is_default = FALSE + WHERE authority_id IS NOT DISTINCT FROM NEW.authority_id + AND id <> NEW.id; + END IF; + RETURN NEW; +END; +$$ LANGUAGE plpgsql; +-- +micrate StatementEnd + +DROP TRIGGER IF EXISTS trg_signage_ai_providers_single_default ON signage_ai_providers; + +CREATE TRIGGER trg_signage_ai_providers_single_default + BEFORE INSERT OR UPDATE OF is_default + ON signage_ai_providers + FOR EACH ROW + WHEN (NEW.is_default = TRUE) + EXECUTE FUNCTION signage_ai_providers_ensure_single_default(); + +-- One row per generate or edit request. Candidate images are written into +-- result->'images' one at a time by the runner, bumping "version" each write so a +-- long polling client can tell what it has already seen. +CREATE TABLE IF NOT EXISTS "signage_ai_jobs" ( + id UUID PRIMARY KEY DEFAULT uuidv7(), + authority_id TEXT NOT NULL, + provider_id UUID, + provider_type TEXT, + model TEXT, + + user_id TEXT, + user_email TEXT, + user_name TEXT, + + -- a refine points at the job it was refined from + parent_job_id UUID, + + version INTEGER NOT NULL DEFAULT 0, + kind TEXT NOT NULL, + state TEXT NOT NULL DEFAULT 'QUEUED', + cancel_requested BOOLEAN NOT NULL DEFAULT FALSE, + idempotency_key TEXT, + + candidates INTEGER NOT NULL DEFAULT 1, + images_produced INTEGER NOT NULL DEFAULT 0, + + prompt TEXT, + request JSONB NOT NULL DEFAULT '{}'::jsonb, + result JSONB NOT NULL DEFAULT '{}'::jsonb, + + error_kind TEXT, + error_message TEXT, + upload_ids TEXT[] NOT NULL DEFAULT '{}'::TEXT[], + + cost_units DOUBLE PRECISION, + latency_ms BIGINT, + + started_at TIMESTAMPTZ, + finished_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL, + updated_at TIMESTAMPTZ NOT NULL, + + CONSTRAINT signage_ai_jobs_kind_check + CHECK (kind IN ('GENERATE', 'EDIT')), + CONSTRAINT signage_ai_jobs_state_check + CHECK (state IN ('QUEUED', 'RUNNING', 'DONE', 'FAILED', 'CANCELLED')), + CONSTRAINT signage_ai_jobs_request_check + CHECK (jsonb_typeof(request) = 'object'), + CONSTRAINT signage_ai_jobs_result_check + CHECK (jsonb_typeof(result) = 'object'), + FOREIGN KEY (authority_id) REFERENCES authority(id) ON DELETE CASCADE, + FOREIGN KEY (provider_id) REFERENCES signage_ai_providers(id) ON DELETE SET NULL, + FOREIGN KEY (user_id) REFERENCES "user"(id) ON DELETE SET NULL, + FOREIGN KEY (parent_job_id) REFERENCES signage_ai_jobs(id) ON DELETE SET NULL +); + +CREATE INDEX IF NOT EXISTS signage_ai_jobs_authority_id_index + ON "signage_ai_jobs" USING BTREE (authority_id); +CREATE INDEX IF NOT EXISTS signage_ai_jobs_state_index + ON "signage_ai_jobs" USING BTREE (state); +CREATE INDEX IF NOT EXISTS signage_ai_jobs_created_at_index + ON "signage_ai_jobs" USING BTREE (created_at); +CREATE INDEX IF NOT EXISTS signage_ai_jobs_user_created_index + ON "signage_ai_jobs" USING BTREE (user_id, created_at DESC); +CREATE INDEX IF NOT EXISTS signage_ai_jobs_parent_job_id_index + ON "signage_ai_jobs" USING BTREE (parent_job_id); + +-- a repeated submission with the same key returns the job that already exists +CREATE UNIQUE INDEX IF NOT EXISTS signage_ai_jobs_idempotency_idx + ON "signage_ai_jobs" (user_id, idempotency_key) + WHERE idempotency_key IS NOT NULL; + +-- +micrate Down +-- SQL section 'Down' is executed when this migration is rolled back + +DROP TABLE IF EXISTS "signage_ai_jobs"; + +DROP TRIGGER IF EXISTS trg_signage_ai_providers_single_default ON signage_ai_providers; +DROP FUNCTION IF EXISTS signage_ai_providers_ensure_single_default() CASCADE; +DROP TABLE IF EXISTS "signage_ai_providers"; diff --git a/spec/generator.cr b/spec/generator.cr index 734a1b50..8dc042fc 100644 --- a/spec/generator.cr +++ b/spec/generator.cr @@ -742,6 +742,61 @@ module PlaceOS::Model ) end + def self.signage_ai_provider( + authority : Authority? = nil, + provider : SignageAIProvider::Provider = SignageAIProvider::Provider::OPENAI, + name : String? = nil, + credentials : String = %({"api_key":"test-key"}), + is_default : Bool = false, + enabled : Bool = true, + quotas : Hash(String, JSON::Any) = {} of String => JSON::Any, + ) + row = SignageAIProvider.new + row.name = name || Faker::Hacker.noun + row.provider = provider + row.authority_id = authority.try(&.id.as(String)) + row.credentials = credentials + row.default_model = provider.google? ? "gemini-3.1-flash-image" : "gpt-image-2" + row.is_default = is_default + row.enabled = enabled + row.quotas = JSON::Any.new(quotas) + row + end + + def self.signage_ai_job( + authority : Authority? = nil, + user : User? = nil, + provider : SignageAIProvider? = nil, + kind : SignageAIJob::Kind = SignageAIJob::Kind::Generate, + state : SignageAIJob::State = SignageAIJob::State::Queued, + candidates : Int32 = 2, + parent_job_id : UUID? = nil, + ) + unless authority + existing = Authority.find_by_domain("localhost") + authority = existing || self.authority.save! + end + user ||= self.user(authority: authority).save! + + job = SignageAIJob.new + job.authority_id = authority.id.as(String) + job.user_id = user.id + job.user_email = user.email.to_s + job.user_name = user.name + job.kind = kind + job.state = state + job.candidates = candidates + job.parent_job_id = parent_job_id + job.prompt = "a poster for the office party" + job.result = JSON::Any.new({"images" => JSON::Any.new(Array(JSON::Any).new(candidates) { JSON::Any.new(nil) })}) + if provider + job.provider_id = provider.id.as(UUID) + job.provider_type = provider.provider.to_s + job.model = provider.default_model + end + job + end + def self.storage(type = Storage::Type::S3, bucket : String? = nil, authority_id : String? = nil) Storage.new(storage_type: type, bucket_name: bucket || Faker::Hacker.noun, access_key: Faker::Hacker.noun, access_secret: Faker::Hacker.noun, diff --git a/spec/helper.cr b/spec/helper.cr index c50b8cbd..def4d5e6 100644 --- a/spec/helper.cr +++ b/spec/helper.cr @@ -26,6 +26,8 @@ Spec.after_suite do # Models that inherit directly from ::PgORM::Base (not ModelBase) — # cleared in dependency order (children first so FKs don't fire). [ + PlaceOS::Model::SignageAIJob, + PlaceOS::Model::SignageAIProvider, PlaceOS::Model::GroupSignageTemplate, PlaceOS::Model::SignageTemplate::SystemTemplate, PlaceOS::Model::SignageTemplate, diff --git a/spec/signage_ai_job_spec.cr b/spec/signage_ai_job_spec.cr new file mode 100644 index 00000000..188f16f1 --- /dev/null +++ b/spec/signage_ai_job_spec.cr @@ -0,0 +1,231 @@ +require "./helper" +require "uuid" + +module PlaceOS::Model + # created_at is written by the Timestamps module on insert, so a row can only + # be aged by moving it afterwards. + def self.backdate_job(job : SignageAIJob, age : Time::Span) + ::PgORM::Database.connection do |db| + db.exec("UPDATE signage_ai_jobs SET created_at = $1 WHERE id = $2", Time.utc - age, job.id.as(UUID)) + end + end + + describe SignageAIJob do + Spec.before_each do + SignageAIJob.clear + SignageAIProvider.clear + end + + test_round_trip(SignageAIJob) + + it "writes candidates concurrently without losing entries" do + job = Generator.signage_ai_job(candidates: 4).save! + id = job.id.as(UUID) + + done = Channel(Nil).new + 4.times do |index| + spawn do + SignageAIJob.bump_image(id, index, {"state" => JSON::Any.new("done"), "upload_id" => JSON::Any.new("uploads-#{index}")}) + done.send(nil) + end + end + 4.times { done.receive } + + found = SignageAIJob.find!(id) + found.images.size.should eq 4 + found.images.each_with_index do |image, index| + image["upload_id"].as_s.should eq "uploads-#{index}" + end + found.images_produced.should eq 4 + found.version.should eq 4 + end + + it "builds the images array when result has no images key" do + job = Generator.signage_ai_job(candidates: 1).save! + job.result = JSON::Any.new({} of String => JSON::Any) + job.save! + + SignageAIJob.bump_image(job.id.as(UUID), 0, {"state" => JSON::Any.new("done")}) + SignageAIJob.find!(job.id.as(UUID)).images.size.should eq 1 + end + + it "bumps the version on its own" do + job = Generator.signage_ai_job.save! + SignageAIJob.bump_version(job.id.as(UUID)).should eq 1 + SignageAIJob.find!(job.id.as(UUID)).version.should eq 1 + end + + it "sums candidates for quotas, counting failed jobs too" do + authority = Generator.localhost_authority + user = Generator.user(authority: authority).save! + Generator.signage_ai_job(authority: authority, user: user, candidates: 3).save! + Generator.signage_ai_job(authority: authority, user: user, candidates: 2, state: SignageAIJob::State::Done).save! + Generator.signage_ai_job(authority: authority, user: user, candidates: 9, state: SignageAIJob::State::Failed).save! + + # a failed job usually still reached the vendor and was billed + since = 1.day.ago + SignageAIJob.sum_candidates(user.id.as(String), since).should eq 14 + SignageAIJob.sum_candidates_for_authority(authority.id.as(String), since).should eq 14 + SignageAIJob.sum_candidates(user.id.as(String), 1.minute.from_now).should eq 0 + end + + it "walks a refine chain oldest first" do + first = Generator.signage_ai_job.save! + second = Generator.signage_ai_job(parent_job_id: first.id.as(UUID)).save! + third = Generator.signage_ai_job(parent_job_id: second.id.as(UUID)).save! + + third.chain.map(&.id).should eq [first.id, second.id] + first.chain.should be_empty + end + + it "reports final states" do + SignageAIJob::State::Queued.final?.should be_false + SignageAIJob::State::Running.final?.should be_false + SignageAIJob::State::Done.final?.should be_true + SignageAIJob::State::Failed.final?.should be_true + SignageAIJob::State::Cancelled.final?.should be_true + end + + it "finds jobs left running by a replica that went away" do + job = Generator.signage_ai_job(state: SignageAIJob::State::Running).save! + job.started_at = 30.minutes.ago + job.save! + + SignageAIJob.stale(10.minutes.ago).map(&.id).should eq [job.id] + SignageAIJob.stale(1.hour.ago).should be_empty + end + + it "groups usage by provider and model" do + authority = Generator.localhost_authority + provider = Generator.signage_ai_provider(authority: authority).save! + 2.times do + job = Generator.signage_ai_job(authority: authority, provider: provider, candidates: 2).save! + job.images_produced = 2 + job.cost_units = 0.5 + job.save! + end + + rows = SignageAIJob.usage(authority.id.as(String), 1.day.ago, 1.day.from_now) + rows.size.should eq 1 + row = rows.first + row.provider.should eq "OPENAI" + row.model.should eq "gpt-image-2" + row.jobs.should eq 2 + row.candidates.should eq 4 + row.images_produced.should eq 4 + row.cost_units.should eq 1.0 + end + + describe "automatic cleanup" do + # cleanup ships off, so every example that exercises it has to turn it on + before_each { SignageAIJob.cleanup_after = 90.days } + after_each { SignageAIJob.cleanup_after = SignageAIJob::CLEANUP_AFTER } + + it "does nothing unless an operator configures a window" do + SignageAIJob::CLEANUP_AFTER.should eq Time::Span.zero + + SignageAIJob.cleanup_after = SignageAIJob::CLEANUP_AFTER + job = Generator.signage_ai_job.save! + PlaceOS::Model.backdate_job(job, 1825.days) + + Generator.signage_ai_job.save! + + SignageAIJob.find?(job.id.as(UUID)).should_not be_nil + end + + it "removes jobs past the window as a new job is created" do + old_job = Generator.signage_ai_job.save! + recent = Generator.signage_ai_job.save! + PlaceOS::Model.backdate_job(old_job, 100.days) + + # nothing has run yet: the row is old but no insert has happened since + SignageAIJob.find?(old_job.id.as(UUID)).should_not be_nil + + # the insert itself is what triggers the sweep + fresh = Generator.signage_ai_job.save! + + SignageAIJob.find?(old_job.id.as(UUID)).should be_nil + SignageAIJob.find?(recent.id.as(UUID)).should_not be_nil + SignageAIJob.find?(fresh.id.as(UUID)).should_not be_nil + end + + it "keeps jobs inside the window" do + job = Generator.signage_ai_job.save! + PlaceOS::Model.backdate_job(job, 89.days) + + Generator.signage_ai_job.save! + + SignageAIJob.find?(job.id.as(UUID)).should_not be_nil + end + + it "keeps everything when the window is explicitly zero" do + SignageAIJob.cleanup_after = Time::Span.zero + job = Generator.signage_ai_job.save! + PlaceOS::Model.backdate_job(job, 1825.days) + + Generator.signage_ai_job.save! + + SignageAIJob.find?(job.id.as(UUID)).should_not be_nil + SignageAIJob.cleanup_expired.should eq 0 + end + + it "reads a negative window as keep everything, not delete everything" do + SignageAIJob.cleanup_after = -30.days + job = Generator.signage_ai_job.save! + PlaceOS::Model.backdate_job(job, 1825.days) + + SignageAIJob.cleanup_expired.should eq 0 + SignageAIJob.find?(job.id.as(UUID)).should_not be_nil + end + + it "bounds how much one insert removes" do + # create them all first: each insert sweeps, so backdating as we go would + # let the second insert remove the first row before the third exists + jobs = Array.new(3) { Generator.signage_ai_job.save! } + jobs.each { |job| PlaceOS::Model.backdate_job(job, 100.days) } + + SignageAIJob.count.should eq 3 + {SignageAIJob.cleanup_expired(limit: 2), SignageAIJob.count}.should eq({2, 1}) + {SignageAIJob.cleanup_expired(limit: 2), SignageAIJob.count}.should eq({1, 0}) + SignageAIJob.cleanup_expired(limit: 2).should eq 0 + end + + it "removes the oldest first" do + oldest = Generator.signage_ai_job.save! + newest = Generator.signage_ai_job.save! + PlaceOS::Model.backdate_job(oldest, 200.days) + PlaceOS::Model.backdate_job(newest, 100.days) + + SignageAIJob.count.should eq 2 + {SignageAIJob.cleanup_expired(limit: 1), SignageAIJob.count}.should eq({1, 1}) + SignageAIJob.find?(oldest.id.as(UUID)).should be_nil + SignageAIJob.find?(newest.id.as(UUID)).should_not be_nil + end + + it "keeps a recent refine when its parent ages out" do + parent = Generator.signage_ai_job.save! + child = Generator.signage_ai_job(parent_job_id: parent.id.as(UUID)).save! + PlaceOS::Model.backdate_job(parent, 100.days) + + SignageAIJob.cleanup_expired.should eq 1 + + found = SignageAIJob.find?(child.id.as(UUID)) + found.should_not be_nil + found.as(SignageAIJob).parent_job_id.should be_nil + end + end + + it "clears the provider link but keeps the job when a provider is deleted" do + authority = Generator.localhost_authority + provider = Generator.signage_ai_provider(authority: authority).save! + job = Generator.signage_ai_job(authority: authority, provider: provider).save! + + provider.destroy + + found = SignageAIJob.find!(job.id.as(UUID)) + found.provider_id.should be_nil + found.provider_type.should eq "OPENAI" + found.model.should eq "gpt-image-2" + end + end +end diff --git a/spec/signage_ai_provider_spec.cr b/spec/signage_ai_provider_spec.cr new file mode 100644 index 00000000..13460a91 --- /dev/null +++ b/spec/signage_ai_provider_spec.cr @@ -0,0 +1,81 @@ +require "./helper" +require "uuid" + +module PlaceOS::Model + describe SignageAIProvider do + Spec.before_each do + SignageAIJob.clear + SignageAIProvider.clear + end + + test_round_trip(SignageAIProvider) + + it "encrypts credentials and keeps them out of as_json" do + authority = Generator.localhost_authority + secret = %({"api_key":"sk-not-a-real-key"}) + row = Generator.signage_ai_provider(authority: authority, credentials: secret).save! + + row.credentials_encrypted?.should be_true + row.credentials.should_not contain "sk-not-a-real-key" + + found = SignageAIProvider.find!(row.id.as(UUID)) + found.decrypt_credentials.should eq secret + found.credentials_json["api_key"].as_s.should eq "sk-not-a-real-key" + + rendered = found.as_json.to_json + rendered.should_not contain "sk-not-a-real-key" + rendered.should_not contain "credentials" + rendered.should contain found.name + end + + it "does not double encrypt on repeated saves" do + row = Generator.signage_ai_provider(credentials: %({"api_key":"one"})).save! + ciphertext = row.credentials + row.name = "renamed" + row.save! + row.credentials.should eq ciphertext + SignageAIProvider.find!(row.id.as(UUID)).credentials_json["api_key"].as_s.should eq "one" + end + + it "keeps one default per authority and leaves other authorities alone" do + authority = Generator.localhost_authority + other = Generator.authority(domain: "http://default-#{UUID.random}.test").save! + + first = Generator.signage_ai_provider(authority: authority, is_default: true).save! + elsewhere = Generator.signage_ai_provider(authority: other, is_default: true).save! + second = Generator.signage_ai_provider(authority: authority, is_default: true).save! + + SignageAIProvider.find!(first.id.as(UUID)).is_default.should be_false + SignageAIProvider.find!(second.id.as(UUID)).is_default.should be_true + SignageAIProvider.find!(elsewhere.id.as(UUID)).is_default.should be_true + end + + it "falls back to the shared row when a domain has none" do + authority = Generator.localhost_authority + shared = Generator.signage_ai_provider(authority: nil, is_default: true).save! + + SignageAIProvider.default_for(authority.id.as(String)).try(&.id).should eq shared.id + + own = Generator.signage_ai_provider(authority: authority, is_default: true).save! + SignageAIProvider.default_for(authority.id.as(String)).try(&.id).should eq own.id + + available = SignageAIProvider.available_for(authority.id.as(String)).map(&.id) + available.should contain own.id + available.should contain shared.id + end + + it "ignores disabled rows when picking a default" do + authority = Generator.localhost_authority + Generator.signage_ai_provider(authority: authority, is_default: true, enabled: false).save! + shared = Generator.signage_ai_provider(authority: nil).save! + + SignageAIProvider.default_for(authority.id.as(String)).try(&.id).should eq shared.id + end + + it "reads quotas" do + row = Generator.signage_ai_provider(quotas: {"user_per_day" => JSON::Any.new(12_i64)}).save! + SignageAIProvider.find!(row.id.as(UUID)).quota("user_per_day").should eq 12 + row.quota("domain_per_month").should be_nil + end + end +end diff --git a/src/placeos-models/signage_ai_job.cr b/src/placeos-models/signage_ai_job.cr new file mode 100644 index 00000000..91c14851 --- /dev/null +++ b/src/placeos-models/signage_ai_job.cr @@ -0,0 +1,353 @@ +require "json" +require "uuid" +require "uuid/json" + +require "./base/model" +require "./authority" +require "./signage_ai_provider" +require "./user" + +module PlaceOS::Model + # One image generation or edit request. + # + # Candidates are written into `result["images"]` one at a time by the runner as + # each vendor call finishes. Every write bumps `version`, which is what the long + # polling endpoint compares against so a client can tell what it has already seen + # (`updated_at` only has whole second resolution and cannot order two writes in + # the same second). + class SignageAIJob < ::PgORM::Base + include PgORM::Timestamps + include Neuroplastic + + Log = ::Log.for(self) + + table :signage_ai_jobs + + # How long a job row is kept, once cleanup is turned on. + # + # Off unless SIGNAGE_AI_CLEANUP_DAYS is set to a positive number of days, so + # nothing starts deleting a site's history because it was upgraded. Somewhere + # that generates seasonally keeps last year's work by doing nothing at all. + # Zero and negatives both read as "keep everything" rather than as a cutoff in + # the future that would take the whole table with it. + # + # Note for operators turning it on: the usage report and the per domain monthly + # quota both read this table, so a window shorter than a month makes both + # under-count. + CLEANUP_AFTER = ((ENV["SIGNAGE_AI_CLEANUP_DAYS"]?.try(&.to_i?) || 0).days) + + # The window actually used, seeded from the constant above. Settable so the + # specs can exercise a configured window, which a constant read at compile + # time cannot express. + class_property cleanup_after : Time::Span = CLEANUP_AFTER + + # Rows removed per insert. The work each job creation does is bounded, so a + # long disabled cleanup being switched back on cannot stall a request while + # it catches up: it takes several inserts instead of one long one. + CLEANUP_BATCH = 1000 + + default_primary_key id : UUID, autogenerated: true + + enum Kind + Generate + Edit + + def to_s : String + super.downcase + end + end + + enum State + Queued + Running + Done + Failed + Cancelled + + def to_s : String + super.downcase + end + + def final? + self == Done || self == Failed || self == Cancelled + end + end + + attribute authority_id : String, es_type: "keyword" + belongs_to :authority, class_name: PlaceOS::Model::Authority + + attribute provider_id : UUID? = nil, es_type: "keyword" + belongs_to :provider, class_name: PlaceOS::Model::SignageAIProvider, foreign_key: provider_id + + # recorded as strings as well so a job survives its provider row being deleted + # (named provider_type because `provider` is the belongs_to association above) + attribute provider_type : String? = nil, es_type: "keyword" + attribute model : String? = nil, es_type: "keyword" + + attribute user_id : String? = nil, es_type: "keyword" + belongs_to :user, class_name: PlaceOS::Model::User, foreign_key: user_id + attribute user_email : String? = nil, es_type: "keyword" + attribute user_name : String? = nil, sanitize: :text + + # a refine points at the job it was refined from + attribute parent_job_id : UUID? = nil, es_type: "keyword" + + attribute version : Int32 = 0 + attribute kind : Kind = Kind::Generate, converter: PlaceOS::Model::PGEnumConverter(PlaceOS::Model::SignageAIJob::Kind), es_type: "keyword" + attribute state : State = State::Queued, converter: PlaceOS::Model::PGEnumConverter(PlaceOS::Model::SignageAIJob::State), es_type: "keyword" + attribute cancel_requested : Bool = false + attribute idempotency_key : String? = nil, es_type: "keyword" + + attribute candidates : Int32 = 1 + attribute images_produced : Int32 = 0 + + attribute prompt : String? = nil, es_ignore: true + attribute request : JSON::Any = JSON::Any.new({} of String => JSON::Any), es_ignore: true + attribute result : JSON::Any = JSON::Any.new({} of String => JSON::Any), es_ignore: true + + attribute error_kind : String? = nil, es_type: "keyword" + attribute error_message : String? = nil, sanitize: :text + attribute upload_ids : Array(String) = [] of String, es_type: "keyword" + + attribute cost_units : Float64? = nil + attribute latency_ms : Int64? = nil + + attribute started_at : Time? = nil + attribute finished_at : Time? = nil + + validates :authority_id, presence: true + + before_create :cleanup_expired_jobs + + def final? : Bool + state.final? + end + + # The candidate slots, in order. Entries are `null` until that candidate lands. + def images : Array(JSON::Any) + result["images"]?.try(&.as_a?) || [] of JSON::Any + end + + # Walk back up a refine chain, oldest first. Bounded so a cycle cannot hang a request. + def chain(limit : Int32 = 25) : Array(SignageAIJob) + seen = [] of SignageAIJob + current = self + limit.times do + parent_id = current.parent_job_id + break unless parent_id + parent = SignageAIJob.find?(parent_id) + break unless parent + seen.unshift(parent) + current = parent + end + seen + end + + # Write one candidate into result->'images'->index and bump the version, in a + # single statement. Candidate fibers must never load, modify and save the row: + # `save` writes the whole column and would drop entries written concurrently. + def self.bump_image(id : UUID, index : Int32, image : JSON::Any | Hash(String, JSON::Any)) : Int32? + payload = image.to_json + ::PgORM::Database.connection do |db| + db.query_one?(<<-SQL, id, index.to_s, payload, as: Int32) + UPDATE signage_ai_jobs + SET + result = jsonb_set( + CASE + WHEN jsonb_typeof(result -> 'images') = 'array' THEN result + ELSE jsonb_set(result, '{images}', '[]'::jsonb, true) + END, + ARRAY['images', $2], + $3::jsonb, + true + ), + images_produced = images_produced + 1, + version = version + 1, + updated_at = now() + WHERE id = $1 + RETURNING version; + SQL + end + end + + # Write one candidate's entry without counting it as newly produced. + # + # `bump_image` has the increment welded into its statement because it is the + # runner recording a candidate that has just landed. Claiming records that + # an image the runner already counted became a media item, so it must not + # count again: it did, and the usage report was inflated by every save. + def self.attach_item(id : UUID, index : Int32, image : JSON::Any | Hash(String, JSON::Any)) : Int32? + payload = image.to_json + ::PgORM::Database.connection do |db| + db.query_one?(<<-SQL, id, index.to_s, payload, as: Int32) + UPDATE signage_ai_jobs + SET + result = jsonb_set( + CASE + WHEN jsonb_typeof(result -> 'images') = 'array' THEN result + ELSE jsonb_set(result, '{images}', '[]'::jsonb, true) + END, + ARRAY['images', $2], + $3::jsonb, + true + ), + version = version + 1, + updated_at = now() + WHERE id = $1 + RETURNING version; + SQL + end + end + + # Bump the version without touching the images, so a state change wakes a + # long polling client too. + def self.bump_version(id : UUID) : Int32? + ::PgORM::Database.connection do |db| + db.query_one?(<<-SQL, id, as: Int32) + UPDATE signage_ai_jobs + SET version = version + 1, updated_at = now() + WHERE id = $1 + RETURNING version; + SQL + end + end + + # Candidates requested by a user since a point in time, for the per user quota. + # + # Failed jobs count. Most of them reached the vendor and were billed, and + # exempting them meant a caller whose requests kept failing had no limit at + # all, which is the one case where a limit matters most. + def self.sum_candidates(user_id : String, since : Time) : Int32 + ::PgORM::Database.connection do |db| + db.query_one(<<-SQL, user_id, since, as: Int64) + SELECT COALESCE(SUM(candidates), 0)::bigint + FROM signage_ai_jobs + WHERE user_id = $1 + AND created_at >= $2 + SQL + end.to_i32 + end + + # Candidates requested across a domain since a point in time, for the domain quota. + def self.sum_candidates_for_authority(authority_id : String, since : Time) : Int32 + ::PgORM::Database.connection do |db| + db.query_one(<<-SQL, authority_id, since, as: Int64) + SELECT COALESCE(SUM(candidates), 0)::bigint + FROM signage_ai_jobs + WHERE authority_id = $1 + AND created_at >= $2 + SQL + end.to_i32 + end + + # Housekeeping, run as each new job is created. Does nothing unless an + # operator has set a retention window. + # + # Deliberately never raises. `before_create` runs inside the transaction that + # inserts the new row, so an error escaping here would abort that transaction + # and the caller would fail to create a job because tidying up went wrong. + protected def cleanup_expired_jobs + self.class.cleanup_expired + nil + end + + # Remove job rows past the retention window. Returns how many went. + # + # Two statements rather than one `DELETE ... WHERE id IN (SELECT ... LIMIT n + # FOR UPDATE SKIP LOCKED)`. That single statement does not reliably delete at + # most n rows. When the planner picks a nested loop semi join it re-runs the + # subquery once per candidate row, and because SKIP LOCKED and the rows already + # deleted by this statement change what comes back each time, the delete lands + # on the union of those evaluations. Confirmed here with auto_explain: a hash + # semi join runs the Limit once and removes 2 of 3, a nested loop runs it three + # times and removes all 3. Which plan you get moves with the table statistics, + # so it is fine until it suddenly is not. + # + # Reading the ids first and deleting exactly those cannot be re-evaluated. + # + # The delete runs in a nested transaction, which PgORM issues as a SAVEPOINT, + # so a failure rolls back to the savepoint and leaves the surrounding insert + # intact. Without it a failed delete would poison the whole transaction and + # take the new job down with it. + def self.cleanup_expired(limit : Int32 = CLEANUP_BATCH) : Int32 + window = cleanup_after + return 0 if window <= Time::Span.zero + return 0 if limit <= 0 + + cutoff = Time.utc - window + deleted = 0 + + begin + ::PgORM::Database.transaction do |tx| + ids = [] of UUID + tx.connection.query(<<-SQL, cutoff, limit) do |rs| + SELECT id + FROM signage_ai_jobs + WHERE created_at < $1 + ORDER BY created_at + LIMIT $2 + SQL + rs.each { ids << rs.read(UUID) } + end + + next if ids.empty? + + # one placeholder per id: an array cannot bind as a single parameter + placeholders = Array.new(ids.size) { |index| "$#{index + 1}" }.join(", ") + args = ids.map(&.as(::PgORM::Value)) + result = tx.connection.exec("DELETE FROM signage_ai_jobs WHERE id IN (#{placeholders})", args: args) + deleted = result.rows_affected.to_i + end + rescue ex + # housekeeping is never worth failing a user's request over + Log.warn(exception: ex) { {message: "signage AI job cleanup failed", cutoff: cutoff.to_s} } + return 0 + end + + Log.info { {message: "signage AI jobs cleaned up", removed: deleted, cutoff: cutoff.to_s} } if deleted > 0 + deleted + end + + # Jobs left running by a replica that went away. + def self.stale(older_than : Time) : Array(SignageAIJob) + where("state = 'RUNNING' AND started_at < ?", older_than).to_a + end + + # Usage per provider and model, for the Backoffice usage tab. + record UsageRow, provider : String, model : String, jobs : Int64, candidates : Int64, + images_produced : Int64, cost_units : Float64 do + include JSON::Serializable + end + + def self.usage(authority_id : String, from : Time, to : Time) : Array(UsageRow) + rows = [] of UsageRow + ::PgORM::Database.connection do |db| + db.query(<<-SQL, authority_id, from, to) do |rs| + SELECT + COALESCE(provider_type, 'unknown') AS provider, + COALESCE(model, 'unknown') AS model, + COUNT(*)::bigint AS jobs, + COALESCE(SUM(candidates), 0)::bigint AS candidates, + COALESCE(SUM(images_produced), 0)::bigint AS images_produced, + COALESCE(SUM(cost_units), 0)::double precision AS cost_units + FROM signage_ai_jobs + WHERE authority_id = $1 AND created_at >= $2 AND created_at < $3 + GROUP BY 1, 2 + ORDER BY 1, 2 + SQL + rs.each do + rows << UsageRow.new( + provider: rs.read(String), + model: rs.read(String), + jobs: rs.read(Int64), + candidates: rs.read(Int64), + images_produced: rs.read(Int64), + cost_units: rs.read(Float64), + ) + end + end + end + rows + end + end +end diff --git a/src/placeos-models/signage_ai_provider.cr b/src/placeos-models/signage_ai_provider.cr new file mode 100644 index 00000000..b6a06f6d --- /dev/null +++ b/src/placeos-models/signage_ai_provider.cr @@ -0,0 +1,140 @@ +require "json" +require "uuid" +require "uuid/json" + +require "./base/model" +require "./authority" +require "./utilities/encryption" + +module PlaceOS::Model + # Vendor credentials used to generate or edit signage artwork. + # + # A row with a `nil` authority_id is the shared fallback, the same arrangement + # `Storage` uses. Credentials are a JSON object whose shape depends on the + # provider, encrypted at rest and never rendered by the API. + class SignageAIProvider < ::PgORM::Base + include PgORM::Timestamps + include Neuroplastic + + Log = ::Log.for(self) + + table :signage_ai_providers + + default_primary_key id : UUID, autogenerated: true + + enum Provider + OPENAI + AZURE_OPENAI + GOOGLE_VERTEX + + def google? + self == GOOGLE_VERTEX + end + + # OpenAI and Azure OpenAI speak the same request and response shape + def openai_compatible? + self == OPENAI || self == AZURE_OPENAI + end + end + + attribute name : String, sanitize: :text, es_subfield: "keyword" + attribute provider : Provider = Provider::OPENAI, converter: PlaceOS::Model::PGEnumConverter(PlaceOS::Model::SignageAIProvider::Provider), es_type: "keyword" + + attribute authority_id : String? = nil, es_type: "keyword" + belongs_to :authority, class_name: PlaceOS::Model::Authority + + # encrypted JSON object, see `credentials_json` + attribute credentials : String, mass_assignment: false, es_ignore: true + + # optional overrides. `endpoint` also lets a deployment point at a gateway or, + # locally, at a stub vendor. + attribute endpoint : String? = nil, sanitize: :text + attribute location : String? = nil, sanitize: :text + attribute default_model : String? = nil, sanitize: :text + attribute allowed_models : Array(String) = [] of String, es_type: "keyword" + + attribute enabled : Bool = true + attribute is_default : Bool = false + + # {"user_per_day": 60, "domain_per_month": 2000} + attribute quotas : JSON::Any = JSON::Any.new({} of String => JSON::Any), es_ignore: true + + validates :name, presence: true + validates :credentials, presence: true + + before_save do + self.credentials = PlaceOS::Encryption.encrypt(credentials, level: level, id: encryption_id) + end + + # The provider used for a domain: its own default, then its own first enabled + # row, then the shared fallback. + def self.default_for(authority_id : String?) : SignageAIProvider? + if authority_id + own = where(authority_id: authority_id, enabled: true).order(is_default: :desc, created_at: :asc).first? + return own if own + end + where(authority_id: nil, enabled: true).order(is_default: :desc, created_at: :asc).first? + end + + # Rows a domain may use: its own, plus the shared fallback. + def self.available_for(authority_id : String?) : Array(SignageAIProvider) + rows = where(authority_id: nil, enabled: true).to_a + rows = where(authority_id: authority_id, enabled: true).to_a + rows if authority_id + rows + end + + def decrypt_credentials : String + PlaceOS::Encryption.decrypt(credentials, level: level, id: encryption_id) + end + + # Parsed credentials. Raises `Model::Error` when the stored value is not an object, + # which is what a wrong `PLACE_SERVER_SECRET` looks like from here. + def credentials_json : Hash(String, JSON::Any) + parsed = JSON.parse(decrypt_credentials) + parsed.as_h + rescue JSON::ParseException | TypeCastError + raise Model::Error.new("credentials for signage AI provider #{id} could not be read") + end + + def credentials_json=(value : Hash(String, JSON::Any) | JSON::Any) + self.credentials = value.to_json + end + + def credentials_encrypted? : Bool + PlaceOS::Encryption.is_encrypted?(credentials) + end + + def quota(key : String) : Int32? + quotas[key]?.try(&.as_i?) + end + + # Rendered by the API in place of `to_json`: the same row without credentials. + def as_json + { + id: id.as(UUID).to_s, + name: name, + provider: provider.to_s, + authority_id: authority_id, + endpoint: endpoint, + location: location, + default_model: default_model, + allowed_models: allowed_models, + enabled: enabled, + is_default: is_default, + quotas: quotas, + created_at: created_at.try(&.to_unix), + updated_at: updated_at.try(&.to_unix), + } + end + + private def level : PlaceOS::Encryption::Level + PlaceOS::Encryption::Level::NeverDisplay + end + + # Keyed on the domain rather than the row id: `id` is generated by the database + # so it is not available while `before_save` runs on a new row. + private def encryption_id : String + authority_id || "shared" + end + end +end