mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-08 15:41:15 +02:00
* fix(filer): don't skip unflushed events on metadata subscription gaps A subscriber that falls behind the in-memory log ring during a write burst could have its read position jumped past events that were evicted but not yet flushed to disk — silently lost for filer.backup/filer.sync/ mount subscribers. Route all three gap-skip sites through resolveDiskGapResume: only skip past windows older than a settled horizon (2*LogFlushInterval); recent gaps wait for the flush and re-read disk. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(filer): harden metadata gap-skip guard Address review findings on the settled-horizon guard: - local subscriptions: gate the skip on the buffer's flush watermark observed before the disk read (resolveLocalGapResume) — a disk miss is then proof the gap is empty, with no wall-clock assumptions - aggregated subscriptions: cap skips strictly below the horizon boundary (persisted reads exclude ts <= cursor) and pace capped advances, so the sliding horizon cannot cause disk-probe spinning - replace unbounded sync.Cond waits with a bounded select on the buffer's subscriber channel + retry timer + ctx cancellation, eliminating the lost-wakeup stall Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(filer): close remaining gap-skip loss paths from re-review - log_buffer: a sentinel-offset (time-based) read below the earliest in-memory entry silently started at earliest, skipping a window that may hold evicted-but-unflushed events. Track the ring's eviction watermark (lastEvictedTsNs) and keep the inclusive fast path only when nothing at/after the position was ever evicted; otherwise return ResumeFromDiskError so the subscription gap guard decides. - capped horizon advances stay on the disk-probe path (never expose a mid-gap position to the memory read) and keep pacing - gap jumps land just below earliest: positions are exclusive, so the earliest entry itself is still delivered - subscriber notification keys include clientId/epoch so a replacement stream never inherits a channel the old stream's cleanup closes Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(log_buffer): generalize the eviction-watermark read gate Third-review findings: - epoch/zero-time reads bypassed the eviction watermark: gate them the same way, so a SinceNs=0 subscriber cannot silently start at the earliest retained entry after an unflushed window was evicted - apply the watermark gate regardless of the cursor's batch offset (batch offsets carry no meaning for time-based reads); this also serves adjacent cursors (earliest == current+1) from memory instead of stalling them in the gap loop - LoopProcessLogData reader names include clientId/epoch, since they are registered as subscriber keys internally (same collision as the outer notification keys) - test: pin the gate (below/at watermark, epoch-after-eviction) and update the slow-consumer test to the sharper contract — complete in-memory history is served from memory; disk only once evicted Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(log_buffer): watermark equality is unsafe for inclusive sentinel cursors Sentinel (Offset <= 0) time-based cursors search from ts-1ns, i.e. they read inclusively of their own timestamp — and the evicted window may end exactly at that timestamp. Allow watermark equality only for exclusive (positive-offset) cursors; sentinel cursors must be strictly above it. Also shut down the test buffer and pin the equality cases. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix(filer): wait on an empty aggregated buffer instead of falling through When a disk read finds nothing for a ResumeFromDiskError gap and the aggregated buffer has no readable entries (zero earliest time), resolveDiskGapResume declines to advance and control fell through to the in-memory read — which returns ResumeFromDiskError again immediately, spinning through the probe cycle without any wait. Fold the case into the existing recent-gap branch so it waits (notification, cancellation, or the retry interval) before re-probing, matching the local subscription path, which already waits unconditionally. * fix(filer): serve the exact eviction-boundary entry without another flush wait When the flush watermark has passed the earliest in-memory entry but that entry sits exactly one nanosecond above the cursor, the exclusive resume target collapses onto the cursor and resolveLocalGapResume declines to advance — while a sentinel cursor at the eviction watermark keeps deferring the in-memory read to disk. Progress then depends on the next flush cycle. Re-arm the cursor with a positive (exclusive) offset instead: ReadFromBuffer explicitly allows positive-offset cursors at the watermark, so the boundary entry is served from memory immediately. Also promote the aggregated path's horizon-capped gap skip to a warning: that skip may pass events a stalled flush lands later (the aggregated ring has no flush watermark to gate on), so operators should see when a flush stall outlasts the settled horizon. * fix(filer): resume a gap skip with an exclusive cursor The resume position landed one nanosecond below the earliest in-memory entry but kept the inclusive sentinel offset. When that timestamp is also the eviction watermark, the read gate answers ResumeFromDiskError for an inclusive cursor, the disk read finds nothing, and the resume target collapses onto the cursor, so neither helper can advance: the subscriber parks on timed waits forever. A skip is only taken once the gap is proven empty, so nothing remains to deliver at the resume timestamp. Resume exclusively instead, and let an inclusive cursor already sitting on the target count as progress. * fix(filer): gate aggregated gap skips on eviction, not wall clock Wall-clock age never proved persistence. The aggregated ring has no flush watermark because peers persist their own local logs, so a horizon of 2 * LogFlushInterval was standing in for one. While the disk stayed empty the horizon kept sliding forward and walked the cursor past evicted-but-unflushed events one window at a time - the volume-outage stall this change set exists to survive. The ring does carry a real proof: its eviction watermark. If nothing at or after the cursor was ever dropped, memory still holds every entry after it and the gap is provably empty. Below the watermark entries were dropped and only the producing peer's flush can supply them, so wait for that flush instead of advancing. The skip that remains is the one the ring can prove, which keeps the infinite-loop guard for a genuinely empty gap. * perf(filer): stop re-probing disk on every metadata append A subscriber parked on a gap woke on the log buffer's subscriber channel, which fires on every append. On the aggregated path each wake re-ran a full ReadPersistedLogBuffer - a ListDirectoryEntries plus a readahead goroutine - so one parked subscriber turned every cluster-wide metadata event into a store query, exactly while the cluster is already struggling with the flush stall that parked it. An append also cannot settle the gap: what this waits on is a peer persisting its own log, which nothing local signals. Wait on the timer alone there. The local buffer does signal its flush on that channel, so keep it, but drain the stale token first: otherwise the appends riding the same channel spin the wait at write rate. The retry interval covers the notification the drain discards. * fix(filer): surface a metadata subscriber parked on a gap Refusing to skip an unresolved gap trades silent loss for a silent stall, and a stall is no easier to diagnose: filer.sync and mount followers just stop advancing, with no error on either end. The only trace was a V(3) line nobody runs with. Warn on entry and once a minute after, and carry a subscribe_gap_stalled gauge for the duration, so a flush that never lands shows up as a stalled subscriber rather than a consumer that mysteriously went quiet. * refactor(filer): drop the metadata listener cond with no waiters left Both listenersCond.Wait() sites are gone, so listenersWaits never leaves zero, the guard in the filer's notify callback never fires, and the two Broadcast calls around client registration wake nobody. Remove the cond, its counter and its lock, and pass a nil notify func: the log buffer already skips a nil one, and the subscriber channels carry these wakeups now. * test(log_buffer): drive the eviction watermark through a real seal Both eviction tests wrote lastEvictedTsNs directly, leaving the one line that sets it uncovered: copyToFlushInternal has to read slot 0 before SealBuffer shifts it out, and moving that read one statement later still passed every test. Append past the ring instead, and check the watermark appears only on the seal that drops a window, matches that window's stop, and then advances. The buffer under test never flushes, matching the aggregated meta ring where the watermark is the only emptiness proof. * refactor(filer): build the subscriber reader name once per stream It was rebuilt on every loop iteration from values that cannot change for the life of the stream. Hoist it next to the notification key it mirrors. * docs(filer): tighten the gap-handling comments Several ran to eight lines restating the same reasoning at each site. Keep the non-obvious why, drop the retelling. * fix(filer): mark a confirmed disk position as an exclusive cursor A disk read returns the timestamp of the last entry it handed to the subscriber, but the cursor built from it stayed inclusive. Land that cursor exactly on the eviction watermark and the read gate sends it back to disk for an entry the disk just delivered; the re-read finds nothing, neither emptiness proof holds for an inclusive cursor there, and the subscriber parks - until some later flush, or forever while one stays stalled. Everything after the watermark was sitting in memory the whole time. Carry the offset instead, so it also covers the ring evicting onto the cursor after the read rather than before it. The constant is no longer gap-specific, so it is now named for what it asserts. * perf(filer): re-park a gap wait woken by an ordinary append Draining one stale token did not bound anything: the subscriber channel carries an append per metadata write, so under continuous writes the next one satisfies the wait immediately. During the write burst these stalls come from, each parked subscriber ran the recovery loop at write rate rather than the intended two-second cadence. Only a flush can settle a gap, so check the flush watermark on wake-up and go back to waiting if an append is all that arrived. The retry timer is created once, so re-parking does not extend the interval. * fix(filer): keep a disk-derived cursor inclusive Marking every disk position exclusive assumed the entry at that timestamp was the only one there. On the aggregated stream it is not: disk can deliver one filer's persisted event at T while another filer's event at the same T is still unflushed in the aggregate ring. The exclusive cursor skipped it, and since the persisted reader also excludes ts <= its start, nothing would ever bring it back. Take the weaker guarantee instead. The reason the exclusive cursor was introduced - an inclusive one landing on the eviction watermark parks forever - is better answered in the resolvers: at the watermark memory holds nothing (retained windows start strictly after it) and the persisted reader cannot return that entry at any later time either, so refusing the gap buys nothing and never ends. Skip it whatever the cursor's inclusivity. That proof does not depend on a flush, so the local resolver takes it as a second, independent disjunct alongside its flush watermark. * perf(filer): wake a parked gap wait on flushes only Re-parking on an append bounded the work but not the wake-ups: the subscriber channel carries one per metadata write, so a parked subscriber still took a scheduling round-trip per write. Worse, draining it kept the channel empty, so every writer's non-blocking notify succeeded instead of falling through - the burst paid for the wake-ups too. Give the log buffer a flush-only subscriber list, notified from loopFlush, and park on that. The append channel now fills once and stays full, which is exactly the state the non-blocking send is designed for. * fix(filer): stop the gap-stall gauge from leaking a series per connection clientName is req.ClientName + "@" + peer address, so it carries the client's ephemeral source port and changes on every reconnect. Labelling the gauge with it minted a new series per connection, and clearing it only ever Set(0), so nothing was ever released - a client in a reconnect loop grows the filer's metric map and /metrics payload without bound. Key it on the stable client-supplied name, as the neighbouring subscribe gauge already does, and delete the series on teardown. Two logging fixes ride along, both in the same reporter: the resume warning was unpaced while the park warning throttles to one a minute, so a burst that parks and resumes every couple of seconds warned on every cycle; and clear() doubled as the teardown path, announcing "resumed after 14m0s parked" for a client that actually gave up and disconnected still behind. * fix(filer): give a parked gap wait the exits the read loop has A park never re-enters the read loop, so every exit that loop relies on stopped working while a subscriber was parked. It kept scanning the filer store every 2s for a client that a higher-epoch reconnect had already superseded, until the TCP connection finally died - hours, on a half-open one. It never reached the only code that honors UntilNs, so bounded callers like `weed shell fs.verify` and `filer.meta.tail -until=` hung instead of exiting. And a notification channel closed out from under it turned the bounded wait into a spin, since a receive on a closed channel returns instantly. Check all three where the subscriber actually waits. Bound the wait too: waiting is productive while the window is still queued for flush, but a peer that never returns, or whose filer store this filer cannot read at all, makes it permanent - and a subscriber that silently stops delivering is no better than one that silently skips. Fail the stream after that instead of hanging, loudly enough to say which gap and for how long. * fix(log_buffer): keep the eviction gate out of the shared read path The gate belonged to the filer's subscribe loops but was installed in ReadFromBuffer, which the message queue shares and which has no gap handling of its own. Three MQ paths broke on it: GetUnflushedMessages asks for everything in memory past the flush watermark and got ResumeFromDiskError instead, so SQL results silently dropped unflushed rows; a RESET_TO_EARLIEST consumer's epoch cursor was sent to disk, and the MQ disk reader resets an empty read back to epoch, so the whole partition replayed on every pass; and a disk cursor landing on the watermark carried offset -2, which the gate refused and the disk reader's ts <= start filter also skips, so neither side could ever serve it. The rewrite also dropped the old Offset <= 0 requirement, letting a stale positive-offset cursor jump to memory over on-disk history, and refused any negative-timestamp cursor even with nothing evicted. Restore the read path exactly as it was and keep only lastEvictedTsNs, which is the useful primitive. The filer loops now consult the watermark themselves before reading memory, which is where the gap handling that makes the refusal actionable already lives - and they check it every pass, not just when the disk came up empty, since a disk read can leave the cursor short of the watermark too. * fix(filer): read the log file whose window spans past its own name A log file is named for the start of the window it holds, but a window runs up to a flush interval longer, so "12-30" can hold entries through 12:31:58. File selection compared the cursor's minute against that name and skipped anything sorting earlier, so a subscriber resuming at 12:31:10 never saw the rest of that file: the read reported nothing on disk while the entries sat in it. That was survivable when a miss only meant "wait and retry", but the gap resolvers now read a miss as proof the range is empty and move the cursor past it, which turns those entries into silent loss. Start the file scan a flush interval early; entries are still filtered against the exact cursor, so this only opens one more file and never re-delivers. * fix(filer): resume a gap with a cursor the memory read will serve Reverting the eviction gate put ReadFromBuffer back to refusing every positive-offset cursor below the in-memory window, but the resolvers still handed back one - earliest-1 marked exclusive. So the resume bounced straight to ResumeFromDiskError, the disk had nothing, and the resolver saw its own cursor as no progress and parked: a subscriber stalled, and after the new bound failed outright, with the whole gap sitting in the ring the entire time. The exclusive marking only existed to dodge a park the gate itself caused, and the gate is gone. Resume at earliest with the sentinel offset, which is the position master used and which case 2.1 reads inclusively. That makes the resume unconditionally ahead of the cursor, so the progress check it needed goes away with it. * fix(filer): stop dropping a log file whose window outruns its name Listing the earlier file was not enough: the iterator then decided whether to read it by comparing the *following* file's name against the cursor, which treats that name as an upper bound on this file's contents. It is not one. A file is named for the start of the window it holds, minute-truncated, and the window runs up to a flush interval longer, so "12-30" can hold an event at 12:31:20 while "12-31" sits right after it - and a cursor at 12:31:10 skipped straight past the event. Bound the decision on the file itself: skip it only when its name plus the minute truncation plus a flush interval still lands at or before the cursor. That also covers the last file in the queue, which the old check never skipped because it had no successor to compare against. * fix(filer): stop a chunk-ref read from rewinding the subscriber CollectLogFileRefs reports the minute-level name of the last file it shipped, and the caller assigns that straight to the read position. Since the scan now reaches back a flush interval to catch a spanning file, a request at 12:31:10 that picks up the 12-30 ref moved the cursor to 12:30:00 - so the memory read that followed replayed events older than the client's own SinceNs and re-sent what the chunk reader had already been handed. The same rewind was reachable before, within a minute, whenever the last file's name sorted behind the request. Clamp the reported position to the one that was asked for. It still under-advances by design, since the server never reads the entries it ships refs for and cannot know where they end. * fix(filer): resume below earliest so a one-entry window is not skipped Moving the resume onto earliest itself was wrong for the smallest window there is. A sealed window holding a single entry has startTime == stopTime, and the sealed-buffer lookup only enters a window whose stopTime is strictly after the cursor, so a cursor sitting exactly on earliest walks past it and its sole event is never delivered. Low-volume metadata windows are routinely one entry, which is precisely when losing it is hardest to notice. Go back to one nanosecond below, which takes the startTime.After branch and returns the whole window, and keep the sentinel offset the memory read requires. The earlier test used an active multi-entry buffer, where both cursors happen to work; the new one seals single-entry windows and compares which entry comes back. * fix(filer): count a persisted read as progress only when it moves A chunk-ref read reports the minute-level name of the last file it shipped, now clamped so it never rewinds, so it comes back non-zero even when it names the position the subscriber already held. The loops read non-zero as progress: they cleared the stall timer, then found the cursor still short of the eviction watermark and parked again. Every retry re-shipped the same refs and reset the timer, so the bound that is supposed to end an unrecoverable stall was never reached - and a chunk-capable client buffers those refs waiting for an event that never comes, so it just accumulates duplicates. Require the reported position to be strictly ahead of the cursor. * fix(filer): count the evicted ranges the aggregated stream cannot prove The eviction watermark belongs to the merged ring, but the disk it gets checked against is the union of every peer's own log and each peer flushes on its own schedule. A read that lifts the cursor from below the watermark to above it may have done so entirely on a peer that is already ahead, while a lagging peer still holds unflushed events inside the range just crossed; when it flushes them they sit behind the cursor and are never delivered. One aggregate maximum is not proof that every peer persisted the range. Nothing available locally separates that from the ordinary case where every peer had in fact persisted it: the aggregator tracks peers by address while log files carry a random per-filer id, so "has this peer flushed through T" cannot be answered here at all. Deciding it needs the source filer's own flush watermark carried on the subscribe stream, which is a wire change this does not make. Count and log the crossing so the window is at least measurable instead of invisible. * fix(filer): send log file refs through the pipelined sender sendLogFileRefs wrote on the raw gRPC stream while pipelinedSender's goroutine concurrently calls Send for queued memory events - two senders on one stream, which gRPC forbids. The window used to open once per fall-behind; the gap retry loop now reopens it every pass. Routing refs through the sender also restores ordering: refs used to overtake up to 1024 queued events, and the client treats any non-ref message as the signal to process buffered refs, so overtaken refs were applied against the wrong position. The reason refs bypassed the sender was the batcher: their TsNs of 0 reads as far behind, and the client recognizes refs by the top-level field alone - a refs envelope would drop its Events tail, and refs inside Events would be applied as an empty event. Teach the batcher instead: refs messages always go solo, and one drained mid-batch is sent solo right after that batch. * fix(filer): count parked subscribers instead of flagging them by name The per-client gauge series could not work. Its label was rebuilt from the peer address at first, which leaks a series per reconnect; keyed on the client-supplied name instead, it collides - every mount registers as "mount" - so one stream's teardown deleted a parked sibling's live series, and the sibling never re-created it because its own park state said the gauge was already set. Either way the alert this gauge exists to drive goes dark. A count needs no identity: Inc on park, Dec on resume or teardown, scope as the only label. Client details stay in the logs. Also start warning only once a stall has outlived the warn interval. park() warned immediately on every first park, so catch-up churn that parks and resumes every couple of seconds logged a warning pair per cycle - exactly the flood the pacing was supposed to prevent, burying the long-stall warnings that matter. * fix(filer): one gap resolver, and re-arm an unservable adjacent cursor The two resolvers were the same function - the aggregated one is the local one with a flush watermark of zero, since its ring never flushes - duplicated down to the comment justifying the resume target. A fix applied to one and not the other is how the two streams drift; merge them. The merge carries the one behavioral fix both copies needed. Timestamp collision bumps make adjacent entries exactly 1ns apart, so an entry ending an evicted window leaves the cursor exactly one below earliest with a positive batch offset. The resume target then equals the cursor and both copies refused it as no progress - but that cursor cannot be served (ReadFromBuffer refuses positive offsets below the window) while the sentinel resume at the same timestamp is, and both deliver exactly the entries after it. Refusing parked a subscriber whose data was entirely in memory: until the next flush locally, and through a 15-minute stall failure on the aggregated path. Advance on equal target when the held cursor is exclusive; a sentinel there is already served, so it still refuses. * fix(filer): a bounded subscription parked exactly on UntilNs is finished The park's UntilNs check was strict while the bound is inclusive and cursors are exclusive: a disk read whose last entry sits exactly on UntilNs leaves the cursor there with everything up to the bound already delivered. If the next range was an unprovable gap, the completed subscription parked anyway and eventually failed - fs.verify hanging and then erroring on a healthy cluster. * fix(filer): re-derive the subscribe loops as one state machine The loops had grown three generations of gap handling - a post-disk-read branch, a post-memory-read branch, and an every-pass guard bolted on in front of the memory read - each consulting state the others mutated. The worst interaction wedged the aggregated stream permanently: diskExhausted compared the disk result against the cursor that result had just updated, and paired it with a ResumeFromDiskError latch that only a resolver advance cleared, so one fall-behind sent every later pass into the gap block and LoopProcessLogData never ran again. Mounts kept a healthy-looking stream and applied nothing. Both loops now run the same derived sequence. One disk pass; progress is the pre-update cursor against the result, freshly each pass. One gap decision before the memory read: a cursor the ring evicted past either keeps draining the disk (it just advanced), resolves forward (the gap is proven empty), or parks - and a cursor memory refused with nothing evicted after it re-arms onto the retained window. The stale-latch branch is gone: the error is only consulted on a pass whose own disk read came up empty. All four parks go through one parkOnGap helper, so the exits live in one place. The eviction watermark is read after the disk read and the same value feeds both the guard and the unproven-crossing report, which previously compared against a snapshot taken before a potentially minutes-long backlog read and missed evictions landing during it. The next-day jump now clears the stall reporter - it used to leave a stale park epoch that could kill the next brief park instantly at the 15-minute bound - and counts its own watermark crossing. Chunk-ref reads get two rules the old loops lacked. Refs are sent once per position: retries re-sent identical batches every two seconds into a client that only drains them on a non-ref message, growing an unbounded pending list. And when refs cannot advance the cursor - they report file start minutes, which sit below the content the client was actually given, pinning the cursor under the watermark forever during bursts - the pass falls through to entry reads, which move the cursor by real timestamps and double as the client's drain signal. * fix(filer): give up on an unprovable gap instead of failing the stream Failing after maxGapStall assumed the client could do something better, but every consumer just reconnects at the same SinceNs and hits the same wall, so an unprovable gap - a dead peer, or a peer whose filer store this filer cannot read at all - turned into a permanent 15-minute fail/reconnect loop delivering nothing. Master handled the same state by skipping instantly and silently. Take the middle: wait the full bound, then abandon the gap and resume at the eviction watermark, where everything retained starts strictly after, so the loss is exactly the range that could not be proven. The skip shares the unproven-crossing counter and logs at error level - loss is bounded, recorded, and the stream keeps working. A stall with nothing evicted past the cursor loses nothing by waiting, so it restarts the clock and keeps parking rather than skipping. * test(filer): pin the file-skip bound through the production predicate The spanning-file test asserted against its own copy of the arithmetic, so a regression in the iterator - restoring the next-file-name comparison or dropping the flush-interval term - would keep CI green while re-introducing the silent loss the fix closed. Extract the bound into logFileMayContainAfter, call it from the iterator, and point the test at it; breaking the production expression now fails the test. * test(log_buffer): pin the flush-subscriber contract The registry the filer's gap parks wait on had no test at all. Cover the observable contract: an append never wakes a flush subscriber, a flush does with the watermark already stored, unregistering closes the channel so an abandoned waiter unblocks, and double or unknown unregisters are harmless. The store-before-notify ordering in loopFlush is what makes the parks' wake-up re-check sound, and it is not black-box testable - reordering leaves a same-goroutine window of nanoseconds that hundreds of tight round-trips never catch. Mark it load-bearing at the site instead; a reorder now at least has to argue with the comment it deletes. * fix(filer): a -1 SinceNs is a position, not the refs-gate sentinel The once-per-position refs gate used -1 as "never sent", but a client may legally subscribe with SinceNs=-1, whose cursor timestamp is exactly -1: the very first pass then believed refs were already sent there and fell back to streaming the whole persisted history entry by entry - the bootstrap load chunks mode exists to avoid. Use MinInt64, which no cursor can carry. * refactor(filer): drop the aggregator's listener cond with no waiters left Same shape as the FilerServer cond already removed: nothing increments ListenersWaits and nothing ever calls Wait, so the three Broadcasts wake nobody and the notify callback's guard is always false. Aggregated subscribers wake through the buffer's subscriber channels now. * fix(filer): make the gap metrics say what they count The crossing counter's help text described only the aggregated peer case, but give-ups increment it for local stalls too - a wedged local flush - which sends an operator chasing peer replication when the problem is the local store. Label it by scope and say both. The stalled gauge counted every park, including waits with nothing evicted and nothing at risk, while its help text promised evicted-but-unpersisted events; describe it as what it is, a count of subscribers parked on a gap. * test(filer): make the flush and stall tests assert what they claim The flush-subscriber rounds were vacuous: the probe entry sat above the round timestamps, so every round entry was collision-bumped and the stored watermark exceeded the local value each assertion compared against - the same bump mistake this test suite already made once. Put the probe below the rounds and guard each round against bumping, so a vacuous setup fails instead of passing. The stall-outcome test wrote the reporter's park epoch directly, bypassing the gauge Inc that gaveUp() later Decs - leaving the shared process gauge at -1 for every test that runs after it. Park through the real path, age the park by hand, release what the test holds, and assert the gauge lands back where it started. * fix(filer): finish the stream checks before marking it parked parkOnGap stamped the reporter before waitOnGap ran its instant done exits, so a bounded subscription completing inside a gap state was marked parked for the microsecond before done fired - a phantom gauge blip and a false "disconnected still behind" warning on every healthy completion. Fold waitOnGap into parkOnGap so the done exits run first and the park mark only ever covers a stream that actually waits. While the park owns its timer, back the retry off as the stall ages - 2s probes growing toward one a minute - since every retry re-reads the persisted log, and probing the store each 2s for 15 minutes per parked subscriber during the very outage that parked them makes the bad time worse. The subtest still named for the old fail-the-stream stall behavior goes with the merge. * fix(log_buffer): gate filer cursors against eviction under the read lock The subscribe loops checked the eviction watermark and then read memory, but a seal can land between the two: the read then served a sentinel cursor from the earliest retained window, silently skipping the window just evicted - the loss class this PR exists to make loud, surviving as a race. The only place the check is atomic with the serve decision is inside ReadFromBuffer, under the lock seals take to evict. Rather than put the policy back into the shared read path - which broke four message-queue readers last time - add a new sentinel offset that opts into it: EvictionGatedOffset reads inclusively exactly like -2, except below the watermark it is refused to disk. The filer loops stamp it on every cursor they hand the memory read; a refusal lands in the same gap machinery the loop-side check feeds, so the race collapses into the handled path. MQ cursors never carry it and keep master behavior byte-for-byte. * fix(filer): gate the aggregated gap on received, not bumped, timestamps The aggregated ring rewrites an out-of-order arrival to its head plus a nanosecond, so after any bump-heavy interval - a peer history replay following a restart is enough - its eviction watermark lives above every timestamp that exists on any peer's disk. Comparing a disk cursor against it parked subscribers that had in fact drained every peer's log: a 15-minute delivery freeze ending in a give-up skip and a false loss alarm, on a healthy cluster where master resumed instantly. Track a second watermark in the received timestamp space - the highest pre-bump timestamp among evicted entries - and gate the aggregated loop on that. Disk cursors and received timestamps are the same space, so the comparison means what it says: at or past it, every evicted entry's original was at or below the cursor, and everything flushed of them was already delivered. The bumped watermark keeps guarding the in-ring read gate, whose cursors live in ring space. The local buffer is untouched: it flushes its own bumped timestamps, so there the two spaces are one. * fix(filer): ship each log chunk once and stop echoing ref'd files inline Chunk mode duplicated data through two doors. Consecutive ref collections overlap by design - the scan backs off a flush interval to catch a spanning file, and a filer appends chunks to its newest file - so the same file was shipped again on every pass that re-listed it: the client re-downloaded its chunks, and a duplicated file mid-batch rewinds timestamps inside the client's per-filer merge, which reads each stream as sorted - transiently resurrecting deleted entries during catch-up. Track per subscription how many chunks of each file were shipped and send only the unsent suffix; state prunes with the scan window, so it holds a few files per filer. The second door was the entry fallback: when refs cannot advance the minute-named cursor, the pass streamed the ref'd file's tail inline, and the client applies inline events unfiltered - the same tail it already applied from chunks. Entry passes for chunk clients now advance the cursor without delivering; everything they skip is covered by the refs already sent or the deltas the next collection ships. * refactor(filer): one gap decision shared by both subscribe loops The post-disk gap tree - guard, drain, resolve, two parks, the re-arm - existed twice, differing only in buffer, watermark space, flush getter, park channels, and reason strings. Four rounds of review fixes have shown the copies drift the moment one is edited alone. gapPass now carries the five differences and the tree lives once; the loops shrink to a three-way switch between reading memory, restarting the pass, and ending the stream. * docs(filer): trim the gap-machinery comments to the why Several blocks had grown to ten-plus lines restating what the tests already pin or retelling one rationale at multiple sites. Keep the non-obvious why - the load-bearing flush ordering, the two timestamp spaces, the refusal-at-equality argument - in a few lines each. * fix(log_buffer): credit an entry's received timestamp to its own window The received-ts capture ran before the rollover check, so an append that sealed the previous window stamped its timestamp onto that window and then lost it in the reset of the new one. The eviction watermark this feeds broke both ways: the sealed window's value was inflated by an entry it does not contain - parking aggregated subscribers on gaps that were drained - and the entry's real window was deflated, proving gaps empty that still held its event on some peer's unflushed path. Credit the timestamp only after the entry lands, when its window is known. * fix(filer): rebase a shipped chunk suffix to logical offset zero A grown file's delta kept the chunks' original file offsets, but the client's chunk reader starts at logical zero and a list opening higher reads as instant EOF - a successfully empty replay, and since chunk clients no longer receive disk entries inline, the appended events were silently dropped. Clone the suffix chunks with offsets rebased to zero; the cut is record-aligned because each append is one uploaded chunk of whole entries, so the suffix decodes as a file of its own. * fix(filer): finish every chunk refs batch with a transition the client acts on Both chunk consumers buffer refs until a non-ref message arrives, so a source with historical logs and a quiet ring - a mount reconnecting after a filer restart is the common case - shipped its backlog and then went silent: the client sat on the refs until the next metadata mutation anywhere in the cluster. The disk step now ends every batch with the empty-notification marker the client already treats as a resume-cursor advance. The same step closes the inline replay: the cursor used to stay at the last file's minute name, so the memory read re-delivered the retained tail of a file the client had just read via chunks - T1..Tn applied twice. The advance-only entry read now runs on every chunk pass, moving the cursor to the true disk content end before memory is consulted. Ordering inside the pass is load-bearing: the entry read can outrun the shipped refs by a chunk appended between collection and read, and the transition timestamp becomes the client's refs filter - stamping it past unshipped content would silently drop that chunk's events on the next delta. The pass therefore re-ships the delta after the entry read, so the transition never exceeds shipped content. Bump-displaced aggregated entries can still arrive inline above the cursor with originals below it; that duplication is bounded and stays within the documented at-least-once residual. * fix(filer): prune ref state at the minute the scan actually stops at The collector compares file names at minute granularity while the prune used the exact-nanosecond scan bound, so for a cursor at 12:31:20 the 12-30 file was still collected but its sent state was already deleted - the next pass reshipped the whole file, re-creating the duplicate-refs class the state exists to prevent. Truncate the bound to the minute the file names live in. * fix(filer): derive the chunk cursor from the shipped refs themselves The advance-only entry read left the three positions that must agree in each other's blind spots. Its snapshot could trail the second delta's, so a chunk appended between them shipped events newer than the cursor and the memory pass sent them again. And it made the filer decode the tail range on every pass, serialized ahead of the client's own reads by the transition marker - re-introducing a slice of the replay work chunk mode exists to offload. Compute the cursor from the shipped set instead: the final entry timestamp of each filer's last shipped chunk, decoded once through the shared chunk cache. Refs coverage, transition marker, and memory start are then the same number by construction - nothing is decoded twice, nothing is dropped, and the per-pass server cost falls to one cached chunk decode per filer. The second delta and the once-per-position refs gate existed to patch the entry read's snapshot races, so both go with it; the range read survives only as a fallback for legacy chunks that do not decode standalone. * fix(filer): keep the chunk-cursor probe inside the shipped snapshot Three holes in the tail probe, all variations of stepping outside what was shipped. A permanently missing chunk failed the stream before the transition marker, so the client discarded its pending refs and reconnected to the same failure forever - blocking all later metadata behind one dead volume, where every other replay path (including the client's own reader) skips such chunks; the probe now walks back to the last readable chunk, and a filer with nothing readable simply contributes no cursor. The legacy fallback re-listed the logs after the refs were collected, so a concurrent append could push the range end over an unshipped chunk and the marker past events the client never received; it now streams the shipped chunk list itself, so no snapshot other than the shipped one is ever consulted. And a file selected before UntilNs can hold entries past it, which the client filters while still adopting the marker as its checkpoint - a later bounded request then skipped them; the marker is clamped to the bound. * fix(filer): make the cursor probe an exact mirror of the client's reader The probe answered from the server's view of the chunks; the marker's correctness depends on the client's. Its backward walk found the last readable chunk, but the client reads forward and stops at the first unreadable one, never resuming within a file - for readable, missing, readable the marker claimed the suffix the client never applied, losing those entries permanently. Keeping only each filer's final file ref discarded the progress of earlier readable files when that file was wholly missing, rewinding the marker to the start cursor. And a torn trailing size prefix - what a crashed writer leaves - failed the probe where the client reads a clean end, blocking the marker forever on data the client accepts. The probe is now shaped like the reader it answers for: per file the readable prefix, per filer the newest file with content, and no condition escapes as an error - understating the marker only re-ships, overstating loses events, and a probe failure must never block the transition the client is waiting on. Each rule is pinned by a test that fails against the previous shape. * fix(filer): judge chunk readability at the volumes, not the decode cache Two ways the probe's answer could drift from what the client experiences. A chunk this server decoded earlier stays warm in the shared cache after its volume dies, so the probe sailed past a chunk the direct-reading client stops at - marker beyond the unread suffix, entries lost. Every chunk now passes a volume lookup before the cache is consulted; the lookup rides the master client's in-memory map, so the probe stays cheap. And a probe stop was treated as harmless understatement, but the delta had already marked the whole ref sent: a transient server-side failure left the cursor stranded behind shipped content for the life of the connection, parking aggregated streams below the watermark for data the client already holds. The pass now rolls back the sent state of every ref above the file that answered the probe, so unreached refs re-ship and re-probe until the cursor gets there. Re-shipped entries at or below the client's checkpoint are filtered client-side, and batches are marker-separated, so a re-shipped file cannot rewind a merge mid-batch. * test(filer): end-to-end subscribe-loop harness and wire-contract tests Every escaped bug across this change's review rounds lived in an interaction the unit tests could not see: the loop state machine, the disk/memory handoff, or the server/client contract. The harness runs the real SubscribeLocalMetadata loop against a real leveldb-backed filer, faking only the volume layer behind the existing test hooks, and asserts the delivered stream itself. Eight scenarios, each pinning a class this change was reviewed for: the headline evicted-unflushed gap parks and then delivers in full; a ring that evicted nothing serves memory promptly; a backlog-to-live handoff with 1ms-adjacent timestamps across every boundary delivers exactly once; a flush-proven gap over vacuumed log files skips to the retained ring including a single-entry window; a bounded subscription terminates at its bound; a permanently wedged flush ends in the give-up skip with the stream still alive; and chunk mode is checked against the real client code - pb.ReadLogFileRefs applied to the shipped refs must cover everything the transition marker claims, with and without a dead volume in the middle. Validated by re-introducing three fixed bugs: the missing eviction guard delivers during the unproven gap, a 2ms cursor error at the handoff drops exactly one event, and resuming at rather than below the earliest retained window loses a single-entry window's sole event - each caught by the scenario built for it. The gap timing knobs become vars so parks run at test speed, a small filer hook swaps the volume-touching read functions, and a sender test pins the refs wire rules the client depends on: never batched, never an envelope, everything in order. * fix(filer): re-ship a partially read answering file, pin the probe's limits The sent-state rollback stopped at files newer than the one that answered the probe. When the answering file itself was only prefix-readable - a dead or transient chunk mid-file - its unread suffix stayed marked sent, and the next append advanced the cursor past it for good. The probe now reports whether the answering file was read through to its end, and a prefix-limited answer re-ships that file too; a torn tail counts as complete, since the client's read ends there as well. The rollback rules live in one predicate with a table test - files below a complete answer stay sent, because the client has moved past them and re-shipping cannot rewind its filter. Two test honesty fixes ride along. The loop harness derived its timestamp base from time.Now() per call, so expectations recomputed across a second boundary drifted by exactly one second; the base is now fixed per harness. And the probe's liveness boundary is pinned as a test instead of a comment: a volume lookup cannot see a dead needle or a stale location inside a resolvable volume, so a warm cache can answer past a chunk the client fails on - accepted because metadata log chunks die volume-at-a-time and the alternative is a real read per probe, which is what the probe exists to avoid. The test states the boundary so changing it is a decision, not an accident. --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: Chris Lu <chris.lu@gmail.com>
1216 lines
48 KiB
Go
1216 lines
48 KiB
Go
package weed_server
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/prometheus/client_golang/prometheus"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/stats"
|
|
|
|
"google.golang.org/protobuf/proto"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/filer"
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
"github.com/seaweedfs/seaweedfs/weed/util/log_buffer"
|
|
)
|
|
|
|
// Vars, not consts: the loop tests shrink them to drive parks and give-ups in
|
|
// test time.
|
|
var (
|
|
// unflushedGapRetryInterval caps the wait of a subscriber parked on a recent
|
|
// (possibly-unflushed) gap, in case the flush notification is missed.
|
|
unflushedGapRetryInterval = 2 * time.Second
|
|
|
|
// gapStallWarnInterval paces the warning for a subscriber that stays parked.
|
|
gapStallWarnInterval = time.Minute
|
|
|
|
// maxGapStall bounds a gap wait before giving up and skipping it, counted
|
|
// and logged: a dead peer makes the wait permanent, and failing the stream
|
|
// only moves the loop into a client that reconnects to the same wall.
|
|
maxGapStall = 15 * time.Minute
|
|
)
|
|
|
|
const (
|
|
// MaxUnsyncedEvents send empty notification with timestamp when certain amount of events have been filtered
|
|
MaxUnsyncedEvents = 1e3
|
|
|
|
// idleHeartbeatInterval bounds how often a caught-up subscriber that asked
|
|
// for idle heartbeats is reminded that the source is alive and has nothing
|
|
// newer. It keeps freshness signals such as filer.sync's sync_offset metric
|
|
// from looking stuck during read-only periods on the source.
|
|
idleHeartbeatInterval = 5 * time.Second
|
|
)
|
|
|
|
// metadataStreamSender is satisfied by both gRPC stream types and pipelinedSender.
|
|
type metadataStreamSender interface {
|
|
Send(*filer_pb.SubscribeMetadataResponse) error
|
|
}
|
|
|
|
const (
|
|
// batchBehindThreshold: when an event's timestamp is older than this
|
|
// relative to wall clock, the sender switches to batch mode for throughput.
|
|
// When events are closer to current time, they are sent one-by-one for
|
|
// low latency.
|
|
batchBehindThreshold = 2 * time.Minute
|
|
maxBatchSize = 256
|
|
)
|
|
|
|
// pipelinedSender decouples event reading from gRPC delivery by buffering
|
|
// messages in a channel. A dedicated goroutine handles stream.Send(), allowing
|
|
// the reader to continue reading ahead without waiting for the client to
|
|
// acknowledge each event.
|
|
//
|
|
// When the client declares support for batching AND events are far behind
|
|
// current time (backlog catch-up), multiple events are packed into a single
|
|
// stream.Send() using the Events field. Otherwise events are sent one-by-one.
|
|
type pipelinedSender struct {
|
|
sendCh chan *filer_pb.SubscribeMetadataResponse
|
|
errCh chan error
|
|
done chan struct{}
|
|
canBatch bool // true only if client set ClientSupportsBatching
|
|
}
|
|
|
|
func newPipelinedSender(stream metadataStreamSender, bufSize int, clientSupportsBatching bool) *pipelinedSender {
|
|
s := &pipelinedSender{
|
|
sendCh: make(chan *filer_pb.SubscribeMetadataResponse, bufSize),
|
|
errCh: make(chan error, 1),
|
|
done: make(chan struct{}),
|
|
canBatch: clientSupportsBatching,
|
|
}
|
|
go s.sendLoop(stream)
|
|
return s
|
|
}
|
|
|
|
func (s *pipelinedSender) sendLoop(stream metadataStreamSender) {
|
|
defer close(s.done)
|
|
for msg := range s.sendCh {
|
|
// LogFileRefs messages are unbatchable: the client recognizes them by
|
|
// the top-level field and skips the rest of the response, so a refs
|
|
// envelope would drop its Events tail and refs inside Events would be
|
|
// applied as an (empty) event. Their TsNs is 0, which the batch
|
|
// heuristic would misread as far behind. Always send them solo.
|
|
shouldBatch := s.canBatch && len(msg.LogFileRefs) == 0 &&
|
|
time.Now().UnixNano()-msg.TsNs > int64(batchBehindThreshold)
|
|
|
|
if !shouldBatch {
|
|
// Real-time: send immediately for low latency
|
|
if err := stream.Send(msg); err != nil {
|
|
s.reportErr(err)
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
|
|
// Backlog: batch multiple events into one Send for throughput.
|
|
// The first event goes in the top-level fields; additional events
|
|
// go in the Events slice. Old clients ignore the Events field.
|
|
batch := make([]*filer_pb.SubscribeMetadataResponse, 0, maxBatchSize)
|
|
batch = append(batch, msg)
|
|
var trailingRefs *filer_pb.SubscribeMetadataResponse
|
|
drain:
|
|
for len(batch) < maxBatchSize {
|
|
select {
|
|
case next, ok := <-s.sendCh:
|
|
if !ok {
|
|
break drain
|
|
}
|
|
if len(next.LogFileRefs) > 0 {
|
|
// already consumed; send it solo right after the batch
|
|
trailingRefs = next
|
|
break drain
|
|
}
|
|
batch = append(batch, next)
|
|
default:
|
|
break drain
|
|
}
|
|
}
|
|
|
|
var toSend *filer_pb.SubscribeMetadataResponse
|
|
if len(batch) == 1 {
|
|
toSend = batch[0]
|
|
} else {
|
|
// Pack batch: first event is the envelope, rest go in Events
|
|
toSend = batch[0]
|
|
toSend.Events = batch[1:]
|
|
}
|
|
if err := stream.Send(toSend); err != nil {
|
|
s.reportErr(err)
|
|
return
|
|
}
|
|
if toSend.Events != nil {
|
|
toSend.Events = nil
|
|
}
|
|
if trailingRefs != nil {
|
|
if err := stream.Send(trailingRefs); err != nil {
|
|
s.reportErr(err)
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *pipelinedSender) reportErr(err error) {
|
|
select {
|
|
case s.errCh <- err:
|
|
default:
|
|
}
|
|
// Don't drain sendCh here — Send() detects the exit via <-s.done
|
|
// and the deferred close(s.done) in sendLoop will fire after this returns.
|
|
}
|
|
|
|
func (s *pipelinedSender) Send(msg *filer_pb.SubscribeMetadataResponse) error {
|
|
select {
|
|
case s.sendCh <- msg:
|
|
return nil
|
|
case err := <-s.errCh:
|
|
return err
|
|
case <-s.done:
|
|
// Sender goroutine exited (stream error or shutdown).
|
|
select {
|
|
case err := <-s.errCh:
|
|
return err
|
|
default:
|
|
return fmt.Errorf("pipelined sender closed")
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *pipelinedSender) Close() error {
|
|
close(s.sendCh)
|
|
<-s.done
|
|
select {
|
|
case err := <-s.errCh:
|
|
return err
|
|
default:
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// reportUnprovenAggregatedCrossing records the residual hole: a disk read that
|
|
// crosses the eviction watermark may have advanced on one peer's log while a
|
|
// lagging peer still holds unflushed events in the crossed range. Locally
|
|
// undecidable (log files carry random filer ids, peers are tracked by address);
|
|
// closing it needs each peer's flush watermark on the subscribe stream.
|
|
func reportUnprovenAggregatedCrossing(cursorBeforeTsNs, cursorAfterTsNs, evictedTsNs int64, clientName, pathPrefix string) {
|
|
if evictedTsNs == 0 || cursorBeforeTsNs >= evictedTsNs || cursorAfterTsNs < evictedTsNs {
|
|
return
|
|
}
|
|
stats.FilerSubscribeUnprovenGapCrossings.WithLabelValues("aggregated").Inc()
|
|
glog.Warningf("aggregated subscriber %s %s crossed an evicted range (%v..%v] on peer disk reads; a peer that flushes into it later will not be re-read",
|
|
clientName, pathPrefix, time.Unix(0, cursorBeforeTsNs), time.Unix(0, evictedTsNs))
|
|
}
|
|
|
|
// diskReadAdvanced reports whether a persisted read moved the subscriber on.
|
|
// A chunk-ref read reports the minute-level name of the last file it shipped,
|
|
// clamped so it never rewinds, so it comes back non-zero even when it names the
|
|
// position that was already current. Treating that as progress clears the stall
|
|
// timer, and a subscriber parked on a gap it re-ships the same refs for would
|
|
// reset the timer every retry and never reach the stall bound.
|
|
func diskReadAdvanced(processedTsNs int64, cursor log_buffer.MessagePosition) bool {
|
|
return processedTsNs != 0 && processedTsNs > cursor.Time.UnixNano()
|
|
}
|
|
|
|
// gapResumeCursorOffset marks every cursor these loops hand to the memory read:
|
|
// gated, so a seal racing the loop's watermark check is refused under the
|
|
// read's own lock instead of silently served from the earliest window.
|
|
const gapResumeCursorOffset = log_buffer.EvictionGatedOffset
|
|
|
|
// memoryHoldsGap reports whether nothing after the cursor was evicted. Equality
|
|
// counts: the evicted window ends on the watermark, retained windows start
|
|
// strictly after it, and the persisted reader skips ts <= cursor, so no wait
|
|
// can ever produce the boundary entry - refusing there never ends.
|
|
func memoryHoldsGap(currentTsNs, lastEvictedTsNs int64) bool {
|
|
if lastEvictedTsNs == 0 {
|
|
return true // nothing was ever dropped from the ring
|
|
}
|
|
return currentTsNs >= lastEvictedTsNs
|
|
}
|
|
|
|
// gapStallReporter makes a parked subscriber visible: a flush that never lands
|
|
// stalls the stream for good, and filer.sync and mount followers just stop
|
|
// advancing with no error on either side.
|
|
//
|
|
// The gauge counts parked subscribers per scope. It deliberately carries no
|
|
// per-client label: clientName embeds the ephemeral source port (a series per
|
|
// reconnect), and the client-supplied name is not unique either - every mount
|
|
// registers as "mount" - so same-named streams would clobber and delete each
|
|
// other's series. A count needs no identity and no cleanup; the logs carry the
|
|
// client details.
|
|
type gapStallReporter struct {
|
|
scope string
|
|
clientName string
|
|
pathPrefix string
|
|
since time.Time
|
|
lastWarnAt time.Time
|
|
}
|
|
|
|
func (r *gapStallReporter) gauge() prometheus.Gauge {
|
|
return stats.FilerSubscribeGapStalledGauge.WithLabelValues(r.scope)
|
|
}
|
|
|
|
// stalledFor reports how long this subscriber has been parked, zero if it is not.
|
|
func (r *gapStallReporter) stalledFor() time.Duration {
|
|
if r.since.IsZero() {
|
|
return 0
|
|
}
|
|
return time.Since(r.since)
|
|
}
|
|
|
|
// park records that the subscriber is waiting on a gap. It stays quiet until
|
|
// the stall has lasted gapStallWarnInterval: during a catch-up burst a
|
|
// subscriber parks and resumes every couple of seconds, and a warning per
|
|
// cycle would bury the long-stall warnings this reporter exists to surface.
|
|
func (r *gapStallReporter) park(cursor time.Time, detail string) {
|
|
now := time.Now()
|
|
if r.since.IsZero() {
|
|
r.since = now
|
|
r.gauge().Inc()
|
|
}
|
|
if now.Sub(r.since) < gapStallWarnInterval {
|
|
return
|
|
}
|
|
if !r.lastWarnAt.IsZero() && now.Sub(r.lastWarnAt) < gapStallWarnInterval {
|
|
return
|
|
}
|
|
r.lastWarnAt = now
|
|
glog.Warningf("%s subscriber %s %s parked %v at %v: %s", r.scope, r.clientName, r.pathPrefix,
|
|
now.Sub(r.since).Truncate(time.Second), cursor, detail)
|
|
}
|
|
|
|
// resumed marks the gap cleared. Only a stall park() had already warned about
|
|
// is worth announcing.
|
|
func (r *gapStallReporter) resumed() {
|
|
if r.since.IsZero() {
|
|
return
|
|
}
|
|
if !r.lastWarnAt.IsZero() {
|
|
glog.Warningf("%s subscriber %s %s resumed after %v parked", r.scope, r.clientName, r.pathPrefix,
|
|
time.Since(r.since).Truncate(time.Second))
|
|
}
|
|
r.since, r.lastWarnAt = time.Time{}, time.Time{}
|
|
r.gauge().Dec()
|
|
}
|
|
|
|
// gaveUp records that the subscriber stopped waiting on an unprovable gap and
|
|
// skipped it. This is the loss the whole gap machinery exists to make loud: it
|
|
// shares the unproven-crossing counter and logs at error level.
|
|
func (r *gapStallReporter) gaveUp(cursor time.Time, skipToTsNs int64, detail string) {
|
|
stats.FilerSubscribeUnprovenGapCrossings.WithLabelValues(r.scope).Inc()
|
|
glog.Errorf("%s subscriber %s %s skipping the gap (%v..%v] after %v parked: %s; events a peer flushes into that range later will not be delivered",
|
|
r.scope, r.clientName, r.pathPrefix, cursor, time.Unix(0, skipToTsNs), r.stalledFor().Truncate(time.Second), detail)
|
|
r.since, r.lastWarnAt = time.Time{}, time.Time{}
|
|
r.gauge().Dec()
|
|
}
|
|
|
|
// restartStall re-arms the stall clock for a park that outlived maxGapStall
|
|
// with nothing to skip to, so the give-up path does not retrigger on every
|
|
// retry while still reporting each full cycle.
|
|
func (r *gapStallReporter) restartStall(cursor time.Time, detail string) {
|
|
glog.Errorf("%s subscriber %s %s still parked after %v at %v with nothing to skip to: %s", r.scope, r.clientName,
|
|
r.pathPrefix, r.stalledFor().Truncate(time.Second), cursor, detail)
|
|
r.since, r.lastWarnAt = time.Now(), time.Time{}
|
|
}
|
|
|
|
// close releases the gauge on teardown. Unlike resumed() it does not claim
|
|
// recovery: a subscriber that disconnects while parked never resumed.
|
|
func (r *gapStallReporter) close() {
|
|
if r.since.IsZero() {
|
|
return
|
|
}
|
|
glog.Warningf("%s subscriber %s %s disconnected after %v parked, still behind", r.scope, r.clientName,
|
|
r.pathPrefix, r.stalledFor().Truncate(time.Second))
|
|
r.gauge().Dec()
|
|
r.since = time.Time{}
|
|
}
|
|
|
|
// parkOnGap parks the subscriber on a gap it cannot read past and reports how
|
|
// to go on. done: the stream is over - the client is gone, a bounded
|
|
// subscription is complete, or the context ended. skip: the park outlived
|
|
// maxGapStall and the caller must resume at skipToTsNs, abandoning the gap
|
|
// (recorded via gaveUp). Otherwise the caller re-probes. notifyChan may be nil,
|
|
// which parks on the retry timer alone - right when no local signal
|
|
// corresponds to the event being waited for. The park is where a stalled
|
|
// subscriber spends all its time, so every exit the read loop relies on has to
|
|
// be checked here too.
|
|
func (fs *FilerServer) parkOnGap(ctx context.Context, req *filer_pb.SubscribeMetadataRequest, gapStall *gapStallReporter, evictedTsNs func() int64, cursor log_buffer.MessagePosition, notifyChan <-chan struct{}, reason string) (skipToTsNs int64, skip bool, done bool) {
|
|
// Done exits run before park(): a finished stream was never parked, and
|
|
// marking it so leaves a false "still behind" trace. A cursor at UntilNs is
|
|
// finished - the bound is inclusive, cursors are exclusive, and
|
|
// LoopProcessLogData (the only place UntilNs ends a stream) is unreachable
|
|
// from a park.
|
|
if req.UntilNs != 0 && cursor.Time.UnixNano() >= req.UntilNs {
|
|
return 0, false, true
|
|
}
|
|
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
|
return 0, false, true
|
|
}
|
|
gapStall.park(cursor.Time, reason)
|
|
if gapStall.stalledFor() >= maxGapStall {
|
|
// Resume at the eviction watermark: everything retained starts
|
|
// strictly after it, so the recorded loss is exactly (cursor, skipTo].
|
|
if evicted := evictedTsNs(); evicted > cursor.Time.UnixNano() {
|
|
gapStall.gaveUp(cursor.Time, evicted, reason)
|
|
return evicted, true, false
|
|
}
|
|
// Nothing was withheld past the cursor - nothing to skip, nothing being
|
|
// lost; keep waiting on a fresh stall cycle.
|
|
gapStall.restartStall(cursor.Time, reason)
|
|
}
|
|
// Re-probes back off as the stall ages: every retry re-reads the persisted
|
|
// log, and probing the store each 2s for 15 minutes - per parked subscriber,
|
|
// during the outage that parked them - makes the bad time worse.
|
|
waitFor := unflushedGapRetryInterval + gapStall.stalledFor()/8
|
|
if waitFor > gapStallWarnInterval {
|
|
waitFor = gapStallWarnInterval
|
|
}
|
|
retry := time.After(waitFor)
|
|
for {
|
|
select {
|
|
case _, ok := <-notifyChan:
|
|
if !ok {
|
|
// Closed out from under us: a receive now returns instantly, so
|
|
// stop watching it rather than spinning until the timer fires.
|
|
notifyChan = nil
|
|
continue
|
|
}
|
|
case <-ctx.Done():
|
|
return 0, false, true
|
|
case <-retry:
|
|
}
|
|
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
|
return 0, false, true
|
|
}
|
|
return 0, false, false
|
|
}
|
|
}
|
|
|
|
// resolveGapResume decides whether a subscriber may skip a gap its disk read
|
|
// found empty. Either proof settles it: nothing after the cursor was evicted,
|
|
// so memory still holds the whole gap; or the flush watermark observed before
|
|
// the read had already passed the earliest in-memory timestamp, so every event
|
|
// in the gap would have been on disk when the read ran and the miss is
|
|
// authoritative. The aggregated ring never flushes - peers persist their own
|
|
// logs - so it passes flushedTsNs 0 and only the eviction proof can hold.
|
|
func resolveGapResume(currentTsNs, currentOffset, earliestMemTsNs, flushedTsNs, lastEvictedTsNs int64) (advanceToTsNs int64, advance bool) {
|
|
// No in-memory data (zero time → negative UnixNano), or memory not ahead of us.
|
|
if earliestMemTsNs <= 0 || earliestMemTsNs <= currentTsNs {
|
|
return 0, false
|
|
}
|
|
// The gap may still hold unflushed events.
|
|
if !memoryHoldsGap(currentTsNs, lastEvictedTsNs) && flushedTsNs < earliestMemTsNs {
|
|
return 0, false
|
|
}
|
|
// Resume just below earliest, not at it. A sealed window holding a single
|
|
// entry has startTime == stopTime == earliest, and the sealed-buffer lookup
|
|
// only enters a window whose stopTime is strictly after the cursor, so a
|
|
// cursor sitting exactly on earliest skips that window entirely and loses
|
|
// its sole event. One nanosecond lower takes the startTime.After branch and
|
|
// returns the whole window.
|
|
target := earliestMemTsNs - 1
|
|
if target < currentTsNs {
|
|
return 0, false
|
|
}
|
|
if target == currentTsNs && currentOffset <= 0 {
|
|
// The sentinel resume would be the position we already hold.
|
|
return 0, false
|
|
}
|
|
// target > cursor is plainly forward. target == cursor with a positive
|
|
// (exclusive) offset is progress too: that cursor cannot be served -
|
|
// ReadFromBuffer refuses positive offsets below the window - while the
|
|
// sentinel one is, and both deliver exactly the entries after target.
|
|
return target, true
|
|
}
|
|
|
|
// gapPass carries what the shared post-disk gap decisions differ by between
|
|
// the two subscribe loops; everything else about them must stay identical, and
|
|
// this PR's history shows they drift when edited separately.
|
|
type gapPass struct {
|
|
fs *FilerServer
|
|
req *filer_pb.SubscribeMetadataRequest
|
|
gapStall *gapStallReporter
|
|
earliest func() time.Time
|
|
evicted func() int64 // gap-proof watermark; aggregated uses the received-ts space
|
|
flushed func() int64 // flush watermark the last disk read observed; aggregated: 0
|
|
gapChan <-chan struct{}
|
|
dataChan <-chan struct{}
|
|
gapReason func(earliest time.Time, evictedTsNs int64) string
|
|
}
|
|
|
|
type gapOutcome int
|
|
|
|
const (
|
|
gapProceed gapOutcome = iota // read memory
|
|
gapContinue // restart the pass
|
|
gapDone // the stream is over
|
|
)
|
|
|
|
// resolve is the gap decision both loops run between the disk pass and the
|
|
// memory read. A cursor the ring evicted past cannot be served from memory
|
|
// without skipping what was dropped: keep draining the disk if it just moved,
|
|
// skip if a proof says the gap is empty, park otherwise. A cursor memory
|
|
// refused with nothing evicted after it re-arms onto the retained window.
|
|
func (p *gapPass) resolve(ctx context.Context, cursor *log_buffer.MessagePosition, latch *error, diskAdvanced bool) gapOutcome {
|
|
earliest := p.earliest()
|
|
evictedTsNs := p.evicted()
|
|
cursorTsNs := cursor.Time.UnixNano()
|
|
if !memoryHoldsGap(cursorTsNs, evictedTsNs) {
|
|
if diskAdvanced {
|
|
return gapContinue // the disk may hold more of the gap
|
|
}
|
|
if advanceToTsNs, advance := resolveGapResume(cursorTsNs, cursor.Offset, earliest.UnixNano(), p.flushed(), evictedTsNs); advance {
|
|
p.gapStall.resumed()
|
|
glog.V(3).Infof("%s subscriber %s: gap proven empty, skipping from %v to earliest memory %v",
|
|
p.gapStall.scope, p.gapStall.clientName, cursor.Time, earliest)
|
|
*cursor = log_buffer.NewMessagePosition(advanceToTsNs, gapResumeCursorOffset)
|
|
*latch = nil
|
|
return gapProceed
|
|
}
|
|
return p.park(ctx, cursor, latch, p.gapChan, p.gapReason(earliest, evictedTsNs))
|
|
}
|
|
if !diskAdvanced && errors.Is(*latch, log_buffer.ResumeFromDiskError) {
|
|
// Memory refused the cursor though nothing after it was evicted: its
|
|
// exclusive offset predates the retained window. Re-arm it onto the
|
|
// window; failing even that, wait for data.
|
|
if advanceToTsNs, advance := resolveGapResume(cursorTsNs, cursor.Offset, earliest.UnixNano(), p.flushed(), evictedTsNs); advance {
|
|
p.gapStall.resumed()
|
|
*cursor = log_buffer.NewMessagePosition(advanceToTsNs, gapResumeCursorOffset)
|
|
*latch = nil
|
|
return gapProceed
|
|
}
|
|
return p.park(ctx, cursor, latch, p.dataChan, "no readable in-memory entries yet")
|
|
}
|
|
return gapProceed
|
|
}
|
|
|
|
func (p *gapPass) park(ctx context.Context, cursor *log_buffer.MessagePosition, latch *error, notifyChan <-chan struct{}, reason string) gapOutcome {
|
|
skipTo, skip, done := p.fs.parkOnGap(ctx, p.req, p.gapStall, p.evicted, *cursor, notifyChan, reason)
|
|
if done {
|
|
return gapDone
|
|
}
|
|
if skip {
|
|
*cursor = log_buffer.NewMessagePosition(skipTo, gapResumeCursorOffset)
|
|
*latch = nil
|
|
}
|
|
return gapContinue
|
|
}
|
|
|
|
func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeMetadataServer) error {
|
|
if fs.filer.MetaAggregator == nil || !fs.filer.MetaAggregator.HasRemotePeers() {
|
|
return fs.SubscribeLocalMetadata(req, stream)
|
|
}
|
|
|
|
ctx := stream.Context()
|
|
peerAddress := findClientAddress(ctx, 0)
|
|
|
|
isReplacing, alreadyKnown, clientName := fs.addClient("", req.ClientName, peerAddress, req.PathPrefix, req.ClientId, req.ClientEpoch)
|
|
if isReplacing {
|
|
} else if alreadyKnown {
|
|
return fmt.Errorf("duplicated subscription detected for client %s id %d", clientName, req.ClientId)
|
|
}
|
|
defer func() {
|
|
glog.V(0).Infof("disconnect %v subscriber %s clientId:%d", clientName, req.PathPrefix, req.ClientId)
|
|
fs.deleteClient("", clientName, req.ClientId, req.ClientEpoch)
|
|
}()
|
|
|
|
lastReadTime := log_buffer.NewMessagePosition(req.SinceNs, gapResumeCursorOffset)
|
|
glog.V(0).Infof(" %v starts to subscribe %s from %+v", clientName, req.PathPrefix, lastReadTime)
|
|
|
|
sender := newPipelinedSender(stream, 1024, req.ClientSupportsBatching)
|
|
defer sender.Close()
|
|
|
|
// Register for instant notification when new data arrives in the aggregated log buffer.
|
|
// Used to replace the 1127ms sleep with event-driven wake-up.
|
|
// Key includes clientId/epoch: a replacement stream may reuse the same
|
|
// clientName (same gRPC conn), and sharing the channel would let the old
|
|
// stream's deferred unregister close it under the new stream.
|
|
aggNotifyName := fmt.Sprintf("aggSubscribe:%s:%d:%d", clientName, req.ClientId, req.ClientEpoch)
|
|
// Same key shape for the reader: LoopProcessLogData registers it as a
|
|
// subscriber internally, once per loop iteration.
|
|
aggReaderName := fmt.Sprintf("aggMeta:%s:%d:%d", clientName, req.ClientId, req.ClientEpoch)
|
|
aggNotifyChan := fs.filer.MetaAggregator.MetaLogBuffer.RegisterSubscriber(aggNotifyName)
|
|
defer fs.filer.MetaAggregator.MetaLogBuffer.UnregisterSubscriber(aggNotifyName)
|
|
|
|
gapStall := &gapStallReporter{scope: "aggregated", clientName: clientName, pathPrefix: req.PathPrefix}
|
|
defer gapStall.close()
|
|
|
|
var unsyncedEvents int64
|
|
eachEventNotificationFn := fs.eachEventNotificationFn(req, sender, clientName, &unsyncedEvents)
|
|
|
|
// lastSeenTsNs tracks how far the subscriber has read so idle heartbeats are
|
|
// only emitted once it is caught up to the buffer head. It is read and
|
|
// written from this single goroutine, so no synchronization is needed.
|
|
var lastSeenTsNs int64
|
|
var lastHeartbeatNs int64
|
|
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
|
|
eachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) {
|
|
lastSeenTsNs = logEntry.TsNs
|
|
return baseEachLogEntryFn(logEntry)
|
|
}
|
|
|
|
var processedTsNs int64
|
|
var readPersistedLogErr error
|
|
var readInMemoryLogErr error
|
|
var isDone bool
|
|
sentRefs := make(map[string]sentRefState)
|
|
|
|
aggBuffer := fs.filer.MetaAggregator.MetaLogBuffer
|
|
gaps := &gapPass{
|
|
fs: fs,
|
|
req: req,
|
|
gapStall: gapStall,
|
|
earliest: aggBuffer.GetEarliestTime,
|
|
evicted: aggBuffer.GetLastEvictedOriginalTsNs,
|
|
flushed: func() int64 { return 0 }, // the aggregated ring never flushes
|
|
gapChan: nil, // nothing local signals a peer's flush; the timer paces it
|
|
dataChan: aggNotifyChan,
|
|
gapReason: func(earliest time.Time, evictedTsNs int64) string {
|
|
return fmt.Sprintf("gap evicted through %v is not on a peer's disk yet (earliest memory %v)",
|
|
time.Unix(0, evictedTsNs), earliest)
|
|
},
|
|
}
|
|
|
|
for {
|
|
|
|
glog.V(4).Infof("read on disk %v aggregated subscribe %s from %+v", clientName, req.PathPrefix, lastReadTime)
|
|
|
|
cursorBeforeDiskTsNs := lastReadTime.Time.UnixNano()
|
|
|
|
if req.ClientSupportsMetadataChunks {
|
|
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, req.UntilNs, sentRefs)
|
|
} else {
|
|
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn)
|
|
}
|
|
if readPersistedLogErr != nil {
|
|
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
|
|
}
|
|
if isDone {
|
|
return nil
|
|
}
|
|
|
|
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
|
|
diskAdvanced := diskReadAdvanced(processedTsNs, lastReadTime)
|
|
// Read after the disk read (an eviction landing mid-read must count) and
|
|
// in received-ts space: the ring's bumped stopTimes exceed anything on
|
|
// any peer's disk, and gating disk cursors on them parks subscribers
|
|
// that drained every peer's log.
|
|
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
|
|
if diskAdvanced {
|
|
gapStall.resumed()
|
|
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, clientName, req.PathPrefix)
|
|
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
|
|
} else if readInMemoryLogErr == nil {
|
|
// Nothing on disk and memory never spoke: scan forward for the next
|
|
// day that has logs.
|
|
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
|
|
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
|
|
found, err := fs.filer.HasPersistedLogFiles(position)
|
|
if err != nil {
|
|
return fmt.Errorf("checking persisted log files: %w", err)
|
|
}
|
|
if found {
|
|
gapStall.resumed()
|
|
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, clientName, req.PathPrefix)
|
|
lastReadTime = position
|
|
}
|
|
}
|
|
|
|
switch gaps.resolve(ctx, &lastReadTime, &readInMemoryLogErr, diskAdvanced) {
|
|
case gapDone:
|
|
return nil
|
|
case gapContinue:
|
|
continue
|
|
}
|
|
|
|
glog.V(4).Infof("read in memory %v aggregated subscribe %s from %+v", clientName, req.PathPrefix, lastReadTime)
|
|
|
|
lastReadTime, isDone, readInMemoryLogErr = fs.filer.MetaAggregator.MetaLogBuffer.LoopProcessLogData(aggReaderName, lastReadTime, req.UntilNs, func() bool {
|
|
select {
|
|
case <-ctx.Done():
|
|
return false
|
|
default:
|
|
}
|
|
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
|
return false
|
|
}
|
|
lastHeartbeatNs = fs.maybeSendIdleHeartbeat(req, sender, fs.filer.MetaAggregator.MetaLogBuffer, lastReadTime.Time.UnixNano(), lastSeenTsNs, lastHeartbeatNs)
|
|
return true
|
|
}, eachLogEntryFn)
|
|
if readInMemoryLogErr != nil {
|
|
if errors.Is(readInMemoryLogErr, log_buffer.ResumeFromDiskError) {
|
|
// Fell behind the ring: back to the disk pass, and from there to
|
|
// the gap resolution above if the disk has nothing either.
|
|
continue
|
|
}
|
|
glog.Errorf("processed to %v: %v", lastReadTime, readInMemoryLogErr)
|
|
if !errors.Is(readInMemoryLogErr, log_buffer.ResumeError) {
|
|
break
|
|
}
|
|
}
|
|
if isDone {
|
|
return nil
|
|
}
|
|
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
|
glog.V(0).Infof("client %v is closed", clientName)
|
|
return nil
|
|
}
|
|
|
|
// Wait for new data (event-driven instead of 1127ms polling).
|
|
// Drain any stale notification first to avoid a spurious wake-up.
|
|
select {
|
|
case <-aggNotifyChan:
|
|
default:
|
|
}
|
|
select {
|
|
case <-aggNotifyChan:
|
|
case <-ctx.Done():
|
|
return nil
|
|
}
|
|
}
|
|
|
|
return readInMemoryLogErr
|
|
|
|
}
|
|
|
|
func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeLocalMetadataServer) error {
|
|
|
|
ctx := stream.Context()
|
|
peerAddress := findClientAddress(ctx, 0)
|
|
|
|
// use negative client id to differentiate from addClient()/deleteClient() used in SubscribeMetadata()
|
|
req.ClientId = -req.ClientId
|
|
|
|
isReplacing, alreadyKnown, clientName := fs.addClient("local", req.ClientName, peerAddress, req.PathPrefix, req.ClientId, req.ClientEpoch)
|
|
if isReplacing {
|
|
} else if alreadyKnown {
|
|
return fmt.Errorf("duplicated local subscription detected for client %s clientId:%d", clientName, req.ClientId)
|
|
}
|
|
defer func() {
|
|
glog.V(0).Infof("disconnect %v local subscriber %s clientId:%d", clientName, req.PathPrefix, req.ClientId)
|
|
fs.deleteClient("local", clientName, req.ClientId, req.ClientEpoch)
|
|
}()
|
|
|
|
lastReadTime := log_buffer.NewMessagePosition(req.SinceNs, gapResumeCursorOffset)
|
|
glog.V(0).Infof(" + %v local subscribe %s from %+v clientId:%d", clientName, req.PathPrefix, lastReadTime, req.ClientId)
|
|
|
|
sender := newPipelinedSender(stream, 1024, req.ClientSupportsBatching)
|
|
defer sender.Close()
|
|
|
|
// Bounded gap waits use the buffer's subscriber notification plus a retry
|
|
// timer, so a flush landing between the disk read and the wait cannot
|
|
// strand the subscriber (no lost-wakeup window). Key includes clientId/
|
|
// epoch so a replacement stream never shares (and loses) the channel.
|
|
localNotifyName := fmt.Sprintf("localGap:%s:%d:%d", clientName, req.ClientId, req.ClientEpoch)
|
|
// Same key shape for the reader: LoopProcessLogData registers it as a
|
|
// subscriber internally, once per loop iteration.
|
|
localReaderName := fmt.Sprintf("localMeta:%s:%d:%d", clientName, req.ClientId, req.ClientEpoch)
|
|
localFlushChan := fs.filer.LocalMetaLogBuffer.RegisterFlushSubscriber(localNotifyName)
|
|
defer fs.filer.LocalMetaLogBuffer.UnregisterFlushSubscriber(localNotifyName)
|
|
|
|
gapStall := &gapStallReporter{scope: "local", clientName: clientName, pathPrefix: req.PathPrefix}
|
|
defer gapStall.close()
|
|
|
|
var unsyncedEvents int64
|
|
eachEventNotificationFn := fs.eachEventNotificationFn(req, sender, clientName, &unsyncedEvents)
|
|
|
|
// lastSeenTsNs tracks how far the subscriber has read so idle heartbeats are
|
|
// only emitted once it is caught up to the buffer head. It is read and
|
|
// written from this single goroutine, so no synchronization is needed.
|
|
var lastSeenTsNs int64
|
|
var lastHeartbeatNs int64
|
|
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
|
|
eachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) {
|
|
lastSeenTsNs = logEntry.TsNs
|
|
return baseEachLogEntryFn(logEntry)
|
|
}
|
|
|
|
var processedTsNs int64
|
|
var readPersistedLogErr error
|
|
var readInMemoryLogErr error
|
|
var isDone bool
|
|
var lastCheckedFlushTsNs int64 = -1 // Track the last flushed time we checked
|
|
var lastDiskReadTsNs int64 = -1 // Track the last read position we used for disk read
|
|
sentRefs := make(map[string]sentRefState)
|
|
|
|
localBuffer := fs.filer.LocalMetaLogBuffer
|
|
gaps := &gapPass{
|
|
fs: fs,
|
|
req: req,
|
|
gapStall: gapStall,
|
|
earliest: localBuffer.GetEarliestTime,
|
|
evicted: localBuffer.GetLastEvictedTsNs, // local disk carries the ring's own timestamps
|
|
flushed: func() int64 { return lastCheckedFlushTsNs },
|
|
gapChan: localFlushChan,
|
|
dataChan: localFlushChan,
|
|
gapReason: func(earliest time.Time, evictedTsNs int64) string {
|
|
return fmt.Sprintf("gap is not flushed yet (earliest memory %v, flushed through %v)",
|
|
earliest, time.Unix(0, lastCheckedFlushTsNs))
|
|
},
|
|
}
|
|
|
|
for {
|
|
// Check if new data has been flushed to disk since last check, or if read position advanced
|
|
currentFlushTsNs := fs.filer.LocalMetaLogBuffer.GetLastFlushTsNs()
|
|
currentReadTsNs := lastReadTime.Time.UnixNano()
|
|
// Read from disk if: first time, new flush observed, or read position advanced (draining backlog)
|
|
shouldReadFromDisk := lastCheckedFlushTsNs == -1 ||
|
|
currentFlushTsNs > lastCheckedFlushTsNs ||
|
|
currentReadTsNs > lastDiskReadTsNs
|
|
|
|
diskAdvanced := false
|
|
if shouldReadFromDisk {
|
|
// Record the position we are about to read from
|
|
lastDiskReadTsNs = currentReadTsNs
|
|
glog.V(4).Infof("read on disk %v local subscribe %s from %+v (lastFlushed: %v)", clientName, req.PathPrefix, lastReadTime, time.Unix(0, currentFlushTsNs))
|
|
if req.ClientSupportsMetadataChunks {
|
|
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, req.UntilNs, sentRefs)
|
|
} else {
|
|
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn)
|
|
}
|
|
if readPersistedLogErr != nil {
|
|
glog.V(0).Infof("read on disk %v local subscribe %s from %+v: %v", clientName, req.PathPrefix, lastReadTime, readPersistedLogErr)
|
|
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
|
|
}
|
|
if isDone {
|
|
return nil
|
|
}
|
|
|
|
// Update the last checked flushed time
|
|
lastCheckedFlushTsNs = currentFlushTsNs
|
|
|
|
diskAdvanced = diskReadAdvanced(processedTsNs, lastReadTime)
|
|
if diskAdvanced {
|
|
gapStall.resumed()
|
|
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
|
|
} else if readInMemoryLogErr == nil {
|
|
// Nothing on disk and memory never spoke: scan forward for the
|
|
// next day that has logs.
|
|
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
|
|
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
|
|
found, err := fs.filer.HasPersistedLogFiles(position)
|
|
if err != nil {
|
|
return fmt.Errorf("checking persisted log files: %w", err)
|
|
}
|
|
if found {
|
|
gapStall.resumed()
|
|
lastReadTime = position
|
|
}
|
|
}
|
|
}
|
|
|
|
switch gaps.resolve(ctx, &lastReadTime, &readInMemoryLogErr, diskAdvanced) {
|
|
case gapDone:
|
|
return nil
|
|
case gapContinue:
|
|
continue
|
|
}
|
|
|
|
glog.V(3).Infof("read in memory %v local subscribe %s from %+v", clientName, req.PathPrefix, lastReadTime)
|
|
|
|
lastReadTime, isDone, readInMemoryLogErr = fs.filer.LocalMetaLogBuffer.LoopProcessLogData(localReaderName, lastReadTime, req.UntilNs, func() bool {
|
|
select {
|
|
case <-ctx.Done():
|
|
return false
|
|
default:
|
|
}
|
|
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
|
return false
|
|
}
|
|
lastHeartbeatNs = fs.maybeSendIdleHeartbeat(req, sender, fs.filer.LocalMetaLogBuffer, lastReadTime.Time.UnixNano(), lastSeenTsNs, lastHeartbeatNs)
|
|
return true
|
|
}, eachLogEntryFn)
|
|
if readInMemoryLogErr != nil {
|
|
if errors.Is(readInMemoryLogErr, log_buffer.ResumeFromDiskError) {
|
|
// Fell behind the ring: back to the disk pass (it re-runs when
|
|
// the flush or the cursor moved), and from there to the gap
|
|
// resolution above if the disk has nothing either.
|
|
continue
|
|
}
|
|
glog.Errorf("processed to %v: %v", lastReadTime, readInMemoryLogErr)
|
|
if !errors.Is(readInMemoryLogErr, log_buffer.ResumeError) {
|
|
break
|
|
}
|
|
}
|
|
if isDone {
|
|
return nil
|
|
}
|
|
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
|
|
return nil
|
|
}
|
|
}
|
|
|
|
return readInMemoryLogErr
|
|
|
|
}
|
|
|
|
func eachLogEntryFn(req *filer_pb.SubscribeMetadataRequest, sender metadataStreamSender, eachEventNotificationFn func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error, filtered *int64) log_buffer.EachLogEntryFuncType {
|
|
// A shallow scan of the path fields skips unmarshaling chunk-heavy events
|
|
// this subscriber would filter out anyway; scan surprises fall back to the
|
|
// full decode. Only a delivery resets the shared unsynced-events counter.
|
|
prefilter := req.PathPrefix != "" || len(req.PathPrefixes) > 0 || len(req.Directories) > 0
|
|
return func(logEntry *filer_pb.LogEntry) (bool, error) {
|
|
if prefilter {
|
|
if skeleton, ok := filer_pb.ScanMetadataEventSkeleton(logEntry.Data); ok &&
|
|
!filer_pb.MetadataEventMatchesSubscription(skeleton, req.PathPrefix, req.PathPrefixes, req.Directories) {
|
|
*filtered++
|
|
if *filtered > MaxUnsyncedEvents {
|
|
if err := sender.Send(&filer_pb.SubscribeMetadataResponse{
|
|
EventNotification: &filer_pb.EventNotification{},
|
|
TsNs: skeleton.TsNs,
|
|
}); err != nil {
|
|
return false, err
|
|
}
|
|
*filtered = 0
|
|
}
|
|
return false, nil
|
|
}
|
|
}
|
|
event := &filer_pb.SubscribeMetadataResponse{}
|
|
// proto.Unmarshal (not UnmarshalVT) validates UTF-8 in string fields, so
|
|
// malformed metadata is rejected here instead of reaching path filtering
|
|
// and subscribers.
|
|
if err := proto.Unmarshal(logEntry.Data, event); err != nil {
|
|
glog.Errorf("unexpected unmarshal filer_pb.SubscribeMetadataResponse: %v", err)
|
|
return false, fmt.Errorf("unexpected unmarshal filer_pb.SubscribeMetadataResponse: %w", err)
|
|
}
|
|
|
|
if err := eachEventNotificationFn(event.Directory, event.EventNotification, event.TsNs); err != nil {
|
|
return false, err
|
|
}
|
|
|
|
return false, nil
|
|
}
|
|
}
|
|
|
|
// maybeSendIdleHeartbeat emits an empty response carrying the current time when
|
|
// the subscriber has consumed everything up to the buffer head. The client uses
|
|
// it to advance freshness signals (e.g. filer.sync's sync_offset) without moving
|
|
// its resume checkpoint, so a restart still re-reads from the last real event.
|
|
//
|
|
// The catch-up floor is the max of two read-progress markers:
|
|
// - readPositionTsNs: how far the read cursor has advanced. It starts at
|
|
// SinceNs and also covers metadata-chunks mode, where persisted entries are
|
|
// replayed as log file refs rather than through eachLogEntryFn.
|
|
// - lastSeenTsNs: the timestamp of the most recent entry streamed in this
|
|
// call. It advances live while reading the in-memory backlog, before the
|
|
// read cursor returned by LoopProcessLogData has been updated.
|
|
//
|
|
// While the buffer head is past that floor the subscriber is still behind (e.g.
|
|
// replaying a backlog) and no heartbeat is sent. Returns the (possibly advanced)
|
|
// lastHeartbeatNs.
|
|
func (fs *FilerServer) maybeSendIdleHeartbeat(req *filer_pb.SubscribeMetadataRequest, sender metadataStreamSender, logBuffer *log_buffer.LogBuffer, readPositionTsNs, lastSeenTsNs, lastHeartbeatNs int64) int64 {
|
|
if !req.ClientSupportsIdleHeartbeat {
|
|
return lastHeartbeatNs
|
|
}
|
|
floorTsNs := lastSeenTsNs
|
|
if readPositionTsNs > floorTsNs {
|
|
floorTsNs = readPositionTsNs
|
|
}
|
|
if logBuffer.LastTsNs.Load() > floorTsNs {
|
|
// the buffer holds data the subscriber has not reached yet
|
|
return lastHeartbeatNs
|
|
}
|
|
now := time.Now().UnixNano()
|
|
if now-lastHeartbeatNs < int64(idleHeartbeatInterval) {
|
|
return lastHeartbeatNs
|
|
}
|
|
if err := sender.Send(&filer_pb.SubscribeMetadataResponse{TsNs: now}); err != nil {
|
|
glog.V(0).Infof("=> idle heartbeat to %s: %v", req.ClientName, err)
|
|
return lastHeartbeatNs
|
|
}
|
|
// A heartbeat is a send too: advance the freshness gauge so an idle but
|
|
// healthy subscriber doesn't look stale. The gauge otherwise only moves on
|
|
// real matching events, which never arrive on a quiet path.
|
|
var sourceFiler string
|
|
if fs.option != nil {
|
|
sourceFiler = fs.option.Host.String()
|
|
}
|
|
stats.FilerServerLastSendTsOfSubscribeGauge.WithLabelValues(sourceFiler, req.ClientName, req.PathPrefix).Set(float64(now))
|
|
return now
|
|
}
|
|
|
|
// chunkDiskPass is the disk step for chunk-capable clients: ship the unsent
|
|
// refs, then advance the cursor to the shipped content's own end - the final
|
|
// entry timestamp of each filer's last shipped chunk, decoded through the
|
|
// shared chunk cache. Deriving the cursor from the shipped set itself keeps
|
|
// the three positions that must agree in lockstep: the client's refs cover
|
|
// exactly up to the cursor, the transition marker (which becomes the client's
|
|
// refs filter) equals it, and the memory pass delivers strictly after it - no
|
|
// range is decoded twice and none is dropped. The transition is the
|
|
// empty-notification marker: both chunk consumers buffer refs until a non-ref
|
|
// message, so an idle source would otherwise strand the backlog in the
|
|
// client's pending list until the next mutation.
|
|
func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamSender, startPos log_buffer.MessagePosition, untilNs int64, sent map[string]sentRefState) (processedTsNs int64, isDone bool, err error) {
|
|
collected, _, err := fs.filer.CollectLogFileRefs(ctx, startPos, untilNs)
|
|
if err != nil {
|
|
return 0, false, err
|
|
}
|
|
refs := deltaLogFileRefs(collected, sent, filer.PersistedLogScanStartTsNs(startPos.Time))
|
|
if len(refs) == 0 {
|
|
return startPos.Time.UnixNano(), false, nil
|
|
}
|
|
if err := fs.sendRefsBatched(sender, refs); err != nil {
|
|
return 0, false, err
|
|
}
|
|
|
|
// Shipped content end, read from the shipped chunks alone - a fresh
|
|
// listing here could see a concurrent append and move the cursor past
|
|
// unshipped content. The probe mirrors the client's reader exactly (per
|
|
// file the readable prefix, per filer the newest file with content), so
|
|
// the marker never claims events the client will not apply, and it cannot
|
|
// fail: a dead volume must not block the transition the client waits on.
|
|
cursorTsNs := startPos.Time.UnixNano()
|
|
refsPerFiler := make(map[string][]*filer_pb.LogFileChunkRef, 2)
|
|
for _, ref := range refs {
|
|
refsPerFiler[ref.FilerId] = append(refsPerFiler[ref.FilerId], ref)
|
|
}
|
|
for filerId, filerRefs := range refsPerFiler {
|
|
tailTsNs, answeredFileTsNs, ok, complete := fs.filer.LastShippedLogEntryTsNsForFiler(filerRefs)
|
|
if ok && tailTsNs > cursorTsNs {
|
|
cursorTsNs = tailTsNs
|
|
}
|
|
// Refs the cursor did not reach must re-ship on a later pass: their
|
|
// content sits above the marker, and sent-state that outlives a
|
|
// transient probe failure would strand the cursor behind them for the
|
|
// life of the connection - parking aggregated streams below the
|
|
// watermark. A prefix-limited answer re-ships the answering file too,
|
|
// or its unread suffix is abandoned the moment a later append advances
|
|
// past it. Re-shipped entries at or below the client's checkpoint are
|
|
// filtered client-side, and batches are marker-separated, so a
|
|
// re-shipped whole file cannot rewind a merge mid-batch.
|
|
for _, ref := range filerRefs {
|
|
if refNeedsReship(ref.FileTsNs, ok, answeredFileTsNs, complete) {
|
|
delete(sent, sentRefKey(filerId, ref.FileTsNs))
|
|
}
|
|
}
|
|
}
|
|
// A file selected before the bound can hold entries past it. The client
|
|
// filters those but adopts the marker as its checkpoint, so an unclamped
|
|
// marker makes a later bounded request skip them.
|
|
if untilNs != 0 && cursorTsNs > untilNs {
|
|
cursorTsNs = untilNs
|
|
}
|
|
if err := sender.Send(&filer_pb.SubscribeMetadataResponse{
|
|
EventNotification: &filer_pb.EventNotification{},
|
|
TsNs: cursorTsNs,
|
|
}); err != nil {
|
|
return 0, false, err
|
|
}
|
|
return cursorTsNs, false, nil
|
|
}
|
|
|
|
// sendRefsBatched sends refs through the pipelined sender, which keeps them
|
|
// out of Events batches; gRPC allows one sending goroutine per stream and the
|
|
// sender's goroutine is it.
|
|
func (fs *FilerServer) sendRefsBatched(sender metadataStreamSender, refs []*filer_pb.LogFileChunkRef) error {
|
|
const maxRefsPerMessage = 64
|
|
for i := 0; i < len(refs); i += maxRefsPerMessage {
|
|
end := i + maxRefsPerMessage
|
|
if end > len(refs) {
|
|
end = len(refs)
|
|
}
|
|
if err := sender.Send(&filer_pb.SubscribeMetadataResponse{LogFileRefs: refs[i:end]}); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// sentRefState tracks, per subscription, how many chunks of each log file have
|
|
// been shipped as refs. Collection re-lists files up to a flush interval behind
|
|
// the cursor (the spanning-file back-off), and a filer appends further chunks
|
|
// to its newest file, so consecutive collections overlap; shipping only each
|
|
// file's unsent chunk suffix keeps every per-filer ref stream duplicate-free
|
|
// and timestamp-sorted - the contract the client's merge reads them under.
|
|
type sentRefState struct {
|
|
chunks int
|
|
fileTsNs int64
|
|
}
|
|
|
|
// refNeedsReship says whether a shipped ref's sent state must be dropped so a
|
|
// later pass re-ships it: everything above the file that answered the probe
|
|
// (the cursor never reached it), the answering file itself when its read was
|
|
// prefix-limited (its unread suffix would otherwise be abandoned the moment a
|
|
// later append advances past it), and everything when nothing answered. Files
|
|
// below a complete answer stay sent: the client has moved past them, and
|
|
// re-shipping cannot rewind its filter.
|
|
func refNeedsReship(fileTsNs int64, answered bool, answeredFileTsNs int64, complete bool) bool {
|
|
if !answered {
|
|
return true
|
|
}
|
|
if fileTsNs > answeredFileTsNs {
|
|
return true
|
|
}
|
|
return fileTsNs == answeredFileTsNs && !complete
|
|
}
|
|
|
|
func sentRefKey(filerId string, fileTsNs int64) string {
|
|
return fmt.Sprintf("%s/%d", filerId, fileTsNs)
|
|
}
|
|
|
|
// deltaLogFileRefs reduces a collection to the chunks not yet shipped, updates
|
|
// the sent state, and prunes files the scan window has moved past.
|
|
//
|
|
// A shipped suffix is rebased to logical offset zero: the client's chunk
|
|
// reader starts at zero, and a chunk list opening at a higher offset reads as
|
|
// instant EOF - an empty replay that would silently drop the appended events.
|
|
// The cut is record-aligned because each append is one chunk of whole entries
|
|
// (logFlushFunc appends one uploaded window per flush), so the rebased suffix
|
|
// decodes as a file of its own.
|
|
func deltaLogFileRefs(refs []*filer_pb.LogFileChunkRef, sent map[string]sentRefState, pruneBeforeTsNs int64) []*filer_pb.LogFileChunkRef {
|
|
out := make([]*filer_pb.LogFileChunkRef, 0, len(refs))
|
|
for _, ref := range refs {
|
|
key := sentRefKey(ref.FilerId, ref.FileTsNs)
|
|
prior := sent[key].chunks
|
|
if len(ref.Chunks) <= prior {
|
|
continue
|
|
}
|
|
chunks := ref.Chunks[prior:]
|
|
if base := chunks[0].Offset; base != 0 {
|
|
rebased := make([]*filer_pb.FileChunk, len(chunks))
|
|
for i, c := range chunks {
|
|
cc := proto.Clone(c).(*filer_pb.FileChunk)
|
|
cc.Offset -= base
|
|
rebased[i] = cc
|
|
}
|
|
chunks = rebased
|
|
}
|
|
out = append(out, &filer_pb.LogFileChunkRef{
|
|
Chunks: chunks,
|
|
FileTsNs: ref.FileTsNs,
|
|
FilerId: ref.FilerId,
|
|
})
|
|
sent[key] = sentRefState{chunks: len(ref.Chunks), fileTsNs: ref.FileTsNs}
|
|
}
|
|
for key, st := range sent {
|
|
if st.fileTsNs < pruneBeforeTsNs {
|
|
delete(sent, key)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (fs *FilerServer) eachEventNotificationFn(req *filer_pb.SubscribeMetadataRequest, sender metadataStreamSender, clientName string, filtered *int64) func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error {
|
|
return func(dirPath string, eventNotification *filer_pb.EventNotification, tsNs int64) error {
|
|
defer func() {
|
|
if *filtered > MaxUnsyncedEvents {
|
|
if err := sender.Send(&filer_pb.SubscribeMetadataResponse{
|
|
EventNotification: &filer_pb.EventNotification{},
|
|
TsNs: tsNs,
|
|
}); err == nil {
|
|
*filtered = 0
|
|
}
|
|
}
|
|
}()
|
|
|
|
*filtered++
|
|
foundSelf := false
|
|
for _, sig := range eventNotification.Signatures {
|
|
if sig == req.Signature && req.Signature != 0 {
|
|
return nil
|
|
}
|
|
if sig == fs.filer.Signature {
|
|
foundSelf = true
|
|
}
|
|
}
|
|
if !foundSelf {
|
|
eventNotification.Signatures = append(eventNotification.Signatures, fs.filer.Signature)
|
|
}
|
|
|
|
// get complete path to the file or directory
|
|
var entryName string
|
|
if eventNotification.OldEntry != nil {
|
|
entryName = eventNotification.OldEntry.Name
|
|
} else if eventNotification.NewEntry != nil {
|
|
entryName = eventNotification.NewEntry.Name
|
|
}
|
|
|
|
fullpath := util.Join(dirPath, entryName)
|
|
|
|
// skip on filer internal meta logs
|
|
if strings.HasPrefix(fullpath, filer.SystemLogDir) {
|
|
return nil
|
|
}
|
|
|
|
message := &filer_pb.SubscribeMetadataResponse{
|
|
Directory: dirPath,
|
|
EventNotification: eventNotification,
|
|
TsNs: tsNs,
|
|
}
|
|
|
|
if !filer_pb.MetadataEventMatchesSubscription(message, req.PathPrefix, req.PathPrefixes, req.Directories) {
|
|
return nil
|
|
}
|
|
|
|
// collect timestamps for path
|
|
stats.FilerServerLastSendTsOfSubscribeGauge.WithLabelValues(fs.option.Host.String(), req.ClientName, req.PathPrefix).Set(float64(tsNs))
|
|
|
|
// println("sending", dirPath, entryName)
|
|
if err := sender.Send(message); err != nil {
|
|
glog.V(0).Infof("=> client %v: %+v", clientName, err)
|
|
return err
|
|
}
|
|
*filtered = 0
|
|
return nil
|
|
}
|
|
}
|
|
|
|
func (fs *FilerServer) addClient(scope string, clientType string, clientAddress string, pathPrefix string, clientId int32, clientEpoch int32) (isReplacing, alreadyKnown bool, clientName string) {
|
|
clientName = clientType + "@" + clientAddress
|
|
glog.V(0).Infof("+ %v listener %v clientId %v clientEpoch %v", scope, clientName, clientId, clientEpoch)
|
|
if clientId != 0 {
|
|
fs.knownListenersLock.Lock()
|
|
defer fs.knownListenersLock.Unlock()
|
|
epoch, found := fs.knownListeners[clientId]
|
|
if !found || epoch < clientEpoch {
|
|
fs.knownListeners[clientId] = clientEpoch
|
|
isReplacing = true
|
|
if fs.subscribers == nil {
|
|
fs.subscribers = make(map[int32]*metadataSubscriber)
|
|
}
|
|
fs.subscribers[clientId] = &metadataSubscriber{
|
|
clientName: clientName,
|
|
clientType: clientType,
|
|
address: clientAddress,
|
|
pathPrefix: pathPrefix,
|
|
clientId: clientId,
|
|
clientEpoch: clientEpoch,
|
|
connectedAtNs: time.Now().UnixNano(),
|
|
}
|
|
} else {
|
|
alreadyKnown = true
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
func (fs *FilerServer) deleteClient(scope string, clientName string, clientId int32, clientEpoch int32) {
|
|
glog.V(0).Infof("- %v listener %v clientId %v clientEpoch %v", scope, clientName, clientId, clientEpoch)
|
|
if clientId != 0 {
|
|
fs.knownListenersLock.Lock()
|
|
defer fs.knownListenersLock.Unlock()
|
|
epoch, found := fs.knownListeners[clientId]
|
|
if found && epoch <= clientEpoch {
|
|
delete(fs.knownListeners, clientId)
|
|
delete(fs.subscribers, clientId)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (fs *FilerServer) hasClient(clientId int32, clientEpoch int32) bool {
|
|
if clientId != 0 {
|
|
fs.knownListenersLock.Lock()
|
|
defer fs.knownListenersLock.Unlock()
|
|
epoch, found := fs.knownListeners[clientId]
|
|
if found && epoch <= clientEpoch {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|