1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
pub use crate::buffer::{SliceBuffer, SliceInfo};
pub use crate::read::StreamReader;
pub use crate::store::StreamStorage;
pub use crate::stream::{DocId, Head, PeerId, SignedHead, Slice, Stream, StreamId};
pub use crate::write::StreamWriter;
pub use bao::Hash;
pub use ed25519_dalek::{Keypair, PublicKey, SecretKey};

mod buffer;
mod read;
mod store;
mod stream;
mod write;

#[cfg(test)]
mod tests {
    use super::*;
    use anyhow::Result;
    use bao::decode::SliceDecoder;
    use rand::RngCore;
    use std::io::{BufReader, Read, Write};
    use tempdir::TempDir;

    fn rand_bytes(size: usize) -> Vec<u8> {
        let mut rng = rand::thread_rng();
        let mut data = Vec::with_capacity(size);
        data.resize(data.capacity(), 0);
        rng.fill_bytes(&mut data);
        data
    }

    fn keypair(secret: [u8; 32]) -> Keypair {
        let secret = SecretKey::from_bytes(&secret).unwrap();
        let public = PublicKey::from(&secret);
        Keypair { secret, public }
    }

    #[test]
    fn test_append_stream() -> Result<()> {
        let tmp = TempDir::new("test_append_stream")?;
        let mut storage = StreamStorage::open(tmp.path(), keypair([0; 32]))?;
        let data = rand_bytes(1_000_000);

        let doc = DocId::unique();
        let mut stream = storage.append(doc)?;
        stream.write_all(&data)?;
        stream.commit()?;

        let stream = storage.slice(stream.id(), 0, data.len() as u64)?;
        let mut reader = BufReader::new(stream);
        let mut data2 = Vec::with_capacity(data.len());
        data2.resize(data2.capacity(), 0);
        reader.read_exact(&mut data2)?;
        assert_eq!(data, data2);

        Ok(())
    }

    #[test]
    fn test_extract_slice() -> Result<()> {
        let tmp = TempDir::new("test_extract_slice")?;
        let mut storage = StreamStorage::open(tmp.path(), keypair([0; 32]))?;
        let data = rand_bytes(1027);

        let doc = DocId::unique();
        let mut stream = storage.append(doc)?;
        stream.write_all(&data)?;
        stream.commit()?;

        let offset = 8;
        let len = 32;
        let slice = data[offset..(offset + len)].to_vec();

        let mut stream = storage.slice(stream.id(), offset as u64, len as u64)?;
        let mut slice2 = vec![];
        stream.read_to_end(&mut slice2)?;
        assert_eq!(slice2, slice);

        let mut vslice = Slice::default();
        storage.extract(stream.id(), offset as u64, len as u64, &mut vslice)?;

        let mut slice2 = vec![];
        vslice.head.verify(stream.id())?;
        let mut decoder = SliceDecoder::new(
            &vslice.data[..],
            &Hash::from(vslice.head.head.hash),
            offset as u64,
            len as u64,
        );
        decoder.read_to_end(&mut slice2)?;
        assert_eq!(slice2, slice);
        Ok(())
    }

    #[test]
    fn test_sync() -> Result<()> {
        let tmp = TempDir::new("test_sync_1")?;
        let mut storage = StreamStorage::open(tmp.path(), keypair([0; 32]))?;
        let data = rand_bytes(8192);

        let doc = DocId::unique();
        let mut stream = storage.append(doc)?;
        stream.write_all(&data[..4096])?;
        let head1 = stream.commit()?;
        stream.write_all(&data[4096..])?;
        let head2 = stream.commit()?;

        let tmp = TempDir::new("test_sync_2")?;
        let mut storage2 = StreamStorage::open(tmp.path(), keypair([1; 32]))?;
        let stream = storage2.subscribe(stream.id())?;
        let mut stream = SliceBuffer::new(stream, 1024);

        let mut slice = Slice::default();
        for head in [head1, head2].iter() {
            head.verify(stream.id())?;
            stream.prepare(head.head().len() - stream.head().head().len);
            for i in 0..stream.slices().len() {
                let info = &stream.slices()[i];
                storage.extract(stream.id(), info.offset, info.len, &mut slice)?;
                stream.add_slice(&slice, i)?;
            }
            stream.commit(*head.sig())?;
        }

        let mut stream = storage2.slice(stream.id(), 0, 8192)?;
        let mut data2 = vec![];
        stream.read_to_end(&mut data2)?;
        assert_eq!(data, data2);

        let tmp = TempDir::new("test_sync_3")?;
        let mut storage = StreamStorage::open(tmp.path(), keypair([1; 32]))?;
        let stream = storage.subscribe(stream.id())?;
        let mut stream = SliceBuffer::new(stream, 1024);

        let mut slice = Slice::default();
        for head in [head1, head2].iter() {
            head.verify(stream.id()).unwrap();
            stream.prepare(head.head().len() - stream.head().head().len);
            for i in 0..stream.slices().len() {
                let info = &stream.slices()[i];
                storage2.extract(stream.id(), info.offset, info.len, &mut slice)?;
                stream.add_slice(&slice, i)?;
            }
            stream.commit(*head.sig())?;
        }

        let mut stream = storage.slice(stream.id(), 0, 8192)?;
        let mut data2 = vec![];
        stream.read_to_end(&mut data2)?;
        assert_eq!(data, data2);

        Ok(())
    }
}