mirror of https://github.com/sbt/sbt.git
The byteStreamStub lazy val applied withDeadlineAfter once, baking in an absolute deadline that was reused for the store's whole (session-long) lifetime. remoteTimeoutInSec after the first blob transfer, every later ByteStream read/write was rejected with DEADLINE_EXCEEDED, so the remote cache could neither upload nor download blobs >chunkSizeBytes for the rest of the sbt server's life. Derive a fresh stub with the deadline per RPC so each call gets its own relative timeout. Co-authored-by: Yannick Heiber <yhe@famly.co> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
3c350df583
commit
35e99d27c1
|
|
@ -151,12 +151,16 @@ class GrpcActionCacheStore(
|
||||||
case _ => casStub0
|
case _ => casStub0
|
||||||
lazy val byteStreamStub0 = ByteStreamGrpc.newStub(channel)
|
lazy val byteStreamStub0 = ByteStreamGrpc.newStub(channel)
|
||||||
lazy val byteStreamStub = remoteHeaders match
|
lazy val byteStreamStub = remoteHeaders match
|
||||||
case x :: xs =>
|
case x :: xs => byteStreamStub0.withCallCredentials(creds)
|
||||||
byteStreamStub0
|
case _ => byteStreamStub0
|
||||||
.withCallCredentials(creds)
|
|
||||||
.withDeadlineAfter(remoteTimeoutInSec, TimeUnit.SECONDS)
|
// The deadline must be attached per call, not on the memoized stub above.
|
||||||
case _ =>
|
// withDeadlineAfter computes an absolute deadline at the moment it is called, so a stub
|
||||||
byteStreamStub0.withDeadlineAfter(remoteTimeoutInSec, TimeUnit.SECONDS)
|
// stored in a (session-lived) lazy val would expire remoteTimeoutInSec after first use and
|
||||||
|
// then reject every later call with DEADLINE_EXCEEDED. Deriving a fresh stub per RPC gives
|
||||||
|
// each call its own relative timeout.
|
||||||
|
private[internal] def byteStreamStubWithDeadline =
|
||||||
|
byteStreamStub.withDeadlineAfter(remoteTimeoutInSec, TimeUnit.SECONDS)
|
||||||
|
|
||||||
override def storeName: String = "remote"
|
override def storeName: String = "remote"
|
||||||
|
|
||||||
|
|
@ -279,7 +283,7 @@ class GrpcActionCacheStore(
|
||||||
def uploadBlob(blob: VirtualFile): Future[HashedVirtualFileRef] =
|
def uploadBlob(blob: VirtualFile): Future[HashedVirtualFileRef] =
|
||||||
val d = Digest(blob)
|
val d = Digest(blob)
|
||||||
withSingleResponse[ByteStreamProto.WriteResponse, HashedVirtualFileRef]: (p, resObs) =>
|
withSingleResponse[ByteStreamProto.WriteResponse, HashedVirtualFileRef]: (p, resObs) =>
|
||||||
val reqObs = byteStreamStub.write(resObs)
|
val reqObs = byteStreamStubWithDeadline.write(resObs)
|
||||||
val un = uploadName(d, UUID.randomUUID())
|
val un = uploadName(d, UUID.randomUUID())
|
||||||
var off: Long = 0L
|
var off: Long = 0L
|
||||||
try
|
try
|
||||||
|
|
@ -327,7 +331,7 @@ class GrpcActionCacheStore(
|
||||||
b.setResourceName(dn)
|
b.setResourceName(dn)
|
||||||
b.setReadOffset(0L)
|
b.setReadOffset(0L)
|
||||||
val req = b.build()
|
val req = b.build()
|
||||||
byteStreamStub.read(req, resObs)
|
byteStreamStubWithDeadline.read(req, resObs)
|
||||||
p.future
|
p.future
|
||||||
|
|
||||||
// helper function for many-to-one gRPC streaming
|
// helper function for many-to-one gRPC streaming
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,11 @@
|
||||||
package sbt
|
package sbt
|
||||||
package internal
|
package internal
|
||||||
|
|
||||||
|
import java.net.URI
|
||||||
|
import java.nio.file.Files
|
||||||
|
import sbt.internal.inc.PlainVirtualFileConverter
|
||||||
|
import sbt.util.DiskActionCacheStore
|
||||||
|
|
||||||
object GrpcActionCacheStoreTest extends verify.BasicTestSuite:
|
object GrpcActionCacheStoreTest extends verify.BasicTestSuite:
|
||||||
test("chunkBytes"):
|
test("chunkBytes"):
|
||||||
val actual = GrpcActionCacheStore.chunkBytes(0L)
|
val actual = GrpcActionCacheStore.chunkBytes(0L)
|
||||||
|
|
@ -15,4 +20,26 @@ object GrpcActionCacheStoreTest extends verify.BasicTestSuite:
|
||||||
|
|
||||||
val actual4 = GrpcActionCacheStore.chunkBytes(meg + 1)
|
val actual4 = GrpcActionCacheStore.chunkBytes(meg + 1)
|
||||||
assert(actual4 == List(meg, 1L))
|
assert(actual4 == List(meg, 1L))
|
||||||
|
|
||||||
|
// Regression test: the ByteStream deadline must be applied per RPC, not baked into the
|
||||||
|
// memoized stub. A deadline on the (session-lived) stub is absolute and expires
|
||||||
|
// remoteTimeoutInSec after first use, after which every later blob transfer fails with
|
||||||
|
// DEADLINE_EXCEEDED for the rest of the sbt server's life.
|
||||||
|
test("byteStream deadline is per-call, not on the memoized stub"):
|
||||||
|
val store = newStore()
|
||||||
|
assert(store.byteStreamStub.getCallOptions.getDeadline == null)
|
||||||
|
|
||||||
|
val deadline1 = store.byteStreamStubWithDeadline.getCallOptions.getDeadline
|
||||||
|
val deadline2 = store.byteStreamStubWithDeadline.getCallOptions.getDeadline
|
||||||
|
assert(deadline1 != null)
|
||||||
|
assert(deadline2 != null)
|
||||||
|
// Distinct Deadline instances derived at call time, not a single shared frozen one.
|
||||||
|
assert(!deadline1.eq(deadline2))
|
||||||
|
|
||||||
|
private def newStore(): GrpcActionCacheStore =
|
||||||
|
val base = Files.createTempDirectory("grpc-action-cache-test")
|
||||||
|
val disk = DiskActionCacheStore(base, PlainVirtualFileConverter.converter)
|
||||||
|
// A plaintext URI is enough to build the stubs; no connection is opened by reading
|
||||||
|
// CallOptions, so no server is required.
|
||||||
|
GrpcActionCacheStore(new URI("grpc://localhost:1"), None, None, None, Nil, disk)
|
||||||
end GrpcActionCacheStoreTest
|
end GrpcActionCacheStoreTest
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue