use crate::http_envelope::unwrap_sync_envelope;
use crate::module::ImModule;
use crate::older_context::LoadOlderState;
use crate::state::CorrelationContext;
use crate::{error::ImError, query::LocalStoreMode};
use helix_core::effect::{Effect, ScanOrder, ScanSpec, SqlValue, StorageOp};
use helix_core::tick::PortOutcome;
use helix_core::EffectSink;
use serde_json::Value;
use std::collections::HashSet;
const OLDER_NAVIGATION_ORDER: &[ScanOrder] = &[
ScanOrder::desc("create_at"),
ScanOrder::desc("temporary_id"),
];
const NEWER_NAVIGATION_ORDER: &[ScanOrder] =
&[ScanOrder::asc("create_at"), ScanOrder::asc("temporary_id")];
impl ImModule {
pub(crate) fn start_timeline_navigation_or_local(
&mut self,
state: crate::timeline_navigation::TimelineNavigationState,
out: &mut EffectSink,
) {
let Some(request_id) = state.request_id() else {
self.start_timeline_navigation_http(state, out);
return;
};
if !self
.state
.timeline_navigation_pending
.insert(request_id.to_string())
{
out.push(crate::read_relay::emit_read_error(
request_id,
"duplicate timeline reqId",
));
return;
}
let key = state.coverage_key();
let Some(coverage) = self.state.timeline_navigation_coverage.get(&key).cloned() else {
self.start_timeline_navigation_http(state, out);
return;
};
let Some(coverage_rows) =
local_timeline_navigation_rows(&state, &coverage, self.config.auth_user_id.as_str())
else {
tracing::debug!(
channel_id = state.channel_id().as_str(),
operation = state.operation(),
"timeline navigation coverage is incomplete; falling back to authority"
);
self.start_timeline_navigation_http(state, out);
return;
};
let shaped = crate::render_ready::shape_message_rows_for_viewer(
&Value::Array(coverage_rows.clone()),
self.config.auth_user_id.as_str(),
);
let mut body = serde_json::json!({
"reqId": request_id,
"channelId": state.channel_id().as_str(),
"operation": state.operation(),
"messages": shaped,
"hasOlder": coverage.has_older,
"hasNewer": coverage.has_newer,
"pageSize": state.page_size(),
"olderCursor": coverage.older_cursor.clone().unwrap_or(Value::Null),
"newerCursor": coverage.newer_cursor.clone().unwrap_or(Value::Null),
});
if matches!(
state.kind(),
crate::timeline_navigation::TimelineNavigationKind::Locate { .. }
) {
body["targetPostId"] = serde_json::json!(state.anchor_message_id());
let Some(target_index) = coverage.target_index else {
self.start_timeline_navigation_http(state, out);
return;
};
body["targetIndex"] = serde_json::json!(target_index);
}
self.clear_timeline_navigation_pending(&state);
out.push(crate::read_relay::emit_read_body(request_id, body));
}
fn clear_timeline_navigation_pending(
&mut self,
state: &crate::timeline_navigation::TimelineNavigationState,
) {
if let Some(request_id) = state.request_id() {
self.state.timeline_navigation_pending.remove(request_id);
}
}
fn remember_timeline_navigation_coverage(
&mut self,
state: &crate::timeline_navigation::TimelineNavigationState,
rows: &[Value],
target_index: Option<usize>,
) {
let page = state.page();
self.state.timeline_navigation_coverage.insert(
state.coverage_key(),
crate::timeline_navigation::TimelineNavigationCoverage {
channel_id: state.channel_id(),
rows: rows.to_vec(),
has_older: page.has_older,
has_newer: page.has_newer,
target_index,
older_cursor: state.older_cursor().cloned(),
newer_cursor: state.newer_cursor().cloned(),
},
);
}
pub(crate) fn start_timeline_navigation_http(
&mut self,
state: crate::timeline_navigation::TimelineNavigationState,
out: &mut EffectSink,
) {
let corr = self.alloc_corr_internal();
out.push(crate::timeline_navigation::navigation_http(
&state,
self.config.api_base_url.as_str(),
self.state.connection_id.as_deref(),
corr,
));
self.state.corr_map.insert(
corr,
CorrelationContext::TimelineNavigationHttp {
state: Box::new(state),
},
);
}
pub(super) fn handle_timeline_navigation_http_reply(
&mut self,
mut state: Box<crate::timeline_navigation::TimelineNavigationState>,
outcome: &PortOutcome,
now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
let reply = match outcome {
PortOutcome::Ok(reply) => reply,
PortOutcome::Err(error) => {
tracing::warn!(error = ?error, "timeline navigation HTTP failed");
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
};
let raw = match unwrap_sync_envelope(reply.0.as_ref()) {
Ok(raw) => raw,
Err(error) => {
tracing::warn!(error = ?error, "timeline navigation envelope invalid");
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
};
if let Err(error) = state.ingest_http_body(&raw) {
tracing::warn!(error = ?error, "timeline navigation body invalid");
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
if matches!(
state.kind(),
crate::timeline_navigation::TimelineNavigationKind::Locate { .. }
) {
let rows = match crate::query::render_ready::locate::filter_visible_remote_rows(
state.rows(),
state.channel_id().as_str(),
state.anchor_message_id(),
self.config.auth_user_id.as_str(),
) {
Ok(rows) => rows,
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline locate target is not visible in HTTP authority"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
};
let target_index = rows.iter().position(|row| {
crate::timeline_navigation::row_matches_message_identity(
row,
state.anchor_message_id(),
)
});
state.set_target_index(target_index);
state.replace_rows(rows);
}
let (rows, ops) = match crate::query::local_first::visible_remote_rows_and_cache_ops(
state.channel_id(),
state.rows().to_vec(),
&[],
self.config.auth_user_id.as_str(),
) {
Ok(result) => result,
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline navigation rows could not become durable messages"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
};
state.replace_rows(rows);
if self.local_store_mode == LocalStoreMode::Disabled {
tracing::warn!(
channel_id = state.channel_id().as_str(),
"timeline navigation requires durable storage before event publication"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
if ops.is_empty() {
self.schedule_timeline_navigation_readback(state, out)?;
return Ok(());
}
let corr = self.alloc_corr_internal();
out.push(Effect::Persist { corr, ops });
self.state
.corr_map
.insert(corr, CorrelationContext::TimelineNavigationCache { state });
Ok(())
}
pub(super) fn emit_timeline_navigation_failed_event(
&mut self,
state: &crate::timeline_navigation::TimelineNavigationState,
_now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
self.clear_timeline_navigation_pending(state);
let page_direction = match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Older {
anchor_post_id, ..
} => Some(("older", anchor_post_id.as_str())),
crate::timeline_navigation::TimelineNavigationKind::Newer {
anchor_post_id, ..
} => Some(("newer", anchor_post_id.as_str())),
crate::timeline_navigation::TimelineNavigationKind::Locate { .. } => None,
};
if let Some((direction, anchor_post_id)) = page_direction {
if let Some(request_id) = state.request_id() {
let message = if direction == "older" {
"timeline older query failed before durable result"
} else {
"timeline newer query failed before durable result"
};
out.push(crate::read_relay::emit_read_error(request_id, message));
return Ok(());
}
let page = state.page();
out.push(
crate::event::timeline::page(
state.channel_id().as_str(),
state.window_token(),
direction,
"failed",
Vec::new(),
page.has_older,
page.has_newer,
Some(anchor_post_id),
)?
.into_effect(),
);
return Ok(());
}
let crate::timeline_navigation::TimelineNavigationKind::Locate {
navigation_token, ..
} = state.kind()
else {
return Ok(());
};
let _ = navigation_token;
if let Some(request_id) = state.request_id() {
out.push(crate::read_relay::emit_read_error(
request_id,
"timeline locate query failed before durable result",
));
}
Ok(())
}
pub(super) fn schedule_timeline_navigation_readback(
&mut self,
state: Box<crate::timeline_navigation::TimelineNavigationState>,
out: &mut EffectSink,
) -> Result<(), ImError> {
let exact_keys = timeline_navigation_readback_keys(&state);
if !exact_keys.is_empty() {
self.schedule_timeline_navigation_exact_readback(state, exact_keys, 1, Vec::new(), out);
return Ok(());
}
if matches!(
state.kind(),
crate::timeline_navigation::TimelineNavigationKind::Locate { .. }
) {
tracing::warn!(
channel_id = state.channel_id().as_str(),
"timeline locate authority page has no durable identity"
);
self.emit_timeline_navigation_failed_event(&state, 0, out)?;
return Ok(());
}
let corr = self.alloc_corr_internal();
let effect = match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Older { .. } => {
Effect::Persist {
corr,
ops: vec![StorageOp::Scan(older_navigation_scan_spec(&state))],
}
}
crate::timeline_navigation::TimelineNavigationKind::Newer { .. } => {
Effect::Persist {
corr,
ops: vec![StorageOp::Scan(newer_navigation_scan_spec(&state))],
}
}
crate::timeline_navigation::TimelineNavigationKind::Locate { .. } => {
unreachable!("locate readback must use exact identity continuation")
}
};
out.push(effect);
self.state.corr_map.insert(
corr,
CorrelationContext::TimelineNavigationReadback { state },
);
Ok(())
}
fn schedule_timeline_navigation_exact_readback(
&mut self,
state: Box<crate::timeline_navigation::TimelineNavigationState>,
accepted_keys: Vec<crate::timeline_navigation::TimelineNavigationReadbackKey>,
next_index: usize,
rows: Vec<Value>,
out: &mut EffectSink,
) {
let key = accepted_keys[next_index - 1].clone();
let corr = self.alloc_corr_internal();
out.push(Effect::Persist {
corr,
ops: vec![StorageOp::Scan(timeline_navigation_exact_scan_spec(&key))],
});
self.state.corr_map.insert(
corr,
CorrelationContext::TimelineNavigationExactReadback {
state,
accepted_keys,
next_index,
rows,
},
);
}
pub(super) fn handle_timeline_navigation_exact_readback_reply(
&mut self,
mut state: Box<crate::timeline_navigation::TimelineNavigationState>,
accepted_keys: Vec<crate::timeline_navigation::TimelineNavigationReadbackKey>,
next_index: usize,
mut rows: Vec<Value>,
outcome: &PortOutcome,
now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
let expected_key = accepted_keys
.get(next_index.saturating_sub(1))
.ok_or_else(|| {
ImError::Parse("timeline exact readback index out of bounds".to_string())
})?;
let durable_rows = match outcome {
PortOutcome::Ok(reply) => {
match crate::query::local_first::parse_local_rows(reply.0.as_ref()) {
Ok(rows) => rows,
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline exact durable read-back malformed"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline exact durable read-back failed"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
};
let mut seen = rows
.iter()
.filter_map(|row| {
accepted_keys
.iter()
.find(|key| durable_row_matches_readback_key(row, key))
.map(|key| (key.column, key.value.clone()))
})
.collect::<HashSet<_>>();
let mut expected_found = false;
for row in durable_rows {
let Some(key) = accepted_keys
.iter()
.find(|key| durable_row_matches_readback_key(&row, key))
else {
continue;
};
expected_found |= key == expected_key;
if seen.insert((key.column, key.value.clone())) {
rows.push(row);
}
}
if !expected_found {
tracing::warn!(
channel_id = state.channel_id().as_str(),
expected_key = expected_key.value.as_str(),
"timeline exact durable read-back missing authority identity"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
let mut next = next_index;
while next < accepted_keys.len() {
if seen.contains(&(
accepted_keys[next].column,
accepted_keys[next].value.clone(),
)) {
next += 1;
} else {
break;
}
}
if next < accepted_keys.len() {
self.schedule_timeline_navigation_exact_readback(
state,
accepted_keys,
next + 1,
rows,
out,
);
return Ok(());
}
let projected_rows = match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Older { .. } => {
older_navigation_rows(&state, rows, self.config.auth_user_id.as_str())
}
crate::timeline_navigation::TimelineNavigationKind::Newer { .. } => {
newer_navigation_rows(&state, rows, self.config.auth_user_id.as_str())
}
crate::timeline_navigation::TimelineNavigationKind::Locate {
target_message_id,
..
} => match crate::query::render_ready::locate::normalize_durable_located_window(
rows,
state.channel_id().as_str(),
target_message_id,
self.config.auth_user_id.as_str(),
state.page_size(),
) {
Ok(rows) => {
let target_index = rows.iter().position(|row| {
crate::timeline_navigation::row_matches_message_identity(
row,
target_message_id,
)
});
if target_index.is_none() {
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
state.set_target_index(target_index);
rows
}
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline locate exact durable projection failed closed"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
},
};
if !matches!(
state.kind(),
crate::timeline_navigation::TimelineNavigationKind::Locate { .. }
) && projected_rows.len() != accepted_keys.len()
{
tracing::warn!(
channel_id = state.channel_id().as_str(),
expected = accepted_keys.len(),
actual = projected_rows.len(),
"timeline exact durable read-back projected an incomplete page"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
let target_index = match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Locate { .. } => {
state.target_index()
}
_ => None,
};
self.remember_timeline_navigation_coverage(&state, &projected_rows, target_index);
match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Newer { .. } => {
match self.emit_timeline_navigation_newer_result(&state, &projected_rows) {
Ok(effect) => out.push(effect),
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline newer exact durable rows could not project"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
}
}
}
crate::timeline_navigation::TimelineNavigationKind::Older { .. } => {
out.push(self.emit_timeline_navigation_result(&state, &projected_rows)?);
}
crate::timeline_navigation::TimelineNavigationKind::Locate { .. } => {
out.push(self.emit_timeline_locate_result(&state, &projected_rows)?);
}
}
Ok(())
}
pub(super) fn handle_timeline_navigation_readback_reply(
&mut self,
state: Box<crate::timeline_navigation::TimelineNavigationState>,
outcome: &PortOutcome,
now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
if matches!(
state.kind(),
crate::timeline_navigation::TimelineNavigationKind::Newer { .. }
) {
let rows = match outcome {
PortOutcome::Ok(reply) => {
match crate::query::local_first::parse_local_rows(reply.0.as_ref()) {
Ok(rows) => {
newer_navigation_rows(&state, rows, self.config.auth_user_id.as_str())
}
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline newer durable read-back malformed"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline newer durable read-back failed"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
return Ok(());
}
};
if newer_http_has_rows(&state) && rows.is_empty() {
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
} else {
self.remember_timeline_navigation_coverage(&state, &rows, None);
match self.emit_timeline_navigation_newer_result(&state, &rows) {
Ok(effect) => out.push(effect),
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline newer durable rows could not project"
);
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
}
}
}
return Ok(());
}
if matches!(
state.kind(),
crate::timeline_navigation::TimelineNavigationKind::Older { .. }
) {
let rows = match outcome {
PortOutcome::Ok(reply) => {
match crate::query::local_first::parse_local_rows(reply.0.as_ref()) {
Ok(rows) => {
older_navigation_rows(&state, rows, self.config.auth_user_id.as_str())
}
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline older durable read-back malformed"
);
Vec::new()
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"timeline older durable read-back failed"
);
Vec::new()
}
};
if !state.authority_rows().is_empty() && rows.is_empty() {
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
} else {
self.remember_timeline_navigation_coverage(&state, &rows, None);
out.push(self.emit_timeline_navigation_result(&state, &rows)?);
}
return Ok(());
}
self.emit_timeline_navigation_failed_event(&state, now_ms, out)?;
Ok(())
}
fn emit_timeline_navigation_result(
&mut self,
state: &crate::timeline_navigation::TimelineNavigationState,
rows: &[Value],
) -> Result<Effect, ImError> {
self.clear_timeline_navigation_pending(state);
let request_id = state
.request_id()
.ok_or_else(|| ImError::Parse("timeline older result missing req_id".to_string()))?;
let shaped = crate::render_ready::shape_message_rows_for_viewer(
&Value::Array(rows.to_vec()),
self.config.auth_user_id.as_str(),
);
let page = state.page();
Ok(crate::read_relay::emit_read_body(
request_id,
serde_json::json!({
"reqId": request_id,
"channelId": state.channel_id().as_str(),
"operation": state.operation(),
"messages": shaped,
"hasOlder": page.has_older,
"hasNewer": page.has_newer,
"pageSize": state.page_size(),
"olderCursor": Value::Null,
"newerCursor": Value::Null,
}),
))
}
fn emit_timeline_navigation_newer_result(
&mut self,
state: &crate::timeline_navigation::TimelineNavigationState,
rows: &[Value],
) -> Result<Effect, ImError> {
self.clear_timeline_navigation_pending(state);
let request_id = state
.request_id()
.ok_or_else(|| ImError::Parse("timeline newer result missing req_id".to_string()))?;
let shaped = crate::render_ready::shape_message_rows_for_viewer(
&Value::Array(rows.to_vec()),
self.config.auth_user_id.as_str(),
);
let shaped_rows = shaped.as_array().ok_or_else(|| {
ImError::Parse("timeline newer durable rows must shape to array".to_string())
})?;
let _ = shaped_rows;
let page = state.page();
Ok(crate::read_relay::emit_read_body(
request_id,
serde_json::json!({
"reqId": request_id,
"channelId": state.channel_id().as_str(),
"operation": state.operation(),
"messages": shaped,
"hasOlder": page.has_older,
"hasNewer": page.has_newer,
"pageSize": state.page_size(),
"olderCursor": Value::Null,
"newerCursor": Value::Null,
}),
))
}
fn emit_timeline_locate_result(
&mut self,
state: &crate::timeline_navigation::TimelineNavigationState,
rows: &[Value],
) -> Result<Effect, ImError> {
self.clear_timeline_navigation_pending(state);
let request_id = state
.request_id()
.ok_or_else(|| ImError::Parse("timeline locate result missing req_id".to_string()))?;
let crate::timeline_navigation::TimelineNavigationKind::Locate {
target_message_id: _,
navigation_token: _,
} = state.kind()
else {
return Err(ImError::Parse(
"timeline locate result requires locate state".to_string(),
));
};
let page = state.page();
Ok(crate::read_relay::emit_read_body(
request_id,
serde_json::json!({
"reqId": request_id,
"channelId": state.channel_id().as_str(),
"operation": state.operation(),
"messages": rows,
"hasOlder": page.has_older,
"hasNewer": page.has_newer,
"pageSize": state.page_size(),
"targetPostId": state.anchor_message_id(),
"targetIndex": state.target_index(),
"olderCursor": state.older_cursor().cloned().unwrap_or(Value::Null),
"newerCursor": state.newer_cursor().cloned().unwrap_or(Value::Null),
}),
))
}
pub(super) fn handle_outbound_pinned_reply(
&mut self,
req_id: String,
account_id: String,
channel_id: crate::state::ChannelId,
projection_key: String,
epoch: u64,
outcome: &PortOutcome,
out: &mut EffectSink,
) -> Result<(), ImError> {
let PortOutcome::Ok(reply) = outcome else {
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"http request failed",
));
return Ok(());
};
let raw_body = match unwrap_sync_envelope(reply.0.as_ref()) {
Ok(body) => body,
Err(error) => {
tracing::warn!(req_id, error = ?error, "pinned read envelope decode failed");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"response envelope decode failed",
));
return Ok(());
}
};
let current_epoch = self
.state
.pinned_projection_epochs
.get(&channel_id)
.copied()
.unwrap_or(0);
if current_epoch != epoch {
out.push(crate::read_relay::emit_read_result(
req_id.as_str(),
raw_body.as_ref(),
));
return Ok(());
}
let corr = self.alloc_corr_internal();
let persist = crate::query::pinned_projection::persist_effect(
projection_key,
account_id,
channel_id,
raw_body.as_ref(),
corr,
)?;
self.state.corr_map.insert(
corr,
crate::state::CorrelationContext::PinnedProjectionPersist {
req_id,
raw_body: raw_body.to_vec(),
},
);
out.push(persist);
Ok(())
}
pub(super) fn handle_outbound_read_reply(
&mut self,
req_id: &String,
command: &str,
channel_id: Option<&str>,
outcome: &PortOutcome,
out: &mut EffectSink,
) -> Result<(), ImError> {
match outcome {
PortOutcome::Ok(reply) => match unwrap_sync_envelope(reply.0.as_ref()) {
Ok(raw_body) => {
if command == "im_get_schedule" {
let body = match serde_json::from_slice::<serde_json::Value>(&raw_body) {
Ok(body)
if body.get("status").is_none_or(|status| status == "SUCCESS") =>
{
body
}
_ => {
out.push(crate::read_relay::emit_read_error(
req_id,
"SCHEDULE_READ_FAILED",
));
return Ok(());
}
};
let schedules = crate::outbound::posts::read::schedule_projections(
channel_id.unwrap_or_default(),
&body,
);
out.push(crate::read_relay::emit_read_body(
req_id.as_str(),
serde_json::json!({
"ok": true,
"schedules": schedules,
}),
));
} else if command == "im_announcement_list" {
self.emit_announcement_list_reply(
req_id,
channel_id.unwrap_or_default(),
&raw_body,
out,
)?;
} else if command == "im_announcement_delete" {
self.emit_announcement_delete_reply(
req_id,
channel_id.unwrap_or_default(),
&raw_body,
out,
)?;
} else {
out.push(crate::read_relay::emit_read_result(
req_id.as_str(),
raw_body.as_ref(),
));
}
}
Err(e) => {
tracing::warn!(
req_id,
error = ?e,
"read reply envelope decode failed"
);
out.push(crate::read_relay::emit_read_error(
req_id,
"response envelope decode failed",
));
}
},
PortOutcome::Err(e) => {
tracing::warn!(req_id = req_id.as_str(), error = ?e, "read outbound http failed");
out.push(crate::read_relay::emit_read_error(
req_id,
"http request failed",
));
}
}
Ok(())
}
fn emit_announcement_list_reply(
&mut self,
req_id: &str,
channel_id: &str,
raw_body: &[u8],
out: &mut EffectSink,
) -> Result<(), ImError> {
let body = match serde_json::from_slice::<serde_json::Value>(raw_body) {
Ok(body) => body,
Err(error) => {
tracing::warn!(req_id, error = ?error, "announcement list response JSON invalid");
out.push(crate::read_relay::emit_read_error(
req_id,
"announcement list response JSON invalid",
));
return Ok(());
}
};
let Some(snapshot) =
crate::outbound::posts::read_ext::announcement_list_projection(channel_id, &body)
else {
out.push(crate::read_relay::emit_read_error(
req_id,
"announcement list response missing canonical versioned snapshot",
));
return Ok(());
};
let version = snapshot
.get("version")
.and_then(serde_json::Value::as_u64)
.unwrap_or_default();
let Some(channel) = crate::state::ChannelId::from_str(channel_id) else {
out.push(crate::read_relay::emit_read_error(
req_id,
"announcement list response has invalid channel",
));
return Ok(());
};
let committed = self
.state
.announcement_versions
.get(&channel)
.copied()
.unwrap_or_default();
let pending = self
.state
.announcement_reload_versions
.get(&channel)
.copied();
if version < committed || pending.is_some_and(|expected| version < expected) {
out.push(crate::read_relay::emit_read_error(
req_id,
"stale announcement version",
));
return Ok(());
}
let should_emit = version > committed || pending == Some(version);
self.state
.announcement_versions
.insert(channel, committed.max(version));
if pending.is_some_and(|expected| expected <= version) {
self.state.announcement_reload_versions.remove(&channel);
}
if should_emit {
out.push(crate::event::announcement::list_updated(snapshot.clone())?.into_effect());
}
out.push(crate::read_relay::emit_read_body(req_id, snapshot));
Ok(())
}
fn emit_announcement_delete_reply(
&mut self,
req_id: &str,
channel_id: &str,
raw_body: &[u8],
out: &mut EffectSink,
) -> Result<(), ImError> {
let body = match serde_json::from_slice::<serde_json::Value>(raw_body) {
Ok(body) => body,
Err(error) => {
tracing::warn!(req_id, error = ?error, "announcement delete response JSON invalid");
out.push(crate::read_relay::emit_read_error(
req_id,
"announcement delete response JSON invalid",
));
return Ok(());
}
};
let Some(result) =
crate::outbound::posts::read_ext::announcement_delete_projection(channel_id, &body)
else {
out.push(crate::read_relay::emit_read_error(
req_id,
"announcement delete response missing canonical version",
));
return Ok(());
};
let Some(channel) = crate::state::ChannelId::from_str(channel_id) else {
out.push(crate::read_relay::emit_read_error(
req_id,
"announcement delete response has invalid channel",
));
return Ok(());
};
let version = result
.get("version")
.and_then(serde_json::Value::as_u64)
.unwrap_or_default();
let committed = self
.state
.announcement_versions
.get(&channel)
.copied()
.unwrap_or_default();
let pending = self
.state
.announcement_reload_versions
.get(&channel)
.copied();
let no_op = result
.get("noOp")
.and_then(serde_json::Value::as_bool)
.unwrap_or(version == 0);
if no_op {
if version > committed {
self.state.announcement_versions.insert(channel, version);
}
out.push(crate::read_relay::emit_read_body(req_id, result));
return Ok(());
}
if version < committed || pending.is_some_and(|expected| expected > version) {
out.push(crate::read_relay::emit_read_error(
req_id,
"stale announcement delete version",
));
return Ok(());
}
self.state
.announcement_versions
.insert(channel, committed.max(version));
out.push(crate::read_relay::emit_read_body(req_id, result));
self.enqueue_announcement_reload(channel, version, out)
}
fn enqueue_announcement_reload(
&mut self,
channel: crate::state::ChannelId,
expected_version: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
let pending = self
.state
.announcement_reload_versions
.get(&channel)
.copied()
.unwrap_or_default();
if pending >= expected_version {
return Ok(());
}
let corr = self.alloc_corr_internal();
let req_id = format!("announcement-refresh-{}", corr.raw());
let payload = serde_json::to_vec(&serde_json::json!({
"channel_id": channel.as_str(),
"req_id": req_id,
}))
.map_err(|error| ImError::Parse(format!("announcement refresh payload: {error}")))?;
let effects = crate::commands::handle_outbound(
"im_announcement_list",
&payload,
self.config.api_base_url.as_str(),
self.config.default_api_base_url.as_str(),
self.state.connection_id.as_deref(),
corr,
)?;
self.state
.announcement_reload_versions
.insert(channel, expected_version);
self.state.corr_map.insert(
corr,
CorrelationContext::OutboundReadReply {
req_id,
command: "im_announcement_list".to_string(),
channel_id: Some(channel.as_str().to_string()),
},
);
for effect in effects {
out.push(effect);
}
Ok(())
}
pub(super) fn handle_exact_posts_http_reply(
&mut self,
req_id: &str,
requested_ids: Vec<String>,
outcome: &PortOutcome,
out: &mut EffectSink,
) {
let raw_body = match outcome {
PortOutcome::Ok(reply) => {
match crate::http_envelope::unwrap_success_envelope(reply.0.as_ref(), "posts/get") {
Ok(raw) => raw,
Err(error) => {
tracing::warn!(req_id, error = ?error, "exact posts HTTP envelope invalid");
out.push(crate::read_relay::emit_read_error(
req_id,
"exact posts response envelope invalid",
));
return;
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(req_id, error = ?error, "exact posts HTTP request failed");
out.push(crate::read_relay::emit_read_error(
req_id,
"exact posts http failed",
));
return;
}
};
let posts = match parse_exact_posts_body(&raw_body) {
Ok(posts) => posts,
Err(error) => {
tracing::warn!(req_id, error = ?error, "exact posts response body invalid");
out.push(crate::read_relay::emit_read_error(
req_id,
"exact posts response body invalid",
));
return;
}
};
let (accepted_keys, ops) = self.collect_exact_posts_cache_ops(&requested_ids, posts);
if !ops.is_empty() && self.local_store_mode == LocalStoreMode::Disabled {
tracing::warn!(req_id, "exact posts requires durable local storage");
out.push(crate::read_relay::emit_read_error(
req_id,
"exact posts durable storage unavailable",
));
return;
}
let corr = self.alloc_corr_internal();
out.push(Effect::Persist { corr, ops });
self.state.corr_map.insert(
corr,
CorrelationContext::ExactPostsPersist {
req_id: req_id.to_string(),
accepted_keys,
},
);
}
fn collect_exact_posts_cache_ops(
&self,
requested_ids: &[String],
posts: Vec<Value>,
) -> (Vec<String>, Vec<StorageOp>) {
let requested = requested_ids.iter().cloned().collect::<HashSet<_>>();
let mut matched = HashSet::new();
let mut accepted_keys = Vec::new();
let mut accepted_storage = HashSet::new();
let mut ops = Vec::new();
for post in posts {
let fields = crate::ws::parser::extract_post_fields(&post);
let candidates = [fields.id.as_str(), fields.temporary_id.as_str()]
.into_iter()
.filter(|id| !id.is_empty() && requested.contains(*id))
.map(str::to_string)
.collect::<Vec<_>>();
if candidates.is_empty() || candidates.iter().all(|id| matched.contains(id)) {
continue;
}
let Some(channel_id) = crate::state::ChannelId::from_str(fields.channel_id.as_str())
else {
tracing::warn!("exact posts row has invalid channel id; ignoring row");
continue;
};
let (visible_rows, mut row_ops) =
match crate::query::local_first::visible_remote_rows_and_cache_ops(
channel_id,
vec![post],
&[],
self.config.auth_user_id.as_str(),
) {
Ok(result) => result,
Err(error) => {
tracing::warn!(error = ?error, "exact posts row failed closed");
continue;
}
};
let Some(normalized) = visible_rows.into_iter().next() else {
continue;
};
let normalized_fields = crate::ws::parser::extract_post_fields(&normalized);
let Some(storage_key) = exact_storage_key(&normalized_fields) else {
continue;
};
if !accepted_storage.insert(storage_key.clone()) {
continue;
}
for id in candidates {
matched.insert(id);
}
accepted_keys.push(storage_key);
protect_exact_read_only_columns(&mut row_ops);
ops.extend(row_ops);
}
(accepted_keys, ops)
}
pub(super) fn handle_exact_posts_persist_reply(
&mut self,
req_id: String,
accepted_keys: Vec<String>,
outcome: &PortOutcome,
out: &mut EffectSink,
) {
if let PortOutcome::Err(error) = outcome {
tracing::warn!(%req_id, error = ?error, "exact posts persist failed");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"exact posts persist failed",
));
return;
}
if accepted_keys.is_empty() {
out.push(crate::read_relay::emit_read_body(
req_id.as_str(),
serde_json::json!({"reqId": req_id, "posts": []}),
));
return;
}
self.schedule_exact_posts_readback(req_id, accepted_keys, 1, Vec::new(), out);
}
pub(super) fn handle_exact_posts_readback_reply(
&mut self,
req_id: String,
accepted_keys: Vec<String>,
next_index: usize,
mut rows: Vec<Value>,
outcome: &PortOutcome,
out: &mut EffectSink,
) {
let row = match outcome {
PortOutcome::Ok(reply) => {
match crate::query::local_first::parse_local_rows(reply.0.as_ref()) {
Ok(rows) => rows.into_iter().next(),
Err(error) => {
tracing::warn!(%req_id, error = ?error, "exact posts durable readback malformed");
None
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(%req_id, error = ?error, "exact posts durable readback failed");
None
}
};
let Some(row) = row else {
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"exact posts durable readback failed",
));
return;
};
rows.push(row);
if next_index < accepted_keys.len() {
self.schedule_exact_posts_readback(req_id, accepted_keys, next_index + 1, rows, out);
return;
}
let projected = project_exact_posts(&accepted_keys, rows);
out.push(crate::read_relay::emit_read_body(
req_id.as_str(),
serde_json::json!({"reqId": req_id, "posts": projected}),
));
}
pub(super) fn handle_initial_window_http_reply(
&mut self,
req_id: String,
post_id: String,
page_size: u32,
outcome: &PortOutcome,
out: &mut EffectSink,
) {
let raw_body = match outcome {
PortOutcome::Ok(reply) => match crate::http_envelope::unwrap_success_envelope(
reply.0.as_ref(),
"posts/getPostsAfterIndex",
) {
Ok(raw) => raw,
Err(error) => {
tracing::warn!(%req_id, error = ?error, "initial window HTTP envelope invalid");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window response envelope invalid",
));
return;
}
},
PortOutcome::Err(error) => {
tracing::warn!(%req_id, error = ?error, "initial window HTTP request failed");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window http failed",
));
return;
}
};
let mut posts =
match crate::query::local_first::parse_initial_window_posts(&raw_body, &post_id) {
Ok(posts) => posts,
Err(error) => {
tracing::warn!(%req_id, error = ?error, "initial window response body invalid");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window target must be first",
));
return;
}
};
posts.truncate(page_size as usize);
let (accepted_keys, ops) = match self.collect_initial_window_cache_ops(&post_id, posts) {
Ok(result) => result,
Err(error) => {
tracing::warn!(%req_id, error = ?error, "initial window rows failed closed");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window rows invalid",
));
return;
}
};
if !ops.is_empty() && self.local_store_mode == LocalStoreMode::Disabled {
tracing::warn!(%req_id, "initial window requires durable local storage");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window durable storage unavailable",
));
return;
}
let corr = self.alloc_corr_internal();
out.push(Effect::Persist { corr, ops });
self.state.corr_map.insert(
corr,
CorrelationContext::InitialWindowPersist {
req_id,
post_id,
accepted_keys,
},
);
}
fn collect_initial_window_cache_ops(
&self,
post_id: &str,
posts: Vec<Value>,
) -> Result<(Vec<String>, Vec<StorageOp>), ImError> {
if posts.is_empty() {
return Ok((Vec::new(), Vec::new()));
}
let first_fields = crate::ws::parser::extract_post_fields(&posts[0]);
let channel_id = crate::state::ChannelId::from_str(first_fields.channel_id.as_str())
.ok_or_else(|| ImError::Parse("initial window row missing channelId".to_string()))?;
let (rows, mut ops) = crate::query::local_first::visible_remote_rows_and_cache_ops(
channel_id,
posts,
&[],
self.config.auth_user_id.as_str(),
)?;
if rows.is_empty() || !crate::query::local_first::post_matches_identity(&rows[0], post_id) {
return Ok((Vec::new(), Vec::new()));
}
let mut accepted_keys = Vec::with_capacity(rows.len());
let mut seen = HashSet::new();
for row in rows {
let fields = crate::ws::parser::extract_post_fields(&row);
let Some(key) = exact_storage_key(&fields) else {
continue;
};
if seen.insert(key.clone()) {
accepted_keys.push(key);
}
}
protect_exact_read_only_columns(&mut ops);
Ok((accepted_keys, ops))
}
pub(super) fn handle_initial_window_persist_reply(
&mut self,
req_id: String,
post_id: String,
accepted_keys: Vec<String>,
outcome: &PortOutcome,
out: &mut EffectSink,
) {
if let PortOutcome::Err(error) = outcome {
tracing::warn!(%req_id, error = ?error, "initial window persist failed");
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window persist failed",
));
return;
}
if accepted_keys.is_empty() {
out.push(crate::read_relay::emit_read_body(
req_id.as_str(),
serde_json::json!({"reqId": req_id, "posts": []}),
));
return;
}
self.schedule_initial_window_readback(req_id, post_id, accepted_keys, 1, Vec::new(), out);
}
pub(super) fn handle_initial_window_readback_reply(
&mut self,
req_id: String,
post_id: String,
accepted_keys: Vec<String>,
next_index: usize,
mut rows: Vec<Value>,
outcome: &PortOutcome,
out: &mut EffectSink,
) {
let row = match outcome {
PortOutcome::Ok(reply) => {
match crate::query::local_first::parse_local_rows(reply.0.as_ref()) {
Ok(rows) => rows.into_iter().next(),
Err(error) => {
tracing::warn!(%req_id, error = ?error, "initial window durable readback malformed");
None
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(%req_id, error = ?error, "initial window durable readback failed");
None
}
};
let Some(row) = row else {
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window durable readback failed",
));
return;
};
rows.push(row);
if next_index < accepted_keys.len() {
self.schedule_initial_window_readback(
req_id,
post_id,
accepted_keys,
next_index + 1,
rows,
out,
);
return;
}
if rows
.first()
.is_none_or(|row| !crate::query::local_first::post_matches_identity(row, &post_id))
{
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window target missing from durable readback",
));
return;
}
let projected = project_exact_posts(&accepted_keys, rows);
if projected
.first()
.is_none_or(|row| !crate::query::local_first::post_matches_identity(row, &post_id))
{
out.push(crate::read_relay::emit_read_error(
req_id.as_str(),
"initial window target is not first in durable result",
));
return;
}
out.push(crate::read_relay::emit_read_body(
req_id.as_str(),
serde_json::json!({"reqId": req_id, "posts": projected}),
));
}
fn schedule_initial_window_readback(
&mut self,
req_id: String,
post_id: String,
accepted_keys: Vec<String>,
next_index: usize,
rows: Vec<Value>,
out: &mut EffectSink,
) {
let key = accepted_keys[next_index - 1].clone();
let corr = self.alloc_corr_internal();
out.push(Effect::Persist {
corr,
ops: vec![StorageOp::Scan(exact_posts_scan_spec(&key))],
});
self.state.corr_map.insert(
corr,
CorrelationContext::InitialWindowReadback {
req_id,
post_id,
accepted_keys,
next_index,
rows,
},
);
}
fn schedule_exact_posts_readback(
&mut self,
req_id: String,
accepted_keys: Vec<String>,
next_index: usize,
rows: Vec<Value>,
out: &mut EffectSink,
) {
let key = accepted_keys[next_index - 1].clone();
let corr = self.alloc_corr_internal();
out.push(Effect::Persist {
corr,
ops: vec![StorageOp::Scan(exact_posts_scan_spec(&key))],
});
self.state.corr_map.insert(
corr,
CorrelationContext::ExactPostsReadback {
req_id,
accepted_keys,
next_index,
rows,
},
);
}
pub(super) fn handle_draft_readback_reply(
&self,
command: Box<crate::draft::QueryDraftCommand>,
outcome: &PortOutcome,
out: &mut EffectSink,
) -> Result<(), ImError> {
let req_id = command.req_id.as_deref().unwrap_or_default();
let PortOutcome::Ok(reply) = outcome else {
out.push(crate::read_relay::emit_read_error(
req_id,
"draft read failed",
));
return Ok(());
};
let rows = match helix_core::port_codec::rows_from_reply_bytes(&reply.0) {
Ok(rows) => rows,
Err(error) => {
tracing::warn!(error = ?error, "draft readback decode failed");
out.push(crate::read_relay::emit_read_error(
req_id,
"draft readback decode failed",
));
return Ok(());
}
};
out.push(command.result_event(&rows)?.into_effect());
Ok(())
}
pub(super) fn handle_todo_query_reply(&self, outcome: &PortOutcome, out: &mut EffectSink) {
match outcome {
PortOutcome::Ok(reply) => match unwrap_sync_envelope(reply.0.as_ref()) {
Ok(raw_body) => {
out.push(crate::todo::emit_todo_updated(raw_body.as_ref()));
}
Err(e) => {
tracing::warn!(error = ?e, "todo query reply envelope decode failed");
}
},
PortOutcome::Err(e) => {
tracing::warn!(error = ?e, "todo query outbound http failed");
}
}
}
pub(super) fn handle_load_older_context_reply(
&mut self,
mut state: Box<LoadOlderState>,
outcome: &PortOutcome,
now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
match outcome {
PortOutcome::Ok(reply) => {
let rows = unwrap_sync_envelope(reply.0.as_ref())
.map(|raw| crate::older_context::extract_post_rows(&raw))
.unwrap_or_default();
match state.ingest_round(&rows) {
crate::older_context::RoundDecision::Continue => {
let next_corr = self.alloc_corr_internal();
let (_, body) = crate::older_context::build_post_context_body(&state);
out.push(crate::older_context::post_context_http_tracked(
&self.config.api_base_url,
&body,
next_corr,
self.state.connection_id.as_deref(),
state.request_id(),
));
self.state
.corr_map
.insert(next_corr, CorrelationContext::LoadOlderContext { state });
}
crate::older_context::RoundDecision::Done => {
self.finish_load_older_http(state, now_ms, out)?;
}
}
}
PortOutcome::Err(e) => {
tracing::warn!(error = ?e, "load_older_context postContext http failed");
self.emit_load_older_failed_event(&state, now_ms, out)?;
}
}
Ok(())
}
fn finish_load_older_http(
&mut self,
state: Box<LoadOlderState>,
now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
if state.failed() {
return self.emit_load_older_failed_event(&state, now_ms, out);
}
let (_, ops) = match crate::query::local_first::visible_remote_rows_and_cache_ops(
*state.channel_id(),
state.older_rows(),
&[],
self.config.auth_user_id.as_str(),
) {
Ok(result) => result,
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"load_older_context response could not become durable messages"
);
return self.emit_load_older_failed_event(&state, now_ms, out);
}
};
if ops.is_empty() {
return self.schedule_load_older_readback(state, out);
}
if self.local_store_mode == LocalStoreMode::Disabled {
tracing::warn!(
channel_id = state.channel_id().as_str(),
"load_older_context requires a writable store for durable timeline publication"
);
return self.emit_load_older_failed_event(&state, now_ms, out);
}
let corr = self.alloc_corr_internal();
out.push(Effect::Persist { corr, ops });
self.state
.corr_map
.insert(corr, CorrelationContext::LoadOlderCache { state });
Ok(())
}
pub(super) fn schedule_load_older_readback(
&mut self,
state: Box<LoadOlderState>,
out: &mut EffectSink,
) -> Result<(), ImError> {
let window_token = state.window_token().ok_or_else(|| {
ImError::Parse("load older missing attached window token".to_string())
})?;
let request = crate::query::MessageQueryRequest {
channel_id: *state.channel_id(),
limit: state.readback_limit(),
window_token: window_token.to_string(),
};
let corr = self.alloc_corr_internal();
out.push(crate::query::build_message_query_from_request(
&request, corr,
));
self.state
.corr_map
.insert(corr, CorrelationContext::LoadOlderReadback { state });
Ok(())
}
pub(super) fn handle_load_older_readback_reply(
&mut self,
state: Box<LoadOlderState>,
outcome: &PortOutcome,
now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
let mut rows_desc = match outcome {
PortOutcome::Ok(reply) => {
match crate::query::local_first::parse_local_rows(reply.0.as_ref()) {
Ok(rows) => rows,
Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"load_older_context readback was malformed"
);
return self.emit_load_older_failed_event(&state, now_ms, out);
}
}
}
PortOutcome::Err(error) => {
tracing::warn!(
channel_id = state.channel_id().as_str(),
error = ?error,
"load_older_context readback failed"
);
return self.emit_load_older_failed_event(&state, now_ms, out);
}
};
crate::query::local_first::sort_recent_rows_desc(&mut rows_desc);
let window_token = state.window_token().ok_or_else(|| {
ImError::Parse("load older missing attached window token".to_string())
})?;
let request = crate::query::MessageQueryRequest {
channel_id: *state.channel_id(),
limit: state.readback_limit(),
window_token: window_token.to_string(),
};
let page = crate::timeline_state::WindowPage {
window_token: window_token.to_string(),
has_older: state.has_more(),
has_newer: false,
has_more: state.has_more(),
};
out.push(self.emit_timeline_snapshot_with_causation(
&request,
&rows_desc,
now_ms,
None,
Some(page),
)?);
Ok(())
}
pub(super) fn emit_load_older_failed_event(
&mut self,
state: &LoadOlderState,
_now_ms: u64,
out: &mut EffectSink,
) -> Result<(), ImError> {
out.push(
crate::event::timeline::page(
state.channel_id().as_str(),
state.window_token().unwrap_or("latest"),
"older",
"failed",
Vec::new(),
true,
false,
Some(state.anchor_post_id()),
)?
.into_effect(),
);
Ok(())
}
}
fn older_navigation_scan_spec(
state: &crate::timeline_navigation::TimelineNavigationState,
) -> ScanSpec {
ScanSpec {
table: "message",
limit: Some(state.page_size().saturating_add(1)),
filter: Some((
"channel_id",
SqlValue::Text(state.channel_id().as_str().to_string()),
)),
order_by: OLDER_NAVIGATION_ORDER,
}
}
fn timeline_navigation_readback_keys(
state: &crate::timeline_navigation::TimelineNavigationState,
) -> Vec<crate::timeline_navigation::TimelineNavigationReadbackKey> {
let mut seen = HashSet::new();
state
.rows()
.iter()
.filter_map(timeline_navigation_readback_key)
.filter(|key| seen.insert((key.column, key.value.clone())))
.collect()
}
fn timeline_navigation_readback_key(
row: &Value,
) -> Option<crate::timeline_navigation::TimelineNavigationReadbackKey> {
let fields = crate::ws::parser::extract_post_fields(row);
if !fields.temporary_id.is_empty() {
Some(crate::timeline_navigation::TimelineNavigationReadbackKey {
column: "temporary_id",
value: fields.temporary_id,
})
} else if !fields.id.is_empty() {
Some(crate::timeline_navigation::TimelineNavigationReadbackKey {
column: "id",
value: fields.id,
})
} else {
None
}
}
fn newer_navigation_scan_spec(
state: &crate::timeline_navigation::TimelineNavigationState,
) -> ScanSpec {
ScanSpec {
table: "message",
limit: Some(state.page_size().saturating_add(1)),
filter: Some((
"channel_id",
SqlValue::Text(state.channel_id().as_str().to_string()),
)),
order_by: NEWER_NAVIGATION_ORDER,
}
}
fn local_timeline_navigation_rows(
state: &crate::timeline_navigation::TimelineNavigationState,
coverage: &crate::timeline_navigation::TimelineNavigationCoverage,
viewer_user_id: &str,
) -> Option<Vec<Value>> {
match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Locate {
target_message_id, ..
} => {
let rows = crate::query::render_ready::locate::normalize_durable_located_window(
coverage.rows.clone(),
state.channel_id().as_str(),
target_message_id,
viewer_user_id,
state.page_size(),
)
.ok()?;
let target_index = rows.iter().position(|row| {
crate::timeline_navigation::row_matches_message_identity(row, target_message_id)
})?;
(coverage.target_index == Some(target_index)).then_some(rows)
}
crate::timeline_navigation::TimelineNavigationKind::Older { .. }
| crate::timeline_navigation::TimelineNavigationKind::Newer { .. } => {
Some(coverage.rows.clone())
}
}
}
fn timeline_navigation_direction(
state: &crate::timeline_navigation::TimelineNavigationState,
) -> &'static str {
match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Older { .. } => "older",
crate::timeline_navigation::TimelineNavigationKind::Newer { .. } => "newer",
crate::timeline_navigation::TimelineNavigationKind::Locate { .. } => "locate",
}
}
fn older_navigation_rows(
state: &crate::timeline_navigation::TimelineNavigationState,
rows: Vec<Value>,
auth_user_id: &str,
) -> Vec<Value> {
if state.rows().is_empty() {
return Vec::new();
}
let cursor = match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Older { cursor, .. }
if !state.is_windowless() =>
{
Some(cursor)
}
crate::timeline_navigation::TimelineNavigationKind::Older { .. } => None,
_ => return Vec::new(),
};
let mut seen = HashSet::new();
let mut filtered = rows
.into_iter()
.filter_map(|row| {
let fields = crate::ws::parser::extract_post_fields(&row);
if fields.channel_id != state.channel_id().as_str()
|| !crate::channel_write::post_updates_from_fields(
state.channel_id(),
&fields,
auth_user_id,
)
.visible
{
return None;
}
let key = durable_navigation_key(&row)?;
if let Some(cursor) = cursor {
if key.0 > cursor.create_at
|| (key.0 == cursor.create_at && key.1.as_str() >= cursor.temporary_id.as_str())
{
return None;
}
}
let identity = if !fields.id.is_empty() {
format!("id:{}", fields.id)
} else if !fields.temporary_id.is_empty() {
format!("tmp:{}", fields.temporary_id)
} else {
return None;
};
seen.insert(identity).then_some((key, row))
})
.collect::<Vec<_>>();
filtered.sort_by(|left, right| left.0.cmp(&right.0));
filtered
.into_iter()
.take(state.page_size() as usize)
.map(|(_, row)| row)
.collect()
}
fn newer_http_has_rows(state: &crate::timeline_navigation::TimelineNavigationState) -> bool {
let Some(cursor) = (match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Newer { cursor, .. } => Some(cursor),
_ => None,
}) else {
return false;
};
state
.authority_rows()
.iter()
.filter_map(durable_navigation_key)
.any(|key| {
key.0 > cursor.create_at
|| (key.0 == cursor.create_at && key.1.as_str() > cursor.temporary_id.as_str())
})
}
fn newer_navigation_rows(
state: &crate::timeline_navigation::TimelineNavigationState,
rows: Vec<Value>,
auth_user_id: &str,
) -> Vec<Value> {
if !newer_http_has_rows(state) {
return Vec::new();
}
let cursor = match state.kind() {
crate::timeline_navigation::TimelineNavigationKind::Newer { cursor, .. }
if !state.is_windowless() =>
{
Some(cursor)
}
crate::timeline_navigation::TimelineNavigationKind::Newer { .. } => None,
_ => return Vec::new(),
};
let mut seen = HashSet::new();
let mut filtered = rows
.into_iter()
.filter_map(|row| {
let fields = crate::ws::parser::extract_post_fields(&row);
if fields.channel_id != state.channel_id().as_str()
|| !crate::channel_write::post_updates_from_fields(
state.channel_id(),
&fields,
auth_user_id,
)
.visible
{
return None;
}
let key = durable_navigation_key(&row)?;
if let Some(cursor) = cursor {
if key.0 < cursor.create_at
|| (key.0 == cursor.create_at && key.1.as_str() <= cursor.temporary_id.as_str())
{
return None;
}
}
let identity = if !fields.id.is_empty() {
format!("id:{}", fields.id)
} else if !fields.temporary_id.is_empty() {
format!("tmp:{}", fields.temporary_id)
} else {
return None;
};
seen.insert(identity).then_some((key, row))
})
.collect::<Vec<_>>();
filtered.sort_by(|left, right| left.0.cmp(&right.0));
filtered
.into_iter()
.take(state.page_size() as usize)
.map(|(_, row)| row)
.collect()
}
fn durable_navigation_key(row: &Value) -> Option<(i64, String)> {
let create_at = row
.get("create_at")
.or_else(|| row.get("createAt"))
.and_then(Value::as_i64)?;
let temporary_id = row
.get("temporary_id")
.or_else(|| row.get("temporaryId"))
.and_then(Value::as_str)
.filter(|value| !value.is_empty())?;
Some((create_at, temporary_id.to_string()))
}
fn durable_row_matches_readback_key(
row: &Value,
key: &crate::timeline_navigation::TimelineNavigationReadbackKey,
) -> bool {
let fields = crate::ws::parser::extract_post_fields(row);
match key.column {
"temporary_id" => fields.temporary_id == key.value,
"id" => fields.id == key.value,
_ => false,
}
}
fn parse_exact_posts_body(raw: &[u8]) -> Result<Vec<Value>, ImError> {
let root: Value = serde_json::from_slice(raw)
.map_err(|error| ImError::Parse(format!("posts/get body: {error}")))?;
let status = root
.get("status")
.and_then(Value::as_str)
.ok_or_else(|| ImError::Parse("posts/get body missing string status".to_string()))?;
if !status.eq_ignore_ascii_case("SUCCESS") {
return Err(ImError::Parse(format!("posts/get backend status {status}")));
}
let rows = root
.get("data")
.and_then(Value::as_array)
.ok_or_else(|| ImError::Parse("posts/get response missing data array".to_string()))?;
if rows.iter().any(|row| !row.is_object()) {
return Err(ImError::Parse(
"posts/get data rows must be objects".to_string(),
));
}
Ok(rows.clone())
}
fn exact_storage_key(fields: &crate::sync_session::PostFields) -> Option<String> {
if !fields.temporary_id.is_empty() {
Some(fields.temporary_id.clone())
} else if !fields.id.is_empty() {
Some(fields.id.clone())
} else {
None
}
}
fn protect_exact_read_only_columns(ops: &mut Vec<StorageOp>) {
for op in ops {
if let StorageOp::BatchUpsert(spec) = op {
for column in ["create_at", "read_bits", "event_seq"] {
if !spec.exclude_from_update.contains(&column) {
spec.exclude_from_update.push(column);
}
}
}
}
}
fn timeline_navigation_exact_scan_spec(
key: &crate::timeline_navigation::TimelineNavigationReadbackKey,
) -> ScanSpec {
ScanSpec {
table: "message",
limit: Some(1),
filter: Some((key.column, SqlValue::Text(key.value.clone()))),
order_by: &[],
}
}
fn exact_posts_scan_spec(key: &str) -> ScanSpec {
ScanSpec {
table: "message",
limit: Some(1),
filter: Some(("temporary_id", SqlValue::Text(key.to_string()))),
order_by: &[],
}
}
fn exact_json_object(raw: &str) -> Value {
serde_json::from_str::<Value>(raw)
.ok()
.filter(Value::is_object)
.unwrap_or_else(|| serde_json::json!({}))
}
fn project_exact_posts(accepted_keys: &[String], rows: Vec<Value>) -> Vec<Value> {
let accepted = accepted_keys.iter().collect::<HashSet<_>>();
let mut seen = HashSet::new();
let mut projected = Vec::new();
for row in rows {
let fields = crate::ws::parser::extract_post_fields(&row);
let Some(key) = exact_storage_key(&fields) else {
continue;
};
if !accepted.contains(&key) || !seen.insert(key) {
continue;
}
projected.push(serde_json::json!({
"id": if fields.id.is_empty() { fields.temporary_id.clone() } else { fields.id.clone() },
"temporaryId": fields.temporary_id,
"channelId": fields.channel_id,
"userId": fields.user_id,
"userSnapshot": exact_json_object(&fields.user_snapshot),
"type": fields.msg_type,
"message": fields.message,
"simpleMessage": fields.simple_message,
"createAt": fields.create_at,
"updateAt": fields.update_at,
"sendStatus": "sent",
"viewers": fields.viewers,
"mentions": fields.mentions,
"props": exact_json_object(&fields.props),
"topic": exact_json_object(&fields.topic),
"replyId": fields.reply_id,
"replyRootId": fields.reply_root_id,
"replyFirstLevelId": fields.reply_first_level_id,
"replyCount": fields.reply_count,
"repliedMessage": exact_json_object(&fields.replied_message),
"replyMessages": exact_json_object(&fields.reply_messages),
}));
}
projected
}