Skip to main content

locus_core_rs/storage/surrealdb/
raw_queries.rs

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}