Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ed71b02b01 | ||
|
|
a90dc85466 | ||
|
|
67a60143cb | ||
|
|
99ac20de35 | ||
|
|
8d4f324717 |
@@ -12,6 +12,8 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/), and this
|
||||
|
||||
### Fixed
|
||||
|
||||
- **Android streams Standard Hermes attachments into its on-disk media cache.** Automatic Dashboard downloads no longer preallocate and duplicate the full response body on the OkHttp dispatcher, while declared and observed size limits remain enforced. (#531)
|
||||
- **Android fresh Gateway chats wait for their upstream runtime before sending.** Android now consumes the lazy `session.create` readiness edge instead of racing the first prompt into compute-host ownership rejection. Genuine ownership refusals keep the exact holder-aware error and retryable local prompt instead of falling into unfinished-message history recovery.
|
||||
- **Windows desktop updates keep the installed CLI and UI on one release.** `hermes-relay update` now detects a colocated management UI, reports both installed versions, and uses the verified bundle installer to replace and restart the affected surfaces together. Explicit CLI-only installations retain the standalone binary updater.
|
||||
|
||||
## [Android 1.15.0] - 2026-08-31
|
||||
|
||||
+36
-26
@@ -52,7 +52,6 @@ import okhttp3.Response
|
||||
import java.io.IOException
|
||||
import java.io.InputStream
|
||||
import java.io.OutputStream
|
||||
import java.io.ByteArrayOutputStream
|
||||
import java.net.URLEncoder
|
||||
import java.util.concurrent.TimeUnit
|
||||
import kotlin.coroutines.resume
|
||||
@@ -294,6 +293,34 @@ internal fun copyBounded(
|
||||
return written
|
||||
}
|
||||
|
||||
internal fun copyMediaBounded(
|
||||
input: InputStream,
|
||||
output: OutputStream,
|
||||
declaredLength: Long?,
|
||||
limitBytes: Long,
|
||||
): Long {
|
||||
require(limitBytes > 0)
|
||||
require(declaredLength == null || declaredLength >= 0)
|
||||
if (declaredLength != null && declaredLength > limitBytes) {
|
||||
throw IOException("File exceeds the configured download limit")
|
||||
}
|
||||
val buffer = ByteArray(DEFAULT_BUFFER_SIZE)
|
||||
var written = 0L
|
||||
while (true) {
|
||||
val read = input.read(buffer)
|
||||
if (read < 0) break
|
||||
written += read
|
||||
if (written > limitBytes) {
|
||||
throw IOException("File exceeds the configured download limit")
|
||||
}
|
||||
output.write(buffer, 0, read)
|
||||
}
|
||||
if (declaredLength != null && written != declaredLength) {
|
||||
throw IOException("Media file changed while it was being downloaded")
|
||||
}
|
||||
return written
|
||||
}
|
||||
|
||||
/** One entry from `GET /api/audio/elevenlabs/voices` — non-secret voice metadata. */
|
||||
data class ElevenLabsVoice(
|
||||
val voiceId: String,
|
||||
@@ -312,14 +339,14 @@ data class ElevenLabsVoices(
|
||||
)
|
||||
|
||||
/**
|
||||
* Bytes fetched from upstream's authenticated managed-files surface.
|
||||
* Metadata for a file streamed from upstream's authenticated managed-files surface.
|
||||
*
|
||||
* The Dashboard applies its own managed-root, sensitive-file, and maximum-size
|
||||
* policy before these bytes leave the Hermes host. Android applies the user's
|
||||
* policy before the file leaves the Hermes host. Android applies the user's
|
||||
* stricter inbound-media cap while reading the response as a second boundary.
|
||||
*/
|
||||
data class DashboardFetchedFile(
|
||||
val bytes: ByteArray,
|
||||
val sizeBytes: Long,
|
||||
val contentType: String,
|
||||
val fileName: String?,
|
||||
)
|
||||
@@ -401,13 +428,14 @@ class DashboardApiClient(
|
||||
* This is the same authenticated `/api/files/download` route official
|
||||
* Desktop uses for remote gateway files. The caller supplies its display
|
||||
* cap so a malicious or stale Content-Length cannot cause an unbounded
|
||||
* allocation. Audio/video are downloaded into Android's local media cache;
|
||||
* allocation. The supplied output is normally Android's on-disk media cache;
|
||||
* local playback supplies seeking, so `/api/files/stream` is unnecessary
|
||||
* on this path.
|
||||
*/
|
||||
suspend fun downloadManagedFile(
|
||||
serverPath: String,
|
||||
maxBytes: Long,
|
||||
output: OutputStream,
|
||||
): Result<DashboardFetchedFile> = withContext(Dispatchers.IO) {
|
||||
if (serverPath.isBlank()) {
|
||||
return@withContext Result.failure(IOException("Media path is empty"))
|
||||
@@ -428,29 +456,11 @@ class DashboardApiClient(
|
||||
if (declaredLength != null && declaredLength > maxBytes) {
|
||||
throw IOException("File exceeds the configured download limit")
|
||||
}
|
||||
val initialSize = declaredLength
|
||||
?.coerceAtMost(Int.MAX_VALUE.toLong())
|
||||
?.toInt()
|
||||
?: DEFAULT_BUFFER_SIZE
|
||||
val output = ByteArrayOutputStream(initialSize)
|
||||
body.byteStream().use { input ->
|
||||
val buffer = ByteArray(DEFAULT_BUFFER_SIZE)
|
||||
var readTotal = 0L
|
||||
while (true) {
|
||||
val count = input.read(buffer)
|
||||
if (count < 0) break
|
||||
readTotal += count
|
||||
if (readTotal > maxBytes) {
|
||||
throw IOException("File exceeds the configured download limit")
|
||||
}
|
||||
output.write(buffer, 0, count)
|
||||
}
|
||||
if (declaredLength != null && readTotal != declaredLength) {
|
||||
throw IOException("Media file changed while it was being downloaded")
|
||||
}
|
||||
val readTotal = body.byteStream().use { input ->
|
||||
copyMediaBounded(input, output, declaredLength, maxBytes)
|
||||
}
|
||||
DashboardFetchedFile(
|
||||
bytes = output.toByteArray(),
|
||||
sizeBytes = readTotal,
|
||||
contentType = response.header("Content-Type")
|
||||
?.substringBefore(';')
|
||||
?.trim()
|
||||
|
||||
@@ -119,6 +119,8 @@ class GatewayChatClient(
|
||||
private val promptSubmitTimeoutMs: Long = PROMPT_SUBMIT_REQUEST_TIMEOUT_MS,
|
||||
/** Test seam — idle-progress watchdog base. Production keeps [TURN_TIMEOUT_MS]. */
|
||||
private val turnIdleTimeoutMs: Long = TURN_TIMEOUT_MS,
|
||||
/** Test seam — lazy session.create/resume readiness barrier. */
|
||||
private val sessionReadyTimeoutMs: Long = SESSION_READY_TIMEOUT_MS,
|
||||
/** Test seam — compaction idle lease. Production keeps [COMPACTING_TIMEOUT_MS]. */
|
||||
private val compactingTimeoutMs: Long = COMPACTING_TIMEOUT_MS,
|
||||
/** Random source for ordinary reconnect full-jitter. */
|
||||
@@ -191,6 +193,7 @@ class GatewayChatClient(
|
||||
* would have been abandoned server-side anyway.
|
||||
*/
|
||||
private const val PROMPT_SUBMIT_REQUEST_TIMEOUT_MS = 1_800_000L
|
||||
private const val SESSION_READY_TIMEOUT_MS = 300_000L
|
||||
private const val CONNECT_TIMEOUT_MS = 20_000L
|
||||
|
||||
/**
|
||||
@@ -421,6 +424,9 @@ class GatewayChatClient(
|
||||
/** Invalidates an older async prewarm when a newer session selection wins. */
|
||||
private val prewarmRequestGeneration = AtomicLong(0)
|
||||
private val pendingRpcs = ConcurrentHashMap<Long, CompletableDeferred<JsonObject>>()
|
||||
private val lazyLiveSessions = ConcurrentHashMap.newKeySet<String>()
|
||||
private val readyLiveSessions = ConcurrentHashMap.newKeySet<String>()
|
||||
private val sessionReadyWaiters = ConcurrentHashMap<String, CompletableDeferred<Unit>>()
|
||||
|
||||
/** Monotonic client-local fence for lazy child watch open/close races. */
|
||||
private val childWatchGeneration = AtomicLong(0)
|
||||
@@ -832,7 +838,7 @@ class GatewayChatClient(
|
||||
// through the normal failed-turn callback instead.
|
||||
cleanupStagedAttachments(stagedImagePaths)
|
||||
turn.tracer.done("submit-rejected")
|
||||
turn.callbacks.onError(
|
||||
turn.callbacks.onSubmitRejected(
|
||||
submitError?.message ?: "Hermes rejected the new session",
|
||||
)
|
||||
return@launch
|
||||
@@ -899,6 +905,9 @@ class GatewayChatClient(
|
||||
liveSessionId = null
|
||||
storedSessionId = null
|
||||
liveSessionProfile = null
|
||||
failSessionReadyWaiters("gateway session cleared")
|
||||
lazyLiveSessions.clear()
|
||||
readyLiveSessions.clear()
|
||||
cancelledTurnDrain = null
|
||||
_serverTools.value = null
|
||||
}
|
||||
@@ -3389,6 +3398,7 @@ class GatewayChatClient(
|
||||
requestedStoredId != null &&
|
||||
liveSessionProfile == requestedProfile
|
||||
) {
|
||||
awaitSessionReadyIfRequired(liveSessionId!!)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -3418,6 +3428,7 @@ class GatewayChatClient(
|
||||
liveSessionProfile = requestedProfile
|
||||
updateCancelledDrainLiveSession(requestedStoredId, live)
|
||||
applySessionResultInfo(result)
|
||||
awaitSessionReadyIfRequired(live, result)
|
||||
return
|
||||
}
|
||||
throw GatewayAuthoritativeResumeException(
|
||||
@@ -3453,6 +3464,45 @@ class GatewayChatClient(
|
||||
if (cancelledTurnDrain?.storedSessionId != stored) cancelledTurnDrain = null
|
||||
if (settledTurnDrain?.storedSessionId != stored) settledTurnDrain = null
|
||||
turn.callbacks.onSessionId(stored)
|
||||
awaitSessionReadyIfRequired(live, created)
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for current upstream's authoritative deferred-build edge instead of
|
||||
* racing a lazy session into compute-host turn isolation. Without this
|
||||
* barrier, the parent runtime can claim the durable lease before the child
|
||||
* rechecks it, causing the child to reject itself as a second live owner.
|
||||
*/
|
||||
private suspend fun awaitSessionReadyIfRequired(
|
||||
liveId: String,
|
||||
sessionResult: JsonObject? = null,
|
||||
) {
|
||||
val lazy = (sessionResult?.get("info") as? JsonObject)?.booleanField("lazy") == true
|
||||
if (lazy) lazyLiveSessions += liveId
|
||||
if (liveId !in lazyLiveSessions) return
|
||||
if (readyLiveSessions.remove(liveId)) {
|
||||
lazyLiveSessions.remove(liveId)
|
||||
return
|
||||
}
|
||||
val waiter = sessionReadyWaiters.computeIfAbsent(liveId) { CompletableDeferred() }
|
||||
if (readyLiveSessions.remove(liveId)) waiter.complete(Unit)
|
||||
try {
|
||||
withTimeout(sessionReadyTimeoutMs) { waiter.await() }
|
||||
lazyLiveSessions.remove(liveId)
|
||||
} catch (error: Exception) {
|
||||
throw GatewayPreflightException(
|
||||
error.message ?: "Hermes session initialization timed out",
|
||||
)
|
||||
} finally {
|
||||
sessionReadyWaiters.remove(liveId, waiter)
|
||||
}
|
||||
}
|
||||
|
||||
private fun failSessionReadyWaiters(message: String) {
|
||||
sessionReadyWaiters.values.forEach {
|
||||
it.completeExceptionally(GatewayRpcException(message))
|
||||
}
|
||||
sessionReadyWaiters.clear()
|
||||
}
|
||||
|
||||
private fun createListener(ready: CompletableDeferred<Unit>) = object : WebSocketListener() {
|
||||
@@ -3631,6 +3681,17 @@ class GatewayChatClient(
|
||||
return
|
||||
}
|
||||
|
||||
// Lazy create/resume returns before its AIAgent exists. The deferred
|
||||
// build publishes exact-session session.info when it is ready. Record
|
||||
// that edge before live-session routing so a fast build cannot race
|
||||
// the RPC response and strand the first submit.
|
||||
if (type == "session.info" && !eventSessionId.isNullOrBlank() &&
|
||||
payload?.booleanField("lazy") != true
|
||||
) {
|
||||
readyLiveSessions += eventSessionId
|
||||
sessionReadyWaiters.remove(eventSessionId)?.complete(Unit)
|
||||
}
|
||||
|
||||
// Upstream emits session.reclaimed process-wide, so it is identified
|
||||
// by payload rather than params.session_id. Retire only an exact live
|
||||
// runtime we own; preserve the durable id so the next send resumes it.
|
||||
@@ -3958,6 +4019,9 @@ class GatewayChatClient(
|
||||
it.completeExceptionally(GatewayRpcException("gateway connection lost"))
|
||||
}
|
||||
pendingRpcs.clear()
|
||||
failSessionReadyWaiters("gateway connection lost")
|
||||
lazyLiveSessions.clear()
|
||||
readyLiveSessions.clear()
|
||||
failChildWatches("Child watch disconnected from the gateway")
|
||||
val turn = activeTurn
|
||||
if (turn == null) {
|
||||
@@ -4135,6 +4199,9 @@ class GatewayChatClient(
|
||||
webSocket = null
|
||||
readySignal = null
|
||||
liveSessionId = null
|
||||
failSessionReadyWaiters("gateway socket closed")
|
||||
lazyLiveSessions.clear()
|
||||
readyLiveSessions.clear()
|
||||
attachMethodForSocket = null
|
||||
commandsCatalogCache = null
|
||||
_processCapability.value = GatewayProcessCapability.Unknown
|
||||
@@ -4811,6 +4878,9 @@ class GatewayChatClient(
|
||||
onStatusClear = { kind -> dispatchIfCurrent(stillCurrent) { callbacks.onStatusClear(kind) } },
|
||||
onNoticeShow = { notice -> dispatchIfCurrent(stillCurrent) { callbacks.onNoticeShow(notice) } },
|
||||
onNoticeClear = { key -> dispatchIfCurrent(stillCurrent) { callbacks.onNoticeClear(key) } },
|
||||
onSubmitRejected = { message ->
|
||||
dispatchIfCurrent(stillCurrent) { callbacks.onSubmitRejected(message) }
|
||||
},
|
||||
)
|
||||
|
||||
private fun dispatchIfCurrent(stillCurrent: () -> Boolean, callback: () -> Unit) {
|
||||
|
||||
@@ -279,7 +279,12 @@ class GatewayEventMapper(
|
||||
|
||||
"error" -> {
|
||||
turnEnded = true
|
||||
callbacks.onError(payload.string("message") ?: "Gateway error")
|
||||
val message = payload.string("message") ?: "Gateway error"
|
||||
if (isSessionOwnershipRejection(message)) {
|
||||
callbacks.onSubmitRejected(message)
|
||||
} else {
|
||||
callbacks.onError(message)
|
||||
}
|
||||
}
|
||||
|
||||
"subagent.spawn_requested", "subagent.start", "subagent.thinking", "subagent.tool",
|
||||
@@ -469,6 +474,19 @@ class GatewayEventMapper(
|
||||
private const val MAX_MOA_REFERENCE_CHARS = 16_000
|
||||
private val OUTPUT_RISK_LEVELS = setOf("low", "medium", "high", "critical")
|
||||
private val TERMINAL_EVENTS = setOf("message.complete", "error")
|
||||
|
||||
/**
|
||||
* Isolated-turn paths can acknowledge `prompt.submit` and then emit
|
||||
* the ownership refusal as a plain terminal `error` event. That event
|
||||
* has no JSON-RPC code or structured reason, so match both stable
|
||||
* clauses from upstream's canonical message.
|
||||
*/
|
||||
internal fun isSessionOwnershipRejection(message: String): Boolean {
|
||||
val normalized = message.lowercase()
|
||||
return "already has a live owner (" in normalized &&
|
||||
"only one surface at a time may run a session" in normalized
|
||||
}
|
||||
|
||||
internal fun isFailedMoaReference(text: String): Boolean {
|
||||
val normalized = text.trimStart().lowercase()
|
||||
return normalized.startsWith("[failed:") || normalized.startsWith("[skipped:")
|
||||
|
||||
@@ -741,6 +741,12 @@ class GatewayTurnCallbacks(
|
||||
val onNoticeShow: (GatewayAgentNotice) -> Unit = { _ -> },
|
||||
/** Exact-key dismissal for an upstream notice. */
|
||||
val onNoticeClear: (key: String) -> Unit = { _ -> },
|
||||
/**
|
||||
* The Gateway authoritatively refused `prompt.submit` before a model turn
|
||||
* began. The composed user row remains local and retryable; callers must
|
||||
* not treat this as a dropped stream or reconcile it from history.
|
||||
*/
|
||||
val onSubmitRejected: (String) -> Unit = onError,
|
||||
)
|
||||
|
||||
/**
|
||||
|
||||
@@ -2,6 +2,7 @@ package com.hermesandroid.relay.ui.components
|
||||
|
||||
import android.content.Intent
|
||||
import android.graphics.BitmapFactory
|
||||
import android.net.Uri
|
||||
import androidx.compose.foundation.background
|
||||
import androidx.compose.foundation.clickable
|
||||
import androidx.compose.foundation.layout.Arrangement
|
||||
@@ -97,9 +98,9 @@ private fun isSensitiveAltText(alt: String): Boolean {
|
||||
|
||||
/**
|
||||
* Resolves a server-local image path — an absolute path the agent put in a
|
||||
* markdown image `` — to raw bytes, via the relay's
|
||||
* `/media/by-path` route when a relay session is paired. Returns null when no
|
||||
* relay is available or the fetch fails, so the renderer falls back to the
|
||||
* markdown image `` — to an on-disk cache URI, via the relay's
|
||||
* authenticated Dashboard route (or Relay compatibility route). Returns a
|
||||
* failure when no route is available or the fetch fails, so the renderer falls back to the
|
||||
* "this image is on the server" notice. Provided by ChatScreen from
|
||||
* [com.hermesandroid.relay.viewmodel.ChatViewModel.resolveServerImage]; the
|
||||
* default is null, which preserves the standard (no-plugin) behavior where a
|
||||
@@ -113,12 +114,12 @@ private fun isSensitiveAltText(alt: String): Boolean {
|
||||
* `RelayServerImage` on app open with a server-local image in history).
|
||||
*/
|
||||
sealed interface ServerImageResult {
|
||||
class Success(val bytes: ByteArray, val sensitive: Boolean = false) : ServerImageResult
|
||||
class Success(val cachedUri: String, val sensitive: Boolean = false) : ServerImageResult
|
||||
class Failure(val reason: String) : ServerImageResult
|
||||
}
|
||||
|
||||
fun interface RelayServerImageResolver {
|
||||
/** Fetch the server-local file's bytes over the relay, or a
|
||||
/** Fetch the server-local file into the media cache, or a
|
||||
* [ServerImageResult.Failure] whose reason explains why (unpaired /
|
||||
* sandboxed / missing / decode) so the UI can surface it instead of a
|
||||
* generic placeholder. */
|
||||
@@ -461,6 +462,7 @@ private fun RelayServerImage(
|
||||
maxWidth: Dp,
|
||||
resolver: RelayServerImageResolver,
|
||||
) {
|
||||
val context = LocalContext.current
|
||||
var phase by remember(image.src) {
|
||||
mutableStateOf<RelayImagePhase>(
|
||||
cachedInlineImage(image.src)
|
||||
@@ -478,15 +480,15 @@ private fun RelayServerImage(
|
||||
}
|
||||
when (result) {
|
||||
is ServerImageResult.Success -> {
|
||||
val bytes = result.bytes
|
||||
val cachedUri = Uri.parse(result.cachedUri)
|
||||
val bmp = runCatching {
|
||||
decodeOrientedBitmap(bytes)
|
||||
decodeOrientedBitmap(context, cachedUri)
|
||||
}.getOrNull()?.asImageBitmap()
|
||||
if (bmp != null) {
|
||||
putInlineImage(image.src, bmp, result.sensitive)
|
||||
RelayImagePhase.Loaded(bmp, result.sensitive)
|
||||
} else {
|
||||
RelayImagePhase.Failed("fetched ${bytes.size} B but couldn't decode the image")
|
||||
RelayImagePhase.Failed("fetched file could not be decoded as an image")
|
||||
}
|
||||
}
|
||||
is ServerImageResult.Failure -> RelayImagePhase.Failed(result.reason)
|
||||
@@ -527,6 +529,7 @@ private fun RelayServerImageContent(
|
||||
maxWidth: Dp,
|
||||
resolver: RelayServerImageResolver,
|
||||
) {
|
||||
val context = LocalContext.current
|
||||
var viewerOpen by remember { mutableStateOf(false) }
|
||||
val blurMode = LocalMediaBlurMode.current
|
||||
var revealed by remember(image.src) { mutableStateOf(false) }
|
||||
@@ -542,7 +545,12 @@ private fun RelayServerImageContent(
|
||||
mime = "image/*",
|
||||
// Save/Share re-fetch the original bytes on demand so we don't
|
||||
// hold them in memory next to the decoded bitmap.
|
||||
bytesProvider = { (resolver.fetch(image.src) as? ServerImageResult.Success)?.bytes },
|
||||
bytesProvider = {
|
||||
(resolver.fetch(image.src) as? ServerImageResult.Success)?.let { fetched ->
|
||||
context.contentResolver.openInputStream(Uri.parse(fetched.cachedUri))
|
||||
?.use { it.readBytes() }
|
||||
}
|
||||
},
|
||||
),
|
||||
onDismiss = { viewerOpen = false },
|
||||
sensitive = sensitive,
|
||||
|
||||
@@ -303,15 +303,17 @@ private fun ImageRender(
|
||||
LaunchedEffect(attachment.cachedUri, attachment.content) {
|
||||
val decoded = withContext(Dispatchers.IO) {
|
||||
runCatching {
|
||||
val bytes: ByteArray? = when {
|
||||
when {
|
||||
!attachment.cachedUri.isNullOrBlank() ->
|
||||
context.contentResolver.openInputStream(Uri.parse(attachment.cachedUri))
|
||||
?.use { it.readBytes() }
|
||||
decodeOrientedBitmap(context, Uri.parse(attachment.cachedUri))
|
||||
?.asImageBitmap()
|
||||
attachment.content.isNotBlank() ->
|
||||
android.util.Base64.decode(attachment.content, android.util.Base64.DEFAULT)
|
||||
android.util.Base64.decode(
|
||||
attachment.content,
|
||||
android.util.Base64.DEFAULT,
|
||||
).let { decodeOrientedBitmap(it)?.asImageBitmap() }
|
||||
else -> null
|
||||
}
|
||||
bytes?.let { decodeOrientedBitmap(it)?.asImageBitmap() }
|
||||
}.getOrNull()
|
||||
}
|
||||
if (decoded != null) bitmap = decoded else decodeFailed = true
|
||||
|
||||
@@ -1,8 +1,10 @@
|
||||
package com.hermesandroid.relay.ui.components
|
||||
|
||||
import android.content.Context
|
||||
import android.graphics.Bitmap
|
||||
import android.graphics.BitmapFactory
|
||||
import android.graphics.Matrix
|
||||
import android.net.Uri
|
||||
import androidx.exifinterface.media.ExifInterface
|
||||
import java.io.ByteArrayInputStream
|
||||
|
||||
@@ -62,3 +64,30 @@ internal fun decodeOrientedBitmap(
|
||||
if (oriented !== decoded) decoded.recycle()
|
||||
} ?: decoded
|
||||
}
|
||||
|
||||
/** Decode a cached content URI without first copying the encoded file into RAM. */
|
||||
internal fun decodeOrientedBitmap(context: Context, uri: Uri): Bitmap? {
|
||||
val decoded = context.contentResolver.openFileDescriptor(uri, "r")?.use { descriptor ->
|
||||
BitmapFactory.decodeFileDescriptor(descriptor.fileDescriptor)
|
||||
} ?: return null
|
||||
val orientation = runCatching {
|
||||
context.contentResolver.openFileDescriptor(uri, "r")?.use { descriptor ->
|
||||
ExifInterface(descriptor.fileDescriptor).getAttributeInt(
|
||||
ExifInterface.TAG_ORIENTATION,
|
||||
ExifInterface.ORIENTATION_NORMAL,
|
||||
)
|
||||
}
|
||||
}.getOrNull() ?: ExifInterface.ORIENTATION_NORMAL
|
||||
val transform = imageOrientationTransform(orientation)
|
||||
if (transform.rotationDegrees == 0f && !transform.flipHorizontal) return decoded
|
||||
|
||||
val matrix = Matrix().apply {
|
||||
if (transform.rotationDegrees != 0f) postRotate(transform.rotationDegrees)
|
||||
if (transform.flipHorizontal) postScale(-1f, 1f)
|
||||
}
|
||||
return runCatching {
|
||||
Bitmap.createBitmap(decoded, 0, 0, decoded.width, decoded.height, matrix, true)
|
||||
}.getOrNull()?.also { oriented ->
|
||||
if (oriented !== decoded) decoded.recycle()
|
||||
} ?: decoded
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import androidx.core.content.FileProvider
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import java.io.File
|
||||
import java.io.InputStream
|
||||
import java.security.MessageDigest
|
||||
|
||||
/**
|
||||
@@ -98,7 +99,7 @@ class MediaCacheWriter(
|
||||
suspend fun cache(bytes: ByteArray, contentType: String, fileName: String?): Uri =
|
||||
withContext(Dispatchers.IO) {
|
||||
val dir = cacheRoot
|
||||
val targetName = buildFilename(bytes, contentType, fileName)
|
||||
val targetName = buildFilename(bytes.sha1(), contentType, fileName)
|
||||
val file = File(dir, targetName)
|
||||
file.writeBytes(bytes)
|
||||
// Best-effort LRU enforcement — never fatal on a cache operation.
|
||||
@@ -113,6 +114,45 @@ class MediaCacheWriter(
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Move a completed streaming download into the media cache without ever
|
||||
* materializing the file as a ByteArray. The caller retains ownership of
|
||||
* [source] when this throws; successful calls consume or copy it.
|
||||
*/
|
||||
suspend fun cache(source: File, contentType: String, fileName: String?): Uri {
|
||||
val target = cacheFile(source, contentType, fileName)
|
||||
return withContext(Dispatchers.IO) {
|
||||
FileProvider.getUriForFile(
|
||||
context,
|
||||
"${context.packageName}.fileprovider",
|
||||
target,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
internal suspend fun cacheFile(source: File, contentType: String, fileName: String?): File =
|
||||
withContext(Dispatchers.IO) {
|
||||
require(source.isFile) { "Media cache source does not exist" }
|
||||
val dir = cacheRoot
|
||||
val target = File(dir, buildFilename(source.sha1(), contentType, fileName))
|
||||
if (source.canonicalFile != target.canonicalFile) {
|
||||
if (target.exists()) {
|
||||
source.delete()
|
||||
target.setLastModified(System.currentTimeMillis())
|
||||
} else if (!source.renameTo(target)) {
|
||||
source.inputStream().use { input ->
|
||||
target.outputStream().use { output -> input.copyTo(output) }
|
||||
}
|
||||
source.delete()
|
||||
}
|
||||
}
|
||||
try {
|
||||
enforceCap()
|
||||
} catch (_: Exception) { /* swallow */ }
|
||||
|
||||
target
|
||||
}
|
||||
|
||||
/**
|
||||
* Wipe the entire `hermes-media/` directory. Returns the number of bytes
|
||||
* freed so the caller can show a "Freed X MB" snackbar/toast.
|
||||
@@ -133,15 +173,8 @@ class MediaCacheWriter(
|
||||
|
||||
// --- internals ---
|
||||
|
||||
private fun buildFilename(bytes: ByteArray, contentType: String, fileName: String?): String {
|
||||
private fun buildFilename(sha1: String, contentType: String, fileName: String?): String {
|
||||
val ext = mimeToExt[contentType.lowercase().substringBefore(';').trim()] ?: "bin"
|
||||
// SHA-1 of the bytes — stable + collision-resistant enough for cache keys.
|
||||
val sha1 = try {
|
||||
MessageDigest.getInstance("SHA-1").digest(bytes).joinToString("") { "%02x".format(it) }
|
||||
} catch (_: Exception) {
|
||||
// fallback to a weaker but always-available hash
|
||||
bytes.contentHashCode().toUInt().toString(16)
|
||||
}
|
||||
|
||||
// If the server supplied a filename, keep the base (sanitized) but
|
||||
// force the extension we derived from the MIME. This keeps the
|
||||
@@ -162,6 +195,33 @@ class MediaCacheWriter(
|
||||
}
|
||||
}
|
||||
|
||||
private fun ByteArray.sha1(): String = try {
|
||||
MessageDigest.getInstance("SHA-1").digest(this).toHexString()
|
||||
} catch (_: Exception) {
|
||||
contentHashCode().toUInt().toString(16)
|
||||
}
|
||||
|
||||
private fun File.sha1(): String = try {
|
||||
inputStream().use { input ->
|
||||
val digest = MessageDigest.getInstance("SHA-1")
|
||||
input.updateDigest(digest)
|
||||
digest.digest().toHexString()
|
||||
}
|
||||
} catch (_: Exception) {
|
||||
"${length().toUInt().toString(16)}-${lastModified().toUInt().toString(16)}"
|
||||
}
|
||||
|
||||
private fun InputStream.updateDigest(digest: MessageDigest) {
|
||||
val buffer = ByteArray(DEFAULT_BUFFER_SIZE)
|
||||
while (true) {
|
||||
val read = read(buffer)
|
||||
if (read < 0) return
|
||||
digest.update(buffer, 0, read)
|
||||
}
|
||||
}
|
||||
|
||||
private fun ByteArray.toHexString(): String = joinToString("") { "%02x".format(it) }
|
||||
|
||||
/**
|
||||
* Delete oldest-mtime files until the cache directory fits under the cap.
|
||||
* Called after every write; safe to call more often.
|
||||
|
||||
@@ -173,6 +173,8 @@ import kotlinx.serialization.json.contentOrNull
|
||||
import kotlinx.serialization.json.intOrNull
|
||||
import kotlinx.serialization.json.longOrNull
|
||||
import okhttp3.sse.EventSource
|
||||
import java.io.File
|
||||
import java.io.IOException
|
||||
import java.io.InterruptedIOException
|
||||
import java.util.UUID
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
@@ -10531,6 +10533,10 @@ class ChatViewModel : ViewModel() {
|
||||
_steerableTurn.value = false
|
||||
onPreflightErrorCb(Exception(reason))
|
||||
},
|
||||
onSubmitRejected = { reason ->
|
||||
_steerableTurn.value = false
|
||||
onPreflightErrorCb(Exception(reason))
|
||||
},
|
||||
onFailure = { failure ->
|
||||
val confirmedModel = gateway.serverModel.value
|
||||
?.takeIf { it.isNotBlank() }
|
||||
@@ -10723,7 +10729,7 @@ class ChatViewModel : ViewModel() {
|
||||
* placeholder as actionable FAILED with [MEDIA_TAP_TO_DOWNLOAD] in
|
||||
* [Attachment.errorMessage]. The UI renders a retry card.
|
||||
* 4. Otherwise → kick off a background fetch, enforce the max size cap,
|
||||
* cache the bytes via [MediaCacheWriter], and update the attachment
|
||||
* stream into [MediaCacheWriter], and update the attachment
|
||||
* to LOADED (or FAILED on any error).
|
||||
*
|
||||
* Note: the matching happens by (messageId + relayToken). If the same
|
||||
@@ -10787,7 +10793,10 @@ class ChatViewModel : ViewModel() {
|
||||
expectedHistoryGeneration = historyGeneration,
|
||||
) { _ ->
|
||||
if (relay.mediaUrlConfigured()) {
|
||||
relay.fetchMedia(token, maxBytes = settings.maxInboundSizeMb.toLong().coerceAtLeast(1L) * 1024L * 1024L)
|
||||
relay.fetchMedia(
|
||||
token,
|
||||
maxBytes = settings.maxInboundSizeMb.toLong().coerceAtLeast(1L) * 1024L * 1024L,
|
||||
).asInboundFetchedMedia()
|
||||
} else {
|
||||
Result.failure(MediaRouteUnavailableException())
|
||||
}
|
||||
@@ -10850,7 +10859,7 @@ class ChatViewModel : ViewModel() {
|
||||
fetchServerPath(fetchKey, maxBytes, upstreamMediaClient)
|
||||
} else {
|
||||
if (relay.mediaUrlConfigured()) {
|
||||
relay.fetchMedia(fetchKey, maxBytes = maxBytes)
|
||||
relay.fetchMedia(fetchKey, maxBytes = maxBytes).asInboundFetchedMedia()
|
||||
} else {
|
||||
Result.failure(MediaRouteUnavailableException())
|
||||
}
|
||||
@@ -10909,7 +10918,22 @@ class ChatViewModel : ViewModel() {
|
||||
return ServerImageResult.Failure("Connection changed while loading image")
|
||||
}
|
||||
return result.fold(
|
||||
onSuccess = { ServerImageResult.Success(it.bytes, it.sensitive) },
|
||||
onSuccess = { fetched ->
|
||||
try {
|
||||
val cachedUri = fetched.cachedUri ?: mediaCacheWriter?.cache(
|
||||
requireNotNull(fetched.bytes),
|
||||
fetched.contentType,
|
||||
fetched.fileName,
|
||||
)?.toString()
|
||||
if (cachedUri == null) {
|
||||
ServerImageResult.Failure("Media cache is unavailable")
|
||||
} else {
|
||||
ServerImageResult.Success(cachedUri, fetched.sensitive)
|
||||
}
|
||||
} catch (error: Exception) {
|
||||
ServerImageResult.Failure(error.message ?: "Image cache failed")
|
||||
}
|
||||
},
|
||||
onFailure = { ServerImageResult.Failure(it.message ?: "Image unavailable") },
|
||||
)
|
||||
}
|
||||
@@ -11031,17 +11055,39 @@ class ChatViewModel : ViewModel() {
|
||||
serverPath: String,
|
||||
maxBytes: Long,
|
||||
upstream: DashboardApiClient?,
|
||||
): Result<RelayHttpClient.FetchedMedia> {
|
||||
): Result<InboundFetchedMedia> {
|
||||
if (upstream != null) {
|
||||
val upstreamResult = upstream.downloadManagedFile(serverPath, maxBytes)
|
||||
.map { fetched ->
|
||||
RelayHttpClient.FetchedMedia(
|
||||
contentType = fetched.contentType,
|
||||
bytes = fetched.bytes,
|
||||
fileName = fetched.fileName,
|
||||
sensitive = false,
|
||||
val cache = mediaCacheWriter
|
||||
?: return Result.failure(IOException("Media cache is unavailable"))
|
||||
val context = appContext
|
||||
?: return Result.failure(IOException("Application context is unavailable"))
|
||||
val staging = File.createTempFile("hermes-media-", ".part", context.cacheDir)
|
||||
val upstreamResult = try {
|
||||
val fetchedResult = staging.outputStream().buffered().use { output ->
|
||||
upstream.downloadManagedFile(serverPath, maxBytes, output)
|
||||
}
|
||||
if (fetchedResult.isFailure) {
|
||||
Result.failure(fetchedResult.exceptionOrNull()!!)
|
||||
} else {
|
||||
val fetched = fetchedResult.getOrThrow()
|
||||
val uri = cache.cache(staging, fetched.contentType, fetched.fileName)
|
||||
Result.success(
|
||||
InboundFetchedMedia(
|
||||
contentType = fetched.contentType,
|
||||
fileName = fetched.fileName,
|
||||
sensitive = false,
|
||||
sizeBytes = fetched.sizeBytes,
|
||||
cachedUri = uri.toString(),
|
||||
),
|
||||
)
|
||||
}
|
||||
} catch (error: CancellationException) {
|
||||
throw error
|
||||
} catch (error: Exception) {
|
||||
Result.failure(error)
|
||||
} finally {
|
||||
staging.delete()
|
||||
}
|
||||
if (upstreamResult.isSuccess) return upstreamResult
|
||||
val upstreamFailure = upstreamResult.exceptionOrNull()
|
||||
if (upstreamFailure?.isDashboardManagedFilesUnsupported() != true) {
|
||||
@@ -11050,13 +11096,14 @@ class ChatViewModel : ViewModel() {
|
||||
val relay = relayHttpClient
|
||||
if (relay?.mediaUrlConfigured() == true) {
|
||||
return relay.fetchMediaByPath(serverPath, maxBytes = maxBytes)
|
||||
.asInboundFetchedMedia()
|
||||
}
|
||||
return Result.failure(MediaRouteUnavailableException(upstreamFailure))
|
||||
}
|
||||
|
||||
val relay = relayHttpClient
|
||||
return if (relay?.mediaUrlConfigured() == true) {
|
||||
relay.fetchMediaByPath(serverPath, maxBytes = maxBytes)
|
||||
relay.fetchMediaByPath(serverPath, maxBytes = maxBytes).asInboundFetchedMedia()
|
||||
} else {
|
||||
Result.failure(MediaRouteUnavailableException())
|
||||
}
|
||||
@@ -11096,7 +11143,7 @@ class ChatViewModel : ViewModel() {
|
||||
expectedRole: MessageRole = MessageRole.ASSISTANT,
|
||||
expectedContextKey: String? = activeProfileContextKey,
|
||||
expectedHistoryGeneration: Int = historyLoadGeneration.get(),
|
||||
fetch: suspend (maxBytes: Long) -> Result<RelayHttpClient.FetchedMedia>,
|
||||
fetch: suspend (maxBytes: Long) -> Result<InboundFetchedMedia>,
|
||||
) {
|
||||
val cache = mediaCacheWriter ?: return
|
||||
val maxBytes = settings.maxInboundSizeMb.toLong().coerceAtLeast(1) * 1024L * 1024L
|
||||
@@ -11117,8 +11164,8 @@ class ChatViewModel : ViewModel() {
|
||||
) return
|
||||
result.fold(
|
||||
onSuccess = { fetched ->
|
||||
if (fetched.bytes.size > maxBytes) {
|
||||
val sizeMb = fetched.bytes.size / (1024.0 * 1024.0)
|
||||
if (fetched.sizeBytes > maxBytes) {
|
||||
val sizeMb = fetched.sizeBytes / (1024.0 * 1024.0)
|
||||
updateAttachmentByToken(handler, messageId, fetchKey, expectedRole = expectedRole) { att ->
|
||||
att.copy(
|
||||
state = AttachmentState.FAILED,
|
||||
@@ -11127,22 +11174,26 @@ class ChatViewModel : ViewModel() {
|
||||
),
|
||||
contentType = fetched.contentType,
|
||||
fileName = fetched.fileName ?: att.fileName,
|
||||
fileSize = fetched.bytes.size.toLong()
|
||||
fileSize = fetched.sizeBytes
|
||||
)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
val uri = cache.cache(fetched.bytes, fetched.contentType, fetched.fileName)
|
||||
val cachedUri = fetched.cachedUri ?: cache.cache(
|
||||
requireNotNull(fetched.bytes),
|
||||
fetched.contentType,
|
||||
fetched.fileName,
|
||||
).toString()
|
||||
updateAttachmentByToken(handler, messageId, fetchKey, expectedRole = expectedRole) { att ->
|
||||
att.copy(
|
||||
state = AttachmentState.LOADED,
|
||||
errorMessage = null,
|
||||
contentType = fetched.contentType,
|
||||
fileName = fetched.fileName ?: att.fileName,
|
||||
fileSize = fetched.bytes.size.toLong(),
|
||||
cachedUri = uri.toString(),
|
||||
fileSize = fetched.sizeBytes,
|
||||
cachedUri = cachedUri,
|
||||
// Carry the relay-authoritative sensitivity bit
|
||||
// (X-Media-Sensitive header) onto the LOADED
|
||||
// attachment so the renderer can blur per the
|
||||
@@ -11188,6 +11239,26 @@ class ChatViewModel : ViewModel() {
|
||||
)
|
||||
}
|
||||
|
||||
private data class InboundFetchedMedia(
|
||||
val contentType: String,
|
||||
val fileName: String?,
|
||||
val sensitive: Boolean,
|
||||
val sizeBytes: Long,
|
||||
val bytes: ByteArray? = null,
|
||||
val cachedUri: String? = null,
|
||||
)
|
||||
|
||||
private fun Result<RelayHttpClient.FetchedMedia>.asInboundFetchedMedia(): Result<InboundFetchedMedia> =
|
||||
map { fetched ->
|
||||
InboundFetchedMedia(
|
||||
contentType = fetched.contentType,
|
||||
fileName = fetched.fileName,
|
||||
sensitive = fetched.sensitive,
|
||||
sizeBytes = fetched.bytes.size.toLong(),
|
||||
bytes = fetched.bytes,
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Append an attachment to a specific assistant message, matched by id.
|
||||
* No-ops if the message can't be found (e.g. it was trimmed from the
|
||||
|
||||
+22
-3
@@ -2042,13 +2042,15 @@ class DashboardApiClientTest {
|
||||
.setBody("audio-bytes"),
|
||||
)
|
||||
|
||||
val output = ByteArrayOutputStream()
|
||||
val fetched = DashboardApiClient(baseUrl = server.url("/").toString())
|
||||
.downloadManagedFile("/tmp/Hermes audio/voice reply.mp3", 1024)
|
||||
.downloadManagedFile("/tmp/Hermes audio/voice reply.mp3", 1024, output)
|
||||
.getOrThrow()
|
||||
|
||||
assertEquals("audio/mpeg", fetched.contentType)
|
||||
assertEquals("voice reply.mp3", fetched.fileName)
|
||||
assertEquals("audio-bytes", fetched.bytes.decodeToString())
|
||||
assertEquals(11L, fetched.sizeBytes)
|
||||
assertEquals("audio-bytes", output.toString(Charsets.UTF_8.name()))
|
||||
val request = server.takeRequest()
|
||||
assertEquals("/api/files/download", request.requestUrl!!.encodedPath)
|
||||
assertEquals("/tmp/Hermes audio/voice reply.mp3", request.requestUrl!!.queryParameter("path"))
|
||||
@@ -2059,11 +2061,28 @@ class DashboardApiClientTest {
|
||||
server.enqueue(MockResponse().setHeader("Content-Length", 2048))
|
||||
|
||||
val result = DashboardApiClient(baseUrl = server.url("/").toString())
|
||||
.downloadManagedFile("/tmp/large.bin", 1024)
|
||||
.downloadManagedFile("/tmp/large.bin", 1024, ByteArrayOutputStream())
|
||||
|
||||
assertTrue(result.isFailure)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun downloadManagedFile_rejectsChunkedBodyAfterStreamingCap() = runTest {
|
||||
server.enqueue(
|
||||
MockResponse()
|
||||
.setHeader("Content-Type", "application/octet-stream")
|
||||
.setChunkedBody("seventeen-byte-doc", 3),
|
||||
)
|
||||
val output = ByteArrayOutputStream()
|
||||
|
||||
val result = DashboardApiClient(baseUrl = server.url("/").toString())
|
||||
.downloadManagedFile("/tmp/growing.bin", 16, output)
|
||||
|
||||
assertTrue(result.isFailure)
|
||||
assertTrue(result.exceptionOrNull()?.message.orEmpty().contains("configured download limit"))
|
||||
assertTrue(output.size() <= 16)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun managedFileFallbackClassificationDistinguishesMissingRouteFromMissingFile() {
|
||||
assertTrue(
|
||||
|
||||
+49
-2
@@ -133,6 +133,9 @@ class GatewayClientHarness(
|
||||
@Volatile
|
||||
var createdSessionProfileName: String? = null
|
||||
|
||||
@Volatile
|
||||
var createdSessionLazy = false
|
||||
|
||||
/** Config keys rejected with the older-gateway unknown-key response. */
|
||||
val unsupportedConfigKeys: MutableSet<String> = ConcurrentHashMap.newKeySet()
|
||||
|
||||
@@ -329,6 +332,7 @@ class GatewayClientHarness(
|
||||
?: (params["profile"] as? JsonPrimitive)?.contentOrNull
|
||||
?: "default",
|
||||
)
|
||||
if (createdSessionLazy) put("lazy", true)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -768,6 +772,7 @@ class GatewayChatClientTest {
|
||||
val thinkingDeltas = ConcurrentLinkedQueue<String>()
|
||||
val sessionIds = ConcurrentLinkedQueue<String>()
|
||||
val errors = ConcurrentLinkedQueue<String>()
|
||||
val submitRejections = ConcurrentLinkedQueue<String>()
|
||||
val resumeFailures = ConcurrentLinkedQueue<String>()
|
||||
val interactions = ConcurrentLinkedQueue<GatewayAsk>()
|
||||
val interactionExpiries = ConcurrentLinkedQueue<GatewayAskExpiry>()
|
||||
@@ -805,6 +810,7 @@ class GatewayChatClientTest {
|
||||
onResumeFailure = { resumeFailures += it; completeLatch.countDown() },
|
||||
onStatusUpdate = { _, _ -> },
|
||||
onStatusClear = { },
|
||||
onSubmitRejected = { submitRejections += it; completeLatch.countDown() },
|
||||
)
|
||||
}
|
||||
|
||||
@@ -812,6 +818,7 @@ class GatewayChatClientTest {
|
||||
rpcTimeoutMs: Long = 15_000L,
|
||||
promptSubmitTimeoutMs: Long = 1_800_000L,
|
||||
turnIdleTimeoutMs: Long = 180_000L,
|
||||
sessionReadyTimeoutMs: Long = 300_000L,
|
||||
compactingTimeoutMs: Long = 600_000L,
|
||||
callbackDispatcher: (block: () -> Unit) -> Unit = { it() },
|
||||
ticketTimeoutMs: Long = 8_000L,
|
||||
@@ -834,6 +841,7 @@ class GatewayChatClientTest {
|
||||
rpcTimeoutMs = rpcTimeoutMs,
|
||||
promptSubmitTimeoutMs = promptSubmitTimeoutMs,
|
||||
turnIdleTimeoutMs = turnIdleTimeoutMs,
|
||||
sessionReadyTimeoutMs = sessionReadyTimeoutMs,
|
||||
compactingTimeoutMs = compactingTimeoutMs,
|
||||
)
|
||||
|
||||
@@ -872,6 +880,7 @@ class GatewayChatClientTest {
|
||||
rpcTimeoutMs: Long = 15_000L,
|
||||
promptSubmitTimeoutMs: Long = 1_800_000L,
|
||||
turnIdleTimeoutMs: Long = 180_000L,
|
||||
sessionReadyTimeoutMs: Long = 300_000L,
|
||||
compactingTimeoutMs: Long = 600_000L,
|
||||
ticketTimeoutMs: Long = 8_000L,
|
||||
) {
|
||||
@@ -881,6 +890,7 @@ class GatewayChatClientTest {
|
||||
rpcTimeoutMs = rpcTimeoutMs,
|
||||
promptSubmitTimeoutMs = promptSubmitTimeoutMs,
|
||||
turnIdleTimeoutMs = turnIdleTimeoutMs,
|
||||
sessionReadyTimeoutMs = sessionReadyTimeoutMs,
|
||||
compactingTimeoutMs = compactingTimeoutMs,
|
||||
ticketTimeoutMs = ticketTimeoutMs,
|
||||
)
|
||||
@@ -3944,6 +3954,42 @@ class GatewayChatClientTest {
|
||||
assertEquals(true, (submit["queued"] as? JsonPrimitive)?.booleanOrNull)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `lazy fresh session waits for authoritative ready edge before submit`() {
|
||||
harness.createdSessionLazy = true
|
||||
val r = Recorder()
|
||||
|
||||
client.sendTurn(null, "hello", null, r.callbacks) { r.preflightFailures += it }
|
||||
harness.awaitRpc("session.create")
|
||||
val serverWs = harness.awaitServerSocket()
|
||||
Thread.sleep(100)
|
||||
assertEquals(0, harness.rpcLog.count { it.first == "prompt.submit" })
|
||||
|
||||
serverWs.send(
|
||||
harness.eventFrame(
|
||||
"session.info",
|
||||
buildJsonObject { put("lazy", false) },
|
||||
"live-1",
|
||||
),
|
||||
)
|
||||
|
||||
harness.awaitRpc("prompt.submit")
|
||||
assertTrue(r.preflightFailures.isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `lazy fresh session readiness timeout never submits`() {
|
||||
rebuildClient(sessionReadyTimeoutMs = 50L)
|
||||
harness.createdSessionLazy = true
|
||||
val r = Recorder()
|
||||
|
||||
client.sendTurn(null, "hello", null, r.callbacks) { r.preflightFailures += it }
|
||||
harness.awaitRpc("session.create")
|
||||
|
||||
waitUntil { r.preflightFailures.isNotEmpty() }
|
||||
assertEquals(0, harness.rpcLog.count { it.first == "prompt.submit" })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `authoritative prompt rejections surface server message without preflight fallback`() {
|
||||
val cases = listOf(
|
||||
@@ -3964,8 +4010,9 @@ class GatewayChatClientTest {
|
||||
client.sendTurn(null, "hello-$code", null, r.callbacks) { r.preflightFailures += it }
|
||||
harness.awaitRpc("prompt.submit")
|
||||
|
||||
waitUntil { r.errors.isNotEmpty() }
|
||||
assertEquals(listOf(message), r.errors.toList())
|
||||
waitUntil { r.submitRejections.isNotEmpty() }
|
||||
assertEquals(listOf(message), r.submitRejections.toList())
|
||||
assertTrue(r.errors.isEmpty())
|
||||
assertTrue("$code must not trigger SSE fallback", r.preflightFailures.isEmpty())
|
||||
assertEquals(index + 1, harness.rpcLog.count { it.first == "prompt.submit" })
|
||||
}
|
||||
|
||||
+23
@@ -44,6 +44,7 @@ class GatewayEventMapperTest {
|
||||
var usage: UsageInfo? = null
|
||||
var usageCalls = 0
|
||||
val errors = mutableListOf<String>()
|
||||
val submitRejections = mutableListOf<String>()
|
||||
|
||||
val callbacks = GatewayTurnCallbacks(
|
||||
onSessionId = { sessionIds += it },
|
||||
@@ -74,6 +75,7 @@ class GatewayEventMapperTest {
|
||||
onStatusClear = { statusClears += it },
|
||||
onNoticeShow = { notices += it },
|
||||
onNoticeClear = { noticeClears += it },
|
||||
onSubmitRejected = { submitRejections += it },
|
||||
)
|
||||
}
|
||||
|
||||
@@ -483,10 +485,31 @@ class GatewayEventMapperTest {
|
||||
val mapper = mapperWith(r)
|
||||
mapper.onEvent("error", obj("""{"message":"model exploded"}"""))
|
||||
assertEquals(listOf("model exploded"), r.errors)
|
||||
assertTrue(r.submitRejections.isEmpty())
|
||||
assertTrue(mapper.turnEnded)
|
||||
assertEquals(0, r.completes)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `session ownership error is an authoritative submit rejection`() {
|
||||
val r = Recorder()
|
||||
val mapper = mapperWith(r)
|
||||
val message =
|
||||
"Session fresh-1 already has a live owner (webui, pid 42, running 0m). " +
|
||||
"Only one surface at a time may run a session, because a second one would reason from stale history."
|
||||
|
||||
mapper.onEvent(
|
||||
"error",
|
||||
obj(
|
||||
"""{"message":"Session fresh-1 already has a live owner (webui, pid 42, running 0m). Only one surface at a time may run a session, because a second one would reason from stale history."}""",
|
||||
),
|
||||
)
|
||||
|
||||
assertEquals(listOf(message), r.submitRejections)
|
||||
assertTrue(r.errors.isEmpty())
|
||||
assertTrue(mapper.turnEnded)
|
||||
}
|
||||
|
||||
// --- Tools ---
|
||||
|
||||
@Test
|
||||
|
||||
+27
@@ -1,10 +1,15 @@
|
||||
package com.hermesandroid.relay.ui.components
|
||||
|
||||
import android.graphics.Bitmap
|
||||
import android.media.ExifInterface
|
||||
import android.net.Uri
|
||||
import androidx.test.ext.junit.runners.AndroidJUnit4
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertNotNull
|
||||
import org.junit.Test
|
||||
import org.junit.runner.RunWith
|
||||
import org.robolectric.RuntimeEnvironment
|
||||
import java.io.File
|
||||
|
||||
@RunWith(AndroidJUnit4::class)
|
||||
class OrientedBitmapDecoderTest {
|
||||
@@ -36,4 +41,26 @@ class OrientedBitmapDecoderTest {
|
||||
imageOrientationTransform(ExifInterface.ORIENTATION_TRANSVERSE),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `content uri decoder reads encoded image without byte array handoff`() {
|
||||
val context = RuntimeEnvironment.getApplication()
|
||||
val source = File.createTempFile("oriented-bitmap-", ".png", context.cacheDir)
|
||||
val bitmap = Bitmap.createBitmap(3, 2, Bitmap.Config.ARGB_8888)
|
||||
try {
|
||||
source.outputStream().use { output ->
|
||||
bitmap.compress(Bitmap.CompressFormat.PNG, 100, output)
|
||||
}
|
||||
} finally {
|
||||
bitmap.recycle()
|
||||
}
|
||||
|
||||
val decoded = decodeOrientedBitmap(context, Uri.fromFile(source))
|
||||
|
||||
assertNotNull(decoded)
|
||||
assertEquals(3, decoded?.width)
|
||||
assertEquals(2, decoded?.height)
|
||||
decoded?.recycle()
|
||||
source.delete()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
package com.hermesandroid.relay.util
|
||||
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertFalse
|
||||
import org.junit.Test
|
||||
import org.junit.runner.RunWith
|
||||
import org.robolectric.RobolectricTestRunner
|
||||
import org.robolectric.RuntimeEnvironment
|
||||
import org.robolectric.annotation.Config
|
||||
import java.io.File
|
||||
|
||||
@RunWith(RobolectricTestRunner::class)
|
||||
@Config(sdk = [34])
|
||||
class MediaCacheWriterTest {
|
||||
|
||||
@Test
|
||||
fun cacheFile_promotesStreamedFileWithoutByteArrayHandoff() = runTest {
|
||||
val context = RuntimeEnvironment.getApplication()
|
||||
val source = File.createTempFile("hermes-media-test-", ".part", context.cacheDir)
|
||||
.apply { writeText("streamed-media") }
|
||||
val writer = MediaCacheWriter(context) { 32 }
|
||||
|
||||
val cached = writer.cacheFile(source, "audio/mpeg", "voice reply.mp3")
|
||||
|
||||
assertFalse(source.exists())
|
||||
assertEquals("streamed-media", cached.readText())
|
||||
writer.clear()
|
||||
}
|
||||
}
|
||||
+34
@@ -517,6 +517,40 @@ class ChatViewModelGatewayInboundTurnTest {
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun ownershipErrorAfterFreshSubmitStaysVisibleWithoutHistoryRecovery() {
|
||||
val ownershipError =
|
||||
"Session 20260612_120000_abc123 already has a live owner (webui, pid 42, running 0m). " +
|
||||
"Only one surface at a time may run a session, because a second one would reason from stale history."
|
||||
val historyReads = AtomicInteger(0)
|
||||
viewModel.setProfileMessageLoader {
|
||||
historyReads.incrementAndGet()
|
||||
Result.success(emptyList())
|
||||
}
|
||||
viewModel.createNewChat()
|
||||
|
||||
viewModel.sendMessage("Keep this retryable")
|
||||
gatewayHarness.awaitRpc("prompt.submit")
|
||||
serverWs.send(
|
||||
gatewayHarness.eventFrame(
|
||||
"error",
|
||||
buildJsonObject { put("message", ownershipError) },
|
||||
"live-1",
|
||||
),
|
||||
)
|
||||
|
||||
awaitCondition { viewModel.chatFailure.value?.rawError == ownershipError }
|
||||
assertFalse(handler.isStreaming.value)
|
||||
assertFalse(viewModel.recoveringAnswer.value)
|
||||
assertEquals(0, historyReads.get())
|
||||
assertTrue(handler.messages.value.any { it.content == "Keep this retryable" })
|
||||
assertFalse(
|
||||
handler.messages.value.any {
|
||||
it.content.contains("unfinished message was not found", ignoreCase = true)
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
DiagnosticsLog.clear()
|
||||
|
||||
+7
-1
@@ -11,6 +11,7 @@ import com.hermesandroid.relay.network.upstream.DashboardApiClient
|
||||
import com.hermesandroid.relay.network.upstream.models.MessageItem
|
||||
import com.hermesandroid.relay.util.MediaCacheWriter
|
||||
import io.mockk.coEvery
|
||||
import io.mockk.coVerify
|
||||
import io.mockk.mockk
|
||||
import kotlinx.serialization.json.JsonPrimitive
|
||||
import okhttp3.OkHttpClient
|
||||
@@ -25,6 +26,7 @@ import org.robolectric.RobolectricTestRunner
|
||||
import org.robolectric.RuntimeEnvironment
|
||||
import org.robolectric.Shadows.shadowOf
|
||||
import org.robolectric.annotation.Config
|
||||
import java.io.File
|
||||
|
||||
@RunWith(RobolectricTestRunner::class)
|
||||
@Config(sdk = [34])
|
||||
@@ -52,7 +54,9 @@ class ChatViewModelMediaStateTest {
|
||||
pairedTokenSnapshot = { "paired-session" },
|
||||
)
|
||||
cache = mockk()
|
||||
coEvery { cache.cache(any(), any(), any()) } returns
|
||||
coEvery { cache.cache(any<ByteArray>(), any(), any()) } returns
|
||||
Uri.parse("content://com.axiomlabs.hermesrelay.fileprovider/hermes-media/photo.jpg")
|
||||
coEvery { cache.cache(any<File>(), any(), any()) } returns
|
||||
Uri.parse("content://com.axiomlabs.hermesrelay.fileprovider/hermes-media/photo.jpg")
|
||||
handler = ChatHandler()
|
||||
viewModel = ChatViewModel().also {
|
||||
@@ -197,6 +201,8 @@ class ChatViewModelMediaStateTest {
|
||||
}.attachments.single()
|
||||
assertEquals("audio/mpeg", loaded.contentType)
|
||||
assertEquals("test-voice-message.mp3", loaded.fileName)
|
||||
coVerify(exactly = 1) { cache.cache(any<File>(), "audio/mpeg", "test-voice-message.mp3") }
|
||||
coVerify(exactly = 0) { cache.cache(any<ByteArray>(), any(), any()) }
|
||||
val request = dashboardServer.takeRequest()
|
||||
assertEquals("/api/files/download", request.requestUrl?.encodedPath)
|
||||
assertEquals(path, request.requestUrl?.queryParameter("path"))
|
||||
|
||||
@@ -68,6 +68,7 @@ the upstream contract identifiers it depends on.
|
||||
|---|---|
|
||||
| `initial_history_bind` | Durable, profile-scoped history is already available when the client resumes and first binds its rendered transcript |
|
||||
| `ordinary_turn` | Normal message start, deltas, completion, and persisted history |
|
||||
| `ownership_rejection` | A submit acknowledged before the defense-in-depth ownership check emits the canonical terminal refusal; no user/model row is persisted and clients must not enter history recovery |
|
||||
| `compaction_status` | Compaction status is client-visible before terminal completion and may repeat as a heartbeat |
|
||||
| `rapid_tools_interims` | Rapid chunks, reasoning, tool activity, and interim assistant boundaries |
|
||||
| `queued_follow_up` | Two explicitly owned turns and ordered queue drainage |
|
||||
|
||||
@@ -34,6 +34,10 @@ Verified upstream source snapshot:
|
||||
`apps/desktop/src/lib/media.ts`,
|
||||
`apps/desktop/src/components/assistant-ui/markdown-text.tsx`, and
|
||||
`hermes_cli/web_server.py`.
|
||||
- Gateway per-session submit exclusivity was rechecked against upstream `main`
|
||||
at `c0495c6bce6c4988f79a2e0355e5aa047a3b03af` in
|
||||
`hermes_cli/active_sessions.py`, `tui_gateway/methods_prompt.py`, and
|
||||
`tui_gateway/server.py`.
|
||||
|
||||
## Ownership
|
||||
|
||||
@@ -46,7 +50,7 @@ Verified upstream source snapshot:
|
||||
| `/api/sessions/*` | Upstream Dashboard/Gateway and API server | No | Primary profile-scoped Dashboard directory/history or API-only compatibility storage | Native upstream session list/create/read/update/delete/messages/fork/chat/chat-stream. The standard Android drawer and stored-history reader use authenticated Dashboard REST independently of `/api/ws` readiness; the Gateway socket owns live chat and activity, not whether persisted rows may be read. Dashboard lists expose profile-stamped `pinned`/`archived`, accept `archived=exclude\|only\|include`, and PATCH either durable flag in the owning profile DB. The API-server resource also exposes and patches both fields, but its current list omits archived rows and has no archive filter; Android therefore offers restart-safe archive/restore only on the Dashboard path while API-only pinning remains valid. Newer Dashboard hosts also expose single-session JSON export and guarded bulk cleanup; Android must dry-run prune first. The bootstrap no longer injects session CRUD/messages/fork routes; only `/api/sessions/search` remains a compatibility route. |
|
||||
| `/v1/skills`, `/v1/toolsets` | Upstream API server | No | Discovery | Authenticated read-only API-server skill/toolset inventory; Android Diagnostics summarizes enabled toolsets and Relay tool visibility. |
|
||||
| Dashboard `/api/status`, `/api/auth/me` | Upstream dashboard | No | Manage auth and post-selection diagnostics | Dashboard cookie/session path; separate from API bearer. Optional status diagnostics include Nous bootstrap validity, resource pressure, and profile/gateway topology; these do not gate transport selection. |
|
||||
| Dashboard `/api/auth/ws-ticket`, `/api/ws` | Upstream dashboard/tui_gateway | No | Preferred chat transport | Vanilla Hermes gateway chat path with live reasoning/thinking events. `message.complete` is the ordinary terminal event; `session.info {running:false}` is the authoritative settle backstop when a replacement socket missed that terminal frame. A reconnect reactivates the exact live runtime with `session.activate`; durable `session.resume` remains the cold-open path and an explicit rejection never creates a replacement context. |
|
||||
| Dashboard `/api/auth/ws-ticket`, `/api/ws` | Upstream dashboard/tui_gateway | No | Preferred chat transport | Vanilla Hermes gateway chat path with live reasoning/thinking events. The active-session surface label `webui` identifies this official Dashboard embedded runtime; it is not evidence of an open browser tab, third-party WebUI, or separate client owner. `session.create` is lazy and its deferred build emits exact-session `session.info` when the runtime is ready; Android awaits that edge before its first prompt instead of racing into Dashboard compute-host isolation. `prompt.submit` atomically claims the durable session owner and rejects conflicts with JSON-RPC `4090` plus reason `SESSION_NOT_OWNED` before persisting a user row or starting a model turn. The defense-in-depth turn chokepoint can also emit the same canonical refusal as a terminal `error` event after submit acknowledgement; Android preserves that error and keeps the user row retryable without history recovery. `message.complete` is the ordinary terminal event; `session.info {running:false}` is the authoritative settle backstop when a replacement socket missed that terminal frame. A reconnect reactivates the exact live runtime with `session.activate`; durable `session.resume` remains the cold-open path and an explicit rejection never creates a replacement context. |
|
||||
| Gateway `session.active_list` | Upstream tui_gateway | No | Authoritative process-wide live activity | Returns attachable runtimes across the Gateway process, with live `id`, durable `session_key`, and `starting`, `working`, `waiting`, or `idle`. The only optional selector is `current_session_id`; rows normally carry no profile metadata. Android attributes a row from exact foreground/detached ownership already held by that client, explicit profile metadata if a future upstream sends it, or a unique match to the currently selected passive session in the current connection directory. Duplicate same-id owners across profiles remain unresolved. An exact live/durable Idle row may settle only the same Android-owned turn and unchanged progress generation when a terminal frame is missing; it never claims a passively observed Desktop/TUI turn. Unresolved rows stay unattributed, and absence settles a scope only after a complete, unambiguously resolved successful snapshot. Method-not-found or refresh failure is Unavailable, not Idle. Pending input outranks running work. |
|
||||
| Dashboard `model.options` / `/api/model/*` | Upstream dashboard/tui_gateway | No | Provider/model inventory and selection | Source of truth for coherent provider/model identities. A reasoning boolean or exact effort list is consumed when present; clients do not infer provider identity from a model string alone. |
|
||||
| Gateway `pet.info`, `pet.gallery`, `pet.select`, `pet.disable` | Upstream tui_gateway | No | Profile-scoped animated companion | `pet.info` supplies bounded PNG/WebP sheet bytes, revision, geometry, real frame counts, loop timing, scale, and row taxonomy. Android passes `knownRevision` to avoid duplicate sheet transfer, renders the active pet through its native activity-aware companion, and keeps phone-local pet packs separate. All four RPCs carry the effective profile. |
|
||||
|
||||
@@ -23,6 +23,8 @@ from typing import Iterable, Sequence
|
||||
|
||||
SERVER = "tui_gateway/server.py"
|
||||
SESSION_METHODS = "tui_gateway/methods_session.py"
|
||||
PROMPT_METHODS = "tui_gateway/methods_prompt.py"
|
||||
ACTIVE_SESSIONS = "hermes_cli/active_sessions.py"
|
||||
API_SERVER = "gateway/platforms/api_server.py"
|
||||
|
||||
GATEWAY_TERMINAL = "gateway.message_complete"
|
||||
@@ -30,6 +32,7 @@ GATEWAY_SETTLED_INFO = "gateway.settled_session_info"
|
||||
SESSION_ACTIVATE = "gateway.session_activate_live"
|
||||
SESSION_RESUME = "gateway.session_resume_durable"
|
||||
SESSION_ACTIVE_LIST = "gateway.session_active_list"
|
||||
SESSION_EXCLUSIVE_SUBMIT = "gateway.session_exclusive_submit"
|
||||
SUBAGENT_CHILD_WATCH = "gateway.subagent_child_watch"
|
||||
API_BOUNDARY = "api.fallback_boundary"
|
||||
ALL_CONTRACTS = (
|
||||
@@ -38,6 +41,7 @@ ALL_CONTRACTS = (
|
||||
SESSION_ACTIVATE,
|
||||
SESSION_RESUME,
|
||||
SESSION_ACTIVE_LIST,
|
||||
SESSION_EXCLUSIVE_SUBMIT,
|
||||
SUBAGENT_CHILD_WATCH,
|
||||
API_BOUNDARY,
|
||||
)
|
||||
@@ -409,6 +413,59 @@ def _check_subagent_child_watch(server: SourceFile, methods: SourceFile) -> Chec
|
||||
return CheckResult(contract, False, (), str(exc))
|
||||
|
||||
|
||||
def _check_session_exclusive_submit(
|
||||
server: SourceFile,
|
||||
session_methods: SourceFile,
|
||||
prompt_methods: SourceFile,
|
||||
active_sessions: SourceFile,
|
||||
) -> CheckResult:
|
||||
contract = SESSION_EXCLUSIVE_SUBMIT
|
||||
try:
|
||||
submit = prompt_methods.method_handler("prompt.submit")
|
||||
create = session_methods.method_handler("session.create")
|
||||
build = server.function("_start_agent_build")
|
||||
run_submit = server.function("_run_prompt_submit")
|
||||
submit_text = prompt_methods.segment(submit)
|
||||
create_text = session_methods.segment(create)
|
||||
build_text = server.segment(build)
|
||||
run_text = server.segment(run_submit)
|
||||
active_text = active_sessions.text
|
||||
missing: list[str] = []
|
||||
if "_ensure_active_session_slot" not in submit_text or "4090" not in submit_text:
|
||||
missing.append("prompt.submit atomic ownership refusal")
|
||||
if '"reason"' not in submit_text:
|
||||
missing.append("machine-readable refusal reason")
|
||||
if "_ensure_active_session_slot" not in run_text or '"error"' not in run_text:
|
||||
missing.append("defense-in-depth terminal error event")
|
||||
if "_schedule_agent_build" not in create_text or '"lazy"' not in create_text:
|
||||
missing.append("lazy session.create readiness contract")
|
||||
if '"session.info"' not in build_text or "ready.set()" not in build_text:
|
||||
missing.append("deferred agent-ready session.info edge")
|
||||
for marker in (
|
||||
'SESSION_NOT_OWNED = "SESSION_NOT_OWNED"',
|
||||
"PER_SESSION_EXCLUSIVE_SUBMIT = True",
|
||||
"already has a live owner",
|
||||
"Only one surface at a time may run a session",
|
||||
):
|
||||
if marker not in active_text:
|
||||
missing.append(marker)
|
||||
if missing:
|
||||
raise ValueError("exclusive submit contract missing: " + ", ".join(missing))
|
||||
return CheckResult(
|
||||
contract,
|
||||
True,
|
||||
(
|
||||
prompt_methods.evidence(submit, "prompt.submit refuses ownership conflicts before turn start"),
|
||||
session_methods.evidence(create, "session.create advertises a lazy deferred build"),
|
||||
server.evidence(build, "deferred build emits session.info before setting ready"),
|
||||
server.evidence(run_submit, "synthesized turns recheck ownership and emit terminal error"),
|
||||
f"{ACTIVE_SESSIONS}: SESSION_NOT_OWNED and canonical holder-aware message",
|
||||
),
|
||||
)
|
||||
except ValueError as exc:
|
||||
return CheckResult(contract, False, (), str(exc))
|
||||
|
||||
|
||||
def load_requirements(manifest: Path | None) -> tuple[str, ...]:
|
||||
if manifest is None:
|
||||
return ALL_CONTRACTS
|
||||
@@ -433,8 +490,10 @@ def load_requirements(manifest: Path | None) -> tuple[str, ...]:
|
||||
def audit_sources(root: Path, requirements: Iterable[str]) -> list[CheckResult]:
|
||||
server = SourceFile(root, SERVER)
|
||||
methods = SourceFile(root, SESSION_METHODS)
|
||||
prompt_methods = SourceFile(root, PROMPT_METHODS)
|
||||
active_sessions = SourceFile(root, ACTIVE_SESSIONS)
|
||||
api = SourceFile(root, API_SERVER)
|
||||
source_files = (server, methods, api)
|
||||
source_files = (server, methods, prompt_methods, active_sessions, api)
|
||||
fork_hits = [
|
||||
f"{source.relative}:{marker}"
|
||||
for source in source_files
|
||||
@@ -450,6 +509,9 @@ def audit_sources(root: Path, requirements: Iterable[str]) -> list[CheckResult]:
|
||||
SESSION_ACTIVATE: lambda: _check_activate(server, methods),
|
||||
SESSION_RESUME: lambda: _check_resume(methods),
|
||||
SESSION_ACTIVE_LIST: lambda: _check_active_list(server, methods),
|
||||
SESSION_EXCLUSIVE_SUBMIT: lambda: _check_session_exclusive_submit(
|
||||
server, methods, prompt_methods, active_sessions,
|
||||
),
|
||||
SUBAGENT_CHILD_WATCH: lambda: _check_subagent_child_watch(server, methods),
|
||||
API_BOUNDARY: lambda: _check_api_boundary(api),
|
||||
}
|
||||
|
||||
@@ -30,7 +30,14 @@ def _live_session_payload(sid, session):
|
||||
"running": False, "status": "idle",
|
||||
}
|
||||
|
||||
def _start_agent_build(sid, session):
|
||||
_emit("session.info", sid, {"lazy": False})
|
||||
ready.set()
|
||||
|
||||
def _run_prompt_submit(sid, session, agent):
|
||||
if _ensure_active_session_slot(sid, session) is not None:
|
||||
_emit("error", sid, {})
|
||||
return False
|
||||
try:
|
||||
_emit("message.complete", sid, {})
|
||||
finally:
|
||||
@@ -68,6 +75,11 @@ METHODS_SOURCE = '''
|
||||
def method(name):
|
||||
return lambda fn: fn
|
||||
|
||||
@method("session.create")
|
||||
def _(rid, params):
|
||||
_schedule_agent_build("live")
|
||||
return {"session_id": "live", "info": {"lazy": True}}
|
||||
|
||||
@method("session.resume")
|
||||
def _(rid, params):
|
||||
target = params.get("session_id", "")
|
||||
@@ -106,6 +118,26 @@ def _(rid, params):
|
||||
return _ok(rid, {"sessions": rows})
|
||||
'''
|
||||
|
||||
PROMPT_METHODS_SOURCE = '''
|
||||
def method(name):
|
||||
return lambda fn: fn
|
||||
|
||||
@method("prompt.submit")
|
||||
def _(rid, params):
|
||||
session = sessions[params["session_id"]]
|
||||
if (refusal := _ensure_active_session_slot(params["session_id"], session)) is not None:
|
||||
return _err(rid, 4090, str(refusal), {"reason": refusal.reason})
|
||||
return _ok(rid, {"ok": True})
|
||||
'''
|
||||
|
||||
ACTIVE_SESSIONS_SOURCE = '''
|
||||
SESSION_NOT_OWNED = "SESSION_NOT_OWNED"
|
||||
PER_SESSION_EXCLUSIVE_SUBMIT = True
|
||||
|
||||
def session_already_owned_message(session_id, entry):
|
||||
return f"Session {session_id} already has a live owner ({entry}). Only one surface at a time may run a session."
|
||||
'''
|
||||
|
||||
API_SOURCE = '''
|
||||
ROUTES = [
|
||||
("POST", "/api/sessions/{session_id}/chat/stream"),
|
||||
@@ -132,6 +164,8 @@ class GatewayScenarioConformanceTest(unittest.TestCase):
|
||||
sources = {
|
||||
module.SERVER: SERVER_SOURCE,
|
||||
module.SESSION_METHODS: METHODS_SOURCE,
|
||||
module.PROMPT_METHODS: PROMPT_METHODS_SOURCE,
|
||||
module.ACTIVE_SESSIONS: ACTIVE_SESSIONS_SOURCE,
|
||||
module.API_SERVER: API_SOURCE,
|
||||
}
|
||||
for relative, text in sources.items():
|
||||
|
||||
@@ -123,6 +123,25 @@ class FixtureTestCase(unittest.IsolatedAsyncioTestCase):
|
||||
self.assertEqual(["user", "assistant"], [row["role"] for row in history["messages"]])
|
||||
self.assertEqual(2, history["pagination"]["returned"])
|
||||
|
||||
async def test_ownership_rejection_is_terminal_without_persisted_turn(self) -> None:
|
||||
fixture, base_url = await self.start("ownership_rejection")
|
||||
ws, _ = await self.connect(base_url)
|
||||
await self.rpc(ws, 1, "session.create", {"profile": "default"})
|
||||
await ws.receive_json()
|
||||
await self.rpc(ws, 2, "prompt.submit", {"text": "not recorded in evidence"})
|
||||
frames = await self.frames_until(
|
||||
ws,
|
||||
lambda frame: frame.get("params", {}).get("type") == "error",
|
||||
)
|
||||
error = frames[-1]["params"]["payload"]["message"]
|
||||
self.assertIn("already has a live owner (webui", error)
|
||||
self.assertIn("Only one surface at a time may run a session", error)
|
||||
async with self.session.get(
|
||||
f"{base_url}/api/sessions/{fixture.scenario.stored_session_id}/messages",
|
||||
) as response:
|
||||
history = await response.json()
|
||||
self.assertEqual([], history["messages"])
|
||||
|
||||
async def test_compaction_status_repeats_before_terminal_completion(self) -> None:
|
||||
fixture, base_url = await self.start("compaction_status")
|
||||
ws, _ = await self.connect(base_url)
|
||||
@@ -381,6 +400,7 @@ class ScenarioTestCase(unittest.TestCase):
|
||||
"cross_client_observation",
|
||||
"initial_history_bind",
|
||||
"ordinary_turn",
|
||||
"ownership_rejection",
|
||||
"rapid_tools_interims",
|
||||
"subagent_child_preview",
|
||||
"terminal_gap_activate",
|
||||
@@ -433,6 +453,10 @@ class ScenarioTestCase(unittest.TestCase):
|
||||
("gateway.message_complete", "gateway.session_active_list"),
|
||||
load_scenario("cross_client_observation").contract_requirements,
|
||||
)
|
||||
self.assertEqual(
|
||||
("gateway.session_exclusive_submit",),
|
||||
load_scenario("ownership_rejection").contract_requirements,
|
||||
)
|
||||
|
||||
def test_tls_arguments_must_be_paired(self) -> None:
|
||||
with contextlib.redirect_stderr(io.StringIO()):
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
{
|
||||
"name": "ownership_rejection",
|
||||
"live_session_id": "fixture-live-owner-conflict",
|
||||
"stored_session_id": "20260901_165658_owner_conflict",
|
||||
"profile": "default",
|
||||
"contract_requirements": [
|
||||
"gateway.session_exclusive_submit"
|
||||
],
|
||||
"initial_history": [],
|
||||
"turns": [
|
||||
{
|
||||
"steps": [
|
||||
{
|
||||
"op": "event",
|
||||
"type": "error",
|
||||
"payload": {
|
||||
"message": "Session 20260901_165658_owner_conflict already has a live owner (webui, pid 42, running 0m). Only one surface at a time may run a session, because a second one would reason from a transcript that does not include the first one's work."
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
Reference in New Issue
Block a user