From 4d8f948552dd70d7991a78b6ca3ea76ed1b61d14 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Carl=20H=C3=B6rberg?= Date: Thu, 1 Oct 2026 19:06:19 +0200 Subject: [PATCH 1/4] Signal memory pressure relief and expose PSI readings Pressure notifications only signal the onset. The monitor block now takes a Bool: true when pressure is detected, false when it's relieved. While under pressure the PSI file is polled every check_interval and relief is signalled once "some avg10" drops below release_below. Adds MemoryPressure.pressure to read the current PSI (watched .pressure file, the cgroup's memory.pressure, or /proc/pressure/memory) and MemoryPressure.parse for PSI data. The FIFO, socket and regular file watch loops are merged into one Watcher. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS --- README.md | 17 ++- spec/memory_pressure_spec.cr | 102 ++++++++++++- src/memory_pressure.cr | 281 ++++++++++++++++++----------------- 3 files changed, 256 insertions(+), 144 deletions(-) diff --git a/README.md b/README.md index 6a68d62..5545b05 100644 --- a/README.md +++ b/README.md @@ -44,12 +44,21 @@ SystemD.watchdog # Monitor memory pressure notifications from systemd # Enable with `MemoryPressureWatch=auto` and `MemoryPressureThresholdSec=1s` under `[Service]` -SystemD::MemoryPressure.monitor do - # Called when memory pressure is detected - # Take action like clearing caches, reducing memory usage, etc. - clear_caches +# The block is called with true when memory pressure is detected, and with +# false once the PSI "some avg10" drops below release_below (percent) +SystemD::MemoryPressure.monitor(release_below: 1.0) do |pressure| + if pressure + # Take action like clearing caches, reducing memory usage, etc. + clear_caches + pause_work + else + resume_work + end end +# Read the current pressure stall information (PSI) of the process' cgroup +SystemD::MemoryPressure.pressure.try &.some.avg10 + # Store FDs with the SystemD, they will be sent back # to the application when it restarts. Requires libsystemd clients = Array(TCPSocket).new diff --git a/spec/memory_pressure_spec.cr b/spec/memory_pressure_spec.cr index edca2a7..d369500 100644 --- a/spec/memory_pressure_spec.cr +++ b/spec/memory_pressure_spec.cr @@ -35,7 +35,7 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WRITE") wg = WaitGroup.new(1) - SystemD::MemoryPressure.monitor { wg.done } + SystemD::MemoryPressure.monitor { |pressure| wg.done if pressure } # Write to the FIFO to trigger memory pressure File.open(fifo_path, "w") do |f| @@ -61,7 +61,7 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WRITE") ch = Channel(Nil).new - SystemD::MemoryPressure.monitor { ch.send(nil) } + SystemD::MemoryPressure.monitor { |pressure| ch.send(nil) if pressure } # Accept the connection and send data client = server.accept @@ -90,7 +90,7 @@ describe SystemD::MemoryPressure do ENV["MEMORY_PRESSURE_WRITE"] = Base64.strict_encode(write_data) ch = Channel(Nil).new - SystemD::MemoryPressure.monitor { ch.send(nil) } + SystemD::MemoryPressure.monitor { |pressure| ch.send(nil) if pressure } # Accept the connection and read the threshold data client = server.accept @@ -124,7 +124,7 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WRITE") call_count = 0 - SystemD::MemoryPressure.monitor { call_count += 1 } + SystemD::MemoryPressure.monitor { |pressure| call_count += 1 if pressure } # Accept first connection and trigger pressure client1 = server.accept @@ -152,4 +152,98 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WATCH") end end + + it "parses PSI data" do + data = "some avg10=1.50 avg60=0.25 avg300=0.00 total=12345\nfull avg10=0.75 avg60=0.10 avg300=0.00 total=678\n" + pressure = SystemD::MemoryPressure.parse(data).should_not be_nil + pressure.some.should eq SystemD::MemoryPressure::Stall.new(1.5, 0.25, 0.0, 12345u64) + pressure.full.try(&.avg10).should eq 0.75 + end + + it "parses PSI data without a full line" do + pressure = SystemD::MemoryPressure.parse("some avg10=2.00 avg60=1.00 avg300=0.50 total=1\n").should_not be_nil + pressure.some.avg10.should eq 2.0 + pressure.full.should be_nil + end + + it "returns nil for malformed PSI data" do + SystemD::MemoryPressure.parse("garbage").should be_nil + SystemD::MemoryPressure.parse("some avg10=x avg60=1 avg300=1 total=1").should be_nil + end + + it "reads pressure from a file" do + path = File.tempname + File.write(path, "some avg10=3.00 avg60=2.00 avg300=1.00 total=99\n") + SystemD::MemoryPressure.pressure(path).try(&.some.total).should eq 99 + SystemD::MemoryPressure.pressure("/nonexistent/memory.pressure").should be_nil + ensure + File.delete?(path) if path + end + + it "prefers a watched .pressure file" do + path = File.tempname(suffix: ".pressure") + File.write(path, "some avg10=0.00 avg60=0.00 avg300=0.00 total=0\n") + ENV["MEMORY_PRESSURE_WATCH"] = path + SystemD::MemoryPressure.pressure_path.should eq path + ensure + ENV.delete("MEMORY_PRESSURE_WATCH") + File.delete?(path) if path + end + + it "calls on_relief once pressure drops below the threshold" do + fifo_path = File.tempname + begin + LibC.mkfifo(fifo_path, 0o600).should eq 0 + ENV["MEMORY_PRESSURE_WATCH"] = fifo_path + ENV.delete("MEMORY_PRESSURE_WRITE") + + pressured = Channel(Nil).new(1) + relieved = Channel(Nil).new(1) + SystemD::MemoryPressure.monitor(Float64::MAX, 10.milliseconds) { |pressure| pressure ? pressured.send(nil) : relieved.send(nil) } + + File.open(fifo_path, "w") do |f| + f.sync = true + f.print "pressure" + end + + pressured.receive + select + when relieved.receive + when timeout(5.seconds) + fail "on_relief was not called" + end + ensure + File.delete(fifo_path) if File.exists?(fifo_path) + ENV.delete("MEMORY_PRESSURE_WATCH") + end + end + + it "does not call on_relief while pressure persists" do + pending! "no PSI" unless File.file?("/proc/pressure/memory") + fifo_path = File.tempname + begin + LibC.mkfifo(fifo_path, 0o600).should eq 0 + ENV["MEMORY_PRESSURE_WATCH"] = fifo_path + ENV.delete("MEMORY_PRESSURE_WRITE") + + pressured = Channel(Nil).new(1) + relieved = Channel(Nil).new(1) + SystemD::MemoryPressure.monitor(0.0, 10.milliseconds) { |pressure| pressure ? pressured.send(nil) : relieved.send(nil) } + + File.open(fifo_path, "w") do |f| + f.sync = true + f.print "pressure" + end + + pressured.receive + select + when relieved.receive + fail "on_relief called while under pressure" + when timeout(200.milliseconds) + end + ensure + File.delete(fifo_path) if File.exists?(fifo_path) + ENV.delete("MEMORY_PRESSURE_WATCH") + end + end end diff --git a/src/memory_pressure.cr b/src/memory_pressure.cr index 1002240..fa53677 100644 --- a/src/memory_pressure.cr +++ b/src/memory_pressure.cr @@ -24,8 +24,20 @@ module SystemD module MemoryPressure Log = ::Log.for("systemd.memory_pressure") - # The block is called when memory pressure is detected - def self.monitor(&block : ->) + # One line of a PSI file, avg values are percentages + record Stall, avg10 : Float64, avg60 : Float64, avg300 : Float64, total : UInt64 + + # Pressure stall information, `full` is missing on older kernels + record Pressure, some : Stall, full : Stall? + + # The block is called with `true` when memory pressure is detected and + # with `false` when it is relieved. Notifications only signal the onset, + # so while under pressure the PSI file is polled every *check_interval* + # and relief is signalled once its "some avg10" drops below + # *release_below* (percent), or at the first check if no PSI file can be + # read. + def self.monitor(release_below : Float64 = 1.0, check_interval : Time::Span = 1.second, + &block : Bool ->) watch_path = ENV["MEMORY_PRESSURE_WATCH"]? return unless watch_path @@ -36,184 +48,181 @@ module SystemD Fiber::ExecutionContext::Isolated.new("Memory Pressure Monitor") do begin - monitor_internal(watch_path, &block) + watcher = Watcher.new(watch_path, decode_write_data, release_below, check_interval, block) + watcher.run rescue ex Log.error(exception: ex) { "Memory pressure monitoring failed" } end end end - private def self.monitor_internal(watch_path : String, &block : ->) - write_data = decode_write_data - file_type = determine_file_type(watch_path) - - case file_type - when :regular - monitor_regular_file(watch_path, write_data, &block) - when :fifo - monitor_fifo(watch_path, write_data, &block) - when :socket - monitor_socket(watch_path, write_data, &block) - else - Log.warn { "Unknown file type for #{watch_path}, attempting as regular file" } - monitor_regular_file(watch_path, write_data, &block) + # Current memory pressure, see `pressure_path` for which file is read + def self.pressure : Pressure? + if path = pressure_path + pressure(path) end end - private def self.decode_write_data : Bytes? - if encoded = ENV["MEMORY_PRESSURE_WRITE"]? - Base64.decode(encoded) - end + def self.pressure(path : String) : Pressure? + parse(File.read(path)) + rescue File::Error end - private def self.determine_file_type(path : String) : Symbol - result = LibC.stat(path, out stat) - raise IO::Error.from_errno("stat failed") if result != 0 - - file_mode = stat.st_mode & LibC::S_IFMT - case file_mode - when LibC::S_IFREG - :regular - when LibC::S_IFIFO - :fifo - when LibC::S_IFSOCK - :socket - else - :unknown + # The watched file if it is a regular PSI file, otherwise the process' + # cgroup v2 memory.pressure, falling back to the system wide one + def self.pressure_path : String? + if (watch_path = ENV["MEMORY_PRESSURE_WATCH"]?) && watch_path.ends_with?(".pressure") && + File.file?(watch_path) + return watch_path + end + if cgroup = File.read("/proc/self/cgroup")[/0::(.*)\n/, 1]? + path = "/sys/fs/cgroup#{cgroup}/memory.pressure" + return path if File.file?(path) end + "/proc/pressure/memory" if File.file?("/proc/pressure/memory") + rescue File::Error + "/proc/pressure/memory" if File.file?("/proc/pressure/memory") end - private def self.monitor_regular_file(path : String, write_data : Bytes?, &block : ->) - Log.info { "Monitoring memory pressure on regular file: #{path}" } - - fd = LibC.open(path, LibC::O_RDWR) - raise IO::Error.from_errno("open failed") if fd < 0 - - begin - # Write the pressure threshold data if provided - if write_data - written = LibC.write(fd, write_data, write_data.size) - raise IO::Error.from_errno("write failed") if written < 0 + def self.parse(data : String) : Pressure? + some = full = nil + data.each_line do |line| + kind, _, rest = line.partition(' ') + stall = parse_stall(rest) || next + case kind + when "some" then some = stall + when "full" then full = stall end + end + Pressure.new(some, full) if some + end - poll_fd = LibC::PollFD.new - poll_fd.fd = fd - poll_fd.events = LibC::POLLPRI - - loop do - result = LibC.poll(pointerof(poll_fd), 1, -1) # -1 = infinite timeout - if result < 0 - next if Errno.value == Errno::EINTR - raise IO::Error.from_errno("poll failed") - elsif result > 0 && (poll_fd.revents & LibC::POLLPRI) != 0 - handle_memory_pressure(&block) - # For regular files, we don't read from the FD - end + private def self.parse_stall(fields : String) : Stall? + avg10 = avg60 = avg300 = nil + total = nil + fields.split(' ', remove_empty: true) do |field| + key, _, value = field.partition('=') + case key + when "avg10" then avg10 = value.to_f64? + when "avg60" then avg60 = value.to_f64? + when "avg300" then avg300 = value.to_f64? + when "total" then total = value.to_u64? end - ensure - LibC.close(fd) end + Stall.new(avg10, avg60, avg300, total) if avg10 && avg60 && avg300 && total end - private def self.monitor_fifo(path : String, write_data : Bytes?, &block : ->) - Log.info { "Monitoring memory pressure on FIFO: #{path}" } + private def self.decode_write_data : Bytes? + if encoded = ENV["MEMORY_PRESSURE_WRITE"]? + Base64.decode(encoded) + end + end - fd = LibC.open(path, LibC::O_RDWR) - raise IO::Error.from_errno("open failed") if fd < 0 + private class Watcher + enum Kind + Regular + FIFO + Socket + end - begin - # Write the pressure threshold data if provided - if write_data - written = LibC.write(fd, write_data, write_data.size) - raise IO::Error.from_errno("write failed") if written < 0 - end + @kind : Kind + @fd : Int32 + @under_pressure = false - poll_fd = LibC::PollFD.new - poll_fd.fd = fd - poll_fd.events = LibC::POLLIN + def initialize(@path : String, @write_data : Bytes?, @release_below : Float64, + @check_interval : Time::Span, @callback : Bool ->) + @kind = determine_kind(@path) + Log.info { "Monitoring memory pressure on #{@kind.to_s.downcase}: #{@path}" } + @fd = open + end + def run + poll_fd = LibC::PollFD.new + poll_fd.fd = @fd + poll_fd.events = @kind.regular? ? LibC::POLLPRI : LibC::POLLIN loop do - result = LibC.poll(pointerof(poll_fd), 1, -1) + result = LibC.poll(pointerof(poll_fd), 1, timeout) if result < 0 next if Errno.value == Errno::EINTR raise IO::Error.from_errno("poll failed") - end - if result > 0 && (poll_fd.revents & LibC::POLLIN) != 0 - handle_memory_pressure(&block) - # Read and discard data from FIFO - buf = uninitialized UInt8[4096] - bytes_read = LibC.read(fd, buf, 4096) - # EOF (0) or error is expected, continue polling + elsif result == 0 + check_relief + elsif (poll_fd.revents & poll_fd.events) != 0 + @under_pressure = true + Log.info { "Memory pressure detected" } + @callback.call(true) + poll_fd.fd = @fd if drain end end ensure - LibC.close(fd) + LibC.close(@fd) end - end - - private def self.monitor_socket(path : String, write_data : Bytes?, &block : ->) - Log.info { "Monitoring memory pressure on Unix socket: #{path}" } - fd = connect_unix_socket(path) + private def timeout : Int32 + @under_pressure ? @check_interval.total_milliseconds.to_i : -1 + end - begin - # Write the pressure threshold data if provided - if write_data - written = LibC.write(fd, write_data, write_data.size) - raise IO::Error.from_errno("write failed") if written < 0 - end + private def check_relief + avg10 = (path = MemoryPressure.pressure_path) && MemoryPressure.pressure(path).try &.some.avg10 + return if avg10 && avg10 >= @release_below + @under_pressure = false + Log.info { "Memory pressure relieved" } + @callback.call(false) + end - poll_fd = LibC::PollFD.new - poll_fd.fd = fd - poll_fd.events = LibC::POLLIN + # Reads and discards the notification, returns true if the fd was replaced + private def drain : Bool + return false if @kind.regular? + buf = uninitialized UInt8[4096] + bytes_read = LibC.read(@fd, buf, buf.size) + # EOF or error is expected on a FIFO, a closed socket is reconnected + return false unless @kind.socket? && bytes_read <= 0 + LibC.close(@fd) + @fd = open + true + end - loop do - result = LibC.poll(pointerof(poll_fd), 1, -1) - if result < 0 - next if Errno.value == Errno::EINTR - raise IO::Error.from_errno("poll failed") - end - if result > 0 && (poll_fd.revents & LibC::POLLIN) != 0 - handle_memory_pressure(&block) - # Read and discard data from socket - buffer = uninitialized UInt8[4096] - bytes_read = LibC.read(fd, buffer, 4096) - if bytes_read <= 0 - # Connection closed, reconnect - LibC.close(fd) - fd = connect_unix_socket(path) - poll_fd.fd = fd - if write_data - written = LibC.write(fd, write_data, write_data.size) - raise IO::Error.from_errno("write failed after reconnect") if written < 0 - end - end + private def open : Int32 + fd = @kind.socket? ? connect_unix_socket(@path) : LibC.open(@path, LibC::O_RDWR) + raise IO::Error.from_errno("open failed") if fd < 0 + if data = @write_data + if LibC.write(fd, data, data.size) < 0 + LibC.close(fd) + raise IO::Error.from_errno("write failed") end end - ensure - LibC.close(fd) + fd end - end - private def self.connect_unix_socket(path : String) : Int32 - fd = LibC.socket(LibC::AF_UNIX, LibC::SOCK_STREAM, 0) - raise IO::Error.from_errno("socket creation failed") if fd < 0 + private def determine_kind(path : String) : Kind + result = LibC.stat(path, out stat) + raise IO::Error.from_errno("stat failed") if result != 0 + + case stat.st_mode & LibC::S_IFMT + when LibC::S_IFIFO then Kind::FIFO + when LibC::S_IFSOCK then Kind::Socket + when LibC::S_IFREG then Kind::Regular + else + Log.warn { "Unknown file type for #{path}, attempting as regular file" } + Kind::Regular + end + end - sockaddr = Pointer(LibC::SockaddrUn).malloc - sockaddr.value.sun_family = LibC::AF_UNIX.to_u16 - sockaddr.value.sun_path.to_unsafe.copy_from(path.to_unsafe, {path.bytesize + 1, sockaddr.value.sun_path.size}.min) + private def connect_unix_socket(path : String) : Int32 + fd = LibC.socket(LibC::AF_UNIX, LibC::SOCK_STREAM, 0) + raise IO::Error.from_errno("socket creation failed") if fd < 0 - if LibC.connect(fd, sockaddr.as(LibC::Sockaddr*), sizeof(LibC::SockaddrUn)) < 0 - LibC.close(fd) - raise IO::Error.from_errno("connect failed") - end + sockaddr = Pointer(LibC::SockaddrUn).malloc + sockaddr.value.sun_family = LibC::AF_UNIX.to_u16 + sockaddr.value.sun_path.to_unsafe.copy_from(path.to_unsafe, {path.bytesize + 1, sockaddr.value.sun_path.size}.min) - fd - end + if LibC.connect(fd, sockaddr.as(LibC::Sockaddr*), sizeof(LibC::SockaddrUn)) < 0 + LibC.close(fd) + raise IO::Error.from_errno("connect failed") + end - private def self.handle_memory_pressure(&block : ->) - Log.info { "Memory pressure detected" } - block.call + fd + end end end end From a4d45e4554e73cd8995ae98bd2cafb5a780188e3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Carl=20H=C3=B6rberg?= Date: Thu, 1 Oct 2026 19:56:55 +0200 Subject: [PATCH 2/4] Keep monitor unchanged, add watch for pressure and relief Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS --- README.md | 21 +++++++++-------- spec/memory_pressure_spec.cr | 45 ++++++++++++++++++++++++++++-------- src/memory_pressure.cr | 20 ++++++++++++---- 3 files changed, 61 insertions(+), 25 deletions(-) diff --git a/README.md b/README.md index 5545b05..1ba5a71 100644 --- a/README.md +++ b/README.md @@ -44,16 +44,17 @@ SystemD.watchdog # Monitor memory pressure notifications from systemd # Enable with `MemoryPressureWatch=auto` and `MemoryPressureThresholdSec=1s` under `[Service]` -# The block is called with true when memory pressure is detected, and with -# false once the PSI "some avg10" drops below release_below (percent) -SystemD::MemoryPressure.monitor(release_below: 1.0) do |pressure| - if pressure - # Take action like clearing caches, reducing memory usage, etc. - clear_caches - pause_work - else - resume_work - end +SystemD::MemoryPressure.monitor do + # Called when memory pressure is detected + # Take action like clearing caches, reducing memory usage, etc. + clear_caches +end + +# Notifications only signal the onset of pressure. watch also polls the PSI +# file while under pressure and calls the block with false once its +# "some avg10" drops below release_below (percent) +SystemD::MemoryPressure.watch(release_below: 1.0) do |pressure| + pressure ? pause_work : resume_work end # Read the current pressure stall information (PSI) of the process' cgroup diff --git a/spec/memory_pressure_spec.cr b/spec/memory_pressure_spec.cr index d369500..7f31cf0 100644 --- a/spec/memory_pressure_spec.cr +++ b/spec/memory_pressure_spec.cr @@ -35,7 +35,7 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WRITE") wg = WaitGroup.new(1) - SystemD::MemoryPressure.monitor { |pressure| wg.done if pressure } + SystemD::MemoryPressure.monitor { wg.done } # Write to the FIFO to trigger memory pressure File.open(fifo_path, "w") do |f| @@ -61,7 +61,7 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WRITE") ch = Channel(Nil).new - SystemD::MemoryPressure.monitor { |pressure| ch.send(nil) if pressure } + SystemD::MemoryPressure.monitor { ch.send(nil) } # Accept the connection and send data client = server.accept @@ -90,7 +90,7 @@ describe SystemD::MemoryPressure do ENV["MEMORY_PRESSURE_WRITE"] = Base64.strict_encode(write_data) ch = Channel(Nil).new - SystemD::MemoryPressure.monitor { |pressure| ch.send(nil) if pressure } + SystemD::MemoryPressure.monitor { ch.send(nil) } # Accept the connection and read the threshold data client = server.accept @@ -124,7 +124,7 @@ describe SystemD::MemoryPressure do ENV.delete("MEMORY_PRESSURE_WRITE") call_count = 0 - SystemD::MemoryPressure.monitor { |pressure| call_count += 1 if pressure } + SystemD::MemoryPressure.monitor { call_count += 1 } # Accept first connection and trigger pressure client1 = server.accept @@ -190,7 +190,32 @@ describe SystemD::MemoryPressure do File.delete?(path) if path end - it "calls on_relief once pressure drops below the threshold" do + it "monitor only calls the block on pressure" do + fifo_path = File.tempname + begin + LibC.mkfifo(fifo_path, 0o600).should eq 0 + ENV["MEMORY_PRESSURE_WATCH"] = fifo_path + ENV.delete("MEMORY_PRESSURE_WRITE") + + calls = Atomic(Int32).new(0) + SystemD::MemoryPressure.monitor { calls.add(1) } + + File.open(fifo_path, "w") do |f| + f.sync = true + f.print "pressure" + end + + sleep 100.milliseconds + # the written notification, plus the event a FIFO reports at startup + calls.get.should be <= 2 + calls.get.should be >= 1 + ensure + File.delete(fifo_path) if File.exists?(fifo_path) + ENV.delete("MEMORY_PRESSURE_WATCH") + end + end + + it "watch signals relief once pressure drops below the threshold" do fifo_path = File.tempname begin LibC.mkfifo(fifo_path, 0o600).should eq 0 @@ -199,7 +224,7 @@ describe SystemD::MemoryPressure do pressured = Channel(Nil).new(1) relieved = Channel(Nil).new(1) - SystemD::MemoryPressure.monitor(Float64::MAX, 10.milliseconds) { |pressure| pressure ? pressured.send(nil) : relieved.send(nil) } + SystemD::MemoryPressure.watch(Float64::MAX, 10.milliseconds) { |pressure| pressure ? pressured.send(nil) : relieved.send(nil) } File.open(fifo_path, "w") do |f| f.sync = true @@ -210,7 +235,7 @@ describe SystemD::MemoryPressure do select when relieved.receive when timeout(5.seconds) - fail "on_relief was not called" + fail "relief was not signalled" end ensure File.delete(fifo_path) if File.exists?(fifo_path) @@ -218,7 +243,7 @@ describe SystemD::MemoryPressure do end end - it "does not call on_relief while pressure persists" do + it "watch does not signal relief while pressure persists" do pending! "no PSI" unless File.file?("/proc/pressure/memory") fifo_path = File.tempname begin @@ -228,7 +253,7 @@ describe SystemD::MemoryPressure do pressured = Channel(Nil).new(1) relieved = Channel(Nil).new(1) - SystemD::MemoryPressure.monitor(0.0, 10.milliseconds) { |pressure| pressure ? pressured.send(nil) : relieved.send(nil) } + SystemD::MemoryPressure.watch(0.0, 10.milliseconds) { |pressure| pressure ? pressured.send(nil) : relieved.send(nil) } File.open(fifo_path, "w") do |f| f.sync = true @@ -238,7 +263,7 @@ describe SystemD::MemoryPressure do pressured.receive select when relieved.receive - fail "on_relief called while under pressure" + fail "relief signalled while under pressure" when timeout(200.milliseconds) end ensure diff --git a/src/memory_pressure.cr b/src/memory_pressure.cr index fa53677..7a0c928 100644 --- a/src/memory_pressure.cr +++ b/src/memory_pressure.cr @@ -30,14 +30,23 @@ module SystemD # Pressure stall information, `full` is missing on older kernels record Pressure, some : Stall, full : Stall? + # The block is called when memory pressure is detected + def self.monitor(&block : ->) + start(nil, Time::Span.zero, ->(pressure : Bool) { block.call if pressure }) + end + # The block is called with `true` when memory pressure is detected and # with `false` when it is relieved. Notifications only signal the onset, # so while under pressure the PSI file is polled every *check_interval* # and relief is signalled once its "some avg10" drops below # *release_below* (percent), or at the first check if no PSI file can be # read. - def self.monitor(release_below : Float64 = 1.0, check_interval : Time::Span = 1.second, - &block : Bool ->) + def self.watch(release_below : Float64 = 1.0, check_interval : Time::Span = 1.second, + &block : Bool ->) + start(release_below, check_interval, block) + end + + private def self.start(release_below : Float64?, check_interval : Time::Span, block : Bool ->) watch_path = ENV["MEMORY_PRESSURE_WATCH"]? return unless watch_path @@ -129,7 +138,7 @@ module SystemD @fd : Int32 @under_pressure = false - def initialize(@path : String, @write_data : Bytes?, @release_below : Float64, + def initialize(@path : String, @write_data : Bytes?, @release_below : Float64?, @check_interval : Time::Span, @callback : Bool ->) @kind = determine_kind(@path) Log.info { "Monitoring memory pressure on #{@kind.to_s.downcase}: #{@path}" } @@ -159,12 +168,13 @@ module SystemD end private def timeout : Int32 - @under_pressure ? @check_interval.total_milliseconds.to_i : -1 + @under_pressure && @release_below ? @check_interval.total_milliseconds.to_i : -1 end private def check_relief + release_below = @release_below || return avg10 = (path = MemoryPressure.pressure_path) && MemoryPressure.pressure(path).try &.some.avg10 - return if avg10 && avg10 >= @release_below + return if avg10 && avg10 >= release_below @under_pressure = false Log.info { "Memory pressure relieved" } @callback.call(false) From 98c6a8fed75ff6578b11784fca90579382e5c488 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Carl=20H=C3=B6rberg?= Date: Thu, 1 Oct 2026 20:01:32 +0200 Subject: [PATCH 3/4] Expect exactly one monitor call per FIFO notification Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS --- spec/memory_pressure_spec.cr | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/spec/memory_pressure_spec.cr b/spec/memory_pressure_spec.cr index 7f31cf0..393ceb1 100644 --- a/spec/memory_pressure_spec.cr +++ b/spec/memory_pressure_spec.cr @@ -206,9 +206,7 @@ describe SystemD::MemoryPressure do end sleep 100.milliseconds - # the written notification, plus the event a FIFO reports at startup - calls.get.should be <= 2 - calls.get.should be >= 1 + calls.get.should eq 1 ensure File.delete(fifo_path) if File.exists?(fifo_path) ENV.delete("MEMORY_PRESSURE_WATCH") From b91ac8918fd90dbeb67223bb445e33a1ecca2157 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Carl=20H=C3=B6rberg?= Date: Thu, 1 Oct 2026 20:22:34 +0200 Subject: [PATCH 4/4] Document memory pressure monitoring in its own README section Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS --- README.md | 89 +++++++++++++++++++++++++++++++++++++++++++------------ 1 file changed, 70 insertions(+), 19 deletions(-) diff --git a/README.md b/README.md index 1ba5a71..b2c29dd 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,7 @@ Man pages: https://man7.org/linux/man-pages/man3/sd_pid_notify.3.html https://man7.org/linux/man-pages/man3/sd_listen_fds.3.html -https://systemd.io/MEMORY_PRESSURE/ +https://systemd.io/PRESSURE/ ## Installation @@ -42,24 +42,6 @@ end # Enable systemd watchdog support with `WatchdogSec=5` under `[Service]` SystemD.watchdog -# Monitor memory pressure notifications from systemd -# Enable with `MemoryPressureWatch=auto` and `MemoryPressureThresholdSec=1s` under `[Service]` -SystemD::MemoryPressure.monitor do - # Called when memory pressure is detected - # Take action like clearing caches, reducing memory usage, etc. - clear_caches -end - -# Notifications only signal the onset of pressure. watch also polls the PSI -# file while under pressure and calls the block with false once its -# "some avg10" drops below release_below (percent) -SystemD::MemoryPressure.watch(release_below: 1.0) do |pressure| - pressure ? pause_work : resume_work -end - -# Read the current pressure stall information (PSI) of the process' cgroup -SystemD::MemoryPressure.pressure.try &.some.avg10 - # Store FDs with the SystemD, they will be sent back # to the application when it restarts. Requires libsystemd clients = Array(TCPSocket).new @@ -82,6 +64,75 @@ SystemD.named_listeners do |socket, name| end ``` +## Memory pressure + +systemd can notify a service when its cgroup is under memory pressure. Enable it under `[Service]`: + +```ini +MemoryPressureWatch=auto +# Optional, how long tasks may stall on memory within a 2s window before a notification (default 200ms) +MemoryPressureThresholdSec=200ms +``` + +systemd then passes `MEMORY_PRESSURE_WATCH` (the cgroup's `memory.pressure` file) and `MEMORY_PRESSURE_WRITE` (the trigger to register on it) to the service. Without `MEMORY_PRESSURE_WATCH`, or when it's `/dev/null`, monitoring is disabled and the blocks are never called. + +### Reacting to pressure + +`monitor` calls its block each time memory pressure is detected: + +```crystal +SystemD::MemoryPressure.monitor do + # Take action like clearing caches, reducing memory usage, etc. + clear_caches +end +``` + +### Pressure and relief + +Notifications only signal the onset of pressure, never that it's over. Use `watch` when the application backs off under pressure and must know when to resume. Its block is called with `true` when pressure is detected, and with `false` when it's relieved: + +```crystal +SystemD::MemoryPressure.watch(release_below: 1.0, check_interval: 1.second) do |pressure| + if pressure + pause_work + else + resume_work + end +end +``` + +While under pressure, `watch` reads the PSI file every `check_interval`. It signals relief once `some avg10` drops below `release_below`, which is the percentage of the last 10 seconds that some task stalled on memory. If no PSI file can be read, relief is signalled at the first check. + +`avg10` is a 10-second running average, so it lags behind the stalls. Together with the onset threshold (200ms of stalls in 2s by default), that gives hysteresis: pressure has to be well past before relief is signalled. + +Both blocks run in a dedicated thread (an isolated execution context), so keep them short. For example, set a flag and act on it from the application's own fibers. + +### Reading PSI values + +```crystal +if pressure = SystemD::MemoryPressure.pressure + pressure.some.avg10 # % of time some task stalled on memory, 10s average + pressure.some.avg60 + pressure.some.avg300 + pressure.some.total # total stall time in microseconds + pressure.full.try &.avg10 # % of time all tasks stalled, nil if there's no "full" line +end +``` + +`pressure` reads the watched file when it's a `.pressure` file, otherwise the process' cgroup v2 `memory.pressure`, falling back to the system-wide `/proc/pressure/memory`. `SystemD::MemoryPressure.parse(string)` parses PSI file contents directly. + +### Testing locally + +`MEMORY_PRESSURE_WATCH` can point at a regular PSI file, a FIFO or a Unix socket. Writing to a FIFO simulates a notification: + +```sh +mkfifo /tmp/pressure +MEMORY_PRESSURE_WATCH=/tmp/pressure MEMORY_PRESSURE_WRITE= ./my-app & +echo > /tmp/pressure +``` + +Processes started from a desktop session usually inherit both variables from the desktop's own service, so a shell may already have them set. Clear `MEMORY_PRESSURE_WRITE` when watching a FIFO, as above. Otherwise its trigger is written into the FIFO and read back as a pressure notification. + ## Contributing 1. Fork it ()