15 Commits

Author SHA1 Message Date
8b9dd0d9e8 Merge pull request 'feat(sidecar): protocol 5 — cutover 2b of 7 (edgemain)' (#35) from edge into main
Some checks failed
sync-project-tree / sync (push) Successful in 11s
SonarQube / analysis (push) Failing after -37s
Release sidecar / release (push) Successful in 11m55s
Reviewed-on: #35
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
2026-09-01 13:55:27 +00:00
f41237392d Merge pull request 'feat(sidecar): protocol 5' (#34) from feature/protocol-v5 into edge
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m16s
Reviewed-on: #34
2026-09-01 00:27:05 +00:00
d0c2e7d6e1 feat(sidecar): protocol 5
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m55s
PROTOCOL_VERSION 4 -> 5, and nothing else.

That is the whole change, and it is worth saying why. Protocol 5 adds fields to
house.decay and vendor.listing and one new kind, account.login.result — and the
sidecar needs no code for any of it. Every frame is persisted whole, the board
tables index only the columns they already had, and there is no kind allowlist, so
the new fields ride inside the stored JSON and the new kind lands in `events` like
any other.

No store migration this time, unlike v4. v4 needed one because it added a column to
a board table that already existed; nothing here does. A bump that touches one
constant is the EXPECTED cost of an additive protocol version in a dumb forwarder —
the sidecar defines no schema for a frame's contents, so it needs no change when
they grow. v4 was the exception.

The doc comment records the three enrichments and why they were bumped together: a
protocol bump costs a sidecar release, a republished bundle and an operator update
on every shard, so a field left out costs a whole second round of that rather than a
follow-up commit.

Verified against the real shard: GET /health reports "protocol": 5, and all three
enrichments arrived through the generic forward path — the decay schedule (with
estimatedCollapse present only on the IDOC frame), the vendor fee block, and both
outcomes of account.login.result.

cargo fmt --check clean, clippy -D warnings clean, 39 tests passing.

Docs: RunicGateway/docs link/v5.md.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-31 19:19:48 -05:00
4b8ea768b6 Merge pull request 'ci(release): show the error body, retry the POST, and sweep for orphan tags' (#33) from ci/release-post-retry-and-error-body into main
All checks were successful
sync-project-tree / sync (push) Successful in 9s
Release sidecar / release (push) Successful in -55s
SonarQube / analysis (push) Successful in 51s
Reviewed-on: #33
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
2026-08-24 19:42:20 +00:00
6fb063818a ci(release): show the error body, retry the POST, and sweep for orphan tags
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m42s
This file is the ancestor of installer's release.yml, and installer#22's
release run found two gaps in it the hard way: the run built every artifact,
pushed its tag, then took a 500 from POST /releases one second later and exited
22, leaving the tag orphaned with no binaries published.

link has not hit that, but it has the same two gaps verbatim.

`curl -sSf` prints no response body on an error status, so the only thing such
a failure leaves in the log is "curl: (22) ... error: 500" and the cause has to
be inferred from timestamps. Every call in the release step now captures the
body and prints it on failure, including the asset uploads.

And nothing retried, so a transient 5xx becomes a permanent orphan. The POST
now retries five times with a 5/10/15/20s backoff. 4xx is deliberately not
retried: a bad token or a malformed body will not improve by being sent again,
and retrying would turn a clear failure into a slow one.

The asset uploads get the same treatment, because a release whose SHA256SUMS
does not cover every binary it advertises is worse than no release -- that file
is the trust anchor for an unsigned download.

The third gap is the one worth reading. The orphan-tag recovery in the plan
step is VERSION-SCOPED: it computes VERSION from the newest tag plus the bump,
then only checks refs/tags/v${VERSION}. That recovers an orphan on the very
next run and is useless afterwards, because once any releasable commit lands
the next run computes a NEW version and never looks at the old tag again.

servuo-plugins v0.1.0 proves it, and the proof is pointed: the commit that
ADDED that recovery was itself typed "fix(release): ... recover the orphaned
v0.1.0 tag", so it bumped to v0.1.1 and the run that introduced the recovery
stepped straight past the tag it was written to rescue. That tag is orphaned to
this day.

So the plan step now sweeps every v* tag and warns about any without a release.
Deliberately warns rather than recovers: publishing an old version would mean
building today's tree and shipping it under a tag whose tree it is not, which
is worse than the inconsistency it fixes. It also never fails the run -- a
sweep that can break a good release is a sweep someone will delete.

Verified by extracting both steps from the YAML and running them: bash -n
clean, the YAML parses, no empty template token in either step, the retry loop
exercised against a stubbed curl across seven cases (first-try success,
500-then-success, two 500s then success, five 500s giving up, 403 and 404
aborting without retrying, and a 000 network failure retried), and the sweep
run against the real repositories -- link clean, servuo-plugins reporting
v0.1.0, installer clean.

Typed ci(...) rather than fix(...) on purpose: the plan step bumps on feat/fix,
and this changes no binary, so a release here would be an empty one. That is
the same rule the fix commit above tripped over.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-24 12:40:42 -05:00
7499e099f4 Merge pull request 'feat(sidecar)!: protocol 4 — guild rosters, and a way to migrate the store (Teams cutover 2/6)' (#32) from edge into main
All checks were successful
sync-project-tree / sync (push) Successful in -47s
SonarQube / analysis (push) Successful in 54s
Release sidecar / release (push) Successful in 10m2s
Reviewed-on: #32
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
2026-08-19 08:55:32 +00:00
2408d31ff7 Merge pull request 'feat(sidecar)!: guild rosters on the guild board, and a way to migrate the store' (#31) from feat/teams-phase1-guild-roster into edge
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 3m11s
Reviewed-on: #31
2026-08-17 19:28:55 +00:00
b00f2719a4 fix(sidecar): reassemble a guild roster that arrived in several frames
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m30s
The shard caps members per `guild.roster` frame, so a guild over that cap emits
several frames carrying `seq`/`more`/`total`. The board's upsert wrote whichever
array it was handed, so each frame overwrote the last and only the final chunk
survived: a live 155-member guild, split 50/50/50/5, landed on the board with 5
members while `guild.update` correctly reported 155 beside it.

Every unit test passed through this, because they all exercised a single-frame
roster. Only the live rig caught it — the case does not arise until a guild
exceeds the cap.

Frames are now reassembled in memory and written once, on the frame that closes
the roster. The alternative — appending to the `members` column per frame — was
rejected twice over: it would make the write a read-modify-write, which is the
exact thing splitting the board across two columns exists to avoid, and it would
publish a torn roster, since a reader hitting GET /guilds between frames would
see a partial member list presented as the whole truth.

Buffering here does not make the sidecar stateful in the sense that matters. This
is transport-level reassembly — the same category of work as turning bytes into a
line — and it holds nothing once a roster is complete.

The ordinary case is unchanged and untouched by the buffer: a guild inside the
cap arrives as `seq` 0 with `more` false and is returned immediately, never
entering the map. What the buffer adds is the handling of everything that can go
wrong around a split roster: a fresh `seq` 0 supersedes an abandoned partial, an
out-of-order frame discards the partial rather than storing one with an
undetectable hole, a continuation with no start is ignored, a reconnect drops
every partial (the shard restarts each roster at 0), and accumulation is bounded
so a shard that never sends a closing frame cannot grow this map without limit.

Re-verified on the live rig: four frames reassembled to 153 entries after two
members were removed, with both departed serials absent.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-17 12:52:26 -05:00
9216006208 feat(sidecar)!: guild rosters on the guild board, and a way to migrate the store
Protocol 4 gives the guild board a real member list instead of the member
*count* that was all Protocol 2 could express. `guild.roster` carries the set;
`guild.leave` is forwarded but deliberately not projected.

The roster lives in its own `members` column rather than as a field folded into
`json`. That column holds the verbatim `guild.update` line, so a roster write
into it would clobber the snapshot — name, abbreviation, leader, online count —
that `guild.update` owns. Two writers across two columns of one row means both
stay plain upserts: neither reads the other's value first, so there is no
read-modify-write and no ordering requirement between the two kinds. `GET
/guilds` folds the roster back in as `roster` at read time.

`guild.leave` gets no board arm on purpose. The shard re-emits `guild.roster`
whenever the member set changes, so the board self-corrects within one sweep,
and keeping the delta out of the projection is what keeps the sidecar a
forwarder rather than a thing that maintains state.

This is also the repo's first store migration, and the reason it needed one:
`SCHEMA` is `CREATE TABLE IF NOT EXISTS`, which can add a table but cannot add a
column to a table that already exists. Every schema change up to and including
Protocol 3.0 happened to add whole tables, so `ALTER TABLE` appears nowhere in
this repo's history and the gap was invisible. `guilds.members` is the first
column added to an existing table, so without a mechanism the column would
simply never reach an installed sidecar and every roster write would fail.

The counter is SQLite's own `PRAGMA user_version` — an integer in the database
header, so it costs no table and cannot drift from the file it describes. Each
step runs in a transaction together with the bump recording it, so a step lands
completely or not at all. A database written by a *newer* sidecar warns and
continues rather than failing: every step is additive, so a newer schema has only
columns an older reader ignores, and refusing to start would turn rolling the
binary back — a recovery path — into a dead end.

A migration failure aborts startup, which was already the behaviour and is the
right one: a half-migrated store answers the website with confusing partial data,
and the shard dials *out*, so a sidecar that refuses to start never stalls the
game.

store.rs had no tests before this. The six added here cover the upgrade path that
matters (an existing pre-Protocol-4 database gains the column and lands at the
current version), that a restart re-running the migration is a no-op, that a
roster does not clobber the snapshot, that the two writers work in either order,
and that a guild with no roster yet has no `roster` key at all — "not known" and
"known to be empty" must not be conflated, or a website renders an empty roster
as fact.

Also gates PRs into `edge`, not just `main`. This workstream lands ten phases
there, and gating only the `main` hop would run these checks for the first time
at the cutover. The precedent and the reasoning are already in
RunicGateway/installer's copy of this workflow.

Refs: docs/website/TEAMS.md Part 12 Phase 1

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-17 12:38:44 -05:00
be0efd348c Merge pull request 'docs(readme): give operators an entry point before the build steps' (#30) from docs/installer-first-setup into main
Some checks failed
Release sidecar / release (push) Successful in 10s
SonarQube / analysis (push) Failing after 13s
sync-project-tree / sync (push) Successful in -30s
Reviewed-on: #30
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
2026-08-07 21:33:00 +00:00
3dbc2f490c docs(readme): give operators an entry point before the build steps
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 1m32s
The README opened straight into cargo build, which is the wrong first
instruction for someone standing up a shard: the installer places this
binary, its config, a service account and the service registration.

Adds a short operator section pointing at the installer (and at INSTALL.md
Appendix A3-A4 for installing by hand, still supported), marks everything
below it as development, and lists the installer under related repos.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-07 16:05:56 -05:00
7b6584006e Merge pull request 'fix(release): tag only, and stop pushing to main' (#28) from fix/release-tag-only into main
All checks were successful
sync-project-tree / sync (push) Successful in 8s
SonarQube / analysis (push) Successful in 43s
Release sidecar / release (push) Successful in 10m22s
Reviewed-on: #28
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
2026-08-07 19:17:11 +00:00
67d7800300 Merge pull request 'feat(sidecar): start as a real Windows service' (#29) from feat/windows-service into main
All checks were successful
sync-project-tree / sync (push) Successful in 1m3s
SonarQube / analysis (push) Successful in 1m10s
Release sidecar / release (push) Successful in 8m21s
Reviewed-on: #29
2026-08-07 18:53:14 +00:00
96af2afa68 feat(sidecar): start as a real Windows service
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 2m22s
`sc.exe start RunicGatewayLink` failed with 1053 on every Windows install:
"a timeout was reached (30000 milliseconds) while waiting for the service to
connect", with SERVICE_EXIT_CODE 0. Nothing had crashed. The sidecar was a
plain console program, and the Windows service control manager only supervises
a process that calls StartServiceCtrlDispatcher and identifies itself within
~30 seconds.

The installer's design assumed symmetry with systemd, which supervises any
foreground process. Windows has no equivalent: it is a service-aware binary or
a shim, and a shim was already rejected as a third binary to keep current.

Split the entry point so the platform only owns starting and stopping:

  systemd --> main --> unix::run ---------------+
                                                +--> app::run
  SCM ------> main --> windows::run --> ServiceMain
                                    \-> console fallback

- app.rs is the whole sidecar, unchanged and shared. No #[cfg] on the data path.
- windows.rs speaks the SCM handshake. The dispatcher is tried first and failing
  is expected: ERROR_FAILED_SERVICE_CONTROLLER_CONNECT (1063) means "not started
  by the SCM" and falls through to a normal foreground run, so one binary does
  both with no --service flag to forget.
- Running is reported only once the shard port is bound and the store is open, so
  a bad config fails the start instead of flapping Running -> Stopped, and a
  failed run leaves a nonzero SERVICE_EXIT_CODE instead of the misleading 0.
- A service has no stdout, so service mode logs to uo-link-sidecar.log.<date>
  beside its config, rolled daily, seven kept.
- unix.rs additionally handles SIGTERM, which is what systemctl stop sends and
  which previously took the default disposition mid-write.

The Windows crates are declared under [target.'cfg(windows)'.dependencies].
Verified: a Linux build in rust:1-slim-bookworm succeeds and resolves neither
windows-service nor tracing-appender.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-07 13:37:28 -05:00
36141a23df fix(release): tag only, and stop pushing to main
All checks were successful
PR Checks / rust-gates (pull_request) Successful in 1m0s
The "Commit version bump and push tag" step had two problems, and the
first hid the second.

It has never once executed. An empty template expression written
literally in one of its comments makes the runner fail to build the
step's script, and a step it cannot build is skipped WITHOUT failing
the job. That is why sidecar/Cargo.toml still says 0.1.0 after six
releases, and why the tag-reuse handling added in #27 was dead on
arrival. The tags exist because Gitea's release API creates one when it
publishes -- the pipeline has been working by accident.

And had it executed, it would have been declined: main is protected, so
the push is rejected by the pre-receive hook. The installer's bundle job
hit exactly that today. A release must not depend on a write to a
protected branch.

So the tag is the version, as it already is in servuo-plugins, whose
release workflow was written this way on purpose and has never needed a
protection exception. The workflow still writes the real version into
Cargo.toml before building, so a released binary self-reports
correctly; what it no longer does is commit that edit back. Nothing
downstream reads the file -- the next version is computed from the
newest tag.

The comment is reworded so the step can actually run, and warns against
writing that token in a comment again. No literal occurrence is left in
this file.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-05 17:16:42 -05:00
10 changed files with 1504 additions and 331 deletions

View File

@@ -17,10 +17,12 @@
# so let this workflow run on one PR first. The `PR Checks / *` glob matches # so let this workflow run on one PR first. The `PR Checks / *` glob matches
# without needing the dropdown. # without needing the dropdown.
# #
# Scope note: this gates PRs into `main` only. Feature work that lands on an # Scope note: `edge` is gated as well as `main`. Multi-phase work lands there
# integration branch first (e.g. `edge`) is still caught on the branch's PR into # first, so gating only the `main` hop would run these checks for the first time
# `main`. To gate that earlier hop too, add the branch to the `branches:` list # at the cutover — the one moment a red build is most expensive to discover. This
# below — nothing else needs to change. # is the same call `RunicGateway/installer` made for the same reason, and it was
# taken here after a nine-PR Android workstream landed on an ungated `edge` with
# no CI at all. Adding a branch to the `branches:` list is the whole change.
# #
# Runner: the same self-hosted `ubuntu-latest` runner release.yml uses. Rust is # Runner: the same self-hosted `ubuntu-latest` runner release.yml uses. Rust is
# not assumed to be preinstalled, so the toolchain step bootstraps it the same # not assumed to be preinstalled, so the toolchain step bootstraps it the same
@@ -31,7 +33,7 @@ name: PR Checks
on: on:
pull_request: pull_request:
branches: [main] branches: [main, edge]
# A newer push to the same PR cancels the in-flight run. # A newer push to the same PR cancels the in-flight run.
concurrency: concurrency:

View File

@@ -149,6 +149,39 @@ jobs:
fi fi
fi fi
# ── Orphan sweep ────────────────────────────────────────────────
#
# The check above is VERSION-SCOPED: it only ever asks about the one
# version this run computed. That is enough to recover an orphan on
# the very next run, and useless afterwards — once any releasable
# commit lands, the next run computes a NEW version, never looks at
# the old tag again, and the orphan becomes permanent and silent.
#
# servuo-plugins v0.1.0 is the proof, and the proof is pointed: the
# commit that ADDED the recovery above was itself typed
# `fix(release): ... recover the orphaned v0.1.0 tag`, so it bumped to
# v0.1.1 — and the run that introduced the recovery stepped straight
# past the tag it was written to rescue. That tag is still orphaned.
#
# So every v* tag is checked, and anything missing a release is
# WARNED about. Deliberately not recovered: publishing an old version
# would mean building today's tree and shipping it under a tag whose
# tree it is not, which is worse than the inconsistency it fixes.
# A human decides whether to recover or drop it.
#
# Never fails the run. A sweep that can break a good release is a
# sweep someone will delete.
ORPHANS=""
for T in $(git tag -l 'v*' --sort=-v:refname); do
T_HTTP="$(curl -s -o /dev/null -w '%{http_code}' \
-H "Authorization: token $(printf '%s' "${REGISTRY_TOKEN:-}" | tr -d '\r\n')" \
"https://${GITEA_HOST}/api/v1/repos/${REPO}/releases/tags/${T}" || echo 000)"
[ "$T_HTTP" = "404" ] && ORPHANS="${ORPHANS} ${T}"
done
if [ -n "${ORPHANS}" ]; then
echo "::warning::Tags with no release:${ORPHANS} — a run failed after tagging. Publish or delete them; this job will not do either."
fi
# Changelog range. A recovery run has nothing after the tag, so # Changelog range. A recovery run has nothing after the tag, so
# summarize what the tag itself contains rather than emitting an empty # summarize what the tag itself contains rather than emitting an empty
# list: the range that produced it, i.e. previous-tag..this-tag. # list: the range that produced it, i.e. previous-tag..this-tag.
@@ -279,33 +312,43 @@ jobs:
ls -l dist && echo "----" && cat dist/SHA256SUMS ls -l dist && echo "----" && cat dist/SHA256SUMS
# ── RELEASE ENGINE: commit the bump, tag, push ─────────────────────── # ── RELEASE ENGINE: commit the bump, tag, push ───────────────────────
- name: Commit version bump and push tag # Tag only — `main` is never pushed to.
#
# This step used to commit the version bump back to main first. Two things
# were wrong with that. It has never once executed: an EMPTY template
# expression written literally in a comment (the `$`+`{{ }}` token, which
# is why it is spelled out here) made the runner fail to build the script
# and skip the whole step silently, which is why sidecar/Cargo.toml still
# says 0.1.0 after six releases (the tags exist because the release API
# creates one when it publishes). And had it executed, it would have been
# declined — main is protected, and a release must not depend on a write
# to a protected branch.
#
# So the tag is the version, as it already is in servuo-plugins. The
# workflow still writes the real version into Cargo.toml before building,
# so a released binary self-reports correctly; what it no longer does is
# commit that edit back. The next version is computed from the newest tag,
# never from Cargo.toml, so nothing downstream depends on the file.
- name: Push the release tag
if: ${{ steps.plan.outputs.release == 'true' }} if: ${{ steps.plan.outputs.release == 'true' }}
env: env:
REGISTRY_USER: ${{ secrets.REGISTRY_USER }} REGISTRY_USER: ${{ secrets.REGISTRY_USER }}
REGISTRY_TOKEN: ${{ secrets.REGISTRY_TOKEN }} REGISTRY_TOKEN: ${{ secrets.REGISTRY_TOKEN }}
run: | run: |
set -euo pipefail set -euo pipefail
VERSION="${{ steps.plan.outputs.version }}"
TAG="${{ steps.plan.outputs.tag }}" TAG="${{ steps.plan.outputs.tag }}"
# Secrets can arrive with a trailing newline (depending on how they were # Secrets can arrive with a trailing newline (depending on how they were
# pasted); a stray CR/LF corrupts the remote URL ("credential url cannot # pasted); a stray CR/LF corrupts the remote URL ("credential url cannot
# be parsed"). Strip line breaks before building the URL. Passing them via # be parsed"). Strip line breaks before building the URL. They are passed
# env (not inline ${{ }}) also keeps a newline from breaking this script. # via env rather than interpolated into this script, so a newline cannot
# break it — do NOT write a template token literally in a comment here,
# or the runner will skip this step without failing the job.
CI_USER="$(printf '%s' "${REGISTRY_USER}" | tr -d '\r\n')" CI_USER="$(printf '%s' "${REGISTRY_USER}" | tr -d '\r\n')"
CI_TOKEN="$(printf '%s' "${REGISTRY_TOKEN}" | tr -d '\r\n')" CI_TOKEN="$(printf '%s' "${REGISTRY_TOKEN}" | tr -d '\r\n')"
git config user.name "uo-link-ci" git config user.name "uo-link-ci"
git config user.email "ci@whitlocktech.com" git config user.email "ci@whitlocktech.com"
git remote set-url origin \ git remote set-url origin \
"https://${CI_USER}:${CI_TOKEN}@${GITEA_HOST}/${REPO}.git" "https://${CI_USER}:${CI_TOKEN}@${GITEA_HOST}/${REPO}.git"
git add "${WORKDIR}/Cargo.toml" "${WORKDIR}/Cargo.lock"
if ! git diff --cached --quiet; then
git commit -m "chore(release): bump version to ${TAG} [skip ci]"
git push origin "HEAD:main"
else
echo "Version unchanged (first release) — no bump commit needed."
fi
# The tag may already exist when finishing a run that died after tagging # The tag may already exist when finishing a run that died after tagging
# (see the plan step). `git tag` on an existing name fails under # (see the plan step). `git tag` on an existing name fails under
# `set -e`; pushing an identical existing tag is a harmless no-op. A # `set -e`; pushing an identical existing tag is a harmless no-op. A
@@ -332,18 +375,75 @@ jobs:
# corrupt the Authorization header. # corrupt the Authorization header.
CI_TOKEN="$(printf '%s' "${REGISTRY_TOKEN}" | tr -d '\r\n')" CI_TOKEN="$(printf '%s' "${REGISTRY_TOKEN}" | tr -d '\r\n')"
REL_ID="$(curl -sSf -X POST "${API}/releases" \ PAYLOAD="$(jq -n --arg tag "$TAG" --arg body "$BODY" \
'{tag_name:$tag, name:$tag, body:$body, draft:false, prerelease:false}')"
# This POST is the step that orphaned tag v0.1.1 (run 75): it landed one
# second after the tag push and Gitea answered 500, having not finished
# processing the pushed tag. Re-running the workflow published the same
# four assets untouched, so the failure was a race, not a bad request.
#
# Two things went wrong there, and both are fixed here.
#
# 1. `curl -sSf` prints NO response body on an error status, so all the
# log carried was "curl: (22) ... error: 500" and the cause had to be
# inferred from timestamps. Capture the body and print it.
# 2. Nothing retried, so a transient 5xx became a permanent orphan tag.
# The plan step CAN recover one, but only on a run that reaches it --
# and a later push with no releasable commits stands down before it
# gets there, so in practice the tag sits until a human notices.
#
# 4xx is deliberately NOT retried: a bad token or a malformed body does
# not improve by being sent again, and retrying only turns a clear
# failure into a slow one.
REL_ID=""
for attempt in 1 2 3 4 5; do
HTTP="$(curl -s -o /tmp/rel.json -w '%{http_code}' -X POST "${API}/releases" \
-H "Authorization: token ${CI_TOKEN}" \ -H "Authorization: token ${CI_TOKEN}" \
-H "Content-Type: application/json" \ -H "Content-Type: application/json" \
-d "$(jq -n --arg tag "$TAG" --arg body "$BODY" \ -d "${PAYLOAD}" || echo 000)"
'{tag_name:$tag, name:$tag, body:$body, draft:false, prerelease:false}')" \
| jq -r '.id')" if [ "$HTTP" = "201" ] || [ "$HTTP" = "200" ]; then
REL_ID="$(jq -r '.id' /tmp/rel.json)"
break
fi
echo "::warning::POST /releases attempt ${attempt} returned HTTP ${HTTP}"
echo "--- response body ---"
cat /tmp/rel.json || true
echo
echo "---------------------"
case "$HTTP" in
4*) echo "::error::HTTP ${HTTP} is a client error - not retrying."; exit 1 ;;
esac
if [ "$attempt" = 5 ]; then
echo "::error::POST /releases still failing after 5 attempts. Tag ${TAG} is pushed but has no release."
echo "::error::Re-run this workflow - the plan step detects the orphan tag and republishes it."
exit 1
fi
sleep $(( attempt * 5 ))
done
if [ -z "$REL_ID" ] || [ "$REL_ID" = "null" ]; then
echo "::error::Release created but no id came back; refusing to upload assets blind."
exit 1
fi
echo "Created release ${TAG} (id=${REL_ID})" echo "Created release ${TAG} (id=${REL_ID})"
for f in "${BIN}-linux-x86_64" "${BIN}-linux-aarch64" "${BIN}-windows-x86_64.exe" SHA256SUMS; do for f in "${BIN}-linux-x86_64" "${BIN}-linux-aarch64" "${BIN}-windows-x86_64.exe" SHA256SUMS; do
curl -sSf -X POST "${API}/releases/${REL_ID}/assets?name=${f}" \ # Same treatment. An upload that fails quietly leaves a release whose
# SHA256SUMS does not cover every binary it advertises, which is worse
# than no release at all -- that file IS the trust anchor.
HTTP="$(curl -s -o /tmp/asset.json -w '%{http_code}' -X POST "${API}/releases/${REL_ID}/assets?name=${f}" \
-H "Authorization: token ${CI_TOKEN}" \ -H "Authorization: token ${CI_TOKEN}" \
-F "attachment=@dist/${f}" >/dev/null -F "attachment=@dist/${f}" || echo 000)"
if [ "$HTTP" != "201" ] && [ "$HTTP" != "200" ]; then
echo "::error::uploading ${f} returned HTTP ${HTTP}"
cat /tmp/asset.json || true
exit 1
fi
echo " uploaded ${f}" echo " uploaded ${f}"
done done

View File

@@ -12,11 +12,34 @@ ServUO plugin (C#, net48) ──loopback TCP, newline-JSON──► Rust sidec
The shard never speaks WebSocket and exposes no port of its own — the sidecar is the only The shard never speaks WebSocket and exposes no port of its own — the sidecar is the only
network-facing component, which is what keeps the game unreachable from the internet. network-facing component, which is what keeps the game unreachable from the internet.
## Running a shard? Don't build this
The [**Runic Gateway installer**](https://gitea.whitlocktech.com/RunicGateway/installer) installs
this sidecar for you — the released binary, its config, a hardened service account and the service
registration — alongside the shard plugin, in one run, on Linux or Windows:
```bash
sudo ./runicgateway-installer-linux-x86_64 install
```
It ends by printing the base URL, WebSocket URL, protocol version and auth token to paste into
**Admin → Shard** on your site. Guide:
[installer/INSTALL.md](https://gitea.whitlocktech.com/RunicGateway/docs/src/branch/main/installer/INSTALL.md).
Installing it yourself is supported too — the release binaries on this repo's
[releases page](https://gitea.whitlocktech.com/RunicGateway/link/releases) are the same ones the
installer fetches, and
[INSTALL.md Appendix A3A4](https://gitea.whitlocktech.com/RunicGateway/docs/src/branch/main/installer/INSTALL.md#a3-install-the-sidecar)
covers placing the binary and registering the service by hand.
Everything below this line is for **developing on the sidecar**.
## Related repos ## Related repos
| Repo | What | | Repo | What |
|------|------| |------|------|
| **this**`RunicGateway/link` | The Rust sidecar (`sidecar/`). | | **this**`RunicGateway/link` | The Rust sidecar (`sidecar/`). |
| [RunicGateway/installer](https://gitea.whitlocktech.com/RunicGateway/installer) | The **installer** — deploys this sidecar and the plugin onto a shard host. The supported way to set one up. |
| [RunicGateway/servuo-plugins](https://gitea.whitlocktech.com/RunicGateway/servuo-plugins) | The **C# ServUO plugin** — the shard side of the bridge (`overlay/`, `patches/`, `deploy.ps1`, test scaffolding). | | [RunicGateway/servuo-plugins](https://gitea.whitlocktech.com/RunicGateway/servuo-plugins) | The **C# ServUO plugin** — the shard side of the bridge (`overlay/`, `patches/`, `deploy.ps1`, test scaffolding). |
| [RunicGateway/docs](https://gitea.whitlocktech.com/RunicGateway/docs) | All project documentation — design docs, protocol spec, integration guide, research. | | [RunicGateway/docs](https://gitea.whitlocktech.com/RunicGateway/docs) | All project documentation — design docs, protocol spec, integration guide, research. |
@@ -28,9 +51,10 @@ network-facing component, which is what keeps the game unreachable from the inte
| `.gitea/workflows/pr-checks.yml` | Gates every PR into `main` on `cargo fmt --check`, `cargo clippy -D warnings`, and `cargo test`. | | `.gitea/workflows/pr-checks.yml` | Gates every PR into `main` on `cargo fmt --check`, `cargo clippy -D warnings`, and `cargo test`. |
| `.gitea/workflows/release.yml` | Builds + releases the sidecar binary (Linux + Windows) on every merge to `main`. | | `.gitea/workflows/release.yml` | Builds + releases the sidecar binary (Linux + Windows) on every merge to `main`. |
## Build & run ## Build & run (development)
The sidecar is a standard cargo crate: Building from source is for working *on* the sidecar; a deployment gets its binary from a release,
via the installer or by hand. The sidecar is a standard cargo crate:
```bash ```bash
cd sidecar cd sidecar
@@ -39,7 +63,7 @@ cp sidecar.toml.example sidecar.toml # then edit
cargo run --release cargo run --release
``` ```
Deploying it rather than developing on it: `--config <PATH>` names the config file (as does Deploying it by hand rather than developing on it: `--config <PATH>` names the config file (as does
`$UOLINK_CONFIG`), and `--print-config` prints the resolved settings — **including the auth token `$UOLINK_CONFIG`), and `--print-config` prints the resolved settings — **including the auth token
the website needs** — as JSON, provisioning the config file on first run. That is the supported way the website needs** — as JSON, provisioning the config file on first run. That is the supported way
to read the token back; it is not meant to be scraped from the log. to read the token back; it is not meant to be scraped from the log.
@@ -48,6 +72,29 @@ to read the token back; it is not meant to be scraped from the log.
uo-link-sidecar --print-config --config /etc/runicgateway/sidecar.toml uo-link-sidecar --print-config --config /etc/runicgateway/sidecar.toml
``` ```
### Running as a service
The same binary runs in the foreground and as a system service — there is no `--service` flag to
remember, because the process can tell how it was started.
- **Linux/systemd** supervises any foreground process, so the unit just runs the binary. `SIGTERM`
(what `systemctl stop` sends) and `SIGINT` both unwind it cleanly; logs go to the journal.
- **Windows** cannot. The service control manager only supervises a process that connects back to
it within ~30 seconds via `StartServiceCtrlDispatcher`; a plain console program registered with
`sc.exe create` is killed with **error 1053** despite running perfectly. So on Windows the sidecar
speaks that handshake: started by the SCM it runs as a service, started from a shell the connect
fails with `ERROR_FAILED_SERVICE_CONTROLLER_CONNECT` and it falls through to an ordinary
foreground run. It reports `Running` only once the shard port is bound and the store is open, and
— having no console — logs to `uo-link-sidecar.<date>.log` beside its config, rolled daily.
Only the starting and stopping is platform-specific: `src/app.rs` is the entire sidecar and is
shared, while `src/windows.rs` and `src/unix.rs` do nothing but start it and tell it when to stop.
The Windows crates are declared under `[target.'cfg(windows)'.dependencies]`, so Cargo neither
resolves nor builds them for a Linux target.
Registering the service is the installer's job; to do it by hand see
[INSTALL.md Appendix A4](https://gitea.whitlocktech.com/RunicGateway/docs/src/branch/main/installer/INSTALL.md).
`.gitea/workflows/release.yml` cross-compiles Linux + Windows binaries and cuts a Gitea release on `.gitea/workflows/release.yml` cross-compiles Linux + Windows binaries and cuts a Gitea release on
every merge to `main` (conventional-commit versioning). See [`sidecar/README.md`](sidecar/README.md) every merge to `main` (conventional-commit versioning). See [`sidecar/README.md`](sidecar/README.md)
for configuration and the wire protocol. for configuration and the wire protocol.

95
sidecar/Cargo.lock generated
View File

@@ -242,6 +242,15 @@ version = "2.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853" checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853"
[[package]]
name = "crossbeam-channel"
version = "0.5.16"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d85363c37faeca707aef026efa9f3b34d077bce547e48f770770625c6013679e"
dependencies = [
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-queue" name = "crossbeam-queue"
version = "0.3.13" version = "0.3.13"
@@ -284,6 +293,12 @@ dependencies = [
"zeroize", "zeroize",
] ]
[[package]]
name = "deranged"
version = "0.5.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7cd812cc2bc1d69d4764bd80df88b4317eaef9e773c75226407d9bc0876b211c"
[[package]] [[package]]
name = "digest" name = "digest"
version = "0.10.7" version = "0.10.7"
@@ -921,6 +936,12 @@ dependencies = [
"zeroize", "zeroize",
] ]
[[package]]
name = "num-conv"
version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "521739c6d2bac4aa25192232afe6841231376b2b26d4d9fae5ecf8ca5772e441"
[[package]] [[package]]
name = "num-integer" name = "num-integer"
version = "0.1.46" version = "0.1.46"
@@ -1048,6 +1069,12 @@ dependencies = [
"zerovec", "zerovec",
] ]
[[package]]
name = "powerfmt"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391"
[[package]] [[package]]
name = "ppv-lite86" name = "ppv-lite86"
version = "0.2.21" version = "0.2.21"
@@ -1565,6 +1592,12 @@ version = "2.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292"
[[package]]
name = "symlink"
version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a7973cce6668464ea31f176d85b13c7ab3bba2cb3b77a2ed26abd7801688010a"
[[package]] [[package]]
name = "syn" name = "syn"
version = "2.0.118" version = "2.0.118"
@@ -1642,6 +1675,36 @@ dependencies = [
"cfg-if", "cfg-if",
] ]
[[package]]
name = "time"
version = "0.3.55"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cdb87b95ec50ddfa440816d227a17b2ccbdda963a316a727fda0fc4334f7d134"
dependencies = [
"deranged",
"num-conv",
"powerfmt",
"serde_core",
"time-core",
"time-macros",
]
[[package]]
name = "time-core"
version = "0.1.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9e1c906769ad99c88eaa54e728060edef082f8e358ff32030cb7c7d315e81109"
[[package]]
name = "time-macros"
version = "0.2.32"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7e689342a48d2ea927c87ea50cabf8594854bf940e9310208848d680d668ed85"
dependencies = [
"num-conv",
"time-core",
]
[[package]] [[package]]
name = "tinystr" name = "tinystr"
version = "0.8.3" version = "0.8.3"
@@ -1798,6 +1861,19 @@ dependencies = [
"tracing-core", "tracing-core",
] ]
[[package]]
name = "tracing-appender"
version = "0.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "050686193eb999b4bb3bc2acfa891a13da00f79734704c4b8b4ef1a10b368a3c"
dependencies = [
"crossbeam-channel",
"symlink",
"thiserror 2.0.18",
"time",
"tracing-subscriber",
]
[[package]] [[package]]
name = "tracing-attributes" name = "tracing-attributes"
version = "0.1.31" version = "0.1.31"
@@ -1913,7 +1989,9 @@ dependencies = [
"tokio", "tokio",
"toml", "toml",
"tracing", "tracing",
"tracing-appender",
"tracing-subscriber", "tracing-subscriber",
"windows-service",
] ]
[[package]] [[package]]
@@ -2025,6 +2103,12 @@ dependencies = [
"wasite", "wasite",
] ]
[[package]]
name = "widestring"
version = "1.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "72069c3113ab32ab29e5584db3c6ec55d416895e60715417b5b883a357c3e471"
[[package]] [[package]]
name = "windows-core" name = "windows-core"
version = "0.62.2" version = "0.62.2"
@@ -2075,6 +2159,17 @@ dependencies = [
"windows-link", "windows-link",
] ]
[[package]]
name = "windows-service"
version = "0.8.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "857224b3b211c6f3616921f081ee54721ee3ad2ace2fac6a6337e032f7b4dcf2"
dependencies = [
"bitflags",
"widestring",
"windows-sys 0.61.2",
]
[[package]] [[package]]
name = "windows-strings" name = "windows-strings"
version = "0.5.1" version = "0.5.1"

View File

@@ -17,5 +17,12 @@ toml = "0.8"
getrandom = "0.2" getrandom = "0.2"
chrono = { version = "0.4", default-features = false, features = ["std", "clock"] } chrono = { version = "0.4", default-features = false, features = ["std", "clock"] }
# Speaking the Windows Service Control Manager's startup handshake, and logging somewhere other
# than the stdout a service does not have. Declared per target so Cargo neither resolves nor builds
# either crate for Linux — the Linux binary is byte-for-byte unaffected by Windows service support.
[target.'cfg(windows)'.dependencies]
windows-service = "0.8"
tracing-appender = "0.2"
[profile.release] [profile.release]
opt-level = 2 opt-level = 2

571
sidecar/src/app.rs Normal file
View File

@@ -0,0 +1,571 @@
//! The sidecar itself: everything that happens between "we have a config path" and "we were told
//! to stop". Identical on every platform.
//!
//! This module exists so that *how the process is started and stopped* — a bare `main` under
//! systemd, or a `ServiceMain` under the Windows SCM — is the only thing that differs between
//! hosts. The shard listener, the config, the store, the web server and the event loop are shared
//! code with no `#[cfg]` in sight.
//!
//! [`run`] is parameterised on the two things a supervisor cares about:
//!
//! - `ready` is called once the sidecar is actually up (listener bound, store open). The Windows
//! service reports `Running` to the SCM there, so a config or bind failure surfaces as a *start*
//! failure rather than a service that reports Running and then dies.
//! - `shutdown` is whatever "stop" means on this host: Ctrl-C and `SIGTERM` on Unix, the SCM's
//! `Stop` control on Windows.
use std::collections::HashMap;
use std::future::Future;
use std::sync::atomic::AtomicI64;
use std::sync::Arc;
use std::time::Instant;
use serde_json::Value;
use tokio::sync::{broadcast, mpsc};
use tracing::info;
use crate::{config, rpc, shard, store, web};
/// Runs the sidecar until `shutdown` resolves.
///
/// `config_path` is the `--config` argument, or `None` to resolve `$UOLINK_CONFIG` and the default
/// as usual.
pub async fn run<R, S>(config_path: Option<&str>, ready: R, shutdown: S) -> anyhow::Result<()>
where
R: FnOnce(),
S: Future<Output = ()>,
{
info!("uo-link sidecar starting");
let loaded = config::Config::load(config_path)?;
let cfg = loaded.cfg;
info!(
config = %loaded.path.display(),
shard = %cfg.shard.bind,
web = %cfg.web.bind,
db = %cfg.store.path,
auth = cfg.auth_required(),
"configuration loaded"
);
// Shard link: events in, commands out.
let (event_tx, mut event_rx) = mpsc::unbounded_channel::<shard::ShardEvent>();
let handle = shard::serve(&cfg.shard.bind, event_tx).await?;
// Live feed: every shard event fans out to all connected website WebSocket clients.
let (bcast_tx, _) = broadcast::channel::<String>(1024);
// Request/reply correlation for REST queries.
let rpc = rpc::Rpc::new();
// Durable store: event history, economy series, cached profiles, link map.
let store = store::Store::open(&cfg.store.path).await?;
// Health/observability state.
let started = Instant::now();
let last_event = Arc::new(AtomicI64::new(0));
// Website-facing HTTP server.
let web_state = web::AppState {
events: bcast_tx.clone(),
shard: handle.clone(),
rpc: rpc.clone(),
store: store.clone(),
token: Arc::new(cfg.web.auth_token.clone()),
started,
last_event: last_event.clone(),
};
let web_bind = cfg.web.bind.clone();
tokio::spawn(async move {
if let Err(e) = web::serve(&web_bind, web_state).await {
tracing::error!(error = %e, "web server exited");
}
});
// Event loop: a line that correlates to a pending REST call is a reply — route it to the
// waiting caller and stop. Everything else is a live event: log it, persist it, broadcast it.
let feed_tx = bcast_tx.clone();
let route_rpc = rpc.clone();
let event_store = store.clone();
let last_event_ts = last_event.clone();
let replay_handle = handle.clone(); // re-push external news to the shard on (re)connect
let mut total: u64 = 0;
// Partly-received guild rosters, keyed by guild id. Lives in the event-loop task, so it needs
// no lock and dies with the loop. See `accumulate_roster`.
let mut roster_parts: HashMap<i64, RosterParts> = HashMap::new();
tokio::spawn(async move {
while let Some(ev) = event_rx.recv().await {
// Any line from the shard — including pong heartbeats — is a sign of life.
last_event_ts.store(now_ms(), std::sync::atomic::Ordering::Relaxed);
if route_rpc.try_route(&ev.value).await {
continue; // consumed as a reply
}
total += 1;
match ev.kind.as_str() {
"server.hello" | "mob.login" | "mob.logout" | "player.death" | "vendor.sale"
| "house.decay" | "link.request" => {
info!(kind = %ev.kind, n = total, "{}", ev.value);
}
_ => tracing::debug!(kind = %ev.kind, n = total, "{}", ev.value),
}
// Persist, then broadcast. `pong` and `ws.hello` are ephemeral chatter, not history.
if ev.kind != "pong" {
let t = ev
.value
.get("t")
.and_then(|v| v.as_i64())
.unwrap_or_else(now_ms);
let text = ev.value.to_string();
if let Err(e) = event_store.insert_event(t, &ev.kind, &text).await {
tracing::warn!(error = %e, "failed to persist event");
}
// The champ board is a live projection: champ.update folds in the latest state (one
// row per spawn), champ.remove drops a spawn that despawned or was slain.
match ev.kind.as_str() {
"champ.update" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store
.upsert_champ(
serial,
ev.value.get("status").and_then(|s| s.as_str()),
ev.value.get("name").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert champ board");
}
}
}
"champ.remove" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store.delete_champ(serial).await {
tracing::warn!(error = %e, "failed to remove champ board row");
}
}
}
// Guild board (Protocol 2.0): guild.update folds in the latest roster (one row
// per guild id); guild.remove drops a disbanded guild.
"guild.update" => {
if let Some(id) = ev.value.get("id").and_then(|v| v.as_i64()) {
if let Err(e) = event_store
.upsert_guild(
id,
ev.value.get("name").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert guild board");
}
}
}
// Guild roster (Protocol 4): the member list that `guild.update`'s counts cannot
// express. It writes a *different column* of the same row, so it never races
// guild.update. `guild.leave` deliberately has no arm here — the sidecar
// forwards it (persisted and broadcast below, like any event) and the board's
// roster self-corrects on the next `guild.roster`, which the shard re-emits
// whenever the member set changes. Keeping the delta out of the board is what
// keeps the sidecar a forwarder rather than a thing that maintains state.
//
// A roster over the shard's per-line cap arrives as several frames, so it is
// reassembled before it is stored — see `accumulate_roster` for why that happens
// here rather than by appending to the column.
"guild.roster" => {
if let Some(id) = ev.value.get("id").and_then(|v| v.as_i64()) {
let seq = ev.value.get("seq").and_then(|v| v.as_i64()).unwrap_or(0);
let more = ev
.value
.get("more")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let members = ev
.value
.get("members")
.and_then(|m| m.as_array())
.cloned()
.unwrap_or_default();
if let Some(complete) =
accumulate_roster(&mut roster_parts, id, seq, more, members)
{
let json = serde_json::Value::Array(complete).to_string();
if let Err(e) = event_store.upsert_guild_roster(id, &json, t).await
{
tracing::warn!(error = %e, "failed to upsert guild roster");
}
}
}
}
"guild.remove" => {
if let Some(id) = ev.value.get("id").and_then(|v| v.as_i64()) {
if let Err(e) = event_store.delete_guild(id).await {
tracing::warn!(error = %e, "failed to remove guild board row");
}
}
}
// Governor board (Protocol 2.0): city.update folds in each city's latest
// governance state (one row per city).
"city.update" => {
if let Some(city) = ev.value.get("city").and_then(|c| c.as_str()) {
if let Err(e) = event_store.upsert_governor(city, &text, t).await {
tracing::warn!(error = %e, "failed to upsert governor board");
}
}
}
// House registry (Protocol 2.0): house.update folds in each house's latest state
// (one row per serial); house.remove drops a demolished/traded house.
"house.update" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store
.upsert_house(
serial,
ev.value.get("name").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert house registry");
}
}
}
"house.remove" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store.delete_house(serial).await {
tracing::warn!(error = %e, "failed to remove house registry row");
}
}
}
// Points/loyalty boards (Protocol 3.0): one row per point system, keyed by
// the shard's own PointsType name. The plugin only emits a system whose top N
// actually moved, so this is a sparse stream of overwrites — and there is no
// `points.remove` to handle, because the shard's set of systems is fixed at
// startup and cannot shrink.
"points.board" => {
if let Some(system) = ev.value.get("system").and_then(|s| s.as_str()) {
if let Err(e) = event_store
.upsert_points_board(
system,
ev.value.get("nameString").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert points board");
}
}
}
// Player-vendor market index (Protocol 3.0). Each frame is authoritative for
// one vendor — the shard's round-robin sweep only emits a shop whose contents,
// prices or location actually moved — so this is a whole-row overwrite.
//
// Unlike the boards above there IS a remove: a vendor is dismissed, expires, or
// its owner switches off the in-game Vendor Search flag, and any of those must
// take the shop off the site. The last of the three is a privacy control, so
// dropping the row promptly is the point rather than housekeeping.
"vendor.listing" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
let loc = ev.value.get("location");
let field = |k: &str| loc.and_then(|l| l.get(k));
if let Err(e) = event_store
.upsert_vendor(
serial,
ev.value.get("shopName").and_then(|v| v.as_str()),
ev.value.get("ownerName").and_then(|v| v.as_str()),
field("map").and_then(|v| v.as_str()),
field("x").and_then(|v| v.as_i64()),
field("y").and_then(|v| v.as_i64()),
field("region").and_then(|v| v.as_str()),
ev.value.get("count").and_then(|v| v.as_i64()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert vendor listing");
}
}
}
"vendor.listing.remove" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store.delete_vendor(serial).await {
tracing::warn!(error = %e, "failed to remove vendor listing");
}
}
}
// Shard ruleset (Protocol 3.0): a singleton projection. The shard re-emits
// world.ruleset on every connect, so this row is simply overwritten; `rev`
// lets a reader tell a re-send from an actual config change.
"world.ruleset" => {
if let Err(e) = event_store
.upsert_ruleset(ev.value.get("rev").and_then(|r| r.as_str()), &text, t)
.await
{
tracing::warn!(error = %e, "failed to upsert ruleset");
}
}
_ => {}
}
}
// On a shard (re)connect, re-push the stored external news: the shard rebuilds
// TownCryerSystem.NewsEntries from scratch each boot and does not persist ours. Replay
// with announce=false so a restart does not re-proclaim every article at once. news.add
// is idempotent by id, so replaying to a still-populated shard is harmless.
if ev.kind == "server.hello" {
// A (re)connected shard restarts every roster from `seq` 0, so any half-received
// one belongs to the previous connection and can never be completed.
roster_parts.clear();
match event_store.news_all().await {
Ok(items) => {
for mut item in items {
if let Some(obj) = item.as_object_mut() {
obj.insert("announce".to_string(), serde_json::json!(false));
}
if !replay_handle.send(item.to_string()).await {
break; // shard went away mid-replay
}
}
}
Err(e) => tracing::warn!(error = %e, "news replay: could not read stored news"),
}
}
let _ = feed_tx.send(ev.value.to_string());
}
});
// Heartbeat to the shard, exercising the command path.
let ping_handle = handle.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(15)).await;
if ping_handle.is_connected().await {
let _ = ping_handle
.send(r#"{"kind":"ping","id":"sidecar-heartbeat"}"#.to_string())
.await;
}
}
});
// Everything that can fail at startup has now either failed or succeeded: the shard port is
// bound and the store is open. A supervisor may call this "started".
ready();
shutdown.await;
info!("shutting down");
Ok(())
}
pub fn now_ms() -> i64 {
use std::time::{SystemTime, UNIX_EPOCH};
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
}
/// A guild roster that has arrived in part: the `seq` expected next, and what has accumulated.
struct RosterParts {
next_seq: i64,
members: Vec<Value>,
}
/// Refuses to accumulate a roster past this many members. The shard caps its own frames, so
/// exceeding this means a shard that is buggy or not what it claims to be — and the one thing this
/// buffer must not do is grow without bound on its say-so.
const MAX_ROSTER_MEMBERS: usize = 50_000;
/// Reassembles a `guild.roster` that the shard split across frames, returning the whole member list
/// once the final frame arrives and `None` while one is still incomplete.
///
/// Reassembly happens **here, in memory, before the store** rather than by appending to the
/// `members` column, for two reasons. Appending would make the write a read-modify-write — the exact
/// thing the two-column board design exists to avoid — and it would publish a torn roster: a reader
/// hitting `GET /guilds` between frames would see a partial member list as though it were the truth.
/// Buffering keeps the store's write a single atomic upsert of a complete roster.
///
/// This is transport-level reassembly, not domain state: it is the same category of work as turning
/// bytes into a line, and it holds nothing once a roster is complete. That is what keeps it
/// compatible with the sidecar being a forwarder.
///
/// The ordinary case — a guild inside the shard's per-line cap, which is every realistic one —
/// arrives as `seq` 0 with `more` false and is returned immediately without ever touching the map.
fn accumulate_roster(
parts: &mut HashMap<i64, RosterParts>,
id: i64,
seq: i64,
more: bool,
members: Vec<Value>,
) -> Option<Vec<Value>> {
if seq == 0 {
// A fresh roster supersedes any partial one: the shard restarts at 0 every time it emits,
// so a leftover buffer is from an emission that was interrupted and will never finish.
parts.remove(&id);
if !more {
return Some(members);
}
parts.insert(
id,
RosterParts {
next_seq: 1,
members,
},
);
return None;
}
let entry = match parts.get_mut(&id) {
Some(entry) => entry,
// A continuation with nothing to continue: the sidecar started, or the shard reconnected,
// midway through an emission. Dropping it is right — the next full roster is complete.
None => {
tracing::debug!(
guild = id,
seq,
"roster continuation with no start; ignoring"
);
return None;
}
};
if entry.next_seq != seq {
tracing::warn!(
guild = id,
expected = entry.next_seq,
got = seq,
"roster frames out of order; discarding the partial roster"
);
parts.remove(&id);
return None;
}
entry.members.extend(members);
if entry.members.len() > MAX_ROSTER_MEMBERS {
tracing::warn!(
guild = id,
len = entry.members.len(),
"roster exceeded the reassembly cap; discarding"
);
parts.remove(&id);
return None;
}
if more {
entry.next_seq = seq + 1;
return None;
}
parts.remove(&id).map(|done| done.members)
}
#[cfg(test)]
mod tests {
use super::*;
fn members(names: &[&str]) -> Vec<Value> {
names
.iter()
.map(|n| serde_json::json!({"name": n}))
.collect()
}
fn names(vs: &[Value]) -> Vec<String> {
vs.iter()
.map(|v| v["name"].as_str().unwrap_or_default().to_string())
.collect()
}
#[test]
fn a_single_frame_roster_is_returned_immediately() {
// Every realistic guild takes this path, and it must not depend on the buffer at all.
let mut parts = HashMap::new();
let out = accumulate_roster(&mut parts, 1, 0, false, members(&["Ada", "Bo"]));
assert_eq!(names(&out.expect("complete")), ["Ada", "Bo"]);
assert!(parts.is_empty(), "nothing should be buffered");
}
#[test]
fn a_chunked_roster_reassembles_in_order() {
// The case the live rig caught: without this, only the final frame survived and a
// 155-member guild appeared on the board with 3 members.
let mut parts = HashMap::new();
assert!(accumulate_roster(&mut parts, 1, 0, true, members(&["Ada"])).is_none());
assert!(accumulate_roster(&mut parts, 1, 1, true, members(&["Bo"])).is_none());
let out = accumulate_roster(&mut parts, 1, 2, false, members(&["Cy"]));
assert_eq!(names(&out.expect("complete")), ["Ada", "Bo", "Cy"]);
assert!(parts.is_empty(), "buffer is released once complete");
}
#[test]
fn a_restarted_roster_supersedes_a_partial_one() {
// A shard that reconnects mid-emission starts again at seq 0. The abandoned frames must not
// end up spliced onto the front of the new roster.
let mut parts = HashMap::new();
assert!(accumulate_roster(&mut parts, 1, 0, true, members(&["Stale"])).is_none());
let out = accumulate_roster(&mut parts, 1, 0, false, members(&["Fresh"]));
assert_eq!(names(&out.expect("complete")), ["Fresh"]);
}
#[test]
fn an_out_of_order_frame_discards_the_partial_roster() {
// Better to publish nothing and wait for the next full emission than to store a roster with
// a hole in it that nothing downstream could detect.
let mut parts = HashMap::new();
assert!(accumulate_roster(&mut parts, 1, 0, true, members(&["Ada"])).is_none());
assert!(accumulate_roster(&mut parts, 1, 2, false, members(&["Skipped"])).is_none());
assert!(parts.is_empty());
}
#[test]
fn a_continuation_with_no_start_is_ignored() {
// The sidecar restarting midway through a shard's emission.
let mut parts = HashMap::new();
assert!(accumulate_roster(&mut parts, 1, 3, false, members(&["Orphan"])).is_none());
assert!(parts.is_empty());
}
#[test]
fn two_guilds_reassemble_independently() {
// Rosters for different guilds interleave freely — the sweep emits one guild after another
// and nothing serialises them on the wire.
let mut parts = HashMap::new();
assert!(accumulate_roster(&mut parts, 1, 0, true, members(&["A1"])).is_none());
assert!(accumulate_roster(&mut parts, 2, 0, true, members(&["B1"])).is_none());
let g2 = accumulate_roster(&mut parts, 2, 1, false, members(&["B2"]));
let g1 = accumulate_roster(&mut parts, 1, 1, false, members(&["A2"]));
assert_eq!(names(&g2.expect("guild 2")), ["B1", "B2"]);
assert_eq!(names(&g1.expect("guild 1")), ["A1", "A2"]);
}
#[test]
fn an_empty_roster_is_a_complete_roster() {
// A guild whose last member left emits one frame with an empty array. Treating that as
// "nothing to store" would leave the board showing the roster it had before.
let mut parts = HashMap::new();
let out = accumulate_roster(&mut parts, 1, 0, false, vec![]);
assert_eq!(out.expect("complete").len(), 0);
}
}

View File

@@ -1,21 +1,37 @@
//! uo-link sidecar. //! uo-link sidecar.
//! //!
//! Terminates the loopback link to the ServUO shard and exposes a website-facing HTTP surface. So //! Terminates the loopback link to the ServUO shard and exposes a website-facing HTTP surface: the
//! far: the shard link (bidirectional) and a WebSocket live feed. REST queries and SQLite come next. //! shard link (bidirectional), a WebSocket live feed, and REST queries backed by SQLite.
//!
//! # Layout
//!
//! `main` does argument handling and nothing else; the sidecar proper lives in [`app`] and is the
//! same code on every platform. Only *how the process is started and stopped* is
//! platform-specific:
//!
//! ```text
//! systemd ──▶ main ──▶ unix::run ─────────────────────────┐
//! ├──▶ app::run
//! SCM ─────▶ main ──▶ windows::run ──▶ ServiceMain ───────┘
//! └─▶ console fallback ──┘
//! ```
//!
//! The Windows half is not optional politeness: the SCM refuses to supervise a program that does
//! not speak its startup handshake (see [`windows`]). The platform modules are gated with `#[cfg]`
//! and their dependencies are declared per target, so none of it reaches a Linux build.
mod app;
mod cli; mod cli;
mod config; mod config;
mod rpc; mod rpc;
mod shard; mod shard;
mod store; mod store;
#[cfg(unix)]
mod unix;
mod web; mod web;
#[cfg(windows)]
mod windows;
use std::sync::atomic::AtomicI64;
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::{broadcast, mpsc};
use tracing::info;
use tracing_subscriber::EnvFilter; use tracing_subscriber::EnvFilter;
/// Wire-protocol version between the website and the sidecar. Bump this whenever an event or /// Wire-protocol version between the website and the sidecar. Bump this whenever an event or
@@ -30,10 +46,38 @@ use tracing_subscriber::EnvFilter;
/// `vendor.listing.remove`, with the `GET /ruleset`, `/points` and `/market` reads that serve them /// `vendor.listing.remove`, with the `GET /ruleset`, `/points` and `/market` reads that serve them
/// from the store. Same shape as the v2 bump — the kinds are additive, the endpoints are not — and /// from the store. Same shape as the v2 bump — the kinds are additive, the endpoints are not — and
/// there is deliberately no feature-negotiation array: v3 implies all three kinds. /// there is deliberately no feature-negotiation array: v3 implies all three kinds.
pub const PROTOCOL_VERSION: u32 = 3; ///
/// v4 (Protocol 4): adds `guild.roster` and `guild.leave`, giving the guild board a real member list
/// instead of the member *count* that was all v2 could express. Additive in the same way again — the
/// kinds are new, `GET /guilds` grows a `roster` key, and nothing existing changed shape. This is the
/// first bump that also needed a **store migration** (`guilds.members`), because it is the first to
/// add a column to a table that already exists rather than a whole new table; see `store::migrate`.
///
/// v5 (Protocol 5): three enrichments that are additive in the same way again, bumped together
/// rather than one at a time because a protocol bump is not cheap here — it costs a sidecar
/// release, a republished bundle and an operator update on every shard, so a field left out costs
/// a whole second round of that rather than a follow-up commit. They are:
///
/// * `house.decay` gains `ownerName` and a decay SCHEDULE — `nextStage`, `decayPeriodSec`,
/// `dynamicDecay`, and `estimatedCollapse` only where it is exactly knowable (at IDOC under
/// dynamic decay; at any stage under static decay, which has no randomness to wait out).
/// * `vendor.listing` gains `ownerAcct` — without which the frame names an owner nobody can
/// resolve to a person — and a `fees` object carrying the charge, the funds, the pay interval
/// and the resolved `dismissalAt`.
/// * `account.login.result` is a NEW kind: the verdict of a login, which the pre-existing
/// `account.login.attempt` structurally cannot carry (its EventSink fires before the auth
/// decision is made).
///
/// **No store migration this time**, unlike v4. Every frame is persisted whole and the board tables
/// index only the columns they already had, so the new fields ride inside the stored JSON and the
/// new kind lands in `events` like any other. That is the dumb-forwarder property doing its job:
/// the sidecar defines no schema for a frame's contents and so needs no change when they grow.
pub const PROTOCOL_VERSION: u32 = 5;
#[tokio::main] // Not `#[tokio::main]`: on Windows the SCM dispatcher takes over this thread and starts the runtime
async fn main() -> anyhow::Result<()> { // itself, on its own thread, once the service actually begins. The runtime is built by whichever
// platform module ends up running.
fn main() -> anyhow::Result<()> {
let args = match cli::parse(std::env::args().skip(1)) { let args = match cli::parse(std::env::args().skip(1)) {
Ok(args) => args, Ok(args) => args,
Err(msg) => { Err(msg) => {
@@ -66,299 +110,15 @@ async fn main() -> anyhow::Result<()> {
cli::Mode::Run => {} cli::Mode::Run => {}
} }
init_tracing(); #[cfg(windows)]
info!("uo-link sidecar starting"); return windows::run(args.config.as_deref());
let loaded = config::Config::load(args.config.as_deref())?; #[cfg(unix)]
let cfg = loaded.cfg; return unix::run(args.config.as_deref());
info!(
config = %loaded.path.display(),
shard = %cfg.shard.bind,
web = %cfg.web.bind,
db = %cfg.store.path,
auth = cfg.auth_required(),
"configuration loaded"
);
// Shard link: events in, commands out.
let (event_tx, mut event_rx) = mpsc::unbounded_channel::<shard::ShardEvent>();
let handle = shard::serve(&cfg.shard.bind, event_tx).await?;
// Live feed: every shard event fans out to all connected website WebSocket clients.
let (bcast_tx, _) = broadcast::channel::<String>(1024);
// Request/reply correlation for REST queries.
let rpc = rpc::Rpc::new();
// Durable store: event history, economy series, cached profiles, link map.
let store = store::Store::open(&cfg.store.path).await?;
// Health/observability state.
let started = Instant::now();
let last_event = Arc::new(AtomicI64::new(0));
// Website-facing HTTP server.
let web_state = web::AppState {
events: bcast_tx.clone(),
shard: handle.clone(),
rpc: rpc.clone(),
store: store.clone(),
token: Arc::new(cfg.web.auth_token.clone()),
started,
last_event: last_event.clone(),
};
let web_bind = cfg.web.bind.clone();
tokio::spawn(async move {
if let Err(e) = web::serve(&web_bind, web_state).await {
tracing::error!(error = %e, "web server exited");
}
});
// Event loop: a line that correlates to a pending REST call is a reply — route it to the
// waiting caller and stop. Everything else is a live event: log it, persist it, broadcast it.
let feed_tx = bcast_tx.clone();
let route_rpc = rpc.clone();
let event_store = store.clone();
let last_event_ts = last_event.clone();
let replay_handle = handle.clone(); // re-push external news to the shard on (re)connect
let mut total: u64 = 0;
tokio::spawn(async move {
while let Some(ev) = event_rx.recv().await {
// Any line from the shard — including pong heartbeats — is a sign of life.
last_event_ts.store(now_ms(), std::sync::atomic::Ordering::Relaxed);
if route_rpc.try_route(&ev.value).await {
continue; // consumed as a reply
} }
total += 1; /// Logging for a foreground run: human-readable, on stdout.
match ev.kind.as_str() { pub fn init_console_tracing() {
"server.hello" | "mob.login" | "mob.logout" | "player.death" | "vendor.sale"
| "house.decay" | "link.request" => {
info!(kind = %ev.kind, n = total, "{}", ev.value);
}
_ => tracing::debug!(kind = %ev.kind, n = total, "{}", ev.value),
}
// Persist, then broadcast. `pong` and `ws.hello` are ephemeral chatter, not history.
if ev.kind != "pong" {
let t = ev
.value
.get("t")
.and_then(|v| v.as_i64())
.unwrap_or_else(now_ms);
let text = ev.value.to_string();
if let Err(e) = event_store.insert_event(t, &ev.kind, &text).await {
tracing::warn!(error = %e, "failed to persist event");
}
// The champ board is a live projection: champ.update folds in the latest state (one
// row per spawn), champ.remove drops a spawn that despawned or was slain.
match ev.kind.as_str() {
"champ.update" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store
.upsert_champ(
serial,
ev.value.get("status").and_then(|s| s.as_str()),
ev.value.get("name").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert champ board");
}
}
}
"champ.remove" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store.delete_champ(serial).await {
tracing::warn!(error = %e, "failed to remove champ board row");
}
}
}
// Guild board (Protocol 2.0): guild.update folds in the latest roster (one row
// per guild id); guild.remove drops a disbanded guild.
"guild.update" => {
if let Some(id) = ev.value.get("id").and_then(|v| v.as_i64()) {
if let Err(e) = event_store
.upsert_guild(
id,
ev.value.get("name").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert guild board");
}
}
}
"guild.remove" => {
if let Some(id) = ev.value.get("id").and_then(|v| v.as_i64()) {
if let Err(e) = event_store.delete_guild(id).await {
tracing::warn!(error = %e, "failed to remove guild board row");
}
}
}
// Governor board (Protocol 2.0): city.update folds in each city's latest
// governance state (one row per city).
"city.update" => {
if let Some(city) = ev.value.get("city").and_then(|c| c.as_str()) {
if let Err(e) = event_store.upsert_governor(city, &text, t).await {
tracing::warn!(error = %e, "failed to upsert governor board");
}
}
}
// House registry (Protocol 2.0): house.update folds in each house's latest state
// (one row per serial); house.remove drops a demolished/traded house.
"house.update" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store
.upsert_house(
serial,
ev.value.get("name").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert house registry");
}
}
}
"house.remove" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store.delete_house(serial).await {
tracing::warn!(error = %e, "failed to remove house registry row");
}
}
}
// Points/loyalty boards (Protocol 3.0): one row per point system, keyed by
// the shard's own PointsType name. The plugin only emits a system whose top N
// actually moved, so this is a sparse stream of overwrites — and there is no
// `points.remove` to handle, because the shard's set of systems is fixed at
// startup and cannot shrink.
"points.board" => {
if let Some(system) = ev.value.get("system").and_then(|s| s.as_str()) {
if let Err(e) = event_store
.upsert_points_board(
system,
ev.value.get("nameString").and_then(|n| n.as_str()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert points board");
}
}
}
// Player-vendor market index (Protocol 3.0). Each frame is authoritative for
// one vendor — the shard's round-robin sweep only emits a shop whose contents,
// prices or location actually moved — so this is a whole-row overwrite.
//
// Unlike the boards above there IS a remove: a vendor is dismissed, expires, or
// its owner switches off the in-game Vendor Search flag, and any of those must
// take the shop off the site. The last of the three is a privacy control, so
// dropping the row promptly is the point rather than housekeeping.
"vendor.listing" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
let loc = ev.value.get("location");
let field = |k: &str| loc.and_then(|l| l.get(k));
if let Err(e) = event_store
.upsert_vendor(
serial,
ev.value.get("shopName").and_then(|v| v.as_str()),
ev.value.get("ownerName").and_then(|v| v.as_str()),
field("map").and_then(|v| v.as_str()),
field("x").and_then(|v| v.as_i64()),
field("y").and_then(|v| v.as_i64()),
field("region").and_then(|v| v.as_str()),
ev.value.get("count").and_then(|v| v.as_i64()),
&text,
t,
)
.await
{
tracing::warn!(error = %e, "failed to upsert vendor listing");
}
}
}
"vendor.listing.remove" => {
if let Some(serial) = ev.value.get("serial").and_then(|s| s.as_str()) {
if let Err(e) = event_store.delete_vendor(serial).await {
tracing::warn!(error = %e, "failed to remove vendor listing");
}
}
}
// Shard ruleset (Protocol 3.0): a singleton projection. The shard re-emits
// world.ruleset on every connect, so this row is simply overwritten; `rev`
// lets a reader tell a re-send from an actual config change.
"world.ruleset" => {
if let Err(e) = event_store
.upsert_ruleset(ev.value.get("rev").and_then(|r| r.as_str()), &text, t)
.await
{
tracing::warn!(error = %e, "failed to upsert ruleset");
}
}
_ => {}
}
}
// On a shard (re)connect, re-push the stored external news: the shard rebuilds
// TownCryerSystem.NewsEntries from scratch each boot and does not persist ours. Replay
// with announce=false so a restart does not re-proclaim every article at once. news.add
// is idempotent by id, so replaying to a still-populated shard is harmless.
if ev.kind == "server.hello" {
match event_store.news_all().await {
Ok(items) => {
for mut item in items {
if let Some(obj) = item.as_object_mut() {
obj.insert("announce".to_string(), serde_json::json!(false));
}
if !replay_handle.send(item.to_string()).await {
break; // shard went away mid-replay
}
}
}
Err(e) => tracing::warn!(error = %e, "news replay: could not read stored news"),
}
}
let _ = feed_tx.send(ev.value.to_string());
}
});
// Heartbeat to the shard, exercising the command path.
let ping_handle = handle.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(15)).await;
if ping_handle.is_connected().await {
let _ = ping_handle
.send(r#"{"kind":"ping","id":"sidecar-heartbeat"}"#.to_string())
.await;
}
}
});
tokio::signal::ctrl_c().await?;
info!("shutting down");
Ok(())
}
fn now_ms() -> i64 {
use std::time::{SystemTime, UNIX_EPOCH};
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
}
fn init_tracing() {
tracing_subscriber::fmt() tracing_subscriber::fmt()
.with_env_filter( .with_env_filter(
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")), EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")),

View File

@@ -45,6 +45,7 @@ impl Store {
.await?; .await?;
sqlx::query(SCHEMA).execute(&pool).await?; sqlx::query(SCHEMA).execute(&pool).await?;
migrate(&pool).await?;
info!(%path, "store ready"); info!(%path, "store ready");
Ok(Self { pool }) Ok(Self { pool })
} }
@@ -236,12 +237,61 @@ impl Store {
Ok(()) Ok(())
} }
/// Upserts one guild's member roster (Protocol 4), keyed by guild id, touching **only** the
/// `members` column.
///
/// Deliberately not a write to `json`. That column holds the verbatim `guild.update` line, and a
/// roster arriving as its own event must not clobber the snapshot — name, abbreviation, leader,
/// online count — that `guild.update` owns. Splitting the two writers across two columns of one
/// row is what lets both be plain upserts: neither needs to read the other's value first, so
/// there is no read-modify-write and no ordering requirement between the two kinds.
///
/// The `INSERT` half is not redundant: a roster can arrive before the first `guild.update` for a
/// guild, and the row it creates then carries `'{}'` until that update fills it in.
pub async fn upsert_guild_roster(
&self,
id: i64,
members_json: &str,
t: i64,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO guilds (id, name, json, updated_t, members) VALUES (?, NULL, '{}', ?, ?)
ON CONFLICT(id) DO UPDATE SET members = excluded.members, updated_t = excluded.updated_t",
)
.bind(id)
.bind(t)
.bind(members_json)
.execute(&self.pool)
.await?;
Ok(())
}
/// The full guild board: every guild's latest snapshot, ordered by name. /// The full guild board: every guild's latest snapshot, ordered by name.
///
/// The roster is stored in its own column (see [`Self::upsert_guild_roster`]) and folded into
/// the projected object as `roster` here, at read time. A guild that has had a `guild.update`
/// but no `guild.roster` yet simply has no `roster` key, which is the honest representation of
/// "not known" and distinct from a guild whose roster is genuinely empty.
pub async fn guilds_all(&self) -> anyhow::Result<Vec<Value>> { pub async fn guilds_all(&self) -> anyhow::Result<Vec<Value>> {
let rows = sqlx::query("SELECT json FROM guilds ORDER BY name, id") let rows = sqlx::query("SELECT json, members FROM guilds ORDER BY name, id")
.fetch_all(&self.pool) .fetch_all(&self.pool)
.await?; .await?;
Ok(parse_json_column(rows))
Ok(rows
.into_iter()
.filter_map(|r| {
let mut v: Value = serde_json::from_str(&r.get::<String, _>("json")).ok()?;
let members: Option<String> = r.get("members");
if let (Some(obj), Some(raw)) = (v.as_object_mut(), members) {
if let Ok(list) = serde_json::from_str::<Value>(&raw) {
obj.insert("roster".into(), list);
}
}
Some(v)
})
.collect())
} }
// ---- governor board (Protocol 2.0) ---- // ---- governor board (Protocol 2.0) ----
@@ -513,6 +563,75 @@ impl Store {
} }
} }
/// The schema version this build expects. Bump it, and add the matching arm to [`migrate`], for
/// every change that `SCHEMA` alone cannot make to a database that already exists.
const SCHEMA_VERSION: i64 = 1;
/// Brings an existing database forward to [`SCHEMA_VERSION`].
///
/// `SCHEMA` is `CREATE TABLE IF NOT EXISTS` only, which is enough to *add a table* but cannot add a
/// column to a table that is already there. Every schema change up to and including Protocol 3.0
/// happened to add whole tables, so this never mattered and `ALTER TABLE` appears nowhere in this
/// repo's history. `guilds.members` (Protocol 4) is the first column added to an existing table, so
/// the mechanism has to exist now.
///
/// The version counter is SQLite's own `PRAGMA user_version`: an integer in the database header, so
/// it needs no table of its own and cannot be separated from the file it describes. Each step runs
/// in a transaction **together with** the bump that records it, so a step either lands completely or
/// not at all, and an interrupted run resumes at the right place rather than re-applying half of one.
///
/// A failure here propagates and aborts startup, deliberately. A half-migrated store answers the
/// website with confusing partial data, which is worse than being plainly absent — and the shard
/// dials *out* to the sidecar, so a sidecar that refuses to start never stalls the game.
async fn migrate(pool: &SqlitePool) -> anyhow::Result<()> {
let mut version: i64 = sqlx::query_scalar("PRAGMA user_version")
.fetch_one(pool)
.await?;
// A database written by a *newer* sidecar than this binary. This is not an error: every step
// here is additive, so a newer schema has only columns and tables an older reader ignores, and
// refusing to start would turn "roll the binary back" — a recovery path — into a dead end.
if version > SCHEMA_VERSION {
tracing::warn!(
found = version,
expected = SCHEMA_VERSION,
"store was written by a newer sidecar; continuing, as migrations are additive"
);
return Ok(());
}
while version < SCHEMA_VERSION {
let next = version + 1;
let mut tx = pool.begin().await?;
match next {
// Protocol 4: the guild board carries a member roster. Its own column rather than a
// field folded into `json`, because `json` holds the verbatim `guild.update` line and
// the two writers must not overwrite each other — see `upsert_guild_roster`.
1 => {
sqlx::query("ALTER TABLE guilds ADD COLUMN members TEXT")
.execute(&mut *tx)
.await?;
}
// Unreachable while SCHEMA_VERSION and this match are edited together, which is the
// point of failing loudly rather than silently leaving the counter short.
n => anyhow::bail!("no migration step defined for schema version {n}"),
}
// `PRAGMA` takes no bind parameters, so this is formatted — safe because `next` is an i64
// this loop produced, never anything from outside the process.
sqlx::query(&format!("PRAGMA user_version = {next}"))
.execute(&mut *tx)
.await?;
tx.commit().await?;
info!(version = next, "schema migration applied");
version = next;
}
Ok(())
}
fn parse_json_column(rows: Vec<sqlx::sqlite::SqliteRow>) -> Vec<Value> { fn parse_json_column(rows: Vec<sqlx::sqlite::SqliteRow>) -> Vec<Value> {
rows.into_iter() rows.into_iter()
.filter_map(|r| serde_json::from_str(&r.get::<String, _>("json")).ok()) .filter_map(|r| serde_json::from_str(&r.get::<String, _>("json")).ok())
@@ -611,3 +730,195 @@ CREATE TABLE IF NOT EXISTS ruleset (
updated_t INTEGER NOT NULL updated_t INTEGER NOT NULL
); );
"#; "#;
#[cfg(test)]
mod tests {
use super::*;
/// A unique scratch database path. Matches `config`'s idiom — `std::env::temp_dir()` plus the
/// test name — so the cases stay independent under the parallel test runner.
fn scratch(name: &str) -> String {
let dir = std::env::temp_dir().join(format!("uo-link-store-test-{name}"));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("create scratch dir");
dir.join("uo-link.db").to_string_lossy().into_owned()
}
/// The `guilds` table exactly as a pre-Protocol-4 sidecar left it: no `members` column, and
/// `user_version` still 0. This is the shape a real operator's database is in before an update,
/// and the only starting point where the migration does anything.
async fn legacy_db(path: &str) -> SqlitePool {
let pool = SqlitePoolOptions::new()
.max_connections(1)
.connect_with(
SqliteConnectOptions::new()
.filename(path)
.create_if_missing(true),
)
.await
.expect("open legacy db");
sqlx::query(
"CREATE TABLE guilds (
id INTEGER PRIMARY KEY,
name TEXT,
json TEXT NOT NULL,
updated_t INTEGER NOT NULL
)",
)
.execute(&pool)
.await
.expect("create legacy guilds table");
pool.close().await;
pool
}
async fn user_version(store: &Store) -> i64 {
sqlx::query_scalar("PRAGMA user_version")
.fetch_one(&store.pool)
.await
.expect("read user_version")
}
async fn guild_columns(store: &Store) -> Vec<String> {
sqlx::query("PRAGMA table_info(guilds)")
.fetch_all(&store.pool)
.await
.expect("table_info")
.into_iter()
.map(|r| r.get::<String, _>("name"))
.collect()
}
#[tokio::test]
async fn an_existing_pre_protocol_4_database_gains_the_members_column() {
// The case that matters: `SCHEMA`'s CREATE TABLE IF NOT EXISTS is a no-op against a table
// that is already there, so without `migrate` this database would never get the column and
// every roster write would fail against a live install.
let path = scratch("legacy-upgrade");
legacy_db(&path).await;
let store = Store::open(&path).await.expect("open migrates");
assert!(
guild_columns(&store).await.contains(&"members".to_string()),
"the migration must add guilds.members to a database that already had the table"
);
assert_eq!(user_version(&store).await, SCHEMA_VERSION);
}
#[tokio::test]
async fn a_fresh_database_lands_at_the_current_version() {
let path = scratch("fresh");
let store = Store::open(&path).await.expect("open");
assert!(guild_columns(&store).await.contains(&"members".to_string()));
assert_eq!(user_version(&store).await, SCHEMA_VERSION);
}
#[tokio::test]
async fn reopening_an_already_migrated_database_is_a_no_op() {
// Every sidecar restart re-runs this path, so a second run must not attempt the ALTER again
// — which would fail with "duplicate column name" and, since a migration failure aborts
// startup, would leave the sidecar unable to start at all after its first upgrade.
let path = scratch("idempotent");
legacy_db(&path).await;
Store::open(&path).await.expect("first open");
let store = Store::open(&path).await.expect("second open must succeed");
assert_eq!(user_version(&store).await, SCHEMA_VERSION);
}
#[tokio::test]
async fn a_roster_does_not_clobber_the_guild_update_snapshot() {
// The invariant the two-column split exists to give. `json` holds the verbatim guild.update
// line; if a roster write touched it, name/abbr/online would vanish from the board.
let path = scratch("no-clobber");
let store = Store::open(&path).await.expect("open");
store
.upsert_guild(
7,
Some("The Cartographers"),
r#"{"kind":"guild.update","id":7,"name":"The Cartographers","abbr":"MAP","members":2,"online":1}"#,
100,
)
.await
.expect("upsert guild");
store
.upsert_guild_roster(
7,
r#"[{"serial":"0x1","name":"Ada"},{"serial":"0x2","name":"Bo"}]"#,
200,
)
.await
.expect("upsert roster");
let guilds = store.guilds_all().await.expect("read board");
assert_eq!(guilds.len(), 1);
let g = &guilds[0];
assert_eq!(
g["name"], "The Cartographers",
"guild.update's name survived"
);
assert_eq!(g["abbr"], "MAP", "guild.update's abbr survived");
assert_eq!(g["online"], 1, "guild.update's online count survived");
assert_eq!(g["roster"].as_array().expect("roster is an array").len(), 2);
assert_eq!(g["roster"][0]["name"], "Ada");
}
#[tokio::test]
async fn the_two_writers_are_order_independent() {
// A roster can arrive before the first guild.update for a guild — on a reconnect the shard
// re-emits both and nothing orders them. Neither write may depend on the other's row.
let path = scratch("either-order");
let store = Store::open(&path).await.expect("open");
store
.upsert_guild_roster(9, r#"[{"serial":"0x3","name":"Cy"}]"#, 100)
.await
.expect("roster first");
store
.upsert_guild(
9,
Some("Late Arrivals"),
r#"{"kind":"guild.update","id":9,"name":"Late Arrivals","abbr":"LTE"}"#,
200,
)
.await
.expect("update second");
let guilds = store.guilds_all().await.expect("read board");
assert_eq!(guilds.len(), 1, "one row, not two");
assert_eq!(guilds[0]["name"], "Late Arrivals");
assert_eq!(guilds[0]["roster"].as_array().expect("roster").len(), 1);
}
#[tokio::test]
async fn a_guild_with_no_roster_yet_has_no_roster_key() {
// "Not known" and "known to be empty" are different, and the board must not conflate them:
// a website reading `roster: []` would render an empty roster as fact.
let path = scratch("absent-roster");
let store = Store::open(&path).await.expect("open");
store
.upsert_guild(
11,
Some("Unswept"),
r#"{"kind":"guild.update","id":11,"name":"Unswept"}"#,
100,
)
.await
.expect("upsert guild");
let guilds = store.guilds_all().await.expect("read board");
assert!(
guilds[0].get("roster").is_none(),
"a guild with no roster event must not grow a roster key"
);
}
}

41
sidecar/src/unix.rs Normal file
View File

@@ -0,0 +1,41 @@
//! Unix startup and shutdown.
//!
//! There is no supervisor protocol to speak: systemd starts the process, and stops it by sending
//! `SIGTERM`. All this module does is translate the two signals that mean "stop" into the future
//! [`crate::app::run`] waits on, so a `systemctl stop` unwinds the same way a Ctrl-C does instead
//! of being killed by the default `SIGTERM` disposition mid-write.
use tokio::signal::unix::{signal, SignalKind};
pub fn run(config_path: Option<&str>) -> anyhow::Result<()> {
crate::init_console_tracing();
tokio::runtime::Runtime::new()?.block_on(crate::app::run(config_path, || {}, shutdown_signal()))
}
/// Resolves on the first `SIGINT` or `SIGTERM`.
async fn shutdown_signal() {
// A failure to install a handler is not worth aborting a running sidecar for: fall back to
// pending, which leaves that signal's default disposition (terminate) in place.
let mut term = match signal(SignalKind::terminate()) {
Ok(s) => s,
Err(e) => {
tracing::warn!(error = %e, "could not listen for SIGTERM");
std::future::pending::<()>().await;
unreachable!()
}
};
let mut int = match signal(SignalKind::interrupt()) {
Ok(s) => s,
Err(e) => {
tracing::warn!(error = %e, "could not listen for SIGINT");
term.recv().await;
return;
}
};
tokio::select! {
_ = term.recv() => tracing::info!("SIGTERM received"),
_ = int.recv() => tracing::info!("SIGINT received"),
}
}

239
sidecar/src/windows.rs Normal file
View File

@@ -0,0 +1,239 @@
//! Windows startup and shutdown: the SCM handshake.
//!
//! Unlike systemd, the Windows Service Control Manager cannot supervise an arbitrary console
//! program. A binary registered with `sc.exe create` has ~30 seconds to call
//! `StartServiceCtrlDispatcher` and connect back to the SCM; one that never does is killed with
//! **error 1053, "the service did not respond to the start request in a timely fashion"** — even
//! though the process itself started perfectly and is sitting there serving traffic. That is the
//! entire reason this module exists.
//!
//! ## One binary, two ways in
//!
//! The dispatcher is tried first and *failing is expected*: when the process was started from a
//! shell rather than by the SCM, the connect fails with `ERROR_FAILED_SERVICE_CONTROLLER_CONNECT`
//! (1063), and that — and only that — falls through to a normal foreground run. So
//! `uo-link-sidecar.exe --config ...` stays an ordinary console app you can Ctrl-C, `cargo run`
//! still works, and the same binary can be registered as a service with no `--service` flag for an
//! operator to forget. Any other dispatcher error is a real failure and is reported.
//!
//! ## Logging goes to a file, because a service has no stdout
//!
//! Under the SCM there is no console attached, so the normal stdout subscriber writes into the
//! void. In service mode the sidecar logs to a daily-rolled file next to its config instead
//! (`uo-link-sidecar.YYYY-MM-DD.log`, seven kept). A service whose start fails leaves a reason
//! behind rather than only an SCM error code.
use std::ffi::OsString;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::OnceLock;
use std::time::Duration;
use tokio::sync::Notify;
use tracing_subscriber::EnvFilter;
use windows_service::service::{
ServiceControl, ServiceControlAccept, ServiceExitCode, ServiceState, ServiceStatus, ServiceType,
};
use windows_service::service_control_handler::{self, ServiceControlHandlerResult};
use windows_service::{define_windows_service, service_dispatcher};
/// Must match the name the installer registers (`installer/src/service.rs::WINDOWS_SERVICE`). For
/// an own-process service the SCM ignores it, but a mismatch would be a trap for whoever converts
/// this to a shared-process service later.
pub const SERVICE_NAME: &str = "RunicGatewayLink";
const SERVICE_TYPE: ServiceType = ServiceType::OWN_PROCESS;
/// `ERROR_FAILED_SERVICE_CONTROLLER_CONNECT` — "this process was not started by the SCM", which is
/// the normal answer when a human runs the binary.
const ERROR_FAILED_SERVICE_CONTROLLER_CONNECT: i32 = 1063;
/// `service_main` is called through an `extern "system"` trampoline and so can capture nothing.
/// The parsed `--config` is handed over here instead of being re-parsed, so the service and a
/// console run resolve their configuration through exactly the same code path.
static CONFIG_PATH: OnceLock<Option<String>> = OnceLock::new();
pub fn run(config_path: Option<&str>) -> anyhow::Result<()> {
let _ = CONFIG_PATH.set(config_path.map(str::to_string));
match service_dispatcher::start(SERVICE_NAME, ffi_service_main) {
Ok(()) => Ok(()),
// Not started by the SCM: this is a foreground run, which is not an error.
Err(windows_service::Error::Winapi(e))
if e.raw_os_error() == Some(ERROR_FAILED_SERVICE_CONTROLLER_CONNECT) =>
{
console_run(config_path)
}
Err(e) => Err(anyhow::Error::new(e)
.context("could not connect to the Windows service control manager")),
}
}
/// A normal foreground run: stdout logging, Ctrl-C to stop.
fn console_run(config_path: Option<&str>) -> anyhow::Result<()> {
crate::init_console_tracing();
tokio::runtime::Runtime::new()?.block_on(crate::app::run(config_path, || {}, async {
let _ = tokio::signal::ctrl_c().await;
}))
}
define_windows_service!(ffi_service_main, service_main);
fn service_main(_arguments: Vec<OsString>) {
// Arguments are deliberately ignored: for an own-process service the `binPath=` arguments
// arrive on the process command line and have already been parsed in `main`. What lands here
// is whatever was typed after `sc start`, which nothing in this deployment uses.
if let Err(e) = serve() {
// Nowhere left to report to but the log: the status handle is gone or was never obtained.
tracing::error!(error = %e, "service exited with an error");
}
}
fn serve() -> anyhow::Result<()> {
let config_path = CONFIG_PATH.get().cloned().flatten();
// Held for the life of the service: dropping the guard stops the background log writer.
let _log_guard = init_service_tracing(config_path.as_deref());
// The SCM calls the control handler on its own thread, so the stop signal crosses a thread
// boundary into the async world. `notify_one` stores a permit if nothing is waiting yet, so a
// stop that arrives during startup is not lost.
let stop = Arc::new(Notify::new());
let handler_stop = stop.clone();
let status_handle =
service_control_handler::register(SERVICE_NAME, move |control| match control {
ServiceControl::Interrogate => ServiceControlHandlerResult::NoError,
ServiceControl::Stop | ServiceControl::Shutdown => {
handler_stop.notify_one();
ServiceControlHandlerResult::NoError
}
_ => ServiceControlHandlerResult::NotImplemented,
})?;
// Registering the handler is the handshake 1053 was about. Everything after this point gets to
// take as long as it credibly needs, as long as the state keeps being reported.
status_handle.set_service_status(ServiceStatus {
service_type: SERVICE_TYPE,
current_state: ServiceState::StartPending,
controls_accepted: ServiceControlAccept::empty(),
exit_code: ServiceExitCode::Win32(0),
checkpoint: 0,
wait_hint: Duration::from_secs(30),
process_id: None,
})?;
let ready_handle = status_handle;
let result = tokio::runtime::Runtime::new()?.block_on(crate::app::run(
config_path.as_deref(),
// Reported only once the shard port is bound and the store is open, so a bad config or a
// taken port fails the *start* instead of flapping Running → Stopped a moment later.
move || {
let _ = ready_handle.set_service_status(ServiceStatus {
service_type: SERVICE_TYPE,
current_state: ServiceState::Running,
controls_accepted: ServiceControlAccept::STOP | ServiceControlAccept::SHUTDOWN,
exit_code: ServiceExitCode::Win32(0),
checkpoint: 0,
wait_hint: Duration::default(),
process_id: None,
});
},
async move { stop.notified().await },
));
// A failed run must leave a nonzero SERVICE_EXIT_CODE behind: `sc query` reporting STOPPED with
// exit code 0 is what made the original failure look like a clean stop.
let exit_code = match &result {
Ok(()) => ServiceExitCode::Win32(0),
Err(e) => {
tracing::error!(error = %e, "sidecar failed");
ServiceExitCode::ServiceSpecific(1)
}
};
status_handle.set_service_status(ServiceStatus {
service_type: SERVICE_TYPE,
current_state: ServiceState::Stopped,
controls_accepted: ServiceControlAccept::empty(),
exit_code,
checkpoint: 0,
wait_hint: Duration::default(),
process_id: None,
})?;
result
}
/// Where the service writes its log: beside the config it was pointed at, which is the directory
/// the installer already provisions and grants the service account write access to.
fn log_dir(config_path: Option<&str>) -> PathBuf {
if let Some(parent) = config_path
.map(PathBuf::from)
.as_deref()
.and_then(|p| p.parent())
.filter(|p| !p.as_os_str().is_empty())
{
return parent.to_path_buf();
}
match std::env::var_os("ProgramData") {
Some(program_data) => PathBuf::from(program_data).join("RunicGateway"),
None => std::env::temp_dir(),
}
}
/// Returns `None` if the log file could not be opened — a service that cannot write a log is still
/// a service worth running, and the SCM start must not fail over it.
fn init_service_tracing(
config_path: Option<&str>,
) -> Option<tracing_appender::non_blocking::WorkerGuard> {
let appender = tracing_appender::rolling::Builder::new()
.rotation(tracing_appender::rolling::Rotation::DAILY)
.filename_prefix("uo-link-sidecar")
.filename_suffix("log")
.max_log_files(7)
.build(log_dir(config_path))
.ok()?;
let (writer, guard) = tracing_appender::non_blocking(appender);
tracing_subscriber::fmt()
.with_env_filter(
EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")),
)
.with_ansi(false) // a log file is not a terminal
.with_writer(writer)
.init();
Some(guard)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn log_dir_follows_the_config_file() {
assert_eq!(
log_dir(Some(r"C:\ProgramData\RunicGateway\sidecar.toml")),
PathBuf::from(r"C:\ProgramData\RunicGateway")
);
}
#[test]
fn a_bare_filename_does_not_become_the_filesystem_root() {
// `--config sidecar.toml` has a parent of "", which as a path means the root of the current
// drive — somewhere a service account cannot write. Fall back instead.
let dir = log_dir(Some("sidecar.toml"));
assert_ne!(dir, PathBuf::from(""));
assert!(dir.is_absolute(), "{}", dir.display());
}
#[test]
fn no_config_falls_back_to_program_data() {
let dir = log_dir(None);
assert!(dir.is_absolute(), "{}", dir.display());
}
#[test]
fn service_name_matches_the_installer() {
// installer/src/service.rs::WINDOWS_SERVICE. Kept as a literal on both sides — the two
// repos are released independently and do not share a crate.
assert_eq!(SERVICE_NAME, "RunicGatewayLink");
}
}