Skip to content
Merged
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
251 changes: 247 additions & 4 deletions crates/event-sorcery-ffi/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,14 +25,15 @@ use url::form_urlencoded;
#[cfg(test)]
use event_sorcery::CommitRequest;
use event_sorcery::{
DeadReason, Engine, EngineError, JobClaimHandle, JobClaimResult, JobId, JobLeaseResult,
JobRuntime, JobRuntimeError, JobSeed, JobSettlementResult, JobStoreError, OpaqueCommitRequest,
OpaqueProposedEvent, OpaqueStoredEvent, SnapshotWrite, SqliteJobError, StreamIdentity,
CompactionPolicy, DeadReason, Engine, EngineError, JobClaimHandle, JobClaimResult, JobId,
JobLeaseResult, JobRuntime, JobRuntimeError, JobSeed, JobSettlementResult, JobStoreError,
OpaqueCommitRequest, OpaqueProposedEvent, OpaqueStoredEvent, ReconcileError,
SchemaReconciliation, SchemaTarget, SnapshotWrite, SqliteJobError, StreamIdentity,
decode_opaque_payload, encode_opaque_payload,
};

const ABI_MAJOR: u32 = 0;
const ABI_MINOR: u32 = 4;
const ABI_MINOR: u32 = 5;
const ES_OK: i32 = 0;
const ES_ERR_DECODE: i32 = 1;
const ES_ERR_CONFLICT: i32 = 2;
Expand Down Expand Up @@ -594,6 +595,76 @@ pub unsafe extern "C" fn es_snapshot_discard(
})
}

#[unsafe(no_mangle)]
/// Reconciles aggregate schema metadata through the shared engine.
///
/// The reported tag is `0` when the stored version changed and any stale
/// snapshots were cleared, and `1` when it already matched.
///
/// # Safety
///
/// `store` must be the original owner cell passed to [`es_open`]. `request`
/// must reference a readable buffer. `out_reconciliation` must be writable, and
/// `out_error` must be null or writable. The request storage may alias either
/// output. `store`, `out_reconciliation`, and `out_error` must all be distinct.
pub unsafe extern "C" fn es_schema_reconcile(
store: *mut *mut EsStore,
request: *const EsBuf,
out_reconciliation: *mut u8,
out_error: *mut EsBuf,
) -> i32 {
let target = decode_schema_target(request);
clear_tag_output(out_reconciliation);
let lease = match EsStore::acquire(store) {
Ok(lease) => lease,
Err(error) => return write_error(out_error, error),
};
ffi_call(Some(&lease.store), out_error, || {
if out_reconciliation.is_null() {
return Err(AbiError::State("out_reconciliation is null"));
}
let target = target?;
let reconciliation = lease
.inner
.runtime
.block_on(lease.inner.engine.reconcile_schema(&target))
.map_err(AbiError::from)?;
let tag = match reconciliation {
SchemaReconciliation::Changed => 0,
SchemaReconciliation::Unchanged => 1,
};
unsafe { out_reconciliation.write(tag) };
Ok(())
})
}

#[unsafe(no_mangle)]
/// Records an aggregate schema version after its recovery completes.
///
/// # Safety
///
/// `store` must be the original owner cell passed to [`es_open`]. `request`
/// must reference a readable buffer. `out_error` must be null or writable. The
/// owner cell and `out_error` must be distinct.
pub unsafe extern "C" fn es_schema_record(
store: *mut *mut EsStore,
request: *const EsBuf,
out_error: *mut EsBuf,
) -> i32 {
let lease = match EsStore::acquire(store) {
Ok(lease) => lease,
Err(error) => return write_error(out_error, error),
};
ffi_call(Some(&lease.store), out_error, || {
let target = decode_schema_target(request)?;
lease
.inner
.runtime
.block_on(lease.inner.engine.record_schema(&target))
.map_err(AbiError::from)
})
}

#[unsafe(no_mangle)]
/// Enqueues an erased payload through the retained event-sourced job runtime.
///
Expand Down Expand Up @@ -1693,6 +1764,11 @@ impl From<EngineError> for AbiError {
// as malformed input: which entry point the job intent arrived
// through must not change how the caller is told to react.
EngineError::JobFlush(JobStoreError::InvalidInstant(_)) => Self::MalformedInput,
// Clearing a compaction policy that snapshots already relied on is
// permanently refused, so it must not surface as retryable storage.
EngineError::Schema(ReconcileError::CompactedSnapshotClear { .. }) => {
Self::State("compacted snapshot retention cannot be cleared")
}
other => Self::Engine(other),
}
}
Expand Down Expand Up @@ -2222,6 +2298,28 @@ const fn require_identifier_limit(observed: usize) -> Result<(), AbiError> {
})
}

/// Decodes `[version, aggregate_type, schema_version, compaction_policy]`,
/// where the policy is `0` to retain events and `1` to compact them after a
/// snapshot. An unknown policy is malformed input rather than a silent retain:
/// the two policies disagree about whether clearing snapshots can lose state.
fn decode_schema_target(buffer: *const EsBuf) -> Result<SchemaTarget, AbiError> {
let (version, aggregate_type, schema_version, compaction): (u8, String, u64, u8) =
decode(buffer)?;
require_version(version)?;
require_identifier_limit(aggregate_type.len())?;
let compaction = match compaction {
0 => CompactionPolicy::Retain,
1 => CompactionPolicy::CompactAfterSnapshot,
_ => return Err(AbiError::MalformedInput),
};

Ok(SchemaTarget::new(
aggregate_type,
schema_version,
compaction,
))
}

fn decode_open_options(buffer: *const EsBuf) -> Result<OpenOptions, AbiError> {
let (version, path, busy_timeout_ms, pool_size, runtime_threads): (u8, String, u64, u32, u32) =
decode(buffer)?;
Expand Down Expand Up @@ -4814,4 +4912,149 @@ mod tests {
unsafe { es_buf_free(&raw mut error) };
assert_eq!(unsafe { es_close(&raw mut store) }, ES_OK);
}

#[test]
fn schema_exports_preserve_the_existing_recovery_protocol() {
let mut store = ptr::null_mut();
open_store(&mut store);
let mut error = empty_buffer();
let mut encoded = encode_request(&(1_u8, "ffi-schema", 1_u64, 0_u8));
let request = caller_buffer(&mut encoded);
let mut reconciliation = u8::MAX;

assert_eq!(
unsafe {
es_schema_reconcile(
&raw mut store,
&raw const request,
&raw mut reconciliation,
&raw mut error,
)
},
ES_OK
);
assert_eq!(reconciliation, 0);
assert!(error.ptr.is_null());

assert_eq!(
unsafe { es_schema_record(&raw mut store, &raw const request, &raw mut error) },
ES_OK
);

reconciliation = u8::MAX;
assert_eq!(
unsafe {
es_schema_reconcile(
&raw mut store,
&raw const request,
&raw mut reconciliation,
&raw mut error,
)
},
ES_OK
);
assert_eq!(reconciliation, 1);
assert_eq!(unsafe { es_close(&raw mut store) }, ES_OK);
}

#[test]
fn schema_exports_refuse_clearing_a_compacted_version_as_invalid_state() {
let mut store = ptr::null_mut();
open_store(&mut store);
let mut error = empty_buffer();
let mut recorded = encode_request(&(1_u8, "ffi-schema", 1_u64, 1_u8));
let record_request = caller_buffer(&mut recorded);

assert_eq!(
unsafe { es_schema_record(&raw mut store, &raw const record_request, &raw mut error) },
ES_OK
);

let mut advanced = encode_request(&(1_u8, "ffi-schema", 2_u64, 1_u8));
let reconcile_request = caller_buffer(&mut advanced);
let mut reconciliation = u8::MAX;
assert_eq!(
unsafe {
es_schema_reconcile(
&raw mut store,
&raw const reconcile_request,
&raw mut reconciliation,
&raw mut error,
)
},
ES_ERR_STATE
);
assert_eq!(reconciliation, 0);
unsafe { es_buf_free(&raw mut error) };
assert_eq!(unsafe { es_close(&raw mut store) }, ES_OK);
}

#[test]
fn schema_exports_reject_unknown_compaction_policy() {
let mut store = ptr::null_mut();
open_store(&mut store);
let mut error = empty_buffer();
let mut encoded = encode_request(&(1_u8, "ffi-schema", 1_u64, 2_u8));
let request = caller_buffer(&mut encoded);
let mut reconciliation = u8::MAX;

assert_eq!(
unsafe {
es_schema_reconcile(
&raw mut store,
&raw const request,
&raw mut reconciliation,
&raw mut error,
)
},
ES_ERR_DECODE
);
assert_eq!(reconciliation, 0);
unsafe { es_buf_free(&raw mut error) };

assert_eq!(
unsafe { es_schema_record(&raw mut store, &raw const request, &raw mut error) },
ES_ERR_DECODE
);
unsafe { es_buf_free(&raw mut error) };
assert_eq!(unsafe { es_close(&raw mut store) }, ES_OK);
}

#[test]
fn schema_exports_bound_the_aggregate_type_identifier() {
let mut store = ptr::null_mut();
open_store(&mut store);
let mut error = empty_buffer();
let mut encoded =
encode_request(&(1_u8, "x".repeat(MAX_IDENTIFIER_BYTES + 1), 1_u64, 0_u8));
let request = caller_buffer(&mut encoded);
let mut reconciliation = u8::MAX;

assert_eq!(
unsafe {
es_schema_reconcile(
&raw mut store,
&raw const request,
&raw mut reconciliation,
&raw mut error,
)
},
ES_ERR_RESOURCE_LIMIT
);
assert_eq!(reconciliation, 0);
assert_eq!(
decode_error(&error),
(
1,
ES_ERR_RESOURCE_LIMIT,
Value::Array(vec![
Value::Text("identifier_text".to_string()),
Value::Integer((MAX_IDENTIFIER_BYTES + 1).into()),
Value::Integer(MAX_IDENTIFIER_BYTES.into()),
]),
)
);
unsafe { es_buf_free(&raw mut error) };
assert_eq!(unsafe { es_close(&raw mut store) }, ES_OK);
}
}