Repository navigation
feat(cortex): org:<uid> root — tenant-relative hosted paths and a retired user: root for the transition #244
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
76d93d2
0de446f
8373e34
2393889
a28a629
1c92521
3b315ed
c5c3348
1110ebb
519bf40
6c00104
e7c229b
b5d96e5
ec7c5b2
45a48b8
1c9e9f0
76ed631
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -13,7 +13,7 @@ use tinymemory_api::{ | |
| use super::CortexEngine; | ||
| use super::scopes::KindScope; | ||
| use crate::cortex::envelope::{Decoded, Envelope, decode_event, labels, rebuild, rebuild_whole}; | ||
| use crate::cortex::error::Result; | ||
| use crate::cortex::error::{Error, Result}; | ||
|
|
||
| /// The kinds `filter` admits, in the fixed order | ||
| /// [`ItemKind::ALL`] lists them. | ||
|
|
@@ -168,7 +168,9 @@ impl CortexEngine { | |
| for (id, events) in self.item_events(&scope, &ids).await? { | ||
| let envelopes: Vec<Envelope> = events.into_iter().map(|d| d.envelope).collect(); | ||
| if let Some(item) = rebuild_whole(&envelopes) { | ||
| found.insert(ItemId::new(id.clone()), hit(&id, &item, 0.0)); | ||
| found | ||
| .entry(ItemId::new(id.clone())) | ||
| .or_insert_with(|| hit(&id, &item, 0.0)); | ||
|
Comment on lines
+171
to
+173
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
For a hosted tenant-root request at Useful? React with 👍 / 👎. |
||
| } | ||
| } | ||
| } | ||
|
|
@@ -199,15 +201,31 @@ impl CortexEngine { | |
| } | ||
| let lookups: Vec<(KindScope, Vec<String>)> = by_node | ||
| .into_iter() | ||
| .map(|(namespace, ids)| (KindScope::new(&self.layout, namespace.clone(), kind), ids)) | ||
| .flat_map(|(namespace, ids)| { | ||
| KindScope::read(&self.layout, namespace, kind) | ||
| .into_iter() | ||
| .map(move |scope| (scope, ids.clone())) | ||
| }) | ||
| .collect(); | ||
| let found: Vec<HashMap<String, Vec<Decoded>>> = stream::iter(lookups) | ||
| .map(|(scope, ids)| async move { self.item_events(&scope, &ids).await }) | ||
| let mut found: Vec<(usize, HashMap<String, Vec<Decoded>>)> = stream::iter(lookups) | ||
| .enumerate() | ||
| .map(|(order, (scope, ids))| async move { | ||
| Ok::<_, Error>((order, self.item_events(&scope, &ids).await?)) | ||
| }) | ||
| .buffer_unordered(LOOKUPS_AT_ONCE) | ||
| .try_collect() | ||
| .await?; | ||
| // An item held below both the root and a retired root is rebuilt from | ||
| // the events of one of them, never from both at once: the root's | ||
| // (lookups are in read order, root first), whichever lookup answered | ||
| // first; the retired root's only when the root's copy cannot be | ||
| // rebuilt. | ||
| found.sort_by_key(|(order, _)| *order); | ||
| let mut out = HashMap::new(); | ||
| for (id, events) in found.into_iter().flatten() { | ||
| for (id, events) in found.into_iter().flat_map(|(_, found)| found) { | ||
| if out.contains_key(&id) { | ||
|
senamakel marked this conversation as resolved.
|
||
| continue; | ||
|
senamakel marked this conversation as resolved.
|
||
| } | ||
| let envelopes: Vec<Envelope> = events.into_iter().map(|d| d.envelope).collect(); | ||
| if let Some(item) = rebuild_whole(&envelopes) { | ||
| out.insert(id, item); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -173,16 +173,22 @@ impl CortexEngine { | |
| if len > 0 { | ||
| filled_pages += 1; | ||
| } | ||
| let shadowed = self.shadowed(scope, &page.items).await?; | ||
| for (position, event) in page.items.iter().enumerate().skip(at.offset) { | ||
| at.offset = position + 1; | ||
| let id = event.get("id").and_then(Value::as_str); | ||
| if id.is_some() && id == at.last.as_deref() { | ||
| continue; | ||
| } | ||
| at.last = id.map(str::to_owned); | ||
| if let Some(found) = | ||
| self.admit(kind, &req, event, &mut seen, walk == Walk::Preview) | ||
| { | ||
| if let Some(found) = self.admit( | ||
| kind, | ||
| &req, | ||
| event, | ||
| &mut seen, | ||
| &shadowed, | ||
| walk == Walk::Preview, | ||
| ) { | ||
| pending.push(found); | ||
| if pending.len() == req.limit { | ||
| let exhausted = at.offset == len | ||
|
|
@@ -213,6 +219,25 @@ impl CortexEngine { | |
| Ok((self.resolve(pending).await?, next)) | ||
| } | ||
|
|
||
| /// The ids of the items on `events` (a page of `scope`) that the root | ||
| /// also holds, when `scope` is below the retired root: the root's copy is | ||
| /// the one listed, so an item held below both is one item on every page, | ||
| /// not one per root. Empty for any other scope. | ||
| async fn shadowed(&self, scope: &KindScope, events: &[Value]) -> Result<HashSet<String>> { | ||
| if !self.layout.is_retired(&scope.path) { | ||
| return Ok(HashSet::new()); | ||
| } | ||
| let ids: Vec<String> = events | ||
| .iter() | ||
| .filter_map(decode_event) | ||
| .map(|decoded| decoded.envelope.id) | ||
| .collect::<HashSet<_>>() | ||
| .into_iter() | ||
| .collect(); | ||
| let active = KindScope::new(&self.layout, scope.namespace.clone(), scope.kind); | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Only shadow retired items after confirming the active copy is complete This treats any event found in the active scope as sufficient to shadow the retired copy. During a partial transition, the active scope can contain only some events for a conversation or chunked document while the retired scope still contains the complete item. The listing then suppresses the retired occurrence, and [RULE] incomplete-fallback · |
||
| Ok(self.item_events(&active, &ids).await?.into_keys().collect()) | ||
| } | ||
|
|
||
| /// Whether one raw event starts an item this listing returns; with | ||
| /// `preview`, the item is taken from that event alone. | ||
| fn admit( | ||
|
|
@@ -221,10 +246,11 @@ impl CortexEngine { | |
| req: &ListRequest, | ||
| event: &Value, | ||
| seen: &mut HashSet<String>, | ||
| shadowed: &HashSet<String>, | ||
| preview: bool, | ||
| ) -> Option<Pending> { | ||
| let envelope = decode_event(event)?.envelope; | ||
| if !keeps(&req.filter, kind, &envelope) { | ||
| if !keeps(&req.filter, kind, &envelope) || shadowed.contains(&envelope.id) { | ||
| return None; | ||
| } | ||
| let starts = envelope.part().is_none_or(|index| index == 0); | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.