Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
216 changes: 111 additions & 105 deletions crates/tracedecay-global-db/src/observation_adapter.rs

Large diffs are not rendered by default.

273 changes: 166 additions & 107 deletions crates/tracedecay-global-db/src/observation_collision_tests.rs

Large diffs are not rendered by default.

85 changes: 3 additions & 82 deletions crates/tracedecay-global-db/src/observation_projection/rebuild.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,8 @@ use super::apply::{
use super::state::{
consume_projection_queue_item, decode_observation_row, decode_sequence,
ensure_projection_output_state_cache, projection_retry_state, queued_sequence, read_checkpoint,
read_message, read_observation, read_session, schedule_projection_retry, storage,
storage_message, write_checkpoint,
read_message, read_observation, read_session, reaggregate_output_state_for_output,
schedule_projection_retry, storage, storage_message, write_checkpoint,
};
use super::transition::{
MessageTransition, MessageTransitionState, WorkflowFactTarget, WorkflowFactTransition,
Expand Down Expand Up @@ -1128,86 +1128,7 @@ async fn reconcile_collided_observation_provenance(
.await
.map_err(|error| storage("remove collided projection provenance", error))?;
for (output_provider, output_message_id) in affected {
conn.execute(
"DELETE FROM temp.observation_projection_output_state
WHERE projector_version = ?1
AND output_provider = ?2 AND output_message_id = ?3",
params![
SESSION_MESSAGE_PROJECTOR_VERSION,
output_provider.as_str(),
output_message_id.as_str(),
],
)
.await
.map_err(|error| storage("reset collided projection output state", error))?;
conn.execute(
"INSERT INTO temp.observation_projection_output_state (
projector_version, output_provider, output_message_id,
canonical_observation_id, latest_observation_id, latest_sequence,
projector_owned, owner_count
)
SELECT groups.projector_version, groups.output_provider, groups.output_message_id,
CASE WHEN groups.projector_owned = 1 THEN (
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
) ELSE (
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence ASC, provenance.observation_id ASC
LIMIT 1
) END,
(
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
),
(
SELECT observation.sequence
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
),
groups.projector_owned, groups.owner_count
FROM (
SELECT projector_version, output_provider, output_message_id,
MAX(message_created) AS projector_owned,
COUNT(*) AS owner_count
FROM observation_projection_provenance
WHERE projector_version = ?1
AND output_provider = ?2 AND output_message_id = ?3
GROUP BY projector_version, output_provider, output_message_id
) AS groups",
params![
SESSION_MESSAGE_PROJECTOR_VERSION,
output_provider.as_str(),
output_message_id.as_str(),
],
)
.await
.map_err(|error| storage("reaggregate collided projection output state", error))?;
reaggregate_output_state_for_output(conn, &output_provider, &output_message_id).await?;
}
Ok(())
}
Expand Down
174 changes: 115 additions & 59 deletions crates/tracedecay-global-db/src/observation_projection/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -360,6 +360,117 @@ pub(super) async fn read_message(
}))
}

/// The one definition of the projected-output ownership aggregation that
/// populates `temp.observation_projection_output_state`: for every
/// `(projector_version, output_provider, output_message_id)` group in the
/// (optionally filtered) provenance authority it derives the canonical owner
/// (newest row when the projector owns the output, oldest otherwise), the
/// newest owner row and its sequence, and the group's ownership counts.
/// Whole-cache initialization and per-output re-aggregation both render
/// their statement from this single spelling so the aggregation cannot
/// drift; `provenance_filter` scopes only the grouped rows (the correlated
/// owner lookups constrain themselves to each group's exact key).
fn output_state_aggregation_sql(provenance_filter: &str) -> String {
format!(
"INSERT INTO temp.observation_projection_output_state (
projector_version, output_provider, output_message_id,
canonical_observation_id, latest_observation_id, latest_sequence,
projector_owned, owner_count
)
SELECT groups.projector_version, groups.output_provider, groups.output_message_id,
CASE WHEN groups.projector_owned = 1 THEN (
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
) ELSE (
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence ASC, provenance.observation_id ASC
LIMIT 1
) END,
(
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
),
(
SELECT observation.sequence
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
),
groups.projector_owned, groups.owner_count
FROM (
SELECT projector_version, output_provider, output_message_id,
MAX(message_created) AS projector_owned,
COUNT(*) AS owner_count
FROM observation_projection_provenance
{provenance_filter}
GROUP BY projector_version, output_provider, output_message_id
) AS groups"
)
}

/// Re-aggregates the ownership cache for one exact output from the
/// provenance authority: the output's cached row is removed and rebuilt
/// through the canonical aggregation ([`output_state_aggregation_sql`]), so
/// convergence paths (e.g. collided-provenance reconciliation) share the
/// initialization's single definition.
pub(super) async fn reaggregate_output_state_for_output(
conn: &impl Executor,
output_provider: &str,
output_message_id: &str,
) -> ProjectionStoreResult<()> {
conn.execute(
"DELETE FROM temp.observation_projection_output_state
WHERE projector_version = ?1
AND output_provider = ?2 AND output_message_id = ?3",
params![
SESSION_MESSAGE_PROJECTOR_VERSION,
output_provider,
output_message_id,
],
)
.await
.map_err(|error| storage("reset collided projection output state", error))?;
conn.execute(
&output_state_aggregation_sql(
"WHERE projector_version = ?1
AND output_provider = ?2 AND output_message_id = ?3",
),
params![
SESSION_MESSAGE_PROJECTOR_VERSION,
output_provider,
output_message_id,
],
)
.await
.map_err(|error| storage("reaggregate collided projection output state", error))?;
Ok(())
}

pub(super) struct ProjectionOutputOwner {
pub(super) sequence: u64,
pub(super) observation: DurableObservationV1,
Expand Down Expand Up @@ -427,68 +538,13 @@ pub(super) async fn ensure_projection_output_state_cache(

conn.execute_batch(
"DELETE FROM temp.observation_projection_output_state;
DELETE FROM temp.observation_projection_output_state_meta;
WITH owner_groups AS (
SELECT projector_version, output_provider, output_message_id,
MAX(message_created) AS projector_owned,
COUNT(*) AS owner_count
FROM observation_projection_provenance
GROUP BY projector_version, output_provider, output_message_id
)
INSERT INTO temp.observation_projection_output_state (
projector_version, output_provider, output_message_id,
canonical_observation_id, latest_observation_id, latest_sequence,
projector_owned, owner_count
)
SELECT groups.projector_version, groups.output_provider, groups.output_message_id,
CASE WHEN groups.projector_owned = 1 THEN (
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
) ELSE (
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence ASC, provenance.observation_id ASC
LIMIT 1
) END,
(
SELECT provenance.observation_id
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
),
(
SELECT observation.sequence
FROM observation_projection_provenance AS provenance
JOIN observations AS observation
ON observation.observation_id = provenance.observation_id
WHERE provenance.projector_version = groups.projector_version
AND provenance.output_provider = groups.output_provider
AND provenance.output_message_id = groups.output_message_id
ORDER BY observation.sequence DESC, provenance.observation_id DESC
LIMIT 1
),
groups.projector_owned, groups.owner_count
FROM owner_groups AS groups;",
DELETE FROM temp.observation_projection_output_state_meta;",
)
.await
.map_err(|error| storage("initialize projection output state cache", error))?;
conn.execute(&output_state_aggregation_sql(""), ())
.await
.map_err(|error| storage("initialize projection output state cache", error))?;
conn.execute(
"INSERT INTO temp.observation_projection_output_state_meta(initialized, data_version)
VALUES (1, ?1)",
Expand Down
16 changes: 16 additions & 0 deletions crates/tracedecay-global-db/src/registered.rs
Original file line number Diff line number Diff line change
Expand Up @@ -581,6 +581,22 @@ impl RegisteredGlobalDb {
crate::GlobalDbObservationStore::new(self.database.clone())
}

/// Test-only [`Self::observation_store`] variant that binds the adapter
/// to an explicit runtime dispatch seam (e.g. a counting wrapper over
/// this client's runtime client), so tests can observe the record work a
/// persist call dispatches without changing the adapter's production
/// shape.
#[cfg(test)]
pub(crate) fn observation_store_with_runtime_dispatch<R>(
&self,
runtime: R,
) -> crate::GlobalDbObservationStore<R>
where
R: crate::observation_adapter::ObservationRuntimeDispatch,
{
crate::GlobalDbObservationStore::with_runtime_dispatch(self.database.clone(), runtime)
}

/// Retains this exact client for closed runtime read/submit requests.
///
/// The returned capability has no raw Store runtime, connection, or
Expand Down
Loading
Loading