windows as Source
This commit is contained in:
parent
f2f07a0f9a
commit
2806d6f8dd
89
README.md
89
README.md
@ -21,7 +21,12 @@ Looks nice, builds fast, runs everywhere with one Binary.
|
|||||||
Block-level sync for a large file or block device between two machines (or
|
Block-level sync for a large file or block device between two machines (or
|
||||||
two paths on the same machine), driven from a third, passive "manager"
|
two paths on the same machine), driven from a third, passive "manager"
|
||||||
machine. Single static Go binary, no runtime dependencies beyond the
|
machine. Single static Go binary, no runtime dependencies beyond the
|
||||||
system `ssh` client for remote endpoints. Linux only.
|
system `ssh` client for SSH endpoints.
|
||||||
|
|
||||||
|
Runs on **Linux and Windows**. Endpoints are reached either over SSH (as
|
||||||
|
below) or, for a host with no SSH server (typically Windows), by starting
|
||||||
|
`clonetool listen` there and connecting to it over an encrypted, password-
|
||||||
|
protected TCP channel — see **Listen mode (no SSH)** below.
|
||||||
|
|
||||||
## How it works
|
## How it works
|
||||||
|
|
||||||
@ -98,6 +103,46 @@ system `ssh` client for remote endpoints. Linux only.
|
|||||||
Shrinking an existing non-empty file prompts for confirmation unless
|
Shrinking an existing non-empty file prompts for confirmation unless
|
||||||
`--yes` is passed.
|
`--yes` is passed.
|
||||||
|
|
||||||
|
## Listen mode (no SSH)
|
||||||
|
|
||||||
|
For a host that has no SSH server — a Windows machine, typically — start
|
||||||
|
clonetool there in **listen mode** and point the manager's location at it
|
||||||
|
with a `tcp://host:port/path` spec instead of `host:path`. There is no
|
||||||
|
self-deploy over this path: you copy the binary to that host yourself and
|
||||||
|
start it by hand.
|
||||||
|
|
||||||
|
```
|
||||||
|
# On the source host (e.g. Windows), started manually:
|
||||||
|
clonetool listen --bind :9000 --password s3cret
|
||||||
|
|
||||||
|
# On the manager, pull that source into a local/SSH destination:
|
||||||
|
clonetool sync --source tcp://WINHOST:9000/C:\data\disk.img \
|
||||||
|
--dest backup:/srv/disk.img --password s3cret
|
||||||
|
```
|
||||||
|
|
||||||
|
- A listening host only ever **accepts** connections; it never dials out.
|
||||||
|
So when the **source** is a `tcp://` endpoint, clonetool always uses the
|
||||||
|
**pull** direction — the destination side connects to the listener and
|
||||||
|
pulls the changed blocks. The destination must therefore have a network
|
||||||
|
route to `host:port`. (Push is skipped for a listener source.)
|
||||||
|
- **Password:** required for any `tcp://` endpoint, on both the listener
|
||||||
|
(`clonetool listen`) and the manager (`clonetool sync`). Provide it with
|
||||||
|
`--password`, `--password-file PATH`, or the `CLONETOOL_PASSWORD`
|
||||||
|
environment variable (env/file keep it out of the process list and shell
|
||||||
|
history). The manager passes it to the destination agent over the already-
|
||||||
|
encrypted control channel, never on a command line.
|
||||||
|
- **Security:** the channel is TLS 1.3 (the listener uses a fresh, in-memory
|
||||||
|
self-signed certificate). The password is verified by a challenge-response
|
||||||
|
bound to the TLS session (HMAC over the connection's exported keying
|
||||||
|
material, key stretched with PBKDF2), so a man-in-the-middle can't
|
||||||
|
authenticate even though the certificate isn't checked against a CA. Both
|
||||||
|
ends authenticate each other; a wrong password fails the connection before
|
||||||
|
any file is touched. Uses only the Go standard library — still one static
|
||||||
|
binary.
|
||||||
|
- The listener serves connections until you stop it with Ctrl-C. A single
|
||||||
|
sync opens two connections to it (one for the stat, one for the data
|
||||||
|
pull), so it must stay running for the whole job.
|
||||||
|
|
||||||
## Build
|
## Build
|
||||||
|
|
||||||
```
|
```
|
||||||
@ -132,8 +177,12 @@ cross-compile if your hosts differ.
|
|||||||
clonetool sync --source LOC --dest LOC [options]
|
clonetool sync --source LOC --dest LOC [options]
|
||||||
```
|
```
|
||||||
|
|
||||||
`LOC` is either a local path (`/dev/sdb`, `./image.bin`) or
|
`LOC` is one of:
|
||||||
`[user@]host:path` for a path reached over SSH.
|
|
||||||
|
- a local path — `/dev/sdb`, `./image.bin`, `C:\data\disk.img`, `\\.\PhysicalDrive0`
|
||||||
|
- `[user@]host:path` — reached over SSH
|
||||||
|
- `tcp://host:port/path` — a host running `clonetool listen` (see **Listen
|
||||||
|
mode** above); no SSH needed there
|
||||||
|
|
||||||
```
|
```
|
||||||
# Same machine
|
# Same machine
|
||||||
@ -145,24 +194,38 @@ clonetool sync --source /dev/sda --dest /srv/sda.img
|
|||||||
# Two remote machines, orchestrated from a third
|
# Two remote machines, orchestrated from a third
|
||||||
clonetool sync --source db1:/dev/vdb --dest backup-host:/srv/db1.img
|
clonetool sync --source db1:/dev/vdb --dest backup-host:/srv/db1.img
|
||||||
|
|
||||||
|
# Windows source with no SSH: it runs `clonetool listen --bind :9000 --password p`
|
||||||
|
clonetool sync --source tcp://winbox:9000/C:\data\disk.img \
|
||||||
|
--dest backup-host:/srv/win.img --password p
|
||||||
|
|
||||||
# Re-run any time; only changed blocks move
|
# Re-run any time; only changed blocks move
|
||||||
clonetool sync --source db1:/dev/vdb --dest backup-host:/srv/db1.img
|
clonetool sync --source db1:/dev/vdb --dest backup-host:/srv/db1.img
|
||||||
```
|
```
|
||||||
|
|
||||||
Options:
|
Options for `sync`:
|
||||||
|
|
||||||
| Flag | Default | Meaning |
|
| Flag | Default | Meaning |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
| `--block-size` | `4M` | Block size (accepts `K`/`M`/`G` suffixes). |
|
| `--block-size` | `4M` | Block size (accepts `K`/`M`/`G` suffixes). |
|
||||||
| `--job` | — | Optional label shown in progress/log output. |
|
| `--job` | — | Optional label shown in progress/log output. |
|
||||||
| `--yes` | off | Don't prompt before shrinking an existing destination file. |
|
| `--yes` | off | Don't prompt before shrinking an existing destination file. |
|
||||||
| `--sudo` | `auto` | Block-device privilege escalation: `auto` (on a permission error), `always`, or `never`. Remote elevation needs passwordless sudo. |
|
| `--sudo` | `auto` | Block-device privilege escalation: `auto` (on a permission error), `always`, or `never`. Remote elevation needs passwordless sudo. No-op on Windows (run as Administrator instead). |
|
||||||
| `--deploy` | `true` | Copy this binary to remote hosts that lack a runnable `clonetool`. `--deploy=false` to disable. |
|
| `--deploy` | `true` | Copy this binary to SSH hosts that lack a runnable `clonetool`. `--deploy=false` to disable. (Never applies to `tcp://` endpoints.) |
|
||||||
| `--connect-timeout` | `8` | SSH connect timeout (seconds) used for the push/pull direction probe. |
|
| `--connect-timeout` | `8` | SSH/TLS connect timeout (seconds) used for the push/pull direction probe. |
|
||||||
| `--ssh` | `ssh` | ssh binary to use. |
|
| `--ssh` | `ssh` | ssh binary to use. |
|
||||||
| `--ssh-opt` | — | Extra `-o OPT` passed to ssh (repeatable). |
|
| `--ssh-opt` | — | Extra `-o OPT` passed to ssh (repeatable). |
|
||||||
| `--remote-bin` | `clonetool` | Path to clonetool on remote hosts. |
|
| `--remote-bin` | `clonetool` | Path to clonetool on remote SSH hosts. |
|
||||||
| `--manager-host` | local hostname | Address a peer should use to reach this machine, needed only when source or dest is local to the manager *and* the other side is remote and ends up needing to dial back in (pull fallback). |
|
| `--manager-host` | local hostname | Address a peer should use to reach this machine, needed only when source or dest is local to the manager *and* the other side is remote and ends up needing to dial back in (pull fallback). |
|
||||||
|
| `--password` | — | Shared password for a `tcp://` listen endpoint. |
|
||||||
|
| `--password-file` | — | Read the listen-endpoint password from a file (or set `CLONETOOL_PASSWORD`). |
|
||||||
|
|
||||||
|
Options for `listen`:
|
||||||
|
|
||||||
|
| Flag | Default | Meaning |
|
||||||
|
|---|---|---|
|
||||||
|
| `--bind` | `:9000` | Address to listen on, e.g. `:9000` or `0.0.0.0:9000`. |
|
||||||
|
| `--password` | — | Shared password (or `--password-file`, or `CLONETOOL_PASSWORD`). Required. |
|
||||||
|
| `--password-file` | — | Read the shared password from a file. |
|
||||||
|
|
||||||
`clonetool version` prints the binary's `GOOS/GOARCH` and build timestamp,
|
`clonetool version` prints the binary's `GOOS/GOARCH` and build timestamp,
|
||||||
e.g. `clonetool linux/amd64 build=2024-06-01T12:00:00Z` (the timestamp is
|
e.g. `clonetool linux/amd64 build=2024-06-01T12:00:00Z` (the timestamp is
|
||||||
@ -176,7 +239,15 @@ and `wr(dst)` is the destination actually writing changed blocks.
|
|||||||
|
|
||||||
## Caveats
|
## Caveats
|
||||||
|
|
||||||
- Block-device size detection uses `BLKGETSIZE64`; the tool is Linux only.
|
- Block-device sizing uses `BLKGETSIZE64` on Linux and
|
||||||
|
`IOCTL_DISK_GET_LENGTH_INFO` on Windows (`\\.\PhysicalDrive0`, `\\.\C:`).
|
||||||
|
On other platforms (e.g. macOS) only regular files are supported; a
|
||||||
|
device path there errors out.
|
||||||
|
- **Windows raw disks:** syncing a `\\.\PhysicalDrive*` requires running
|
||||||
|
clonetool as **Administrator**, and the disk should be offline/unmounted
|
||||||
|
(a live, mounted volume can refuse writes or give inconsistent reads).
|
||||||
|
Raw-disk I/O is sector-aligned automatically; `--block-size` must be a
|
||||||
|
multiple of the sector size. `--sudo` does nothing on Windows.
|
||||||
- If a destination path doesn't exist yet, it's created as a regular
|
- If a destination path doesn't exist yet, it's created as a regular
|
||||||
file — clonetool won't create device nodes, so double-check device
|
file — clonetool won't create device nodes, so double-check device
|
||||||
paths for typos before running.
|
paths for typos before running.
|
||||||
|
|||||||
178
agent.go
178
agent.go
@ -25,13 +25,15 @@ func cmdAgent(args []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
in := NewFrameReader(os.Stdin)
|
||||||
|
out := NewFrameWriter(os.Stdout)
|
||||||
switch *role {
|
switch *role {
|
||||||
case "control":
|
case roleControl:
|
||||||
return runControlAgent()
|
return runControlAgent(in, out)
|
||||||
case "sink":
|
case roleSink:
|
||||||
return runSinkRole(*path, *base, *size, *blockSize)
|
return runSinkRole(in, out, *path, *base, *size, *blockSize)
|
||||||
case "source-stream":
|
case roleSourceStream:
|
||||||
return runSourceStreamRole(*path, *base, *size, *blockSize)
|
return runSourceStreamRole(in, out, *path, *base, *size, *blockSize)
|
||||||
default:
|
default:
|
||||||
return fmt.Errorf("agent: unknown or missing --role %q", *role)
|
return fmt.Errorf("agent: unknown or missing --role %q", *role)
|
||||||
}
|
}
|
||||||
@ -42,10 +44,7 @@ func cmdAgent(args []string) error {
|
|||||||
// manager over stdin/stdout with CtrlMsg frames.
|
// manager over stdin/stdout with CtrlMsg frames.
|
||||||
// ---------------------------------------------------------------------
|
// ---------------------------------------------------------------------
|
||||||
|
|
||||||
func runControlAgent() error {
|
func runControlAgent(in *FrameReader, out *FrameWriter) error {
|
||||||
in := NewFrameReader(os.Stdin)
|
|
||||||
out := NewFrameWriter(os.Stdout)
|
|
||||||
|
|
||||||
for {
|
for {
|
||||||
typ, payload, err := in.ReadFrame()
|
typ, payload, err := in.ReadFrame()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@ -126,38 +125,15 @@ func peerAgentCommand(req CtrlMsg, tailArgs []string) *exec.Cmd {
|
|||||||
|
|
||||||
func runPushDriver(req CtrlMsg, out *FrameWriter) {
|
func runPushDriver(req CtrlMsg, out *FrameWriter) {
|
||||||
tailArgs := []string{
|
tailArgs := []string{
|
||||||
"agent", "--role", "sink",
|
"agent", "--role", roleSink,
|
||||||
"--path", req.PeerPath,
|
"--path", req.PeerPath,
|
||||||
"--base", strconv.FormatInt(req.PeerBase, 10),
|
"--base", strconv.FormatInt(req.PeerBase, 10),
|
||||||
"--size", strconv.FormatInt(req.Size, 10),
|
"--size", strconv.FormatInt(req.Size, 10),
|
||||||
"--block-size", strconv.FormatInt(req.BlockSize, 10),
|
"--block-size", strconv.FormatInt(req.BlockSize, 10),
|
||||||
}
|
}
|
||||||
cmd := peerAgentCommand(req, tailArgs)
|
link, reason := openPeerLink(req, roleSink, tailArgs)
|
||||||
stdin, err := cmd.StdinPipe()
|
if reason != "" {
|
||||||
if err != nil {
|
_ = out.WriteJSON(CtrlMsg{Type: msgPushFailed, Reason: reason})
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPushFailed, Reason: err.Error()})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
stdout, err := cmd.StdoutPipe()
|
|
||||||
if err != nil {
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPushFailed, Reason: err.Error()})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
stderrBuf := newLimitedBuffer(4096)
|
|
||||||
cmd.Stderr = stderrBuf
|
|
||||||
|
|
||||||
if err := cmd.Start(); err != nil {
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPushFailed, Reason: fmt.Sprintf("start ssh: %v", err)})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
fw := NewFrameWriter(stdin)
|
|
||||||
fr := NewFrameReader(stdout)
|
|
||||||
|
|
||||||
timeout := time.Duration(req.ConnectTimeoutSec+2) * time.Second
|
|
||||||
if err := waitReady(fr, timeout); err != nil {
|
|
||||||
_ = cmd.Process.Kill()
|
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPushFailed, Reason: fmt.Sprintf("%v (remote stderr: %s)", err, stderrBuf.String())})
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -167,26 +143,82 @@ func runPushDriver(req CtrlMsg, out *FrameWriter) {
|
|||||||
// as they arrive, so the two scans and the transfer all overlap.
|
// as they arrive, so the two scans and the transfer all overlap.
|
||||||
srcFile, err := os.Open(req.Path)
|
srcFile, err := os.Open(req.Path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = cmd.Process.Kill()
|
link.kill()
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: err.Error(), NeedPriv: isPermErr(err)})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: err.Error(), NeedPriv: isPermErr(err)})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
defer srcFile.Close()
|
defer srcFile.Close()
|
||||||
|
|
||||||
if fatalErr := pumpPush(req, alignmentFor(req.Path), srcFile, fw, fr, out); fatalErr != nil {
|
if fatalErr := pumpPush(req, alignmentFor(req.Path), srcFile, link.fw, link.fr, out); fatalErr != nil {
|
||||||
_ = cmd.Process.Kill()
|
link.kill()
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fatalErr.Error()})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fatalErr.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if err := cmd.Wait(); err != nil {
|
if err := link.finish(); err != nil {
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fmt.Sprintf("sink process: %v (stderr: %s)", err, stderrBuf.String())})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fmt.Sprintf("sink: %v (stderr: %s)", err, link.stderr())})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPushOK})
|
_ = out.WriteJSON(CtrlMsg{Type: msgPushOK})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// peerLink is a framed, READY-confirmed connection to a peer helper (sink or
|
||||||
|
// source-stream), reached either by spawning it over ssh/local (spawnPeerLink)
|
||||||
|
// or by dialing a listening clonetool over TLS (dialPeerLink). It hides which
|
||||||
|
// mechanism is in use from the push/pull drivers.
|
||||||
|
type peerLink struct {
|
||||||
|
fw *FrameWriter
|
||||||
|
fr *FrameReader
|
||||||
|
stderr func() string // bounded tail of the peer's stderr, "" for a net peer
|
||||||
|
finish func() error // wait for clean completion (cmd.Wait / conn close)
|
||||||
|
kill func() // force teardown on error
|
||||||
|
closeSend func() error // close only the send side, to unblock a parked peer
|
||||||
|
}
|
||||||
|
|
||||||
|
// openPeerLink establishes a peerLink for role, dialing a listener when the
|
||||||
|
// request targets a net peer and otherwise spawning the helper. A non-empty
|
||||||
|
// reason means setup/handshake failed softly, so the caller reports
|
||||||
|
// <push|pull>_failed and the manager can try the other direction.
|
||||||
|
func openPeerLink(req CtrlMsg, role string, tailArgs []string) (*peerLink, string) {
|
||||||
|
if req.PeerNet {
|
||||||
|
return dialPeerLink(req, role)
|
||||||
|
}
|
||||||
|
return spawnPeerLink(req, tailArgs)
|
||||||
|
}
|
||||||
|
|
||||||
|
// spawnPeerLink runs the helper over ssh (or as a local subprocess) and waits
|
||||||
|
// for its READY handshake.
|
||||||
|
func spawnPeerLink(req CtrlMsg, tailArgs []string) (*peerLink, string) {
|
||||||
|
cmd := peerAgentCommand(req, tailArgs)
|
||||||
|
stdin, err := cmd.StdinPipe()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err.Error()
|
||||||
|
}
|
||||||
|
stdout, err := cmd.StdoutPipe()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err.Error()
|
||||||
|
}
|
||||||
|
stderrBuf := newLimitedBuffer(4096)
|
||||||
|
cmd.Stderr = stderrBuf
|
||||||
|
if err := cmd.Start(); err != nil {
|
||||||
|
return nil, fmt.Sprintf("start ssh: %v", err)
|
||||||
|
}
|
||||||
|
fr := NewFrameReader(stdout)
|
||||||
|
timeout := time.Duration(req.ConnectTimeoutSec+2) * time.Second
|
||||||
|
if err := waitReady(fr, timeout); err != nil {
|
||||||
|
_ = cmd.Process.Kill()
|
||||||
|
_ = cmd.Wait()
|
||||||
|
return nil, fmt.Sprintf("%v (remote stderr: %s)", err, stderrBuf.String())
|
||||||
|
}
|
||||||
|
return &peerLink{
|
||||||
|
fw: NewFrameWriter(stdin),
|
||||||
|
fr: fr,
|
||||||
|
stderr: stderrBuf.String,
|
||||||
|
finish: cmd.Wait,
|
||||||
|
kill: func() { _ = cmd.Process.Kill(); _ = cmd.Wait() },
|
||||||
|
closeSend: stdin.Close,
|
||||||
|
}, ""
|
||||||
|
}
|
||||||
|
|
||||||
func waitReady(fr *FrameReader, timeout time.Duration) error {
|
func waitReady(fr *FrameReader, timeout time.Duration) error {
|
||||||
type result struct {
|
type result struct {
|
||||||
typ frameType
|
typ frameType
|
||||||
@ -428,40 +460,19 @@ func pumpPush(req CtrlMsg, align int64, srcFile *os.File, fw *FrameWriter, fr *F
|
|||||||
|
|
||||||
func runPullDriver(req CtrlMsg, out *FrameWriter) {
|
func runPullDriver(req CtrlMsg, out *FrameWriter) {
|
||||||
tailArgs := []string{
|
tailArgs := []string{
|
||||||
"agent", "--role", "source-stream",
|
"agent", "--role", roleSourceStream,
|
||||||
"--path", req.PeerPath,
|
"--path", req.PeerPath,
|
||||||
"--base", strconv.FormatInt(req.PeerBase, 10),
|
"--base", strconv.FormatInt(req.PeerBase, 10),
|
||||||
"--size", strconv.FormatInt(req.Size, 10),
|
"--size", strconv.FormatInt(req.Size, 10),
|
||||||
"--block-size", strconv.FormatInt(req.BlockSize, 10),
|
"--block-size", strconv.FormatInt(req.BlockSize, 10),
|
||||||
}
|
}
|
||||||
cmd := peerAgentCommand(req, tailArgs)
|
link, reason := openPeerLink(req, roleSourceStream, tailArgs)
|
||||||
stdin, err := cmd.StdinPipe()
|
if reason != "" {
|
||||||
if err != nil {
|
_ = out.WriteJSON(CtrlMsg{Type: msgPullFailed, Reason: reason})
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPullFailed, Reason: err.Error()})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
stdout, err := cmd.StdoutPipe()
|
|
||||||
if err != nil {
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPullFailed, Reason: err.Error()})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
stderrBuf := newLimitedBuffer(4096)
|
|
||||||
cmd.Stderr = stderrBuf
|
|
||||||
|
|
||||||
if err := cmd.Start(); err != nil {
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPullFailed, Reason: fmt.Sprintf("start ssh: %v", err)})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
fw := NewFrameWriter(stdin)
|
|
||||||
fr := NewFrameReader(stdout)
|
|
||||||
|
|
||||||
timeout := time.Duration(req.ConnectTimeoutSec+2) * time.Second
|
|
||||||
if err := waitReady(fr, timeout); err != nil {
|
|
||||||
_ = cmd.Process.Kill()
|
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPullFailed, Reason: fmt.Sprintf("%v (remote stderr: %s)", err, stderrBuf.String())})
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
fw := link.fw
|
||||||
|
fr := link.fr
|
||||||
|
|
||||||
// Handshake succeeded. Open the destination and fingerprint its current
|
// Handshake succeeded. Open the destination and fingerprint its current
|
||||||
// content block by block, streaming each hash to the source stream the
|
// content block by block, streaming each hash to the source stream the
|
||||||
@ -469,8 +480,7 @@ func runPullDriver(req CtrlMsg, out *FrameWriter) {
|
|||||||
// whatever it streams back into the same file concurrently.
|
// whatever it streams back into the same file concurrently.
|
||||||
dstFile, err := os.OpenFile(req.Path, os.O_RDWR, 0)
|
dstFile, err := os.OpenFile(req.Path, os.O_RDWR, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = cmd.Process.Kill()
|
link.kill()
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: err.Error(), NeedPriv: isPermErr(err)})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: err.Error(), NeedPriv: isPermErr(err)})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@ -488,7 +498,7 @@ func runPullDriver(req CtrlMsg, out *FrameWriter) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
// Unblock the source stream (waiting for more hashes) so the
|
// Unblock the source stream (waiting for more hashes) so the
|
||||||
// dest loop below can unwind instead of hanging.
|
// dest loop below can unwind instead of hanging.
|
||||||
_ = stdin.Close()
|
_ = link.closeSend()
|
||||||
}
|
}
|
||||||
hashErrCh <- err
|
hashErrCh <- err
|
||||||
}()
|
}()
|
||||||
@ -514,19 +524,17 @@ func runPullDriver(req CtrlMsg, out *FrameWriter) {
|
|||||||
hashErr := <-hashErrCh
|
hashErr := <-hashErrCh
|
||||||
|
|
||||||
if hashErr != nil {
|
if hashErr != nil {
|
||||||
_ = cmd.Process.Kill()
|
link.kill()
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fmt.Sprintf("scan destination: %v", hashErr), NeedPriv: isPermErr(hashErr)})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fmt.Sprintf("scan destination: %v", hashErr), NeedPriv: isPermErr(hashErr)})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if loopErr != nil {
|
if loopErr != nil {
|
||||||
_ = cmd.Process.Kill()
|
link.kill()
|
||||||
_ = cmd.Wait()
|
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: loopErr.Error()})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: loopErr.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if err := cmd.Wait(); err != nil {
|
if err := link.finish(); err != nil {
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fmt.Sprintf("source-stream process: %v (stderr: %s)", err, stderrBuf.String())})
|
_ = out.WriteJSON(CtrlMsg{Type: msgError, Message: fmt.Sprintf("source-stream: %v (stderr: %s)", err, link.stderr())})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
_ = out.WriteJSON(CtrlMsg{Type: msgPullOK})
|
_ = out.WriteJSON(CtrlMsg{Type: msgPullOK})
|
||||||
@ -537,10 +545,7 @@ func runPullDriver(req CtrlMsg, out *FrameWriter) {
|
|||||||
// destination host. Dumb write endpoint: verify+pwrite+ack per block.
|
// destination host. Dumb write endpoint: verify+pwrite+ack per block.
|
||||||
// ---------------------------------------------------------------------
|
// ---------------------------------------------------------------------
|
||||||
|
|
||||||
func runSinkRole(path string, base, size, blockSize int64) error {
|
func runSinkRole(in *FrameReader, out *FrameWriter, path string, base, size, blockSize int64) error {
|
||||||
out := NewFrameWriter(os.Stdout)
|
|
||||||
in := NewFrameReader(os.Stdin)
|
|
||||||
|
|
||||||
f, err := os.OpenFile(path, os.O_RDWR, 0)
|
f, err := os.OpenFile(path, os.O_RDWR, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Fprintf(os.Stderr, "sink: open %s: %v\n", path, err)
|
fmt.Fprintf(os.Stderr, "sink: open %s: %v\n", path, err)
|
||||||
@ -622,10 +627,7 @@ func runSinkRole(path string, base, size, blockSize int64) error {
|
|||||||
// writing straight to its own stdout.
|
// writing straight to its own stdout.
|
||||||
// ---------------------------------------------------------------------
|
// ---------------------------------------------------------------------
|
||||||
|
|
||||||
func runSourceStreamRole(path string, base, size, blockSize int64) error {
|
func runSourceStreamRole(in *FrameReader, out *FrameWriter, path string, base, size, blockSize int64) error {
|
||||||
out := NewFrameWriter(os.Stdout)
|
|
||||||
in := NewFrameReader(os.Stdin)
|
|
||||||
|
|
||||||
f, err := os.Open(path)
|
f, err := os.Open(path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Fprintf(os.Stderr, "source-stream: open %s: %v\n", path, err)
|
fmt.Fprintf(os.Stderr, "source-stream: open %s: %v\n", path, err)
|
||||||
|
|||||||
2
build.sh
2
build.sh
@ -15,6 +15,8 @@ build() {
|
|||||||
|
|
||||||
build linux amd64 clonetool-linux-amd64
|
build linux amd64 clonetool-linux-amd64
|
||||||
build linux arm64 clonetool-linux-arm64
|
build linux arm64 clonetool-linux-arm64
|
||||||
|
build windows amd64 clonetool-windows-amd64.exe
|
||||||
|
build windows arm64 clonetool-windows-arm64.exe
|
||||||
|
|
||||||
echo "done:"
|
echo "done:"
|
||||||
ls -la dist/
|
ls -la dist/
|
||||||
|
|||||||
25
control.go
25
control.go
@ -6,8 +6,11 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
|
"net"
|
||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
|
"strconv"
|
||||||
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
// errNeedPriv is wrapped into the error from a control-agent call that
|
// errNeedPriv is wrapped into the error from a control-agent call that
|
||||||
@ -187,6 +190,9 @@ func (c *Controller) connectAndPump(req CtrlMsg, cb transferCallbacks, okType, f
|
|||||||
func (c *Controller) Close() error {
|
func (c *Controller) Close() error {
|
||||||
_ = c.fw.WriteJSON(CtrlMsg{Type: msgClose})
|
_ = c.fw.WriteJSON(CtrlMsg{Type: msgClose})
|
||||||
_ = c.in.Close()
|
_ = c.in.Close()
|
||||||
|
if c.cmd == nil {
|
||||||
|
return nil // net transport: closing the connection is enough
|
||||||
|
}
|
||||||
err := c.cmd.Wait()
|
err := c.cmd.Wait()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
var exitErr *exec.ExitError
|
var exitErr *exec.ExitError
|
||||||
@ -197,3 +203,22 @@ func (c *Controller) Close() error {
|
|||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// startNetworkController connects to a listening clonetool (spec must be a net
|
||||||
|
// endpoint) over TLS, authenticates with password, and requests the control
|
||||||
|
// role. The returned Controller behaves like an ssh/local one for the manager;
|
||||||
|
// it has no child process (cmd == nil).
|
||||||
|
func startNetworkController(spec Spec, tag, password string, connectTimeoutSec int) (*Controller, error) {
|
||||||
|
_ = tag // errors are prefixed by the caller (source:/dest:)
|
||||||
|
addr := net.JoinHostPort(spec.Host, strconv.Itoa(spec.Port))
|
||||||
|
timeout := time.Duration(connectTimeoutSec+2) * time.Second
|
||||||
|
conn, fr, fw, err := dialAuth(addr, password, timeout)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("connect to listener %s: %w", addr, err)
|
||||||
|
}
|
||||||
|
if err := fw.WriteJSON(CtrlMsg{Type: msgRole, Role: roleControl}); err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, fmt.Errorf("request control role on %s: %w", addr, err)
|
||||||
|
}
|
||||||
|
return &Controller{tag: tag, in: conn, fw: fw, fr: fr}, nil
|
||||||
|
}
|
||||||
|
|||||||
24
ctrlmsg.go
24
ctrlmsg.go
@ -38,6 +38,18 @@ type CtrlMsg struct {
|
|||||||
// Sudo tells a push/pull driver to run the peer helper it spawns
|
// Sudo tells a push/pull driver to run the peer helper it spawns
|
||||||
// (sink / source-stream) under "sudo -n" (remote) or "sudo" (local).
|
// (sink / source-stream) under "sudo -n" (remote) or "sudo" (local).
|
||||||
Sudo bool `json:"sudo,omitempty"`
|
Sudo bool `json:"sudo,omitempty"`
|
||||||
|
// Peer listen endpoint: when PeerNet is set, the push/pull driver reaches
|
||||||
|
// the peer by dialing a listening clonetool over TLS (authenticated with
|
||||||
|
// PeerPassword) instead of spawning it over ssh. PeerHost carries the
|
||||||
|
// listener host in that case; PeerPort its port. PeerPassword travels only
|
||||||
|
// over the manager<->control channel (ssh-encrypted or TLS), never argv.
|
||||||
|
PeerNet bool `json:"peerNet,omitempty"`
|
||||||
|
PeerPort int `json:"peerPort,omitempty"`
|
||||||
|
PeerPassword string `json:"peerPassword,omitempty"`
|
||||||
|
|
||||||
|
// Role is the endpoint requested on a listen connection's role-request
|
||||||
|
// frame (control | sink | source-stream), see nettransport.go / listen.go.
|
||||||
|
Role string `json:"role,omitempty"`
|
||||||
|
|
||||||
// failure/error detail
|
// failure/error detail
|
||||||
Reason string `json:"reason,omitempty"`
|
Reason string `json:"reason,omitempty"`
|
||||||
@ -78,4 +90,16 @@ const (
|
|||||||
msgError = "error"
|
msgError = "error"
|
||||||
msgClose = "close"
|
msgClose = "close"
|
||||||
msgBye = "bye"
|
msgBye = "bye"
|
||||||
|
// msgRole is the first frame on a listen connection after authentication:
|
||||||
|
// it names the role to run (control | sink | source-stream) and, for the
|
||||||
|
// one-shot data roles, the path/window to operate on.
|
||||||
|
msgRole = "role"
|
||||||
|
msgRoleOK = "role_ok"
|
||||||
|
msgRoleErr = "role_err"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
roleControl = "control"
|
||||||
|
roleSink = "sink"
|
||||||
|
roleSourceStream = "source-stream"
|
||||||
)
|
)
|
||||||
|
|||||||
20
device.go
20
device.go
@ -14,15 +14,21 @@ type PathInfo struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// alignmentFor returns the offset/length alignment a path's handle requires
|
// alignmentFor returns the offset/length alignment a path's handle requires
|
||||||
// for positioned reads and writes. Linux block devices accept ordinary
|
// for positioned reads and writes. On Linux this is always 1 (block devices
|
||||||
// buffered pread/pwrite at any alignment, so there is nothing to round to.
|
// accept buffered pread/pwrite at any alignment); on Windows a raw disk handle
|
||||||
func alignmentFor(string) int64 { return 1 }
|
// requires sector-aligned I/O, so deviceAlignment reports the sector size.
|
||||||
|
func alignmentFor(path string) int64 { return deviceAlignment(path) }
|
||||||
// canElevate reports whether a permission failure opening a device is worth
|
|
||||||
// retrying under `sudo` (see --sudo).
|
|
||||||
func canElevate() bool { return true }
|
|
||||||
|
|
||||||
func statPath(path string) (PathInfo, error) {
|
func statPath(path string) (PathInfo, error) {
|
||||||
|
// Windows raw-disk paths (\\.\PhysicalDrive0) aren't reported as devices by
|
||||||
|
// os.Stat, so recognize them by name first and size them via ioctl.
|
||||||
|
if isRawDevice(path) {
|
||||||
|
sz, err := blockDeviceSize(path)
|
||||||
|
if err != nil {
|
||||||
|
return PathInfo{}, fmt.Errorf("stat device %s: %w", path, err)
|
||||||
|
}
|
||||||
|
return PathInfo{Exists: true, IsDevice: true, Size: sz}, nil
|
||||||
|
}
|
||||||
fi, err := os.Stat(path)
|
fi, err := os.Stat(path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if os.IsNotExist(err) {
|
if os.IsNotExist(err) {
|
||||||
|
|||||||
@ -26,3 +26,16 @@ func blockDeviceSize(path string) (int64, error) {
|
|||||||
}
|
}
|
||||||
return int64(size), nil
|
return int64(size), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// isRawDevice reports paths that must be sized via ioctl rather than os.Stat.
|
||||||
|
// On Linux block devices are recognized through os.ModeDevice instead, so this
|
||||||
|
// is always false.
|
||||||
|
func isRawDevice(string) bool { return false }
|
||||||
|
|
||||||
|
// deviceAlignment is 1 on Linux: buffered pread/pwrite on a block device accept
|
||||||
|
// any alignment.
|
||||||
|
func deviceAlignment(string) int64 { return 1 }
|
||||||
|
|
||||||
|
// canElevate reports whether a permission failure opening a device is worth
|
||||||
|
// retrying under sudo (see --sudo).
|
||||||
|
func canElevate() bool { return true }
|
||||||
|
|||||||
22
device_other.go
Normal file
22
device_other.go
Normal file
@ -0,0 +1,22 @@
|
|||||||
|
//go:build !linux && !windows
|
||||||
|
|
||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"runtime"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Stubs so clonetool still builds (and runs for regular-file syncs) on
|
||||||
|
// platforms without raw block-device support wired up (e.g. macOS for local
|
||||||
|
// development). Block-device paths error out clearly.
|
||||||
|
|
||||||
|
func isRawDevice(string) bool { return false }
|
||||||
|
|
||||||
|
func blockDeviceSize(string) (int64, error) {
|
||||||
|
return 0, fmt.Errorf("block devices are not supported on %s", runtime.GOOS)
|
||||||
|
}
|
||||||
|
|
||||||
|
func deviceAlignment(string) int64 { return 1 }
|
||||||
|
|
||||||
|
func canElevate() bool { return true }
|
||||||
98
device_windows.go
Normal file
98
device_windows.go
Normal file
@ -0,0 +1,98 @@
|
|||||||
|
//go:build windows
|
||||||
|
|
||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/binary"
|
||||||
|
"strings"
|
||||||
|
"syscall"
|
||||||
|
)
|
||||||
|
|
||||||
|
// IOCTL codes (winioctl.h):
|
||||||
|
//
|
||||||
|
// IOCTL_DISK_GET_LENGTH_INFO -> GET_LENGTH_INFORMATION { LARGE_INTEGER Length }
|
||||||
|
// IOCTL_DISK_GET_DRIVE_GEOMETRY -> DISK_GEOMETRY (BytesPerSector at offset 20)
|
||||||
|
const (
|
||||||
|
ioctlDiskGetLengthInfo = 0x0007405C
|
||||||
|
ioctlDiskGetDriveGeometry = 0x00070000
|
||||||
|
)
|
||||||
|
|
||||||
|
// isRawDevice recognizes the raw-disk paths that os.Stat won't report as
|
||||||
|
// devices: \\.\PhysicalDrive<N> and the volume form \\.\X:.
|
||||||
|
func isRawDevice(path string) bool {
|
||||||
|
l := strings.ToLower(path)
|
||||||
|
if strings.HasPrefix(l, `\\.\physicaldrive`) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
// \\.\C: — a single drive letter behind the device namespace.
|
||||||
|
if strings.HasPrefix(path, `\\.\`) && len(path) == 6 && path[5] == ':' {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// openDevice opens a raw device handle with the sharing a live disk needs.
|
||||||
|
func openDevice(path string, write bool) (syscall.Handle, error) {
|
||||||
|
p, err := syscall.UTF16PtrFromString(path)
|
||||||
|
if err != nil {
|
||||||
|
return syscall.InvalidHandle, err
|
||||||
|
}
|
||||||
|
access := uint32(syscall.GENERIC_READ)
|
||||||
|
if write {
|
||||||
|
access |= syscall.GENERIC_WRITE
|
||||||
|
}
|
||||||
|
return syscall.CreateFile(p, access,
|
||||||
|
syscall.FILE_SHARE_READ|syscall.FILE_SHARE_WRITE, nil,
|
||||||
|
syscall.OPEN_EXISTING, 0, 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
func blockDeviceSize(path string) (int64, error) {
|
||||||
|
h, err := openDevice(path, false)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
defer syscall.CloseHandle(h)
|
||||||
|
|
||||||
|
var out [8]byte
|
||||||
|
var ret uint32
|
||||||
|
if err := syscall.DeviceIoControl(h, ioctlDiskGetLengthInfo, nil, 0,
|
||||||
|
&out[0], uint32(len(out)), &ret, nil); err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return int64(binary.LittleEndian.Uint64(out[:])), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// deviceSectorSize queries the disk's physical sector size (bytes-per-sector),
|
||||||
|
// falling back to 512 if the geometry can't be read.
|
||||||
|
func deviceSectorSize(path string) int64 {
|
||||||
|
h, err := openDevice(path, false)
|
||||||
|
if err != nil {
|
||||||
|
return 512
|
||||||
|
}
|
||||||
|
defer syscall.CloseHandle(h)
|
||||||
|
|
||||||
|
var geom [24]byte
|
||||||
|
var ret uint32
|
||||||
|
if err := syscall.DeviceIoControl(h, ioctlDiskGetDriveGeometry, nil, 0,
|
||||||
|
&geom[0], uint32(len(geom)), &ret, nil); err != nil {
|
||||||
|
return 512
|
||||||
|
}
|
||||||
|
bps := binary.LittleEndian.Uint32(geom[20:24])
|
||||||
|
if bps == 0 {
|
||||||
|
return 512
|
||||||
|
}
|
||||||
|
return int64(bps)
|
||||||
|
}
|
||||||
|
|
||||||
|
// deviceAlignment reports the sector alignment a raw disk handle requires;
|
||||||
|
// regular files need none.
|
||||||
|
func deviceAlignment(path string) int64 {
|
||||||
|
if isRawDevice(path) {
|
||||||
|
return deviceSectorSize(path)
|
||||||
|
}
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
|
||||||
|
// canElevate is false on Windows: there is no sudo. Raw-disk access requires
|
||||||
|
// running clonetool as Administrator, which the operator arranges outside the tool.
|
||||||
|
func canElevate() bool { return false }
|
||||||
122
listen.go
Normal file
122
listen.go
Normal file
@ -0,0 +1,122 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"crypto/tls"
|
||||||
|
"flag"
|
||||||
|
"fmt"
|
||||||
|
"net"
|
||||||
|
"os"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// cmdListen runs clonetool in listen mode: it waits for inbound TLS
|
||||||
|
// connections (see nettransport.go), authenticates each with the shared
|
||||||
|
// password, and serves whatever role the peer requests (control for the
|
||||||
|
// manager's stat calls, source-stream/sink for the data transfer). This is the
|
||||||
|
// no-ssh path for a host — typically a Windows source — that this machine and
|
||||||
|
// the destination connect *to* instead of being reached over ssh.
|
||||||
|
func cmdListen(args []string) error {
|
||||||
|
fs := flag.NewFlagSet("listen", flag.ContinueOnError)
|
||||||
|
bind := fs.String("bind", ":9000", "address to listen on, e.g. :9000 or 0.0.0.0:9000")
|
||||||
|
password := fs.String("password", "", "shared password (or set CLONETOOL_PASSWORD, or --password-file)")
|
||||||
|
passwordFile := fs.String("password-file", "", "read the shared password from this file")
|
||||||
|
if err := fs.Parse(args); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
pw, err := resolvePassword(*password, *passwordFile)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if pw == "" {
|
||||||
|
return fmt.Errorf("listen: a password is required (--password, --password-file, or CLONETOOL_PASSWORD)")
|
||||||
|
}
|
||||||
|
|
||||||
|
tlsCfg, err := newServerTLSConfig()
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("listen: tls setup: %w", err)
|
||||||
|
}
|
||||||
|
ln, err := net.Listen("tcp", *bind)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("listen on %s: %w", *bind, err)
|
||||||
|
}
|
||||||
|
defer ln.Close()
|
||||||
|
fmt.Fprintf(os.Stderr, "clonetool listening on %s (%s); press Ctrl-C to stop\n", ln.Addr(), versionString())
|
||||||
|
|
||||||
|
for {
|
||||||
|
raw, err := ln.Accept()
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("accept: %w", err)
|
||||||
|
}
|
||||||
|
go serveListenConn(raw, tlsCfg, pw)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// serveListenConn handles one accepted connection: TLS + password auth, then a
|
||||||
|
// single role request that it dispatches. Each connection is one role; the
|
||||||
|
// manager and the destination open separate connections.
|
||||||
|
func serveListenConn(raw net.Conn, tlsCfg *tls.Config, password string) {
|
||||||
|
remote := raw.RemoteAddr().String()
|
||||||
|
defer raw.Close()
|
||||||
|
|
||||||
|
conn := tls.Server(raw, tlsCfg)
|
||||||
|
const handshakeTimeout = 30 * time.Second
|
||||||
|
_ = conn.SetDeadline(time.Now().Add(handshakeTimeout))
|
||||||
|
if err := conn.Handshake(); err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: tls handshake: %v\n", remote, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
fr, fw, err := serverAuth(conn, password, handshakeTimeout)
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: %v\n", remote, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
_ = conn.SetDeadline(time.Time{}) // transfers set their own pacing
|
||||||
|
|
||||||
|
req, err := readCtrlFrame(fr)
|
||||||
|
if err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: read role request: %v\n", remote, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if req.Type != msgRole {
|
||||||
|
_ = fw.WriteJSON(CtrlMsg{Type: msgRoleErr, Message: fmt.Sprintf("expected a role request, got %q", req.Type)})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := dispatchRole(req, fr, fw, remote); err != nil {
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: role %s: %v\n", remote, req.Role, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func dispatchRole(req CtrlMsg, fr *FrameReader, fw *FrameWriter, remote string) error {
|
||||||
|
switch req.Role {
|
||||||
|
case roleControl:
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: control session\n", remote)
|
||||||
|
return runControlAgent(fr, fw)
|
||||||
|
case roleSourceStream:
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: source-stream %s\n", remote, req.Path)
|
||||||
|
return runSourceStreamRole(fr, fw, req.Path, req.Base, req.Size, req.BlockSize)
|
||||||
|
case roleSink:
|
||||||
|
fmt.Fprintf(os.Stderr, "listen: %s: sink %s\n", remote, req.Path)
|
||||||
|
return runSinkRole(fr, fw, req.Path, req.Base, req.Size, req.BlockSize)
|
||||||
|
default:
|
||||||
|
_ = fw.WriteJSON(CtrlMsg{Type: msgRoleErr, Message: fmt.Sprintf("unknown role %q", req.Role)})
|
||||||
|
return fmt.Errorf("unknown role %q", req.Role)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// resolvePassword picks the shared password from, in order: --password-file, an
|
||||||
|
// inline --password, then the CLONETOOL_PASSWORD environment variable. A
|
||||||
|
// trailing newline in a password file is stripped.
|
||||||
|
func resolvePassword(inline, file string) (string, error) {
|
||||||
|
if file != "" {
|
||||||
|
b, err := os.ReadFile(file)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("read password file: %w", err)
|
||||||
|
}
|
||||||
|
return strings.TrimRight(string(b), "\r\n"), nil
|
||||||
|
}
|
||||||
|
if inline != "" {
|
||||||
|
return inline, nil
|
||||||
|
}
|
||||||
|
return os.Getenv("CLONETOOL_PASSWORD"), nil
|
||||||
|
}
|
||||||
28
main.go
28
main.go
@ -16,6 +16,8 @@ func main() {
|
|||||||
switch os.Args[1] {
|
switch os.Args[1] {
|
||||||
case "sync":
|
case "sync":
|
||||||
err = cmdSync(os.Args[2:])
|
err = cmdSync(os.Args[2:])
|
||||||
|
case "listen":
|
||||||
|
err = cmdListen(os.Args[2:])
|
||||||
case "agent":
|
case "agent":
|
||||||
err = cmdAgent(os.Args[2:])
|
err = cmdAgent(os.Args[2:])
|
||||||
case "version", "--version":
|
case "version", "--version":
|
||||||
@ -39,13 +41,19 @@ func usage() {
|
|||||||
|
|
||||||
Usage:
|
Usage:
|
||||||
clonetool sync --source LOC --dest LOC [options]
|
clonetool sync --source LOC --dest LOC [options]
|
||||||
|
clonetool listen --bind ADDR --password SECRET
|
||||||
clonetool version
|
clonetool version
|
||||||
clonetool agent --role {control|sink|source-stream} ... (internal, spawned automatically)
|
clonetool agent --role {control|sink|source-stream} ... (internal, spawned automatically)
|
||||||
|
|
||||||
LOC is either a local path, or [user@]host:path for a path reached over SSH.
|
LOC is one of:
|
||||||
|
a local path (e.g. /dev/sda, ./image.bin, C:\data\disk.img, \\.\PhysicalDrive0)
|
||||||
|
[user@]host:path a path reached over SSH
|
||||||
|
tcp://host:port/path a clonetool started with "listen" on that host (no SSH needed)
|
||||||
|
|
||||||
Source, destination, and the machine running "sync" (the manager) may all be
|
Source, destination, and the machine running "sync" (the manager) may all be
|
||||||
different machines: the manager only orchestrates, it never reads or writes
|
different machines: the manager only orchestrates, it never reads or writes
|
||||||
a single block itself.
|
a single block itself. A tcp:// source is a host running "clonetool listen"
|
||||||
|
(e.g. a Windows box with no SSH server); the destination pulls from it directly.
|
||||||
|
|
||||||
clonetool keeps no state between runs: every sync re-reads and re-hashes both
|
clonetool keeps no state between runs: every sync re-reads and re-hashes both
|
||||||
the source and the destination and transfers only the blocks that differ.
|
the source and the destination and transfers only the blocks that differ.
|
||||||
@ -64,6 +72,14 @@ Options for sync:
|
|||||||
--remote-bin PATH path to clonetool on remote hosts (default "clonetool")
|
--remote-bin PATH path to clonetool on remote hosts (default "clonetool")
|
||||||
--manager-host HOST address peers should use to reach this machine, when
|
--manager-host HOST address peers should use to reach this machine, when
|
||||||
source or dest has no host part (defaults to the local hostname)
|
source or dest has no host part (defaults to the local hostname)
|
||||||
|
--password SECRET shared password for a tcp:// listen endpoint
|
||||||
|
--password-file PATH read the listen-endpoint password from a file
|
||||||
|
(or set CLONETOOL_PASSWORD)
|
||||||
|
|
||||||
|
Options for listen:
|
||||||
|
--bind ADDR address to listen on (default ":9000")
|
||||||
|
--password SECRET shared password (or --password-file, or CLONETOOL_PASSWORD)
|
||||||
|
--password-file PATH read the shared password from a file
|
||||||
`)
|
`)
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -88,12 +104,19 @@ func cmdSync(args []string) error {
|
|||||||
sshBin := fs.String("ssh", "ssh", "ssh binary")
|
sshBin := fs.String("ssh", "ssh", "ssh binary")
|
||||||
remoteBin := fs.String("remote-bin", "clonetool", "clonetool path on remote hosts")
|
remoteBin := fs.String("remote-bin", "clonetool", "clonetool path on remote hosts")
|
||||||
managerHost := fs.String("manager-host", "", "address peers use to reach this machine")
|
managerHost := fs.String("manager-host", "", "address peers use to reach this machine")
|
||||||
|
password := fs.String("password", "", "shared password for a tcp:// listen endpoint")
|
||||||
|
passwordFile := fs.String("password-file", "", "read the listen-endpoint password from a file")
|
||||||
var sshOpts stringSlice
|
var sshOpts stringSlice
|
||||||
fs.Var(&sshOpts, "ssh-opt", `extra "-o OPT" passed to ssh (repeatable)`)
|
fs.Var(&sshOpts, "ssh-opt", `extra "-o OPT" passed to ssh (repeatable)`)
|
||||||
if err := fs.Parse(args); err != nil {
|
if err := fs.Parse(args); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pw, err := resolvePassword(*password, *passwordFile)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
if *source == "" || *dest == "" {
|
if *source == "" || *dest == "" {
|
||||||
fs.Usage()
|
fs.Usage()
|
||||||
return fmt.Errorf("--source and --dest are required")
|
return fmt.Errorf("--source and --dest are required")
|
||||||
@ -112,6 +135,7 @@ func cmdSync(args []string) error {
|
|||||||
Job: *job, Source: *source, Dest: *dest, BlockSize: blockSize,
|
Job: *job, Source: *source, Dest: *dest, BlockSize: blockSize,
|
||||||
Yes: *yes, Sudo: *sudoMode, Deploy: *deploy, ConnectTimeoutSec: *connectTimeout,
|
Yes: *yes, Sudo: *sudoMode, Deploy: *deploy, ConnectTimeoutSec: *connectTimeout,
|
||||||
SSHBin: *sshBin, SSHOpts: sshOpts, RemoteBin: *remoteBin, ManagerHost: *managerHost,
|
SSHBin: *sshBin, SSHOpts: sshOpts, RemoteBin: *remoteBin, ManagerHost: *managerHost,
|
||||||
|
Password: pw,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
129
manager.go
129
manager.go
@ -28,6 +28,9 @@ type SyncConfig struct {
|
|||||||
SSHOpts []string
|
SSHOpts []string
|
||||||
RemoteBin string
|
RemoteBin string
|
||||||
ManagerHost string
|
ManagerHost string
|
||||||
|
// Password authenticates connections to a listen ("tcp://…") endpoint.
|
||||||
|
// Required when the source or dest is such an endpoint; unused for ssh/local.
|
||||||
|
Password string
|
||||||
}
|
}
|
||||||
|
|
||||||
func runSync(cfg SyncConfig) error {
|
func runSync(cfg SyncConfig) error {
|
||||||
@ -42,18 +45,22 @@ func runSync(cfg SyncConfig) error {
|
|||||||
if err := checkNotSame(srcSpec, dstSpec); err != nil {
|
if err := checkNotSame(srcSpec, dstSpec); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
if (srcSpec.IsNet() || dstSpec.IsNet()) && cfg.Password == "" {
|
||||||
|
return fmt.Errorf("a tcp:// listen endpoint needs a password (--password, --password-file, or CLONETOOL_PASSWORD)")
|
||||||
|
}
|
||||||
|
|
||||||
// Make sure each remote endpoint has a runnable clonetool, copying this
|
// Make sure each ssh endpoint has a runnable clonetool, copying this binary
|
||||||
// binary over if not (unless --deploy=false). A host that appears on
|
// over if not (unless --deploy=false). A host that appears on both sides is
|
||||||
// both sides is only probed once.
|
// only probed once. Listen ("tcp://…") endpoints are skipped: the user
|
||||||
|
// started clonetool there manually, so there is nothing to deploy.
|
||||||
srcRemoteBin, dstRemoteBin := cfg.RemoteBin, cfg.RemoteBin
|
srcRemoteBin, dstRemoteBin := cfg.RemoteBin, cfg.RemoteBin
|
||||||
if !srcSpec.IsLocal() {
|
if !srcSpec.IsLocal() && !srcSpec.IsNet() {
|
||||||
if srcRemoteBin, err = resolveRemoteBin(&cfg, srcSpec, "source"); err != nil {
|
if srcRemoteBin, err = resolveRemoteBin(&cfg, srcSpec, "source"); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !dstSpec.IsLocal() {
|
if !dstSpec.IsLocal() && !dstSpec.IsNet() {
|
||||||
if !srcSpec.IsLocal() && sameHost(srcSpec, dstSpec) {
|
if !srcSpec.IsLocal() && !srcSpec.IsNet() && sameHost(srcSpec, dstSpec) {
|
||||||
dstRemoteBin = srcRemoteBin
|
dstRemoteBin = srcRemoteBin
|
||||||
} else if dstRemoteBin, err = resolveRemoteBin(&cfg, dstSpec, "dest"); err != nil {
|
} else if dstRemoteBin, err = resolveRemoteBin(&cfg, dstSpec, "dest"); err != nil {
|
||||||
return err
|
return err
|
||||||
@ -90,42 +97,49 @@ func runSync(cfg SyncConfig) error {
|
|||||||
cb := transferCallbacks{onProgress: pp.print}
|
cb := transferCallbacks{onProgress: pp.print}
|
||||||
|
|
||||||
bothLocal := srcSpec.IsLocal() && dstSpec.IsLocal()
|
bothLocal := srcSpec.IsLocal() && dstSpec.IsLocal()
|
||||||
srcHost, srcUser := resolveConnectHost(srcSpec, &cfg)
|
|
||||||
dstHost, dstUser := resolveConnectHost(dstSpec, &cfg)
|
|
||||||
|
|
||||||
pushReq := CtrlMsg{
|
pushReq := CtrlMsg{Path: srcSpec.Path, Size: targetSize, BlockSize: cfg.BlockSize}
|
||||||
Path: srcSpec.Path, Size: targetSize, BlockSize: cfg.BlockSize,
|
fillPeer(&pushReq, dstSpec, &cfg, dstRemoteBin, bothLocal, dstSudo)
|
||||||
PeerHost: dstHost, PeerUser: dstUser, PeerPath: dstSpec.Path, PeerLocal: bothLocal,
|
pullReq := CtrlMsg{Path: dstSpec.Path, Size: targetSize, BlockSize: cfg.BlockSize}
|
||||||
RemoteBin: dstRemoteBin, SSHBin: cfg.SSHBin, SSHOpts: cfg.SSHOpts, ConnectTimeoutSec: cfg.ConnectTimeoutSec,
|
fillPeer(&pullReq, srcSpec, &cfg, srcRemoteBin, bothLocal, srcSudo)
|
||||||
Sudo: dstSudo,
|
|
||||||
}
|
if srcSpec.IsNet() {
|
||||||
if bothLocal {
|
// A listening source only accepts connections; it can't be driven to
|
||||||
fmt.Fprintf(os.Stderr, "both source and dest are local; syncing directly (no ssh) ...\n")
|
// dial out for a push. The destination pulls from it directly.
|
||||||
} else {
|
fmt.Fprintf(os.Stderr, "source is a listener; pulling %s <- %s ...\n", dstSpec, srcSpec)
|
||||||
fmt.Fprintf(os.Stderr, "attempting push %s -> %s ...\n", srcSpec, dstSpec)
|
ok, reason, err := dstCtrl.ConnectPull(pullReq, cb)
|
||||||
}
|
if err != nil {
|
||||||
ok, reason, err := srcCtrl.ConnectPush(pushReq, cb)
|
return err
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if !ok {
|
|
||||||
pp.finish()
|
|
||||||
fmt.Fprintf(os.Stderr, "push not possible (%s); trying pull %s <- %s ...\n", reason, dstSpec, srcSpec)
|
|
||||||
pullReq := CtrlMsg{
|
|
||||||
Path: dstSpec.Path, Size: targetSize, BlockSize: cfg.BlockSize,
|
|
||||||
PeerHost: srcHost, PeerUser: srcUser, PeerPath: srcSpec.Path, PeerLocal: bothLocal,
|
|
||||||
RemoteBin: srcRemoteBin, SSHBin: cfg.SSHBin, SSHOpts: cfg.SSHOpts, ConnectTimeoutSec: cfg.ConnectTimeoutSec,
|
|
||||||
Sudo: srcSudo,
|
|
||||||
}
|
}
|
||||||
ok2, reason2, err2 := dstCtrl.ConnectPull(pullReq, cb)
|
if !ok {
|
||||||
if err2 != nil {
|
pp.finish()
|
||||||
return err2
|
|
||||||
}
|
|
||||||
if !ok2 {
|
|
||||||
return fmt.Errorf(
|
return fmt.Errorf(
|
||||||
"could not establish a direct connection in either direction (push: %s; pull: %s); "+
|
"could not pull from the listening source (%s); check that the destination host can reach %s:%d and that --password matches the listener",
|
||||||
"run the manager on the source or destination host, or set up SSH connectivity in at least one direction",
|
reason, srcSpec.Host, srcSpec.Port)
|
||||||
reason, reason2)
|
}
|
||||||
|
} else {
|
||||||
|
if bothLocal {
|
||||||
|
fmt.Fprintf(os.Stderr, "both source and dest are local; syncing directly (no ssh) ...\n")
|
||||||
|
} else {
|
||||||
|
fmt.Fprintf(os.Stderr, "attempting push %s -> %s ...\n", srcSpec, dstSpec)
|
||||||
|
}
|
||||||
|
ok, reason, err := srcCtrl.ConnectPush(pushReq, cb)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if !ok {
|
||||||
|
pp.finish()
|
||||||
|
fmt.Fprintf(os.Stderr, "push not possible (%s); trying pull %s <- %s ...\n", reason, dstSpec, srcSpec)
|
||||||
|
ok2, reason2, err2 := dstCtrl.ConnectPull(pullReq, cb)
|
||||||
|
if err2 != nil {
|
||||||
|
return err2
|
||||||
|
}
|
||||||
|
if !ok2 {
|
||||||
|
return fmt.Errorf(
|
||||||
|
"could not establish a direct connection in either direction (push: %s; pull: %s); "+
|
||||||
|
"run the manager on the source or destination host, or set up SSH connectivity in at least one direction",
|
||||||
|
reason, reason2)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -144,6 +158,20 @@ func runSync(cfg SyncConfig) error {
|
|||||||
// restart under sudo. The returned bool reports whether the agent (and any
|
// restart under sudo. The returned bool reports whether the agent (and any
|
||||||
// peer helper it later spawns for this side) is running elevated.
|
// peer helper it later spawns for this side) is running elevated.
|
||||||
func bringUpController(spec Spec, tag string, cfg *SyncConfig, remoteBin, probePath string) (*Controller, PathInfo, bool, error) {
|
func bringUpController(spec Spec, tag string, cfg *SyncConfig, remoteBin, probePath string) (*Controller, PathInfo, bool, error) {
|
||||||
|
if spec.IsNet() {
|
||||||
|
// A listen endpoint is reached over TLS; sudo/self-deploy don't apply.
|
||||||
|
c, err := startNetworkController(spec, tag, cfg.Password, cfg.ConnectTimeoutSec)
|
||||||
|
if err != nil {
|
||||||
|
return nil, PathInfo{}, false, err
|
||||||
|
}
|
||||||
|
info, err := c.Stat(probePath)
|
||||||
|
if err != nil {
|
||||||
|
c.Close()
|
||||||
|
return nil, PathInfo{}, false, err
|
||||||
|
}
|
||||||
|
return c, info, false, nil
|
||||||
|
}
|
||||||
|
|
||||||
sudo := cfg.Sudo == "always" && canElevate()
|
sudo := cfg.Sudo == "always" && canElevate()
|
||||||
c, err := startController(spec, tag, cfg, remoteBin, sudo)
|
c, err := startController(spec, tag, cfg, remoteBin, sudo)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@ -184,6 +212,29 @@ func hostLabel(spec Spec) string {
|
|||||||
return spec.Host
|
return spec.Host
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// fillPeer populates the peer-connection fields of a connect_push/connect_pull
|
||||||
|
// request from the peer's Spec: a net (listen) peer is reached by TLS dial with
|
||||||
|
// the shared password; an ssh/local peer keeps the ssh/local fields.
|
||||||
|
func fillPeer(req *CtrlMsg, peer Spec, cfg *SyncConfig, peerRemoteBin string, bothLocal, sudo bool) {
|
||||||
|
req.PeerPath = peer.Path
|
||||||
|
req.PeerLocal = bothLocal
|
||||||
|
req.RemoteBin = peerRemoteBin
|
||||||
|
req.SSHBin = cfg.SSHBin
|
||||||
|
req.SSHOpts = cfg.SSHOpts
|
||||||
|
req.ConnectTimeoutSec = cfg.ConnectTimeoutSec
|
||||||
|
req.Sudo = sudo
|
||||||
|
if peer.IsNet() {
|
||||||
|
req.PeerNet = true
|
||||||
|
req.PeerHost = peer.Host
|
||||||
|
req.PeerPort = peer.Port
|
||||||
|
req.PeerPassword = cfg.Password
|
||||||
|
return
|
||||||
|
}
|
||||||
|
host, user := resolveConnectHost(peer, cfg)
|
||||||
|
req.PeerHost = host
|
||||||
|
req.PeerUser = user
|
||||||
|
}
|
||||||
|
|
||||||
func resolveConnectHost(spec Spec, cfg *SyncConfig) (host, user string) {
|
func resolveConnectHost(spec Spec, cfg *SyncConfig) (host, user string) {
|
||||||
if !spec.IsLocal() {
|
if !spec.IsLocal() {
|
||||||
return spec.Host, spec.User
|
return spec.Host, spec.User
|
||||||
|
|||||||
267
nettransport.go
Normal file
267
nettransport.go
Normal file
@ -0,0 +1,267 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"crypto/ecdsa"
|
||||||
|
"crypto/elliptic"
|
||||||
|
"crypto/hmac"
|
||||||
|
"crypto/pbkdf2"
|
||||||
|
"crypto/rand"
|
||||||
|
"crypto/sha256"
|
||||||
|
"crypto/tls"
|
||||||
|
"crypto/x509"
|
||||||
|
"crypto/x509/pkix"
|
||||||
|
"encoding/hex"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"math/big"
|
||||||
|
"net"
|
||||||
|
"strconv"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// The network ("listen") transport is an alternative to ssh for reaching an
|
||||||
|
// endpoint that has no ssh server — typically a Windows source started with
|
||||||
|
// `clonetool listen`. It is:
|
||||||
|
//
|
||||||
|
// - TLS 1.3 for confidentiality. The listener uses a fresh, self-signed,
|
||||||
|
// in-memory certificate generated at startup; the dialer does not verify
|
||||||
|
// it (there is no CA). Authentication is by shared password, below.
|
||||||
|
// - A password challenge-response *bound to the TLS session*. Both sides
|
||||||
|
// derive a key from the password with PBKDF2 and prove knowledge of it by
|
||||||
|
// HMAC'ing the connection's RFC 5705 exported keying material, which is
|
||||||
|
// unique to this specific TLS session. A man-in-the-middle terminating TLS
|
||||||
|
// sees a different exporter value on each leg, so it cannot produce a valid
|
||||||
|
// HMAC for either side even though the certificate is unverified. Both
|
||||||
|
// directions are checked, so each end authenticates the other.
|
||||||
|
//
|
||||||
|
// Once authenticated, the connection carries the ordinary frame protocol
|
||||||
|
// (frame.go): first a msgRole request naming the endpoint to run, then that
|
||||||
|
// role's traffic.
|
||||||
|
|
||||||
|
const (
|
||||||
|
authExporterLabel = "EXPORTER-clonetool-auth-v1"
|
||||||
|
authKDFIterations = 200_000
|
||||||
|
authKeyLen = 32
|
||||||
|
authExporterLen = 32
|
||||||
|
authMsgType = "auth"
|
||||||
|
)
|
||||||
|
|
||||||
|
// authKDFSalt is a fixed application salt. Per-connection uniqueness comes from
|
||||||
|
// the TLS exporter that the HMAC is computed over, not from this salt; it only
|
||||||
|
// domain-separates the derived key.
|
||||||
|
var authKDFSalt = []byte("clonetool-listen-auth-salt-v1")
|
||||||
|
|
||||||
|
func deriveAuthKey(password string) ([]byte, error) {
|
||||||
|
return pbkdf2.Key(sha256.New, password, authKDFSalt, authKDFIterations, authKeyLen)
|
||||||
|
}
|
||||||
|
|
||||||
|
// authTag is HMAC(key, side || exporter): "side" separates the client's proof
|
||||||
|
// from the server's so the same value can't simply be reflected.
|
||||||
|
func authTag(key []byte, side string, exporter []byte) []byte {
|
||||||
|
m := hmac.New(sha256.New, key)
|
||||||
|
m.Write([]byte(side))
|
||||||
|
m.Write(exporter)
|
||||||
|
return m.Sum(nil)
|
||||||
|
}
|
||||||
|
|
||||||
|
func exporterFor(cs tls.ConnectionState) ([]byte, error) {
|
||||||
|
return cs.ExportKeyingMaterial(authExporterLabel, nil, authExporterLen)
|
||||||
|
}
|
||||||
|
|
||||||
|
// newServerTLSConfig builds a TLS config with a fresh self-signed cert.
|
||||||
|
func newServerTLSConfig() (*tls.Config, error) {
|
||||||
|
cert, err := selfSignedCert()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return &tls.Config{
|
||||||
|
Certificates: []tls.Certificate{cert},
|
||||||
|
MinVersion: tls.VersionTLS13,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func selfSignedCert() (tls.Certificate, error) {
|
||||||
|
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
|
||||||
|
if err != nil {
|
||||||
|
return tls.Certificate{}, err
|
||||||
|
}
|
||||||
|
serial, err := rand.Int(rand.Reader, new(big.Int).Lsh(big.NewInt(1), 128))
|
||||||
|
if err != nil {
|
||||||
|
return tls.Certificate{}, err
|
||||||
|
}
|
||||||
|
tmpl := x509.Certificate{
|
||||||
|
SerialNumber: serial,
|
||||||
|
Subject: pkix.Name{CommonName: "clonetool"},
|
||||||
|
NotBefore: time.Now().Add(-time.Hour),
|
||||||
|
NotAfter: time.Now().Add(24 * 365 * time.Hour),
|
||||||
|
KeyUsage: x509.KeyUsageDigitalSignature | x509.KeyUsageKeyEncipherment,
|
||||||
|
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth},
|
||||||
|
BasicConstraintsValid: true,
|
||||||
|
}
|
||||||
|
der, err := x509.CreateCertificate(rand.Reader, &tmpl, &tmpl, &key.PublicKey, key)
|
||||||
|
if err != nil {
|
||||||
|
return tls.Certificate{}, err
|
||||||
|
}
|
||||||
|
return tls.Certificate{Certificate: [][]byte{der}, PrivateKey: key}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// dialAuth dials addr, completes the TLS handshake and the mutual password
|
||||||
|
// challenge, and returns the authenticated connection plus its frame
|
||||||
|
// reader/writer. On any failure it closes the connection.
|
||||||
|
func dialAuth(addr, password string, timeout time.Duration) (*tls.Conn, *FrameReader, *FrameWriter, error) {
|
||||||
|
d := net.Dialer{Timeout: timeout}
|
||||||
|
raw, err := d.Dial("tcp", addr)
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, nil, err
|
||||||
|
}
|
||||||
|
conn := tls.Client(raw, &tls.Config{InsecureSkipVerify: true, MinVersion: tls.VersionTLS13})
|
||||||
|
_ = conn.SetDeadline(time.Now().Add(timeout))
|
||||||
|
if err := conn.Handshake(); err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, nil, nil, fmt.Errorf("tls handshake: %w", err)
|
||||||
|
}
|
||||||
|
key, err := deriveAuthKey(password)
|
||||||
|
if err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, nil, nil, err
|
||||||
|
}
|
||||||
|
exporter, err := exporterFor(conn.ConnectionState())
|
||||||
|
if err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, nil, nil, fmt.Errorf("tls keying material: %w", err)
|
||||||
|
}
|
||||||
|
fw := NewFrameWriter(conn)
|
||||||
|
fr := NewFrameReader(conn)
|
||||||
|
|
||||||
|
// Prove ourselves to the server, then verify the server back.
|
||||||
|
if err := fw.WriteJSON(CtrlMsg{Type: authMsgType, Message: hex.EncodeToString(authTag(key, "client", exporter))}); err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, nil, nil, err
|
||||||
|
}
|
||||||
|
m, err := readCtrlFrame(fr)
|
||||||
|
if err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, nil, nil, fmt.Errorf("authentication failed (bad password?): %w", err)
|
||||||
|
}
|
||||||
|
want := authTag(key, "server", exporter)
|
||||||
|
got, _ := hex.DecodeString(m.Message)
|
||||||
|
if m.Type != authMsgType || !hmac.Equal(got, want) {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, nil, nil, fmt.Errorf("server authentication failed (password mismatch)")
|
||||||
|
}
|
||||||
|
_ = conn.SetDeadline(time.Time{}) // clear; per-op timeouts are handled by callers
|
||||||
|
return conn, fr, fw, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// serverAuth runs the server side of the password challenge on an
|
||||||
|
// already-TLS-handshaken connection.
|
||||||
|
func serverAuth(conn *tls.Conn, password string, timeout time.Duration) (*FrameReader, *FrameWriter, error) {
|
||||||
|
key, err := deriveAuthKey(password)
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
exporter, err := exporterFor(conn.ConnectionState())
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, fmt.Errorf("tls keying material: %w", err)
|
||||||
|
}
|
||||||
|
fw := NewFrameWriter(conn)
|
||||||
|
fr := NewFrameReader(conn)
|
||||||
|
|
||||||
|
_ = conn.SetReadDeadline(time.Now().Add(timeout))
|
||||||
|
m, err := readCtrlFrame(fr)
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
want := authTag(key, "client", exporter)
|
||||||
|
got, _ := hex.DecodeString(m.Message)
|
||||||
|
if m.Type != authMsgType || !hmac.Equal(got, want) {
|
||||||
|
return nil, nil, fmt.Errorf("client authentication failed (password mismatch)")
|
||||||
|
}
|
||||||
|
if err := fw.WriteJSON(CtrlMsg{Type: authMsgType, Message: hex.EncodeToString(authTag(key, "server", exporter))}); err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
_ = conn.SetReadDeadline(time.Time{})
|
||||||
|
return fr, fw, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// readCtrlFrame reads one frameCtrlJSON frame and decodes it into a CtrlMsg.
|
||||||
|
func readCtrlFrame(fr *FrameReader) (CtrlMsg, error) {
|
||||||
|
typ, payload, err := fr.ReadFrame()
|
||||||
|
if err != nil {
|
||||||
|
return CtrlMsg{}, err
|
||||||
|
}
|
||||||
|
if typ != frameCtrlJSON {
|
||||||
|
return CtrlMsg{}, fmt.Errorf("expected control frame, got type %d", typ)
|
||||||
|
}
|
||||||
|
var m CtrlMsg
|
||||||
|
if err := json.Unmarshal(payload, &m); err != nil {
|
||||||
|
return CtrlMsg{}, err
|
||||||
|
}
|
||||||
|
return m, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// dialPeerLink connects to a listening clonetool over TLS, authenticates,
|
||||||
|
// requests role for the given window, and waits for its READY handshake. It is
|
||||||
|
// the net counterpart of spawnPeerLink (see agent.go).
|
||||||
|
func dialPeerLink(req CtrlMsg, role string) (*peerLink, string) {
|
||||||
|
addr := net.JoinHostPort(req.PeerHost, strconv.Itoa(req.PeerPort))
|
||||||
|
timeout := time.Duration(req.ConnectTimeoutSec+2) * time.Second
|
||||||
|
conn, fr, fw, err := dialAuth(addr, req.PeerPassword, timeout)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err.Error()
|
||||||
|
}
|
||||||
|
if err := fw.WriteJSON(CtrlMsg{
|
||||||
|
Type: msgRole, Role: role,
|
||||||
|
Path: req.PeerPath, Base: req.PeerBase, Size: req.Size, BlockSize: req.BlockSize,
|
||||||
|
}); err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, err.Error()
|
||||||
|
}
|
||||||
|
if err := waitPeerReady(fr, timeout); err != nil {
|
||||||
|
_ = conn.Close()
|
||||||
|
return nil, err.Error()
|
||||||
|
}
|
||||||
|
return &peerLink{
|
||||||
|
fw: fw,
|
||||||
|
fr: fr,
|
||||||
|
stderr: func() string { return "" },
|
||||||
|
finish: conn.Close,
|
||||||
|
kill: func() { _ = conn.Close() },
|
||||||
|
closeSend: conn.Close,
|
||||||
|
}, ""
|
||||||
|
}
|
||||||
|
|
||||||
|
// waitPeerReady waits for a data role's frameReady, surfacing a msgRoleErr the
|
||||||
|
// listener may send instead when it can't start the requested role.
|
||||||
|
func waitPeerReady(fr *FrameReader, timeout time.Duration) error {
|
||||||
|
type result struct {
|
||||||
|
typ frameType
|
||||||
|
payload []byte
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
ch := make(chan result, 1)
|
||||||
|
go func() {
|
||||||
|
typ, payload, err := fr.ReadFrame()
|
||||||
|
ch <- result{typ, payload, err}
|
||||||
|
}()
|
||||||
|
select {
|
||||||
|
case r := <-ch:
|
||||||
|
if r.err != nil {
|
||||||
|
return fmt.Errorf("handshake failed: %w", r.err)
|
||||||
|
}
|
||||||
|
switch r.typ {
|
||||||
|
case frameReady:
|
||||||
|
return nil
|
||||||
|
case frameCtrlJSON:
|
||||||
|
var m CtrlMsg
|
||||||
|
if json.Unmarshal(r.payload, &m) == nil && m.Type == msgRoleErr {
|
||||||
|
return fmt.Errorf("remote refused role: %s", m.Message)
|
||||||
|
}
|
||||||
|
return fmt.Errorf("handshake failed: unexpected control frame")
|
||||||
|
default:
|
||||||
|
return fmt.Errorf("handshake failed: unexpected frame type %d", r.typ)
|
||||||
|
}
|
||||||
|
case <-time.After(timeout):
|
||||||
|
return fmt.Errorf("handshake timed out after %s", timeout)
|
||||||
|
}
|
||||||
|
}
|
||||||
83
spec.go
83
spec.go
@ -3,21 +3,32 @@ package main
|
|||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Spec is a parsed source/dest location: [user@]host:path, or a bare local
|
// Spec is a parsed source/dest location. It is one of:
|
||||||
// path (Host == "").
|
// - a bare local path (Host == "", Net == false)
|
||||||
|
// - an ssh endpoint [user@]host:path (Host != "", Net == false)
|
||||||
|
// - a listen endpoint tcp://host:port/path (Net == true) — a clonetool
|
||||||
|
// started in "listen" mode that this side connects to over TLS instead
|
||||||
|
// of ssh. Only the source is expected to be a listener today.
|
||||||
type Spec struct {
|
type Spec struct {
|
||||||
Raw string
|
Raw string
|
||||||
User string
|
User string
|
||||||
Host string
|
Host string
|
||||||
Path string
|
Path string
|
||||||
|
Net bool // tcp:// listen endpoint (connect to a listening clonetool)
|
||||||
|
Port int // listener port, when Net
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s Spec) IsLocal() bool { return s.Host == "" }
|
func (s Spec) IsLocal() bool { return s.Host == "" && !s.Net }
|
||||||
|
func (s Spec) IsNet() bool { return s.Net }
|
||||||
|
|
||||||
func (s Spec) String() string {
|
func (s Spec) String() string {
|
||||||
|
if s.Net {
|
||||||
|
return fmt.Sprintf("tcp://%s:%d/%s", s.Host, s.Port, s.Path)
|
||||||
|
}
|
||||||
if s.IsLocal() {
|
if s.IsLocal() {
|
||||||
return s.Path
|
return s.Path
|
||||||
}
|
}
|
||||||
@ -27,14 +38,38 @@ func (s Spec) String() string {
|
|||||||
return fmt.Sprintf("%s:%s", s.Host, s.Path)
|
return fmt.Sprintf("%s:%s", s.Host, s.Path)
|
||||||
}
|
}
|
||||||
|
|
||||||
// parseSpec parses "[user@]host:path" or a local "path". A leading "/",
|
// looksLikeLocalPath reports whether raw is an ordinary local path that must
|
||||||
// "./" or "../", or the absence of any colon, is treated as a local path so
|
// never be mistaken for a "host:path" spec: a unix absolute/relative path, a
|
||||||
// that ordinary absolute/relative paths are never mistaken for a host spec.
|
// Windows drive path (C:\... or C:/...), or a Windows UNC / raw-device path
|
||||||
|
// (\\server\share, \\.\PhysicalDrive0).
|
||||||
|
func looksLikeLocalPath(raw string) bool {
|
||||||
|
if strings.HasPrefix(raw, "/") || strings.HasPrefix(raw, "./") || strings.HasPrefix(raw, "../") {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
if strings.HasPrefix(raw, `\\`) { // UNC or \\.\ device
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
// Windows drive path: a single letter, a colon, then a separator.
|
||||||
|
if len(raw) >= 3 && raw[1] == ':' &&
|
||||||
|
((raw[0] >= 'A' && raw[0] <= 'Z') || (raw[0] >= 'a' && raw[0] <= 'z')) &&
|
||||||
|
(raw[2] == '\\' || raw[2] == '/') {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// parseSpec parses a location: "tcp://host:port/path" (a listen endpoint),
|
||||||
|
// "[user@]host:path" (ssh), or a local path. Local paths (unix or Windows,
|
||||||
|
// including drive letters and \\.\ devices) and anything without a colon are
|
||||||
|
// treated as local so an ordinary path is never mistaken for a host spec.
|
||||||
func parseSpec(raw string) (Spec, error) {
|
func parseSpec(raw string) (Spec, error) {
|
||||||
if raw == "" {
|
if raw == "" {
|
||||||
return Spec{}, fmt.Errorf("empty location")
|
return Spec{}, fmt.Errorf("empty location")
|
||||||
}
|
}
|
||||||
if strings.HasPrefix(raw, "/") || strings.HasPrefix(raw, "./") || strings.HasPrefix(raw, "../") || !strings.Contains(raw, ":") {
|
if strings.HasPrefix(raw, "tcp://") {
|
||||||
|
return parseNetSpec(raw)
|
||||||
|
}
|
||||||
|
if looksLikeLocalPath(raw) || !strings.Contains(raw, ":") {
|
||||||
return Spec{Raw: raw, Path: raw}, nil
|
return Spec{Raw: raw, Path: raw}, nil
|
||||||
}
|
}
|
||||||
idx := strings.Index(raw, ":")
|
idx := strings.Index(raw, ":")
|
||||||
@ -55,15 +90,47 @@ func parseSpec(raw string) (Spec, error) {
|
|||||||
return Spec{Raw: raw, User: user, Host: host, Path: path}, nil
|
return Spec{Raw: raw, User: user, Host: host, Path: path}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// parseNetSpec parses "tcp://host:port/path". Everything after the first "/"
|
||||||
|
// that follows the host:port authority is taken literally as the path, so a
|
||||||
|
// Windows path (C:\dir\file, \\.\PhysicalDrive0) survives unmangled.
|
||||||
|
func parseNetSpec(raw string) (Spec, error) {
|
||||||
|
rest := strings.TrimPrefix(raw, "tcp://")
|
||||||
|
slash := strings.IndexByte(rest, '/')
|
||||||
|
if slash < 0 {
|
||||||
|
return Spec{}, fmt.Errorf("cannot parse %q (expected tcp://host:port/path)", raw)
|
||||||
|
}
|
||||||
|
authority := rest[:slash]
|
||||||
|
path := rest[slash+1:]
|
||||||
|
if path == "" {
|
||||||
|
return Spec{}, fmt.Errorf("cannot parse %q: empty path", raw)
|
||||||
|
}
|
||||||
|
colon := strings.LastIndexByte(authority, ':')
|
||||||
|
if colon <= 0 || colon == len(authority)-1 {
|
||||||
|
return Spec{}, fmt.Errorf("cannot parse %q (expected tcp://host:port/path)", raw)
|
||||||
|
}
|
||||||
|
host := authority[:colon]
|
||||||
|
port, err := strconv.Atoi(authority[colon+1:])
|
||||||
|
if err != nil || port <= 0 || port > 65535 {
|
||||||
|
return Spec{}, fmt.Errorf("cannot parse %q: invalid port %q", raw, authority[colon+1:])
|
||||||
|
}
|
||||||
|
return Spec{Raw: raw, Host: host, Port: port, Path: path, Net: true}, nil
|
||||||
|
}
|
||||||
|
|
||||||
// checkNotSame does a best-effort local check that source and dest don't
|
// checkNotSame does a best-effort local check that source and dest don't
|
||||||
// refer to the exact same path, to avoid an obviously destructive mistake.
|
// refer to the exact same path, to avoid an obviously destructive mistake.
|
||||||
// It cannot resolve whether two different remote hostnames are actually the
|
// It cannot resolve whether two different remote hostnames are actually the
|
||||||
// same machine.
|
// same machine.
|
||||||
func checkNotSame(src, dst Spec) error {
|
func checkNotSame(src, dst Spec) error {
|
||||||
|
if src.IsNet() != dst.IsNet() {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if src.IsNet() && (!strings.EqualFold(src.Host, dst.Host) || src.Port != dst.Port) {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
if src.IsLocal() != dst.IsLocal() {
|
if src.IsLocal() != dst.IsLocal() {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
if !src.IsLocal() && !strings.EqualFold(src.Host, dst.Host) {
|
if !src.IsLocal() && !src.IsNet() && !strings.EqualFold(src.Host, dst.Host) {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
if !src.IsLocal() && src.User != dst.User {
|
if !src.IsLocal() && src.User != dst.User {
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user