1pub const INIT_SCHEMA_QUERY: &str = r#"
2 DEFINE TABLE IF NOT EXISTS temporal_node SCHEMAFULL;
3 DEFINE FIELD IF NOT EXISTS tenant_id ON temporal_node TYPE string;
4 DEFINE FIELD IF NOT EXISTS session_id ON temporal_node TYPE string;
5 DEFINE FIELD IF NOT EXISTS raw ON temporal_node TYPE string;
6 DEFINE FIELD IF NOT EXISTS tier ON temporal_node TYPE string;
7 DEFINE FIELD IF NOT EXISTS timestamp ON temporal_node TYPE datetime;
8 DEFINE FIELD IF NOT EXISTS compression_depth ON temporal_node TYPE int;
9 DEFINE FIELD IF NOT EXISTS parent_node_id ON temporal_node TYPE option<string>;
10 DEFINE FIELD IF NOT EXISTS sync_key ON temporal_node TYPE string;
11 DEFINE FIELD IF NOT EXISTS updated_at ON temporal_node TYPE datetime;
12 DEFINE FIELD IF NOT EXISTS source_metadata ON temporal_node TYPE option<object>;
13 DEFINE FIELD IF NOT EXISTS context_summary ON temporal_node TYPE option<string>;
14 DEFINE FIELD IF NOT EXISTS semantic_tags ON temporal_node TYPE option<array<string>>;
15 DEFINE FIELD IF NOT EXISTS semantic_links ON temporal_node TYPE option<array<object>>;
16 DEFINE FIELD IF NOT EXISTS embedding ON temporal_node TYPE option<array<float>>;
17 DEFINE FIELD IF NOT EXISTS embedding_model ON temporal_node TYPE option<string>;
18 DEFINE FIELD IF NOT EXISTS embedding_dimensions ON temporal_node TYPE option<int>;
19 DEFINE FIELD IF NOT EXISTS embedded_at ON temporal_node TYPE option<datetime>;
20 DEFINE FIELD IF NOT EXISTS psi ON temporal_node TYPE float;
21 DEFINE FIELD IF NOT EXISTS rho ON temporal_node TYPE float;
22 DEFINE FIELD IF NOT EXISTS kappa ON temporal_node TYPE float;
23 DEFINE FIELD IF NOT EXISTS user_stability ON temporal_node TYPE float;
24 DEFINE FIELD IF NOT EXISTS user_friction ON temporal_node TYPE float;
25 DEFINE FIELD IF NOT EXISTS user_logic ON temporal_node TYPE float;
26 DEFINE FIELD IF NOT EXISTS user_autonomy ON temporal_node TYPE float;
27 DEFINE FIELD IF NOT EXISTS user_psi ON temporal_node TYPE float;
28 DEFINE FIELD IF NOT EXISTS model_stability ON temporal_node TYPE float;
29 DEFINE FIELD IF NOT EXISTS model_friction ON temporal_node TYPE float;
30 DEFINE FIELD IF NOT EXISTS model_logic ON temporal_node TYPE float;
31 DEFINE FIELD IF NOT EXISTS model_autonomy ON temporal_node TYPE float;
32 DEFINE FIELD IF NOT EXISTS model_psi ON temporal_node TYPE float;
33 DEFINE FIELD IF NOT EXISTS comp_stability ON temporal_node TYPE float;
34 DEFINE FIELD IF NOT EXISTS comp_friction ON temporal_node TYPE float;
35 DEFINE FIELD IF NOT EXISTS comp_logic ON temporal_node TYPE float;
36 DEFINE FIELD IF NOT EXISTS comp_autonomy ON temporal_node TYPE float;
37 DEFINE FIELD IF NOT EXISTS comp_psi ON temporal_node TYPE float;
38
39 DEFINE TABLE IF NOT EXISTS semantic_tag_index SCHEMAFULL;
40 DEFINE FIELD IF NOT EXISTS tenant_id ON semantic_tag_index TYPE string;
41 DEFINE FIELD IF NOT EXISTS session_id ON semantic_tag_index TYPE string;
42 DEFINE FIELD IF NOT EXISTS node_id ON semantic_tag_index TYPE string;
43 DEFINE FIELD IF NOT EXISTS sync_key ON semantic_tag_index TYPE string;
44 DEFINE FIELD IF NOT EXISTS tag ON semantic_tag_index TYPE string;
45 DEFINE FIELD IF NOT EXISTS embedding ON semantic_tag_index TYPE option<array<float>>;
46 DEFINE FIELD IF NOT EXISTS embedding_model ON semantic_tag_index TYPE option<string>;
47 DEFINE FIELD IF NOT EXISTS embedding_dimensions ON semantic_tag_index TYPE option<int>;
48 DEFINE FIELD IF NOT EXISTS embedded_at ON semantic_tag_index TYPE option<datetime>;
49 DEFINE FIELD IF NOT EXISTS updated_at ON semantic_tag_index TYPE datetime;
50
51 DEFINE TABLE IF NOT EXISTS calibration SCHEMAFULL;
52 DEFINE FIELD IF NOT EXISTS tenant_id ON calibration TYPE string;
53 DEFINE FIELD IF NOT EXISTS session_id ON calibration TYPE string;
54 DEFINE FIELD IF NOT EXISTS stability ON calibration TYPE float;
55 DEFINE FIELD IF NOT EXISTS friction ON calibration TYPE float;
56 DEFINE FIELD IF NOT EXISTS logic ON calibration TYPE float;
57 DEFINE FIELD IF NOT EXISTS autonomy ON calibration TYPE float;
58 DEFINE FIELD IF NOT EXISTS psi ON calibration TYPE float;
59 DEFINE FIELD IF NOT EXISTS trigger ON calibration TYPE string;
60 DEFINE FIELD IF NOT EXISTS created_at ON calibration TYPE datetime;
61
62 DEFINE TABLE IF NOT EXISTS sync_checkpoint SCHEMAFULL;
63 DEFINE FIELD IF NOT EXISTS tenant_id ON sync_checkpoint TYPE string;
64 DEFINE FIELD IF NOT EXISTS session_id ON sync_checkpoint TYPE string;
65 DEFINE FIELD IF NOT EXISTS connector_id ON sync_checkpoint TYPE string;
66 DEFINE FIELD IF NOT EXISTS cursor_updated_at ON sync_checkpoint TYPE option<datetime>;
67 DEFINE FIELD IF NOT EXISTS cursor_sync_key ON sync_checkpoint TYPE option<string>;
68 DEFINE FIELD IF NOT EXISTS metadata ON sync_checkpoint TYPE option<object>;
69 DEFINE FIELD IF NOT EXISTS updated_at ON sync_checkpoint TYPE datetime;
70
71 DEFINE INDEX IF NOT EXISTS idx_node_session ON temporal_node FIELDS session_id;
72 DEFINE INDEX IF NOT EXISTS idx_node_tenant_session ON temporal_node FIELDS tenant_id, session_id;
73 DEFINE INDEX IF NOT EXISTS idx_node_tier ON temporal_node FIELDS tier;
74 DEFINE INDEX IF NOT EXISTS idx_node_timestamp ON temporal_node FIELDS timestamp;
75 DEFINE INDEX IF NOT EXISTS idx_node_change_cursor ON temporal_node FIELDS tenant_id, session_id, updated_at, sync_key;
76 DEFINE INDEX IF NOT EXISTS idx_node_sync_identity ON temporal_node FIELDS tenant_id, session_id, sync_key UNIQUE;
77 DEFINE INDEX IF NOT EXISTS idx_tag_lookup ON semantic_tag_index FIELDS tenant_id, tag;
78 DEFINE INDEX IF NOT EXISTS idx_tag_node ON semantic_tag_index FIELDS tenant_id, node_id;
79 DEFINE INDEX IF NOT EXISTS idx_tag_sync_identity ON semantic_tag_index FIELDS tenant_id, sync_key, tag UNIQUE;
80 DEFINE INDEX IF NOT EXISTS idx_cal_session ON calibration FIELDS session_id;
81 DEFINE INDEX IF NOT EXISTS idx_cal_tenant_session ON calibration FIELDS tenant_id, session_id;
82 DEFINE INDEX IF NOT EXISTS idx_checkpoint_scope ON sync_checkpoint FIELDS tenant_id, session_id, connector_id UNIQUE;
83 SELECT * FROM calibration LIMIT 0;
84 "#;
85
86pub fn query_nodes_query(where_clause: &str, capped_limit: usize) -> String {
87 format!(
88 r#"
89 SELECT
90 tenant_id AS TenantId,
91 session_id AS SessionId,
92 raw AS Raw,
93 tier AS Tier,
94 timestamp AS Timestamp,
95 compression_depth AS CompressionDepth,
96 parent_node_id AS ParentNodeId,
97 sync_key AS SyncKey,
98 updated_at AS UpdatedAt,
99 source_metadata AS SourceMetadata,
100 context_summary AS ContextSummary,
101 semantic_tags AS SemanticTags,
102 semantic_links AS SemanticLinks,
103 embedding AS Embedding,
104 embedding_model AS EmbeddingModel,
105 embedding_dimensions AS EmbeddingDimensions,
106 embedded_at AS EmbeddedAt,
107 psi AS Psi,
108 rho AS Rho,
109 kappa AS Kappa,
110 user_stability AS UserStability,
111 user_friction AS UserFriction,
112 user_logic AS UserLogic,
113 user_autonomy AS UserAutonomy,
114 user_psi AS UserPsi,
115 model_stability AS ModelStability,
116 model_friction AS ModelFriction,
117 model_logic AS ModelLogic,
118 model_autonomy AS ModelAutonomy,
119 model_psi AS ModelPsi,
120 comp_stability AS CompStability,
121 comp_friction AS CompFriction,
122 comp_logic AS CompLogic,
123 comp_autonomy AS CompAutonomy,
124 comp_psi AS CompPsi,
125 0 AS ResonanceDelta
126 FROM temporal_node
127 {where_clause}
128 ORDER BY Timestamp DESC
129 LIMIT {capped_limit};
130 "#
131 )
132}
133
134pub fn create_temporal_node_query(
135 record_id: &str,
136 include_parent_assignment: bool,
137 include_source_metadata_assignment: bool,
138 include_embedding_assignment: bool,
139 include_context_summary_assignment: bool,
140 include_embedding_vector_assignment: bool,
141 include_embedding_model_assignment: bool,
142 include_embedding_dimensions_assignment: bool,
143 include_embedded_at_assignment: bool,
144 include_semantic_tags_assignment: bool,
145 include_semantic_links_assignment: bool,
146) -> String {
147 let parent_assignment = if include_parent_assignment {
148 "\n parent_node_id = $parent_node_id,"
149 } else {
150 ""
151 };
152
153 let source_metadata_assignment = if include_source_metadata_assignment {
154 "\n source_metadata = $source_metadata,"
155 } else {
156 ""
157 };
158
159 let context_summary_assignment = if include_embedding_assignment {
160 let context_summary_value = if include_context_summary_assignment {
161 "$context_summary"
162 } else {
163 "NONE"
164 };
165 let embedding_value = if include_embedding_vector_assignment {
166 "$embedding"
167 } else {
168 "NONE"
169 };
170 let embedding_model_value = if include_embedding_model_assignment {
171 "$embedding_model"
172 } else {
173 "NONE"
174 };
175 let embedding_dimensions_value = if include_embedding_dimensions_assignment {
176 "$embedding_dimensions"
177 } else {
178 "NONE"
179 };
180 let embedded_at_assignment = if include_embedded_at_assignment {
181 "<datetime>$embedded_at"
182 } else {
183 "NONE"
184 };
185
186 format!(
187 "\n context_summary = {context_summary_value},\n embedding = {embedding_value},\n embedding_model = {embedding_model_value},\n embedding_dimensions = {embedding_dimensions_value},\n embedded_at = {embedded_at_assignment},"
188 )
189 } else {
190 "\n context_summary = NONE,\n embedding = NONE,\n embedding_model = NONE,\n embedding_dimensions = NONE,\n embedded_at = NONE,".to_string()
191 };
192
193 let semantic_tags_value = if include_semantic_tags_assignment {
194 "$semantic_tags"
195 } else {
196 "NONE"
197 };
198 let semantic_links_value = if include_semantic_links_assignment {
199 "$semantic_links"
200 } else {
201 "NONE"
202 };
203
204 format!(
205 r#"
206 CREATE temporal_node:`{record_id}` SET
207 tenant_id = $tenant_id,
208 session_id = $session_id,
209 raw = $raw,
210 tier = $tier,
211 timestamp = <datetime>$timestamp,
212 compression_depth = $compression_depth,{parent_assignment}
213 sync_key = $sync_key,
214 updated_at = <datetime>$updated_at,{source_metadata_assignment}
215 {context_summary_assignment}
216 psi = $psi,
217 rho = $rho,
218 kappa = $kappa,
219 user_stability = $user_stability,
220 user_friction = $user_friction,
221 user_logic = $user_logic,
222 user_autonomy = $user_autonomy,
223 user_psi = $user_psi,
224 model_stability = $model_stability,
225 model_friction = $model_friction,
226 model_logic = $model_logic,
227 model_autonomy = $model_autonomy,
228 model_psi = $model_psi,
229 comp_stability = $comp_stability,
230 comp_friction = $comp_friction,
231 comp_logic = $comp_logic,
232 comp_autonomy = $comp_autonomy,
233 comp_psi = $comp_psi,
234 semantic_tags = {semantic_tags_value},
235 semantic_links = {semantic_links_value};
236 "#
237 )
238}
239
240pub fn get_by_resonance_query(
241 current_stability: f32,
242 current_friction: f32,
243 current_logic: f32,
244 current_autonomy: f32,
245 additional_predicate: &str,
246 limit: usize,
247) -> String {
248 let stability = format!("{current_stability:.4}");
249 let friction = format!("{current_friction:.4}");
250 let logic = format!("{current_logic:.4}");
251 let autonomy = format!("{current_autonomy:.4}");
252 let where_suffix = if additional_predicate.trim().is_empty() {
253 String::new()
254 } else {
255 format!(" AND {}", additional_predicate.trim())
256 };
257
258 format!(
259 r#"
260 SELECT
261 tenant_id AS TenantId,
262 session_id AS SessionId,
263 raw AS Raw,
264 tier AS Tier,
265 timestamp AS Timestamp,
266 compression_depth AS CompressionDepth,
267 parent_node_id AS ParentNodeId,
268 sync_key AS SyncKey,
269 updated_at AS UpdatedAt,
270 source_metadata AS SourceMetadata,
271 context_summary AS ContextSummary,
272 semantic_tags AS SemanticTags,
273 semantic_links AS SemanticLinks,
274 embedding AS Embedding,
275 embedding_model AS EmbeddingModel,
276 embedding_dimensions AS EmbeddingDimensions,
277 embedded_at AS EmbeddedAt,
278 psi AS Psi,
279 rho AS Rho,
280 kappa AS Kappa,
281 user_stability AS UserStability,
282 user_friction AS UserFriction,
283 user_logic AS UserLogic,
284 user_autonomy AS UserAutonomy,
285 user_psi AS UserPsi,
286 model_stability AS ModelStability,
287 model_friction AS ModelFriction,
288 model_logic AS ModelLogic,
289 model_autonomy AS ModelAutonomy,
290 model_psi AS ModelPsi,
291 comp_stability AS CompStability,
292 comp_friction AS CompFriction,
293 comp_logic AS CompLogic,
294 comp_autonomy AS CompAutonomy,
295 comp_psi AS CompPsi,
296 (
297 math::abs(model_stability - {stability})
298 + math::abs(model_friction - {friction})
299 + math::abs(model_logic - {logic})
300 + math::abs(model_autonomy - {autonomy})
301 ) / 4.0 AS ResonanceDelta
302 FROM temporal_node
303 WHERE session_id = $session_id
304 AND (tenant_id = $tenant_id OR tenant_id = NONE OR tenant_id = '')
305 {where_suffix}
306 ORDER BY ResonanceDelta ASC
307 LIMIT {limit};
308 "#
309 )
310}
311
312pub fn get_by_resonance_global_query(
313 current_stability: f32,
314 current_friction: f32,
315 current_logic: f32,
316 current_autonomy: f32,
317 additional_predicate: &str,
318 limit: usize,
319) -> String {
320 let stability = format!("{current_stability:.4}");
321 let friction = format!("{current_friction:.4}");
322 let logic = format!("{current_logic:.4}");
323 let autonomy = format!("{current_autonomy:.4}");
324 let where_clause = if additional_predicate.trim().is_empty() {
325 String::new()
326 } else {
327 format!("WHERE {}", additional_predicate.trim())
328 };
329
330 format!(
331 r#"
332 SELECT
333 tenant_id AS TenantId,
334 session_id AS SessionId,
335 raw AS Raw,
336 tier AS Tier,
337 timestamp AS Timestamp,
338 compression_depth AS CompressionDepth,
339 parent_node_id AS ParentNodeId,
340 sync_key AS SyncKey,
341 updated_at AS UpdatedAt,
342 source_metadata AS SourceMetadata,
343 context_summary AS ContextSummary,
344 semantic_tags AS SemanticTags,
345 semantic_links AS SemanticLinks,
346 embedding AS Embedding,
347 embedding_model AS EmbeddingModel,
348 embedding_dimensions AS EmbeddingDimensions,
349 embedded_at AS EmbeddedAt,
350 psi AS Psi,
351 rho AS Rho,
352 kappa AS Kappa,
353 user_stability AS UserStability,
354 user_friction AS UserFriction,
355 user_logic AS UserLogic,
356 user_autonomy AS UserAutonomy,
357 user_psi AS UserPsi,
358 model_stability AS ModelStability,
359 model_friction AS ModelFriction,
360 model_logic AS ModelLogic,
361 model_autonomy AS ModelAutonomy,
362 model_psi AS ModelPsi,
363 comp_stability AS CompStability,
364 comp_friction AS CompFriction,
365 comp_logic AS CompLogic,
366 comp_autonomy AS CompAutonomy,
367 comp_psi AS CompPsi,
368 (
369 math::abs(model_stability - {stability})
370 + math::abs(model_friction - {friction})
371 + math::abs(model_logic - {logic})
372 + math::abs(model_autonomy - {autonomy})
373 ) / 4.0 AS ResonanceDelta
374 FROM temporal_node
375 {where_clause}
376 ORDER BY ResonanceDelta ASC
377 LIMIT {limit};
378 "#
379 )
380}
381
382pub const GET_LAST_AVEC_QUERY: &str = r#"
383 SELECT stability, friction, logic, autonomy, psi, created_at
384 FROM calibration
385 WHERE session_id = $session_id
386 AND (tenant_id = $tenant_id OR tenant_id = NONE OR tenant_id = '')
387 ORDER BY created_at DESC
388 LIMIT 1;
389 "#;
390
391pub const GET_TRIGGER_HISTORY_QUERY: &str = r#"
392 SELECT trigger, created_at FROM calibration
393 WHERE session_id = $session_id
394 AND (tenant_id = $tenant_id OR tenant_id = NONE OR tenant_id = '')
395 ORDER BY created_at ASC;
396 "#;
397
398pub const STORE_CALIBRATION_QUERY: &str = r#"
399 CREATE calibration SET
400 tenant_id = $tenant_id,
401 session_id = $session_id,
402 stability = $stability,
403 friction = $friction,
404 logic = $logic,
405 autonomy = $autonomy,
406 psi = $psi,
407 trigger = $trigger,
408 created_at = <datetime>$created_at;
409 "#;
410
411pub const SELECT_TEMPORAL_NODE_LEGACY_SYNC_QUERY: &str = r#"
412 SELECT id, session_id, timestamp, sync_key, updated_at
413 FROM temporal_node
414 WHERE tenant_id = NONE OR tenant_id = '' OR sync_key = NONE OR sync_key = '' OR updated_at = NONE;
415 "#;
416
417pub fn update_temporal_node_legacy_sync_query(record_id: &str) -> String {
418 format!(
419 r#"
420 UPDATE temporal_node:`{record_id}` SET
421 tenant_id = $tenant_id,
422 sync_key = $sync_key,
423 updated_at = <datetime>$updated_at;
424 "#
425 )
426}
427
428pub const FIND_EXISTING_NODE_BY_SYNC_KEY_QUERY: &str = r#"
429 SELECT
430 id AS Id,
431 source_metadata AS SourceMetadata,
432 context_summary AS ContextSummary,
433 semantic_tags AS SemanticTags,
434 semantic_links AS SemanticLinks,
435 embedding AS Embedding,
436 embedding_model AS EmbeddingModel,
437 embedding_dimensions AS EmbeddingDimensions,
438 embedded_at AS EmbeddedAt
439 FROM temporal_node
440 WHERE session_id = $session_id
441 AND sync_key = $sync_key
442 AND (tenant_id = $tenant_id OR tenant_id = NONE OR tenant_id = '')
443 LIMIT 1;
444 "#;
445
446pub const FIND_EXISTING_NODE_BY_SYNC_KEY_EXACT_QUERY: &str = r#"
447 SELECT
448 id AS Id,
449 source_metadata AS SourceMetadata,
450 context_summary AS ContextSummary,
451 semantic_tags AS SemanticTags,
452 semantic_links AS SemanticLinks,
453 embedding AS Embedding,
454 embedding_model AS EmbeddingModel,
455 embedding_dimensions AS EmbeddingDimensions,
456 embedded_at AS EmbeddedAt
457 FROM temporal_node
458 WHERE session_id = $session_id
459 AND sync_key = $sync_key
460 AND tenant_id = $tenant_id
461 LIMIT 1;
462 "#;
463
464pub const FIND_EXISTING_NODE_BY_SYNC_KEY_ANY_TENANT_QUERY: &str = r#"
465 SELECT
466 id AS Id,
467 source_metadata AS SourceMetadata,
468 context_summary AS ContextSummary,
469 semantic_tags AS SemanticTags,
470 semantic_links AS SemanticLinks,
471 embedding AS Embedding,
472 embedding_model AS EmbeddingModel,
473 embedding_dimensions AS EmbeddingDimensions,
474 embedded_at AS EmbeddedAt
475 FROM temporal_node
476 WHERE session_id = $session_id
477 AND sync_key = $sync_key
478 LIMIT 1;
479 "#;
480
481pub fn update_temporal_node_sync_metadata_query(
482 record_id: &str,
483 clear_source_metadata: bool,
484) -> String {
485 let source_metadata_assignment = if clear_source_metadata {
486 "source_metadata = NONE,"
487 } else {
488 "source_metadata = $source_metadata,"
489 };
490
491 format!(
492 r#"
493 UPDATE temporal_node:`{record_id}` SET
494 {source_metadata_assignment}
495 updated_at = <datetime>$updated_at;
496 "#
497 )
498}
499
500pub fn update_temporal_node_query(
501 record_id: &str,
502 include_parent_assignment: bool,
503 include_source_metadata_assignment: bool,
504 include_embedding_assignment: bool,
505 include_context_summary_assignment: bool,
506 include_embedding_vector_assignment: bool,
507 include_embedding_model_assignment: bool,
508 include_embedding_dimensions_assignment: bool,
509 include_embedded_at_assignment: bool,
510 include_semantic_tags_assignment: bool,
511 include_semantic_links_assignment: bool,
512) -> String {
513 let parent_assignment = if include_parent_assignment {
514 "\n parent_node_id = $parent_node_id,"
515 } else {
516 "\n parent_node_id = NONE,"
517 };
518
519 let source_metadata_assignment = if include_source_metadata_assignment {
520 "\n source_metadata = $source_metadata,"
521 } else {
522 "\n source_metadata = NONE,"
523 };
524
525 let context_summary_assignment = if include_embedding_assignment {
526 let context_summary_value = if include_context_summary_assignment {
527 "$context_summary"
528 } else {
529 "NONE"
530 };
531 let embedding_value = if include_embedding_vector_assignment {
532 "$embedding"
533 } else {
534 "NONE"
535 };
536 let embedding_model_value = if include_embedding_model_assignment {
537 "$embedding_model"
538 } else {
539 "NONE"
540 };
541 let embedding_dimensions_value = if include_embedding_dimensions_assignment {
542 "$embedding_dimensions"
543 } else {
544 "NONE"
545 };
546 let embedded_at_assignment = if include_embedded_at_assignment {
547 "<datetime>$embedded_at"
548 } else {
549 "NONE"
550 };
551
552 format!(
553 "\n context_summary = {context_summary_value},\n embedding = {embedding_value},\n embedding_model = {embedding_model_value},\n embedding_dimensions = {embedding_dimensions_value},\n embedded_at = {embedded_at_assignment},"
554 )
555 } else {
556 "\n context_summary = NONE,\n embedding = NONE,\n embedding_model = NONE,\n embedding_dimensions = NONE,\n embedded_at = NONE,"
557 .to_string()
558 };
559
560 let semantic_tags_value = if include_semantic_tags_assignment {
561 "$semantic_tags"
562 } else {
563 "NONE"
564 };
565 let semantic_links_value = if include_semantic_links_assignment {
566 "$semantic_links"
567 } else {
568 "NONE"
569 };
570
571 format!(
572 r#"
573 UPDATE temporal_node:`{record_id}` SET
574 tenant_id = $tenant_id,
575 session_id = $session_id,
576 raw = $raw,
577 tier = $tier,
578 timestamp = <datetime>$timestamp,
579 compression_depth = $compression_depth,{parent_assignment}
580 sync_key = $sync_key,
581 updated_at = <datetime>$updated_at,{source_metadata_assignment}
582 {context_summary_assignment}
583 psi = $psi,
584 rho = $rho,
585 kappa = $kappa,
586 user_stability = $user_stability,
587 user_friction = $user_friction,
588 user_logic = $user_logic,
589 user_autonomy = $user_autonomy,
590 user_psi = $user_psi,
591 model_stability = $model_stability,
592 model_friction = $model_friction,
593 model_logic = $model_logic,
594 model_autonomy = $model_autonomy,
595 model_psi = $model_psi,
596 comp_stability = $comp_stability,
597 comp_friction = $comp_friction,
598 comp_logic = $comp_logic,
599 comp_autonomy = $comp_autonomy,
600 comp_psi = $comp_psi,
601 semantic_tags = {semantic_tags_value},
602 semantic_links = {semantic_links_value};
603 "#
604 )
605}
606
607pub fn query_changes_since_query(limit: usize) -> String {
608 format!(
609 r#"
610 SELECT
611 tenant_id AS TenantId,
612 session_id AS SessionId,
613 raw AS Raw,
614 tier AS Tier,
615 timestamp AS Timestamp,
616 compression_depth AS CompressionDepth,
617 parent_node_id AS ParentNodeId,
618 sync_key AS SyncKey,
619 updated_at AS UpdatedAt,
620 source_metadata AS SourceMetadata,
621 context_summary AS ContextSummary,
622 semantic_tags AS SemanticTags,
623 semantic_links AS SemanticLinks,
624 embedding AS Embedding,
625 embedding_model AS EmbeddingModel,
626 embedding_dimensions AS EmbeddingDimensions,
627 embedded_at AS EmbeddedAt,
628 psi AS Psi,
629 rho AS Rho,
630 kappa AS Kappa,
631 user_stability AS UserStability,
632 user_friction AS UserFriction,
633 user_logic AS UserLogic,
634 user_autonomy AS UserAutonomy,
635 user_psi AS UserPsi,
636 model_stability AS ModelStability,
637 model_friction AS ModelFriction,
638 model_logic AS ModelLogic,
639 model_autonomy AS ModelAutonomy,
640 model_psi AS ModelPsi,
641 comp_stability AS CompStability,
642 comp_friction AS CompFriction,
643 comp_logic AS CompLogic,
644 comp_autonomy AS CompAutonomy,
645 comp_psi AS CompPsi,
646 0 AS ResonanceDelta
647 FROM temporal_node
648 WHERE session_id = $session_id
649 AND (tenant_id = $tenant_id OR tenant_id = NONE OR tenant_id = '')
650 AND (
651 NOT $include_cursor
652 OR updated_at > <datetime>$cursor_updated_at
653 OR (
654 updated_at = <datetime>$cursor_updated_at
655 AND sync_key > $cursor_sync_key
656 )
657 )
658 ORDER BY updated_at ASC, sync_key ASC
659 LIMIT {limit};
660 "#
661 )
662}
663
664pub const GET_SYNC_CHECKPOINT_QUERY: &str = r#"
665 SELECT
666 session_id AS SessionId,
667 connector_id AS ConnectorId,
668 cursor_updated_at AS CursorUpdatedAt,
669 cursor_sync_key AS CursorSyncKey,
670 updated_at AS UpdatedAt,
671 metadata AS Metadata
672 FROM sync_checkpoint
673 WHERE session_id = $session_id
674 AND connector_id = $connector_id
675 AND (tenant_id = $tenant_id OR tenant_id = NONE OR tenant_id = '')
676 LIMIT 1;
677 "#;
678
679pub fn upsert_sync_checkpoint_query(record_id: &str, include_metadata_assignment: bool) -> String {
680 let metadata_assignment = if include_metadata_assignment {
681 "metadata = $metadata,"
682 } else {
683 "metadata = NONE,"
684 };
685
686 format!(
687 r#"
688 UPSERT sync_checkpoint:`{record_id}` SET
689 tenant_id = $tenant_id,
690 session_id = $session_id,
691 connector_id = $connector_id,
692 cursor_updated_at = <datetime>$cursor_updated_at,
693 cursor_sync_key = $cursor_sync_key,
694 {metadata_assignment}
695 updated_at = <datetime>$updated_at;
696 "#
697 )
698}
699
700pub const SELECT_CALIBRATION_MISSING_TENANT_QUERY: &str = r#"
701 SELECT id, session_id
702 FROM calibration
703 WHERE tenant_id = NONE OR tenant_id = '';
704 "#;
705
706pub const SELECT_SCOPE_BY_NODE_ID_QUERY: &str = r#"
707 SELECT
708 tenant_id AS TenantId,
709 session_id AS SessionId
710 FROM temporal_node
711 WHERE id = type::record('temporal_node', $node_id)
712 LIMIT 1;
713 "#;
714
715pub const COUNT_TEMPORAL_SCOPE_QUERY: &str = r#"
716 SELECT count() AS Count
717 FROM temporal_node
718 WHERE session_id = $session_id
719 AND (tenant_id = $tenant_id OR ($include_legacy AND (tenant_id = NONE OR tenant_id = '')))
720 LIMIT 1;
721 "#;
722
723pub const COUNT_CALIBRATION_SCOPE_QUERY: &str = r#"
724 SELECT count() AS Count
725 FROM calibration
726 WHERE session_id = $session_id
727 AND (tenant_id = $tenant_id OR ($include_legacy AND (tenant_id = NONE OR tenant_id = '')))
728 LIMIT 1;
729 "#;
730
731pub const COUNT_CHECKPOINT_SCOPE_QUERY: &str = r#"
732 SELECT count() AS Count
733 FROM sync_checkpoint
734 WHERE session_id = $session_id
735 AND (tenant_id = $tenant_id OR ($include_legacy AND (tenant_id = NONE OR tenant_id = '')))
736 LIMIT 1;
737 "#;
738
739pub const APPLY_SCOPE_REKEY_QUERY: &str = r#"
740 BEGIN TRANSACTION;
741
742 UPDATE temporal_node
743 SET
744 tenant_id = $target_tenant_id,
745 session_id = $target_session_id
746 WHERE session_id = $source_session_id
747 AND (tenant_id = $source_tenant_id OR ($source_include_legacy AND (tenant_id = NONE OR tenant_id = '')));
748
749 UPDATE calibration
750 SET
751 tenant_id = $target_tenant_id,
752 session_id = $target_session_id
753 WHERE session_id = $source_session_id
754 AND (tenant_id = $source_tenant_id OR ($source_include_legacy AND (tenant_id = NONE OR tenant_id = '')));
755
756 COMMIT TRANSACTION;
757 "#;
758
759pub const DELETE_TAG_ROWS_FOR_SYNC_KEY_QUERY: &str = r#"
760 DELETE semantic_tag_index
761 WHERE tenant_id = $tenant_id AND sync_key = $sync_key;
762 "#;
763
764pub const DELETE_TAG_ROWS_FOR_SESSION_QUERY: &str = r#"
765 DELETE semantic_tag_index
766 WHERE tenant_id = $tenant_id AND session_id = $session_id;
767 "#;
768
769pub const DELETE_TEMPORAL_NODE_BY_SYNC_KEY_QUERY: &str = r#"
770 DELETE temporal_node
771 WHERE tenant_id = $tenant_id AND session_id = $session_id AND sync_key = $sync_key;
772 "#;
773
774pub const DELETE_TEMPORAL_NODE_BY_ID_QUERY: &str = r#"
775 DELETE type::thing('temporal_node', $node_id);
776 "#;
777
778pub fn purge_temporal_nodes_query(tier_clause: Option<&str>) -> String {
779 let tier_filter = tier_clause
780 .map(|clause| format!(" AND {clause}"))
781 .unwrap_or_default();
782 format!(
783 r#"
784 DELETE temporal_node
785 WHERE tenant_id = $tenant_id AND session_id = $session_id{tier_filter};
786 "#
787 )
788}
789
790pub const PURGE_CALIBRATION_SESSION_QUERY: &str = r#"
791 DELETE calibration
792 WHERE tenant_id = $tenant_id AND session_id = $session_id;
793 "#;
794
795pub const PURGE_CHECKPOINT_SESSION_QUERY: &str = r#"
796 DELETE sync_checkpoint
797 WHERE tenant_id = $tenant_id AND session_id = $session_id;
798 "#;
799
800pub const SELECT_NODE_BY_SYNC_KEY_QUERY: &str = r#"
801 SELECT
802 meta::id(id) AS NodeId,
803 sync_key AS SyncKey
804 FROM temporal_node
805 WHERE tenant_id = $tenant_id AND session_id = $session_id AND sync_key = $sync_key
806 LIMIT 1;
807 "#;
808
809pub const SELECT_NODE_BY_ID_QUERY: &str = r#"
810 SELECT
811 meta::id(id) AS NodeId,
812 sync_key AS SyncKey,
813 session_id AS SessionId
814 FROM type::thing('temporal_node', $node_id)
815 LIMIT 1;
816 "#;
817
818pub const UPSERT_TAG_ROW_QUERY: &str = r#"
819 UPSERT semantic_tag_index
820 SET
821 tenant_id = $tenant_id,
822 session_id = $session_id,
823 node_id = $node_id,
824 sync_key = $sync_key,
825 tag = $tag,
826 embedding = $embedding,
827 embedding_model = $embedding_model,
828 embedding_dimensions = $embedding_dimensions,
829 embedded_at = $embedded_at,
830 updated_at = $updated_at
831 WHERE tenant_id = $tenant_id AND sync_key = $sync_key AND tag = $tag;
832 "#;
833
834pub const UPSERT_TAG_ROW_META_QUERY: &str = r#"
835 UPSERT semantic_tag_index
836 SET
837 tenant_id = $tenant_id,
838 session_id = $session_id,
839 node_id = $node_id,
840 sync_key = $sync_key,
841 tag = $tag,
842 updated_at = $updated_at
843 WHERE tenant_id = $tenant_id AND sync_key = $sync_key AND tag = $tag;
844 "#;
845
846pub const LIST_TAGS_FOR_SYNC_KEY_QUERY: &str = r#"
847 SELECT tag AS Tag FROM semantic_tag_index
848 WHERE tenant_id = $tenant_id AND sync_key = $sync_key;
849 "#;
850
851pub fn query_tag_records_query(where_clause: &str, capped_limit: usize) -> String {
852 format!(
853 r#"
854 SELECT
855 tenant_id AS TenantId,
856 session_id AS SessionId,
857 node_id AS NodeId,
858 sync_key AS SyncKey,
859 tag AS Tag,
860 embedding AS Embedding,
861 embedding_model AS EmbeddingModel,
862 embedding_dimensions AS EmbeddingDimensions,
863 embedded_at AS EmbeddedAt,
864 updated_at AS UpdatedAt
865 FROM semantic_tag_index
866 WHERE {where_clause}
867 ORDER BY updated_at DESC
868 LIMIT {capped_limit};
869 "#
870 )
871}
872
873pub fn find_sync_keys_by_tags_query(match_all: bool) -> String {
874 if match_all {
875 r#"
876 SELECT sync_key AS SyncKey, count() AS TagCount
877 FROM semantic_tag_index
878 WHERE tenant_id = $tenant_id
879 AND tag IN $tags
880 AND ($session_id IS NONE OR session_id = $session_id)
881 GROUP BY sync_key
882 HAVING TagCount = array::len($tags)
883 LIMIT $limit;
884 "#
885 .to_string()
886 } else {
887 r#"
888 SELECT DISTINCT sync_key AS SyncKey
889 FROM semantic_tag_index
890 WHERE tenant_id = $tenant_id
891 AND tag IN $tags
892 AND ($session_id IS NONE OR session_id = $session_id)
893 LIMIT $limit;
894 "#
895 .to_string()
896 }
897}
898
899pub fn find_tags_vocabulary_query(where_clause: &str, capped_limit: usize) -> String {
900 format!(
901 r#"
902 SELECT DISTINCT tag AS Tag
903 FROM semantic_tag_index
904 WHERE {where_clause}
905 ORDER BY tag ASC
906 LIMIT {capped_limit};
907 "#
908 )
909}
910
911pub fn update_record_tenant_query(record_id: &str) -> String {
912 format!(
913 r#"
914 UPDATE {record_id}
915 SET tenant_id = $tenant_id;
916 "#
917 )
918}
919
920#[cfg(test)]
921mod tests {
922 use super::create_temporal_node_query;
923
924 #[test]
925 fn create_temporal_node_query_uses_none_for_missing_embedded_at() {
926 let query = create_temporal_node_query(
927 "abc123", false, false, true, true, false, false, false, false, false, false,
928 );
929
930 assert!(query.contains("embedded_at = NONE"));
931 assert!(!query.contains("embedded_at = <datetime>$embedded_at"));
932 assert!(query.contains("embedding = NONE"));
933 assert!(query.contains("embedding_model = NONE"));
934 assert!(query.contains("embedding_dimensions = NONE"));
935 assert!(query.contains("semantic_tags = NONE"));
936 assert!(query.contains("semantic_links = NONE"));
937 }
938
939 #[test]
940 fn create_temporal_node_query_uses_datetime_cast_when_embedded_at_present() {
941 let query = create_temporal_node_query(
942 "abc123", false, false, true, true, true, true, true, true, false, false,
943 );
944
945 assert!(query.contains("embedded_at = <datetime>$embedded_at"));
946 }
947}