Author SHA1 Message Date
jcoffey-dev 22d8ad8572 Publish amd64 first, then arm64, on a builder that keeps its cache
ci / fork-checks (pull_request) Successful in 49s
ci / build (pull_request) Successful in 4m30s
Two release builds side by side on one machine each take twice as long,
and production only needs amd64. publish-amd64 now pushes :<version> as
soon as the amd64 build is done; publish-arm64 builds arm64 afterwards,
then replaces :<version> with the two-platform index and moves :latest.

Both jobs use one named BuildKit builder whose container outlives the
job, so the dependency layer (cargo chef cook) is reused until the
dependencies change. The release is created after amd64; the binaries
are attached once arm64 is in.
2026-09-24 08:04:53 -07:00
jcoffey-dev 7109e67f07 Merge pull request 'Trace search: index event type and queue id as integers' (#33) from fix/pg-index-trace-types into main
ci / fork-checks (push) Successful in 30s
ci / build (push) Successful in 37m5s
2026-09-24 14:57:18 +00:00
jcoffey-dev 52b5a5f909 Mark tests/src/store/query.rs as modified by the fork
ci / fork-checks (pull_request) Successful in 1m4s
ci / build (pull_request) Successful in 4m20s
The trace document test changed an upstream file, so it carries the
AGPL section 5(a) notice (tools/fork/notice-check.py).
2026-09-24 07:39:04 -07:00
jcoffey-dev 9232662913 Trace search: index event type and queue id as integers
ci / fork-checks (pull_request) Failing after 47s
ci / build (pull_request) Successful in 4m55s
The trace index task wrote the event type (its name) and the queue id as
text, but the tracing search index types both as integers on every
backend: BIGINT on PostgreSQL and MySQL, long on Elasticsearch. On
PostgreSQL every batch holding a trace document failed with "cannot
convert between the Rust type String and the Postgres type int8", and
since a batch writes trace and email documents together, email indexing
stalled behind it.

The document is now built by trace_search_document(), which writes:

- the event type as the opening event's numeric id, the event
  x:Trace/query's event filter already matches on;
- the queue id as an integer, the first one the trace names;
- every queue id into the keywords as well, since the column holds one
  value and an SMTP session can queue several messages.

index_keyword() replaced the field on every call, so before this only the
last event type and queue id survived anyway.

x:Trace/query's queueId filter parses the id (a string, or now a number)
and matches the column or the keywords, so a session is found by any of
its queue ids on every backend. The monitoring spec says what is indexed.

Traces indexed before this on the built-in index keep their text values;
the reindexTelemetry maintenance task rebuilds them.

Tests: the search store suite builds trace documents with the index
task's code, indexes them and finds them by queue id, event type and
keyword (Sqlite, PostgreSQL, MySQL); the monitoring suite finds a real
trace by queueId through x:Trace/query.
2026-09-24 07:29:08 -07:00
jcoffey-dev 57d1c5b074 Merge pull request 'Broadcast subscriber: fix the inverted subscribe retry backoff' (#32) from fix/subscriber-backoff into main
ci / fork-checks (push) Successful in 43s
ci / build (push) Canceled after 30m24s
2026-09-24 14:26:52 +00:00
6 changed files with 381 additions and 40 deletions
+69 -15
View File
@@ -3,11 +3,28 @@
# whether a person pushed it or weekly-release.yml created it through the # whether a person pushed it or weekly-release.yml created it through the
# releases API. # releases API.
# #
# The image is multi-arch (linux/amd64, linux/arm64) as before, but built in # The image is multi-arch (linux/amd64, linux/arm64), built by two jobs on
# one buildx run on host1 instead of one native runner per architecture: the # the image-build runner rather than one buildx run for both. The Dockerfile's
# Dockerfile's builder stage runs on the build platform and cross-compiles # builder stage runs on the build platform and cross-compiles with an aarch64
# with an aarch64 linker, so only the small final stage (apt, setcap) goes # linker, so only the small final stage (apt, setcap) goes through QEMU for
# through QEMU for arm64. No digest-joining job is needed. # arm64 -- but two release builds (LTO, one codegen unit) side by side on one
# machine each take twice as long. Production runs amd64, so amd64 goes first
# and on its own:
# * publish-amd64 pushes :<version>-amd64 and :<version>, a plain amd64
# image, as soon as its build is done. A deploy can start from it.
# * publish-arm64 then builds arm64, pushes :<version>-arm64, and replaces
# :<version> with the two-platform index. :latest moves only here, so it
# never names an image without arm64.
#
# Both jobs use one BuildKit builder, `gitea-builder`, whose container
# (buildx_buildkit_gitea-builder0) and state volume stay on the runner's host
# between jobs: a job container's `buildx create` finds the existing container
# and reuses it and its cache. The dependency build (`cargo chef cook`) is
# keyed on the recipe, which only a dependency change alters, so a release
# normally compiles just the workspace. Removing that container or its volume
# costs the next release a cold build, nothing more. The planner and dependency
# layers for the build platform are shared, so arm64 also reuses what amd64
# just did where it can.
# #
# Two guards before anything is pushed: # Two guards before anything is pushed:
# * the tag must be v<brand_version!>. The version is a string in # * the tag must be v<brand_version!>. The version is a string in
@@ -62,7 +79,7 @@ jobs:
echo "version=$V" >> "$GITHUB_OUTPUT" echo "version=$V" >> "$GITHUB_OUTPUT"
echo "version $V" echo "version $V"
publish: publish-amd64:
needs: [version] needs: [version]
runs-on: docker runs-on: docker
container: container:
@@ -81,16 +98,15 @@ jobs:
test -n "$REGISTRY" && test -n "$VERSION" test -n "$REGISTRY" && test -n "$VERSION"
test -n "$PACKAGE_TOKEN" || { echo "PACKAGE_TOKEN secret is not set on this repository" >&2; exit 1; } test -n "$PACKAGE_TOKEN" || { echo "PACKAGE_TOKEN secret is not set on this repository" >&2; exit 1; }
echo "$PACKAGE_TOKEN" | docker login -u jcoffey-dev --password-stdin "$REGISTRY" echo "$PACKAGE_TOKEN" | docker login -u jcoffey-dev --password-stdin "$REGISTRY"
docker run --privileged --rm tonistiigi/binfmt --install arm64
docker buildx create --use --name gitea-builder --driver docker-container || docker buildx use gitea-builder docker buildx create --use --name gitea-builder --driver docker-container || docker buildx use gitea-builder
# Attestations off, as before: they add manifests of their own to the # Attestations off, as before: they add manifests of their own, and the
# index, and the index should hold the two images and nothing else. # index should hold the two images and nothing else.
- run: | - run: |
docker buildx build \ docker buildx build \
--platform linux/amd64,linux/arm64 \ --platform linux/amd64 \
--provenance=false --sbom=false \ --provenance=false --sbom=false \
--tag "$IMAGE:$VERSION-amd64" \
--tag "$IMAGE:$VERSION" \ --tag "$IMAGE:$VERSION" \
--tag "$IMAGE:latest" \
--push . --push .
docker buildx imagetools inspect "$IMAGE:$VERSION" docker buildx imagetools inspect "$IMAGE:$VERSION"
# Gitea keeps a container package on its owner; linking it shows it on # Gitea keeps a container package on its owner; linking it shows it on
@@ -103,11 +119,47 @@ jobs:
- if: always() - if: always()
run: docker logout "$REGISTRY" || true run: docker logout "$REGISTRY" || true
publish-arm64:
needs: [version, publish-amd64]
runs-on: docker
container:
image: docker:28-cli@sha256:625d9431a9f54c5a2bc90f24f0e1c3d55b1349fd857dd85035f98c2c9acbdd4d # 28-cli
volumes:
- /var/run/docker.sock:/var/run/docker.sock
env:
DOCKER_BUILDKIT: "1"
REGISTRY: ${{ vars.REGISTRY }}
IMAGE: ${{ vars.REGISTRY }}/${{ github.repository }}
VERSION: ${{ needs.version.outputs.version }}
PACKAGE_TOKEN: ${{ secrets.PACKAGE_TOKEN }}
steps:
- uses: coffey-labs/actions/checkout@fab0c4d45e0162963965f1555df27b7bed5e20ec
- run: |
echo "$PACKAGE_TOKEN" | docker login -u jcoffey-dev --password-stdin "$REGISTRY"
docker run --privileged --rm tonistiigi/binfmt --install arm64
docker buildx create --use --name gitea-builder --driver docker-container || docker buildx use gitea-builder
# The index is built from the two per-architecture tags rather than from
# :<version>, which by now is the amd64 image and would be read as such.
- run: |
docker buildx build \
--platform linux/arm64 \
--provenance=false --sbom=false \
--tag "$IMAGE:$VERSION-arm64" \
--push .
docker buildx imagetools create \
--tag "$IMAGE:$VERSION" \
--tag "$IMAGE:latest" \
"$IMAGE:$VERSION-amd64" "$IMAGE:$VERSION-arm64"
docker buildx imagetools inspect "$IMAGE:$VERSION"
- if: always()
run: docker logout "$REGISTRY" || true
# The weekly release creates its Release (and so the tag) first; a tag # The weekly release creates its Release (and so the tag) first; a tag
# pushed by hand has none. Either way the tag ends up with exactly one # pushed by hand has none. Either way the tag ends up with exactly one
# Release, created after the image exists so its pull instructions work. # Release, created once the amd64 image exists so its pull instructions
# work; arm64 and the binaries follow.
release: release:
needs: [version, publish] needs: [version, publish-amd64]
runs-on: light runs-on: light
container: container:
image: python:3.13-slim@sha256:8d9d0b8bcf6506481eae4907c18f5e3e7902e629f5f6d684f9e7c32e85e3ddf0 # 3.13-slim image: python:3.13-slim@sha256:8d9d0b8bcf6506481eae4907c18f5e3e7902e629f5f6d684f9e7c32e85e3ddf0 # 3.13-slim
@@ -131,7 +183,9 @@ jobs:
except urllib.error.HTTPError as e: except urllib.error.HTTPError as e:
if e.code != 404: raise if e.code != 404: raise
image = f"{os.environ['REGISTRY']}/{os.environ['REPO']}:{version}" image = f"{os.environ['REGISTRY']}/{os.environ['REPO']}:{version}"
body = (f"Container image: `{image}` (linux/amd64, linux/arm64); also `:latest`.\n\n" body = (f"Container image: `{image}` (linux/amd64, linux/arm64); also `:latest`. "
"amd64 is published first; arm64 is added to the same tag when its build "
"finishes, and `:latest` moves then.\n\n"
"Binaries for a host install are attached: `inbuxa-linux-amd64.tar.gz` and " "Binaries for a host install are attached: `inbuxa-linux-amd64.tar.gz` and "
"`inbuxa-linux-arm64.tar.gz`, with `SHA256SUMS`. Each is the binary out of this " "`inbuxa-linux-arm64.tar.gz`, with `SHA256SUMS`. Each is the binary out of this "
"release's image for that architecture, so it is the same build. The image " "release's image for that architecture, so it is the same build. The image "
@@ -154,7 +208,7 @@ jobs:
# `docker create` does not start anything, so pulling an arm64 image on an # `docker create` does not start anything, so pulling an arm64 image on an
# amd64 runner and copying a file out of it needs no emulation. # amd64 runner and copying a file out of it needs no emulation.
binaries: binaries:
needs: [version, publish, release] needs: [version, publish-arm64, release]
runs-on: docker runs-on: docker
container: container:
image: docker:28-cli@sha256:625d9431a9f54c5a2bc90f24f0e1c3d55b1349fd857dd85035f98c2c9acbdd4d # 28-cli image: docker:28-cli@sha256:625d9431a9f54c5a2bc90f24f0e1c3d55b1349fd857dd85035f98c2c9acbdd4d # 28-cli
+17 -2
View File
@@ -427,9 +427,24 @@ pub(crate) async fn trace_query(
} }
None => false, None => false,
}, },
Property::QueueId => match value.as_str() { // The queue id column is an integer on every search backend, and
// holds a trace's first queue id; the keywords carry all of them
Property::QueueId => match value
.as_str()
.and_then(|v| v.trim().parse::<u64>().ok())
.or_else(|| value.as_u64())
{
Some(queue_id) => { Some(queue_id) => {
search.push(SearchFilter::eq(TracingSearchField::QueueId, queue_id.to_string())); search.extend([
SearchFilter::Or,
SearchFilter::eq(TracingSearchField::QueueId, queue_id),
SearchFilter::has_text(
TracingSearchField::Keywords,
queue_id.to_string(),
nlp::language::Language::None,
),
SearchFilter::End,
]);
true true
} }
None => false, None => false,
+59 -22
View File
@@ -567,20 +567,14 @@ async fn build_contact_document(
} }
// inbuxa: MON-16: a trace's search document, when trace search is on: // inbuxa: MON-16: a trace's search document, when trace search is on
// its event types, queue ids, and addresses, their domains, hosts, IPs,
// message ids and account names as keywords
async fn build_tracing_span_document( async fn build_tracing_span_document(
server: &Server, server: &Server,
span_id: u64, span_id: u64,
) -> trc::Result<Option<IndexDocument>> { ) -> trc::Result<Option<IndexDocument>> {
use common::telemetry::tracers::store::MaybeTrace; use common::telemetry::tracers::store::MaybeTrace;
use registry::schema::{enums::SearchTracingField, structs::Search}; use registry::schema::structs::Search;
use store::{ use store::write::{TelemetryClass, ValueClass};
search::TracingSearchField,
write::{TelemetryClass, ValueClass},
};
use trc::Key;
let settings = server let settings = server
.registry() .registry()
@@ -590,7 +584,6 @@ async fn build_tracing_span_document(
if !settings.index_telemetry { if !settings.index_telemetry {
return Ok(None); return Ok(None);
} }
let wants = |field: SearchTracingField| settings.index_tracing_fields.iter().any(|f| *f == field);
let Some(MaybeTrace(Some(trace))) = server let Some(MaybeTrace(Some(trace))) = server
.tracing_store() .tracing_store()
.get_value::<MaybeTrace>(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span( .get_value::<MaybeTrace>(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(
@@ -601,23 +594,67 @@ async fn build_tracing_span_document(
return Ok(None); return Ok(None);
}; };
let mut document = IndexDocument::new(SearchIndex::Tracing).with_id(span_id); Ok(Some(trace_search_document(
let mut seen = store::ahash::AHashSet::new(); span_id,
for event in trace.events.iter() { &trace,
if wants(SearchTracingField::EventType) && seen.insert(event.event.as_str().to_string()) { &settings
document.index_keyword(TracingSearchField::EventType, event.event.as_str()); .index_tracing_fields
.iter()
.copied()
.collect::<Vec<_>>(),
)))
} }
/// inbuxa: MON-16: the search document for a stored trace.
///
/// The event type and queue id columns are integers on every search backend
/// (BIGINT on PostgreSQL and MySQL, long on Elasticsearch), and each holds a
/// single value per trace: the event type is the trace's opening event, the
/// one `x:Trace/query` filters on, and the queue id is the first queue id the
/// trace mentions. Every queue id also goes into the keywords, so a session
/// that queued several messages is found by any of them.
pub fn trace_search_document(
span_id: u64,
trace: &registry::schema::structs::Trace,
fields: &[registry::schema::enums::SearchTracingField],
) -> IndexDocument {
use registry::schema::{enums::SearchTracingField, structs::TraceValue};
use store::search::TracingSearchField;
use trc::Key;
let wants = |field: SearchTracingField| fields.contains(&field);
let mut document = IndexDocument::new(SearchIndex::Tracing).with_id(span_id);
if wants(SearchTracingField::EventType)
&& let Some(first) = trace.events.iter().next()
{
document.index_unsigned(TracingSearchField::EventType, first.event.to_id() as u64);
}
let mut seen = store::ahash::AHashSet::new();
let mut queue_id_indexed = false;
for event in trace.events.iter() {
for kv in event.key_values.iter() { for kv in event.key_values.iter() {
let text = match &kv.value { let text = match &kv.value {
registry::schema::structs::TraceValue::String(v) => v.value.clone(), TraceValue::String(v) => v.value.clone(),
registry::schema::structs::TraceValue::UnsignedInt(v) => v.value.to_string(), TraceValue::UnsignedInt(v) => v.value.to_string(),
registry::schema::structs::TraceValue::IpAddr(v) => v.value.to_string(), TraceValue::IpAddr(v) => v.value.to_string(),
_ => continue, _ => continue,
}; };
match kv.key { match kv.key {
Key::QueueId if wants(SearchTracingField::QueueId) => { Key::QueueId => {
if seen.insert(format!("q:{text}")) { let Ok(queue_id) = text.parse::<u64>() else {
document.index_keyword(TracingSearchField::QueueId, &text); continue;
};
if wants(SearchTracingField::QueueId) && !queue_id_indexed {
document.index_unsigned(TracingSearchField::QueueId, queue_id);
queue_id_indexed = true;
}
if wants(SearchTracingField::Keywords) && seen.insert(format!("k:{text}")) {
document.index_text(
TracingSearchField::Keywords,
&text,
nlp::language::Language::None,
);
} }
} }
Key::From Key::From
@@ -648,7 +685,7 @@ async fn build_tracing_span_document(
} }
} }
} }
Ok(Some(document)) document
} }
// inbuxa: UD-1, UD-4: archives a deleted file, event or contact noted at // inbuxa: UD-1, UD-4: archives a deleted file, event or contact noted at
+8 -3
View File
@@ -212,10 +212,15 @@ unchanged.
- **MON-16.** With `indexTelemetry` on, storing a trace schedules an - **MON-16.** With `indexTelemetry` on, storing a trace schedules an
`IndexTrace` task. The task builds one document for `SearchIndex::Tracing` `IndexTrace` task. The task builds one document for `SearchIndex::Tracing`
with the fields named in `indexTracingFields`: with the fields named in `indexTracingFields`:
- `eventType`: every event type in the trace; - `eventType`: the trace's opening event, as its numeric id;
- `queueId`: every `queueId` value; - `queueId`: the first `queueId` value, as an integer;
- `keywords`: every address in `from` and `to`, each address's domain, every - `keywords`: every address in `from` and `to`, each address's domain, every
`domain`, `hostname`, `remoteIp`, `messageId` and `accountName` value. `domain`, `hostname`, `remoteIp`, `messageId` and `accountName` value,
and every `queueId` value.
The event type and queue id are single integer columns on every search
backend (BIGINT on PostgreSQL and MySQL), so the `queueId` filter matches
the column or any queue id in the keywords, and a session that queued
several messages is found by each of them.
So searching `example.org` finds every trace to or from that domain, as the So searching `example.org` finds every trace to or from that domain, as the
upstream suite expects. With `indexTelemetry` off nothing is indexed, and upstream suite expects. With `indexTelemetry` off nothing is indexed, and
the `text` and `queueId` filters are refused (see "Interfaces"). the `text` and `queueId` filters are refused (see "Interfaces").
+163
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]> * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* *
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use crate::{store::deflate_test_resource, utils::server::TestServer}; use crate::{store::deflate_test_resource, utils::server::TestServer};
@@ -122,6 +124,10 @@ pub async fn test(test: &TestServer) {
println!("Running global id filtering tests..."); println!("Running global id filtering tests...");
test_global(store.clone()).await; test_global(store.clone()).await;
// inbuxa: trace documents as the index task builds them
println!("Running trace document tests...");
test_trace_documents(store.clone()).await;
// Large document insert test // Large document insert test
println!("Running large document insert tests..."); println!("Running large document insert tests...");
let mut large_text = String::with_capacity(20 * 1024 * 1024); let mut large_text = String::with_capacity(20 * 1024 * 1024);
@@ -809,3 +815,160 @@ async fn test_global(store: SearchStore) {
AHashSet::from_iter([3, 4, 5]) AHashSet::from_iter([3, 4, 5])
); );
} }
// inbuxa: MON-16: documents built by the index task from stored traces go
// into every search backend (the SQL backends type etyp and qid as BIGINT)
// and are found again by queue id and keyword.
async fn test_trace_documents(store: SearchStore) {
use registry::schema::{
enums::SearchTracingField,
structs::{
Trace, TraceEvent, TraceKeyValue, TraceValue, TraceValueString,
TraceValueUnsignedInt,
},
};
use services::task_manager::index::trace_search_document;
use trc::{DeliveryEvent, EventType, Key, SmtpEvent};
let kv_u = |key: Key, value: u64| TraceKeyValue {
key,
value: TraceValue::UnsignedInt(TraceValueUnsignedInt { value }),
};
let kv_s = |key: Key, value: &str| TraceKeyValue {
key,
value: TraceValue::String(TraceValueString {
value: value.to_string(),
}),
};
let event = |event: EventType, key_values: Vec<TraceKeyValue>| TraceEvent {
event,
key_values: key_values.into(),
..Default::default()
};
let fields = [
SearchTracingField::EventType,
SearchTracingField::QueueId,
SearchTracingField::Keywords,
];
// An SMTP session that queued two messages, and a delivery attempt
let session = Trace {
events: vec![
event(
EventType::Smtp(SmtpEvent::ConnectionStart),
vec![kv_s(Key::RemoteIp, "192.0.2.7")],
),
event(
EventType::Smtp(SmtpEvent::MailFrom),
vec![kv_s(Key::From, "[email protected]")],
),
event(
EventType::Smtp(SmtpEvent::RcptTo),
vec![kv_u(Key::QueueId, 9_000_000_001), kv_s(Key::To, "[email protected]")],
),
event(
EventType::Smtp(SmtpEvent::RcptTo),
vec![kv_u(Key::QueueId, 9_000_000_002)],
),
]
.into(),
};
let delivery = Trace {
events: vec![event(
EventType::Delivery(DeliveryEvent::AttemptStart),
vec![kv_u(Key::QueueId, 9_000_000_003), kv_s(Key::Hostname, "relay.example.net")],
)]
.into(),
};
let documents = vec![
trace_search_document(100, &session, &fields),
trace_search_document(101, &delivery, &fields),
];
assert!(
documents
.iter()
.all(|d| d.has_field(&SearchField::Tracing(TracingSearchField::QueueId))
&& d.has_field(&SearchField::Tracing(TracingSearchField::EventType))),
"trace documents carry a queue id and an event type"
);
store.index(documents).await.unwrap();
if let SearchStore::ElasticSearch(store) = &store {
store.refresh_index(SearchIndex::Tracing).await.unwrap();
}
let query = |filters: Vec<SearchFilter>| {
let store = store.clone();
async move {
store
.query_global(
SearchQuery::new(SearchIndex::Tracing)
.with_filter(SearchFilter::ge(SearchField::Id, 100u64))
.with_filters(filters),
)
.await
.unwrap()
.into_iter()
.collect::<AHashSet<_>>()
}
};
// By queue id, the way x:Trace/query asks: the queue id column, or any
// queue id in the keywords
let by_queue_id = |queue_id: u64| {
vec![
SearchFilter::Or,
SearchFilter::eq(TracingSearchField::QueueId, queue_id),
SearchFilter::has_text(
TracingSearchField::Keywords,
queue_id.to_string(),
Language::None,
),
SearchFilter::End,
]
};
assert_eq!(query(by_queue_id(9_000_000_001)).await, AHashSet::from_iter([100]));
assert_eq!(query(by_queue_id(9_000_000_002)).await, AHashSet::from_iter([100]));
assert_eq!(query(by_queue_id(9_000_000_003)).await, AHashSet::from_iter([101]));
assert_eq!(query(by_queue_id(9_000_000_004)).await, AHashSet::new());
assert_eq!(
query(vec![SearchFilter::eq(TracingSearchField::QueueId, 9_000_000_003u64)]).await,
AHashSet::from_iter([101])
);
// By opening event type
assert_eq!(
query(vec![SearchFilter::eq(
TracingSearchField::EventType,
EventType::Delivery(DeliveryEvent::AttemptStart).to_id() as u64,
)])
.await,
AHashSet::from_iter([101])
);
// By keyword: an address, lowercased, and its domain
assert_eq!(
query(vec![SearchFilter::has_text(
TracingSearchField::Keywords,
"example.org",
Language::None,
)])
.await,
AHashSet::from_iter([100])
);
assert_eq!(
query(vec![SearchFilter::has_text(
TracingSearchField::Keywords,
"relay.example.net",
Language::None,
)])
.await,
AHashSet::from_iter([101])
);
for id in [100u64, 101] {
store
.unindex(
SearchQuery::new(SearchIndex::Tracing)
.with_filter(SearchFilter::eq(SearchField::Id, id)),
)
.await
.unwrap();
}
}
+67
View File
@@ -148,6 +148,73 @@ pub async fn test(test: &mut TestServer) {
"test 9: to" "test 9: to"
); );
// MON-16: the queueId filter finds the traces that name a queue id (the
// session that queued the message and its delivery attempt) through the
// search index, given as a string or a number (the index column is an
// integer)
fn queue_ids(value: &Value, out: &mut Vec<u64>) {
match value {
Value::Object(map) => {
if map.get("key").and_then(|k| k.as_str()) == Some("queueId")
&& let Some(id) = map
.get("value")
.and_then(|v| v.get("value").unwrap_or(v).as_u64())
{
out.push(id);
}
map.values().for_each(|v| queue_ids(v, out));
}
Value::Array(list) => list.iter().for_each(|v| queue_ids(v, out)),
_ => {}
}
}
let with_ids = traces
.iter()
.map(|t| {
let mut ids = Vec::new();
queue_ids(t, &mut ids);
(t["id"].as_str().unwrap().to_string(), ids)
})
.collect::<Vec<_>>();
let queue_id = with_ids
.iter()
.find_map(|(_, ids)| ids.first().copied())
.expect("MON-16: a trace with a queue id");
let mut expected = with_ids
.iter()
.filter(|(_, ids)| ids.contains(&queue_id))
.map(|(id, _)| id.clone())
.collect::<Vec<_>>();
expected.sort();
for filter in [json!(queue_id.to_string()), json!(queue_id)] {
let response = admin
.jmap_method_call("x:Trace/query", json!({"filter": {"queueId": filter}}))
.await;
let mut found = response
.0
.pointer("/methodResponses/0/1/ids")
.and_then(|ids| ids.as_array())
.map(|ids| {
ids.iter()
.filter_map(|id| id.as_str().map(str::to_string))
.collect::<Vec<_>>()
})
.unwrap_or_default();
found.sort();
assert_eq!(found, expected, "MON-16: queueId {filter}: {response:?}");
}
let response = admin
.jmap_method_call(
"x:Trace/query",
json!({"filter": {"queueId": (queue_id ^ 0x5a5a_5a5a).to_string()}}),
)
.await;
assert_eq!(
response.0.pointer("/methodResponses/0/1/ids"),
Some(&json!([])),
"MON-16: an unknown queue id"
);
// Acceptance test 24: destroy removes a trace; create is refused // Acceptance test 24: destroy removes a trace; create is refused
let trace_id = traces[0]["id"].as_str().unwrap().to_string(); let trace_id = traces[0]["id"].as_str().unwrap().to_string();
let response = admin let response = admin