use std::sync::Arc;
use tracing::warn;
use super::common::AppliedPosition;
use crate::control::distributed_applier::{AppliedWrite, ProposeTracker};
use crate::control::state::SharedState;
pub(crate) struct ArraySchemaPayload<'a> {
pub array: &'a str,
pub snapshot_payload: &'a [u8],
pub schema_hlc_bytes: [u8; 18],
}
pub(crate) fn apply_array_schema(
state: &Arc<SharedState>,
tracker: &Arc<ProposeTracker>,
pos: AppliedPosition,
payload: ArraySchemaPayload<'_>,
) -> bool {
let AppliedPosition {
group_id,
log_index,
applied_key,
} = pos;
use nodedb_array::sync::hlc::Hlc;
let ArraySchemaPayload {
array,
snapshot_payload,
schema_hlc_bytes,
} = payload;
let remote_hlc = Hlc::from_bytes(&schema_hlc_bytes);
if let Err(e) =
state
.array_sync_schemas
.import_snapshot_replicated(array, snapshot_payload, remote_hlc)
{
warn!(
group_id, index = log_index, array = %array, error = %e,
"apply_array_schema: import_snapshot_replicated failed"
);
tracker.complete(
group_id,
log_index,
applied_key,
Err(crate::Error::Internal {
detail: format!("schema import: {e}"),
}),
);
return false;
}
if let Err(e) =
crate::control::array_sync::catalog_register::register_array_catalog_entry(state, array)
{
warn!(
group_id, index = log_index, array = %array, error = %e,
"apply_array_schema: register_array_catalog_entry failed"
);
tracker.complete(
group_id,
log_index,
applied_key,
Err(crate::Error::Internal {
detail: format!("array catalog register: {e}"),
}),
);
return false;
}
tracker.complete(
group_id,
log_index,
applied_key,
Ok(AppliedWrite::unversioned(Vec::new())),
);
true
}