1use crate::encoding::{decode_inline, decode_row, make_triple, tag_of, Tag, DEFAULT_GRAPH_ID};
10use crate::error::{Error, Result};
11use crate::sql::{col, Capabilities, Request, Response, Statement};
12use oxrdf::Term;
13use std::collections::{BTreeSet, HashMap};
14
15const CHUNK: usize = 400;
16
17#[derive(Debug, Clone, Copy)]
18enum Lookup {
19 Terms,
20 Triples,
21}
22
23#[derive(Debug, Default)]
25pub struct TermResolver {
26 known: HashMap<i64, Term>,
27 pending: BTreeSet<i64>,
28 triples: HashMap<i64, [i64; 3]>,
29 in_flight: Vec<Lookup>,
30}
31
32impl TermResolver {
33 pub fn with_constants(constants: HashMap<i64, Term>) -> Self {
34 Self {
35 known: constants,
36 ..Self::default()
37 }
38 }
39
40 pub fn want(&mut self, id: i64) {
42 if id == DEFAULT_GRAPH_ID || self.known.contains_key(&id) || self.triples.contains_key(&id)
43 {
44 return;
45 }
46 if let Some(t) = decode_inline(id) {
47 self.known.insert(id, t);
48 return;
49 }
50 self.pending.insert(id);
51 }
52
53 fn id_col(caps: &Capabilities, c: &str) -> String {
54 if caps.int64_as_text {
55 format!("CAST({c} AS TEXT)")
56 } else {
57 c.into()
58 }
59 }
60
61 pub fn request(&mut self, caps: &Capabilities) -> Option<Request> {
63 if self.pending.is_empty() {
64 return None;
65 }
66 let pending = std::mem::take(&mut self.pending);
67 let (triples, terms): (Vec<i64>, Vec<i64>) = pending
68 .into_iter()
69 .partition(|id| tag_of(*id) == Some(Tag::Triple));
70 let mut stmts = Vec::new();
71 self.in_flight.clear();
72 for chunk in terms.chunks(CHUNK) {
73 stmts.push(Statement::new(format!(
74 "SELECT {}, lex, dt, lang, dir FROM terms WHERE id IN ({})",
75 Self::id_col(caps, "id"),
76 join(chunk)
77 )));
78 self.in_flight.push(Lookup::Terms);
79 }
80 for chunk in triples.chunks(CHUNK) {
81 stmts.push(Statement::new(format!(
82 "SELECT {}, {}, {}, {} FROM triple_terms WHERE id IN ({})",
83 Self::id_col(caps, "id"),
84 Self::id_col(caps, "s"),
85 Self::id_col(caps, "p"),
86 Self::id_col(caps, "o"),
87 join(chunk)
88 )));
89 self.in_flight.push(Lookup::Triples);
90 }
91 Some(Request::read(stmts))
92 }
93
94 pub fn absorb(&mut self, response: Response) -> Result<()> {
96 let kinds = std::mem::take(&mut self.in_flight);
97 for (kind, rs) in kinds.into_iter().zip(response) {
98 for row in rs.rows {
99 let id = col(&row, 0)?
100 .as_i64()
101 .ok_or_else(|| Error::corrupted("bad id in lookup"))?;
102 match kind {
103 Lookup::Terms => {
104 let mut it = row.into_iter().skip(1);
105 let lex = it.next().and_then(|v| v.into_string()).unwrap_or_default();
106 let dt = it.next().and_then(|v| v.into_string());
107 let lang = it.next().and_then(|v| v.into_string());
108 let dir = it.next().and_then(|v| v.as_i64());
109 self.known.insert(id, decode_row(id, lex, dt, lang, dir)?);
110 }
111 Lookup::Triples => {
112 let get = |i| {
113 col(&row, i)?
114 .as_i64()
115 .ok_or_else(|| Error::corrupted("bad triple term row"))
116 };
117 let parts = [get(1)?, get(2)?, get(3)?];
118 self.triples.insert(id, parts);
119 for p in parts {
120 self.want(p);
121 }
122 }
123 }
124 }
125 }
126 self.assemble();
127 Ok(())
128 }
129
130 fn assemble(&mut self) {
131 loop {
132 let ready: Vec<i64> = self
133 .triples
134 .iter()
135 .filter(|(_, parts)| parts.iter().all(|p| self.known.contains_key(p)))
136 .map(|(id, _)| *id)
137 .collect();
138 if ready.is_empty() {
139 return;
140 }
141 for id in ready {
142 let [s, p, o] = self.triples.remove(&id).expect("present");
143 if let Ok(t) = make_triple(
144 self.known[&s].clone(),
145 self.known[&p].clone(),
146 self.known[&o].clone(),
147 ) {
148 self.known.insert(id, t.into());
149 }
150 }
151 }
152 }
153
154 pub fn is_complete(&self) -> bool {
156 self.pending.is_empty()
157 }
158
159 pub fn get(&self, id: i64) -> Result<Term> {
161 self.known
162 .get(&id)
163 .cloned()
164 .or_else(|| decode_inline(id))
165 .ok_or_else(|| Error::corrupted(format!("term {id} not found in the dictionary")))
166 }
167}
168
169pub(crate) fn join(ids: &[i64]) -> String {
170 ids.iter()
171 .map(ToString::to_string)
172 .collect::<Vec<_>>()
173 .join(",")
174}
175
176pub(crate) fn ids_of(response: &Response, i: usize) -> Vec<i64> {
178 response
179 .get(i)
180 .map(|rs| {
181 rs.rows
182 .iter()
183 .filter_map(|r| r.first().and_then(|v| v.as_i64()))
184 .collect()
185 })
186 .unwrap_or_default()
187}