diff --git a/AGENTS.md b/AGENTS.md index 4fa85c55..eb9755c6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -679,13 +679,36 @@ says (`--box` picks the boxes whose changes are reported; every box is followed) posting recorded as soon as it is classified. The start is the Date header translated back to when the request was made (mail that lands while the server answers is later than the start), it is taken before the box list, and each -box's cursor starts no later than it (`noLaterThan`): the server bakes the box's last posting -activity into the cursor, so mail that landed in between would otherwise sit behind the -cursor, read by nothing. That +box's cursor starts at it (`watchStartSince`), keeping only the version from the box's +`posting_changes_url`. The since HEY puts there is not its clock but the box's last posting +activity (`Box#last_posting_activity_at`: unbundled postings only, the box's own `updated_at` +when it has none), and the feed answers deletions and bundled postings later than that; and +`/boxes.json` comes through the SDK's ETag cache with an ETag of the box rows alone, which +posting activity does not touch, so a 304 serves the since as it was when the list was +cached. A read from HEY's since reported history as news on every start, and a since later +than the start would leave mail that landed in between behind it, read by nothing. A watch +that cannot read HEY's clock does not start: the workstation's clock is no stand-in for a +cutoff every feed and new mail are measured against, since a fast one would skip changes and +a slow one would report history. That is HEY's semantics and state across events, so the CLI decides it once; what to do about -it is the reader's. A 409 skip-ahead sets that box's floor at the cursor it skipped to +it is the reader's. The Date header is whole seconds, so the start can be up to a second +(plus the clock request's whole time, retries included) early and a change from that window +is reported. Nothing HEY serves on demand says the time finer — Action Cable pings are +whole seconds too; a posting doorbell's `at` and the feeds' cursors carry microseconds, but +only once something has changed, never as a "now" before the watch starts — and rounding +the other way would skip changes. A 409 skip-ahead moves the cursor to HEY's clock when it +answered (`serverNowAnswered` — not taken back by the request's time, since a resync has +no gap to catch and a slow request could leave a busy feed still behind), keeping the feed +version from a list read past the SDK's cache (`newUncachedSDKClient`: the list's ETag is +its rows, which neither posting activity nor a new feed version changes, and HEY answers +409 for a version it no longer speaks), and sets that box's floor there (`newMail.skippedTo`): activity at or before it is never new there, known thread or not, -because the watch never read the gap. `resync` is an event of its own — reported by default, +because the watch never read the gap. The first skip is read from straight away; a 409 +after it is the same recovery (`feedRecovery`), and the next skip waits on the retry +backoff rather than every doorbell. The recovery's one resync goes out with the clean read +that ends it, at the last skip, so a reader that re-reads on it has missed nothing a later +skip passed. A list or clock read that fails is retried on that backoff, and an +interrupt during one ends quietly (`skipFailed`). A calendar's 409 skips the same way. `resync` is an event of its own — reported by default, left out by `--events new` — so a script for new mail never runs on one. The Omarchy bar plugin toasts from those lines itself (app-name, glyph, click-to-focus and the replace-not-stack id all live in the plugin), and nothing desktop-shaped lives in `watch*.go`. @@ -785,11 +808,11 @@ watch that is down costs staleness, not a notice. `hey watch` follows the same streams on its own connection and reports the changes themselves (`internal/cmd/watch_calendar.go`). Rings are coalesced per calendar for `calendarCoalesceDelay`, then the calendar's recording feed is read from its cursor -(`Calendars().AllRecordingChanges`, cursors capped at the watch's start like the boxes' — -`calendarCursorNoLaterThan`) and each recording is a `recording_added`, `recording_updated` +(`Calendars().AllRecordingChanges`, cursors starting at the watch's start like the boxes' — +`calendarCursor`) and each recording is a `recording_added`, `recording_updated` or `recording_deleted` line naming its calendar where a mail line names its box. The poll reports `calendar_added`, `calendar_updated` and `calendar_deleted`, and a recording feed's -409 is `calendar_resync` after skipping ahead to a fresh cursor from the list. The +409 is `calendar_resync` after skipping ahead to HEY's clock, as a box does. The email-specific flags switch all of it off — `--box`, or an `--events` list naming only mail changes (`watchingCalendars` in watch_calendar.go) — and `ready` waits for the calendars' catch-up exactly as it waits for the boxes', on the same retry backoff and the diff --git a/docs/cli.md b/docs/cli.md index 3cdb14eb..63a8e8a7 100644 --- a/docs/cli.md +++ b/docs/cli.md @@ -336,7 +336,14 @@ hey watch --box imbox --events new --run-async 'notify-send -a HEY "New mail in hey watch --run-sync ./triage.sh # one at a time, waiting for each ``` -Runs until interrupted, printing changes as they happen, one line each: +Runs until interrupted, printing changes as they happen, one line each. What changed before +the watch began is not reported unless `--since` reads back to it first, so +`--exit-on-first` waits for a change rather than stopping on an old one. When the watch began +is read off HEY's clock, and a watch that cannot read it exits with an error rather than +guess. That clock is read to the whole second and taken back by however long the request +took (retries included), so a change from up to a second before the watch began, plus that +request's time, can still be reported (and be new, and end `--exit-on-first`); its `at` says +when it happened. Rounding the other way would skip changes that came after: ```json {"change":"added","at":"2026-08-18T09:14:22.031Z","box":{"id":24088,"kind":"imbox","name":"Imbox"},"posting_id":98765,"thread_id":54321,"new":true,"posting":{}} @@ -344,7 +351,7 @@ Runs until interrupted, printing changes as they happen, one line each: Every `added` and `updated` line says whether the posting is new mail: unseen, not muted, and active since the watch last saw the thread — or since the watch began, for a thread it -has not seen, so a box's backlog is never new. Reading, muting or moving a thread is not new +has not seen, so the backlog `--since` reads is never new. Reading, muting or moving a thread is not new activity; a reply on a known thread is. `--events new` selects the new ones, alone or in a union with `added`, `updated`, `deleted` and `resync` — the default is everything but `new`, and `new` alone leaves a `resync` out, so a script for new mail never runs on one. `--box` @@ -359,7 +366,9 @@ calendar — `{"change":"recording_added","calendar":{"id":512,"name":"Household "recording_id":88001,"recording_type":"Calendar::Event","recording":{}}` — and a calendar arriving, changing or leaving is `calendar_added`, `calendar_updated` or `calendar_deleted`. A calendar whose feed fell too far behind is skipped ahead and says so -with `calendar_resync`, the way a box says `resync`. The email-specific flags switch the +with `calendar_resync`, the way a box says `resync`. Either is said once per catch-up, when +the feed can be followed again: a feed still too busy after the skip is skipped again on the +retry backoff, and the line's `at` is the last skip. The email-specific flags switch the calendars off: `--box` scopes the watch to mail, and an `--events` list naming only mail changes does the same. diff --git a/docs/omarchy.md b/docs/omarchy.md index ae892f48..979476f8 100644 --- a/docs/omarchy.md +++ b/docs/omarchy.md @@ -191,15 +191,21 @@ and `--events new` selects the true ones. The rule: watch last recorded for the thread — or later than the watch's start, when it has no record. That start is read off HEY's own clock (the `Date` header of one request, translated back to the moment the request was made), the clock every `active_at` is on, - so a workstation running fast or slow neither calls the backlog new nor sits on new mail; + so a workstation running fast or slow neither calls the backlog new nor sits on new mail + (a watch that cannot read HEY's clock exits with an error rather than use its own); whole seconds, rounded down, so the doubt falls on the side of calling mail a moment old - new, and every box's cursor starts no later than that start, so mail that lands while the + new — mail from up to a second before the watch began, plus however long the clock request + took (retries included), can read as new — and every + box's cursor starts at that start, so mail that lands while the watch is starting up is read and is new. `active_at` moves on new mail only, not on a seen flip, a mute or a move, so reading a thread, marking it unseen again or moving it - into a box is never new, and a reply on a known thread is. A box's first read is its - catch-up from the server's cursor — the box's last activity, not this moment — so it - carries backlog, which the start-time rule keeps out, alongside anything that arrived - while the watch was starting, which is new. + into a box is never new, and a reply on a known thread is. A box's first read starts at + the watch's start, not at the server's cursor — the box's last posting activity, which a + deletion or a bundled posting can postdate, and which a cached box list can serve days + old — so it carries what arrived while the watch was starting, which is new, and at most + the whole-second window above: a change from just before the start that the start could + not be told apart from. `--since` reads backlog first, which the start-time rule keeps + out. - **Every posting the watch reads is recorded**, in every box and whatever `--events` or `--box` reports — `--box` picks what is reported, every box is followed — so a thread known from a filtered-out change, or from another box, is never mistaken for new when its @@ -211,8 +217,8 @@ and `--events new` selects the true ones. The rule: is new. `hey watch --box imbox --events new --exit-on-first` is "block until new mail", and a `--run-*` script sees `HEY_NEW=1` or `HEY_NEW=0`. - **A skip-ahead sets a floor.** A box that answered 409 was never read across the gap, so - the cursor it skips to — the box's last posting activity — becomes that box's floor: - activity at or before it is never new there, on a thread the watch knows or one it does + the cursor it skips to — HEY's clock when the skip read it — becomes that + box's floor: activity at or before it is never new there, on a thread the watch knows or one it does not, so a reply the watch missed and then a move while still unseen is not new mail. Mail after the floor is. The floor is the box's own; a gap thread moved to another box is measured there and may read as new once — the `resync` line is the cue to re-read. diff --git a/internal/cmd/sdk.go b/internal/cmd/sdk.go index c064d2d4..8d3637c1 100644 --- a/internal/cmd/sdk.go +++ b/internal/cmd/sdk.go @@ -136,6 +136,23 @@ func newSDKClient(extra ...hey.ClientOption) *hey.Client { return hey.NewClient(sdkClientCfg, nil, opts...) } +// newUncachedSDKClient builds a sibling of sdk — same configuration, same account — +// without the response cache, for a read a 304 must not answer. The cache revalidates +// on HEY's ETag, and an ETag can cover less than the body: /boxes.json's is the box +// rows, so a box's posting_changes_url, feed version and all, can change under it. +// The SDK has no per-request way past its cache, so this is a client of its own. +func newUncachedSDKClient(ctx context.Context) (*hey.Client, error) { + uncached := *sdkClientCfg + uncached.CacheEnabled = false + uncached.CacheDir = "" + client := hey.NewClient(&uncached, nil, sdkClientOpts...) + + if accountID, scoped := sdk.AccountID(); scoped { + return client.ForAccount(ctx, accountID) + } + return client, nil +} + func selectConfiguredAccount(ctx context.Context) error { client, err := clientForAccountSelection(ctx, rootSDK, cfg.AccountID) if err != nil { diff --git a/internal/cmd/watch.go b/internal/cmd/watch.go index a2a91935..ce0b17c7 100644 --- a/internal/cmd/watch.go +++ b/internal/cmd/watch.go @@ -73,7 +73,11 @@ func newWatchCommand() *watchCommand { Use: "watch", Short: "Follow email threads and calendars as they change", Long: `Print email threads and calendar changes as they happen: piped or with --json, one -JSON object per line; at a terminal, one text line each. Runs until interrupted. +JSON object per line; at a terminal, one text line each. Runs until interrupted. What +changed before the watch began is not reported, unless --since reads back to it first — +save a change from just before it: HEY's clock is read to the whole second, and taken back +by however long reading it took, so a change up to a second before the watch asked, plus +that request's time, may be reported. Its "at" says when it happened. Changes can drive a command instead of being printed, and that is a choice between two behaviours: --run-async spawns the command per change and moves on, so a slow one never @@ -82,7 +86,7 @@ Pass one or the other. Every added and updated line says whether the thread is new mail: unseen, not muted, and active since the watch last saw it — or since the watch began, for a thread it has not -seen, so the backlog a box's first read carries is not new. Reading a thread, muting or +seen, so the backlog --since reads is not new. Reading a thread, muting or moving it is not new activity; a reply on a known thread is. --events new selects the new ones, alone or alongside added, updated and deleted, and a script sees HEY_NEW=1 for them. @@ -98,7 +102,9 @@ Besides the thread changes, three lines describe the watch itself: "ready" once and calendar is caught up and the subscription is live (again after every reconnect's catch-up), "disconnected" when the connection drops, and "resync" when a box changed more than the feed can list one change at a time and the watch skipped ahead — re-read that -box. A resync is an event of its own: reported by default, scripts run for it and +box. One resync covers the whole catch-up and comes once the box is followed again: a box +still too busy after the skip is skipped again on the retry backoff, and its "at" is the +last skip. A resync is an event of its own: reported by default, scripts run for it and --exit-on-first counts it, and --events can leave it out, as --events new does. A calendar's feed falls behind the same way, and calendar_resync is the same word for it. Ready and disconnected are written to stdout only.`, @@ -150,9 +156,12 @@ func (c *watchCommand) run(cmd *cobra.Command, args []string) error { } // New mail is measured against the watch's start, so that is taken before - // the boxes' cursors are read — and the cursors start no later than it, so - // nothing that lands between the two sits behind a cursor, read by nothing. - started := serverNow(ctx) + // the boxes' cursors are read — and the cursors start at it, so nothing + // that lands between the two sits behind a cursor, read by nothing. + started, err := serverNow(ctx) + if err != nil { + return err + } newMail := trackNewMail(started) boxes, err := c.watchedBoxes(ctx, started) @@ -276,7 +285,7 @@ func (c *watchCommand) watchedBoxes(ctx context.Context, started time.Time) (map continue } if c.since == "" { - cursor = noLaterThan(cursor, started) + cursor.Since = watchStartSince(started) } watched[box.Id] = &watchedBox{id: box.Id, kind: box.Kind, name: box.Name, cursor: cursor, reported: c.watching(box)} @@ -306,8 +315,10 @@ func boxIs(box generated.Box, wanted string) bool { wanted == strconv.FormatInt(box.Id, 10) } -// watchCursor is where a box's changes feed should be read from. The server bakes its own -// clock into the box's changes URL, so that's the cursor unless --since moves it. +// watchCursor reads the cursor out of a box's changes URL, moved by --since. Its since +// is HEY's, which neither caller keeps without --since: a watch's first read replaces it +// with the watch's start (watchStartSince), and a skip-ahead with HEY's clock at the +// skip. Both keep the feed version it names. func watchCursor(changesURL, since string) (hey.PostingChangesCursor, error) { if changesURL == "" { return hey.PostingChangesCursor{}, nil @@ -332,22 +343,6 @@ func watchCursor(changesURL, since string) (hey.PostingChangesCursor, error) { const watchCursorTimeLayout = "2006-01-02T15:04:05.000Z" -// noLaterThan moves a box's cursor back to the watch's start when the box's -// own is later. The server bakes the box's last posting activity into its -// cursor, so mail that landed after the watch read HEY's clock and before it -// read the box list is already behind the cursor: the feed would start after -// it, and nothing would ever report it. Starting from the watch's own start -// reads it as part of the catch-up instead — and it is new, since it is later -// than the start. A cursor that cannot be read is left as it is. -func noLaterThan(cursor hey.PostingChangesCursor, started time.Time) hey.PostingChangesCursor { - at, err := time.Parse(watchCursorTimeLayout, cursor.Since) - if err == nil && at.After(started) { - cursor.Since = started.UTC().Format(watchCursorTimeLayout) - } - - return cursor -} - func parseWatchSince(since string) (time.Time, error) { if at, err := time.Parse(time.RFC3339, since); err == nil { return at, nil @@ -364,6 +359,20 @@ type watchedBox struct { name string cursor hey.PostingChangesCursor reported bool // --box named it, or named nothing + recovery feedRecovery +} + +// feedRecovery is how far a box's or a calendar's feed is into getting back from a 409. +// One skip-ahead usually does it; a feed busier than that — more than an increment's +// worth of changes after HEY's clock at the skip — answers 409 again straight after. +// That is one episode, not several. Its skips after the first wait on the retry backoff +// rather than following every doorbell, which would ring as fast as the feed is +// changing, and its one resync goes out with the clean read that ends it: after every +// skip it took, so a reader that re-reads on it is not left stale by a later one. +type feedRecovery struct { + skipped bool // the episode has skipped ahead + skippedTo time.Time // where its last skip landed: the resync's at + holding bool // the retry, not a doorbell, reads the feed next } // watchEvent is one changed posting or calendar recording, as a line of NDJSON or as a @@ -644,7 +653,9 @@ func (w *postingsWatch) read(ctx context.Context, message actioncable.Message) e return nil } - if box, watching := w.boxes[notification.BoxID]; watching { + // A box holding for its retry after a repeated 409 is read by the retry: a doorbell + // would only skip it ahead again. + if box, watching := w.boxes[notification.BoxID]; watching && !box.recovery.holding { if err := w.readBox(ctx, box); err != nil { return err } @@ -669,20 +680,14 @@ func (w *postingsWatch) readBox(ctx context.Context, box *watchedBox) error { return nil } } - w.wasRead(box) - if changes.FullSyncRequired { - fmt.Fprintf(w.errOut, "notice: too much changed in %s to follow one change at a time — skipping ahead, read the box with `hey box view %s`\n", box.name, box.kind) - skipped, err := w.skipAhead(ctx, box) - if err != nil { - return err - } - // A resync says the box is worth re-reading; a box that is gone is not. - if skipped { - w.report(ctx, watchEvent{Change: watchResync, At: watchTime(time.Now())}, box, nil) - } - return nil + return w.recoverBox(ctx, box) + } + w.wasRead(box) + if box.recovery.skipped { + w.report(ctx, watchEvent{Change: watchResync, At: watchTime(box.recovery.skippedTo)}, box, nil) } + box.recovery = feedRecovery{} if changes.NextCursor != nil { box.cursor = *changes.NextCursor @@ -706,6 +711,32 @@ func (w *postingsWatch) readBox(ctx context.Context, box *watchedBox) error { return nil } +// recoverBox gets a box that answered 409 back onto its feed by skipping it ahead. The +// first skip of an episode is announced on stderr and read from straight away, since +// one skip usually lands on a feed the watch can follow; the resync line, the reader's +// cue to re-read the box, goes out with the clean read (readBox). A 409 after a skip +// skips again without a word and leaves the box on the retry backoff, which doubles +// while the feed stays too busy to follow. A box that is gone is not worth re-reading, +// and a skip that has not happened — its retry is on the backoff — says nothing. +func (w *postingsWatch) recoverBox(ctx context.Context, box *watchedBox) error { + skippedTo, skipped, err := w.skipAhead(ctx, box) + if err != nil || !skipped { + return err + } + + first := !box.recovery.skipped + box.recovery.skipped = true + box.recovery.skippedTo = skippedTo + if !first { + box.recovery.holding = true + w.readAgainLater(box) + return nil + } + + fmt.Fprintf(w.errOut, "notice: too much changed in %s to follow one change at a time — skipped ahead, read the box with `hey box view %s`\n", box.name, box.kind) + return w.readBox(ctx, box) +} + // classify decides whether a posting is new mail and records it, in that order. func (w *postingsWatch) classify(box *watchedBox, posting generated.Posting) *bool { isNew := w.newMail.isNew(box.id, posting) @@ -763,38 +794,84 @@ func (w *postingsWatch) settleBackoff() { } } -// skipAhead moves a box's cursor to the server's current one, which is the only way -// back once a box has changed more than an increment can carry, and says whether it -// did. +// skipAhead moves a box's cursor to HEY's clock when it answered (serverNowAnswered), +// which is the only way back once a box has changed more than an increment can carry, +// and says where it skipped to and whether it did. +// +// Not to the since in the box's posting_changes_url, which is what it used to take: +// that is the box's last posting activity rather than HEY's clock, and a deletion or a +// bundled posting can come later than it. The list is still read, for the feed's +// version and to learn whether the box is still there — and read past the SDK's ETag +// cache (newUncachedSDKClient — a skip is rare enough to build one each time), because HEY's ETag for it is the box rows, which neither +// posting activity nor a new feed version touches. A 304 would hand back the version +// HEY had just refused, and a version HEY refuses answers 409 on every read. // // A box the server no longer lists, or no longer serves a changes feed for, has no // cursor to skip to: keeping the one it had would answer 409 on every read, and // installing an empty one would be a usage error on every read instead. Either way the // box can't be followed any more, so it stops being watched — and nothing was skipped. -func (w *postingsWatch) skipAhead(ctx context.Context, box *watchedBox) (bool, error) { - listed, err := sdk.Boxes().List(ctx) +// A clock that cannot be read leaves the cursor where it was, to be tried again on the +// retry's backoff like any read that failed. +func (w *postingsWatch) skipAhead(ctx context.Context, box *watchedBox) (time.Time, bool, error) { + // A skip that could not be made is tried again by the retry, not by the next doorbell, + // which would only meet the same 409 and the same failing read. + later := func() { + box.recovery.holding = true + w.readAgainLater(box) + } + client, err := newUncachedSDKClient(ctx) + if err != nil { + return time.Time{}, false, w.skipFailed(ctx, box.name, err, later) + } + listed, err := client.Boxes().List(ctx) if err != nil { - return false, apierr.FromSDK(err) + return time.Time{}, false, w.skipFailed(ctx, box.name, err, later) } if listed == nil { - return false, apierr.ErrAPI(0, "could not list boxes") + return time.Time{}, false, apierr.ErrAPI(0, "could not list boxes") } for _, listedBox := range *listed { - if listedBox.Id == box.id { - cursor, err := watchCursor(listedBox.PostingChangesUrl, "") - if err != nil { - return false, err - } - if cursor.Since != "" { - box.cursor = cursor - w.newMail.skippedTo(box.id, cursor) - return true, nil - } + if listedBox.Id != box.id { + continue + } + cursor, err := watchCursor(listedBox.PostingChangesUrl, "") + if err != nil { + return time.Time{}, false, err + } + if cursor.Since == "" { + break + } + + now, err := serverNowAnswered(ctx) + if err != nil { + return time.Time{}, false, w.skipFailed(ctx, box.name, err, later) } + cursor.Since = watchStartSince(now) + box.cursor = cursor + w.newMail.skippedTo(box.id, cursor) + return now, true, nil } - return false, w.stopWatching(box) + return time.Time{}, false, w.stopWatching(box) +} + +// skipFailed says what a skip-ahead's failed read — its client, its list or HEY's +// clock — comes to: nothing, when the watch is being interrupted or timed out, which is +// how a watch is meant to end; the error, when waiting will not help; and otherwise a +// warning and the retry's backoff (later), with the cursor where it was. A 500 on the +// list is a reason to try again, not to stop watching. +func (w *postingsWatch) skipFailed(ctx context.Context, name string, err error, later func()) error { + switch { + case ctx.Err() != nil: + return nil //nolint:nilerr // an interrupt or a --timeout is how a watch is meant to end + case permanentReadError(err): + return apierr.FromSDK(err) + default: + fmt.Fprintf(w.errOut, "warning: could not skip %s ahead: %v\n", name, err) + later() + return nil + } } // stopWatching drops a box the watch can't follow any longer. When it was the last one diff --git a/internal/cmd/watch_calendar.go b/internal/cmd/watch_calendar.go index 35b26fd5..3dc53053 100644 --- a/internal/cmd/watch_calendar.go +++ b/internal/cmd/watch_calendar.go @@ -46,11 +46,12 @@ const ( // watchedCalendar is one calendar the watch follows: how far its recording feed has been // read, and the stream that rings when it changes. type watchedCalendar struct { - id int64 - name string - cursor hey.CalendarChangesCursor - stream string - stop context.CancelFunc + id int64 + name string + cursor hey.CalendarChangesCursor + stream string + stop context.CancelFunc + recovery feedRecovery } // calendarsWatch is the calendar half of a watch: the calendars and their cursors, the @@ -79,9 +80,10 @@ func (c *watchCommand) watchingCalendars(changes map[string]bool) bool { } // watchedCalendars reads the calendars and where each one's feed should be read from. The -// cursors are capped at the watch's start the way the boxes' are — the server bakes each -// calendar's last activity into its URL, so a write that lands between reading the clock -// and reading the list would otherwise sit behind the cursor, read by nothing. +// cursors start at the watch's start the way the boxes' do (watchStartSince): the since +// HEY serves is a calendar's own updated_at — for the list, the latest of them — so a +// calendar deleted after the rest last changed was reported by the first poll of every +// watch. func (c *watchCommand) watchedCalendars(ctx context.Context, started time.Time) (*calendarsWatch, error) { list, err := sdk.Calendars().ListWithChanges(ctx) if err != nil { @@ -129,15 +131,16 @@ func (c *watchCommand) followedCalendar(listed hey.ListedCalendar, started time. }, nil } -// calendarCursor is where a feed should be read from: the URL HEY served, moved by -// --since, and otherwise capped at the watch's start. +// calendarCursor is where a feed should be read from: the feed and version the URL HEY +// served names, from --since or else from the watch's start. func (c *watchCommand) calendarCursor(changesURL string, started time.Time) (hey.CalendarChangesCursor, error) { cursor, err := hey.CalendarChangesCursorFrom(changesURL) if err != nil { return hey.CalendarChangesCursor{}, apierr.FromSDK(err) } if c.since == "" { - return calendarCursorNoLaterThan(cursor, started), nil + cursor.Since = watchStartSince(started) + return cursor, nil } at, err := parseWatchSince(c.since) @@ -149,17 +152,6 @@ func (c *watchCommand) calendarCursor(changesURL string, started time.Time) (hey return cursor, nil } -// calendarCursorNoLaterThan is noLaterThan for a calendar feed's cursor, and exists for -// the same race. A cursor that cannot be read is left as it is. -func calendarCursorNoLaterThan(cursor hey.CalendarChangesCursor, started time.Time) hey.CalendarChangesCursor { - at, err := time.Parse(watchCursorTimeLayout, cursor.Since) - if err == nil && at.After(started) { - cursor.Since = started.UTC().Format(watchCursorTimeLayout) - } - - return cursor -} - // calendarDisplayName is what a calendar is called on a watch line. The personal calendar // has no name of its own — HEY leaves the field empty and the web app labels it from the // identity — so the watch says what it is rather than nothing. @@ -277,7 +269,8 @@ func (w *postingsWatch) coalesceCalendarRing() { func (w *postingsWatch) readRungCalendars(ctx context.Context) error { w.calendar.due = nil for _, id := range w.calendar.take() { - if calendar, watching := w.calendar.calendars[id]; watching { + // A calendar holding for its retry after a repeated 409 is read by the retry. + if calendar, watching := w.calendar.calendars[id]; watching && !calendar.recovery.holding { if err := w.readCalendar(ctx, calendar); err != nil { return err } @@ -331,20 +324,14 @@ func (w *postingsWatch) readCalendar(ctx context.Context, calendar *watchedCalen return nil } } - delete(w.calendar.unread, calendar.id) - w.settleBackoff() - if changes.FullSyncRequired { - fmt.Fprintf(w.errOut, "notice: too much changed in %s to follow one change at a time — skipping ahead, re-read the calendar\n", calendar.name) - skipped, err := w.skipCalendarAhead(ctx, calendar) - if err != nil { - return err - } - if skipped { - w.reportCalendar(ctx, watchEvent{Change: watchCalendarResync, At: watchTime(time.Now())}, calendar.id, calendar.name) - } - return nil + return w.recoverCalendar(ctx, calendar) + } + w.calendarWasRead(calendar) + if calendar.recovery.skipped { + w.reportCalendar(ctx, watchEvent{Change: watchCalendarResync, At: watchTime(calendar.recovery.skippedTo)}, calendar.id, calendar.name) } + calendar.recovery = feedRecovery{} if changes.NextCursor != nil { calendar.cursor = *changes.NextCursor @@ -364,6 +351,35 @@ func (w *postingsWatch) readCalendar(ctx context.Context, calendar *watchedCalen return nil } +// recoverCalendar is recoverBox for a calendar's recording feed: the first skip is read +// from straight away, a 409 after a skip waits on the retry backoff, and the episode's +// one calendar_resync goes out with the clean read that ends it. +func (w *postingsWatch) recoverCalendar(ctx context.Context, calendar *watchedCalendar) error { + skippedTo, skipped, err := w.skipCalendarAhead(ctx, calendar) + if err != nil || !skipped { + return err + } + + first := !calendar.recovery.skipped + calendar.recovery.skipped = true + calendar.recovery.skippedTo = skippedTo + if !first { + calendar.recovery.holding = true + w.calendar.unread[calendar.id] = true + w.armRetry() + return nil + } + + fmt.Fprintf(w.errOut, "notice: too much changed in %s to follow one change at a time — skipped ahead, re-read the calendar\n", calendar.name) + return w.readCalendar(ctx, calendar) +} + +// calendarWasRead takes a calendar off the retry list once it is caught up. +func (w *postingsWatch) calendarWasRead(calendar *watchedCalendar) { + delete(w.calendar.unread, calendar.id) + w.settleBackoff() +} + // reportRecordings walks one bucket in type order, so a read reports the same changes in // the same order every time. func (w *postingsWatch) reportRecordings(ctx context.Context, calendar *watchedCalendar, change string, bucket map[string][]generated.Recording, at func(generated.Recording) time.Time) { @@ -380,32 +396,54 @@ func (w *postingsWatch) reportRecordings(ctx context.Context, calendar *watchedC } } -// skipCalendarAhead moves a calendar's cursor to the server's current one, which is the -// only way back once its feed has fallen too far behind, and says whether it did. A -// calendar the server no longer lists cannot be followed any more and stops being watched -// — the calendar-level feed reports its deletion in its own time. -func (w *postingsWatch) skipCalendarAhead(ctx context.Context, calendar *watchedCalendar) (bool, error) { - list, err := sdk.Calendars().ListWithChanges(ctx) +// skipCalendarAhead moves a calendar's cursor to HEY's clock when it answered, the only way +// back once its feed has fallen too far behind, and says where it skipped to and whether +// it did — the rule skipAhead follows for a box, and for the same reason: the since in +// the calendar's URL is its updated_at, not HEY's clock. The list is read for the feed's +// version and to learn whether the calendar is still there, past the SDK's cache as a +// box's is: a new feed version need not change the list's ETag. A calendar the server no +// longer lists cannot be followed any more and stops being watched — the calendar-level +// feed reports its deletion in its own time. A clock that cannot be read leaves the +// cursor where it was, to be tried again on the retry's backoff. +func (w *postingsWatch) skipCalendarAhead(ctx context.Context, calendar *watchedCalendar) (time.Time, bool, error) { + // Tried again by the retry, not by the next ring — as for a box. + later := func() { + calendar.recovery.holding = true + w.calendar.unread[calendar.id] = true + w.armRetry() + } + client, err := newUncachedSDKClient(ctx) if err != nil { - return false, apierr.FromSDK(err) + return time.Time{}, false, w.skipFailed(ctx, calendar.name, err, later) + } + list, err := client.Calendars().ListWithChanges(ctx) + if err != nil { + return time.Time{}, false, w.skipFailed(ctx, calendar.name, err, later) } if list == nil { - return false, apierr.ErrAPI(0, "could not list calendars") + return time.Time{}, false, apierr.ErrAPI(0, "could not list calendars") } for _, listed := range list.Calendars { - if listed.Calendar.Id == calendar.id { - cursor, err := hey.CalendarChangesCursorFrom(listed.RecordingChangesURL) - if err != nil { - return false, apierr.FromSDK(err) - } - calendar.cursor = cursor - return true, nil + if listed.Calendar.Id != calendar.id { + continue + } + cursor, err := hey.CalendarChangesCursorFrom(listed.RecordingChangesURL) + if err != nil { + return time.Time{}, false, apierr.FromSDK(err) + } + + now, err := serverNowAnswered(ctx) + if err != nil { + return time.Time{}, false, w.skipFailed(ctx, calendar.name, err, later) } + cursor.Since = watchStartSince(now) + calendar.cursor = cursor + return now, true, nil } w.stopWatchingCalendar(calendar.id) - return false, nil + return time.Time{}, false, nil } // pollCalendarList reads the calendar-level feed: calendars that arrived are followed diff --git a/internal/cmd/watch_calendar_test.go b/internal/cmd/watch_calendar_test.go index e1fc710e..90b67325 100644 --- a/internal/cmd/watch_calendar_test.go +++ b/internal/cmd/watch_calendar_test.go @@ -6,6 +6,7 @@ import ( "net/http" "net/http/httptest" "strings" + "sync" "testing" "time" @@ -59,8 +60,17 @@ func TestCalendarCursor(t *testing.T) { if err != nil { t.Fatalf("unexpected error: %v", err) } - if cursor.Since != "2026-08-18T09:00:00.000Z" || cursor.Version != "1" { - t.Errorf("cursor = %+v, want the server's own since and version", cursor) + if cursor.Since != "2026-08-18T10:00:00.000Z" || cursor.Version != "1" { + t.Errorf("cursor = %+v, want the watch's start and the server's version", cursor) + } + + later := "https://app.hey.com/calendars/512/recording/changes.json?since=2026-08-18T10%3A30%3A00.518496Z&v=1" + cursor, err = command.calendarCursor(later, started) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if cursor.Since != "2026-08-18T10:00:00.000Z" { + t.Errorf("since = %q, want a cursor later than the start moved back to it", cursor.Since) } command.since = "2026-08-17T08:30:00Z" @@ -81,22 +91,69 @@ func TestCalendarCursor(t *testing.T) { } } -func TestCalendarCursorNoLaterThan(t *testing.T) { - started := time.Date(2026, 8, 18, 9, 0, 0, 0, time.UTC) +// The list's cursor is the latest updated_at of the calendars still on it, so a calendar +// deleted after the rest last changed is later than it — and history by the time the +// watch starts. The feed here answers what is strictly later than its cursor, as HEY's +// does. +func TestWatchPollDoesNotReportHistoryAsItStarts(t *testing.T) { + t.Setenv("HEY_TOKEN", "test-token") - late := hey.CalendarChangesCursor{Since: "2026-08-18T09:30:00.000Z", Version: "1"} - if capped := calendarCursorNoLaterThan(late, started); capped.Since != "2026-08-18T09:00:00.000Z" { - t.Errorf("since = %q, want a cursor later than the start moved back to it", capped.Since) + var mu sync.Mutex + deletedAt := "2026-08-20T16:42:07.204613Z" + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + defer mu.Unlock() + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/calendars.json": + _, _ = w.Write([]byte(`{ + "calendars": [{"calendar": {"id": 512, "name": "Household"}, + "recording_changes_url": "/calendars/512/recording/changes.json?since=2026-08-18T11%3A00%3A00.000000Z&v=1", + "signed_stream_name": "signed-household"}], + "calendar_changes_url": "/calendar/changes.json?since=2026-08-18T11%3A00%3A00.000000Z" + }`)) + case "/calendar/changes.json": + since, err := time.Parse(time.RFC3339Nano, r.URL.Query().Get("since")) + deleted, _ := time.Parse(time.RFC3339Nano, deletedAt) + if err != nil || !deleted.After(since) { + _, _ = w.Write([]byte(`{}`)) + return + } + w.Header().Set("Link", `; rel="next"`) + _, _ = w.Write([]byte(`{"deleted": [{"id": 513, "deleted_at": "` + deletedAt + `"}]}`)) + default: + http.NotFound(w, r) + } + })) + defer server.Close() + initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) + + calendars, err := newWatchCommand().watchedCalendars(context.Background(), watchStarted) + if err != nil { + t.Fatalf("unexpected error: %v", err) } + t.Cleanup(calendars.poll.Stop) + watch, out := newTestWatch(defaultChanges...) + watch.exitOnFirst = true + watch.calendar = calendars - early := hey.CalendarChangesCursor{Since: "2026-08-18T08:00:00.000Z"} - if kept := calendarCursorNoLaterThan(early, started); kept.Since != "2026-08-18T08:00:00.000Z" { - t.Errorf("since = %q, want a cursor before the start left alone", kept.Since) + if err := watch.pollCalendarList(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if out.Len() != 0 || watch.finished() { + t.Fatalf("wrote %q, want nothing for a calendar deleted before the watch began", out.String()) } - unreadable := hey.CalendarChangesCursor{Since: "whenever"} - if kept := calendarCursorNoLaterThan(unreadable, started); kept.Since != "whenever" { - t.Errorf("since = %q, want a cursor we can't read left as it is", kept.Since) + // A calendar deleted after the watch began is a change. + mu.Lock() + deletedAt = "2026-08-21T09:04:31.880112Z" + mu.Unlock() + if err := watch.pollCalendarList(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } + lines := watchLines(t, out) + if len(lines) != 1 || lines[0]["change"] != watchCalendarDeleted || !watch.finished() { + t.Errorf("wrote %v, want the deletion after the start, ending the watch", lines) } } @@ -207,7 +264,12 @@ func TestWatchCalendarSkipsAheadOnAFullSync(t *testing.T) { }`)) return } - w.WriteHeader(http.StatusConflict) + if r.URL.Path == "/identity.json" { + w.Header().Set("Date", skipDate) + _, _ = w.Write([]byte(`{"id":1}`)) + return + } + answerTooFarBehindBefore(w, r, time.Date(2026, 8, 21, 11, 0, 0, 0, time.UTC)) })) defer server.Close() initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) @@ -226,8 +288,85 @@ func TestWatchCalendarSkipsAheadOnAFullSync(t *testing.T) { if event.Change != watchCalendarResync || event.Calendar == nil || event.Calendar.ID != 512 { t.Errorf("event = %+v, want a calendar_resync naming the calendar", event) } - if watch.calendar.calendars[512].cursor.Since != "2026-08-18T11:00:00.000Z" { - t.Errorf("cursor = %+v, want the server's own fresh cursor", watch.calendar.calendars[512].cursor) + skippedTo, err := time.Parse(time.RFC3339Nano, event.At) + if err != nil { + t.Fatalf("resync at %q: %v", event.At, err) + } + wantSkippedToHEYsClock(t, skippedTo) + if got := watch.calendar.calendars[512].cursor; got.Since != watchStartSince(skippedTo) || got.Version != "1" { + t.Errorf("cursor = %+v, want HEY's clock at the skip the resync names, and the feed's version", got) + } +} + +// A calendar's feed can move to a new version without the calendar list's ETag — its +// calendars and the selection — changing, so the skip-ahead reads the list past the cache. +func TestWatchCalendarSkipAheadReadsTheListPastTheCache(t *testing.T) { + t.Setenv("HEY_TOKEN", "test-token") + t.Setenv("XDG_CACHE_HOME", t.TempDir()) + + var mu sync.Mutex + var notModified, conflicts int + version := "1" + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + defer mu.Unlock() + switch r.URL.Path { + case "/identity.json": + w.Header().Set("Date", skipDate) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"id":1}`)) + case "/calendars.json": + w.Header().Set("ETag", `W/"calendars-unchanged"`) + if r.Header.Get("If-None-Match") == `W/"calendars-unchanged"` { + notModified++ + w.WriteHeader(http.StatusNotModified) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"calendars": [{"calendar": {"id": 512, "name": "Household"}, + "recording_changes_url": "/calendars/512/recording/changes.json?since=2026-08-18T11%3A00%3A00.000000Z&v=` + version + `"}], + "calendar_changes_url": "/calendar/changes.json?since=2026-08-18T11%3A00%3A00.000000Z"}`)) + default: + if r.URL.Query().Get("v") != "2" { + conflicts++ + w.WriteHeader(http.StatusConflict) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{}`)) + } + })) + defer server.Close() + initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) + + // The list as the watch read it at its start, now in the SDK's cache; then HEY + // moves the recording feed to version 2. + if _, err := sdk.Calendars().ListWithChanges(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } + mu.Lock() + version = "2" + mu.Unlock() + + watch, out := newTestWatch("recording_added", "calendar_resync") + watch.calendar = newTestCalendarsWatch(t, &watchedCalendar{id: 512, name: "Household", cursor: hey.CalendarChangesCursor{Since: "2026-08-18T11:00:00.000Z", Version: "1"}}) + + for range 2 { + if err := watch.readCalendar(context.Background(), watch.calendar.calendars[512]); err != nil { + t.Fatalf("unexpected error: %v", err) + } + } + + mu.Lock() + defer mu.Unlock() + if notModified != 0 || conflicts != 1 { + t.Errorf("list served from the cache %d times and the feed answered 409 %d times, want never and once", notModified, conflicts) + } + if got := watch.calendar.calendars[512].cursor.Version; got != "2" { + t.Errorf("version = %q, want the one HEY speaks now", got) + } + if lines := watchLines(t, out); len(lines) != 1 || lines[0]["change"] != watchCalendarResync { + t.Errorf("wrote %v, want one calendar_resync", lines) } } @@ -240,6 +379,11 @@ func TestWatchCalendarStopsWatchingAGoneCalendar(t *testing.T) { _, _ = w.Write([]byte(`{"calendars": [], "calendar_changes_url": "/calendar/changes.json?since=2026-08-18T11%3A00%3A00.000Z"}`)) return } + if r.URL.Path == "/identity.json" { + w.Header().Set("Date", skipDate) + _, _ = w.Write([]byte(`{"id":1}`)) + return + } w.WriteHeader(http.StatusConflict) })) defer server.Close() diff --git a/internal/cmd/watch_new.go b/internal/cmd/watch_new.go index a29775a9..e77b87d1 100644 --- a/internal/cmd/watch_new.go +++ b/internal/cmd/watch_new.go @@ -8,6 +8,8 @@ import ( "github.com/basecamp/hey-sdk/go/pkg/generated" hey "github.com/basecamp/hey-sdk/go/pkg/hey" + + "github.com/basecamp/hey-cli/internal/apierr" ) // New mail is a watch event: every added and updated line says whether the @@ -17,7 +19,7 @@ import ( // what to do about it (a toast, a bell, a re-read) is the reader's. // // There is no state file: new means active since the watch began, so the -// backlog a box's first read carries is not new, and what the watch remembers — +// backlog --since reads is not new, and what the watch remembers — // when each thread was last active — lives and dies with it. // newMail keeps up with every thread the watch reads, in every box and whatever @@ -42,12 +44,14 @@ func trackNewMail(started time.Time) *newMail { // one it does not, updated later while still unseen — moved, say — would // otherwise measure its gap activity against an older record and read as // new. Activity at or before the floor is never new in that box; the cursor is -// the box's last posting activity, which bounds every thread in it. The floor +// HEY's clock at the skip, which bounds every thread in it. The floor // is the box's alone — a gap thread that moves to another box is measured // there, and may still read as new once. The resync line is the reader's cue -// to re-read the box either way. +// to re-read the box either way. HEY writes the cursor to the microsecond, so +// it is read as RFC 3339 with any fraction rather than in the watch's own +// millisecond layout, which refuses it. func (n *newMail) skippedTo(boxID int64, cursor hey.PostingChangesCursor) { - if at, err := time.Parse(watchCursorTimeLayout, cursor.Since); err == nil { + if at, err := time.Parse(time.RFC3339Nano, cursor.Since); err == nil { n.floors[boxID] = at } } @@ -64,22 +68,88 @@ func (n *newMail) skippedTo(boxID int64, cursor hey.PostingChangesCursor) { // the request took — the local monotonic clock, which a wrong wall clock does // not touch — and a slow request, or one the SDK retried, only moves the start // earlier. Date is whole seconds, rounded down, which errs the same way: towards -// calling mail a moment old new rather than mail a moment new old. The SDK -// caches GETs by URL, so a query the server ignores keeps this one out of the -// cache; and when the server's clock can't be read, the local clock at the -// start stands in. Either way the start is handed out as a cutoff: a whole -// millisecond, strictly before the instant it stands for. -func serverNow(ctx context.Context) time.Time { +// calling mail a moment old new rather than mail a moment new old. That is a +// window of up to a second, plus the request's own time: a change from that +// long before the watch began is after the start, so it is read, reported and +// can be new, and --exit-on-first can stop on it. Nothing HEY serves on demand +// says the time any finer: Action Cable's pings are whole seconds, and no JSON +// answer carries a server "now". Finer times do exist — a posting doorbell's at +// and a feed's cursors are to the microsecond — but only once something has +// changed, which is no use for the moment a watch starts. And rounding the +// other way would skip up to +// a second of changes that did come after the start, which is worse than +// repeating one that did not. A reader who cares can tell from the line's at. +// The SDK caches GETs by URL, so a query the server ignores keeps this one out +// of the cache. The start is handed out as a cutoff: a whole millisecond, +// strictly before the instant it stands for. +// +// A watch that cannot read HEY's clock does not start. The workstation's clock +// is no stand-in: every feed starts at the start (watchStartSince) and new mail +// is measured against it, so a fast clock would put both in HEY's future — +// changes unread and new mail called old until the clocks met — and a slow one +// would report history. A request that fails here would fail at the box list +// next anyway. +func serverNow(ctx context.Context) (time.Time, error) { + answered, took, err := readHEYsClock(ctx) + if err != nil { + return time.Time{}, err + } + + return cutoffBefore(answered.Add(-took)), nil +} + +// serverNowAnswered is HEY's clock when it answered, not taken back by the +// request's time: where a skip-ahead resumes. A skip has already given up on +// the gap — the resync line says so — so there is nothing to catch between +// asking and the answer, and a skip point taken back by a slow or retried +// request could leave a busy feed still too far behind to follow. +func serverNowAnswered(ctx context.Context) (time.Time, error) { + answered, _, err := readHEYsClock(ctx) + if err != nil { + return time.Time{}, err + } + + return cutoffBefore(answered), nil +} + +// readHEYsClock asks HEY the time: the Date header of its answer, and how long +// the request took on the local monotonic clock. +func readHEYsClock(ctx context.Context) (time.Time, time.Duration, error) { started := time.Now() response, err := rootSDK.Get(ctx, "/identity.json?clock="+strconv.FormatInt(started.UnixNano(), 10)) - if err != nil || response == nil || response.FromCache { - return cutoffBefore(started) + if err != nil { + return time.Time{}, 0, apierr.FromSDK(err) + } + if response == nil || response.FromCache { + return time.Time{}, 0, apierr.ErrAPI(0, "could not read HEY's clock — the watch needs it to tell what happened after it began") } - if at, err := http.ParseTime(response.Headers.Get("Date")); err == nil { - return cutoffBefore(at.Add(-time.Since(started))) + answered, err := http.ParseTime(response.Headers.Get("Date")) + if err != nil { + return time.Time{}, 0, apierr.ErrAPI(0, "HEY's answer carried no Date header — the watch needs HEY's clock to tell what happened after it began") } - return cutoffBefore(started) + return answered, time.Since(started), nil +} + +// watchStartSince is where a feed's first read begins without --since: the +// watch's start, in place of the since HEY's changes URL carries. That since is +// not HEY's clock. A box's is its last posting activity — the latest updated_at +// among its unbundled postings, or the box's own when it has none — and a +// calendar's is its updated_at, the list's the latest of them; the feeds answer +// changes later than that which are history by now: a deletion, a bundled +// posting, a calendar deleted after the rest last changed. Nor is it always +// even that: the box list comes through the SDK's ETag cache, HEY's ETag for it +// is the box rows, and posting activity does not touch them, so a 304 hands +// back the since as it stood when the list was cached — hours or days behind. +// Read from there, the catch-up reported what came after as news on every +// start, and --exit-on-first stopped on the first of it. +// +// From the start a feed reports what happened after it and nothing before — +// including mail that landed after the watch read HEY's clock and before it +// read the box list, which a since later than the start would leave behind it, +// read by nothing; that mail is new, too. +func watchStartSince(started time.Time) string { + return started.UTC().Format(watchCursorTimeLayout) } // cutoffBefore makes an instant usable as the watch's start: a cursor is @@ -95,8 +165,7 @@ func cutoffBefore(at time.Time) time.Time { // isNew says whether a posting is new mail: unseen, not muted, and active since // this watch last saw the thread — or since the watch began, for a thread it // has no record of, because anything active before that was already there: the -// backlog a box's first read carries from the server's cursor, or a thread that -// merely moved in. active_at moves on new mail only, not when a thread is read, +// backlog --since reads, or a thread that merely moved in. active_at moves on new mail only, not when a thread is read, // muted or moved, so none of those is new and a reply on a known thread is. func (n *newMail) isNew(boxID int64, posting generated.Posting) bool { last, known := n.activeAt[posting.Id] diff --git a/internal/cmd/watch_new_test.go b/internal/cmd/watch_new_test.go index 120afa57..27ee263b 100644 --- a/internal/cmd/watch_new_test.go +++ b/internal/cmd/watch_new_test.go @@ -20,6 +20,16 @@ import ( var watchStarted = time.Date(2026, 8, 21, 9, 0, 0, 0, time.UTC) +// readServerNow is serverNow where the test's server answers with its clock. +func readServerNow(t *testing.T) time.Time { + t.Helper() + started, err := serverNow(context.Background()) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + return started +} + func newPosting(id int64, sender, subject string, activeAt time.Time) generated.Posting { return generated.Posting{Id: id, Name: subject, ActiveAt: activeAt, Creator: generated.Contact{Name: sender}} } @@ -47,9 +57,8 @@ func TestNewMailIsSinceTheWatchBeganNotTheBacklog(t *testing.T) { before := watchStarted.Add(-time.Hour) after := watchStarted.Add(30 * time.Second) - // A box's first read carries everything since the server's cursor — the - // box's last activity — which may be hours of backlog, plus mail that - // arrived while the watch was starting up. + // A read from --since carries backlog, alongside mail that arrived while + // the watch was starting up. backlog := newPosting(101, "Maria Delgado", "Lunch on Thursday?", before) arrived := newPosting(102, "Northwind Invoicing", "Invoice #4021", after) if fresh := classifyRead(tracker, backlog, arrived); len(fresh) != 1 || fresh[0] != 102 { @@ -162,6 +171,12 @@ func TestNewMailAfterASkipAheadIsSinceTheSkip(t *testing.T) { if _, has := tracker.floors[24089]; has { t.Error("an unreadable cursor must not become a floor") } + + // HEY writes its cursors to the microsecond. + tracker.skippedTo(24090, hey.PostingChangesCursor{Since: "2026-08-21T10:00:00.518496Z", Version: "2"}) + if got := tracker.floors[24090]; !got.Equal(time.Date(2026, 8, 21, 10, 0, 0, 518496000, time.UTC)) { + t.Errorf("floor = %v, want the cursor HEY wrote, to the microsecond", got) + } } func TestNewMailCarriedTwiceByOneReadIsNewOnce(t *testing.T) { @@ -220,7 +235,7 @@ func TestServerNowReadsTheServersClock(t *testing.T) { initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) date := time.Date(2026, 8, 21, 9, 0, 5, 0, time.UTC) - got := serverNow(context.Background()) + got := readServerNow(t) // The Date header, less the instant the request took: on the server's // clock, whatever the local one says, and no later than the header. if got.After(date) || date.Sub(got) > time.Second { @@ -230,7 +245,7 @@ func TestServerNowReadsTheServersClock(t *testing.T) { t.Errorf("requested %v, want one uncacheable identity request", requested) } - second := serverNow(context.Background()) + second := readServerNow(t) if len(requested) != 2 || requested[1] == requested[0] { t.Errorf("requested %v, want a fresh request each time, never the cache", requested) } @@ -249,16 +264,15 @@ func TestCutoffBeforeIsAWholeMillisecondStrictlyBefore(t *testing.T) { t.Errorf("cutoffBefore(%v) = %v, want strictly before even on a boundary", exact, got) } - // So mail in the start's own millisecond is new, and a cursor at that - // millisecond is moved back to before it. + // So mail in the start's own millisecond is new, and a cursor started at + // the watch's start reads it: the feed answers what is strictly later. tracker := trackNewMail(cutoffBefore(within)) landed := within.Truncate(time.Millisecond) if !tracker.isNew(24088, newPosting(101, "Maria Delgado", "Lunch on Thursday?", landed)) { t.Error("mail in the same millisecond as the watch's start is new") } - cursor := noLaterThan(hey.PostingChangesCursor{Since: landed.Format(watchCursorTimeLayout)}, cutoffBefore(within)) - if cursor.Since != "2026-08-21T09:00:05.122Z" { - t.Errorf("cursor = %q, want it moved back to before the millisecond the mail landed in", cursor.Since) + if since := watchStartSince(cutoffBefore(within)); since != "2026-08-21T09:00:05.122Z" { + t.Errorf("since = %q, want it before the millisecond the mail landed in", since) } } @@ -274,7 +288,7 @@ func TestServerNowIsTheClockWhenTheRequestBeganNotWhenItWasAnswered(t *testing.T defer server.Close() initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) - got := serverNow(context.Background()) + got := readServerNow(t) // Mail that lands while the server is answering is later than the start; // a start taken at the Date header would put it before. @@ -320,7 +334,7 @@ func TestWatchReadsMailThatLandedWhileItReadTheClock(t *testing.T) { initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) // As run does: the clock, then the boxes, then the catch-up. - started := serverNow(context.Background()) + started := readServerNow(t) command := newWatchCommand() boxes, err := command.watchedBoxes(context.Background(), started) if err != nil { @@ -347,23 +361,29 @@ func TestWatchReadsMailThatLandedWhileItReadTheClock(t *testing.T) { } } -func TestServerNowFallsBackToTheLocalClock(t *testing.T) { +// A watch that cannot read HEY's clock does not start: the workstation's clock is +// no stand-in for the cutoff every feed and new mail are measured against. +func TestServerNowRefusesWithoutHEYsClock(t *testing.T) { t.Setenv("HEY_TOKEN", "test-token") - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + + down := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { http.Error(w, "down", http.StatusBadGateway) })) - defer server.Close() - initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) - - before := time.Now() - got := serverNow(context.Background()) - // The local clock at the request's start, as a cutoff: a whole millisecond, - // strictly before — so up to two milliseconds before the instant itself. - if got.Before(before.Add(-2*time.Millisecond)) || got.After(time.Now()) { - t.Errorf("serverNow = %v, want the local clock at the start when the server's can't be read", got) + defer down.Close() + initSDK(auth.NewManager(down.URL, down.Client(), t.TempDir()), down.URL) + if _, err := serverNow(context.Background()); err == nil { + t.Error("expected an error when HEY cannot be reached") } - if got.Nanosecond()%int(time.Millisecond) != 0 { - t.Errorf("serverNow = %v, want a whole millisecond", got) + + undated := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header()["Date"] = nil + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"id":1}`)) + })) + defer undated.Close() + initSDK(auth.NewManager(undated.URL, undated.Client(), t.TempDir()), undated.URL) + if _, err := serverNow(context.Background()); err == nil { + t.Error("expected an error when HEY's answer carries no Date header") } } diff --git a/internal/cmd/watch_test.go b/internal/cmd/watch_test.go index 969a37ae..560d3bea 100644 --- a/internal/cmd/watch_test.go +++ b/internal/cmd/watch_test.go @@ -12,6 +12,8 @@ import ( "os" "slices" "strings" + "sync" + "sync/atomic" "testing" "time" @@ -421,13 +423,18 @@ func TestWatchStopsOnAReadThatCannotWork(t *testing.T) { } } -// boxesAndChanges answers the two reads a skip-ahead makes: a changes feed that is too far -// behind to follow, and the box list it then looks for a fresh cursor in. +// boxesAndChanges answers the reads a skip-ahead makes: a changes feed that is too far +// behind to follow, the box list it looks for the box and its feed's version in, and +// HEY's clock, which it skips to. func boxesAndChanges(t *testing.T, boxes string) *httptest.Server { t.Helper() return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch r.URL.Path { + case "/identity.json": + w.Header().Set("Date", skipDate) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"id":1}`)) case "/boxes.json": w.Header().Set("Content-Type", "application/json") _, _ = w.Write([]byte(boxes)) @@ -437,9 +444,22 @@ func boxesAndChanges(t *testing.T, boxes string) *httptest.Server { })) } -func TestWatchSkipsAheadToTheBoxesOwnCursor(t *testing.T) { +// HEY's clock when a skip-ahead reads it. +const skipDate = "Fri, 21 Aug 2026 11:05:00 GMT" + +// wantSkippedToHEYsClock checks a skip-ahead's point: HEY's clock when it answered — the +// millisecond before the Date header, not taken back by the request's time, which a +// skip has no gap to catch in and which could leave a busy feed still behind. +func wantSkippedToHEYsClock(t *testing.T, skippedTo time.Time) { + t.Helper() + if want := time.Date(2026, 8, 21, 11, 4, 59, 999000000, time.UTC); !skippedTo.Equal(want) { + t.Errorf("skipped to %v, want HEY's clock when it answered, %v", skippedTo, want) + } +} + +func TestWatchSkipsAheadToHEYsClock(t *testing.T) { t.Setenv("HEY_TOKEN", "test-token") - server := boxesAndChanges(t, `[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-21T11%3A02%3A00.000Z&v=2"}]`) + server := boxesAndChanges(t, `[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-21T11%3A02%3A00.518496Z&v=2"}]`) defer server.Close() initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) @@ -454,11 +474,384 @@ func TestWatchSkipsAheadToTheBoxesOwnCursor(t *testing.T) { if watch.boxes[24088] == nil { t.Fatal("the box should still be watched") } - if got := watch.boxes[24088].cursor.Since; got != "2026-08-21T11:02:00.000Z" { - t.Errorf("cursor = %q, want the server's current one", got) + floor := watch.newMail.floors[24088] + wantSkippedToHEYsClock(t, floor) + if got := watch.boxes[24088].cursor; got.Since != watchStartSince(floor) || got.Version != "2" { + t.Errorf("cursor = %+v, want HEY's clock at the skip, the new-mail floor, and the box's feed version", got) + } +} + +// skipHEY is HEY as a skip-ahead meets it: a feed read from before 11:00 too far behind +// to follow and one from after it clean, and the box list, the calendar list and the +// clock answering — save the reads fail answers itself, which it says it did by +// returning true. +func skipHEY(t *testing.T, fail func(http.ResponseWriter, *http.Request) bool) { + t.Helper() + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if fail(w, r) { + return + } + switch r.URL.Path { + case "/identity.json": + w.Header().Set("Date", skipDate) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"id":1}`)) + case "/boxes.json": + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-21T11%3A02%3A00.518496Z&v=2"}]`)) + case "/calendars.json": + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"calendars": [{"calendar": {"id": 512, "name": "Household"}, + "recording_changes_url": "/calendars/512/recording/changes.json?since=2026-08-18T11%3A00%3A00.000000Z&v=1"}]}`)) + default: + answerTooFarBehindBefore(w, r, time.Date(2026, 8, 21, 11, 0, 0, 0, time.UTC)) + } + })) + t.Cleanup(server.Close) + t.Setenv("HEY_TOKEN", "test-token") + initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) +} + +// answerTooFarBehindBefore answers a changes feed read as HEY would for a feed that +// changed too much before behind to follow — `head :conflict`, no body — and not at +// all since. +func answerTooFarBehindBefore(w http.ResponseWriter, r *http.Request, behind time.Time) { + since, err := time.Parse(time.RFC3339Nano, r.URL.Query().Get("since")) + if err != nil || since.Before(behind) { + w.WriteHeader(http.StatusConflict) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{}`)) +} + +// behindWatch is a watch whose Imbox and Household calendar have both fallen too far +// behind to follow. +func behindWatch(t *testing.T) (*postingsWatch, *bytes.Buffer, *bytes.Buffer) { + t.Helper() + watch, out := newTestWatch("added", "resync", "calendar_resync") + errOut := &bytes.Buffer{} + watch.errOut = errOut + watch.boxes[24088].cursor.Since = "2026-08-01T00:00:00.000Z" + watch.calendar = newTestCalendarsWatch(t, &watchedCalendar{id: 512, name: "Household", cursor: hey.CalendarChangesCursor{Since: "2026-08-01T00:00:00.000Z", Version: "1"}}) + return watch, out, errOut +} + +// readBehind reads the box or the calendar behindWatch left behind. +func readBehind(ctx context.Context, watch *postingsWatch, feed string) error { + if feed == "box" { + return watch.readBox(ctx, watch.boxes[24088]) + } + return watch.readCalendar(ctx, watch.calendar.calendars[512]) +} + +// ringFeed rings the doorbell for the box or the calendar behindWatch left behind, and +// answers it the way the watch's loop does. +func ringFeed(t *testing.T, watch *postingsWatch, feed string) { + t.Helper() + var err error + if feed == "box" { + err = watch.read(context.Background(), actioncable.Message(`{"change":"upsert","box_id":24088}`)) + } else { + watch.calendar.ring(512) + <-watch.calendar.wake + err = watch.readRungCalendars(context.Background()) + } + if err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +// listOf is the list a skip-ahead reads for a feed. +func listOf(feed string) string { + if feed == "box" { + return "/boxes.json" + } + return "/calendars.json" +} + +// An interrupt or --timeout while a skip-ahead reads its list or HEY's clock is how a +// watch is meant to end, not a failed read: nothing is warned about, no retry is armed, +// and nothing says it skipped. +func TestWatchSkipAheadEndsQuietlyWhenInterrupted(t *testing.T) { + for _, feed := range []string{"box", "calendar"} { + for _, at := range []string{"/identity.json", listOf(feed)} { + t.Run(feed+at, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + skipHEY(t, func(w http.ResponseWriter, r *http.Request) bool { + if r.URL.Path != at { + return false + } + cancel() + <-r.Context().Done() + return true + }) + watch, out, errOut := behindWatch(t) + + if err := readBehind(ctx, watch, feed); err != nil { + t.Fatalf("read = %v, want an interrupted skip-ahead to end quietly", err) + } + if errOut.Len() != 0 { + t.Errorf("stderr = %q, want nothing for an interrupt", errOut.String()) + } + if watch.retry != nil || len(watch.unread) != 0 || len(watch.calendar.unread) != 0 { + t.Error("an interrupted skip-ahead should not arm a retry") + } + if out.Len() != 0 { + t.Errorf("wrote %q, want no resync for a skip that did not happen", out.String()) + } + if watch.boxes[24088].cursor.Since != "2026-08-01T00:00:00.000Z" || watch.calendar.calendars[512].cursor.Since != "2026-08-01T00:00:00.000Z" { + t.Error("an interrupted skip-ahead should leave the cursor where it was") + } + }) + } + } +} + +// A list that fails to read for a while is a reason to try the skip again, not to stop +// watching: the watch warns, keeps the cursor, retries on the backoff, and says it +// skipped only once it has. +func TestWatchSkipAheadRetriesAListThatFailed(t *testing.T) { + for _, feed := range []string{"box", "calendar"} { + t.Run(feed, func(t *testing.T) { + var down atomic.Bool + var listReads, feedReads atomic.Int32 + down.Store(true) + skipHEY(t, func(w http.ResponseWriter, r *http.Request) bool { + if strings.Contains(r.URL.Path, "/changes") { + feedReads.Add(1) + } + if r.URL.Path != listOf(feed) { + return false + } + listReads.Add(1) + if !down.Load() { + return false + } + http.Error(w, "down for maintenance", http.StatusInternalServerError) + return true + }) + watch, out, errOut := behindWatch(t) + + if err := readBehind(context.Background(), watch, feed); err != nil { + t.Fatalf("read = %v, want a failed list read retried, not the watch ended", err) + } + if !strings.Contains(errOut.String(), "warning: could not skip") || strings.Contains(errOut.String(), "notice") { + t.Errorf("stderr = %q, want a warning and no word of a skip that has not happened", errOut.String()) + } + if watch.retry == nil || len(watch.unread)+len(watch.calendar.unread) != 1 { + t.Fatal("a failed list read should be retried on the backoff") + } + if out.Len() != 0 { + t.Errorf("wrote %q, want no resync before the skip", out.String()) + } + + // Doorbells wait for the retry rather than trying the failing list again. + ringFeed(t, watch, feed) + ringFeed(t, watch, feed) + if got := listReads.Load(); got != 1 { + t.Errorf("read the list %d times, want once — the retry, not a doorbell, tries again", got) + } + + // The list is back: the retry skips, reads the feed from there, and says so + // once. + down.Store(false) + errOut.Reset() + if err := watch.retryUnread(context.Background()); err != nil { + t.Fatalf("retry = %v", err) + } + if lines := watchLines(t, out); len(lines) != 1 || !strings.HasSuffix(lines[0]["change"].(string), "resync") { + t.Errorf("wrote %v, want one resync once the skip happened", lines) + } + if !strings.Contains(errOut.String(), "notice: too much changed") { + t.Errorf("stderr = %q, want the skip announced once it happened", errOut.String()) + } + + // Skipped, the feed follows its doorbells again. + reads := feedReads.Load() + ringFeed(t, watch, feed) + if feedReads.Load() == reads { + t.Error("a doorbell after the skip should read the feed again") + } + }) + } +} + +// A feed busier than one skip can outrun answers 409 again straight after the skip. That +// is one recovery: one notice, skips after the first waiting on the retry backoff — +// doubling — rather than following every doorbell, and one resync line, when a clean +// read ends it, so a reader that re-reads on it has missed nothing a later skip passed. +func TestWatchRecoversFromARepeated409Once(t *testing.T) { + for _, feed := range []string{"box", "calendar"} { + t.Run(feed, func(t *testing.T) { + var feedReads atomic.Int32 + var quiet atomic.Bool + skipHEY(t, func(w http.ResponseWriter, r *http.Request) bool { + if !strings.Contains(r.URL.Path, "/changes") { + return false + } + feedReads.Add(1) + if quiet.Load() { + return false + } + w.WriteHeader(http.StatusConflict) // as HEY answers for a feed this busy + return true + }) + watch, out, errOut := behindWatch(t) + ring := func() { ringFeed(t, watch, feed) } + + for range 5 { + ring() + } + if got := feedReads.Load(); got != 2 { + t.Errorf("read the feed %d times for five doorbells, want twice — the skip's own read, then the rest held for the retry", got) + } + if lines := watchLines(t, out); len(lines) != 0 { + t.Errorf("wrote %v, want the resync kept for the clean read that ends the recovery", lines) + } + if got := strings.Count(errOut.String(), "notice: too much changed"); got != 1 { + t.Errorf("stderr = %q, want the skip announced once", errOut.String()) + } + if watch.retry == nil || watch.backoff != firstWatchRetry { + t.Fatalf("backoff = %v, want the repeat held for the first retry", watch.backoff) + } + + // The retry comes round, as the loop runs it: another 409, another skip, + // the same recovery, and a longer wait before the next. + watch.retry = nil + if err := watch.retryUnread(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if got := feedReads.Load(); got != 3 { + t.Errorf("read the feed %d times, want the retry's read too", got) + } + if out.Len() != 0 { + t.Errorf("wrote %q, want still nothing while the feed stays too busy", out.String()) + } + if watch.backoff != 2*firstWatchRetry { + t.Errorf("backoff = %v, want it doubled while the feed stays too busy", watch.backoff) + } + + // The feed quietens and a clean read ends the recovery with its one resync, + // at the last skip; the next time it falls behind is a recovery of its own. + retryClean := func() { + t.Helper() + quiet.Store(true) + watch.retry = nil + if err := watch.retryUnread(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } + quiet.Store(false) + } + retryClean() + if len(watch.unread)+len(watch.calendar.unread) != 0 { + t.Error("a clean read should leave nothing behind") + } + lines := watchLines(t, out) + if len(lines) != 1 || !strings.HasSuffix(lines[0]["change"].(string), "resync") { + t.Fatalf("wrote %v, want one resync for the whole recovery", lines) + } + skippedTo, err := time.Parse(time.RFC3339Nano, lines[0]["at"].(string)) + if err != nil { + t.Fatalf("resync at %v: %v", lines[0]["at"], err) + } + wantSkippedToHEYsClock(t, skippedTo) + ring() + retryClean() + if lines := watchLines(t, out); len(lines) != 2 { + t.Errorf("wrote %v, want a second resync for a second recovery", lines) + } + }) + } +} + +// HEY's ETag for /boxes.json is the box rows, which neither posting activity nor a new +// feed version touches, so a list the SDK revalidates answers 304 with the since the +// watch fell behind from and the version HEY now refuses — and HEY answers 409 for a +// version it no longer speaks as surely as for too many changes. The skip-ahead reads +// the list past the cache: the next read is on HEY's clock and on the new version. +func TestWatchSkipAheadReadsTheBoxListPastTheCache(t *testing.T) { + t.Setenv("HEY_TOKEN", "test-token") + t.Setenv("XDG_CACHE_HOME", t.TempDir()) + + behind := time.Date(2026, 8, 21, 11, 0, 0, 0, time.UTC) // a since before this is too far behind + var mu sync.Mutex + var notModified, conflicts int + listed := `[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-01T09%3A00%3A00.000000Z&v=2"}]` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + defer mu.Unlock() + switch r.URL.Path { + case "/identity.json": + w.Header().Set("Date", skipDate) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"id":1}`)) + case "/boxes.json": + w.Header().Set("ETag", `W/"boxes-unchanged"`) + if r.Header.Get("If-None-Match") == `W/"boxes-unchanged"` { + notModified++ + w.WriteHeader(http.StatusNotModified) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(listed)) + default: + since, _ := time.Parse(time.RFC3339Nano, r.URL.Query().Get("since")) + if since.Before(behind) || r.URL.Query().Get("v") != "3" { + // As HEY answers: `head :conflict`, no body. + conflicts++ + w.WriteHeader(http.StatusConflict) + return + } + w.Header().Set("Content-Type", "application/json") + w.Header().Set("Link", `<`+r.URL.Path+`?since=2026-08-21T11%3A06%3A12.250000Z&v=3>; rel="next"`) + _, _ = w.Write([]byte(`{"added":[{"id":9004,"kind":"topic","box_id":24088,"name":"Re: Lunch on Thursday?","active_at":"2026-08-21T11:06:12.250Z","created_at":"2026-08-21T11:06:12.250Z","creator":{"name":"Maria Delgado"}}]}`)) + } + })) + defer server.Close() + initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) + + // The box list as the watch read it at its start, now in the SDK's cache. Then HEY + // moves the feed to version 3, and the list's rows — its ETag — stay as they were. + if _, err := sdk.Boxes().List(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } + mu.Lock() + listed = `[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-21T11%3A04%3A58.310442Z&v=3"}]` + mu.Unlock() + + watch, out := newTestWatch("added", "resync") + cursor, err := watchCursor(server.URL+"/boxes/24088/postings/changes.json?since=2026-08-01T09%3A00%3A00.000000Z&v=2", "") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + watch.boxes[24088].cursor = cursor + + ringBox(t, watch) + + mu.Lock() + defer mu.Unlock() + if notModified != 0 { + t.Errorf("the skip-ahead read the box list from the cache %d times, want never — a 304 hands back the version HEY refused", notModified) + } + if conflicts != 1 { + t.Errorf("the feed answered 409 %d times, want once — the skip-ahead should get past it", conflicts) + } + if got := watch.boxes[24088].cursor.Version; got != "3" { + t.Errorf("version = %q, want the one HEY speaks now", got) + } + lines := watchLines(t, out) + if len(lines) != 2 || lines[0]["change"] != watchResync || lines[1]["change"] != "added" || lines[1]["posting_id"] != float64(9004) { + t.Fatalf("wrote %v, want one resync and then the change after it", lines) } - if got := watch.newMail.floors[24088]; !got.Equal(time.Date(2026, 8, 21, 11, 2, 0, 0, time.UTC)) { - t.Errorf("new-mail floor = %v, want the box's floor at the cursor it skipped to", got) + skippedTo, err := time.Parse(time.RFC3339Nano, lines[0]["at"].(string)) + if err != nil { + t.Fatalf("resync at %v: %v", lines[0]["at"], err) + } + wantSkippedToHEYsClock(t, skippedTo) + if floor := watch.newMail.floors[24088]; watchTime(floor) != lines[0]["at"] { + t.Errorf("new-mail floor = %v, want the skip point the resync names", floor) } } @@ -851,14 +1244,15 @@ func TestWatchDoesNotSayReadyOnItsWayOut(t *testing.T) { } } -func TestWatchedBoxesStartNoLaterThanTheWatchDid(t *testing.T) { +func TestWatchedBoxesStartAtTheWatchsStart(t *testing.T) { t.Setenv("HEY_TOKEN", "test-token") server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") // The Imbox's last activity is after the watch read HEY's clock — mail - // landed in between; The Feed's is before. - _, _ = w.Write([]byte(`[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-21T09%3A00%3A30.000Z&v=2"},` + - `{"id":24089,"kind":"feedbox","name":"The Feed","posting_changes_url":"/boxes/24089/postings/changes.json?since=2026-08-21T08%3A00%3A00.000Z&v=2"}]`)) + // landed in between; The Feed's is before, and its feed may still hold + // changes later than that which are history all the same. + _, _ = w.Write([]byte(`[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=2026-08-21T09%3A00%3A30.183221Z&v=2"},` + + `{"id":24089,"kind":"feedbox","name":"The Feed","posting_changes_url":"/boxes/24089/postings/changes.json?since=2026-08-21T08%3A00%3A00.518496Z&v=2"}]`)) })) defer server.Close() initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) @@ -868,11 +1262,10 @@ func TestWatchedBoxesStartNoLaterThanTheWatchDid(t *testing.T) { if err != nil { t.Fatalf("unexpected error: %v", err) } - if got := boxes[24088].cursor.Since; got != "2026-08-21T09:00:00.000Z" { - t.Errorf("Imbox cursor = %q, want it moved back to the watch's start so the mail in between is read", got) - } - if got := boxes[24089].cursor.Since; got != "2026-08-21T08:00:00.000Z" { - t.Errorf("Feed cursor = %q, want the box's own when it is earlier", got) + for _, id := range []int64{24088, 24089} { + if got := boxes[id].cursor; got.Since != "2026-08-21T09:00:00.000Z" || got.Version != "2" { + t.Errorf("%s cursor = %+v, want the watch's start and HEY's version, whatever since HEY served", boxes[id].name, got) + } } if !boxes[24088].reported || !boxes[24089].reported { t.Error("without --box every box is reported") @@ -894,7 +1287,7 @@ func TestWatchedBoxesStartNoLaterThanTheWatchDid(t *testing.T) { } command.boxes = nil - // --since is the reader's choice and wins over both. + // --since is the reader's choice and wins over the start. command.since = "2026-08-21T09:30:00Z" boxes, err = command.watchedBoxes(context.Background(), watchStarted) if err != nil { @@ -905,6 +1298,209 @@ func TestWatchedBoxesStartNoLaterThanTheWatchDid(t *testing.T) { } } +// heyHistory is HEY as a watch's startup meets it: its clock, the box list, and a +// changes feed that answers what is strictly later than its cursor, the way HEY's does. +// Each box's posting_changes_url carries what HEY puts there — the box's last posting +// activity, the latest updated_at among its unbundled postings, or the box's own +// updated_at when it has none — not the time, and a cached box list serves it as it +// was when cached. So the feed can answer a change later than that cursor that is +// history all the same. +type heyHistory struct { + mu sync.Mutex + cursors map[int64]string + changes map[int64][]historyChange +} + +type historyChange struct { + at string + deleted bool + id int64 + subject string +} + +// The clock the server answers with; a watch's start is taken a moment before it. +const heyHistoryDate = "Fri, 21 Aug 2026 09:00:05 GMT" + +func newHEYHistory(t *testing.T) *heyHistory { + t.Helper() + history := &heyHistory{ + cursors: map[int64]string{ + // The Imbox's since as a box list cached before its latest posting serves it. + 24088: "2026-08-20T23:45:00.000000Z", + // Reply Later is empty, so its cursor is the box's own updated_at — years + // before the posting deleted from it last week. + 24091: "2020-06-16T11:22:18.853469Z", + }, + changes: map[int64][]historyChange{ + 24088: {{at: "2026-08-20T23:45:39.083350Z", id: 9001, subject: "Lunch on Thursday?"}}, + 24091: {{at: "2026-08-19T01:07:15.496840Z", id: 9002, deleted: true}}, + }, + } + + server := httptest.NewServer(http.HandlerFunc(history.serve)) + t.Cleanup(server.Close) + t.Setenv("HEY_TOKEN", "test-token") + initSDK(auth.NewManager(server.URL, server.Client(), t.TempDir()), server.URL) + + return history +} + +// land is a change arriving in the Imbox, whose cursor is now its last activity. +func (h *heyHistory) land(change historyChange) { + h.mu.Lock() + defer h.mu.Unlock() + h.changes[24088] = append(h.changes[24088], change) + h.cursors[24088] = change.at +} + +func (h *heyHistory) serve(w http.ResponseWriter, r *http.Request) { + h.mu.Lock() + defer h.mu.Unlock() + + w.Header().Set("Content-Type", "application/json") + var boxID int64 + switch { + case r.URL.Path == "/identity.json": + w.Header().Set("Date", heyHistoryDate) + _, _ = w.Write([]byte(`{"id":1}`)) + case r.URL.Path == "/boxes.json": + _, _ = fmt.Fprintf(w, `[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"/boxes/24088/postings/changes.json?since=%s&v=2"},`+ + `{"id":24091,"kind":"laterbox","name":"Reply Later","posting_changes_url":"/boxes/24091/postings/changes.json?since=%s&v=2"}]`, + h.cursors[24088], h.cursors[24091]) + case scanBox(r.URL.Path, &boxID): + h.serveChanges(w, boxID, r.URL.Query().Get("since")) + default: + http.NotFound(w, r) + } +} + +func scanBox(path string, boxID *int64) bool { + _, err := fmt.Sscanf(path, "/boxes/%d/postings/changes.json", boxID) + return err == nil +} + +func (h *heyHistory) serveChanges(w http.ResponseWriter, boxID int64, since string) { + from, err := time.Parse(time.RFC3339Nano, since) + if err != nil { + w.WriteHeader(http.StatusBadRequest) + return + } + + var added, deleted []string + var last string + for _, change := range h.changes[boxID] { + at, _ := time.Parse(time.RFC3339Nano, change.at) + if !at.After(from) { + continue + } + last = change.at + if change.deleted { + deleted = append(deleted, fmt.Sprintf(`{"id":%d,"deleted_at":%q}`, change.id, change.at)) + } else { + added = append(added, fmt.Sprintf(`{"id":%d,"kind":"topic","box_id":%d,"name":%q,"created_at":%q,"updated_at":%q,"active_at":%q,"creator":{"name":"Maria Delgado"}}`, + change.id, boxID, change.subject, change.at, change.at, change.at)) + } + } + if last != "" { + w.Header().Set("Link", fmt.Sprintf(`; rel="next"`, boxID, last)) + } + _, _ = fmt.Fprintf(w, `{"added":[%s],"deleted":[%s]}`, strings.Join(added, ","), strings.Join(deleted, ",")) +} + +// startWatch begins a watch the way run does: HEY's clock, then the boxes, then the +// catch-up. +func startWatch(t *testing.T, command *watchCommand, watch *postingsWatch) { + t.Helper() + started := readServerNow(t) + boxes, err := command.watchedBoxes(context.Background(), started) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + watch.boxes = boxes + watch.newMail = trackNewMail(started) + if err := watch.catchUp(context.Background()); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestWatchDoesNotReportHistoryAsItStarts(t *testing.T) { + history := newHEYHistory(t) + watch, out := newTestWatch(defaultChanges...) + watch.exitOnFirst = true + + startWatch(t, newWatchCommand(), watch) + + lines := watchLines(t, out) + if len(lines) != 1 || lines[0]["change"] != "ready" { + t.Fatalf("wrote %v, want ready alone — a deletion and a posting from before the watch began are history", lines) + } + if watch.finished() { + t.Fatal("--exit-on-first should still be waiting for a change") + } + + // A reply lands after the watch began: that is the change it was waiting for. + history.land(historyChange{at: "2026-08-21T09:00:12.250000Z", id: 9003, subject: "Re: Lunch on Thursday?"}) + ringBox(t, watch) + + lines = watchLines(t, out) + if len(lines) != 2 || lines[1]["change"] != "added" || lines[1]["posting_id"] != float64(9003) || lines[1]["new"] != true { + t.Errorf("wrote %v, want the reply that landed after the start, new", lines) + } + if !watch.finished() { + t.Error("--exit-on-first should end on the change that landed after the start") + } +} + +func TestWatchReadsMailThatLandedBeforeItReadTheBoxes(t *testing.T) { + history := newHEYHistory(t) + // The watch reads HEY's clock a moment before the Date header; this lands on the + // header's second, after the start and before the box list, so the Imbox's cursor is + // already past it. + history.land(historyChange{at: "2026-08-21T09:00:05.000000Z", id: 9003, subject: "Invoice #4021"}) + watch, out := newTestWatch(defaultChanges...) + + startWatch(t, newWatchCommand(), watch) + + lines := watchLines(t, out) + if len(lines) != 2 || lines[0]["posting_id"] != float64(9003) || lines[0]["new"] != true || lines[1]["change"] != "ready" { + t.Errorf("wrote %v, want the mail that landed during startup, new, then ready — and no history", lines) + } +} + +// HEY's clock is read to the whole second, so the watch cannot tell a change from the +// Date header's own second that came before it from one that came after, and it reads +// both: skipping one that came after would be worse than repeating one that did not. +// Anything earlier than the second the watch read — less the request's time — is +// behind the start and not reported. +func TestWatchReadsTheWholeSecondHEYsClockWasReadIn(t *testing.T) { + history := newHEYHistory(t) + history.land(historyChange{at: "2026-08-21T09:00:03.990000Z", id: 9003, subject: "Invoice #4021"}) + history.land(historyChange{at: "2026-08-21T09:00:05.200000Z", id: 9004, subject: "Re: Invoice #4021"}) + watch, out := newTestWatch(defaultChanges...) + + startWatch(t, newWatchCommand(), watch) + + lines := watchLines(t, out) + if len(lines) != 2 || lines[0]["posting_id"] != float64(9004) || lines[0]["at"] != "2026-08-21T09:00:05.200Z" || lines[0]["new"] != true || lines[1]["change"] != "ready" { + t.Errorf("wrote %v, want the change from the Date header's second, with when it happened, and not the one before it", lines) + } +} + +func TestWatchSinceReadsTheHistoryFirst(t *testing.T) { + newHEYHistory(t) + command := newWatchCommand() + command.since = "2026-08-19" + watch, out := newTestWatch(defaultChanges...) + + startWatch(t, command, watch) + + lines := watchLines(t, out) + if len(lines) != 3 || lines[0]["posting_id"] != float64(9001) || lines[0]["new"] != false || + lines[1]["change"] != "deleted" || lines[1]["posting_id"] != float64(9002) || lines[2]["change"] != "ready" { + t.Errorf("wrote %v, want both changes since --since, not new, then ready", lines) + } +} + func TestWatchAnnouncesADropBeforeTheReconnectThatFollowedIt(t *testing.T) { server := changesServer(t, `{}`) watch, _ := newTestWatch("added") @@ -998,10 +1594,10 @@ func TestWatchReportsAResyncAfterSkippingAhead(t *testing.T) { var server *httptest.Server server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if strings.Contains(r.URL.Path, "/postings/changes") { - // As haystack answers: `head :conflict`, no body. - w.WriteHeader(http.StatusConflict) + answerTooFarBehindBefore(w, r, time.Date(2026, 8, 21, 11, 0, 0, 0, time.UTC)) return } + w.Header().Set("Date", skipDate) w.Header().Set("Content-Type", "application/json") _, _ = w.Write([]byte(`[{"id":24088,"kind":"imbox","name":"Imbox","posting_changes_url":"` + server.URL + `/boxes/24088/postings/changes.json?since=2026-08-21T12%3A00%3A00.000Z&v=2"}]`)) })) @@ -1026,8 +1622,13 @@ func TestWatchReportsAResyncAfterSkippingAhead(t *testing.T) { if event.Change != watchResync || event.Box == nil || event.Box.ID != 24088 { t.Errorf("event = %+v, want a resync for the Imbox", event) } - if watch.boxes[24088].cursor.Since != "2026-08-21T12:00:00.000Z" { - t.Errorf("cursor = %+v, want it moved to the server's current one", watch.boxes[24088].cursor) + skippedTo, err := time.Parse(time.RFC3339Nano, event.At) + if err != nil { + t.Fatalf("resync at %q: %v", event.At, err) + } + wantSkippedToHEYsClock(t, skippedTo) + if watch.boxes[24088].cursor.Since != watchStartSince(skippedTo) { + t.Errorf("cursor = %+v, want it moved to HEY's clock at the skip the resync names", watch.boxes[24088].cursor) } } diff --git a/skills/hey/SKILL.md b/skills/hey/SKILL.md index 4ed1d5ee..32e67f4e 100644 --- a/skills/hey/SKILL.md +++ b/skills/hey/SKILL.md @@ -714,14 +714,18 @@ stdout, one per line, instead of the usual envelope (at a terminal, one text lin "name"}, "posting_id": ..., "thread_id": ..., "new": true|false, "posting": {...}}`. Use `thread_id` with `hey thread read` (absent for a bundle row that names no single thread). `new` is on every `added` and `updated` line and says whether the posting is new mail — unseen, not muted, and active since the watch last saw the thread, -or since the watch began for a thread it has not seen; the backlog a watch starts with is +or since the watch began for a thread it has not seen; the backlog `--since` reads is never new, nor is reading, muting or moving a thread, and a reply on a known thread is. +Without `--since`, nothing from before the watch began is reported, save a change from up to +a second before it plus however long reading HEY's clock took (that clock is read to the whole +second, and taken back by the request's time); its `at` says when it happened. `--events new` selects the new ones, alone or in a union with the other three. A deleted posting carries no `posting`, `thread_id` or `new`. Three more lines describe the watch itself: `{"change": "ready"}` once every box and calendar is caught up and the subscription is live (again after every reconnect's catch-up), `{"change": "disconnected"}` when the connection drops, and `{"change": "resync", "box": {...}}` when a box changed more than the -feed can list and the watch skipped ahead — re-read that box. A resync is an event of its +feed can list and the watch skipped ahead — re-read that box; one resync covers the whole +catch-up, however many skips it takes, and comes once the box is followed again. A resync is an event of its own: reported by default (`--run-*` scripts run for it, `--exit-on-first` counts it) and left out by `--events new`, so a script for new mail never runs on one. `ready` and `disconnected` carry no `box`, are written only when no `--run-*` command is given, and never count for