From 6e7aac090cfb2c41b390b672c69d3c581e2ec9c1 Mon Sep 17 00:00:00 2001 From: 28allday Date: Wed, 27 May 2026 21:20:30 +0100 Subject: [PATCH] Stop a cancelled transfer from carrying on into the next one MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When the receiver cancelled a transfer partway through, the sender kept uploading the remaining files of that batch (the loop continued past the failure), and a newly started transfer never stopped the old goroutine — so the old, cancelled transfer appeared to "carry on". Fixes: - Each send now runs under a cancellable context. A new transfer to a peer supersedes (cancels) any still-running send to that same peer. - A network/session error mid-batch aborts the rest of the batch, emitting a clean Cancel for the un-sent files instead of pushing them anyway. A local file-open error still skips only that file. - Bounded connection timeouts (dial 10s, TLS 10s, response-header 30s) so a vanished peer fails fast; the overall timeout stays off for large files. Test: a receiver whose session is gone (returns 403) causes the sender to make exactly one upload attempt and then abort, reporting the rest cancelled. Co-Authored-By: Claude Opus 4.7 (1M context) --- internal/client/cancel_test.go | 94 ++++++++++++++++++++++++++++++++++ internal/client/client.go | 90 +++++++++++++++++++++++++++----- 2 files changed, 172 insertions(+), 12 deletions(-) create mode 100644 internal/client/cancel_test.go diff --git a/internal/client/cancel_test.go b/internal/client/cancel_test.go new file mode 100644 index 0000000..fde9597 --- /dev/null +++ b/internal/client/cancel_test.go @@ -0,0 +1,94 @@ +package client + +import ( + "encoding/json" + "net" + "net/http" + "net/http/httptest" + "net/url" + "os" + "path/filepath" + "strconv" + "sync/atomic" + "testing" + "time" + + "omarchy-send/internal/discovery" + "omarchy-send/internal/protocol" + "omarchy-send/internal/transfer" +) + +// When the receiver cancels mid-batch (its session is gone, so /upload returns +// 403), the sender must stop pushing the remaining files rather than carry the +// old transfer on. Regression test for "a cancelled transfer carries on". +func TestSendAbortsBatchWhenSessionGone(t *testing.T) { + var uploads int32 + + mux := http.NewServeMux() + mux.HandleFunc(protocol.PathPrepareUpload, func(w http.ResponseWriter, r *http.Request) { + var req protocol.PrepareUploadRequest + _ = json.NewDecoder(r.Body).Decode(&req) + tokens := make(map[string]string, len(req.Files)) + for id := range req.Files { + tokens[id] = "tok-" + id + } + _ = json.NewEncoder(w).Encode(protocol.PrepareUploadResponse{SessionID: "s1", Files: tokens}) + }) + mux.HandleFunc(protocol.PathUpload, func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&uploads, 1) + http.Error(w, "forbidden", http.StatusForbidden) // session cancelled on remote + }) + ts := httptest.NewServer(mux) + defer ts.Close() + + u, _ := url.Parse(ts.URL) + srcDir := t.TempDir() + var paths []string + for _, n := range []string{"a.bin", "b.bin", "c.bin"} { + p := filepath.Join(srcDir, n) + if err := os.WriteFile(p, []byte("payload-"+n), 0o644); err != nil { + t.Fatal(err) + } + paths = append(paths, p) + } + + host, portStr, err := net.SplitHostPort(u.Host) + if err != nil { + t.Fatalf("split %q: %v", u.Host, err) + } + port, _ := strconv.Atoi(portStr) + + sender := New(protocol.DeviceInfo{Alias: "snd", Protocol: "http"}) + peer := discovery.Peer{Info: protocol.DeviceInfo{Protocol: "http", Port: port}, IP: host} + sender.Send(peer, paths, "") + + // Drain events until things go quiet. + errs, cancels := 0, 0 + timeout := time.After(3 * time.Second) + for done := false; !done; { + select { + case ev := <-sender.Events(): + switch ev.Kind { + case transfer.Error: + errs++ + case transfer.Cancel: + cancels++ + } + case <-time.After(300 * time.Millisecond): + done = true + case <-timeout: + done = true + } + } + + if got := atomic.LoadInt32(&uploads); got != 1 { + t.Fatalf("expected exactly 1 upload attempt before aborting, got %d", got) + } + // One file errored; the other two should be reported as cancelled, not retried. + if errs != 1 { + t.Errorf("expected 1 error event, got %d", errs) + } + if cancels != 2 { + t.Errorf("expected 2 cancel events for the un-sent files, got %d", cancels) + } +} diff --git a/internal/client/client.go b/internal/client/client.go index 040f965..e5a7203 100644 --- a/internal/client/client.go +++ b/internal/client/client.go @@ -21,6 +21,7 @@ import ( "strconv" "strings" "sync" + "time" "omarchy-send/internal/dbg" "omarchy-send/internal/discovery" @@ -28,12 +29,24 @@ import ( "omarchy-send/internal/transfer" ) +// errOpen wraps a failure to open a source file, so the send loop can skip just +// that file rather than aborting the whole batch (which it does for peer/network +// errors, where the shared session is dead). +var errOpen = errors.New("open source file") + +// inflight tracks one running send so it can be cancelled — e.g. when a newer +// transfer to the same peer supersedes it. +type inflight struct { + cancel context.CancelFunc +} + // Sender uploads files to peers. Events are delivered on Events(). type Sender struct { mu sync.Mutex self protocol.DeviceInfo http *http.Client events chan transfer.Event + active map[string]*inflight // in-flight sends keyed by peer IP } // New returns a Sender advertising self. TLS chain validation is disabled (we @@ -42,10 +55,19 @@ func New(self protocol.DeviceInfo) *Sender { return &Sender{ self: self, http: &http.Client{ - Timeout: 0, // large files: no overall timeout - Transport: &http.Transport{TLSClientConfig: &tls.Config{InsecureSkipVerify: true}}, + // No overall timeout (large files), but bound the parts that can + // silently wedge on a vanished peer: connecting, the TLS handshake, + // and waiting for response headers after the body is sent. + Timeout: 0, + Transport: &http.Transport{ + TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, + DialContext: (&net.Dialer{Timeout: 10 * time.Second}).DialContext, + TLSHandshakeTimeout: 10 * time.Second, + ResponseHeaderTimeout: 30 * time.Second, + }, }, events: make(chan transfer.Event, 256), + active: make(map[string]*inflight), } } @@ -73,6 +95,26 @@ func (s *Sender) Send(peer discovery.Peer, paths []string, pin string) { } func (s *Sender) send(peer discovery.Peer, paths []string, pin string) { + // A new transfer to a peer supersedes any still-running one to the same + // peer: cancel it so a half-finished old batch can't carry on once the user + // starts something new. + ctx, cancel := context.WithCancel(context.Background()) + h := &inflight{cancel: cancel} + s.mu.Lock() + if prev := s.active[peer.IP]; prev != nil { + prev.cancel() + } + s.active[peer.IP] = h + s.mu.Unlock() + defer func() { + cancel() + s.mu.Lock() + if s.active[peer.IP] == h { + delete(s.active, peer.IP) + } + s.mu.Unlock() + }() + // Expand any directories into their files, then build metadata keyed by a // generated fileId. A directory's files carry a relative FileName (e.g. // "Trip/day1/img.jpg") so the receiver can recreate the folder structure. @@ -97,7 +139,7 @@ func (s *Sender) send(peer discovery.Peer, paths []string, pin string) { dbg.Logf("SEND prepare-upload to %s: files=%s", peer.IP, string(meta)) } base := s.url(peer) - prepResp, err := s.prepareUpload(base, files, pin) + prepResp, err := s.prepareUpload(ctx, base, files, pin) if err != nil { dbg.Logf("send prepare-upload to %s failed: %v", peer.IP, err) if errors.Is(err, transfer.ErrPinRequired) { @@ -114,12 +156,31 @@ func (s *Sender) send(peer discovery.Peer, paths []string, pin string) { for id, token := range prepResp.Files { meta := files[id] key := prepResp.SessionID + ":" + id - if err := s.uploadFile(base, prepResp.SessionID, id, token, key, pathByID[id], meta); err != nil { - dbg.Logf("send upload %q failed: %v", meta.FileName, err) - s.emit(transfer.Event{Dir: transfer.Outgoing, Kind: transfer.Error, ID: key, FileName: meta.FileName, Err: err}) + + // If the transfer was cancelled (superseded, or aborted after an earlier + // failure), don't push the rest of the batch — report a clean cancel. + if ctx.Err() != nil { + s.emit(transfer.Event{Dir: transfer.Outgoing, Kind: transfer.Cancel, ID: key, FileName: meta.FileName}) continue } - s.emit(transfer.Event{Dir: transfer.Outgoing, Kind: transfer.FileDone, ID: key, FileName: meta.FileName, Received: meta.Size, Total: meta.Size}) + + err := s.uploadFile(ctx, base, prepResp.SessionID, id, token, key, pathByID[id], meta) + if err == nil { + s.emit(transfer.Event{Dir: transfer.Outgoing, Kind: transfer.FileDone, ID: key, FileName: meta.FileName, Received: meta.Size, Total: meta.Size}) + continue + } + dbg.Logf("send upload %q failed: %v", meta.FileName, err) + if ctx.Err() != nil { + s.emit(transfer.Event{Dir: transfer.Outgoing, Kind: transfer.Cancel, ID: key, FileName: meta.FileName}) + } else { + s.emit(transfer.Event{Dir: transfer.Outgoing, Kind: transfer.Error, ID: key, FileName: meta.FileName, Err: err}) + } + // A failure to open a local file is specific to that file — skip it and + // keep going. Any other failure means the peer/session is gone, so abort + // the rest of the batch (the shared session can't be resumed). + if !errors.Is(err, errOpen) { + cancel() + } } } @@ -174,13 +235,18 @@ func (s *Sender) expand(paths []string) []fileItem { return items } -func (s *Sender) prepareUpload(base string, files map[string]protocol.FileMetadata, pin string) (protocol.PrepareUploadResponse, error) { +func (s *Sender) prepareUpload(ctx context.Context, base string, files map[string]protocol.FileMetadata, pin string) (protocol.PrepareUploadResponse, error) { reqBody, _ := json.Marshal(protocol.PrepareUploadRequest{Info: s.selfCopy(), Files: files}) url := base + protocol.PathPrepareUpload if pin != "" { url += "?pin=" + neturl.QueryEscape(pin) } - resp, err := s.http.Post(url, "application/json", bytes.NewReader(reqBody)) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(reqBody)) + if err != nil { + return protocol.PrepareUploadResponse{}, err + } + req.Header.Set("Content-Type", "application/json") + resp, err := s.http.Do(req) if err != nil { return protocol.PrepareUploadResponse{}, err } @@ -198,10 +264,10 @@ func (s *Sender) prepareUpload(base string, files map[string]protocol.FileMetada return pr, nil } -func (s *Sender) uploadFile(base, sessionID, fileID, token, key, path string, meta protocol.FileMetadata) error { +func (s *Sender) uploadFile(ctx context.Context, base, sessionID, fileID, token, key, path string, meta protocol.FileMetadata) error { f, err := os.Open(path) if err != nil { - return err + return fmt.Errorf("%w: %v", errOpen, err) } defer f.Close() @@ -216,7 +282,7 @@ func (s *Sender) uploadFile(base, sessionID, fileID, token, key, path string, me } url := fmt.Sprintf("%s%s?sessionId=%s&fileId=%s&token=%s", base, protocol.PathUpload, sessionID, fileID, token) - req, err := http.NewRequestWithContext(context.Background(), http.MethodPost, url, pr) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, pr) if err != nil { return err }