Line data Source code
1 : use std::{ffi::CStr, sync::Arc};
2 :
3 : use parking_lot::{Mutex, MutexGuard};
4 : use postgres_ffi::v16::wal_generator::{LogicalMessageGenerator, WalGenerator};
5 : use utils::lsn::Lsn;
6 :
7 : use super::block_storage::BlockStorage;
8 :
9 : /// Simulation implementation of walproposer WAL storage.
10 : pub struct DiskWalProposer {
11 : state: Mutex<State>,
12 : }
13 :
14 : impl DiskWalProposer {
15 9606 : pub fn new() -> Arc<DiskWalProposer> {
16 9606 : Arc::new(DiskWalProposer {
17 9606 : state: Mutex::new(State {
18 9606 : internal_available_lsn: Lsn(0),
19 9606 : prev_lsn: Lsn(0),
20 9606 : disk: BlockStorage::new(),
21 9606 : wal_generator: WalGenerator::new(LogicalMessageGenerator::new(c"", &[])),
22 9606 : }),
23 9606 : })
24 9606 : }
25 :
26 15282 : pub fn lock(&self) -> MutexGuard<State> {
27 15282 : self.state.lock()
28 15282 : }
29 : }
30 :
31 : pub struct State {
32 : // flush_lsn
33 : internal_available_lsn: Lsn,
34 : // needed for WAL generation
35 : prev_lsn: Lsn,
36 : // actual WAL storage
37 : disk: BlockStorage,
38 : // WAL record generator
39 : wal_generator: WalGenerator<LogicalMessageGenerator>,
40 : }
41 :
42 : impl State {
43 1051 : pub fn read(&self, pos: u64, buf: &mut [u8]) {
44 1051 : self.disk.read(pos, buf);
45 1051 : // TODO: fail on reading uninitialized data
46 1051 : }
47 :
48 179 : pub fn write(&mut self, pos: u64, buf: &[u8]) {
49 179 : self.disk.write(pos, buf);
50 179 : }
51 :
52 : /// Update the internal available LSN to the given value.
53 358 : pub fn reset_to(&mut self, lsn: Lsn) {
54 358 : self.internal_available_lsn = lsn;
55 358 : self.prev_lsn = Lsn(0); // Safekeeper doesn't care if this is omitted
56 358 : self.wal_generator.lsn = self.internal_available_lsn;
57 358 : self.wal_generator.prev_lsn = self.prev_lsn;
58 358 : }
59 :
60 : /// Get current LSN.
61 1968 : pub fn flush_rec_ptr(&self) -> Lsn {
62 1968 : self.internal_available_lsn
63 1968 : }
64 :
65 : /// Inserts a logical record in the WAL at the current LSN.
66 11726 : pub fn insert_logical_message(&mut self, prefix: &CStr, msg: &[u8]) {
67 11726 : let (_, record) = self.wal_generator.append_logical_message(prefix, msg);
68 11726 : self.disk.write(self.internal_available_lsn.into(), &record);
69 11726 : self.prev_lsn = self.internal_available_lsn;
70 11726 : self.internal_available_lsn += record.len() as u64;
71 11726 : }
72 : }
|