The pgrx extension: SQL entry points gluing jev-protocol to
postjevsql-pg through their safe APIs (contract §0).
3#![forbid(unsafe_code)]
5use std::cell::{Cell, RefCell}; 6use std::collections::HashMap; 7use std::ffi::{CStr, CString}; 8use std::rc::Rc; 9use std::task::{Poll, Waker}; 10use std::time::Duration; 11 12use pgrx::datum::Internal; 13use pgrx::prelude::*; 14use pgrx::{AnyElement, GucContext, GucFlags, GucRegistry, GucSetting, PostgresGucEnum}; 15use jev_client::{Client, ClientError}; 16use jev_protocol::{ 17 Choice, ChoiceAnswer, Json, Key, ModelId, Noul, NoulAnswer, Question, Questions, Response, Score, ScoreAnswer, 18 question_bytes, 19}; 20use postjevsql_core::label_tree::{self, Branch, Describe, Outcome, Params, Search, Step, Tau, Tree}; 21use postjevsql_core::shortlist::{self, Shortlist}; 22use postjevsql_core::{CacheKey, KeyScope, Layout, PROMPT_VERSION, Placement, Planned, Planner, Sample, TokenRatio}; 23use postjevsql_pg::scan::{ 24 self, Arg, Call, Estimate, Field, Function, Judge, Judgement, Out, Property, Refusal, Returns, Shown, Statement, Value, 25}; 26use postjevsql_pg::owner::as_owner_of; 27use postjevsql_pg::stats::{self, Telemetry}; 28use postjevsql_pg::{Backend, Endpoint, HttpsTransport, KeepAlive, RateLimit, StringCheck, define_checked_string, join_all, stream_limit}; 29 30mod failure; 31use failure::Failure; 32 33pgrx::pg_module_magic!(name, version); 34 35static ENDPOINT: GucSetting<Option<CString>> = 36 GucSetting::<Option<CString>>::new(Some(c"https://api.typesafe.ai/v1/systemone")); 37static MODEL: GucSetting<Option<CString>> = GucSetting::<Option<CString>>::new(None); 38static CA_FILE: GucSetting<Option<CString>> = GucSetting::<Option<CString>>::new(None); 39static DNS_SERVERS: GucSetting<Option<CString>> = GucSetting::<Option<CString>>::new(None); 40static KEEPALIVE_INTERVAL: GucSetting<i32> = GucSetting::<i32>::new(20_000); 41static KEEPALIVE_TIMEOUT: GucSetting<i32> = GucSetting::<i32>::new(10_000); 42static API_KEY_FILE: GucSetting<Option<CString>> = GucSetting::<Option<CString>>::new(None); 43static CONCURRENCY: GucSetting<i32> = GucSetting::<i32>::new(100); 44static PRICE_PER_MTOK: GucSetting<f64> = GucSetting::<f64>::new(0.042); 45static MAX_TOKENS_PER_SECOND: GucSetting<i32> = GucSetting::<i32>::new(250_000); 46static MAX_REQUESTS_PER_MINUTE: GucSetting<i32> = GucSetting::<i32>::new(1_200); 47static MAX_ROWS: GucSetting<i32> = GucSetting::<i32>::new(-1); 48static MAX_COST: GucSetting<f64> = GucSetting::<f64>::new(-1.0); 49static CACHE_NAMESPACE: GucSetting<Option<CString>> = GucSetting::<Option<CString>>::new(None); 50static CACHE_ONLY: GucSetting<bool> = GucSetting::<bool>::new(false); 51static ON_ERROR: GucSetting<OnError> = GucSetting::<OnError>::new(OnError::Error); 52static THRESHOLD: GucSetting<f64> = GucSetting::<f64>::new(0.5);
What a failed judgment does (contract SQL surface).
The call is NULL: jev_prob NULL, and jev() false through
its qual, so a fault skips a row and never passes one.
66#[pg_guard] 67pub extern "C-unwind" fn _PG_init() { 68 // Superuser-only: whoever sets the endpoint receives the API key. 69 GucRegistry::define_string_guc( 70 c"jev.endpoint", 71 c"TypeSafe evaluation endpoint.", 72 c"Must be https. Requests carry the server's API key.", 73 &ENDPOINT, 74 GucContext::Suset, 75 GucFlags::default(), 76 ); 77 // Refused at SET, not first at query time: an alias moves under you. 78 define_checked_string::<PinnedModel>( 79 c"jev.model", 80 c"The pinned Jev model version, e.g. jev-1.13.0.", 81 c"Aliases such as jev-latest are refused. Answers from any other model are rejected.", 82 &MODEL, 83 GucContext::Userset, 84 GucFlags::default(), 85 ); 86 GucRegistry::define_string_guc( 87 c"jev.ca_file", 88 c"PEM file of extra CA certificates to trust.", 89 c"Added to the Mozilla roots; for tests and private endpoints.", 90 &CA_FILE, 91 GucContext::Suset, 92 GucFlags::default(), 93 ); 94 // Superuser-only: it chooses whose account pays. 95 GucRegistry::define_string_guc( 96 c"jev.api_key_file", 97 c"File holding the TypeSafe API key.", 98 c"Exclusive with TYPESAFE_API_KEY in the server environment.", 99 &API_KEY_FILE, 100 GucContext::Suset, 101 GucFlags::default(), 102 ); 103 // Superuser-only, like the endpoint: it decides where names point. 104 GucRegistry::define_string_guc( 105 c"jev.dns_servers", 106 c"DNS servers for jev.endpoint's host, as ip[:port], comma-separated.", 107 c"Empty uses /etc/resolv.conf.", 108 &DNS_SERVERS, 109 GucContext::Suset, 110 GucFlags::default(), 111 ); 112 // Superuser-only, like the endpoint: they shape the one connection 113 // every role's requests share. Tests shorten them. 114 GucRegistry::define_int_guc( 115 c"jev.keepalive_interval", 116 c"Quiet time on the connection before a keep-alive ping.", 117 c"Pings are sent only while a jev scan runs.", 118 &KEEPALIVE_INTERVAL, 119 1, 120 i32::MAX, 121 GucContext::Suset, 122 GucFlags::UNIT_MS, 123 ); 124 GucRegistry::define_int_guc( 125 c"jev.keepalive_timeout", 126 c"How long a keep-alive ping may go unanswered before the connection is closed.", 127 c"The next request opens a new connection.", 128 &KEEPALIVE_TIMEOUT, 129 1, 130 i32::MAX, 131 GucContext::Suset, 132 GucFlags::UNIT_MS, 133 ); 134 GucRegistry::define_int_guc( 135 c"jev.concurrency", 136 c"Rows a jev scan judges at once.", 137 c"Each is one HTTP/2 stream on the backend's connection. A LIMIT stops judging after this many rows past the last one returned.", 138 &CONCURRENCY, 139 1, 140 1000, 141 GucContext::Userset, 142 GucFlags::default(), 143 ); 144 // Superuser-only: spend is reckoned from it. 145 GucRegistry::define_float_guc( 146 c"jev.price_per_mtok", 147 c"Dollars per million input tokens for jev.model.", 148 c"Output tokens are not billed. The default is jev-1.13.0's launch price.", 149 &PRICE_PER_MTOK, 150 0.0, 151 f64::MAX, 152 GucContext::Suset, 153 GucFlags::default(), 154 ); 155 // Sighup: one account, one limit, so every backend must agree on it. 156 // The defaults are jev-1.13.0's published account limits. 157 GucRegistry::define_int_guc( 158 c"jev.max_tokens_per_second", 159 c"Ceiling on input tokens per second across the whole cluster; 0 is no limit.", 160 c"Enforced before sending, in shared memory. The rate in force halves on each 429 and climbs back as answers arrive.", 161 &MAX_TOKENS_PER_SECOND, 162 0, 163 i32::MAX, 164 GucContext::Sighup, 165 GucFlags::default(), 166 ); 167 GucRegistry::define_int_guc( 168 c"jev.max_requests_per_minute", 169 c"Ceiling on requests per minute across the whole cluster; 0 is no limit.", 170 c"Enforced before sending, in shared memory. The rate in force halves on each 429 and climbs back as answers arrive.", 171 &MAX_REQUESTS_PER_MINUTE, 172 0, 173 i32::MAX, 174 GucContext::Sighup, 175 GucFlags::default(), 176 ); 177 // Like temp_file_limit, -1 is no limit, so 0 can mean "send nothing". 178 GucRegistry::define_int_guc( 179 c"jev.max_rows", 180 c"Rows a statement may send to Jev; -1 for no limit.", 181 c"Checked against the planner's estimate before the first request, and against the rows sent while it runs.", 182 &MAX_ROWS, 183 -1, 184 i32::MAX, 185 GucContext::Userset, 186 GucFlags::default(), 187 ); 188 GucRegistry::define_float_guc( 189 c"jev.max_cost", 190 c"Dollars a statement may spend on Jev; -1 for no limit.", 191 c"Each request is budgeted at three billed attempts, at jev.price_per_mtok. Checked against the estimate before the first request, and while it runs against the input tokens each answer reports plus the worst case of the requests in flight.", 192 &MAX_COST, 193 -1.0, 194 f64::MAX, 195 GucContext::Userset, 196 GucFlags::default(), 197 ); 198 GucRegistry::define_string_guc( 199 c"jev.cache_namespace", 200 c"Namespace of the answer cache.", 201 c"Part of every cache key: changing it re-judges every input.", 202 &CACHE_NAMESPACE, 203 GucContext::Userset, 204 GucFlags::default(), 205 ); 206 // Keeps the cache usable once a pinned model is withdrawn: nothing 207 // is sent, so no API key is needed and nothing is spent. 208 GucRegistry::define_bool_guc( 209 c"jev.cache_only", 210 c"Serve jev answers from the cache only.", 211 c"A judgment with no cached answer fails instead of being sent.", 212 &CACHE_ONLY, 213 GucContext::Userset, 214 GucFlags::default(), 215 ); 216 // Userset: it decides only what this session's own failures do. 217 GucRegistry::define_enum_guc( 218 c"jev.on_error", 219 c"What a failed jev judgment does: error or unsure.", 220 c"unsure answers NULL for a call whose request failed at the remote or on the way to it, so the row is skipped by a jev qual. Settings, key, budget and cache_only failures still raise.", 221 &ON_ERROR, 222 GucContext::Userset, 223 GucFlags::default(), 224 ); 225 // Noul only: thresholds do not transfer between question types 226 // (contract *SQL surface*). Its default is the last step of jev()'s 227 // precedence, so RESET falls back to 0.5. 228 GucRegistry::define_float_guc( 229 c"jev.threshold", 230 c"Probability at or above which jev() is true, when the call gives none.", 231 c"A threshold passed to jev() takes precedence. Record a tuned value against the pinned jev.model.", 232 &THRESHOLD, 233 0.0, 234 1.0, 235 GucContext::Userset, 236 GucFlags::default(), 237 ); 238 scan::register::<JevScan>(); 239}
The durable cache (contract Cache). Each row is an audit receipt for
one judgment: key is postjevsql_core::CacheKey, the hash of the
model pin, prompt version, namespace, layout, state and question,
which are the exact bytes sent.
245pgrx::extension_sql!( 246 "CREATE TABLE jev_cache ( 247 key bytea PRIMARY KEY, 248 model text NOT NULL, 249 answered_model text NOT NULL, 250 namespace text NOT NULL, 251 layout text NOT NULL, 252 state text NOT NULL, 253 question bytea NOT NULL, 254 request_id text, 255 input_tokens bigint NOT NULL, 256 output_tokens bigint NOT NULL, 257 answer jsonb NOT NULL, 258 at timestamptz NOT NULL DEFAULT now() 259 ); 260 REVOKE ALL ON TABLE jev_cache FROM PUBLIC;", 261 name = "jev_cache", 262);
The typed results of jev_eval and the *_full functions (contract
SQL surface): each judgment's answer with the model that answered it
and the tokens of the request it was asked in, from its jev_cache
row. A Noul answer carries no confidence. probabilities keeps the
order asked: an array, of [label, p] pairs for a Choice.
269pgrx::extension_sql!( 270 "CREATE TYPE jev_noul_result AS ( 271 probability float8, model text, input_tokens int, output_tokens int 272 ); 273 CREATE TYPE jev_score_result AS ( 274 score float8, confidence float8, probabilities jsonb, legend jsonb, 275 model text, input_tokens int, output_tokens int 276 ); 277 CREATE TYPE jev_choice_result AS ( 278 choice text, confidence float8, probabilities jsonb, 279 model text, input_tokens int, output_tokens int 280 ); 281 CREATE TYPE jev_tree_result AS ( 282 path text, probability float8, separation float8, depth int, leaf bool, fit float8 283 ); 284 CREATE TYPE jev_shortlist_result AS ( 285 choice text, confidence float8, probabilities jsonb, 286 model text, input_tokens int, output_tokens int, fit float8 287 );", 288 name = "result_types", 289);
Spending the API key takes an explicit GRANT (contract Cost and
safety). finalize runs after every function exists.
293pgrx::extension_sql!( 294 "REVOKE ALL ON FUNCTION jev_prob(anyelement, text) FROM PUBLIC; 295 REVOKE ALL ON FUNCTION jev(anyelement, text) FROM PUBLIC; 296 REVOKE ALL ON FUNCTION jev(anyelement, text, float8) FROM PUBLIC; 297 REVOKE ALL ON FUNCTION jev_score(anyelement, text, text[]) FROM PUBLIC; 298 REVOKE ALL ON FUNCTION jev_score_norm(anyelement, text, text[]) FROM PUBLIC; 299 REVOKE ALL ON FUNCTION jev_choice(anyelement, text, text[]) FROM PUBLIC; 300 REVOKE ALL ON FUNCTION jev_eval(anyelement, text) FROM PUBLIC; 301 REVOKE ALL ON FUNCTION jev_score_full(anyelement, text, text[]) FROM PUBLIC; 302 REVOKE ALL ON FUNCTION jev_choice_full(anyelement, text, text[]) FROM PUBLIC; 303 REVOKE ALL ON FUNCTION jev_confidence(anyelement, text, text, text[]) FROM PUBLIC; 304 REVOKE ALL ON FUNCTION jev_choice_tree(anyelement, text, text[]) FROM PUBLIC; 305 REVOKE ALL ON FUNCTION jev_choice_tree(anyelement, text, text[], float8) FROM PUBLIC; 306 REVOKE ALL ON FUNCTION jev_choice_tree_full(anyelement, text, text[]) FROM PUBLIC; 307 REVOKE ALL ON FUNCTION jev_choice_tree_full(anyelement, text, text[], float8) FROM PUBLIC; 308 REVOKE ALL ON FUNCTION jev_choice_shortlist(anyelement, text, text[]) FROM PUBLIC; 309 REVOKE ALL ON FUNCTION jev_choice_shortlist(anyelement, text, text[], text[]) FROM PUBLIC; 310 REVOKE ALL ON FUNCTION jev_choice_shortlist_full(anyelement, text, text[]) FROM PUBLIC; 311 REVOKE ALL ON FUNCTION jev_choice_shortlist_full(anyelement, text, text[], text[]) FROM PUBLIC; 312 REVOKE ALL ON FUNCTION jev_support(internal) FROM PUBLIC; 313 REVOKE ALL ON FUNCTION jev_stats() FROM PUBLIC;", 314 name = "revoke_from_public", 315 finalize, 316);
The probability that the answer to yes/no question about row is
yes. The jev scan evaluates every call it can claim; this body runs
only for a call it could not, and refuses.
Whether the answer to yes/no question about row is yes: its
probability is at least jev.threshold (0.5 unless set). Evaluated by
the jev scan, like [jev_prob], which it shares a judgment with.
Whether the probability that the answer to question about row is
yes is at least threshold, which takes precedence over
jev.threshold.
The probability-weighted level, from 0 to levels − 1, of row on
the rubric question whose levels are given lowest first. Evaluated
by the jev scan, which sends the levels in the caller's order.
[jev_score] scaled to 0 to 1, so rubrics of different lengths
compare. The two share one judgment and one cache row.
The label of options that best answers question about row.
Evaluated by the jev scan, which sends the options in the caller's
order; a single option is the answer without asking.
The Noul [jev_prob] asks, in full: its probability, the model that
answered and the request's tokens. Shares a judgment with jev_prob
and jev().
[jev_score] in full: the level, the vendor's confidence, each
level's probability and description, the model and the tokens.
383#[pg_extern(volatile, parallel_restricted, support = jev_support, requires = ["result_types"])] 384fn jev_score_full( 385 _row: AnyElement, 386 _question: &str, 387 _levels: Vec<Option<String>>, 388) -> pgrx::composite_type!('static, "jev_score_result") { 389 Failure::Unbatched("jev_score_full").raise() 390}
[jev_choice] in full: the label, the vendor's confidence, each
option's probability in the order given, the model and the tokens.
394#[pg_extern(volatile, parallel_restricted, support = jev_support, requires = ["result_types"])] 395fn jev_choice_full( 396 _row: AnyElement, 397 _question: &str, 398 _options: Vec<Option<String>>, 399) -> pgrx::composite_type!('static, "jev_choice_result") { 400 Failure::Unbatched("jev_choice_full").raise() 401}
The vendor's confidence in the answer to the Choice or Score question
about row, as it returned it and never recomputed. kind is
'choice' or 'score', and options the options or levels, as
pg-jev's signature has them. Shares a judgment with [jev_choice] or
[jev_score] on the same arguments. A Noul carries no confidence, so
'noul' is refused rather than derived from the probability.
The leaf of the label tree paths that best answers question about
row, as its path. Each path is a leaf's labels from the top down,
joined by >; siblings are sent in the order they first appear.
Evaluated by the jev scan as a best-first search, one Choice per node,
so more than 255 labels can be chosen from (contract SQL surface).
[jev_choice_tree] stopping early: the path of the deepest node whose
probability is at least tau, which may be a group rather than a
leaf; NULL when no label is that likely.
[jev_choice_tree] in full: the chosen path, its probability (the
product of its edges'), its separation from the runner-up leaf, its
depth and whether it is a leaf. Shares every node's judgment with
jev_choice_tree on the same arguments.
436#[pg_extern(volatile, parallel_restricted, support = jev_support, requires = ["result_types"])] 437fn jev_choice_tree_full( 438 _row: AnyElement, 439 _question: &str, 440 _paths: Vec<Option<String>>, 441) -> pgrx::composite_type!('static, "jev_tree_result") { 442 Failure::Unbatched("jev_choice_tree_full").raise() 443}
[jev_choice_tree_tau] in full. The separation is NULL: stopping at τ
proves no runner-up. A stop at the root has a NULL path and depth 0.
447#[pg_extern( 448 name = "jev_choice_tree_full", 449 volatile, 450 parallel_restricted, 451 support = jev_support, 452 requires = ["result_types"] 453)] 454fn jev_choice_tree_full_tau( 455 _row: AnyElement, 456 _question: &str, 457 _paths: Vec<Option<String>>, 458 _tau: f64, 459) -> pgrx::composite_type!('static, "jev_tree_result") { 460 Failure::Unbatched("jev_choice_tree_full").raise() 461}
The option of options that best answers question about row, by
chunk and shortlist, for label sets with no meaningful hierarchy: one
request asks a Choice over each run of up to 254 options, labels only,
and a second asks one Choice over the top 3 of each, in the caller's
order (contract SQL surface, "More than 255 options").
[jev_choice_shortlist] with a description per option, sent with the
finalists only; a NULL description sends none.
475#[pg_extern(name = "jev_choice_shortlist", volatile, parallel_restricted, support = jev_support)] 476fn jev_choice_shortlist_described( 477 _row: AnyElement, 478 _question: &str, 479 _options: Vec<Option<String>>, 480 _descriptions: Vec<Option<String>>, 481) -> String { 482 Failure::Unbatched("jev_choice_shortlist").raise() 483}
[jev_choice_shortlist] in full, as a jev_shortlist_result: the
chosen label, and the finalist Choice's confidence, probabilities as
[label, p] pairs in the caller's order (among the finalists, not over
every option), model and request tokens; then fit, the probability
that any option fits, from a Noul asked in the first request. Shares
every Choice's judgment with jev_choice_shortlist on the same
arguments. A single option is 1.0, with a NULL model, 0 tokens and a
NULL fit, since nothing is sent.
493#[pg_extern(volatile, parallel_restricted, support = jev_support, requires = ["result_types"])] 494fn jev_choice_shortlist_full( 495 _row: AnyElement, 496 _question: &str, 497 _options: Vec<Option<String>>, 498) -> pgrx::composite_type!('static, "jev_shortlist_result") { 499 Failure::Unbatched("jev_choice_shortlist_full").raise() 500}
[jev_choice_shortlist_described] in full.
503#[pg_extern( 504 name = "jev_choice_shortlist_full", 505 volatile, 506 parallel_restricted, 507 support = jev_support, 508 requires = ["result_types"] 509)] 510fn jev_choice_shortlist_full_described( 511 _row: AnyElement, 512 _question: &str, 513 _options: Vec<Option<String>>, 514 _descriptions: Vec<Option<String>>, 515) -> pgrx::composite_type!('static, "jev_shortlist_result") { 516 Failure::Unbatched("jev_choice_shortlist_full").raise() 517}
The planner support function of every jev* function.
This session's work since the backend started.
526#[pg_extern(volatile, parallel_restricted)] 527#[allow(clippy::type_complexity, reason = "pgrx reads the column names from this signature")] 528fn jev_stats() -> TableIterator< 529 'static, 530 ( 531 name!(requests, i64), 532 name!(retries, i64), 533 name!(redials, i64), 534 name!(connections, i64), 535 name!(cache_hits, i64), 536 name!(cache_misses, i64), 537 name!(dedupe_hits, i64), 538 name!(input_tokens, i64), 539 name!(output_tokens, i64), 540 name!(cost, f64), 541 name!(in_flight, i64), 542 ), 543> { 544 let s = stats::snapshot(); 545 let n = |v: u64| i64::try_from(v).unwrap_or(i64::MAX); 546 TableIterator::once(( 547 n(s.requests), 548 n(s.retries), 549 n(s.redials), 550 n(s.connections), 551 n(s.cache_hits), 552 n(s.cache_misses), 553 n(s.dedupe_hits), 554 n(s.input_tokens), 555 n(s.output_tokens), 556 s.cost, 557 n(s.in_flight), 558 )) 559}
One scan's settings and client, read once in the executor's startup
(contract Execution: a generic plan outlives a SET).
563struct JevScan {
None under jev.cache_only, which never sends, and in an
EvalPlanQual recheck.
566 client: Option<Rc<Client<HttpsTransport, Backend, Telemetry>>>,
An EvalPlanQual recheck: a miss is a row version the statement did not judge, refused rather than sent (contract SQL surface).
jev.on_error = unsure.
573 unsure: bool,
jev.threshold, for a jev() call that gives none.
Answers received and not yet stored: judgements cannot call
Postgres, so settle writes them.
One scan's work, as EXPLAIN ANALYZE shows it. jev_stats() counts
the same things for the whole backend.
The judgments the statement has asked, by cache key: an equal call in
any of its scans waits on the first instead of sending (contract
Execution, "Deduplicate before sending"). One per statement, like
the [Budget].
One judgment's answer, shared by every equal call in the statement.
It failed; the scan that sent it raises the error.
615 Failed,
620impl Shared { 621 fn resolve(&self, answer: Result<Judged, Unanswered>) { 622 if self.answer.borrow().is_none() { 623 *self.answer.borrow_mut() = Some(answer); 624 for waker in self.waiting.take() { 625 waker.wake(); 626 } 627 } 628 } 629 630 async fn wait(&self) -> Result<Judged, Unanswered> { 631 std::future::poll_fn(|cx| match self.answer.borrow().clone() { 632 Some(answer) => Poll::Ready(answer), 633 None => { 634 self.waiting.borrow_mut().push(cx.waker().clone()); 635 Poll::Pending 636 } 637 }) 638 .await 639 } 640}
The judgments a request owes to their waiters. Those it drops
unanswered leave asking, so a later equal call asks again, and wake
their waiters with why.
652impl Drop for Owed { 653 fn drop(&mut self) { 654 // Never panics: this runs while a cancel unwinds. 655 let why = if self.failed { Unanswered::Failed } else { Unanswered::Dropped }; 656 for (key, shared) in self.shared.drain(..) { 657 if shared.answer.borrow().is_none() { 658 if let Ok(mut asking) = self.asking.0.try_borrow_mut() 659 && asking.get(&key).is_some_and(|s| Rc::ptr_eq(s, &shared)) 660 { 661 asking.remove(&key); 662 } 663 shared.resolve(Err(why)); 664 } 665 } 666 } 667}
One answered judgment, as jev_cache stores it.
The spend guards, and what the statement has spent against them. One per statement, shared by its scans (contract Cost and safety).
Dollars billed by the answers received, from their usage.
692 spent: Cell<f64>,
697impl Budget { 698 fn read() -> Budget { 699 Budget { 700 max_rows: u64::try_from(MAX_ROWS.get()).ok(), 701 max_cost: Some(MAX_COST.get()).filter(|&c| c >= 0.0), 702 price_per_mtok: PRICE_PER_MTOK.get(), 703 ratio: learned_ratio(), 704 rows: Cell::new(0), 705 spent: Cell::new(0.0), 706 in_flight: Cell::new(0.0), 707 } 708 }
Dollars chars of request would cost at worst.
Refuses before sending when this row would take the statement over a guard, counting what it has spent and the worst case of what it has in flight; otherwise commits the row's worst case.
718 fn commit(self: &Rc<Self>, requests: u64, chars: usize) -> Result<Committed, Failure> { 719 let rows = self.rows.get() + 1; 720 if let Some(max) = self.max_rows.filter(|&max| rows > max) { 721 return Err(Failure::Budget(format!("this statement would send more than jev.max_rows ({max}) rows"))); 722 } 723 let worst = self.price(requests as f64, chars as f64); 724 if let Some(max) = self.max_cost.filter(|&max| self.spent.get() + self.in_flight.get() + worst > max) { 725 return Err(Failure::Budget(format!("this statement would spend more than jev.max_cost (${max})"))); 726 } 727 self.rows.set(rows); 728 self.in_flight.set(self.in_flight.get() + worst); 729 Ok(Committed { budget: self.clone(), worst: Some(worst) }) 730 } 731}
A row's worst case, held in flight until its answers report what was billed.
741impl Committed {
Replaces the worst case with the input tokens the answers report, for every attempt each took: with no idempotency key, an attempt that timed out may have been billed too.
745 fn settle(&mut self, answers: &[jev_client::Answered]) { 746 let tokens: f64 = answers.iter().map(|a| a.response.usage().input_tokens as f64 * f64::from(a.attempts)).sum(); 747 if let Some(worst) = self.worst.take() { 748 let b = &self.budget; 749 b.in_flight.set(b.in_flight.get() - worst); 750 b.spent.set(b.spent.get() + tokens * b.price_per_mtok / 1e6); 751 } 752 } 753}
755impl Drop for Committed {
Unanswered (failed, cancelled, cut off by a LIMIT): what it was billed is unknown, so the worst case stays spent.
Indexed by [Call::function]; [JevScan::lookup] matches on it.
771 const FUNCTIONS: &'static [Function] = &[ 772 Function { name: "jev_prob", args: &[Arg::Row, Arg::Text], returns: Returns::Float8 }, 773 Function { name: "jev", args: &[Arg::Row, Arg::Text], returns: Returns::Bool }, 774 Function { name: "jev", args: &[Arg::Row, Arg::Text, Arg::Float8], returns: Returns::Bool }, 775 Function { name: "jev_score", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Float8 }, 776 Function { name: "jev_score_norm", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Float8 }, 777 Function { name: "jev_choice", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Text }, 778 Function { name: "jev_eval", args: &[Arg::Row, Arg::Text], returns: Returns::Record }, 779 Function { name: "jev_score_full", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Record }, 780 Function { name: "jev_choice_full", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Record }, 781 Function { 782 name: "jev_confidence", 783 args: &[Arg::Row, Arg::Text, Arg::Text, Arg::TextArray], 784 returns: Returns::Float8, 785 }, 786 Function { name: "jev_choice_tree", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Text }, 787 Function { 788 name: "jev_choice_tree", 789 args: &[Arg::Row, Arg::Text, Arg::TextArray, Arg::Float8], 790 returns: Returns::Text, 791 }, 792 Function { name: "jev_choice_tree_full", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Record }, 793 Function { 794 name: "jev_choice_tree_full", 795 args: &[Arg::Row, Arg::Text, Arg::TextArray, Arg::Float8], 796 returns: Returns::Record, 797 }, 798 Function { name: "jev_choice_shortlist", args: &[Arg::Row, Arg::Text, Arg::TextArray], returns: Returns::Text }, 799 Function { 800 name: "jev_choice_shortlist", 801 args: &[Arg::Row, Arg::Text, Arg::TextArray, Arg::TextArray], 802 returns: Returns::Text, 803 }, 804 Function { 805 name: "jev_choice_shortlist_full", 806 args: &[Arg::Row, Arg::Text, Arg::TextArray], 807 returns: Returns::Record, 808 }, 809 Function { 810 name: "jev_choice_shortlist_full", 811 args: &[Arg::Row, Arg::Text, Arg::TextArray, Arg::TextArray], 812 returns: Returns::Record, 813 }, 814 ];
816 fn begin(statement: &Statement) -> Result<Self, Refusal> { 817 let setting_err = |e: &dyn std::fmt::Display| Failure::Setting(e.to_string()).refusal(); 818 let model = setting(&MODEL).ok_or_else(|| setting_err(&"jev.model is not set"))?; 819 let model = ModelId::pinned(&model).map_err(|e| setting_err(&e))?; 820 let endpoint = setting(&ENDPOINT).ok_or_else(|| setting_err(&"jev.endpoint is not set"))?; 821 let endpoint = Endpoint::parse(&endpoint, setting(&CA_FILE).as_deref(), setting(&DNS_SERVERS).as_deref()) 822 .map_err(|e| setting_err(&e))? 823 .with_keep_alive(KeepAlive { 824 interval: Duration::from_millis(KEEPALIVE_INTERVAL.get() as u64), 825 timeout: Duration::from_millis(KEEPALIVE_TIMEOUT.get() as u64), 826 }); 827 // The token ratio is learned once, with the statement's budget, 828 // and prices the rate limit's admissions and the tree's lookahead. 829 let budget = statement.shared(Budget::read); 830 let client = if CACHE_ONLY.get() || statement.rechecking() { 831 None 832 } else { 833 let key = api_key().map_err(|e| Failure::Key(e).refusal())?; 834 let limit = RateLimit::attach( 835 MAX_TOKENS_PER_SECOND.get() as f64, 836 MAX_REQUESTS_PER_MINUTE.get() as f64, 837 budget.ratio, 838 ); 839 let transport = HttpsTransport::new(endpoint.clone()).rate_limited(limit); 840 let client = Client::new(transport, Backend, model.clone(), &key) 841 .map_err(|e| Failure::Client(e).refusal())? 842 .observed(Telemetry { price_per_mtok: PRICE_PER_MTOK.get() }); 843 Some(Rc::new(client)) 844 }; 845 let scope = KeyScope { 846 model, 847 prompt: PROMPT_VERSION, 848 namespace: setting(&CACHE_NAMESPACE).unwrap_or_default(), 849 }; 850 Ok(JevScan { 851 client, 852 rechecking: statement.rechecking(), 853 endpoint, 854 concurrency: CONCURRENCY.get() as usize, 855 unsure: ON_ERROR.get() == OnError::Unsure, 856 threshold: THRESHOLD.get(), 857 budget, 858 scope, 859 answered: Rc::default(), 860 asking: statement.shared(Asking::default), 861 counts: Rc::default(), 862 }) 863 } 864 865 fn afford(estimate: &Estimate) -> Result<(), Refusal> { 866 // Nothing to price; and the cache query learning the token ratio 867 // passes through here itself, with no jev scan. 868 if CACHE_ONLY.get() || estimate.rows <= 0.0 { 869 return Ok(()); 870 } 871 let budget = Budget::read(); 872 let refuse = |message: String| { 873 Err(Failure::Budget(message).refusal()) 874 }; 875 if let Some(max) = budget.max_rows.filter(|&max| estimate.rows > max as f64) { 876 return refuse(format!("this statement is estimated to send {:.0} rows, over jev.max_rows ({max})", estimate.rows)); 877 } 878 let cost = budget.price(estimate.rows, estimate.bytes); 879 if let Some(max) = budget.max_cost.filter(|&max| cost > max) { 880 return refuse(format!("this statement is estimated to spend up to ${cost:.6}, over jev.max_cost (${max})")); 881 } 882 Ok(()) 883 }
The price before it is paid, and under ANALYZE what was paid
(contract Cost and safety). Plain EXPLAIN runs no rows, so it
cannot know the cache hits: it prices every candidate as a miss,
the same bound [Judge::afford] refuses on.
889 fn explain(estimate: Option<&Estimate>, run: Option<&Self>) -> Vec<Property> { 890 let estimated = |label, v: f64| Property { label, unit: None, value: Shown::Integer(v.round() as i64) }; 891 let counted = |label, c: &Cell<u64>| Property { 892 label, 893 unit: None, 894 value: Shown::Integer(i64::try_from(c.get()).unwrap_or(i64::MAX)), 895 }; 896 let dollars = |label, v| Property { label, unit: Some(c"USD"), value: Shown::Float(v, 6) }; 897 let mut shown = Vec::new(); 898 if let Some(e) = estimate { 899 let budget = Budget::read(); 900 shown.push(estimated(c"Candidate Rows", e.rows)); 901 shown.push(estimated(c"Estimated Requests", e.rows)); 902 shown.push(estimated(c"Estimated Input Tokens", budget.ratio.tokens(e.rows, e.bytes))); 903 shown.push(dollars(c"Worst-Case Cost", budget.price(e.rows, e.bytes))); 904 } 905 if let Some(scan) = run { 906 let c = &scan.counts; 907 shown.push(counted(c"Cache Hits", &c.hits)); 908 shown.push(counted(c"Cache Misses", &c.misses)); 909 shown.push(counted(c"Shared Judgments", &c.deduped)); 910 shown.push(counted(c"Requests", &c.requests)); 911 shown.push(counted(c"Input Tokens", &c.input_tokens)); 912 shown.push(dollars(c"Cost", c.input_tokens.get() as f64 * scan.budget.price_per_mtok / 1e6)); 913 } 914 shown 915 }
jev.concurrency, clamped to the streams the server accepts.
A call equal to one the statement has already asked waits for that answer; the rest are looked up in the cache. Each distinct row argument with a miss is one request, with the row as the state and each distinct missed question about it asked once (contract Execution, "row as state").
A jev_choice_tree call is a label-tree search of several rounds:
each round's Choices are looked up on the backend and the misses
sent in one request, with the row as the state. A
jev_choice_shortlist call is the same with two rounds.
932 fn judge(&self, calls: Vec<Call>) -> Judgement { 933 let (lists, calls): (Vec<_>, Vec<_>) = 934 calls.into_iter().enumerate().partition(|(_, c)| SHORTLIST_FUNCTIONS.contains(&c.function)); 935 let (trees, calls): (Vec<_>, Vec<_>) = 936 calls.into_iter().partition(|(_, c)| TREE_FUNCTIONS.contains(&c.function)); 937 let (at, calls): (Vec<usize>, Vec<Call>) = calls.into_iter().unzip(); 938 let trees = match trees.into_iter().map(|(i, c)| TreeCall::new(&c).map(|t| (i, t))).collect::<Result<Vec<_>, _>>() { 939 Ok(trees) => trees, 940 Err(e) => return Box::pin(async move { Err(e) }), 941 }; 942 let lists = match lists 943 .into_iter() 944 .map(|(i, c)| ShortlistCall::new(&c).map(|l| (i, l))) 945 .collect::<Result<Vec<_>, _>>() 946 { 947 Ok(lists) => lists, 948 Err(e) => return Box::pin(async move { Err(e) }), 949 }; 950 let flat = self.judge_flat(calls); 951 let (lookup, sender, unsure) = (self.lookup_state(), self.sender(), self.unsure); 952 Box::pin(async move { 953 let flat = flat.await?; 954 let mut outs: Vec<Option<Out>> = (0..at.len() + trees.len() + lists.len()).map(|_| None).collect(); 955 for (i, out) in at.into_iter().zip(flat) { 956 outs[i] = out; 957 } 958 for (i, tree) in trees { 959 outs[i] = tree.search(&lookup, &sender, unsure).await?; 960 } 961 for (i, list) in lists { 962 outs[i] = list.search(&lookup, &sender, unsure).await?; 963 } 964 Ok(outs) 965 }) 966 }
Stores the answers received since the last call.
968 fn settle(&self) -> Result<(), Refusal> { 969 let receipts = std::mem::take(&mut *self.answered.borrow_mut()); 970 let (model, namespace) = (self.scope.model.as_str(), &self.scope.namespace); 971 if receipts.is_empty() { 972 return Ok(()); 973 } 974 let store = || -> Result<(), pgrx::spi::Error> { 975 for r in receipts { 976 let answer = pgrx::JsonB(r.answer); 977 Spi::run_with_args( 978 "INSERT INTO jev_cache (key, model, answered_model, namespace, layout, state, question, 979 request_id, input_tokens, output_tokens, answer) 980 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) 981 ON CONFLICT (key) DO NOTHING", 982 &[ 983 r.key.as_bytes().to_vec().into(), 984 model.into(), 985 r.answered_model.into(), 986 namespace.as_str().into(), 987 r.layout.tag().into(), 988 r.state.into(), 989 r.question.into(), 990 r.request_id.into(), 991 i64::try_from(r.usage.input_tokens).unwrap_or(i64::MAX).into(), 992 i64::try_from(r.usage.output_tokens).unwrap_or(i64::MAX).into(), 993 answer.into(), 994 ], 995 )?; 996 } 997 Ok(()) 998 }; 999 as_owner_of(c"jev_cache", store) 1000 .and_then(|r| r.map_err(|e| e.to_string())) 1001 .map_err(|e| Failure::Cache(e).refusal()) 1002 } 1003}
1005impl JevScan {
Every call but a label-tree search's.
1007 fn judge_flat(&self, calls: Vec<Call>) -> Judgement { 1008 let asked = match self.lookup(&calls) { 1009 Ok(asked) => asked, 1010 Err(e) => return Box::pin(async move { Err(e) }), 1011 }; 1012 let sender = self.sender(); 1013 // Made now, not in the future: a judgement dropped before its 1014 // first poll still owes its keys. 1015 let owed = sender.owe(&asked.requests); 1016 let unsure = self.unsure; 1017 // Under `jev.on_error = unsure` a failed judgment is NULL; any 1018 // other failure, and every failure under `error`, raises. 1019 let judged = move |r: Result<Judged, Failure>| match r { 1020 Ok(p) => Ok(Some(p)), 1021 Err(e) if unsure && e.is_failed_judgment() => Ok(None), 1022 Err(e) => Err(e.refusal()), 1023 }; 1024 Box::pin(async move { 1025 let Asked { requests, answers, shapes } = asked; 1026 let responses = match sender.send(&requests, owed).await { 1027 Ok(responses) => Some(responses), 1028 Err(e) => judged(Err(e)).map(|_| None)?, 1029 }; 1030 let mut outs = Vec::with_capacity(answers.len()); 1031 for (answer, shape) in answers.into_iter().zip(shapes) { 1032 let p = match answer { 1033 Answer::Cached(p) => Some(p), 1034 Answer::Asked(request, key) => responses 1035 .as_ref() 1036 .map(|r| requests[request].keys[&key].key.read(&r[request].response).0), 1037 Answer::Shared(shared, again) => judged(sender.wait(shared, again).await)?, 1038 }; 1039 outs.push(p.map(|j| shape.out(j))); 1040 } 1041 Ok(outs) 1042 }) 1043 }
1045 fn lookup_state(&self) -> Lookup { 1046 Lookup { 1047 scope: self.scope.clone(), 1048 planner: Planner::new(), 1049 asking: self.asking.clone(), 1050 counts: self.counts.clone(), 1051 sends: self.client.is_some(), 1052 rechecking: self.rechecking, 1053 } 1054 } 1055 1056 fn sender(&self) -> Sender { 1057 Sender { 1058 client: self.client.clone(), 1059 budget: self.budget.clone(), 1060 answered: self.answered.clone(), 1061 asking: self.asking.clone(), 1062 counts: self.counts.clone(), 1063 } 1064 }
Answers each call from the cache where it can, and groups the rest into requests. Runs on the backend, so it may query.
1068 fn lookup(&self, calls: &[Call]) -> Result<Asked, Refusal> { 1069 let mut requests: Vec<Request> = Vec::new(); 1070 let mut answers = Vec::with_capacity(calls.len()); 1071 let mut shapes = Vec::with_capacity(calls.len()); 1072 let lookup = self.lookup_state(); 1073 // Answered at once, or asked once every row's questions are known, 1074 // since those choose its layout. 1075 let mut pending: Vec<Option<(&String, &String, Question)>> = Vec::with_capacity(calls.len()); 1076 for call in calls { 1077 let noul = |q: &String| Question::Noul(Noul::new(Json::text(q))); 1078 let (row, question, asked, shape) = match (call.function, call.args.as_slice()) { 1079 (0, [Value::Json(row), Value::Text(q)]) => (row, q, noul(q), Shape::Probability), 1080 (6, [Value::Json(row), Value::Text(q)]) => (row, q, noul(q), Shape::Full), 1081 (1, [Value::Json(row), Value::Text(q)]) => (row, q, noul(q), Shape::AtLeast(self.threshold)), 1082 (2, [Value::Json(row), Value::Text(q), Value::Float8(t)]) => { 1083 if !(0.0..=1.0).contains(t) { 1084 return Err(Failure::Setting(format!("jev() threshold must be between 0 and 1, not {t}")).refusal()); 1085 } 1086 (row, q, noul(q), Shape::AtLeast(*t)) 1087 } 1088 (f @ (3 | 4 | 7), [Value::Json(row), Value::Text(q), Value::TextArray(levels)]) => { 1089 let shape = match f { 1090 3 => Shape::Level, 1091 4 => Shape::Normalized, 1092 _ => Shape::Full, 1093 }; 1094 (row, q, Question::Score(score(q, levels)?), shape) 1095 } 1096 (f @ (5 | 8), [Value::Json(row), Value::Text(q), Value::TextArray(options)]) => { 1097 let shape = if f == 5 { Shape::Label } else { Shape::Full }; 1098 match choice(q, options)? { 1099 Choosing::Ask(c) => (row, q, Question::Choice(c), shape), 1100 Choosing::Only(label) => { 1101 // Nothing to judge: no request, no cache row. 1102 shapes.push(shape); 1103 answers.push(Some(Answer::Cached(Judged::only(label)))); 1104 pending.push(None); 1105 continue; 1106 } 1107 } 1108 } 1109 (9, [Value::Json(row), Value::Text(q), Value::Text(kind), Value::TextArray(options)]) => { 1110 match kind.as_str() { 1111 "score" => (row, q, Question::Score(score(q, options)?), Shape::Confidence), 1112 "choice" => match choice(q, options)? { 1113 Choosing::Ask(c) => (row, q, Question::Choice(c), Shape::Confidence), 1114 Choosing::Only(label) => { 1115 shapes.push(Shape::Confidence); 1116 answers.push(Some(Answer::Cached(Judged::only(label)))); 1117 pending.push(None); 1118 continue; 1119 } 1120 }, 1121 other => { 1122 return Err(Failure::Setting(format!( 1123 "jev_confidence: kind must be 'choice' or 'score', not '{other}'; a Noul answer carries no confidence" 1124 )) 1125 .refusal()); 1126 } 1127 } 1128 } 1129 (f, args) => unreachable!("FUNCTIONS[{f}] does not take {args:?}"), 1130 }; 1131 shapes.push(shape); 1132 answers.push(None); 1133 pending.push(Some((row, question, asked))); 1134 } 1135 let layouts = lookup.layouts(pending.iter().flatten().map(|(row, _, asked)| (*row, asked)))?; 1136 let answers = answers 1137 .into_iter() 1138 .zip(pending) 1139 .map(|(answer, asked)| match (answer, asked) { 1140 (Some(answer), _) => Ok(answer), 1141 (None, Some((row, question, asked))) => lookup.one(&mut requests, row, layouts[row], question, asked), 1142 (None, None) => unreachable!("every call is answered or asked"), 1143 }) 1144 .collect::<Result<Vec<_>, _>>()?; 1145 Ok(Asked { requests, answers, shapes }) 1146 } 1147}
[JevScan::FUNCTIONS]' indexes of jev_choice_tree, without and with
tau, then of jev_choice_tree_full.
1151const TREE_FUNCTIONS: [usize; 4] = [10, 11, 12, 13];
Of [TREE_FUNCTIONS], those returning jev_tree_result.
1154const TREE_FULL_FUNCTIONS: [usize; 2] = [12, 13];
How a jev_choice_tree path separates its labels, as Within: in each
node's instructions shows the parent path.
1158const PATH_SEPARATOR: &str = " > ";
One jev_choice_tree call: its row, question and tree, built on the
backend before anything is sent.
1171impl TreeCall {
Refused with 22023 before anything is sent: a NULL or empty path, a
repeated one, a path that is also a group, or tau outside [0, 1].
1174 fn new(call: &Call) -> Result<TreeCall, Refusal> { 1175 let refuse = |m: String| Failure::Setting(format!("jev_choice_tree: {m}")).refusal(); 1176 let (row, question, paths, tau) = match call.args.as_slice() { 1177 [Value::Json(row), Value::Text(q), Value::TextArray(paths)] => (row, q, paths, None), 1178 [Value::Json(row), Value::Text(q), Value::TextArray(paths), Value::Float8(t)] => { 1179 let tau = Tau::new(*t).ok_or_else(|| refuse(format!("tau must be between 0 and 1, not {t}")))?; 1180 (row, q, paths, Some(tau)) 1181 } 1182 args => unreachable!("jev_choice_tree does not take {args:?}"), 1183 }; 1184 let paths = paths 1185 .iter() 1186 .map(|p| p.as_deref().ok_or_else(|| refuse("a path is NULL".into()))) 1187 .collect::<Result<Vec<_>, _>>()?; 1188 if paths.is_empty() { 1189 return Err(refuse("no paths".into())); 1190 } 1191 let branch = Branch::from_paths(paths.iter().map(|p| p.split(PATH_SEPARATOR))) 1192 .map_err(|e| refuse(e.to_string()))?; 1193 let tree = Tree::new(branch).map_err(|e| refuse(e.to_string()))?; 1194 let full = TREE_FULL_FUNCTIONS.contains(&call.function); 1195 Ok(TreeCall { row: row.clone(), question: question.clone(), tree, tau, full }) 1196 }
The label-tree search: each round's Choices are looked up on the
backend, and the misses sent together in one request with the row
as the state. Each node's Choice is its own judgment and its own
jev_cache row. Returns the path of the leaf, or under tau of the
deepest node at or above it (NULL when that is the root).
1203 async fn search(self, lookup: &Lookup, sender: &Sender, unsure: bool) -> Result<Option<Out>, Refusal> { 1204 let TreeCall { row, question, tree, tau, full } = self; 1205 let describe = Describe::default(); 1206 let choice = |a: label_tree::Ask| { 1207 tree.choice(a.node(), a.order(), &question, describe).expect("the search asks only nodes of two or more children") 1208 }; 1209 let mut search = Search::new(&tree, Params { tau, ..Params::default() }); 1210 let ratio = sender.budget.ratio; 1211 // Under `full`, the "does any label fit" Noul rides in the first 1212 // round's request; NULL when the search sends nothing. 1213 let mut gate = full.then(|| tree.gate(&question, describe)).flatten(); 1214 let mut fit = None; 1215 loop { 1216 let asks = match search.step(|a| { 1217 question_bytes(&Question::Choice(choice(a))).map_or(0, |b| (b.len() as f64 / ratio.chars_per_token()) as usize) 1218 }) { 1219 Step::Done(outcome) => return Ok(tree_out(&tree, &outcome, full, fit)), 1220 Step::Ask(asks) => asks, 1221 }; 1222 let choices: Vec<Question> = asks.iter().map(|&a| Question::Choice(choice(a))).collect(); 1223 // A failed Choice under `unsure` leaves the search unfinished. 1224 let Some((answered, gated)) = round(lookup, sender, unsure, &row, &question, choices, gate.take()).await? else { 1225 return Ok(None); 1226 }; 1227 fit = fit.or(gated); 1228 for (ask, judged) in asks.into_iter().zip(answered) { 1229 search.answer(ask, &probabilities(&judged)?).map_err(|e| wrong_answer(&e))?; 1230 } 1231 } 1232 } 1233}
A Choice judgment's probabilities, in the order its options were sent.
One round of a multi-round search about row: choices looked up on
the backend, and the misses sent together in one request with the row
as the state. Each Choice is its own judgment and its own jev_cache
row. gate, the "does any label fit" Noul, is one more judgment in
the same request. Returns each Choice's judgment in the order sent and
the gate's probability, or None when one failed under
jev.on_error = unsure.
1247async fn round( 1248 lookup: &Lookup, 1249 sender: &Sender, 1250 unsure: bool, 1251 row: &str, 1252 question: &str, 1253 choices: Vec<Question>, 1254 gate: Option<Noul>, 1255) -> Result<Option<(Vec<Judged>, Option<f64>)>, Refusal> { 1256 let (l, r, q) = (lookup.clone(), row.to_owned(), question.to_owned()); 1257 let gated = gate.is_some(); 1258 let (requests, answers) = scan::on_backend(move || { 1259 let mut requests = Vec::new(); 1260 let asked: Vec<Question> = choices.into_iter().chain(gate.map(Question::Noul)).collect(); 1261 let layouts = l.layouts(asked.iter().map(|c| (&r, c)))?; 1262 let answers = asked 1263 .into_iter() 1264 .map(|c| l.one(&mut requests, &r, layouts[&r], &q, c)) 1265 .collect::<Result<Vec<_>, _>>()?; 1266 Ok::<_, Refusal>((requests, answers)) 1267 }) 1268 .await?; 1269 let judged = |r: Result<Judged, Failure>| match r { 1270 Ok(p) => Ok(Some(p)), 1271 Err(e) if unsure && e.is_failed_judgment() => Ok(None), 1272 Err(e) => Err(e.refusal()), 1273 }; 1274 let owed = sender.owe(&requests); 1275 let responses = match sender.send(&requests, owed).await { 1276 Ok(responses) => Some(responses), 1277 Err(e) => judged(Err(e)).map(|_| None)?, 1278 }; 1279 let mut answered = Vec::with_capacity(answers.len()); 1280 for answer in answers { 1281 let p = match answer { 1282 Answer::Cached(p) => Some(p), 1283 Answer::Asked(request, key) => { 1284 responses.as_ref().map(|r| requests[request].keys[&key].key.read(&r[request].response).0) 1285 } 1286 Answer::Shared(shared, again) => judged(sender.wait(shared, again).await)?, 1287 }; 1288 let Some(p) = p else { return Ok(None) }; 1289 answered.push(p); 1290 } 1291 let fit = if gated { 1292 let gate = answered.pop().expect("the gate was asked last"); 1293 Some(gate.noul_probability().map_err(|e| Failure::Cache(e).refusal())?) 1294 } else { 1295 None 1296 }; 1297 Ok(Some((answered, fit))) 1298}
An answer a search cannot take: the server answered a question other than the one asked.
[JevScan::FUNCTIONS]' indexes of jev_choice_shortlist, without and
with descriptions, then of jev_choice_shortlist_full.
1309const SHORTLIST_FUNCTIONS: [usize; 4] = [14, 15, 16, 17];
Of [SHORTLIST_FUNCTIONS], those returning jev_choice_result.
1312const SHORTLIST_FULL_FUNCTIONS: [usize; 2] = [16, 17];
One jev_choice_shortlist call, validated on the backend before
anything is sent.
jev_choice_shortlist_full: the finalist Choice's judgment, not
only its label.
1325impl ShortlistCall {
Refused with 22023 before anything is sent: no options, a NULL, empty or repeated one, descriptions not one per option, or more options than one finalist Choice can take.
1329 fn new(call: &Call) -> Result<ShortlistCall, Refusal> { 1330 let refuse = |m: String| Failure::Setting(format!("jev_choice_shortlist: {m}")).refusal(); 1331 let (row, question, options, descriptions) = match call.args.as_slice() { 1332 [Value::Json(row), Value::Text(q), Value::TextArray(o)] => (row, q, o, None), 1333 [Value::Json(row), Value::Text(q), Value::TextArray(o), Value::TextArray(d)] => (row, q, o, Some(d)), 1334 args => unreachable!("jev_choice_shortlist does not take {args:?}"), 1335 }; 1336 if let Some(d) = descriptions.filter(|d| d.len() != options.len()) { 1337 return Err(refuse(format!("{} options but {} descriptions", options.len(), d.len()))); 1338 } 1339 let options = options 1340 .iter() 1341 .enumerate() 1342 .map(|(i, o)| { 1343 let label = o.clone().ok_or_else(|| refuse(format!("option {} is NULL", i + 1)))?; 1344 let description = descriptions.and_then(|d| d[i].as_deref()).map(Json::text); 1345 Ok((label, description)) 1346 }) 1347 .collect::<Result<Vec<_>, Refusal>>()?; 1348 let list = Shortlist::new(options, shortlist::Params::default()).map_err(|e| refuse(e.to_string()))?; 1349 let full = SHORTLIST_FULL_FUNCTIONS.contains(&call.function); 1350 Ok(ShortlistCall { row: row.clone(), question: question.clone(), list, full }) 1351 }
Two rounds, each one request: the chunks, then the finalists.
Returns the chosen label, or under full the jev_choice_result
of the finalist Choice with that label.
1356 async fn search(self, lookup: &Lookup, sender: &Sender, unsure: bool) -> Result<Option<Out>, Refusal> { 1357 let ShortlistCall { row, question, list, full } = self; 1358 let mut search = shortlist::Search::new(&list); 1359 let mut finalists = None; 1360 // Under `full`, the "does any label fit" Noul rides in the first 1361 // round's request; NULL when the search sends nothing. 1362 let mut gate = full.then(|| list.gate(&question)).flatten(); 1363 let mut fit = None; 1364 loop { 1365 let asks = match search.step() { 1366 shortlist::Step::Done(outcome) => { 1367 let label = list.label(outcome.choice); 1368 if !full { 1369 return Ok(Some(Out::Text(label.to_owned()))); 1370 } 1371 let judged = finalists.unwrap_or_else(|| Judged::only(label.into())); 1372 let mut record = judged.record(); 1373 // The label the search chose, as `jev_choice_shortlist` 1374 // returns it, even where a tie let the vendor name another. 1375 record[0] = Some(Field::Text(label.to_owned())); 1376 record.push(fit.map(Field::Float8)); 1377 return Ok(Some(Out::Record(record))); 1378 } 1379 shortlist::Step::Ask(asks) => asks, 1380 }; 1381 let choices = asks.iter().map(|&a| Question::Choice(list.choice(&search, a, &question))).collect(); 1382 let Some((answered, gated)) = round(lookup, sender, unsure, &row, &question, choices, gate.take()).await? else { 1383 return Ok(None); 1384 }; 1385 fit = fit.or(gated); 1386 for (ask, judged) in asks.into_iter().zip(answered) { 1387 search.answer(ask, &probabilities(&judged)?).map_err(|e| wrong_answer(&e))?; 1388 if ask == shortlist::Ask::Finalists { 1389 finalists = Some(judged); 1390 } 1391 } 1392 } 1393 } 1394}
A finished search as its function returns it: the path, NULL at the
root, or under full the jev_tree_result record, with fit the
gate's probability.
1399fn tree_out(tree: &Tree, outcome: &Outcome, full: bool, fit: Option<f64>) -> Option<Out> { 1400 let (node, probability, separation, depth, leaf) = match *outcome { 1401 Outcome::Leaf { leaf, probability, separation } => (leaf, probability, separation, tree.depth(leaf), true), 1402 Outcome::Stop { node, depth, leaf, probability } => (node, probability, None, depth, leaf), 1403 }; 1404 let path = tree.path(node); 1405 let path = (!path.is_empty()).then(|| path.join(PATH_SEPARATOR)); 1406 if !full { 1407 return path.map(Out::Text); 1408 } 1409 let depth = i32::try_from(depth).expect("a label path is shorter than i32::MAX"); 1410 Some(Out::Record(vec![ 1411 path.map(Field::Text), 1412 Some(Field::Float8(probability)), 1413 separation.map(Field::Float8), 1414 Some(Field::Int4(depth)), 1415 Some(Field::Bool(leaf)), 1416 fit.map(Field::Float8), 1417 ])) 1418}
What a judgement needs to look a question up: the cache, and the
statement's judgments already asked. Cloned into [scan::on_backend]
work by a label-tree search, which looks each round up.
Chooses each row's layout and state (contract Execution).
Whether a miss may be sent: not under jev.cache_only, nor in a
recheck.
1432 sends: bool,
An EvalPlanQual recheck, whose miss is refused as a serialization failure.
1438impl Lookup {
Each row's layout, from the distinct questions asked about it.
1440 fn layouts<'r, 'q>(&self, asks: impl Iterator<Item = (&'r String, &'q Question)>) -> Result<HashMap<&'r String, Planned>, Refusal> { 1441 let invalid = |e| Failure::Client(ClientError::Invalid(e)).refusal(); 1442 let asks = asks.map(|(row, q)| Ok((row, question_bytes(q).map_err(invalid)?))).collect::<Result<Vec<_>, Refusal>>()?; 1443 Ok(self.planner.layouts(asks.iter().map(|(row, q)| (*row, q.as_slice())))) 1444 }
Answers asked about row from the statement's judgments or the
cache, or adds it to row's request, placed under the layout the
planner chose for it. Runs on the backend, so it may query.
question is its text, for errors.
1450 fn one( 1451 &self, 1452 requests: &mut Vec<Request>, 1453 row: &String, 1454 layout: Planned, 1455 question: &str, 1456 asked: Question, 1457 ) -> Result<Answer, Refusal> { 1458 let invalid = |e| Failure::Client(ClientError::Invalid(e)).refusal(); 1459 let placed = self.planner.place(layout, row).map_err(invalid)?; 1460 let bytes = question_bytes(&asked).map_err(invalid)?; 1461 let key = CacheKey::new(&self.scope, &placed, &asked).map_err(invalid)?; 1462 if let Some(shared) = self.asking.0.borrow().get(&key) { 1463 stats::deduped(); 1464 self.counts.deduped.set(self.counts.deduped.get() + 1); 1465 let again = Again { row: row.clone(), placed, text: question.to_owned(), question: asked, cache_key: key, bytes }; 1466 return Ok(Answer::Shared(shared.clone(), again)); 1467 } 1468 let hit = as_owner_of(c"jev_cache", || Spi::get_one_with_args::<pgrx::JsonB>( 1469 // One row always: pgrx's get_one refuses an empty result. 1470 "SELECT (SELECT jsonb_build_object('answer', answer, 'model', answered_model, 1471 'input_tokens', input_tokens, 'output_tokens', output_tokens) 1472 FROM jev_cache WHERE key = $1)", 1473 &[key.as_bytes().to_vec().into()], 1474 )) 1475 .and_then(|r| r.map_err(|e| e.to_string())) 1476 .and_then(|hit| hit.map(|a| Judged::stored(&a.0)).transpose()) 1477 .map_err(|e| Failure::Cache(e).refusal())?; 1478 stats::cached(hit.is_some()); 1479 let counted = if hit.is_some() { &self.counts.hits } else { &self.counts.misses }; 1480 counted.set(counted.get() + 1); 1481 if let Some(p) = hit { 1482 return Ok(Answer::Cached(p)); 1483 } 1484 if self.rechecking { 1485 return Err(Failure::Recheck(question.to_owned()).refusal()); 1486 } 1487 if !self.sends { 1488 return Err(Failure::CacheOnly(question.to_owned()).refusal()); 1489 } 1490 let request = match requests.iter().position(|r| r.row == *row && r.placed.layout() == layout.layout()) { 1491 Some(i) => i, 1492 None => { 1493 requests.push(Request { row: row.clone(), placed, questions: Questions::new(), keys: HashMap::new() }); 1494 requests.len() - 1 1495 } 1496 }; 1497 let r = &mut requests[request]; 1498 let id = if r.keys.is_empty() { "q".to_owned() } else { format!("q{}", r.keys.len()) }; 1499 let k = AnswerKey::ask(&mut r.questions, &id, asked).map_err(invalid)?; 1500 let shared = Rc::new(Shared::default()); 1501 self.asking.0.borrow_mut().insert(key, shared.clone()); 1502 r.keys.insert(key, Pending { key: k, cache_key: key, text: question.to_owned(), bytes, shared }); 1503 Ok(Answer::Asked(request, key)) 1504 } 1505}
A Score of levels, lowest first, in the caller's order. Refused
with 22023 before anything is sent: a NULL level, or a count outside
2 to 10.
1510fn score(question: &str, levels: &[Option<String>]) -> Result<Score, Refusal> { 1511 let refuse = |m: String| Failure::Setting(format!("jev_score: {m}")).refusal(); 1512 let levels = levels 1513 .iter() 1514 .map(|l| l.as_deref().map(Json::text).ok_or_else(|| refuse("a level is NULL".into()))) 1515 .collect::<Result<Vec<_>, _>>()?; 1516 Score::new(Json::text(question), levels).map_err(|e| refuse(e.to_string())) 1517}
A Choice of options, in the caller's order, or the one option there
is. Refused with 22023 before anything is sent: a NULL, empty or
repeated option, or a count outside 1 to 255.
1522fn choice(question: &str, options: &[Option<String>]) -> Result<Choosing, Refusal> { 1523 let refuse = |m: String| Failure::Setting(format!("jev_choice: {m}")).refusal(); 1524 let options = options 1525 .iter() 1526 .map(|o| o.clone().ok_or_else(|| refuse("an option is NULL".into()))) 1527 .collect::<Result<Vec<_>, _>>()?; 1528 match options.as_slice() { 1529 [only] if !only.is_empty() => Ok(Choosing::Only(only.as_str().into())), 1530 _ => Choice::new(Json::text(question), options.into_iter().map(|o| (o, None))) 1531 .map(Choosing::Ask) 1532 .map_err(|e| refuse(e.to_string())), 1533 } 1534}
A single option, answered with probability 1 and never sent, as the vendor's hierarchical cookbook does.
A judgment's answer, as its calls read it: one per cache row, shared by every call asking the same question about the same row.
The probability that the answer is yes.
1555 Noul(f64),
The probability-weighted level, how many levels there are, and the vendor's confidence.
1558 Score { score: f64, levels: usize, confidence: f64 },
The label chosen, and the vendor's confidence.
jev_cache.answer.
1566 answer: serde_json::Value,
None for a judgment nobody answered (a single option).
1573impl Judged {
A Noul's probability of yes.
Reads the cache lookup: jev_cache.answer, as [AnswerKey::read]
writes it, with its row's model and tokens.
1584 fn stored(row: &serde_json::Value) -> Result<Judged, String> { 1585 let answer = &row["answer"]; 1586 let verdict = if let Some(p) = answer["noul"].as_f64() { 1587 Verdict::Noul(p) 1588 } else { 1589 let unreadable = || format!("unreadable cached answer {answer}"); 1590 let confidence = answer["confidence"].as_f64().ok_or_else(unreadable)?; 1591 if let Some(label) = answer["choice"].as_str() { 1592 Verdict::Choice { label: label.into(), confidence } 1593 } else { 1594 match (answer["score"].as_f64(), answer["probabilities"].as_array()) { 1595 (Some(score), Some(ps)) if ps.len() >= 2 => Verdict::Score { score, levels: ps.len(), confidence }, 1596 _ => return Err(unreadable()), 1597 } 1598 } 1599 }; 1600 let full = Full { 1601 answer: answer.clone(), 1602 model: row["model"].as_str().map(str::to_owned), 1603 input_tokens: row["input_tokens"].as_u64().unwrap_or(0), 1604 output_tokens: row["output_tokens"].as_u64().unwrap_or(0), 1605 }; 1606 Ok(Judged { verdict, full: Rc::new(full) }) 1607 }
A Choice's probabilities, in the order its options were sent.
1610 fn choice_probabilities(&self) -> Result<Vec<f64>, String> { 1611 let unreadable = || format!("unreadable Choice probabilities {}", self.full.answer); 1612 self.full.answer["probabilities"] 1613 .as_array() 1614 .ok_or_else(unreadable)? 1615 .iter() 1616 .map(|pair| pair[1].as_f64().ok_or_else(unreadable)) 1617 .collect() 1618 }
The one option there is: probability 1, never sent, so no model and no tokens.
1622 fn only(label: Rc<str>) -> Judged { 1623 let answer = serde_json::json!({ "choice": &*label, "confidence": 1.0, "probabilities": [[&*label, 1.0]] }); 1624 let full = Full { answer, model: None, input_tokens: 0, output_tokens: 0 }; 1625 Judged { verdict: Verdict::Choice { label, confidence: 1.0 }, full: Rc::new(full) } 1626 }
The *_full record, in its composite type's column order.
1629 fn record(&self) -> Vec<Option<Field>> { 1630 let Full { answer, model, input_tokens, output_tokens } = &*self.full; 1631 let json = |k: &str| (!answer[k].is_null()).then(|| Field::Jsonb(answer[k].to_string())); 1632 let tokens = |t: u64| Some(Field::Int4(i32::try_from(t).unwrap_or(i32::MAX))); 1633 let mut fields = match &self.verdict { 1634 Verdict::Noul(p) => vec![Some(Field::Float8(*p))], 1635 Verdict::Score { score, confidence, .. } => { 1636 vec![Some(Field::Float8(*score)), Some(Field::Float8(*confidence)), json("probabilities"), json("legend")] 1637 } 1638 Verdict::Choice { label, confidence } => { 1639 vec![Some(Field::Text(label.to_string())), Some(Field::Float8(*confidence)), json("probabilities")] 1640 } 1641 }; 1642 fields.extend([model.clone().map(Field::Text), tokens(*input_tokens), tokens(*output_tokens)]); 1643 fields 1644 } 1645}
jev_prob.
1651 Probability,
jev(): true at or above the threshold, as pg-jev's >=.
1653 AtLeast(f64),
jev_score.
1655 Level,
jev_score_norm.
1657 Normalized,
jev_choice.
1659 Label,
jev_confidence, of a Choice or a Score.
1661 Confidence,
1666impl Shape {
The cache key covers the question's type, so a call's judgment is always of its own kind.
1669 fn out(self, judged: Judged) -> Out { 1670 match (self, &judged.verdict) { 1671 (Shape::Full, _) => Out::Record(judged.record()), 1672 (Shape::Probability, Verdict::Noul(p)) => Out::Float8(*p), 1673 (Shape::AtLeast(t), Verdict::Noul(p)) => Out::Bool(*p >= t), 1674 (Shape::Level, Verdict::Score { score, .. }) => Out::Float8(*score), 1675 (Shape::Normalized, Verdict::Score { score, levels, .. }) => Out::Float8(score / (levels - 1) as f64), 1676 (Shape::Label, Verdict::Choice { label, .. }) => Out::Text(label.to_string()), 1677 (Shape::Confidence, Verdict::Score { confidence, .. } | Verdict::Choice { confidence, .. }) => { 1678 Out::Float8(*confidence) 1679 } 1680 (_, verdict) => unreachable!("{verdict:?} answered a call of another kind"), 1681 } 1682 } 1683}
A question's key in its response, by kind.
1693impl AnswerKey { 1694 fn ask(questions: &mut Questions, id: &str, question: Question) -> Result<AnswerKey, jev_protocol::ProtocolError> { 1695 match question { 1696 Question::Noul(q) => questions.noul(id, q).map(AnswerKey::Noul), 1697 Question::Score(q) => questions.score(id, q).map(AnswerKey::Score), 1698 Question::Choice(q) => questions.choice(id, q).map(AnswerKey::Choice), 1699 } 1700 }
The judgment, and the full answer jev_cache stores.
1703 fn read(self, response: &Response) -> (Judged, serde_json::Value) { 1704 let (verdict, stored) = match self { 1705 AnswerKey::Noul(k) => { 1706 let p = response.get(k).noul; 1707 (Verdict::Noul(p), serde_json::json!({ "noul": p })) 1708 } 1709 AnswerKey::Score(k) => { 1710 let a = response.get(k); 1711 let verdict = Verdict::Score { score: a.score, levels: a.probabilities.len(), confidence: a.confidence }; 1712 let stored = serde_json::json!({ 1713 "score": a.score, 1714 "confidence": a.confidence, 1715 "probabilities": a.probabilities, 1716 "legend": a.legend, 1717 }); 1718 (verdict, stored) 1719 } 1720 AnswerKey::Choice(k) => { 1721 let a = response.get(k); 1722 // Probabilities as pairs, in the order asked: a jsonb 1723 // object would sort them. 1724 let stored = serde_json::json!({ 1725 "choice": a.choice, 1726 "confidence": a.confidence, 1727 "probabilities": a.probabilities, 1728 }); 1729 (Verdict::Choice { label: a.choice.as_str().into(), confidence: a.confidence }, stored) 1730 } 1731 }; 1732 let usage = response.usage(); 1733 let full = Full { 1734 answer: stored.clone(), 1735 model: Some(response.model().to_owned()), 1736 input_tokens: usage.input_tokens, 1737 output_tokens: usage.output_tokens, 1738 }; 1739 (Judged { verdict, full: Rc::new(full) }, stored) 1740 } 1741}
A row's calls after the cache: the requests still to send, and where each call's answer comes from.
Request index, and the judgment within it.
1755 Asked(usize, CacheKey),
An equal judgment's, asked earlier in this statement; and how to ask it again if that request is dropped.
A shared call's own judgment, to send if the request it waited on is dropped unanswered.
The question's text, for errors.
What a scan's judgements send with: none of it calls Postgres.
1774struct Sender {
None under jev.cache_only.
Sends requests within the budget, serves their waiters, and
queues their receipts.
1794 async fn send(&self, requests: &[Request], mut owed: Owed) -> Result<Vec<jev_client::Answered>, Failure> { 1795 let responses = match self.exchange(requests).await { 1796 Ok(responses) => responses, 1797 Err(e) => { 1798 owed.failed = true; 1799 return Err(e); 1800 } 1801 }; 1802 for (r, answer) in requests.iter().zip(&responses) { 1803 for Pending { key, cache_key, bytes, shared, .. } in r.keys.values() { 1804 let (judged, stored) = key.read(&answer.response); 1805 shared.resolve(Ok(judged)); 1806 self.answered.borrow_mut().push(Receipt { 1807 key: *cache_key, 1808 layout: r.placed.layout(), 1809 state: r.placed.state().as_str().to_owned(), 1810 question: bytes.clone(), 1811 answered_model: answer.response.model().to_owned(), 1812 request_id: answer.request_id.clone(), 1813 usage: answer.response.usage(), 1814 answer: stored, 1815 }); 1816 } 1817 } 1818 // Answered, so its waiters are served. 1819 drop(owed); 1820 Ok(responses) 1821 }
1823 async fn exchange(&self, requests: &[Request]) -> Result<Vec<jev_client::Answered>, Failure> { 1824 if requests.is_empty() { 1825 return Ok(Vec::new()); 1826 } 1827 let Some(client) = &self.client else { 1828 // A miss is refused in lookup; this is a waiter whose request 1829 // was dropped, which under cache_only no scan sends again. 1830 let question = requests[0].keys.values().next().map(|p| p.text.clone()).unwrap_or_default(); 1831 return Err(Failure::CacheOnly(question)); 1832 }; 1833 let chars: usize = requests.iter().map(|r| r.row.len() + r.keys.values().map(|p| p.bytes.len()).sum::<usize>()).sum(); 1834 let mut committed = self.budget.commit(requests.len() as u64, chars)?; 1835 let counts = &self.counts; 1836 counts.requests.set(counts.requests.get() + requests.len() as u64); 1837 let responses = join_all(requests.iter().map(|r| client.ask(r.placed.state(), &r.questions)).collect()).await; 1838 let responses = 1839 responses.into_iter().collect::<Result<Vec<_>, _>>().map_err(Failure::Client)?; 1840 committed.settle(&responses); 1841 let tokens: u64 = responses.iter().map(|a| a.response.usage().input_tokens).sum(); 1842 counts.input_tokens.set(counts.input_tokens.get() + tokens); 1843 Ok(responses) 1844 }
The shared answer; or, when its request was dropped (another scan's rescan, say), this call's own, waiting again on whichever equal call asked first.
1849 async fn wait(&self, mut shared: Rc<Shared>, again: Again) -> Result<Judged, Failure> { 1850 loop { 1851 match shared.wait().await { 1852 Ok(p) => return Ok(p), 1853 Err(Unanswered::Failed) => return Err(Failure::Shared), 1854 Err(Unanswered::Dropped) => {} 1855 } 1856 let asked = self.asking.0.borrow().get(&again.cache_key).cloned(); 1857 if let Some(asked) = asked { 1858 shared = asked; 1859 continue; 1860 } 1861 let invalid = |e| Failure::Client(ClientError::Invalid(e)); 1862 let mut questions = Questions::new(); 1863 let key = AnswerKey::ask(&mut questions, "q", again.question.clone()).map_err(invalid)?; 1864 let mine = Rc::new(Shared::default()); 1865 self.asking.0.borrow_mut().insert(again.cache_key, mine.clone()); 1866 let pending = 1867 Pending { key, cache_key: again.cache_key, text: again.text.clone(), bytes: again.bytes.clone(), shared: mine }; 1868 let request = Request { 1869 row: again.row.clone(), 1870 placed: again.placed.clone(), 1871 questions, 1872 keys: HashMap::from([(again.cache_key, pending)]), 1873 }; 1874 let requests = [request]; 1875 let owed = self.owe(&requests); 1876 let responses = self.send(&requests, owed).await?; 1877 return Ok(key.read(&responses[0].response).0); 1878 } 1879 } 1880}
The row's layout and the request's state, from the planner.
One question of a request.
1893struct Pending {
The question's text, for errors.
1898 text: String,
Its bytes as sent.
1900 bytes: Vec<u8>,
The API key from exactly one source. Both set is an error rather than a precedence rule, which would silently bill the wrong account. Messages name the sources, never the key.
1908fn api_key() -> Result<String, String> { 1909 let env = std::env::var("TYPESAFE_API_KEY").ok().map(|k| k.trim().to_owned()).filter(|k| !k.is_empty()); 1910 match (env, setting(&API_KEY_FILE)) { 1911 (Some(_), Some(_)) => Err( 1912 "the API key is set twice, by TYPESAFE_API_KEY and jev.api_key_file; unset one".into(), 1913 ), 1914 (Some(key), None) => Ok(key), 1915 (None, Some(path)) => { 1916 let contents = 1917 std::fs::read_to_string(&path).map_err(|e| format!("jev.api_key_file {path}: {e}"))?; 1918 let key = contents.trim(); 1919 if key.is_empty() { 1920 return Err(format!("jev.api_key_file {path} is empty")); 1921 } 1922 Ok(key.to_owned()) 1923 } 1924 (None, None) => Err( 1925 "no API key: set TYPESAFE_API_KEY in the server's environment, or jev.api_key_file".into(), 1926 ), 1927 } 1928}
A string setting's value; empty means unset, as for Postgres's own
string GUCs.
Characters per input token for jev.model, learned from what its
cached answers reported (contract Execution): each request's
characters, its state once and each question it asked, against its
usage.input_tokens. Rows without a request id are their own request.
The fallback when the model is unset or nothing is cached.
1937fn learned_ratio() -> TokenRatio { 1938 let Some(model) = setting(&MODEL) else { return TokenRatio::FALLBACK }; 1939 let samples = as_owner_of(c"jev_cache", || Spi::get_one_with_args::<pgrx::JsonB>( 1940 "SELECT coalesce(jsonb_agg(jsonb_build_array(chars, tokens)), '[]') 1941 FROM (SELECT max(octet_length(state)) + sum(octet_length(question)) AS chars, 1942 max(input_tokens) AS tokens 1943 FROM jev_cache WHERE model = $1 1944 GROUP BY coalesce(request_id, encode(key, 'hex'))) r", 1945 &[model.into()], 1946 )) 1947 .and_then(|r| r.map_err(|e| e.to_string())); 1948 let Ok(Some(samples)) = samples else { return TokenRatio::FALLBACK }; 1949 let samples = serde_json::from_value::<Vec<(u64, u64)>>(samples.0).unwrap_or_default(); 1950 TokenRatio::learn(samples.into_iter().map(|(chars, input_tokens)| Sample { chars, input_tokens })) 1951}