mirror of
https://github.com/zapstore/zapstore.git
synced 2026-10-05 12:38:24 +00:00
Implement persistent query max-age
This commit is contained in:
+30
-1
@@ -1,6 +1,6 @@
|
||||
# purplequartz
|
||||
|
||||
`purplequartz` is an Android API 29+ query wrapper for Quartz 1.12.6. Local and local-and-remote queries expose events read from Quartz's persistent event store. Remote-only queries expose verified relay events directly and save persistent events in the background.
|
||||
`purplequartz` is an Android API 29+ query wrapper for Quartz 1.12.6. Local and local-and-remote queries expose events read from Quartz's persistent event store. Remote-only queries expose verified relay events directly and save persistent events in the background. Local-and-remote queries can use a persistent max-age policy to avoid unnecessary relay requests.
|
||||
|
||||
## Setup
|
||||
|
||||
@@ -67,11 +67,40 @@ purpleQuartz.query(
|
||||
|
||||
Every network query requires an explicit, nonempty relay set. The library does not add outbox, fallback, discovery, or hinted relays.
|
||||
|
||||
### Cached local-and-remote queries
|
||||
|
||||
Set `maxAge` when data may be served from the local store without immediately opening a relay request. For example, this profile query refreshes at most once every six hours:
|
||||
|
||||
```kotlin
|
||||
purpleQuartz.query(
|
||||
filter = Filter(
|
||||
authors = listOf(profilePubkey),
|
||||
kinds = listOf(0),
|
||||
limit = 1,
|
||||
),
|
||||
source = QuerySource.LocalAndRemote(
|
||||
relays = relays,
|
||||
mode = RemoteMode.OneShot(),
|
||||
maxAge = 6.hours,
|
||||
),
|
||||
).collect(::renderState)
|
||||
```
|
||||
|
||||
Freshness is recorded only after every requested relay reaches EOSE and preceding events are committed. Empty successful responses are cached too. Timeouts, failures, partial responses, cancellation, and local-only changes do not refresh the timestamp.
|
||||
|
||||
The cache key contains the complete filter set and exact normalized relay set; `maxAge` and remote mode are policy and are not part of the key. Only a SHA-256 fingerprint and refresh timestamp are persisted, not filter contents. Freshness survives process restarts and is reset when the Quartz database file is deleted or replaced. Out-of-band in-place mutation of the owned Quartz database is unsupported.
|
||||
|
||||
- Fresh `OneShot` queries emit the local projection as `Complete` and open no request.
|
||||
- Fresh `Stream` queries report `Cached`, keep observing local commits, and defer their relay subscription until the max-age window expires. A successful concurrent refresh extends that delay.
|
||||
- `Remote` always opens its requested subscription and bypasses this local cache policy.
|
||||
- Omitting `maxAge` preserves the existing always-refresh behavior.
|
||||
|
||||
## State and failure handling
|
||||
|
||||
`QueryState.items` survives connection failures and synchronization transitions. Inspect `sync` for query progress and `error` for the latest failure. PurpleQuartz observes Android's default network callback and requests an immediate retry when a network becomes available, without waiting for relay backoff. Hosts should also call `purpleQuartz.refreshConnections()` when their app returns to the foreground.
|
||||
|
||||
- `LocalOnly`: the query reads and observes the local store.
|
||||
- `Cached`: a local-and-remote stream is observing the fresh local projection and has deferred its relay request.
|
||||
- `Connecting`: at least one relay has no active request.
|
||||
- `CatchingUp`: every relay is connected and at least one is waiting for EOSE.
|
||||
- `Live`: all relays reached EOSE and a streaming request remains active.
|
||||
|
||||
+56
-11
@@ -2,7 +2,7 @@
|
||||
|
||||
**Status:** Implementation-ready
|
||||
**Created:** 2026-07-11
|
||||
**Revised:** 2026-07-11
|
||||
**Revised:** 2026-07-12
|
||||
**Target:** Version 1
|
||||
**Platform:** Android API 29+
|
||||
**Library/module name:** `purplequartz`
|
||||
@@ -24,7 +24,8 @@ Local / LocalAndRemote:
|
||||
query
|
||||
-> install local observation
|
||||
-> emit the permitted local seed
|
||||
-> optionally start relay REQ
|
||||
-> consult optional max-age freshness
|
||||
-> start relay REQ now, defer it, or complete from cache
|
||||
-> verify each relay event
|
||||
-> commit it to the local event store
|
||||
-> observe the committed store change
|
||||
@@ -139,6 +140,12 @@ Kinds `20000..29999` have no persistent local projection:
|
||||
- `Remote` may emit them directly after validation when they are not already expired;
|
||||
- ephemeral events are never inserted into the canonical event store.
|
||||
|
||||
### 3.6 Query freshness is synchronization metadata
|
||||
|
||||
`LocalAndRemote.maxAge` controls when a local projection may suppress or defer a relay request. Freshness is not inferred from event `createdAt`: that timestamp is author-controlled and cannot represent when an empty or nonempty query last synchronized.
|
||||
|
||||
Freshness metadata contains only a versioned SHA-256 query fingerprint and a successful-refresh timestamp. Quartz remains the sole persistent event source of truth. Deleting or replacing the Quartz database file changes its filesystem identity and rotates the metadata namespace so an old marker cannot suppress synchronization against a new store. Out-of-band in-place mutation of the façade-owned database file is unsupported.
|
||||
|
||||
## 4. Public API
|
||||
|
||||
The declarations below are normative. Implementation imports the corresponding public types from the pinned Quartz artifact and `kotlinx.coroutines`.
|
||||
@@ -175,7 +182,7 @@ Rules:
|
||||
- `query` is the primary v1 API.
|
||||
- The returned `Flow` is cold. Each collection owns one local observer and, when applicable, one remote request.
|
||||
- Collection cancellation closes that request and removes all listeners owned by the collection.
|
||||
- For `LocalAndRemote`, the local observer is active and the first local seed is sent to the collector before the remote request starts.
|
||||
- For `LocalAndRemote`, the local observer is active and the first local seed is sent before a remote request starts. A fresh max-age marker may suppress one-shot startup or defer stream startup.
|
||||
- Multiple filters produce one logical query state.
|
||||
- The façade does not expose `NostrClient`, `subscribeAsFlow`, relay event callbacks, the SQLite connection pool, or raw SQL.
|
||||
- `create` stores `context.applicationContext`, resolves `databaseName` with `Context.getDatabasePath`, creates `EventStore(dbName = absolutePath, relay = null)` with Quartz's default indexing strategy and published fixed reader count, wraps it in exactly one `ObservableEventStore`, and creates exactly one `NostrClient`.
|
||||
@@ -198,8 +205,12 @@ sealed interface QuerySource {
|
||||
data class LocalAndRemote(
|
||||
val relays: Set<NormalizedRelayUrl>,
|
||||
val mode: RemoteMode = RemoteMode.Stream,
|
||||
val maxAge: Duration? = null,
|
||||
) : QuerySource {
|
||||
init { require(relays.isNotEmpty()) }
|
||||
init {
|
||||
require(relays.isNotEmpty())
|
||||
require(maxAge == null || (maxAge.isFinite() && maxAge.isPositive()))
|
||||
}
|
||||
}
|
||||
|
||||
data class Remote(
|
||||
@@ -221,7 +232,7 @@ sealed interface RemoteMode {
|
||||
}
|
||||
```
|
||||
|
||||
`LocalAndRemote` and `Remote` reject an empty relay set. `OneShot.timeout == null` uses `PurpleQuartzConfig.oneShotTimeout`; an explicit timeout must be finite and greater than zero.
|
||||
`LocalAndRemote` and `Remote` reject an empty relay set. `OneShot.timeout == null` uses `PurpleQuartzConfig.oneShotTimeout`; an explicit timeout must be finite and greater than zero. A non-null `maxAge` must be finite and greater than zero; null preserves always-refresh behavior.
|
||||
|
||||
At collection start, the implementation snapshots the relay set and deep-copies each filter's lists and tag maps. That immutable snapshot is used for both store queries and the Quartz subscription, so caller mutation after collection starts cannot change an active query.
|
||||
|
||||
@@ -235,7 +246,9 @@ At collection start, the implementation snapshots the relay set and deep-copies
|
||||
#### `QuerySource.LocalAndRemote`
|
||||
|
||||
- Emits the current local result set without waiting for the network.
|
||||
- Starts the remote REQ after local observation is installed.
|
||||
- With no fresh max-age marker, starts the remote REQ after local observation is installed.
|
||||
- A fresh one-shot emits the local projection as terminal `Complete` without opening a REQ.
|
||||
- A fresh stream reports `Cached`, keeps its local observer active, and starts its REQ when the marker expires. It rechecks before starting so another successful synchronization can extend the delay.
|
||||
- Re-runs the local projection after every observable store change and emits only changed results.
|
||||
- Includes all locally stored events matching the filters, regardless of which request inserted them.
|
||||
- Retains cached results through network errors and reconnects.
|
||||
@@ -247,6 +260,17 @@ At collection start, the implementation snapshots the relay set and deep-copies
|
||||
- Saves persistent incoming events to the local store without delaying direct emission.
|
||||
- Preserves relay arrival order within the query.
|
||||
- Does not promise durable result provenance.
|
||||
- Never consults or establishes local-and-remote freshness; it is the explicit cache-bypass source.
|
||||
|
||||
#### Query freshness identity and success boundary
|
||||
|
||||
- The cache identity is a canonical, versioned fingerprint of every snapshotted filter field and the exact normalized relay set.
|
||||
- Filter and relay ordering are canonicalized where query semantics are order-independent. Null and empty fields remain distinct.
|
||||
- `maxAge`, `RemoteMode`, timeout, subscription ID, and connection generation are excluded because they do not change query coverage.
|
||||
- Filter values are hashed before persistence; search text, tag values, authors, IDs, and relay URLs are not stored in plaintext metadata.
|
||||
- A marker advances only when all current relay generations reach EOSE after every preceding accepted EVENT reaches a terminal verification/persistence outcome. One-shot local-and-remote queries also complete their final local projection first.
|
||||
- Empty successful responses establish freshness. Timeout, failure, partial EOSE, cancellation, local mutation, and clock rollback do not.
|
||||
- Markers survive façade and process recreation. A missing or replaced Quartz database file rotates its metadata identity. Metadata corruption or storage failure causes an extra fetch rather than a stale cache hit.
|
||||
|
||||
### 4.3 Relay selection
|
||||
|
||||
@@ -271,6 +295,7 @@ data class QueryState(
|
||||
|
||||
sealed interface QuerySync {
|
||||
data object LocalOnly : QuerySync
|
||||
data object Cached : QuerySync
|
||||
data object Connecting : QuerySync
|
||||
data object CatchingUp : QuerySync
|
||||
data object Live : QuerySync
|
||||
@@ -333,18 +358,20 @@ State rules:
|
||||
|
||||
- For `Local` and `LocalAndRemote`, `items` comes from an event-store query or projection.
|
||||
- For `Remote`, `items` contains validated events emitted by that query's relay subscription.
|
||||
- `Cached` means a fresh local-and-remote stream is observing local data while its relay request is intentionally deferred.
|
||||
- `Connecting` means at least one requested relay has no active sent REQ for its current connection.
|
||||
- `CatchingUp` means every requested relay has an active generation and at least one has not reached EOSE.
|
||||
- `Live` means every currently routed relay reached EOSE and a streaming query remains subscribed.
|
||||
- `Complete` is terminal success for a one-shot remote query.
|
||||
- `TimedOut` and `Failed` stop the remote request but do not erase `items`.
|
||||
- Every requested relay is present in `relays` from the first network state. Generation `1` is allocated immediately before the initial `subscribe` call.
|
||||
- Every requested relay is present in `relays` from the first network state. Cached states have an empty relay map because no request generation exists yet. Generation `1` is allocated for the transition to `Connecting` before the initial `subscribe` call.
|
||||
- The first `onSubscriptionStarted` for that generation sets `connection = Connected`, resets `eose = false`, and clears transient `lastError`.
|
||||
- After a generation has started, a reconnect's `onConnecting` increments that relay's generation, sets `connection = Connecting`, and resets `eose`; `onSubscriptionStarted` then marks the replacement REQ active. Initial connection attempts before generation 1 starts do not increment it.
|
||||
- EOSE marks only the generation current when its callback entered the ordered ingestion boundary.
|
||||
- A source with zero matching events must still progress out of its initial loading/catching-up state.
|
||||
- `Local` emits exactly one initial `LocalOnly` state even when empty, then emits only when a relevant store change produces a different item list.
|
||||
- `LocalAndRemote` first emits its local seed with `Connecting`, starts the REQ, and then updates synchronization independently of data.
|
||||
- A stale or uncached `LocalAndRemote` first emits its local seed with `Connecting`, starts the REQ, and then updates synchronization independently of data.
|
||||
- A fresh `LocalAndRemote(OneShot)` emits one terminal `Complete`; a fresh `LocalAndRemote(Stream)` first emits `Cached` and transitions to `Connecting` only when freshness expires.
|
||||
- `Remote` first emits an empty `Connecting` state and never reads the store to seed `items`.
|
||||
- A one-shot flow emits exactly one terminal `Complete`, `TimedOut`, or `Failed` state and then completes. Before `LocalAndRemote` emits `Complete`, it performs a final local query so all preceding successful commits are represented.
|
||||
- `LocalAndRemote(Stream)` remains collected after `Live`; `LocalAndRemote(OneShot)` stops local observation after its terminal state.
|
||||
@@ -397,6 +424,7 @@ Validation:
|
||||
- one Quartz client;
|
||||
- one SQLite event store;
|
||||
- one periodic expiration-sweep job;
|
||||
- one bounded, versioned SharedPreferences freshness namespace for its canonical database identity;
|
||||
- request and ingestion coordinators.
|
||||
|
||||
Each remotely backed collection owns one bounded FIFO channel of `ingestionCapacity`, one ingestion worker, one Quartz subscription ID, one subscription listener, and one connection listener. Store access may run on the façade child scope, but collection cancellation remains linked to the collecting coroutine.
|
||||
@@ -405,7 +433,7 @@ An internal lifecycle gate rejects new event-store operations after shutdown sta
|
||||
|
||||
Every admitted store operation decrements the in-flight count from a `NonCancellable` `finally` block. Cancellation cannot strand the count above zero or make `close` wait forever.
|
||||
|
||||
The event store must use Quartz's supported bundled SQLite driver and configuration. Production code must not access Quartz's connection pool, mutate Quartz's schema, or maintain a second canonical event store.
|
||||
The event store must use Quartz's supported bundled SQLite driver and configuration. Production code must not access Quartz's connection pool, mutate Quartz's schema, or maintain a second canonical event store. Query freshness metadata is not an event store: it is capped at 1,024 hashed timestamps per database generation and may be discarded without affecting correctness.
|
||||
|
||||
While open, the façade calls `ObservableEventStore.deleteExpiredEvents()` once per `expirationSweepInterval`. The first sweep occurs after one full interval. Query-list distinctness suppresses no-op sweep emissions.
|
||||
|
||||
@@ -433,6 +461,8 @@ An unexpected persistence failure fails the affected query. It must not be conve
|
||||
|
||||
Streaming queries continue ingesting events after EOSE. A reconnect starts a new generation, resets EOSE, and resends the active filters.
|
||||
|
||||
When all relays reach EOSE, a local-and-remote stream records freshness before reporting `Live`. A local-and-remote one-shot first performs and emits its final store projection, then records freshness and completes. Metadata-write failure does not convert a successful query into failure; absence of a durable marker causes a later collection to fetch again.
|
||||
|
||||
The ingestion boundary is bounded. Saturation must fail and close the whole collected subscription; EVENT and EOSE messages must never be silently dropped.
|
||||
|
||||
Verification and persistence are serialized by the collection's ingestion worker:
|
||||
@@ -504,8 +534,8 @@ Amethyst's outbox implementation combines application caches, relay hints, defau
|
||||
### R4. Source semantics
|
||||
|
||||
- **Current:** no source modes exist.
|
||||
- **Target:** all three modes implement the seed, membership, ordering, completion, and network behavior in section 4.
|
||||
- **Acceptance:** the source-mode verification cases in section 10 pass, including no local seed for `Remote` and background persistence of its persistent events.
|
||||
- **Target:** all three modes implement the seed, membership, ordering, completion, network, and max-age behavior in section 4.
|
||||
- **Acceptance:** the source-mode verification cases in section 10 pass, including no local seed for `Remote`, background persistence of its persistent events, and correct suppression/deferment of fresh local-and-remote requests.
|
||||
|
||||
### R5. Validation and persistence visibility
|
||||
|
||||
@@ -594,6 +624,19 @@ Amethyst's outbox implementation combines application caches, relay hints, defau
|
||||
- [ ] Collection after façade closure emits one `Lifecycle` failure and completes.
|
||||
- [ ] Mutating caller-owned relay/filter collections after collection starts does not alter the active local query or REQ.
|
||||
|
||||
### Query freshness
|
||||
|
||||
- [ ] Null `maxAge` preserves always-refresh local-and-remote behavior.
|
||||
- [ ] A successful all-relay EOSE caches empty and nonempty local-and-remote results.
|
||||
- [ ] A fresh one-shot emits `Complete` without subscribing.
|
||||
- [ ] A fresh stream emits `Cached`, keeps observing the store, and subscribes after expiry.
|
||||
- [ ] A concurrent successful refresh extends a deferred stream's wait.
|
||||
- [ ] Fingerprints are stable across order-equivalent filters/relays and isolate different filters or relay sets.
|
||||
- [ ] Timeout, failure, partial EOSE, cancellation, and local-only mutation do not advance freshness.
|
||||
- [ ] Freshness survives façade/process recreation and is invalidated when the Quartz database file is deleted or replaced.
|
||||
- [ ] Clock rollback, corrupt metadata, and metadata-write failure cause a remote fetch rather than a stale hit.
|
||||
- [ ] `Remote` bypasses local-and-remote freshness.
|
||||
|
||||
### Relay selection
|
||||
|
||||
- [ ] REQs go only to the specified nonempty relay set.
|
||||
@@ -633,6 +676,7 @@ Amethyst's outbox implementation combines application caches, relay hints, defau
|
||||
- Local NIP-50 FTS through Quartz filters.
|
||||
- Ordinary Nostr tag queries.
|
||||
- Deterministic compatibility, lifecycle, and invariant tests.
|
||||
- Persistent bounded max-age metadata for local-and-remote queries.
|
||||
|
||||
### Out of scope
|
||||
|
||||
@@ -690,6 +734,7 @@ Version 1 is complete only when:
|
||||
- Default source: local; every network query requires an explicit nonempty relay set.
|
||||
- Remote-only meaning: validated relay events emit directly and persistent events are also saved.
|
||||
- Database: Quartz event store is the sole canonical persistent source.
|
||||
- Query freshness: persistent hashed synchronization metadata, rotated with database recreation; event timestamps are never used as cache age.
|
||||
- Relays: exact caller-specified sets only.
|
||||
- NIP-65/outbox routing: deferred because it is Amethyst application policy, not a generic Quartz 1.12.6 query capability.
|
||||
- Dependency: published `com.vitorpamplona.quartz:quartz:1.12.6`; `reference/` remains research-only.
|
||||
|
||||
@@ -2,6 +2,7 @@ plugins {
|
||||
alias(libs.plugins.android.library) apply false
|
||||
alias(libs.plugins.android.application) apply false
|
||||
alias(libs.plugins.kotlin.android) apply false
|
||||
alias(libs.plugins.kotlin.compose) apply false
|
||||
}
|
||||
|
||||
subprojects {
|
||||
|
||||
@@ -5,6 +5,10 @@ quartz = "1.12.6"
|
||||
coroutines = "1.11.0"
|
||||
coil = "3.5.0"
|
||||
okhttp = "5.4.0"
|
||||
compose-bom = "2026.05.01"
|
||||
activity = "1.13.0"
|
||||
lifecycle = "2.11.0"
|
||||
navigation = "2.9.8"
|
||||
androidx-test = "1.7.0"
|
||||
junit = "4.13.2"
|
||||
|
||||
@@ -13,8 +17,22 @@ quartz = { module = "com.vitorpamplona.quartz:quartz", version.ref = "quartz" }
|
||||
coroutines-android = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-android", version.ref = "coroutines" }
|
||||
coroutines-test = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-test", version.ref = "coroutines" }
|
||||
coil = { module = "io.coil-kt.coil3:coil", version.ref = "coil" }
|
||||
coil-compose = { module = "io.coil-kt.coil3:coil-compose", version.ref = "coil" }
|
||||
coil-network-okhttp = { module = "io.coil-kt.coil3:coil-network-okhttp", version.ref = "coil" }
|
||||
okhttp = { module = "com.squareup.okhttp3:okhttp", version.ref = "okhttp" }
|
||||
compose-bom = { module = "androidx.compose:compose-bom", version.ref = "compose-bom" }
|
||||
compose-foundation = { module = "androidx.compose.foundation:foundation" }
|
||||
compose-material3 = { module = "androidx.compose.material3:material3" }
|
||||
compose-ui = { module = "androidx.compose.ui:ui" }
|
||||
compose-ui-tooling = { module = "androidx.compose.ui:ui-tooling" }
|
||||
compose-ui-tooling-preview = { module = "androidx.compose.ui:ui-tooling-preview" }
|
||||
compose-ui-test-junit4 = { module = "androidx.compose.ui:ui-test-junit4" }
|
||||
compose-ui-test-manifest = { module = "androidx.compose.ui:ui-test-manifest" }
|
||||
activity-compose = { module = "androidx.activity:activity-compose", version.ref = "activity" }
|
||||
lifecycle-runtime-compose = { module = "androidx.lifecycle:lifecycle-runtime-compose", version.ref = "lifecycle" }
|
||||
lifecycle-viewmodel-compose = { module = "androidx.lifecycle:lifecycle-viewmodel-compose", version.ref = "lifecycle" }
|
||||
navigation-compose = { module = "androidx.navigation:navigation-compose", version.ref = "navigation" }
|
||||
navigation-testing = { module = "androidx.navigation:navigation-testing", version.ref = "navigation" }
|
||||
junit = { module = "junit:junit", version.ref = "junit" }
|
||||
androidx-test-runner = { module = "androidx.test:runner", version.ref = "androidx-test" }
|
||||
androidx-test-junit = { module = "androidx.test.ext:junit", version = "1.3.0" }
|
||||
@@ -23,3 +41,4 @@ androidx-test-junit = { module = "androidx.test.ext:junit", version = "1.3.0" }
|
||||
android-library = { id = "com.android.library", version.ref = "agp" }
|
||||
android-application = { id = "com.android.application", version.ref = "agp" }
|
||||
kotlin-android = { id = "org.jetbrains.kotlin.android", version.ref = "kotlin" }
|
||||
kotlin-compose = { id = "org.jetbrains.kotlin.plugin.compose", version.ref = "kotlin" }
|
||||
|
||||
+33
-1
@@ -14,7 +14,9 @@ import com.vitorpamplona.quartz.nip01Core.signers.EventTemplate
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
|
||||
import com.vitorpamplona.quartz.nip01Core.store.ObservableEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import java.io.File
|
||||
import java.util.UUID
|
||||
import kotlin.time.Duration.Companion.hours
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.CoroutineStart
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
@@ -34,6 +36,36 @@ import org.junit.runner.RunWith
|
||||
@RunWith(AndroidJUnit4::class)
|
||||
@SdkSuppress(minSdkVersion = 29)
|
||||
class QuartzCompatibilityInstrumentedTest {
|
||||
@Test
|
||||
fun queryRefreshMetadataPersistsAndRotatesWithDatabaseIdentity() {
|
||||
val context = InstrumentationRegistry.getInstrumentation().targetContext
|
||||
val path = context.getDatabasePath("refresh-cache-${UUID.randomUUID()}.db").absolutePath
|
||||
val fingerprint = "profile-query"
|
||||
val refreshedAt = System.currentTimeMillis()
|
||||
val database = File(path)
|
||||
val replacement = File("$path.replacement")
|
||||
database.parentFile?.mkdirs()
|
||||
database.writeText("first database")
|
||||
|
||||
try {
|
||||
val first = SharedPreferencesQueryRefreshCache.create(context, path, databaseExisted = false)
|
||||
assertTrue(first.recordRefresh(fingerprint, refreshedAt))
|
||||
|
||||
val reopened = SharedPreferencesQueryRefreshCache.create(context, path, databaseExisted = true)
|
||||
assertEquals(refreshedAt, reopened.lastRefresh(fingerprint))
|
||||
|
||||
replacement.writeText("replacement database")
|
||||
assertTrue(database.delete())
|
||||
assertTrue(replacement.renameTo(database))
|
||||
|
||||
val recreated = SharedPreferencesQueryRefreshCache.create(context, path, databaseExisted = true)
|
||||
assertEquals(null, recreated.lastRefresh(fingerprint))
|
||||
} finally {
|
||||
database.delete()
|
||||
replacement.delete()
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun publishedStoreEmitsAfterCommitAndSurvivesReopen() = runBlocking {
|
||||
val context = InstrumentationRegistry.getInstrumentation().targetContext
|
||||
@@ -87,7 +119,7 @@ class QuartzCompatibilityInstrumentedTest {
|
||||
QuerySync.Connecting,
|
||||
purpleQuartz.query(
|
||||
Filter(kinds = listOf(1)),
|
||||
QuerySource.LocalAndRemote(setOf(relay)),
|
||||
QuerySource.LocalAndRemote(setOf(relay), maxAge = 6.hours),
|
||||
).take(1).toList().single().sync,
|
||||
)
|
||||
assertEquals(
|
||||
|
||||
+145
-18
@@ -23,6 +23,7 @@ import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.nip40Expiration.isExpired
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.CoroutineDispatcher
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.CoroutineStart
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
@@ -52,11 +53,14 @@ class PurpleQuartz private constructor(
|
||||
private val config: PurpleQuartzConfig,
|
||||
private val eventVerifier: (Event) -> Boolean,
|
||||
private val connectivityManager: ConnectivityManager?,
|
||||
private val refreshCache: QueryRefreshCache,
|
||||
private val clock: EpochMillisClock,
|
||||
) : AutoCloseable {
|
||||
private val closed = AtomicBoolean(false)
|
||||
private val networkCallbackRegistered = AtomicBoolean(false)
|
||||
private val sessions = mutableSetOf<QuerySession>()
|
||||
private val sessionsLock = Any()
|
||||
private val refreshLock = Any()
|
||||
private val closeResult = CompletableDeferred<Result<Unit>>()
|
||||
|
||||
private val networkCallback = object : ConnectivityManager.NetworkCallback() {
|
||||
@@ -204,8 +208,10 @@ class PurpleQuartz private constructor(
|
||||
private val stopped = AtomicBoolean(false)
|
||||
private val acceptingCallbacks = AtomicBoolean(true)
|
||||
private val timeoutRequested = AtomicBoolean(false)
|
||||
private val remotePrepared = AtomicBoolean(false)
|
||||
private val remoteStarted = AtomicBoolean(false)
|
||||
private val remoteSubscribed = AtomicBoolean(false)
|
||||
private val remoteDeferred = AtomicBoolean(false)
|
||||
private val stateMutex = Mutex()
|
||||
private val emitMutex = Mutex()
|
||||
private val subId = newSubId()
|
||||
@@ -219,11 +225,15 @@ class PurpleQuartz private constructor(
|
||||
is QuerySource.Remote -> source.mode
|
||||
QuerySource.Local -> null
|
||||
}
|
||||
private val relayStates = relays.associateWith {
|
||||
RelayQueryState(1, RelayConnectionState.Connecting, eose = false)
|
||||
}.toMutableMap()
|
||||
private val callbackGenerations = relays.associateWith { AtomicLong(1) }
|
||||
private val callbackGenerationStarted = relays.associateWith { AtomicBoolean(false) }
|
||||
private val maxAge = (source as? QuerySource.LocalAndRemote)?.maxAge
|
||||
private val queryFingerprint = if (source is QuerySource.LocalAndRemote) {
|
||||
runCatching { QueryFingerprint.create(filters, relays) }.getOrNull()
|
||||
} else {
|
||||
null
|
||||
}
|
||||
private val relayStates = mutableMapOf<NormalizedRelayUrl, RelayQueryState>()
|
||||
private val callbackGenerations = mutableMapOf<NormalizedRelayUrl, AtomicLong>()
|
||||
private val callbackGenerationStarted = mutableMapOf<NormalizedRelayUrl, AtomicBoolean>()
|
||||
private var items: List<Event> = emptyList()
|
||||
private var error: QueryError? = null
|
||||
@Volatile
|
||||
@@ -233,6 +243,7 @@ class PurpleQuartz private constructor(
|
||||
private var observerJob: Job? = null
|
||||
private var workerJob: Job? = null
|
||||
private var timeoutJob: Job? = null
|
||||
private var remoteDelayJob: Job? = null
|
||||
|
||||
private val subscriptionListener = object : SubscriptionListener {
|
||||
override fun onSubscriptionStarted(relay: String, forFilters: List<Filter>) {
|
||||
@@ -283,18 +294,81 @@ class PurpleQuartz private constructor(
|
||||
}
|
||||
|
||||
private fun startLocalOnly() {
|
||||
observerJob = startObserver(QuerySync.LocalOnly)
|
||||
observerJob = startObserver(initialSync = { QuerySync.LocalOnly })
|
||||
}
|
||||
|
||||
private fun startLocalAndRemote() {
|
||||
observerJob = startObserver(QuerySync.Connecting) {
|
||||
startRemote()
|
||||
observerJob = startObserver(
|
||||
initialSync = {
|
||||
val fresh = maxAge?.let(::freshness)?.isFresh == true
|
||||
remoteDeferred.set(fresh)
|
||||
if (!fresh) prepareRemote()
|
||||
when {
|
||||
!fresh -> QuerySync.Connecting
|
||||
mode is RemoteMode.OneShot -> QuerySync.Complete
|
||||
else -> QuerySync.Cached
|
||||
}
|
||||
},
|
||||
afterSeed = { seedSync ->
|
||||
when (seedSync) {
|
||||
QuerySync.Complete -> finish()
|
||||
QuerySync.Cached -> scheduleRemoteAfterCache()
|
||||
else -> startRemote()
|
||||
}
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
private fun scheduleRemoteAfterCache() {
|
||||
val cacheDuration = maxAge ?: return
|
||||
remoteDelayJob = scope.launch {
|
||||
while (!stopped.get()) {
|
||||
val status = freshness(cacheDuration)
|
||||
if (status.isFresh) {
|
||||
delay(status.remainingMillis.coerceIn(1, CACHE_RECHECK_INTERVAL_MILLIS))
|
||||
continue
|
||||
}
|
||||
var shouldStart = false
|
||||
stateMutex.withLock {
|
||||
shouldStart = synchronized(refreshLock) {
|
||||
!freshnessUnlocked(cacheDuration).isFresh &&
|
||||
remoteDeferred.compareAndSet(true, false)
|
||||
}
|
||||
if (!stopped.get() && shouldStart) {
|
||||
prepareRemote()
|
||||
emitState(QuerySync.Connecting)
|
||||
startRemote()
|
||||
}
|
||||
}
|
||||
if (shouldStart || stopped.get()) return@launch
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun freshness(cacheDuration: Duration): QueryFreshness = synchronized(refreshLock) {
|
||||
freshnessUnlocked(cacheDuration)
|
||||
}
|
||||
|
||||
private fun freshnessUnlocked(cacheDuration: Duration): QueryFreshness {
|
||||
val fingerprint = queryFingerprint
|
||||
?: return QueryFreshness(isFresh = false, remainingMillis = 0)
|
||||
return runCatching {
|
||||
refreshCache.freshness(fingerprint, cacheDuration, clock.now())
|
||||
}.getOrDefault(QueryFreshness(isFresh = false, remainingMillis = 0))
|
||||
}
|
||||
|
||||
private fun recordRefresh() {
|
||||
if (source !is QuerySource.LocalAndRemote) return
|
||||
val fingerprint = queryFingerprint ?: return
|
||||
synchronized(refreshLock) {
|
||||
runCatching { refreshCache.recordRefresh(fingerprint, clock.now()) }
|
||||
}
|
||||
}
|
||||
|
||||
private fun startRemoteOnly() {
|
||||
scope.launch {
|
||||
stateMutex.withLock {
|
||||
prepareRemote()
|
||||
emitState(QuerySync.Connecting)
|
||||
startRemote()
|
||||
}
|
||||
@@ -305,16 +379,21 @@ class PurpleQuartz private constructor(
|
||||
* Starts collecting changes before the seed query. A change that races the seed is
|
||||
* queued by the collector and causes a post-seed projection.
|
||||
*/
|
||||
private fun startObserver(initialSync: QuerySync, afterSeed: (() -> Unit)? = null): Job =
|
||||
private fun startObserver(
|
||||
initialSync: () -> QuerySync,
|
||||
afterSeed: ((QuerySync) -> Unit)? = null,
|
||||
): Job =
|
||||
scope.launch(start = CoroutineStart.UNDISPATCHED) {
|
||||
val changes = Channel<Unit>(Channel.CONFLATED)
|
||||
val changesJob = launch(start = CoroutineStart.UNDISPATCHED) {
|
||||
store.changes.collect { changes.trySend(Unit) }
|
||||
}
|
||||
try {
|
||||
stateMutex.withLock { requery(initialSync, force = true) }
|
||||
val seedSync = initialSync()
|
||||
stateMutex.withLock { requery(seedSync, force = true) }
|
||||
while (changes.tryReceive().isSuccess) stateMutex.withLock { requery(sync(), force = false) }
|
||||
afterSeed?.invoke()
|
||||
afterSeed?.invoke(seedSync)
|
||||
if (stopped.get()) return@launch
|
||||
for (ignored in changes) stateMutex.withLock { requery(sync(), force = false) }
|
||||
} catch (cancelled: CancellationException) {
|
||||
throw cancelled
|
||||
@@ -328,6 +407,7 @@ class PurpleQuartz private constructor(
|
||||
|
||||
private fun startRemote() {
|
||||
if (stopped.get() || !remoteStarted.compareAndSet(false, true)) return
|
||||
prepareRemote()
|
||||
workerJob = scope.launch {
|
||||
for (message in messages) {
|
||||
stateMutex.withLock {
|
||||
@@ -356,6 +436,15 @@ class PurpleQuartz private constructor(
|
||||
}
|
||||
}
|
||||
|
||||
private fun prepareRemote() {
|
||||
if (!remotePrepared.compareAndSet(false, true)) return
|
||||
relays.forEach { relay ->
|
||||
relayStates[relay] = RelayQueryState(1, RelayConnectionState.Connecting, eose = false)
|
||||
callbackGenerations[relay] = AtomicLong(1)
|
||||
callbackGenerationStarted[relay] = AtomicBoolean(false)
|
||||
}
|
||||
}
|
||||
|
||||
private fun enqueue(message: Inbound) {
|
||||
if (message.relay !in relays || stopped.get() || !acceptingCallbacks.get()) return
|
||||
val result = messages.trySend(message)
|
||||
@@ -401,11 +490,29 @@ class PurpleQuartz private constructor(
|
||||
is Inbound.EventReceived -> ingest(message.relay, message.event)
|
||||
is Inbound.Eose -> {
|
||||
updateRelay(message.relay) { it.copy(eose = true) }
|
||||
if (timeoutRequested.get()) return
|
||||
error = null
|
||||
if (!timeoutRequested.get() && mode is RemoteMode.OneShot && relayStates.values.all { it.eose }) {
|
||||
if (source is QuerySource.LocalAndRemote) requery(QuerySync.Complete, force = true)
|
||||
else emitState(QuerySync.Complete)
|
||||
finish()
|
||||
val allRelaysCaughtUp = relayStates.values.all { it.eose }
|
||||
if (allRelaysCaughtUp) {
|
||||
if (mode is RemoteMode.OneShot) {
|
||||
if (source is QuerySource.LocalAndRemote) {
|
||||
try {
|
||||
requery(QuerySync.Complete, force = true)
|
||||
} catch (cancelled: CancellationException) {
|
||||
throw cancelled
|
||||
} catch (_: Throwable) {
|
||||
fail(QueryError.UnsupportedLocalProjection("Local store projection failed"))
|
||||
return
|
||||
}
|
||||
recordRefresh()
|
||||
} else {
|
||||
emitState(QuerySync.Complete)
|
||||
}
|
||||
finish()
|
||||
} else {
|
||||
recordRefresh()
|
||||
emitState(sync())
|
||||
}
|
||||
} else {
|
||||
emitState(sync())
|
||||
}
|
||||
@@ -469,6 +576,7 @@ class PurpleQuartz private constructor(
|
||||
}
|
||||
|
||||
private fun sync(): QuerySync = when {
|
||||
remoteDeferred.get() -> QuerySync.Cached
|
||||
relays.isEmpty() -> QuerySync.LocalOnly
|
||||
relayStates.values.any { it.connection != RelayConnectionState.Connected } -> QuerySync.Connecting
|
||||
relayStates.values.any { !it.eose } -> QuerySync.CatchingUp
|
||||
@@ -479,7 +587,8 @@ class PurpleQuartz private constructor(
|
||||
private suspend fun emitState(sync: QuerySync) {
|
||||
emitMutex.withLock {
|
||||
if (!stopped.get()) {
|
||||
val state = QueryState(items, sync, relayStates.toMap(), error)
|
||||
val visibleRelays = if (remoteDeferred.get()) emptyMap() else relayStates.toMap()
|
||||
val state = QueryState(items, sync, visibleRelays, error)
|
||||
lastState = state
|
||||
producer.send(state)
|
||||
}
|
||||
@@ -553,6 +662,7 @@ class PurpleQuartz private constructor(
|
||||
private fun cleanup() {
|
||||
acceptingCallbacks.set(false)
|
||||
timeoutJob?.cancel()
|
||||
remoteDelayJob?.cancel()
|
||||
observerJob?.cancel()
|
||||
messages.close()
|
||||
workerJob?.cancel()
|
||||
@@ -578,7 +688,9 @@ class PurpleQuartz private constructor(
|
||||
config.validate()
|
||||
val parentJob = requireActiveParentJob(parentScope)
|
||||
val appContext = context.applicationContext
|
||||
val path = appContext.getDatabasePath(config.databaseName).absoluteFile.canonicalPath
|
||||
val databaseFile = appContext.getDatabasePath(config.databaseName).absoluteFile
|
||||
val databaseExisted = databaseFile.exists()
|
||||
val path = databaseFile.canonicalPath
|
||||
synchronized(openDatabases) {
|
||||
check(path !in openDatabases) { "A PurpleQuartz instance already owns this database" }
|
||||
openDatabases += path
|
||||
@@ -589,6 +701,11 @@ class PurpleQuartz private constructor(
|
||||
var client: INostrClient? = null
|
||||
try {
|
||||
store = ObservableEventStore(EventStore(dbName = path, relay = null))
|
||||
val refreshCache = SharedPreferencesQueryRefreshCache.create(
|
||||
context = appContext,
|
||||
databasePath = path,
|
||||
databaseExisted = databaseExisted,
|
||||
)
|
||||
client = NostrClient(websocketBuilder, scope)
|
||||
return PurpleQuartz(
|
||||
store,
|
||||
@@ -599,6 +716,8 @@ class PurpleQuartz private constructor(
|
||||
config,
|
||||
DEFAULT_EVENT_VERIFIER,
|
||||
appContext.getSystemService(Context.CONNECTIVITY_SERVICE) as? ConnectivityManager,
|
||||
refreshCache,
|
||||
SYSTEM_CLOCK,
|
||||
)
|
||||
} catch (failure: Throwable) {
|
||||
runCatching { client?.close() }.exceptionOrNull()?.let(failure::addSuppressed)
|
||||
@@ -621,10 +740,13 @@ class PurpleQuartz private constructor(
|
||||
config: PurpleQuartzConfig = PurpleQuartzConfig(),
|
||||
databasePath: String = "test-${System.nanoTime()}",
|
||||
eventVerifier: (Event) -> Boolean = DEFAULT_EVENT_VERIFIER,
|
||||
refreshCache: QueryRefreshCache = InMemoryQueryRefreshCache(),
|
||||
clock: EpochMillisClock = SYSTEM_CLOCK,
|
||||
dispatcher: CoroutineDispatcher = Dispatchers.IO,
|
||||
): PurpleQuartz {
|
||||
config.validate()
|
||||
val job = SupervisorJob(requireActiveParentJob(parentScope))
|
||||
val scope = CoroutineScope(parentScope.coroutineContext + job + Dispatchers.IO)
|
||||
val scope = CoroutineScope(parentScope.coroutineContext + job + dispatcher)
|
||||
return PurpleQuartz(
|
||||
store = ObservableEventStore(eventStore),
|
||||
client = client,
|
||||
@@ -634,6 +756,8 @@ class PurpleQuartz private constructor(
|
||||
config = config,
|
||||
eventVerifier = eventVerifier,
|
||||
connectivityManager = null,
|
||||
refreshCache = refreshCache,
|
||||
clock = clock,
|
||||
)
|
||||
}
|
||||
|
||||
@@ -662,5 +786,8 @@ class PurpleQuartz private constructor(
|
||||
private val DEFAULT_EVENT_VERIFIER: (Event) -> Boolean = { event ->
|
||||
event.verifyId() && event.verifySignature()
|
||||
}
|
||||
|
||||
private val SYSTEM_CLOCK = EpochMillisClock(System::currentTimeMillis)
|
||||
private const val CACHE_RECHECK_INTERVAL_MILLIS = 60_000L
|
||||
}
|
||||
}
|
||||
|
||||
+113
@@ -0,0 +1,113 @@
|
||||
package dev.zapstore.purplequartz
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import java.io.ByteArrayOutputStream
|
||||
import java.io.DataOutputStream
|
||||
import java.security.MessageDigest
|
||||
|
||||
internal object QueryFingerprint {
|
||||
private const val FORMAT_VERSION = 1
|
||||
|
||||
fun create(
|
||||
filters: List<Filter>,
|
||||
relays: Set<NormalizedRelayUrl>,
|
||||
): String {
|
||||
val encodedFilters = filters.map(::encodeFilter).sortedWith(::compareBytes)
|
||||
val canonical = ByteArrayOutputStream().use { bytes ->
|
||||
DataOutputStream(bytes).use { output ->
|
||||
output.writeInt(FORMAT_VERSION)
|
||||
output.writeInt(encodedFilters.size)
|
||||
encodedFilters.forEach { filter ->
|
||||
output.writeInt(filter.size)
|
||||
output.write(filter)
|
||||
}
|
||||
output.writeStrings(relays.map { it.url }.sorted())
|
||||
}
|
||||
bytes.toByteArray()
|
||||
}
|
||||
return MessageDigest.getInstance("SHA-256")
|
||||
.digest(canonical)
|
||||
.joinToString(separator = "") { byte -> "%02x".format(byte.toInt() and 0xff) }
|
||||
}
|
||||
|
||||
private fun encodeFilter(filter: Filter): ByteArray =
|
||||
ByteArrayOutputStream().use { bytes ->
|
||||
DataOutputStream(bytes).use { output ->
|
||||
output.writeNullableStrings(filter.ids)
|
||||
output.writeNullableStrings(filter.authors)
|
||||
output.writeNullableInts(filter.kinds)
|
||||
output.writeNullableStringMap(filter.tags)
|
||||
output.writeNullableStringMap(filter.tagsAll)
|
||||
output.writeNullableLong(filter.since)
|
||||
output.writeNullableLong(filter.until)
|
||||
output.writeNullableInt(filter.limit)
|
||||
output.writeNullableString(filter.search)
|
||||
}
|
||||
bytes.toByteArray()
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeNullableStrings(values: List<String>?) {
|
||||
if (values == null) {
|
||||
writeInt(-1)
|
||||
} else {
|
||||
writeStrings(values.sorted())
|
||||
}
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeStrings(values: List<String>) {
|
||||
writeInt(values.size)
|
||||
values.forEach { value -> writeString(value) }
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeNullableInts(values: List<Int>?) {
|
||||
if (values == null) {
|
||||
writeInt(-1)
|
||||
} else {
|
||||
writeInt(values.size)
|
||||
values.sorted().forEach(::writeInt)
|
||||
}
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeNullableStringMap(values: Map<String, List<String>>?) {
|
||||
if (values == null) {
|
||||
writeInt(-1)
|
||||
} else {
|
||||
writeInt(values.size)
|
||||
values.toSortedMap().forEach { (key, entries) ->
|
||||
writeString(key)
|
||||
writeStrings(entries.sorted())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeNullableLong(value: Long?) {
|
||||
writeBoolean(value != null)
|
||||
if (value != null) writeLong(value)
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeNullableInt(value: Int?) {
|
||||
writeBoolean(value != null)
|
||||
if (value != null) writeInt(value)
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeNullableString(value: String?) {
|
||||
writeBoolean(value != null)
|
||||
if (value != null) writeString(value)
|
||||
}
|
||||
|
||||
private fun DataOutputStream.writeString(value: String) {
|
||||
val encoded = value.toByteArray(Charsets.UTF_8)
|
||||
writeInt(encoded.size)
|
||||
write(encoded)
|
||||
}
|
||||
|
||||
private fun compareBytes(left: ByteArray, right: ByteArray): Int {
|
||||
val commonLength = minOf(left.size, right.size)
|
||||
for (index in 0 until commonLength) {
|
||||
val comparison = (left[index].toInt() and 0xff).compareTo(right[index].toInt() and 0xff)
|
||||
if (comparison != 0) return comparison
|
||||
}
|
||||
return left.size.compareTo(right.size)
|
||||
}
|
||||
}
|
||||
+149
@@ -0,0 +1,149 @@
|
||||
package dev.zapstore.purplequartz
|
||||
|
||||
import android.content.Context
|
||||
import android.content.SharedPreferences
|
||||
import java.io.File
|
||||
import java.nio.file.Files
|
||||
import java.nio.file.attribute.BasicFileAttributes
|
||||
import java.security.MessageDigest
|
||||
import java.util.UUID
|
||||
import kotlin.time.Duration
|
||||
|
||||
internal fun interface EpochMillisClock {
|
||||
fun now(): Long
|
||||
}
|
||||
|
||||
internal interface QueryRefreshCache {
|
||||
fun lastRefresh(fingerprint: String): Long?
|
||||
|
||||
fun recordRefresh(fingerprint: String, epochMillis: Long): Boolean
|
||||
}
|
||||
|
||||
internal data class QueryFreshness(
|
||||
val isFresh: Boolean,
|
||||
val remainingMillis: Long,
|
||||
)
|
||||
|
||||
internal fun QueryRefreshCache.freshness(
|
||||
fingerprint: String,
|
||||
maxAge: Duration,
|
||||
now: Long,
|
||||
): QueryFreshness {
|
||||
val refreshedAt = lastRefresh(fingerprint)
|
||||
?: return QueryFreshness(isFresh = false, remainingMillis = 0)
|
||||
if (refreshedAt <= 0 || now < refreshedAt) {
|
||||
return QueryFreshness(isFresh = false, remainingMillis = 0)
|
||||
}
|
||||
|
||||
val maxAgeMillis = maxOf(1, maxAge.inWholeMilliseconds)
|
||||
val age = now - refreshedAt
|
||||
val remaining = maxAgeMillis - age
|
||||
return QueryFreshness(
|
||||
isFresh = remaining > 0,
|
||||
remainingMillis = remaining.coerceAtLeast(0),
|
||||
)
|
||||
}
|
||||
|
||||
internal class InMemoryQueryRefreshCache : QueryRefreshCache {
|
||||
private val entries = mutableMapOf<String, Long>()
|
||||
|
||||
override fun lastRefresh(fingerprint: String): Long? = synchronized(entries) {
|
||||
entries[fingerprint]
|
||||
}
|
||||
|
||||
override fun recordRefresh(fingerprint: String, epochMillis: Long): Boolean = synchronized(entries) {
|
||||
if (epochMillis <= 0) return@synchronized false
|
||||
entries[fingerprint] = epochMillis
|
||||
true
|
||||
}
|
||||
}
|
||||
|
||||
internal class SharedPreferencesQueryRefreshCache private constructor(
|
||||
private val preferences: SharedPreferences,
|
||||
private val entryPrefix: String,
|
||||
) : QueryRefreshCache {
|
||||
private val lock = Any()
|
||||
|
||||
override fun lastRefresh(fingerprint: String): Long? = runCatching {
|
||||
preferences.takeIf { it.contains(entryPrefix + fingerprint) }
|
||||
?.getLong(entryPrefix + fingerprint, 0L)
|
||||
?.takeIf { it > 0 }
|
||||
}.getOrNull()
|
||||
|
||||
override fun recordRefresh(fingerprint: String, epochMillis: Long): Boolean {
|
||||
if (epochMillis <= 0) return false
|
||||
synchronized(lock) {
|
||||
val key = entryPrefix + fingerprint
|
||||
val currentEntries = preferences.all.asSequence()
|
||||
.filter { (entryKey, value) -> entryKey.startsWith(entryPrefix) && value is Long }
|
||||
.sortedBy { (_, value) -> value as Long }
|
||||
.toList()
|
||||
val removals = if (preferences.contains(key)) {
|
||||
0
|
||||
} else {
|
||||
(currentEntries.size - MAX_ENTRIES + 1).coerceAtLeast(0)
|
||||
}
|
||||
val editor = preferences.edit()
|
||||
currentEntries.take(removals).forEach { (entryKey, _) -> editor.remove(entryKey) }
|
||||
val committed = editor.putLong(key, epochMillis).commit()
|
||||
if (!committed) {
|
||||
// commit updates SharedPreferences memory before attempting disk I/O.
|
||||
// Remove the optimistic value so this process also fails open.
|
||||
preferences.edit().remove(key).apply()
|
||||
}
|
||||
return committed
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val PREFERENCES_NAME = "dev.zapstore.purplequartz.query-refresh-v1"
|
||||
private const val MAX_ENTRIES = 1_024
|
||||
|
||||
fun create(
|
||||
context: Context,
|
||||
databasePath: String,
|
||||
databaseExisted: Boolean = File(databasePath).exists(),
|
||||
): SharedPreferencesQueryRefreshCache {
|
||||
val preferences = context.getSharedPreferences(PREFERENCES_NAME, Context.MODE_PRIVATE)
|
||||
val pathHash = sha256(databasePath)
|
||||
val generationKey = "database:$pathHash:generation"
|
||||
val identityKey = "database:$pathHash:identity"
|
||||
val previousGeneration = runCatching { preferences.getString(generationKey, null) }.getOrNull()
|
||||
val previousIdentity = runCatching { preferences.getString(identityKey, null) }.getOrNull()
|
||||
val currentIdentity = databaseIdentity(databasePath)
|
||||
val canReuseGeneration =
|
||||
databaseExisted &&
|
||||
currentIdentity != null &&
|
||||
currentIdentity == previousIdentity &&
|
||||
previousGeneration != null
|
||||
val generation = previousGeneration?.takeIf { canReuseGeneration } ?: UUID.randomUUID().toString()
|
||||
|
||||
if (!canReuseGeneration) {
|
||||
val oldPrefix = "refresh:$pathHash:"
|
||||
val editor = preferences.edit()
|
||||
.putString(generationKey, generation)
|
||||
.putString(identityKey, currentIdentity)
|
||||
preferences.all.keys.filter { it.startsWith(oldPrefix) }.forEach(editor::remove)
|
||||
editor.commit()
|
||||
}
|
||||
|
||||
return SharedPreferencesQueryRefreshCache(
|
||||
preferences = preferences,
|
||||
entryPrefix = "refresh:$pathHash:$generation:",
|
||||
)
|
||||
}
|
||||
|
||||
private fun sha256(value: String): String =
|
||||
MessageDigest.getInstance("SHA-256")
|
||||
.digest(value.toByteArray(Charsets.UTF_8))
|
||||
.joinToString(separator = "") { byte -> "%02x".format(byte.toInt() and 0xff) }
|
||||
|
||||
private fun databaseIdentity(databasePath: String): String? = runCatching {
|
||||
val attributes = Files.readAttributes(File(databasePath).toPath(), BasicFileAttributes::class.java)
|
||||
val fileKey = attributes.fileKey()?.toString().orEmpty()
|
||||
val createdAt = attributes.creationTime().toMillis()
|
||||
if (fileKey.isEmpty() && createdAt <= 0) return@runCatching null
|
||||
sha256("$fileKey:$createdAt")
|
||||
}.getOrNull()
|
||||
}
|
||||
}
|
||||
@@ -10,9 +10,13 @@ sealed interface QuerySource {
|
||||
data class LocalAndRemote(
|
||||
val relays: Set<NormalizedRelayUrl>,
|
||||
val mode: RemoteMode = RemoteMode.Stream,
|
||||
val maxAge: Duration? = null,
|
||||
) : QuerySource {
|
||||
init {
|
||||
require(relays.isNotEmpty()) { "At least one relay is required" }
|
||||
require(maxAge == null || (maxAge.isFinite() && maxAge.isPositive())) {
|
||||
"Cache max age must be finite and positive"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -49,6 +53,7 @@ data class QueryState(
|
||||
|
||||
sealed interface QuerySync {
|
||||
data object LocalOnly : QuerySync
|
||||
data object Cached : QuerySync
|
||||
data object Connecting : QuerySync
|
||||
data object CatchingUp : QuerySync
|
||||
data object Live : QuerySync
|
||||
|
||||
+5
-1
@@ -12,7 +12,9 @@ import com.vitorpamplona.quartz.nip01Core.store.ObservableEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.nip40Expiration.isExpired
|
||||
import java.lang.reflect.Modifier
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.hours
|
||||
import org.junit.Assert.assertTrue
|
||||
import org.junit.Test
|
||||
|
||||
@@ -32,8 +34,10 @@ class PublicContractsTest {
|
||||
assertFails { RemoteMode.OneShot(0.milliseconds) }
|
||||
|
||||
val relay = "wss://relay.example".normalizeRelayUrl()
|
||||
assertFails { QuerySource.LocalAndRemote(setOf(relay), maxAge = 0.milliseconds) }
|
||||
assertFails { QuerySource.LocalAndRemote(setOf(relay), maxAge = Duration.INFINITE) }
|
||||
assertTrue(QuerySource.Remote(setOf(relay)).relays.contains(relay))
|
||||
assertTrue(QuerySource.LocalAndRemote(setOf(relay)).relays.contains(relay))
|
||||
assertTrue(QuerySource.LocalAndRemote(setOf(relay), maxAge = 6.hours).relays.contains(relay))
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+214
-10
@@ -18,6 +18,7 @@ import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.collect
|
||||
import kotlinx.coroutines.flow.toList
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
@@ -26,12 +27,189 @@ import org.junit.Assert.assertFalse
|
||||
import org.junit.Assert.assertTrue
|
||||
import org.junit.Test
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.hours
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
class QueryInvariantTest {
|
||||
private val relay = "wss://relay.example".normalizeRelayUrl()
|
||||
private val eventSequence = AtomicInteger()
|
||||
|
||||
@Test
|
||||
fun `successful empty one-shot is cached across facades and expires`() = runBlocking {
|
||||
val cache = InMemoryQueryRefreshCache()
|
||||
val clock = MutableEpochClock(1_000)
|
||||
val filter = Filter(kinds = listOf(1))
|
||||
val source = QuerySource.LocalAndRemote(
|
||||
relays = setOf(relay),
|
||||
mode = RemoteMode.OneShot(5.seconds),
|
||||
maxAge = 6.hours,
|
||||
)
|
||||
val firstClient = ControlledClient()
|
||||
val first = PurpleQuartz.createForTesting(
|
||||
ControlledEventStore(),
|
||||
firstClient,
|
||||
this,
|
||||
databasePath = "persistent-cache-test",
|
||||
eventVerifier = { true },
|
||||
refreshCache = cache,
|
||||
clock = clock,
|
||||
)
|
||||
val firstStates = CopyOnWriteArrayList<QueryState>()
|
||||
val firstCollection = launch { first.query(filter, source).collect(firstStates::add) }
|
||||
|
||||
firstClient.awaitSubscription()
|
||||
firstClient.started(relay)
|
||||
firstClient.eose(relay)
|
||||
awaitState(firstStates) { it.sync == QuerySync.Complete }
|
||||
withTimeout(2.seconds) { firstCollection.join() }
|
||||
first.close()
|
||||
|
||||
val secondClient = ControlledClient()
|
||||
val second = PurpleQuartz.createForTesting(
|
||||
ControlledEventStore(),
|
||||
secondClient,
|
||||
this,
|
||||
databasePath = "persistent-cache-test",
|
||||
eventVerifier = { true },
|
||||
refreshCache = cache,
|
||||
clock = clock,
|
||||
)
|
||||
val cachedStates = second.query(filter, source).toList()
|
||||
|
||||
assertEquals(listOf(QuerySync.Complete), cachedStates.map(QueryState::sync))
|
||||
assertTrue(cachedStates.single().relays.isEmpty())
|
||||
assertTrue(secondClient.requests.isEmpty())
|
||||
|
||||
clock.advance(6.hours.inWholeMilliseconds)
|
||||
val staleCollection = launch { second.query(filter, source).collect { } }
|
||||
secondClient.awaitSubscription()
|
||||
assertEquals(1, secondClient.requests.size)
|
||||
|
||||
staleCollection.cancel()
|
||||
staleCollection.join()
|
||||
second.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `fresh stream defers remote and observes an extended refresh window`() = runBlocking {
|
||||
val cache = InMemoryQueryRefreshCache()
|
||||
val filter = Filter(kinds = listOf(1))
|
||||
val source = QuerySource.LocalAndRemote(
|
||||
relays = setOf(relay),
|
||||
mode = RemoteMode.Stream,
|
||||
maxAge = 1.seconds,
|
||||
)
|
||||
val fingerprint = QueryFingerprint.create(listOf(filter), setOf(relay))
|
||||
cache.recordRefresh(fingerprint, System.currentTimeMillis())
|
||||
val client = ControlledClient()
|
||||
val purpleQuartz = PurpleQuartz.createForTesting(
|
||||
ControlledEventStore(),
|
||||
client,
|
||||
this,
|
||||
eventVerifier = { true },
|
||||
refreshCache = cache,
|
||||
)
|
||||
val states = CopyOnWriteArrayList<QueryState>()
|
||||
val collection = collect(purpleQuartz, source, states)
|
||||
|
||||
val cached = awaitState(states) { it.sync == QuerySync.Cached }
|
||||
assertTrue(cached.relays.isEmpty())
|
||||
assertTrue(client.requests.isEmpty())
|
||||
|
||||
cache.recordRefresh(fingerprint, System.currentTimeMillis())
|
||||
delay(700)
|
||||
cache.recordRefresh(fingerprint, System.currentTimeMillis())
|
||||
delay(500)
|
||||
assertTrue(client.requests.isEmpty())
|
||||
|
||||
client.awaitSubscription()
|
||||
assertEquals(1, client.requests.size)
|
||||
|
||||
collection.cancel()
|
||||
collection.join()
|
||||
purpleQuartz.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `timeout and persistence failure do not populate freshness`() = runBlocking {
|
||||
val filter = Filter(kinds = listOf(1))
|
||||
val source = QuerySource.LocalAndRemote(
|
||||
relays = setOf(relay),
|
||||
mode = RemoteMode.OneShot(50.milliseconds),
|
||||
maxAge = 1.hours,
|
||||
)
|
||||
val fingerprint = QueryFingerprint.create(listOf(filter), setOf(relay))
|
||||
|
||||
val timeoutCache = InMemoryQueryRefreshCache()
|
||||
val timeoutClient = ControlledClient()
|
||||
val timeoutQuartz = PurpleQuartz.createForTesting(
|
||||
ControlledEventStore(),
|
||||
timeoutClient,
|
||||
this,
|
||||
config = PurpleQuartzConfig(oneShotTimeout = 50.milliseconds),
|
||||
eventVerifier = { true },
|
||||
refreshCache = timeoutCache,
|
||||
)
|
||||
val timeoutStates = CopyOnWriteArrayList<QueryState>()
|
||||
val timeoutCollection = collect(timeoutQuartz, source, timeoutStates)
|
||||
timeoutClient.awaitSubscription()
|
||||
awaitState(timeoutStates) { it.sync == QuerySync.TimedOut }
|
||||
withTimeout(2.seconds) { timeoutCollection.join() }
|
||||
assertEquals(null, timeoutCache.lastRefresh(fingerprint))
|
||||
timeoutQuartz.close()
|
||||
|
||||
val failureCache = InMemoryQueryRefreshCache()
|
||||
val failureClient = ControlledClient()
|
||||
val failureQuartz = PurpleQuartz.createForTesting(
|
||||
ControlledEventStore(failInserts = true),
|
||||
failureClient,
|
||||
this,
|
||||
eventVerifier = { true },
|
||||
refreshCache = failureCache,
|
||||
)
|
||||
val failureStates = CopyOnWriteArrayList<QueryState>()
|
||||
val failureCollection = collect(failureQuartz, source, failureStates)
|
||||
failureClient.awaitSubscription()
|
||||
failureClient.started(relay)
|
||||
failureClient.event(relay, signedEvent(content = "failed-refresh"))
|
||||
awaitState(failureStates) { it.sync == QuerySync.Failed }
|
||||
withTimeout(2.seconds) { failureCollection.join() }
|
||||
assertEquals(null, failureCache.lastRefresh(fingerprint))
|
||||
failureQuartz.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `failed final local projection fails query without caching`() = runBlocking {
|
||||
val cache = InMemoryQueryRefreshCache()
|
||||
val client = ControlledClient()
|
||||
val filter = Filter(kinds = listOf(1))
|
||||
val source = QuerySource.LocalAndRemote(
|
||||
relays = setOf(relay),
|
||||
mode = RemoteMode.OneShot(5.seconds),
|
||||
maxAge = 1.hours,
|
||||
)
|
||||
val fingerprint = QueryFingerprint.create(listOf(filter), setOf(relay))
|
||||
val purpleQuartz = PurpleQuartz.createForTesting(
|
||||
ControlledEventStore(failQueriesAfter = 1),
|
||||
client,
|
||||
this,
|
||||
eventVerifier = { true },
|
||||
refreshCache = cache,
|
||||
)
|
||||
val states = CopyOnWriteArrayList<QueryState>()
|
||||
val collection = collect(purpleQuartz, source, states)
|
||||
|
||||
client.awaitSubscription()
|
||||
client.started(relay)
|
||||
client.eose(relay)
|
||||
|
||||
val failed = awaitState(states) { it.sync == QuerySync.Failed }
|
||||
assertTrue(failed.error is QueryError.UnsupportedLocalProjection)
|
||||
assertEquals(null, cache.lastRefresh(fingerprint))
|
||||
withTimeout(2.seconds) { collection.join() }
|
||||
purpleQuartz.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `reconnect advances generation and accepts replacement EOSE`() = runBlocking {
|
||||
val client = ControlledClient()
|
||||
@@ -128,16 +306,19 @@ class QueryInvariantTest {
|
||||
client.started(relay)
|
||||
client.event(relay, event)
|
||||
withTimeout(2.seconds) { store.insertionStarted.await() }
|
||||
client.eose(relay)
|
||||
withTimeout(2.seconds) {
|
||||
while (client.unsubscribeCount.get() == 0) delay(10)
|
||||
}
|
||||
|
||||
assertFalse(states.any { it.sync == QuerySync.TimedOut })
|
||||
assertFalse(states.any { it.sync == QuerySync.Complete })
|
||||
assertTrue(states.any { state -> state.items.any { it.id == event.id } })
|
||||
|
||||
store.releaseInsert.complete(Unit)
|
||||
val timedOut = awaitState(states) { it.sync == QuerySync.TimedOut }
|
||||
assertTrue(timedOut.items.any { it.id == event.id })
|
||||
assertFalse(states.any { it.sync == QuerySync.Complete })
|
||||
withTimeout(2.seconds) { collection.join() }
|
||||
purpleQuartz.close()
|
||||
}
|
||||
@@ -298,26 +479,26 @@ private class ControlledClient(
|
||||
closeCount.incrementAndGet()
|
||||
}
|
||||
|
||||
suspend fun awaitSubscription() {
|
||||
suspend fun awaitSubscription(expectedCount: Int = 1) {
|
||||
withTimeout(2.seconds) {
|
||||
while (subscriptionListener == null) delay(10)
|
||||
while (subscriptionListener == null || requests.size < expectedCount) delay(10)
|
||||
}
|
||||
}
|
||||
|
||||
fun started(relay: NormalizedRelayUrl) {
|
||||
subscriptionListener?.onSubscriptionStarted(relay.url, requests.single().getValue(relay))
|
||||
subscriptionListener?.onSubscriptionStarted(relay.url, requests.last().getValue(relay))
|
||||
}
|
||||
|
||||
fun event(relay: NormalizedRelayUrl, event: Event) {
|
||||
subscriptionListener?.onEvent(event, false, relay, requests.single().getValue(relay))
|
||||
subscriptionListener?.onEvent(event, false, relay, requests.last().getValue(relay))
|
||||
}
|
||||
|
||||
fun eose(relay: NormalizedRelayUrl) {
|
||||
subscriptionListener?.onEose(relay, requests.single().getValue(relay))
|
||||
subscriptionListener?.onEose(relay, requests.last().getValue(relay))
|
||||
}
|
||||
|
||||
fun closed(relay: NormalizedRelayUrl) {
|
||||
subscriptionListener?.onClosed("controlled close", relay, requests.single().getValue(relay))
|
||||
subscriptionListener?.onClosed("controlled close", relay, requests.last().getValue(relay))
|
||||
}
|
||||
|
||||
fun connecting(relay: NormalizedRelayUrl) {
|
||||
@@ -331,6 +512,16 @@ private class ControlledClient(
|
||||
}
|
||||
}
|
||||
|
||||
private class MutableEpochClock(
|
||||
private var current: Long,
|
||||
) : EpochMillisClock {
|
||||
override fun now(): Long = current
|
||||
|
||||
fun advance(milliseconds: Long) {
|
||||
current += milliseconds
|
||||
}
|
||||
}
|
||||
|
||||
private class TestRelayClient(
|
||||
override val url: NormalizedRelayUrl,
|
||||
) : IRelayClient {
|
||||
@@ -346,11 +537,13 @@ private class TestRelayClient(
|
||||
private class ControlledEventStore(
|
||||
private val blockInserts: Boolean = false,
|
||||
private val failInserts: Boolean = false,
|
||||
private val failQueriesAfter: Int? = null,
|
||||
) : IEventStore {
|
||||
override val relay: NormalizedRelayUrl? = null
|
||||
val insertionStarted = CompletableDeferred<Unit>()
|
||||
val releaseInsert = CompletableDeferred<Unit>()
|
||||
val closeCount = AtomicInteger()
|
||||
private val queryCount = AtomicInteger()
|
||||
private val events = mutableListOf<Event>()
|
||||
|
||||
override suspend fun insert(event: Event) {
|
||||
@@ -373,14 +566,18 @@ private class ControlledEventStore(
|
||||
}
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
override suspend fun <T : Event> query(filter: Filter): List<T> =
|
||||
synchronized(events) { events.filter(filter::match).map { it as T } }
|
||||
override suspend fun <T : Event> query(filter: Filter): List<T> {
|
||||
beforeQuery()
|
||||
return synchronized(events) { events.filter(filter::match).map { it as T } }
|
||||
}
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
override suspend fun <T : Event> query(filters: List<Filter>): List<T> =
|
||||
synchronized(events) {
|
||||
override suspend fun <T : Event> query(filters: List<Filter>): List<T> {
|
||||
beforeQuery()
|
||||
return synchronized(events) {
|
||||
events.filter { event -> filters.any { it.match(event) } }.map { it as T }
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun <T : Event> query(filter: Filter, onEach: (T) -> Unit) {
|
||||
query<T>(filter).forEach(onEach)
|
||||
@@ -414,4 +611,11 @@ private class ControlledEventStore(
|
||||
override fun close() {
|
||||
closeCount.incrementAndGet()
|
||||
}
|
||||
|
||||
private fun beforeQuery() {
|
||||
val allowed = failQueriesAfter ?: return
|
||||
if (queryCount.incrementAndGet() > allowed) {
|
||||
throw IllegalStateException("controlled query failure")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+72
@@ -0,0 +1,72 @@
|
||||
package dev.zapstore.purplequartz
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertFalse
|
||||
import org.junit.Assert.assertNotEquals
|
||||
import org.junit.Assert.assertTrue
|
||||
import org.junit.Test
|
||||
|
||||
class QueryRefreshCacheTest {
|
||||
@Test
|
||||
fun `fingerprint canonicalizes filter and relay ordering`() {
|
||||
val relayOne = "wss://one.example".normalizeRelayUrl()
|
||||
val relayTwo = "wss://two.example".normalizeRelayUrl()
|
||||
val idOne = "a".repeat(64)
|
||||
val idTwo = "b".repeat(64)
|
||||
val first = listOf(
|
||||
Filter(
|
||||
ids = listOf(idOne, idTwo),
|
||||
kinds = listOf(1, 3),
|
||||
tags = linkedMapOf("t" to listOf("one", "two"), "x" to listOf("value")),
|
||||
),
|
||||
Filter(search = "profile", limit = 1),
|
||||
)
|
||||
val reordered = listOf(
|
||||
Filter(search = "profile", limit = 1),
|
||||
Filter(
|
||||
ids = listOf(idOne, idTwo),
|
||||
kinds = listOf(1, 3),
|
||||
tags = linkedMapOf("x" to listOf("value"), "t" to listOf("one", "two")),
|
||||
),
|
||||
)
|
||||
|
||||
assertEquals(
|
||||
QueryFingerprint.create(first, linkedSetOf(relayOne, relayTwo)),
|
||||
QueryFingerprint.create(reordered, linkedSetOf(relayTwo, relayOne)),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `fingerprint isolates filters and relay sets`() {
|
||||
val relayOne = "wss://one.example".normalizeRelayUrl()
|
||||
val relayTwo = "wss://two.example".normalizeRelayUrl()
|
||||
val authorOne = "a".repeat(64)
|
||||
val authorTwo = "b".repeat(64)
|
||||
val profile = listOf(Filter(authors = listOf(authorOne), kinds = listOf(0)))
|
||||
|
||||
val baseline = QueryFingerprint.create(profile, setOf(relayOne))
|
||||
|
||||
assertNotEquals(
|
||||
baseline,
|
||||
QueryFingerprint.create(listOf(Filter(authors = listOf(authorTwo), kinds = listOf(0))), setOf(relayOne)),
|
||||
)
|
||||
assertNotEquals(baseline, QueryFingerprint.create(profile, setOf(relayTwo)))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `freshness expires and treats a backwards clock as stale`() {
|
||||
val cache = InMemoryQueryRefreshCache()
|
||||
assertTrue(cache.recordRefresh("profile", 1_000))
|
||||
|
||||
val fresh = cache.freshness("profile", 60.seconds, 31_000)
|
||||
assertTrue(fresh.isFresh)
|
||||
assertEquals(30_000, fresh.remainingMillis)
|
||||
|
||||
assertFalse(cache.freshness("profile", 60.seconds, 61_000).isFresh)
|
||||
assertFalse(cache.freshness("profile", 60.seconds, 999).isFresh)
|
||||
assertFalse(cache.freshness("missing", 60.seconds, 1_000).isFresh)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user