-- broker.broker_jobs -- Job queue for broker execution. -- tenant_id / RLS: rows are only visible/writable when tenant_id matches -- current_setting('broker.tenant_id', true) for the current transaction. -- Callers must invoke broker.broker_set_tenant(...) before enqueue/claim; -- if they don't, tenant_id defaults to 'default' and current_setting -- also defaults to NULL -> broker_add_job coalesces to 'default' so -- single-tenant use keeps working unmodified. CREATE TABLE IF NOT EXISTS broker.broker_jobs ( id_broker_jobs BIGSERIAL PRIMARY KEY, job_name VARCHAR(255) NOT NULL, job_priority INTEGER NOT NULL DEFAULT 0, job_queue INTEGER NOT NULL DEFAULT 1, job_language VARCHAR(50) NOT NULL DEFAULT 'sql', execute_str TEXT NOT NULL, execute_result TEXT, error_msg TEXT, complete_status INTEGER NOT NULL DEFAULT 0, run_as VARCHAR(100), rid_broker_schedule BIGINT, rid_broker_queueinstance BIGINT, tenant_id TEXT NOT NULL DEFAULT 'default', -- Lease / retry / idempotency attempt_count INTEGER NOT NULL DEFAULT 0, max_attempts INTEGER NOT NULL DEFAULT 1, available_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(), leased_at TIMESTAMP WITH TIME ZONE, lease_expires_at TIMESTAMP WITH TIME ZONE, lease_token UUID, idempotency_key TEXT, created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(), updated_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(), started_at TIMESTAMP WITH TIME ZONE, completed_at TIMESTAMP WITH TIME ZONE, CONSTRAINT broker_jobs_complete_status_check CHECK (complete_status IN (0, 1, 2, 3, 4)), CONSTRAINT broker_jobs_job_queue_check CHECK (job_queue > 0), CONSTRAINT fk_schedule FOREIGN KEY (rid_broker_schedule) REFERENCES broker.broker_schedule(id_broker_schedule) ON DELETE SET NULL, CONSTRAINT fk_instance FOREIGN KEY (rid_broker_queueinstance) REFERENCES broker.broker_queueinstance(id_broker_queueinstance) ON DELETE SET NULL ); -- General-purpose indexes CREATE INDEX IF NOT EXISTS idx_broker_jobs_status ON broker.broker_jobs(complete_status); CREATE INDEX IF NOT EXISTS idx_broker_jobs_schedule ON broker.broker_jobs(rid_broker_schedule); CREATE INDEX IF NOT EXISTS idx_broker_jobs_instance ON broker.broker_jobs(rid_broker_queueinstance); CREATE INDEX IF NOT EXISTS idx_broker_jobs_created ON broker.broker_jobs(created_at); CREATE INDEX IF NOT EXISTS idx_broker_jobs_name ON broker.broker_jobs(job_name, complete_status); CREATE INDEX IF NOT EXISTS idx_broker_jobs_tenant ON broker.broker_jobs(tenant_id); -- Claim index: exactly what broker_get's WHERE/ORDER BY needs, partial on pending rows only. CREATE INDEX IF NOT EXISTS idx_broker_jobs_claim ON broker.broker_jobs (job_queue, job_priority DESC, created_at, id_broker_jobs) WHERE complete_status = 0; -- Idempotency: at most one pending/any job per (queue, key) when a key is supplied. CREATE UNIQUE INDEX IF NOT EXISTS idx_broker_jobs_idempotency ON broker.broker_jobs (job_queue, idempotency_key) WHERE idempotency_key IS NOT NULL; COMMENT ON TABLE broker.broker_jobs IS 'Job queue for broker execution'; COMMENT ON COLUMN broker.broker_jobs.complete_status IS '0=pending, 1=running, 2=completed, 3=failed (terminal or dead-lettered once attempt_count>=max_attempts), 4=cancelled'; COMMENT ON COLUMN broker.broker_jobs.tenant_id IS 'RLS tenant scaffold; defaults to ''default'' for single-tenant use'; COMMENT ON COLUMN broker.broker_jobs.attempt_count IS 'Number of times this job has been claimed/executed'; COMMENT ON COLUMN broker.broker_jobs.max_attempts IS 'Job is dead-lettered (failed) once attempt_count reaches this value'; COMMENT ON COLUMN broker.broker_jobs.available_at IS 'Job is not claimable until now() >= available_at (used for retry backoff)'; COMMENT ON COLUMN broker.broker_jobs.lease_token IS 'Token handed out by broker_get; broker_run requires a matching token to execute, so an expired/reclaimed lease cannot be double-processed'; COMMENT ON COLUMN broker.broker_jobs.idempotency_key IS 'Optional caller-supplied key; unique per (job_queue, idempotency_key)'; CREATE OR REPLACE FUNCTION broker.tf_broker_jobs_update_timestamp() RETURNS TRIGGER AS $$ BEGIN NEW.updated_at = NOW(); RETURN NEW; END; $$ LANGUAGE plpgsql; DROP TRIGGER IF EXISTS t_broker_jobs_updated_at ON broker.broker_jobs; CREATE TRIGGER t_broker_jobs_updated_at BEFORE UPDATE ON broker.broker_jobs FOR EACH ROW EXECUTE FUNCTION broker.tf_broker_jobs_update_timestamp(); -- Row Level Security: tenant isolation ALTER TABLE broker.broker_jobs ENABLE ROW LEVEL SECURITY; ALTER TABLE broker.broker_jobs FORCE ROW LEVEL SECURITY; DROP POLICY IF EXISTS broker_jobs_tenant_isolation ON broker.broker_jobs; CREATE POLICY broker_jobs_tenant_isolation ON broker.broker_jobs USING (tenant_id = COALESCE(NULLIF(current_setting('broker.tenant_id', true), ''), 'default')) WITH CHECK (tenant_id = COALESCE(NULLIF(current_setting('broker.tenant_id', true), ''), 'default'));