1use std::collections::BTreeMap;
2
3use crate::{
4 Result,
5 backend::{CacheScope, RecordView, StoreStateLifecycle, StoreStateRequest},
6 entry::ID,
7 instance::backend::Backend,
8};
9
10use super::ProjectionDescriptor;
11use super::RecordProjection;
12use crate::crdt::Doc;
13
14pub const OPAQUE_STATE_KEY: &[u8] = &[0x00];
16
17pub(crate) fn opaque_request(
18 database: &ID,
19 store: &str,
20 descriptor: ProjectionDescriptor,
21 source_key: Vec<u8>,
22 scope: CacheScope,
23) -> StoreStateRequest {
24 StoreStateRequest {
25 database: database.clone(),
26 store: store.to_string(),
27 lifecycle: StoreStateLifecycle::Derived,
28 scope,
29 projection: descriptor,
30 source_key,
31 }
32}
33
34pub(crate) fn records_request(
35 database: &ID,
36 store: &str,
37 descriptor: ProjectionDescriptor,
38 source_key: Vec<u8>,
39 scope: CacheScope,
40) -> StoreStateRequest {
41 opaque_request(database, store, descriptor, source_key, scope)
42}
43
44pub(crate) async fn publish_records<'a>(
45 backend: &dyn Backend,
46 request: StoreStateRequest,
47 deltas: impl Iterator<Item = &'a [u8]>,
48 projection: &dyn RecordProjection<Doc>,
49) -> Result<RecordView> {
50 let token = backend.begin_store_state_staging(request).await?;
51 let result = async {
52 let mut records = BTreeMap::new();
53 for bytes in deltas {
54 let delta: Doc = serde_json::from_slice(bytes)?;
55 projection.project_delta(&delta, &mut records)?;
56 }
57 records.retain(|_, value| value.is_some());
58 if !records.is_empty() {
59 backend.stage_store_state_records(&token, records).await?;
60 }
61 backend.publish_store_state(token.clone()).await
62 }
63 .await;
64 if result.is_err() {
65 let _ = backend.abort_store_state(token).await;
66 }
67 result
68}
69
70pub(crate) async fn load_opaque(
71 backend: &dyn Backend,
72 view: &RecordView,
73) -> Result<Option<Vec<u8>>> {
74 backend.store_state_record_get(view, OPAQUE_STATE_KEY).await
75}
76
77pub(crate) async fn publish_opaque(
78 backend: &dyn Backend,
79 request: StoreStateRequest,
80 bytes: Vec<u8>,
81) -> Result<RecordView> {
82 let token = backend.begin_store_state_staging(request).await?;
83 let result = async {
84 backend
85 .stage_store_state_records(
86 &token,
87 BTreeMap::from([(OPAQUE_STATE_KEY.to_vec(), Some(bytes))]),
88 )
89 .await?;
90 backend.publish_store_state(token.clone()).await
91 }
92 .await;
93 if result.is_err() {
94 let _ = backend.abort_store_state(token).await;
95 }
96 result
97}
98
99pub(crate) async fn resolve_cached(
103 backend: &dyn Backend,
104 request: &StoreStateRequest,
105) -> Result<Option<RecordView>> {
106 match backend.resolve_store_state(request).await {
107 Err(err) if err.is_unsupported_store_state() => Ok(None),
108 result => result,
109 }
110}
111
112pub(crate) async fn load_cached(
120 backend: &dyn Backend,
121 request: &StoreStateRequest,
122) -> Result<Option<Vec<u8>>> {
123 for _ in 0..2 {
124 let Some(view) = resolve_cached(backend, request).await? else {
125 return Ok(None);
126 };
127 match load_opaque(backend, &view).await {
128 Err(err) if err.is_invalid_store_state_view() => continue,
129 result => return result,
130 }
131 }
132 Ok(None)
133}
134
135pub(crate) async fn store_cached(
140 backend: &dyn Backend,
141 request: StoreStateRequest,
142 bytes: Vec<u8>,
143) -> Result<()> {
144 match publish_opaque(backend, request, bytes).await {
145 Err(err) if err.is_unsupported_store_state() => Ok(()),
146 result => result.map(|_| ()),
147 }
148}