feat: Implement full network functionality with HTTP and P2P modes
Go Build and Test / build-and-test (pull_request) Successful in 37s
Go Build and Test / build-and-test (pull_request) Successful in 37s
This major update introduces a robust, multi-faceted networking capability, allowing backups to be managed through a central server or via direct peer-to-peer connections. - **Add `serve` command**: Implemented a new `serve` command that runs a persistent HTTP server. This allows for a centralized backup management workflow where multiple clients can upload and download archives. - **HTTP Client Mode**: The `create` and `restore` commands can now function as HTTP clients, targeting a `serve` instance using the `ip:port/filename` address format. - **Restore P2P Mode**: The original peer-to-peer streaming functionality (where `restore` acts as a TCP server and `create` as a client) has been fully restored and integrated into the new, more robust command architecture. - **Accurate Network Progress**: Fixed a critical bug where network progress bars would not display correctly. The `serve` command now sends a `Content-Length` header, and `restore` is refactored to handle asynchronous state updates correctly. - **Smart Address Parsing**: The application now automatically distinguishes between HTTP mode (`ip:port/filename`) and P2P mode (`ip:port`), providing a seamless user experience. - **Update `README.md`**: The project's README has been completely rewritten to document all three operating modes (Local, HTTP Server, P2P) with clear usage examples. - **Add Gitea Actions Workflow**: A new CI pipeline (`.gitea/workflows/build-and-test.yml`) has been added to automatically build and test the application on every push and pull request.
This commit is contained in:
+93
-84
@@ -7,9 +7,11 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
"unsafe"
|
||||
@@ -46,31 +48,30 @@ var createCmd = &cobra.Command{
|
||||
}
|
||||
|
||||
type model struct {
|
||||
progress progress.Model
|
||||
source string
|
||||
target string
|
||||
totalBytes int64
|
||||
processed int64
|
||||
done bool
|
||||
err error
|
||||
startTime time.Time
|
||||
progressCh chan int64
|
||||
isDirectory bool
|
||||
isNetworkTarget bool
|
||||
progress progress.Model
|
||||
source string
|
||||
target string
|
||||
totalBytes int64
|
||||
processed int64
|
||||
done bool
|
||||
err error
|
||||
startTime time.Time
|
||||
progressCh chan int64
|
||||
isDirectory bool
|
||||
httpTargetInfo *ServerTargetInfo
|
||||
isP2PTarget bool // Флаг для старого P2P-режима
|
||||
}
|
||||
|
||||
type progressMsg int64
|
||||
type doneMsg struct{}
|
||||
type errorMsg struct{ err error }
|
||||
|
||||
// getPathSize возвращает размер файла, диска или директории
|
||||
func getPathSize(path string) (int64, error) {
|
||||
fileInfo, err := os.Stat(path)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
// Для Linux: проверка, является ли источник блочным устройством
|
||||
if runtime.GOOS == "linux" {
|
||||
stat, ok := fileInfo.Sys().(*syscall.Stat_t)
|
||||
if ok && (stat.Mode&syscall.S_IFMT) == syscall.S_IFBLK {
|
||||
@@ -94,7 +95,6 @@ func getPathSize(path string) (int64, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// Для директорий
|
||||
if fileInfo.IsDir() {
|
||||
var total int64
|
||||
err := filepath.Walk(path, func(_ string, info os.FileInfo, err error) error {
|
||||
@@ -109,7 +109,6 @@ func getPathSize(path string) (int64, error) {
|
||||
return total, err
|
||||
}
|
||||
|
||||
// Для обычных файлов
|
||||
return fileInfo.Size(), nil
|
||||
}
|
||||
|
||||
@@ -129,15 +128,30 @@ func initialModel(src, dst string) (*model, error) {
|
||||
return nil, err
|
||||
}
|
||||
isDirectory := srcInfo.IsDir()
|
||||
networkTarget := isNetworkAddress(dst)
|
||||
|
||||
targetInfo, isHttp := parseServerTarget(dst)
|
||||
isP2P := false
|
||||
targetPath := dst
|
||||
|
||||
if !networkTarget {
|
||||
if isHttp {
|
||||
// Новый HTTP-режим
|
||||
if isDirectory && !strings.HasSuffix(targetInfo.Filename, ".tar.gz") {
|
||||
targetInfo.Filename += ".tar.gz"
|
||||
} else if !isDirectory && !strings.HasSuffix(targetInfo.Filename, ".gz") {
|
||||
targetInfo.Filename += ".gz"
|
||||
}
|
||||
targetInfo.URL = fmt.Sprintf("http://%s/backup/%s", targetInfo.Address, targetInfo.Filename)
|
||||
targetPath = targetInfo.URL
|
||||
} else if isNetworkAddress(dst) {
|
||||
// Старый P2P-режим
|
||||
isP2P = true
|
||||
} else {
|
||||
// Локальный файл
|
||||
if isDirectory {
|
||||
if filepath.Ext(targetPath) != ".tar.gz" {
|
||||
targetPath += ".tar.gz"
|
||||
}
|
||||
} else { // Это файл или диск
|
||||
} else {
|
||||
if filepath.Ext(targetPath) != ".gz" {
|
||||
targetPath += ".gz"
|
||||
}
|
||||
@@ -145,14 +159,15 @@ func initialModel(src, dst string) (*model, error) {
|
||||
}
|
||||
|
||||
return &model{
|
||||
progress: p,
|
||||
source: src,
|
||||
target: targetPath,
|
||||
totalBytes: totalBytes,
|
||||
startTime: time.Now(),
|
||||
progressCh: make(chan int64, 100),
|
||||
isDirectory: isDirectory,
|
||||
isNetworkTarget: networkTarget,
|
||||
progress: p,
|
||||
source: src,
|
||||
target: targetPath,
|
||||
totalBytes: totalBytes,
|
||||
startTime: time.Now(),
|
||||
progressCh: make(chan int64, 100),
|
||||
isDirectory: isDirectory,
|
||||
httpTargetInfo: targetInfo,
|
||||
isP2PTarget: isP2P,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -163,21 +178,47 @@ func (m *model) Init() tea.Cmd {
|
||||
)
|
||||
}
|
||||
|
||||
// getTargetWriter создает io.WriteCloser для файла или сетевого подключения
|
||||
func (m *model) getTargetWriter() (io.WriteCloser, error) {
|
||||
if m.isNetworkTarget {
|
||||
if m.httpTargetInfo != nil {
|
||||
// Новый HTTP-режим
|
||||
pipeReader, pipeWriter := io.Pipe()
|
||||
req, err := http.NewRequest("POST", m.httpTargetInfo.URL, pipeReader)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create http request: %w", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/octet-stream")
|
||||
req.ContentLength = -1 // Stream upload
|
||||
|
||||
go func() {
|
||||
client := &http.Client{} // No timeout for uploads
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
pipeWriter.CloseWithError(fmt.Errorf("http request failed: %w", err))
|
||||
return
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
bodyBytes, _ := io.ReadAll(resp.Body)
|
||||
err := fmt.Errorf("server returned non-200 status: %s\n%s", resp.Status, string(bodyBytes))
|
||||
pipeWriter.CloseWithError(err)
|
||||
}
|
||||
}()
|
||||
|
||||
return pipeWriter, nil
|
||||
|
||||
} else if m.isP2PTarget {
|
||||
// Старый P2P-режим
|
||||
conn, err := net.DialTimeout("tcp", m.target, 10*time.Second)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to connect to %s: %w", m.target, err)
|
||||
}
|
||||
|
||||
// Протокол: сначала отправляем 8 байт (размер) и 1 байт (тип)
|
||||
// 1. Общий размер (int64)
|
||||
// Отправляем бинарный заголовок (размер + тип)
|
||||
if err := binary.Write(conn, binary.BigEndian, m.totalBytes); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("failed to send backup size: %w", err)
|
||||
}
|
||||
// 2. Тип бэкапа (1 = директория, 0 = файл/диск)
|
||||
var typeByte byte = 0
|
||||
if m.isDirectory {
|
||||
typeByte = 1
|
||||
@@ -186,10 +227,10 @@ func (m *model) getTargetWriter() (io.WriteCloser, error) {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("failed to send backup type: %w", err)
|
||||
}
|
||||
|
||||
return conn, nil
|
||||
}
|
||||
// Логика для локального файла
|
||||
|
||||
// Локальный файл
|
||||
return os.Create(m.target)
|
||||
}
|
||||
|
||||
@@ -216,7 +257,7 @@ func (m *model) startBackup() tea.Msg {
|
||||
return
|
||||
}
|
||||
|
||||
m.progressCh <- -2 // Сигнал завершения
|
||||
m.progressCh <- -2
|
||||
}()
|
||||
|
||||
return nil
|
||||
@@ -229,30 +270,17 @@ func (m *model) backupFileOrDisk(w io.Writer) error {
|
||||
}
|
||||
defer srcFile.Close()
|
||||
|
||||
gzipWriter := gzip.NewWriter(w)
|
||||
progressWriter := &progressTracker{Writer: w, progressCh: m.progressCh}
|
||||
gzipWriter := gzip.NewWriter(progressWriter)
|
||||
defer gzipWriter.Close()
|
||||
|
||||
buf := make([]byte, 32*1024) // 32KB buffer
|
||||
for {
|
||||
n, err := srcFile.Read(buf)
|
||||
if n > 0 {
|
||||
if _, writeErr := gzipWriter.Write(buf[:n]); writeErr != nil {
|
||||
return writeErr
|
||||
}
|
||||
m.progressCh <- int64(n)
|
||||
}
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
_, err = io.Copy(gzipWriter, srcFile)
|
||||
return err
|
||||
}
|
||||
|
||||
func (m *model) backupDirectory(w io.Writer) error {
|
||||
gzipWriter := gzip.NewWriter(w)
|
||||
progressWriter := &progressTracker{Writer: w, progressCh: m.progressCh}
|
||||
gzipWriter := gzip.NewWriter(progressWriter)
|
||||
defer gzipWriter.Close()
|
||||
|
||||
tarWriter := tar.NewWriter(gzipWriter)
|
||||
@@ -288,39 +316,22 @@ func (m *model) backupDirectory(w io.Writer) error {
|
||||
}
|
||||
defer srcFile.Close()
|
||||
|
||||
buf := make([]byte, 32*1024)
|
||||
for {
|
||||
n, err := srcFile.Read(buf)
|
||||
if n > 0 {
|
||||
if _, err := tarWriter.Write(buf[:n]); err != nil {
|
||||
return err
|
||||
}
|
||||
m.progressCh <- int64(n)
|
||||
}
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
_, err = io.Copy(tarWriter, srcFile)
|
||||
return err
|
||||
})
|
||||
}
|
||||
|
||||
func (m *model) progressListener() tea.Msg {
|
||||
select {
|
||||
case n := <-m.progressCh:
|
||||
switch {
|
||||
case n == -1: // Ошибка
|
||||
n := <-m.progressCh
|
||||
if n < 0 {
|
||||
if n == -1 {
|
||||
return errorMsg{m.err}
|
||||
case n == -2: // Завершение
|
||||
return doneMsg{}
|
||||
default: // Прогресс
|
||||
m.processed += n
|
||||
return progressMsg(m.processed)
|
||||
}
|
||||
return doneMsg{}
|
||||
}
|
||||
|
||||
m.processed += n
|
||||
return progressMsg(m.processed)
|
||||
}
|
||||
|
||||
func (m *model) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
||||
@@ -352,7 +363,6 @@ func (m *model) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
||||
m.err = msg.err
|
||||
return m, tea.Quit
|
||||
}
|
||||
|
||||
return m, nil
|
||||
}
|
||||
|
||||
@@ -363,12 +373,11 @@ func (m *model) View() string {
|
||||
|
||||
if m.done {
|
||||
duration := time.Since(m.startTime)
|
||||
// ratio := float64(m.processed) / float64(m.totalBytes) * 100
|
||||
backupType := "File/disk"
|
||||
if m.isDirectory {
|
||||
backupType = "Directory"
|
||||
}
|
||||
if m.isNetworkTarget {
|
||||
if m.httpTargetInfo != nil || m.isP2PTarget {
|
||||
backupType += " network"
|
||||
}
|
||||
|
||||
@@ -386,7 +395,7 @@ func (m *model) View() string {
|
||||
if m.isDirectory {
|
||||
operation = "Archiving"
|
||||
}
|
||||
if m.isNetworkTarget {
|
||||
if m.httpTargetInfo != nil || m.isP2PTarget {
|
||||
operation = "Sending"
|
||||
}
|
||||
|
||||
@@ -408,7 +417,7 @@ func (m *model) View() string {
|
||||
func init() {
|
||||
rootCmd.AddCommand(createCmd)
|
||||
createCmd.Flags().StringVarP(&source, "source", "s", "", "Source file, directory or disk (e.g., /dev/sda) (required)")
|
||||
createCmd.Flags().StringVarP(&target, "target", "t", "", "Target backup file or network address (e.g., /path/to/backup.tar.gz or 127.0.0.1:8080) (required)")
|
||||
createCmd.Flags().StringVarP(&target, "target", "t", "", "Target backup file or network address (e.g., /path/to/backup.tar.gz, 127.0.0.1:8080/backup.gz, or 127.0.0.1:8080) (required)")
|
||||
createCmd.MarkFlagRequired("source")
|
||||
createCmd.MarkFlagRequired("target")
|
||||
}
|
||||
|
||||
+147
-111
@@ -7,6 +7,7 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
@@ -28,7 +29,7 @@ var (
|
||||
var restoreCmd = &cobra.Command{
|
||||
Use: "restore",
|
||||
Short: "Restore a backup",
|
||||
Long: "Restore a backup from a local archive or network stream to a file, disk, or directory.",
|
||||
Long: "Restore a backup from a local archive, P2P stream, or HTTP server.",
|
||||
Run: func(cmd *cobra.Command, args []string) {
|
||||
model, err := initialRestoreModel(restoreSource, restoreTarget)
|
||||
if err != nil {
|
||||
@@ -45,18 +46,27 @@ var restoreCmd = &cobra.Command{
|
||||
}
|
||||
|
||||
type restoreModel struct {
|
||||
progress progress.Model
|
||||
source string
|
||||
target string
|
||||
totalBytes int64
|
||||
processed int64
|
||||
done bool
|
||||
err error
|
||||
startTime time.Time
|
||||
progressCh chan int64
|
||||
isDirectory bool
|
||||
isNetworkSource bool
|
||||
statusMessage string
|
||||
progress progress.Model
|
||||
source string
|
||||
target string
|
||||
totalBytes int64
|
||||
processed int64
|
||||
done bool
|
||||
err error
|
||||
startTime time.Time
|
||||
progressCh chan int64
|
||||
isDirectory bool
|
||||
httpSourceInfo *ServerTargetInfo
|
||||
isP2PSource bool
|
||||
statusMessage string
|
||||
}
|
||||
|
||||
// restoreSourceReadyMsg is sent when the source is ready to be read.
|
||||
// It carries the reader, the total size, and whether it's a directory backup.
|
||||
type restoreSourceReadyMsg struct {
|
||||
reader io.ReadCloser
|
||||
totalBytes int64
|
||||
isDir bool
|
||||
}
|
||||
|
||||
type restoreProgressMsg int64
|
||||
@@ -64,21 +74,24 @@ type restoreDoneMsg struct{}
|
||||
type restoreErrorMsg struct{ err error }
|
||||
|
||||
func initialRestoreModel(src, dst string) (*restoreModel, error) {
|
||||
isNetwork := isNetworkAddress(src)
|
||||
p := progress.New(
|
||||
progress.WithDefaultGradient(),
|
||||
progress.WithWidth(40),
|
||||
)
|
||||
|
||||
var totalBytes int64 = 1 // Placeholder for network source, to avoid division by zero.
|
||||
var isDir bool
|
||||
var totalBytes int64 = 1 // Placeholder, will be updated.
|
||||
var isDir, isP2P bool
|
||||
targetPath := dst
|
||||
status := "Initializing..."
|
||||
|
||||
if isNetwork {
|
||||
if targetPath == "" {
|
||||
return nil, fmt.Errorf("target flag -t is required for a network source")
|
||||
}
|
||||
httpInfo, isHttp := parseServerTarget(src)
|
||||
|
||||
if isHttp {
|
||||
isDir = strings.HasSuffix(httpInfo.Filename, ".tar.gz")
|
||||
status = fmt.Sprintf("Connecting to %s...", httpInfo.Address)
|
||||
src = httpInfo.URL
|
||||
} else if isNetworkAddress(src) {
|
||||
isP2P = true
|
||||
status = fmt.Sprintf("Listening on %s...", src)
|
||||
} else {
|
||||
fileInfo, err := os.Stat(src)
|
||||
@@ -87,89 +100,96 @@ func initialRestoreModel(src, dst string) (*restoreModel, error) {
|
||||
}
|
||||
totalBytes = fileInfo.Size()
|
||||
isDir = strings.HasSuffix(src, ".tar.gz")
|
||||
}
|
||||
|
||||
if isDir && targetPath == "" {
|
||||
targetPath, err = os.Getwd()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to get current directory: %w", err)
|
||||
}
|
||||
} else if !isDir && targetPath == "" {
|
||||
return nil, fmt.Errorf("target flag -t is required for file or disk restoration")
|
||||
if isDir && targetPath == "" {
|
||||
// For directory restores, we can default to the current directory.
|
||||
// We do this check here, but also again after getting the P2P type.
|
||||
var err error
|
||||
targetPath, err = os.Getwd()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to get current directory: %w", err)
|
||||
}
|
||||
} else if !isDir && targetPath == "" && !isP2P {
|
||||
// Target is mandatory for file/disk restores (except P2P, where we can't know the type yet).
|
||||
return nil, fmt.Errorf("target flag -t is required for this restoration")
|
||||
}
|
||||
|
||||
return &restoreModel{
|
||||
progress: p,
|
||||
source: src,
|
||||
target: targetPath,
|
||||
totalBytes: totalBytes,
|
||||
startTime: time.Now(),
|
||||
progressCh: make(chan int64, 100),
|
||||
isDirectory: isDir,
|
||||
isNetworkSource: isNetwork,
|
||||
statusMessage: status,
|
||||
progress: p,
|
||||
source: src,
|
||||
target: targetPath,
|
||||
totalBytes: totalBytes,
|
||||
startTime: time.Now(),
|
||||
progressCh: make(chan int64, 100),
|
||||
isDirectory: isDir,
|
||||
httpSourceInfo: httpInfo,
|
||||
isP2PSource: isP2P,
|
||||
statusMessage: status,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (m *restoreModel) Init() tea.Cmd {
|
||||
return tea.Batch(
|
||||
m.startRestore,
|
||||
m.restoreProgressListener,
|
||||
)
|
||||
return m.fetchSourceCmd
|
||||
}
|
||||
|
||||
func (m *restoreModel) getSourceReader() (io.ReadCloser, error) {
|
||||
if m.isNetworkSource {
|
||||
func (m *restoreModel) fetchSourceCmd() tea.Msg {
|
||||
if m.httpSourceInfo != nil {
|
||||
// New HTTP Mode
|
||||
client := &http.Client{}
|
||||
resp, err := client.Get(m.source)
|
||||
if err != nil {
|
||||
return restoreErrorMsg{fmt.Errorf("failed to connect to server: %w", err)}
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
bodyBytes, _ := io.ReadAll(resp.Body)
|
||||
resp.Body.Close()
|
||||
return restoreErrorMsg{fmt.Errorf("server returned non-200 status: %s\n%s", resp.Status, string(bodyBytes))}
|
||||
}
|
||||
return restoreSourceReadyMsg{reader: resp.Body, totalBytes: resp.ContentLength, isDir: m.isDirectory}
|
||||
|
||||
} else if m.isP2PSource {
|
||||
// Old P2P Mode
|
||||
listener, err := net.Listen("tcp", m.source)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to listen on %s: %w", m.source, err)
|
||||
return restoreErrorMsg{fmt.Errorf("failed to listen on %s: %w", m.source, err)}
|
||||
}
|
||||
defer listener.Close() // Close listener after accepting one connection
|
||||
|
||||
m.statusMessage = fmt.Sprintf("Waiting for connection on %s", m.source)
|
||||
// This is a blocking call
|
||||
// Accept is blocking, so it must be in a command.
|
||||
conn, err := listener.Accept()
|
||||
listener.Close() // Close listener after one connection.
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to accept connection: %w", err)
|
||||
return restoreErrorMsg{fmt.Errorf("failed to accept connection: %w", err)}
|
||||
}
|
||||
m.statusMessage = "Client connected. Receiving data..."
|
||||
|
||||
// Protocol: Read 8 bytes for size, 1 byte for type
|
||||
if err := binary.Read(conn, binary.BigEndian, &m.totalBytes); err != nil {
|
||||
var size int64
|
||||
if err := binary.Read(conn, binary.BigEndian, &size); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("failed to read backup size: %w", err)
|
||||
return restoreErrorMsg{fmt.Errorf("failed to read backup size: %w", err)}
|
||||
}
|
||||
if m.totalBytes == 0 {
|
||||
m.totalBytes = 1 // Avoid division by zero for empty files
|
||||
}
|
||||
|
||||
typeByte := make([]byte, 1)
|
||||
if _, err := io.ReadFull(conn, typeByte); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("failed to read backup type: %w", err)
|
||||
return restoreErrorMsg{fmt.Errorf("failed to read backup type: %w", err)}
|
||||
}
|
||||
m.isDirectory = (typeByte[0] == 1)
|
||||
|
||||
return conn, nil
|
||||
return restoreSourceReadyMsg{reader: conn, totalBytes: size, isDir: typeByte[0] == 1}
|
||||
}
|
||||
|
||||
// Logic for local file
|
||||
m.statusMessage = fmt.Sprintf("Opening archive %s", m.source)
|
||||
return os.Open(m.source)
|
||||
// Local File Mode
|
||||
fileInfo, err := os.Stat(m.source)
|
||||
if err != nil {
|
||||
return restoreErrorMsg{err}
|
||||
}
|
||||
file, err := os.Open(m.source)
|
||||
if err != nil {
|
||||
return restoreErrorMsg{err}
|
||||
}
|
||||
return restoreSourceReadyMsg{reader: file, totalBytes: fileInfo.Size(), isDir: m.isDirectory}
|
||||
}
|
||||
|
||||
func (m *restoreModel) startRestore() tea.Msg {
|
||||
go func() {
|
||||
sourceReader, err := m.getSourceReader()
|
||||
if err != nil {
|
||||
m.err = err
|
||||
m.progressCh <- -1
|
||||
return
|
||||
}
|
||||
defer sourceReader.Close()
|
||||
|
||||
progressReader := &restoreProgressReader{
|
||||
reader: sourceReader,
|
||||
func (m *restoreModel) processRestoreCmd(r io.ReadCloser) tea.Cmd {
|
||||
return func() tea.Msg {
|
||||
defer r.Close()
|
||||
progressReader := &progressTracker{
|
||||
Reader: r,
|
||||
progressCh: m.progressCh,
|
||||
}
|
||||
|
||||
@@ -183,12 +203,11 @@ func (m *restoreModel) startRestore() tea.Msg {
|
||||
if restoreErr != nil {
|
||||
m.err = restoreErr
|
||||
m.progressCh <- -1
|
||||
return
|
||||
} else {
|
||||
m.progressCh <- -2 // Done signal
|
||||
}
|
||||
|
||||
m.progressCh <- -2 // Done signal
|
||||
}()
|
||||
return nil
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func (m *restoreModel) restoreFileOrDisk(r io.Reader) error {
|
||||
@@ -212,7 +231,8 @@ func (m *restoreModel) restoreFileOrDisk(r io.Reader) error {
|
||||
}
|
||||
defer dstFile.Close()
|
||||
|
||||
if _, err = io.Copy(dstFile, gzipReader); err != nil {
|
||||
_, err = io.Copy(dstFile, gzipReader)
|
||||
if err != nil && err != io.EOF {
|
||||
return err
|
||||
}
|
||||
return dstFile.Sync()
|
||||
@@ -264,37 +284,19 @@ func (m *restoreModel) restoreDirectory(r io.Reader) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
type restoreProgressReader struct {
|
||||
reader io.Reader
|
||||
progressCh chan int64
|
||||
read int64
|
||||
}
|
||||
|
||||
func (r *restoreProgressReader) Read(p []byte) (int, error) {
|
||||
n, err := r.reader.Read(p)
|
||||
if n > 0 {
|
||||
r.read += int64(n)
|
||||
r.progressCh <- r.read
|
||||
}
|
||||
return n, err
|
||||
}
|
||||
|
||||
func (m *restoreModel) restoreProgressListener() tea.Msg {
|
||||
select {
|
||||
case n := <-m.progressCh:
|
||||
switch {
|
||||
case n == -1: // Ошибка
|
||||
n := <-m.progressCh
|
||||
if n < 0 {
|
||||
if n == -1 {
|
||||
if m.err == nil {
|
||||
m.err = fmt.Errorf("unknown restoration error")
|
||||
}
|
||||
return restoreErrorMsg{m.err}
|
||||
case n == -2: // Завершение
|
||||
return restoreDoneMsg{}
|
||||
default: // Прогресс
|
||||
m.processed = n
|
||||
return restoreProgressMsg(n)
|
||||
}
|
||||
return restoreDoneMsg{}
|
||||
}
|
||||
m.processed = n
|
||||
return restoreProgressMsg(n)
|
||||
}
|
||||
|
||||
func (m *restoreModel) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
||||
@@ -309,6 +311,31 @@ func (m *restoreModel) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
||||
m.progress.Width = msg.Width - 4
|
||||
return m, nil
|
||||
|
||||
case restoreSourceReadyMsg:
|
||||
m.totalBytes = msg.totalBytes
|
||||
if m.totalBytes <= 0 {
|
||||
m.totalBytes = 1
|
||||
}
|
||||
m.isDirectory = msg.isDir // Update directory type, especially for P2P.
|
||||
|
||||
// Final check for target directory now that we know the type for sure.
|
||||
if m.isDirectory && m.target == "" {
|
||||
wd, err := os.Getwd()
|
||||
if err != nil {
|
||||
return m, func() tea.Msg { return restoreErrorMsg{err} }
|
||||
}
|
||||
m.target = wd
|
||||
} else if !m.isDirectory && m.target == "" {
|
||||
err := fmt.Errorf("target flag -t is required for file/disk restoration")
|
||||
return m, func() tea.Msg { return restoreErrorMsg{err} }
|
||||
}
|
||||
|
||||
m.statusMessage = "Receiving data..."
|
||||
return m, tea.Batch(
|
||||
m.processRestoreCmd(msg.reader),
|
||||
m.restoreProgressListener,
|
||||
)
|
||||
|
||||
case restoreProgressMsg:
|
||||
progressVal := float64(msg) / float64(m.totalBytes)
|
||||
if progressVal > 1.0 {
|
||||
@@ -340,6 +367,9 @@ func (m *restoreModel) View() string {
|
||||
if m.isDirectory {
|
||||
restoreType = "Directory"
|
||||
}
|
||||
if m.httpSourceInfo != nil || m.isP2PSource {
|
||||
restoreType += " network"
|
||||
}
|
||||
return fmt.Sprintf("\n✅ %s restoration complete!\n\n"+
|
||||
"Source: %s\n"+
|
||||
"Target: %s\n"+
|
||||
@@ -348,12 +378,18 @@ func (m *restoreModel) View() string {
|
||||
duration.Round(time.Millisecond))
|
||||
}
|
||||
|
||||
// Show status message before progress starts
|
||||
if m.processed == 0 && (m.isNetworkSource || m.statusMessage != "") {
|
||||
if m.processed == 0 {
|
||||
return m.statusMessage + "\n\nPress Ctrl+C to cancel"
|
||||
}
|
||||
|
||||
title := fmt.Sprintf("Restoring from %s → %s", m.source, m.target)
|
||||
operation := "Restoring"
|
||||
if m.httpSourceInfo != nil {
|
||||
operation = "Downloading"
|
||||
} else if m.isP2PSource {
|
||||
operation = "Receiving"
|
||||
}
|
||||
|
||||
title := fmt.Sprintf("%s from %s → %s", operation, m.source, m.target)
|
||||
progressVal := float64(m.processed) / float64(m.totalBytes)
|
||||
progressView := m.progress.ViewAs(progressVal)
|
||||
stats := fmt.Sprintf("%s / %s (%.1f%%)", formatBytes(m.processed), formatBytes(m.totalBytes), progressVal*100)
|
||||
@@ -385,7 +421,7 @@ func isBlockDevice(path string) bool {
|
||||
|
||||
func init() {
|
||||
rootCmd.AddCommand(restoreCmd)
|
||||
restoreCmd.Flags().StringVarP(&restoreSource, "source", "s", "", "Source archive or network address (e.g., /path/to/archive.tar.gz or 0.0.0.0:8080) (required)")
|
||||
restoreCmd.Flags().StringVarP(&restoreTarget, "target", "t", "", "Target file, disk, or directory. Required for network sources.")
|
||||
restoreCmd.Flags().StringVarP(&restoreSource, "source", "s", "", "Source archive, P2P address (e.g., 0.0.0.0:8080), or HTTP URL (e.g., 127.0.0.1:8080/backup.gz) (required)")
|
||||
restoreCmd.Flags().StringVarP(&restoreTarget, "target", "t", "", "Target file, disk, or directory. Required for non-directory restores.")
|
||||
restoreCmd.MarkFlagRequired("source")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
var (
|
||||
serveAddress string
|
||||
serveDirectory string
|
||||
)
|
||||
|
||||
var serveCmd = &cobra.Command{
|
||||
Use: "serve",
|
||||
Short: "Run a backup server",
|
||||
Long: "Run an HTTP server to upload and download backups.",
|
||||
Run: func(cmd *cobra.Command, args []string) {
|
||||
if serveDirectory == "" {
|
||||
home, err := os.UserHomeDir()
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to get user home directory: %v", err)
|
||||
}
|
||||
serveDirectory = filepath.Join(home, "backups")
|
||||
}
|
||||
|
||||
if err := os.MkdirAll(serveDirectory, 0755); err != nil {
|
||||
log.Fatalf("Failed to create backup directory: %v", err)
|
||||
}
|
||||
|
||||
r := mux.NewRouter()
|
||||
r.HandleFunc("/backup/{filename}", uploadHandler).Methods("POST")
|
||||
r.HandleFunc("/backup/{filename}", downloadHandler).Methods("GET")
|
||||
|
||||
log.Printf("Starting server on %s", serveAddress)
|
||||
log.Printf("Using backup directory: %s", serveDirectory)
|
||||
if err := http.ListenAndServe(serveAddress, r); err != nil {
|
||||
log.Fatalf("Server failed: %v", err)
|
||||
}
|
||||
},
|
||||
}
|
||||
|
||||
func uploadHandler(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
filename := vars["filename"]
|
||||
filePath := filepath.Join(serveDirectory, filename)
|
||||
|
||||
file, err := os.Create(filePath)
|
||||
if err != nil {
|
||||
http.Error(w, "Failed to create file", http.StatusInternalServerError)
|
||||
log.Printf("Error creating file %s: %v", filename, err)
|
||||
return
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
_, err = io.Copy(file, r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, "Failed to write to file", http.StatusInternalServerError)
|
||||
log.Printf("Error writing to file %s: %v", filename, err)
|
||||
return
|
||||
}
|
||||
|
||||
w.WriteHeader(http.StatusOK)
|
||||
fmt.Fprintf(w, "File %s uploaded successfully.", filename)
|
||||
log.Printf("Uploaded %s", filename)
|
||||
}
|
||||
|
||||
func downloadHandler(w http.ResponseWriter, r *http.Request) {
|
||||
vars := mux.Vars(r)
|
||||
filename := vars["filename"]
|
||||
filePath := filepath.Join(serveDirectory, filename)
|
||||
|
||||
// Проверяем, существует ли файл, перед отправкой
|
||||
if _, err := os.Stat(filePath); os.IsNotExist(err) {
|
||||
http.NotFound(w, r)
|
||||
log.Printf("File not found: %s", filename)
|
||||
return
|
||||
}
|
||||
|
||||
// http.ServeFile - это идиоматический способ отправки файлов в Go.
|
||||
// Он автоматически устанавливает Content-Type, Content-Length и другие заголовки.
|
||||
http.ServeFile(w, r, filePath)
|
||||
log.Printf("Downloaded %s", filename)
|
||||
}
|
||||
|
||||
func init() {
|
||||
rootCmd.AddCommand(serveCmd)
|
||||
serveCmd.Flags().StringVarP(&serveAddress, "address", "a", "localhost:8080", "Address and port for the server")
|
||||
serveCmd.Flags().StringVarP(&serveDirectory, "directory", "d", "", "Directory to store backups (defaults to ~/backups)")
|
||||
}
|
||||
@@ -2,6 +2,8 @@ package cmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -22,3 +24,81 @@ func formatBytes(b int64) string {
|
||||
}
|
||||
return fmt.Sprintf("%.1f %ciB", float64(b)/float64(div), "KMGTPE"[exp])
|
||||
}
|
||||
|
||||
// ServerTargetInfo holds the parsed information from a server target string.
|
||||
type ServerTargetInfo struct {
|
||||
// Full HTTP URL for the request, e.g., http://127.0.0.1:8080/backup/mybackup.tar.gz
|
||||
URL string
|
||||
// The address of the server, e.g., 127.0.0.1:8080
|
||||
Address string
|
||||
// The filename for the backup, e.g., mybackup.tar.gz
|
||||
Filename string
|
||||
}
|
||||
|
||||
// parseServerTarget parses a target string like `127.0.0.1:8080/backup.gz`.
|
||||
// It returns a struct with the full URL and filename, or an error if the format is invalid.
|
||||
func parseServerTarget(target string) (*ServerTargetInfo, bool) {
|
||||
if !strings.Contains(target, "/") || !strings.Contains(target, ":") {
|
||||
return nil, false
|
||||
}
|
||||
|
||||
parts := strings.SplitN(target, "/", 2)
|
||||
if len(parts) != 2 {
|
||||
return nil, false // Invalid format
|
||||
}
|
||||
|
||||
address := parts[0]
|
||||
filename := parts[1]
|
||||
|
||||
if filename == "" {
|
||||
return nil, false // Filename cannot be empty
|
||||
}
|
||||
|
||||
// Validate that the first part is a host:port
|
||||
host, port, err := net.SplitHostPort(address)
|
||||
if err != nil {
|
||||
return nil, false // Not a valid host:port
|
||||
}
|
||||
|
||||
if host == "" || port == "" {
|
||||
return nil, false
|
||||
}
|
||||
|
||||
// It looks like a valid server target.
|
||||
fullURL := fmt.Sprintf("http://%s/backup/%s", address, filename)
|
||||
|
||||
return &ServerTargetInfo{
|
||||
URL: fullURL,
|
||||
Address: address,
|
||||
Filename: filename,
|
||||
}, true
|
||||
}
|
||||
|
||||
// progressTracker реализует io.Reader и io.Writer для отслеживания прогресса
|
||||
type progressTracker struct {
|
||||
Reader io.Reader
|
||||
Writer io.Writer
|
||||
progressCh chan int64
|
||||
// Для restore, где io.Copy может вызываться много раз (в tar),
|
||||
// нам нужно отслеживать общий прогресс.
|
||||
processed int64
|
||||
}
|
||||
|
||||
// Write отслеживает прогресс записи (для create)
|
||||
func (pt *progressTracker) Write(p []byte) (int, error) {
|
||||
n, err := pt.Writer.Write(p)
|
||||
if n > 0 {
|
||||
pt.progressCh <- int64(n) // Отправляем дельту
|
||||
}
|
||||
return n, err
|
||||
}
|
||||
|
||||
// Read отслеживает прогресс чтения (для restore)
|
||||
func (pt *progressTracker) Read(p []byte) (int, error) {
|
||||
n, err := pt.Reader.Read(p)
|
||||
if n > 0 {
|
||||
pt.processed += int64(n)
|
||||
pt.progressCh <- pt.processed // Отправляем общий обработанный объем
|
||||
}
|
||||
return n, err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user