Wire up the rsync code in the send kitten
This commit is contained in:
parent
1523424458
commit
063c39ea12
3 changed files with 72 additions and 20 deletions
|
|
@ -22,6 +22,7 @@ import (
|
|||
|
||||
"kitty"
|
||||
"kitty/tools/cli/markup"
|
||||
"kitty/tools/rsync"
|
||||
"kitty/tools/tui"
|
||||
"kitty/tools/tui/loop"
|
||||
"kitty/tools/utils"
|
||||
|
|
@ -101,6 +102,9 @@ type File struct {
|
|||
actual_file *os.File
|
||||
transmitted_bytes, reported_progress int64
|
||||
transmit_started_at, transmit_ended_at, done_at time.Time
|
||||
differ *rsync.Differ
|
||||
delta_loader func() error
|
||||
deltabuf *bytes.Buffer
|
||||
}
|
||||
|
||||
func get_remote_path(local_path string, remote_base string) string {
|
||||
|
|
@ -729,7 +733,7 @@ func (self *SendManager) on_file_status_update(ftc *FileTransmissionCommand) err
|
|||
file.state = TRANSMITTING
|
||||
}
|
||||
if file.state == WAITING_FOR_DATA {
|
||||
panic("TODO: Implement rsync support")
|
||||
file.differ = rsync.NewDiffer()
|
||||
}
|
||||
self.update_collective_statuses()
|
||||
}
|
||||
|
|
@ -764,12 +768,38 @@ func (self *SendManager) on_file_status_update(ftc *FileTransmissionCommand) err
|
|||
return nil
|
||||
}
|
||||
|
||||
func (self *File) start_delta_calculation() (err error) {
|
||||
self.state = TRANSMITTING
|
||||
if self.actual_file == nil {
|
||||
self.actual_file, err = os.Open(self.expanded_local_path)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
self.deltabuf = bytes.NewBuffer(make([]byte, 0, 32+rsync.DataSizeMultiple*self.differ.BlockSize()))
|
||||
self.delta_loader = self.differ.CreateDelta(self.actual_file, self.deltabuf)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SendManager) on_signature_data_received(ftc *FileTransmissionCommand) error {
|
||||
file := self.fid_map[ftc.File_id]
|
||||
if file == nil || file.state != WAITING_FOR_DATA {
|
||||
return nil
|
||||
}
|
||||
panic("TODO: Implement rsync support")
|
||||
if file.differ == nil {
|
||||
file.differ = rsync.NewDiffer()
|
||||
}
|
||||
if err := file.differ.AddSignatureData(ftc.Data); err != nil {
|
||||
return err
|
||||
}
|
||||
self.progress_tracker.signature_bytes += len(ftc.Data)
|
||||
if ftc.Action == Action_end_data {
|
||||
if err := file.differ.FinishSignatureData(); err != nil {
|
||||
return err
|
||||
}
|
||||
return file.start_delta_calculation()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SendManager) on_file_transfer_response(ftc *FileTransmissionCommand) error {
|
||||
|
|
@ -918,24 +948,38 @@ func (self *File) next_chunk() (ans string, asz int, err error) {
|
|||
return
|
||||
}
|
||||
is_last := false
|
||||
// TODO: self.delta_loader rsync support
|
||||
if self.actual_file == nil {
|
||||
self.actual_file, err = os.Open(self.expanded_local_path)
|
||||
if err != nil {
|
||||
var chunk []byte
|
||||
if self.delta_loader != nil {
|
||||
self.deltabuf.Reset()
|
||||
if err = self.delta_loader(); err != nil {
|
||||
if err == io.EOF {
|
||||
is_last = true
|
||||
} else {
|
||||
return
|
||||
}
|
||||
}
|
||||
chunk = make([]byte, len(self.deltabuf.Bytes()))
|
||||
copy(chunk, self.deltabuf.Bytes())
|
||||
} else {
|
||||
if self.actual_file == nil {
|
||||
self.actual_file, err = os.Open(self.expanded_local_path)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
chunk = make([]byte, sz)
|
||||
var n int
|
||||
n, err = self.actual_file.Read(chunk)
|
||||
if err != nil && !errors.Is(err, io.EOF) {
|
||||
return
|
||||
}
|
||||
if n <= 0 {
|
||||
is_last = true
|
||||
} else if pos, _ := self.actual_file.Seek(0, os.SEEK_CUR); pos >= self.file_size {
|
||||
is_last = true
|
||||
}
|
||||
chunk = chunk[:n]
|
||||
}
|
||||
chunk := make([]byte, sz)
|
||||
n, err := self.actual_file.Read(chunk)
|
||||
if err != nil && !errors.Is(err, io.EOF) {
|
||||
return
|
||||
}
|
||||
if n <= 0 {
|
||||
is_last = true
|
||||
} else if pos, _ := self.actual_file.Seek(0, os.SEEK_CUR); pos >= self.file_size {
|
||||
is_last = true
|
||||
}
|
||||
chunk = chunk[:n]
|
||||
uncompressed_sz := len(chunk)
|
||||
cchunk := self.compressor.Compress(chunk)
|
||||
if is_last {
|
||||
|
|
@ -951,6 +995,8 @@ func (self *File) next_chunk() (ans string, asz int, err error) {
|
|||
return
|
||||
}
|
||||
}
|
||||
self.delta_loader = nil
|
||||
self.deltabuf = nil
|
||||
}
|
||||
ans, asz = utils.UnsafeBytesToString(cchunk), uncompressed_sz
|
||||
return
|
||||
|
|
|
|||
|
|
@ -608,9 +608,11 @@ func (r *rsync) CreateDelta(source io.Reader, signature []BlockHash) ([]Operatio
|
|||
}
|
||||
}
|
||||
|
||||
const DataSizeMultiple int = 8
|
||||
|
||||
func (r *rsync) CreateDiff(source io.Reader, signature []BlockHash, output io.Writer) func() error {
|
||||
ans := &diff{
|
||||
block_size: r.BlockSize, buffer: make([]byte, 0, (r.BlockSize * 8)),
|
||||
block_size: r.BlockSize, buffer: make([]byte, 0, (r.BlockSize * DataSizeMultiple)),
|
||||
hash_lookup: make(map[uint32][]BlockHash, len(signature)),
|
||||
source: source, hasher: r.hasher_constructor(),
|
||||
checksummer: r.checksummer_constructor(), output: output,
|
||||
|
|
|
|||
|
|
@ -119,7 +119,7 @@ func (self *Api) read_signature_blocks(data []byte) (consumed int) {
|
|||
return
|
||||
}
|
||||
|
||||
func (self *Differ) finish_signature_data() (err error) {
|
||||
func (self *Differ) FinishSignatureData() (err error) {
|
||||
if len(self.unconsumed_signature_data) > 0 {
|
||||
return fmt.Errorf("There were %d leftover bytes in the signature data", len(self.unconsumed_signature_data))
|
||||
}
|
||||
|
|
@ -225,7 +225,7 @@ func (self *Patcher) CreateSignatureIterator(src io.Reader, output io.Writer) fu
|
|||
|
||||
// Create a serialized delta based on the previously loaded signature
|
||||
func (self *Differ) CreateDelta(src io.Reader, output io.Writer) func() error {
|
||||
if err := self.finish_signature_data(); err != nil {
|
||||
if err := self.FinishSignatureData(); err != nil {
|
||||
return func() error { return err }
|
||||
}
|
||||
if self.signature == nil {
|
||||
|
|
@ -236,6 +236,10 @@ func (self *Differ) CreateDelta(src io.Reader, output io.Writer) func() error {
|
|||
return self.rsync.CreateDiff(src, self.signature, output)
|
||||
}
|
||||
|
||||
func (self *Differ) BlockSize() int {
|
||||
return self.rsync.BlockSize
|
||||
}
|
||||
|
||||
// Add more external signature data
|
||||
func (self *Differ) AddSignatureData(data []byte) (err error) {
|
||||
self.unconsumed_signature_data = append(self.unconsumed_signature_data, data...)
|
||||
|
|
|
|||
Loading…
Reference in a new issue