Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion docs/pipeline/runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -2912,7 +2912,10 @@ header. Grouped by cause, largest first:
`Main.run_rack` in `GzipCache` (Static stays outside so `/cable`
hijack is never compressed; CSS/JS stay identity; identical HTML is
not deflated on every request). spinel tep gzips inline bodies when
`Accept-Encoding` includes gzip, from the same identity-body cache.
`Accept-Encoding` includes gzip, from a cache keyed by SHA-256 of
the identity body (CRuby keys by the body itself — MRI's string
hash is cheaper than SHA-256 here). Gzip runs outside the lock on
both lanes.
Re-run `scripts/campfire-http-shape` before treating the 65
Content-Encoding misses as current.
- **Rails' `Rack::ETag` / `Rack::ConditionalGet` are absent**: no weak
Expand Down
21 changes: 15 additions & 6 deletions runtime/spinel/scaffold/ruby_overlay/runtime/gzip_cache.rb
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,16 @@
#
# Rack::Deflater compresses every response. A campfire room page is the
# same ~420 KB HTML for every wrk GET that shares a session, so that is
# the same deflate over and over. Key by the identity bytes; a hit is
# the compressed copy. Misses gzip once and store.
# the same deflate over and over. Key by the identity bytes: MRI's
# string hash of a 420 KB body is cheaper than SHA-256 of the same
# bytes (measured: digest-keyed cache dropped /rooms/1 from ~1725 to
# ~1140 req/s). The Spinel twin keys by digest because its Hash hashes
# the whole key under the lock and has no GVL.
#
# Gzip itself runs outside the lock. Holding Mutex across Zlib.gzip
# serialized every miss onto one core. Two threads that miss the same
# body both gzip and one write wins — a duplicate deflate, not a
# wrong body.
#
# HTML only, same skips as tep: 1xx/204/304, HEAD, already-encoded,
# small, listed binary types. Wraps run_rack only so /cable's hijack
Expand Down Expand Up @@ -50,16 +58,17 @@ def self.maybe_gzip(env, status, headers, body)
end

def self.compress(raw)
hit = nil
@mutex.synchronize { hit = @store[raw] }
return hit unless hit.nil?
gz = Zlib.gzip(raw)
@mutex.synchronize do
hit = @store[raw]
return hit unless hit.nil?
if @store.size >= MAX_ENTRIES
@store.clear
end
gz = Zlib.gzip(raw)
@store[raw] = gz
gz
end
gz
end

def self.join_body(body)
Expand Down
46 changes: 34 additions & 12 deletions runtime/spinel/scaffold/ruby_overlay/runtime/rails_cache.rb
Original file line number Diff line number Diff line change
Expand Up @@ -68,19 +68,37 @@ def fetch_str(key, ttl, &block)
# have to capture the accumulator, and on the AOT lane a captured
# block dissolves into a heap poly proc (matz/spinel#4245).
#
# Thin delegates, as `fetch_str` above is — `read`/`write` already
# hold the Mutex and already dup Strings on both sides. `read_str`
# narrows to String because the caller appends the answer to a
# string builder; a non-String under that key is a MISS rather than
# a TypeError at the append, and the write that follows corrects it.
# Thin delegates, as `fetch_str` above is — `write` still dups the
# String into the store. `read_str` does NOT dup on the way out:
# a view's `<% cache %>` only appends the hit (`io << hit`), and a
# campfire room page is ~40 message fragments. DupCoder's read-side
# copy was 40 extra 2–5 KB allocations per wrk GET that shared
# nothing with mutation safety, because the stored copy is already
# isolated by the write-side dup. A non-String under that key is a
# MISS rather than a TypeError at the append, and the write that
# follows corrects it. `read` (the untyped half) still dups, so a
# caller that mutates a fetched String cannot corrupt the store.
def read_str(key)
value = read(key)
value.is_a?(String) ? value : nil
@mutex.synchronize do
entry = @data[key.to_s]
return nil if entry.nil?
if expired?(entry)
@data.delete(key.to_s)
return nil
end
encoded = entry[0]
encoded.is_a?(String) ? encoded : nil
Comment thread
coderabbitai[bot] marked this conversation as resolved.
end
end

def write_str(key, value, ttl)
write(key, value, ttl.to_i > 0 ? { expires_in: ttl.to_i } : {})
value
# Dup into the store, then freeze that copy. The caller's
# accumulator stays mutable; the stored fragment is shared
# across hits without a read-side dup.
s = (value.is_a?(String) ? value.dup : value.to_s).freeze
expires_at = ttl.to_i > 0 ? monotonic_now + ttl.to_i : nil
@mutex.synchronize { @data[key.to_s] = [s, expires_at] }
s
end

# The counter behind `rate_limit` (`ActionController::RateLimiter`),
Expand All @@ -95,10 +113,10 @@ def increment_str(key, ttl)
entry = @data[k]
if entry && !expired?(entry) && entry[0].is_a?(String)
n = entry[0].to_i + 1
@data[k] = [n.to_s, entry[1]]
@data[k] = [n.to_s.freeze, entry[1]]
n
else
@data[k] = ["1", ttl.to_i > 0 ? monotonic_now + ttl.to_i : nil]
@data[k] = ["1".freeze, ttl.to_i > 0 ? monotonic_now + ttl.to_i : nil]
1
end
end
Expand All @@ -120,7 +138,11 @@ def write(key, value, opts = {})
expires_at = nil
ttl = opts[:expires_in]
expires_at = monotonic_now + ttl.to_i if ttl
encoded = value.is_a?(String) ? value.dup : [Marshal.dump(value)]
# Freeze the stored String so `read_str` can hand it back without
# a copy. `read` still dups (decode), so an untyped caller that
# mutates what it fetched cannot corrupt the store — the same
# contract `write_str` already keeps.
encoded = value.is_a?(String) ? value.dup.freeze : [Marshal.dump(value)]
@mutex.synchronize { @data[key.to_s] = [encoded, expires_at] }
value
end
Expand Down
1 change: 1 addition & 0 deletions runtime/spinel/tep/tep.rb
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
require "digest"
require "zlib"
require_relative "tep_core"
require_relative "url"
Expand Down
44 changes: 28 additions & 16 deletions runtime/spinel/tep/tep_core.rb
Original file line number Diff line number Diff line change
@@ -1,3 +1,6 @@
require "digest"
require "zlib"

module Tep
# The name the server announces itself by. scaffold/main.rb sets
# Tep::APP.name to the app's own name (the underscored module that
Expand Down Expand Up @@ -172,29 +175,38 @@ def self.q_is_zero?(s, from, to)
true
end

# Gzip of an identity body, keyed by the identity bytes. A campfire
# room page is the same HTML for every wrk GET that shares a session;
# without this, Zlib.gzip runs on every request and is the measured
# cliff (1984 → 694 req/s). Cap is a COUNT so a bound does not need
# an LRU touch on the read path. The lock is the fragment-cache one:
# a green thread can be descheduled inside Hash#[]=.
# Gzip of an identity body, keyed by SHA-256 of the identity bytes.
# A campfire room page is the same HTML for every wrk GET that shares
# a session; without this, Zlib.gzip runs on every request and is the
# measured cliff (1984 → 694 req/s). Keying on the raw 420 KB body
# hashed and compared that whole string under the lock on every hit;
# the digest is 64 hex chars and is computed outside the lock.
#
# Gzip itself also runs outside the lock. Holding Mutex across
# Zlib.gzip serialized every miss onto one green thread — Spinel's
# lane has no GVL, so that was the whole CPU. Two threads that miss
# the same body both gzip and one write wins.
#
# Cap is a COUNT so a bound does not need an LRU touch on the read
# path. The lock is still required: a green thread can be descheduled
# inside Hash#[]=.
GZIP_CACHE_MAX = 64
GZIP_LOCK = Mutex.new
@gzip_bodies = Hash.new("")

def self.gzip_cached(raw)
gz = ""
key = Digest::SHA256.hexdigest(raw)
hit = ""
GZIP_LOCK.synchronize do
hit = @gzip_bodies[raw]
if hit.length > 0
gz = hit
else
if @gzip_bodies.size >= GZIP_CACHE_MAX
@gzip_bodies = Hash.new("")
end
gz = Zlib.gzip(raw)
@gzip_bodies[raw] = gz
hit = @gzip_bodies[key]
end
return hit if hit.length > 0
gz = Zlib.gzip(raw)
GZIP_LOCK.synchronize do
if @gzip_bodies.size >= GZIP_CACHE_MAX
@gzip_bodies = Hash.new("")
end
@gzip_bodies[key] = gz
end
gz
end
Expand Down
126 changes: 126 additions & 0 deletions tests/gzip_cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,3 +61,129 @@ puts "ALL OK"
);
assert!(out.status.success(), "driver exited {:?}", out.status.code());
}

#[test]
fn distinct_bodies_do_not_share_a_gzip() {
let root = Path::new(env!("CARGO_MANIFEST_DIR"));
let script = r#"
require_relative "runtime/spinel/scaffold/ruby_overlay/runtime/gzip_cache"

a_body = "a" * 128
b_body = "b" * 128
a = GzipCache.wrap(lambda { |_e| [200, { "content-type" => "text/html" }, [a_body]] })
b = GzipCache.wrap(lambda { |_e| [200, { "content-type" => "text/html" }, [b_body]] })
env = { "REQUEST_METHOD" => "GET", "HTTP_ACCEPT_ENCODING" => "gzip" }
ga = a.call(env)
gb = b.call(env)
raise "same gzip" if ga[2][0] == gb[2][0]
raise "a not gzip" unless ga[1]["content-encoding"] == "gzip"
raise "b not gzip" unless gb[1]["content-encoding"] == "gzip"
puts "ALL OK"
"#;
let out = Command::new("ruby")
.arg("-e")
.arg(script)
.current_dir(root)
.output()
.expect("ruby is on PATH");
let stdout = String::from_utf8_lossy(&out.stdout);
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(
stdout.contains("ALL OK"),
"distinct-body gzip cache failed\n=== stdout ===\n{stdout}\n=== stderr ===\n{stderr}"
);
assert!(out.status.success(), "driver exited {:?}", out.status.code());
}

#[test]
fn tep_gzip_cached_hits_on_digest_not_body_pointer() {
let root = Path::new(env!("CARGO_MANIFEST_DIR"));
let script = r#"
require "digest"
require "zlib"
require_relative "runtime/spinel/tep/tep_core"

n = 0
orig = Zlib.method(:gzip)
Zlib.define_singleton_method(:gzip) do |raw|
n += 1
orig.call(raw)
end

a = "y" * 128
b = "y" * 128
raise "same object" if a.equal?(b)
ga = Tep.gzip_cached(a)
gb = Tep.gzip_cached(b)
raise "gzipped #{n} times" unless n == 1
raise "bodies differ" unless ga == gb
gc = Tep.gzip_cached("z" * 128)
raise "distinct collided" if gc == ga
raise "second body not gzipped" unless n == 2
puts "ALL OK"
"#;
let out = Command::new("ruby")
.arg("-e")
.arg(script)
.current_dir(root)
.output()
.expect("ruby is on PATH");
let stdout = String::from_utf8_lossy(&out.stdout);
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(
stdout.contains("ALL OK"),
"tep gzip cache failed\n=== stdout ===\n{stdout}\n=== stderr ===\n{stderr}"
);
assert!(out.status.success(), "driver exited {:?}", out.status.code());
}

#[test]
fn overlay_read_str_does_not_dup_a_cached_fragment() {
let root = Path::new(env!("CARGO_MANIFEST_DIR"));
let script = r#"
require_relative "runtime/spinel/scaffold/ruby_overlay/runtime/rails_cache"

store = Rails::MemoryStore.new
frag = "message-html" * 32
store.write_str("k", frag, 0)
hit = store.read_str("k")
raise "miss" if hit.nil?
raise "duped" unless hit.equal?(store.read_str("k"))
raise "mutated store" unless hit.frozen?
begin
hit << "x"
raise "frozen fragment was mutable"
rescue FrozenError
end
other = store.read("k")
raise "untyped read must still dup" if other.equal?(hit)
other << "x"
raise "store corrupted" unless store.read_str("k") == frag
# write (untyped) also freezes, so a later read_str cannot mutate
# the shared entry — CodeRabbit on #432.
store.write("k2", "plain")
hit2 = store.read_str("k2")
raise "write miss" if hit2.nil?
raise "write not frozen" unless hit2.frozen?
begin
hit2 << "x"
raise "write-path fragment was mutable"
rescue FrozenError
end
raise "write store corrupted" unless store.read_str("k2") == "plain"
puts "ALL OK"
"#;
let out = Command::new("ruby")
.arg("-e")
.arg(script)
.current_dir(root)
.output()
.expect("ruby is on PATH");
let stdout = String::from_utf8_lossy(&out.stdout);
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(
stdout.contains("ALL OK"),
"overlay read_str failed\n=== stdout ===\n{stdout}\n=== stderr ===\n{stderr}"
);
assert!(out.status.success(), "driver exited {:?}", out.status.code());
}
Loading