LCOV - code coverage report
Current view: top level - libs/walproposer/src - api_bindings.rs (source / functions) Coverage Total Hit
Test: 190869232aac3a234374e5bb62582e91cf5f5818.info Lines: 95.9 % 369 354
Test Date: 2024-02-23 13:21:27 Functions: 97.3 % 37 36

            Line data    Source code
       1              : //! A C-Rust shim: defines implementation of C walproposer API, assuming wp
       2              : //! callback_data stores Box to some Rust implementation.
       3              : 
       4              : #![allow(dead_code)]
       5              : 
       6              : use std::ffi::CStr;
       7              : use std::ffi::CString;
       8              : 
       9              : use crate::bindings::uint32;
      10              : use crate::bindings::walproposer_api;
      11              : use crate::bindings::NeonWALReadResult;
      12              : use crate::bindings::PGAsyncReadResult;
      13              : use crate::bindings::PGAsyncWriteResult;
      14              : use crate::bindings::Safekeeper;
      15              : use crate::bindings::Size;
      16              : use crate::bindings::StringInfoData;
      17              : use crate::bindings::TimestampTz;
      18              : use crate::bindings::WalProposer;
      19              : use crate::bindings::WalProposerConnStatusType;
      20              : use crate::bindings::WalProposerConnectPollStatusType;
      21              : use crate::bindings::WalProposerExecStatusType;
      22              : use crate::bindings::WalproposerShmemState;
      23              : use crate::bindings::XLogRecPtr;
      24              : use crate::walproposer::ApiImpl;
      25              : use crate::walproposer::StreamingCallback;
      26              : use crate::walproposer::WaitResult;
      27              : 
      28         1328 : extern "C" fn get_shmem_state(wp: *mut WalProposer) -> *mut WalproposerShmemState {
      29         1328 :     unsafe {
      30         1328 :         let callback_data = (*(*wp).config).callback_data;
      31         1328 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      32         1328 :         (*api).get_shmem_state()
      33         1328 :     }
      34         1328 : }
      35              : 
      36         1278 : extern "C" fn start_streaming(wp: *mut WalProposer, startpos: XLogRecPtr) {
      37         1278 :     unsafe {
      38         1278 :         let callback_data = (*(*wp).config).callback_data;
      39         1278 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      40         1278 :         let callback = StreamingCallback::new(wp);
      41         1278 :         (*api).start_streaming(startpos, &callback);
      42         1278 :     }
      43         1278 : }
      44              : 
      45          859 : extern "C" fn get_flush_rec_ptr(wp: *mut WalProposer) -> XLogRecPtr {
      46          859 :     unsafe {
      47          859 :         let callback_data = (*(*wp).config).callback_data;
      48          859 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      49          859 :         (*api).get_flush_rec_ptr()
      50          859 :     }
      51          859 : }
      52              : 
      53      2513362 : extern "C" fn get_current_timestamp(wp: *mut WalProposer) -> TimestampTz {
      54      2513362 :     unsafe {
      55      2513362 :         let callback_data = (*(*wp).config).callback_data;
      56      2513362 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      57      2513362 :         (*api).get_current_timestamp()
      58      2513362 :     }
      59      2513362 : }
      60              : 
      61        86767 : extern "C" fn conn_error_message(sk: *mut Safekeeper) -> *mut ::std::os::raw::c_char {
      62        86767 :     unsafe {
      63        86767 :         let callback_data = (*(*(*sk).wp).config).callback_data;
      64        86767 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      65        86767 :         let msg = (*api).conn_error_message(&mut (*sk));
      66        86767 :         let msg = CString::new(msg).unwrap();
      67        86767 :         // TODO: fix leaking error message
      68        86767 :         msg.into_raw()
      69        86767 :     }
      70        86767 : }
      71              : 
      72       266484 : extern "C" fn conn_status(sk: *mut Safekeeper) -> WalProposerConnStatusType {
      73       266484 :     unsafe {
      74       266484 :         let callback_data = (*(*(*sk).wp).config).callback_data;
      75       266484 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      76       266484 :         (*api).conn_status(&mut (*sk))
      77       266484 :     }
      78       266484 : }
      79              : 
      80       266484 : extern "C" fn conn_connect_start(sk: *mut Safekeeper) {
      81       266484 :     unsafe {
      82       266484 :         let callback_data = (*(*(*sk).wp).config).callback_data;
      83       266484 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      84       266484 :         (*api).conn_connect_start(&mut (*sk))
      85       266484 :     }
      86       266484 : }
      87              : 
      88       239782 : extern "C" fn conn_connect_poll(sk: *mut Safekeeper) -> WalProposerConnectPollStatusType {
      89       239782 :     unsafe {
      90       239782 :         let callback_data = (*(*(*sk).wp).config).callback_data;
      91       239782 :         let api = callback_data as *mut Box<dyn ApiImpl>;
      92       239782 :         (*api).conn_connect_poll(&mut (*sk))
      93       239782 :     }
      94       239782 : }
      95              : 
      96       239782 : extern "C" fn conn_send_query(sk: *mut Safekeeper, query: *mut ::std::os::raw::c_char) -> bool {
      97       239782 :     let query = unsafe { CStr::from_ptr(query) };
      98       239782 :     let query = query.to_str().unwrap();
      99       239782 : 
     100       239782 :     unsafe {
     101       239782 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     102       239782 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     103       239782 :         (*api).conn_send_query(&mut (*sk), query)
     104       239782 :     }
     105       239782 : }
     106              : 
     107       239782 : extern "C" fn conn_get_query_result(sk: *mut Safekeeper) -> WalProposerExecStatusType {
     108       239782 :     unsafe {
     109       239782 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     110       239782 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     111       239782 :         (*api).conn_get_query_result(&mut (*sk))
     112       239782 :     }
     113       239782 : }
     114              : 
     115            0 : extern "C" fn conn_flush(sk: *mut Safekeeper) -> ::std::os::raw::c_int {
     116            0 :     unsafe {
     117            0 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     118            0 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     119            0 :         (*api).conn_flush(&mut (*sk))
     120            0 :     }
     121            0 : }
     122              : 
     123        88975 : extern "C" fn conn_finish(sk: *mut Safekeeper) {
     124        88975 :     unsafe {
     125        88975 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     126        88975 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     127        88975 :         (*api).conn_finish(&mut (*sk))
     128        88975 :     }
     129        88975 : }
     130              : 
     131       139414 : extern "C" fn conn_async_read(
     132       139414 :     sk: *mut Safekeeper,
     133       139414 :     buf: *mut *mut ::std::os::raw::c_char,
     134       139414 :     amount: *mut ::std::os::raw::c_int,
     135       139414 : ) -> PGAsyncReadResult {
     136       139414 :     unsafe {
     137       139414 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     138       139414 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     139       139414 : 
     140       139414 :         // This function has guarantee that returned buf will be valid until
     141       139414 :         // the next call. So we can store a Vec in each Safekeeper and reuse
     142       139414 :         // it on the next call.
     143       139414 :         let mut inbuf = take_vec_u8(&mut (*sk).inbuf).unwrap_or_default();
     144       139414 :         inbuf.clear();
     145       139414 : 
     146       139414 :         let result = (*api).conn_async_read(&mut (*sk), &mut inbuf);
     147       139414 : 
     148       139414 :         // Put a Vec back to sk->inbuf and return data ptr.
     149       139414 :         *amount = inbuf.len() as i32;
     150       139414 :         *buf = store_vec_u8(&mut (*sk).inbuf, inbuf);
     151       139414 : 
     152       139414 :         result
     153       139414 :     }
     154       139414 : }
     155              : 
     156        35254 : extern "C" fn conn_async_write(
     157        35254 :     sk: *mut Safekeeper,
     158        35254 :     buf: *const ::std::os::raw::c_void,
     159        35254 :     size: usize,
     160        35254 : ) -> PGAsyncWriteResult {
     161        35254 :     unsafe {
     162        35254 :         let buf = std::slice::from_raw_parts(buf as *const u8, size);
     163        35254 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     164        35254 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     165        35254 :         (*api).conn_async_write(&mut (*sk), buf)
     166        35254 :     }
     167        35254 : }
     168              : 
     169       269802 : extern "C" fn conn_blocking_write(
     170       269802 :     sk: *mut Safekeeper,
     171       269802 :     buf: *const ::std::os::raw::c_void,
     172       269802 :     size: usize,
     173       269802 : ) -> bool {
     174       269802 :     unsafe {
     175       269802 :         let buf = std::slice::from_raw_parts(buf as *const u8, size);
     176       269802 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     177       269802 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     178       269802 :         (*api).conn_blocking_write(&mut (*sk), buf)
     179       269802 :     }
     180       269802 : }
     181              : 
     182         5957 : extern "C" fn recovery_download(wp: *mut WalProposer, sk: *mut Safekeeper) -> bool {
     183         5957 :     unsafe {
     184         5957 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     185         5957 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     186         5957 : 
     187         5957 :         // currently `recovery_download` is always called right after election
     188         5957 :         (*api).after_election(&mut (*wp));
     189         5957 : 
     190         5957 :         (*api).recovery_download(&mut (*wp), &mut (*sk))
     191         5957 :     }
     192         5957 : }
     193              : 
     194         7960 : extern "C" fn wal_reader_allocate(sk: *mut Safekeeper) {
     195         7960 :     unsafe {
     196         7960 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     197         7960 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     198         7960 :         (*api).wal_reader_allocate(&mut (*sk));
     199         7960 :     }
     200         7960 : }
     201              : 
     202              : #[allow(clippy::unnecessary_cast)]
     203        27294 : extern "C" fn wal_read(
     204        27294 :     sk: *mut Safekeeper,
     205        27294 :     buf: *mut ::std::os::raw::c_char,
     206        27294 :     startptr: XLogRecPtr,
     207        27294 :     count: Size,
     208        27294 :     _errmsg: *mut *mut ::std::os::raw::c_char,
     209        27294 : ) -> NeonWALReadResult {
     210        27294 :     unsafe {
     211        27294 :         let buf = std::slice::from_raw_parts_mut(buf as *mut u8, count);
     212        27294 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     213        27294 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     214        27294 :         // TODO: errmsg is not forwarded
     215        27294 :         (*api).wal_read(&mut (*sk), buf, startptr)
     216        27294 :     }
     217        27294 : }
     218              : 
     219       115200 : extern "C" fn wal_reader_events(sk: *mut Safekeeper) -> uint32 {
     220       115200 :     unsafe {
     221       115200 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     222       115200 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     223       115200 :         (*api).wal_reader_events(&mut (*sk))
     224       115200 :     }
     225       115200 : }
     226              : 
     227        72207 : extern "C" fn init_event_set(wp: *mut WalProposer) {
     228        72207 :     unsafe {
     229        72207 :         let callback_data = (*(*wp).config).callback_data;
     230        72207 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     231        72207 :         (*api).init_event_set(&mut (*wp));
     232        72207 :     }
     233        72207 : }
     234              : 
     235       531393 : extern "C" fn update_event_set(sk: *mut Safekeeper, events: uint32) {
     236       531393 :     unsafe {
     237       531393 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     238       531393 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     239       531393 :         (*api).update_event_set(&mut (*sk), events);
     240       531393 :     }
     241       531393 : }
     242              : 
     243        40414 : extern "C" fn active_state_update_event_set(sk: *mut Safekeeper) {
     244        40414 :     unsafe {
     245        40414 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     246        40414 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     247        40414 :         (*api).active_state_update_event_set(&mut (*sk));
     248        40414 :     }
     249        40414 : }
     250              : 
     251       479564 : extern "C" fn add_safekeeper_event_set(sk: *mut Safekeeper, events: uint32) {
     252       479564 :     unsafe {
     253       479564 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     254       479564 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     255       479564 :         (*api).add_safekeeper_event_set(&mut (*sk), events);
     256       479564 :     }
     257       479564 : }
     258              : 
     259       302055 : extern "C" fn rm_safekeeper_event_set(sk: *mut Safekeeper) {
     260       302055 :     unsafe {
     261       302055 :         let callback_data = (*(*(*sk).wp).config).callback_data;
     262       302055 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     263       302055 :         (*api).rm_safekeeper_event_set(&mut (*sk));
     264       302055 :     }
     265       302055 : }
     266              : 
     267       700058 : extern "C" fn wait_event_set(
     268       700058 :     wp: *mut WalProposer,
     269       700058 :     timeout: ::std::os::raw::c_long,
     270       700058 :     event_sk: *mut *mut Safekeeper,
     271       700058 :     events: *mut uint32,
     272       700058 : ) -> ::std::os::raw::c_int {
     273       700058 :     unsafe {
     274       700058 :         let callback_data = (*(*wp).config).callback_data;
     275       700058 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     276       700058 :         let result = (*api).wait_event_set(&mut (*wp), timeout);
     277       700058 :         match result {
     278              :             WaitResult::Latch => {
     279        71269 :                 *event_sk = std::ptr::null_mut();
     280        71269 :                 *events = crate::bindings::WL_LATCH_SET;
     281        71269 :                 1
     282              :             }
     283              :             WaitResult::Timeout => {
     284        21809 :                 *event_sk = std::ptr::null_mut();
     285        21809 :                 // WaitEventSetWait returns 0 for timeout.
     286        21809 :                 *events = 0;
     287        21809 :                 0
     288              :             }
     289       606980 :             WaitResult::Network(sk, event_mask) => {
     290       606980 :                 *event_sk = sk;
     291       606980 :                 *events = event_mask;
     292       606980 :                 1
     293              :             }
     294              :         }
     295              :     }
     296       700058 : }
     297              : 
     298        72207 : extern "C" fn strong_random(
     299        72207 :     wp: *mut WalProposer,
     300        72207 :     buf: *mut ::std::os::raw::c_void,
     301        72207 :     len: usize,
     302        72207 : ) -> bool {
     303        72207 :     unsafe {
     304        72207 :         let buf = std::slice::from_raw_parts_mut(buf as *mut u8, len);
     305        72207 :         let callback_data = (*(*wp).config).callback_data;
     306        72207 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     307        72207 :         (*api).strong_random(buf)
     308        72207 :     }
     309        72207 : }
     310              : 
     311         2224 : extern "C" fn get_redo_start_lsn(wp: *mut WalProposer) -> XLogRecPtr {
     312         2224 :     unsafe {
     313         2224 :         let callback_data = (*(*wp).config).callback_data;
     314         2224 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     315         2224 :         (*api).get_redo_start_lsn()
     316         2224 :     }
     317         2224 : }
     318              : 
     319         2940 : extern "C" fn finish_sync_safekeepers(wp: *mut WalProposer, lsn: XLogRecPtr) {
     320         2940 :     unsafe {
     321         2940 :         let callback_data = (*(*wp).config).callback_data;
     322         2940 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     323         2940 :         (*api).finish_sync_safekeepers(lsn)
     324         2940 :     }
     325         2940 : }
     326              : 
     327        13652 : extern "C" fn process_safekeeper_feedback(wp: *mut WalProposer, commit_lsn: XLogRecPtr) {
     328        13652 :     unsafe {
     329        13652 :         let callback_data = (*(*wp).config).callback_data;
     330        13652 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     331        13652 :         (*api).process_safekeeper_feedback(&mut (*wp), commit_lsn)
     332        13652 :     }
     333        13652 : }
     334              : 
     335       790105 : extern "C" fn log_internal(
     336       790105 :     wp: *mut WalProposer,
     337       790105 :     level: ::std::os::raw::c_int,
     338       790105 :     line: *const ::std::os::raw::c_char,
     339       790105 : ) {
     340       790105 :     unsafe {
     341       790105 :         let callback_data = (*(*wp).config).callback_data;
     342       790105 :         let api = callback_data as *mut Box<dyn ApiImpl>;
     343       790105 :         let line = CStr::from_ptr(line);
     344       790105 :         let line = line.to_str().unwrap();
     345       790105 :         (*api).log_internal(&mut (*wp), Level::from(level as u32), line)
     346       790105 :     }
     347       790105 : }
     348              : 
     349      1579638 : #[derive(Debug, PartialEq)]
     350              : pub enum Level {
     351              :     Debug5,
     352              :     Debug4,
     353              :     Debug3,
     354              :     Debug2,
     355              :     Debug1,
     356              :     Log,
     357              :     Info,
     358              :     Notice,
     359              :     Warning,
     360              :     Error,
     361              :     Fatal,
     362              :     Panic,
     363              :     WPEvent,
     364              : }
     365              : 
     366              : impl Level {
     367       790105 :     pub fn from(elevel: u32) -> Level {
     368       790105 :         use crate::bindings::*;
     369       790105 : 
     370       790105 :         match elevel {
     371        27294 :             DEBUG5 => Level::Debug5,
     372            0 :             DEBUG4 => Level::Debug4,
     373            0 :             DEBUG3 => Level::Debug3,
     374        84990 :             DEBUG2 => Level::Debug2,
     375            0 :             DEBUG1 => Level::Debug1,
     376       588281 :             LOG => Level::Log,
     377            0 :             INFO => Level::Info,
     378            0 :             NOTICE => Level::Notice,
     379        88975 :             WARNING => Level::Warning,
     380            0 :             ERROR => Level::Error,
     381          544 :             FATAL => Level::Fatal,
     382           21 :             PANIC => Level::Panic,
     383            0 :             WPEVENT => Level::WPEvent,
     384            0 :             _ => panic!("unknown log level {}", elevel),
     385              :         }
     386       790105 :     }
     387              : }
     388              : 
     389        72207 : pub(crate) fn create_api() -> walproposer_api {
     390        72207 :     walproposer_api {
     391        72207 :         get_shmem_state: Some(get_shmem_state),
     392        72207 :         start_streaming: Some(start_streaming),
     393        72207 :         get_flush_rec_ptr: Some(get_flush_rec_ptr),
     394        72207 :         get_current_timestamp: Some(get_current_timestamp),
     395        72207 :         conn_error_message: Some(conn_error_message),
     396        72207 :         conn_status: Some(conn_status),
     397        72207 :         conn_connect_start: Some(conn_connect_start),
     398        72207 :         conn_connect_poll: Some(conn_connect_poll),
     399        72207 :         conn_send_query: Some(conn_send_query),
     400        72207 :         conn_get_query_result: Some(conn_get_query_result),
     401        72207 :         conn_flush: Some(conn_flush),
     402        72207 :         conn_finish: Some(conn_finish),
     403        72207 :         conn_async_read: Some(conn_async_read),
     404        72207 :         conn_async_write: Some(conn_async_write),
     405        72207 :         conn_blocking_write: Some(conn_blocking_write),
     406        72207 :         recovery_download: Some(recovery_download),
     407        72207 :         wal_reader_allocate: Some(wal_reader_allocate),
     408        72207 :         wal_read: Some(wal_read),
     409        72207 :         wal_reader_events: Some(wal_reader_events),
     410        72207 :         init_event_set: Some(init_event_set),
     411        72207 :         update_event_set: Some(update_event_set),
     412        72207 :         active_state_update_event_set: Some(active_state_update_event_set),
     413        72207 :         add_safekeeper_event_set: Some(add_safekeeper_event_set),
     414        72207 :         rm_safekeeper_event_set: Some(rm_safekeeper_event_set),
     415        72207 :         wait_event_set: Some(wait_event_set),
     416        72207 :         strong_random: Some(strong_random),
     417        72207 :         get_redo_start_lsn: Some(get_redo_start_lsn),
     418        72207 :         finish_sync_safekeepers: Some(finish_sync_safekeepers),
     419        72207 :         process_safekeeper_feedback: Some(process_safekeeper_feedback),
     420        72207 :         log_internal: Some(log_internal),
     421        72207 :     }
     422        72207 : }
     423              : 
     424              : impl std::fmt::Display for Level {
     425         3142 :     fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
     426         3142 :         write!(f, "{:?}", self)
     427         3142 :     }
     428              : }
     429              : 
     430              : /// Take ownership of `Vec<u8>` from StringInfoData.
     431              : #[allow(clippy::unnecessary_cast)]
     432       356019 : pub(crate) fn take_vec_u8(pg: &mut StringInfoData) -> Option<Vec<u8>> {
     433       356019 :     if pg.data.is_null() {
     434       216615 :         return None;
     435       139404 :     }
     436       139404 : 
     437       139404 :     let ptr = pg.data as *mut u8;
     438       139404 :     let length = pg.len as usize;
     439       139404 :     let capacity = pg.maxlen as usize;
     440       139404 : 
     441       139404 :     pg.data = std::ptr::null_mut();
     442       139404 :     pg.len = 0;
     443       139404 :     pg.maxlen = 0;
     444       139404 : 
     445       139404 :     unsafe { Some(Vec::from_raw_parts(ptr, length, capacity)) }
     446       356019 : }
     447              : 
     448              : /// Store `Vec<u8>` in StringInfoData.
     449       139414 : fn store_vec_u8(pg: &mut StringInfoData, vec: Vec<u8>) -> *mut ::std::os::raw::c_char {
     450       139414 :     let ptr = vec.as_ptr() as *mut ::std::os::raw::c_char;
     451       139414 :     let length = vec.len();
     452       139414 :     let capacity = vec.capacity();
     453       139414 : 
     454       139414 :     assert!(pg.data.is_null());
     455              : 
     456       139414 :     pg.data = ptr;
     457       139414 :     pg.len = length as i32;
     458       139414 :     pg.maxlen = capacity as i32;
     459       139414 : 
     460       139414 :     std::mem::forget(vec);
     461       139414 : 
     462       139414 :     ptr
     463       139414 : }
        

Generated by: LCOV version 2.1-beta