Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
93c2b8520b
|
||
|
|
c7c45b26f8
|
||
|
|
ea68cb30c4
|
||
|
|
a8dfb16f96
|
||
|
|
b8823c850a
|
||
|
|
78fc717bc2
|
||
|
|
759431b316
|
||
|
|
44a8e58f48
|
||
|
|
f573c8e742
|
||
|
|
26241cb534
|
||
|
|
dd13a01bc9 | ||
|
|
da609cabb9 | ||
|
|
54cd4aab56 | ||
|
|
6709e9dea3 | ||
|
|
e9ce97df85 | ||
|
|
e8407f1719 | ||
|
|
0d0b237df7 |
@@ -0,0 +1,123 @@
|
|||||||
|
name: Publish NuGet packages
|
||||||
|
|
||||||
|
on:
|
||||||
|
push:
|
||||||
|
tags:
|
||||||
|
- 'v*'
|
||||||
|
workflow_dispatch:
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
publish:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
permissions:
|
||||||
|
contents: read
|
||||||
|
packages: write
|
||||||
|
env:
|
||||||
|
NUGET_AUTH_TOKEN: ${{ secrets.SHRINKSDK_PACKAGE_TOKEN }}
|
||||||
|
DOTNET_SYSTEM_GLOBALIZATION_INVARIANT: '1'
|
||||||
|
DOTNET_CLI_TELEMETRY_OPTOUT: '1'
|
||||||
|
LD_LIBRARY_PATH: /opt/dotnet-libs/usr/lib/x86_64-linux-gnu
|
||||||
|
SSL_CERT_FILE: /opt/ca-certificates.crt
|
||||||
|
steps:
|
||||||
|
- name: Fetch exact tagged source
|
||||||
|
env:
|
||||||
|
GITEA_REF: ${{ gitea.ref }}
|
||||||
|
GITEA_REPOSITORY: ${{ gitea.repository }}
|
||||||
|
shell: bash
|
||||||
|
run: |
|
||||||
|
set -euo pipefail
|
||||||
|
tag="${GITEA_REF#refs/tags/}"
|
||||||
|
case "$tag" in
|
||||||
|
v[0-9]*) ;;
|
||||||
|
*) echo "Expected a version tag ref, got: $GITEA_REF" >&2; exit 1 ;;
|
||||||
|
esac
|
||||||
|
export SHRINKSDK_ARCHIVE_URL="https://git.crash.work/${GITEA_REPOSITORY}/archive/${tag}.tar.gz"
|
||||||
|
node --input-type=module <<'NODE'
|
||||||
|
import { writeFile } from 'node:fs/promises';
|
||||||
|
|
||||||
|
const response = await fetch(process.env.SHRINKSDK_ARCHIVE_URL);
|
||||||
|
if (!response.ok) {
|
||||||
|
throw new Error(`Release archive download failed: ${response.status} ${response.statusText}`);
|
||||||
|
}
|
||||||
|
await writeFile('release.tar.gz', new Uint8Array(await response.arrayBuffer()));
|
||||||
|
NODE
|
||||||
|
mkdir release
|
||||||
|
tar -xzf release.tar.gz --strip-components=1 -C release
|
||||||
|
rm -f release.tar.gz
|
||||||
|
printf '%s' "$tag" > release/.shrink-sdk-release-tag
|
||||||
|
|
||||||
|
- name: Install .NET 8 SDK
|
||||||
|
shell: bash
|
||||||
|
run: |
|
||||||
|
set -euo pipefail
|
||||||
|
node --input-type=module <<'NODE'
|
||||||
|
import { writeFile } from 'node:fs/promises';
|
||||||
|
import { rootCertificates } from 'node:tls';
|
||||||
|
|
||||||
|
await writeFile('/opt/ca-certificates.crt', rootCertificates.join('\n'));
|
||||||
|
|
||||||
|
const metadataResponse = await fetch('https://dotnetcli.blob.core.windows.net/dotnet/release-metadata/8.0/releases.json');
|
||||||
|
if (!metadataResponse.ok) throw new Error(`Release metadata download failed: ${metadataResponse.status}`);
|
||||||
|
const metadata = await metadataResponse.json();
|
||||||
|
const sdkVersion = metadata['latest-sdk'];
|
||||||
|
const release = metadata.releases.find(item => item.sdk?.version === sdkVersion);
|
||||||
|
const file = release?.sdk?.files?.find(item => item.rid === 'linux-x64' && item.name.endsWith('.tar.gz'));
|
||||||
|
if (!file) throw new Error(`Linux x64 SDK archive not found for ${sdkVersion}`);
|
||||||
|
const archiveResponse = await fetch(file.url);
|
||||||
|
if (!archiveResponse.ok) throw new Error(`SDK download failed: ${archiveResponse.status}`);
|
||||||
|
await writeFile('/tmp/dotnet-sdk.tar.gz', new Uint8Array(await archiveResponse.arrayBuffer()));
|
||||||
|
|
||||||
|
const poolUrl = 'https://deb.debian.org/debian-security/pool/updates/main/o/openssl/';
|
||||||
|
const poolResponse = await fetch(poolUrl);
|
||||||
|
if (!poolResponse.ok) throw new Error(`OpenSSL package index download failed: ${poolResponse.status}`);
|
||||||
|
const poolIndex = await poolResponse.text();
|
||||||
|
const packages = [...poolIndex.matchAll(/href="(libssl3_[^"]+_amd64\.deb)"/g)].map(match => match[1]).sort();
|
||||||
|
const packageName = packages.at(-1);
|
||||||
|
if (!packageName) throw new Error('Debian libssl3 package was not found');
|
||||||
|
const packageResponse = await fetch(poolUrl + packageName);
|
||||||
|
if (!packageResponse.ok) throw new Error(`OpenSSL package download failed: ${packageResponse.status}`);
|
||||||
|
await writeFile('/tmp/libssl3.deb', new Uint8Array(await packageResponse.arrayBuffer()));
|
||||||
|
NODE
|
||||||
|
mkdir -p /opt/dotnet
|
||||||
|
tar -xzf /tmp/dotnet-sdk.tar.gz -C /opt/dotnet
|
||||||
|
mkdir -p /opt/dotnet-libs
|
||||||
|
dpkg-deb -x /tmp/libssl3.deb /opt/dotnet-libs
|
||||||
|
rm -f /tmp/dotnet-sdk.tar.gz
|
||||||
|
rm -f /tmp/libssl3.deb
|
||||||
|
/opt/dotnet/dotnet --info
|
||||||
|
|
||||||
|
- name: Validate, pack and publish
|
||||||
|
shell: bash
|
||||||
|
run: |
|
||||||
|
set -euo pipefail
|
||||||
|
: "${NUGET_AUTH_TOKEN:?SHRINKSDK_PACKAGE_TOKEN is required}"
|
||||||
|
export PATH="/opt/dotnet:$PATH"
|
||||||
|
cd release
|
||||||
|
tag="$(cat .shrink-sdk-release-tag)"
|
||||||
|
projects=()
|
||||||
|
if [[ -d DotNet~ ]]; then
|
||||||
|
while IFS= read -r -d '' project; do projects+=("$project"); done < <(find DotNet~ -type f -name '*.csproj' -print0)
|
||||||
|
fi
|
||||||
|
if [[ -d Godot~ ]]; then
|
||||||
|
while IFS= read -r -d '' project; do projects+=("$project"); done < <(find Godot~ -type f -name '*.csproj' -print0)
|
||||||
|
fi
|
||||||
|
if [[ "${#projects[@]}" -eq 0 ]]; then
|
||||||
|
while IFS= read -r -d '' project; do projects+=("$project"); done < <(find . -maxdepth 1 -type f -name '*.csproj' -print0)
|
||||||
|
fi
|
||||||
|
test "${#projects[@]}" -gt 0
|
||||||
|
if [[ -f package.json ]]; then
|
||||||
|
version="$(node -p "require('./package.json').version")"
|
||||||
|
else
|
||||||
|
version="$(dotnet msbuild "${projects[0]}" -getProperty:Version -nologo)"
|
||||||
|
fi
|
||||||
|
test "$tag" = "v$version"
|
||||||
|
mkdir -p packages
|
||||||
|
for project in "${projects[@]}"; do
|
||||||
|
dotnet restore "$project" --configfile NuGet.Config
|
||||||
|
dotnet pack "$project" --configuration Release --no-restore --output "$PWD/packages" --include-symbols --include-source
|
||||||
|
done
|
||||||
|
find packages -maxdepth 1 -name '*.nupkg' -type f | grep -q .
|
||||||
|
dotnet nuget push 'packages/*.nupkg' --api-key "$NUGET_AUTH_TOKEN" --source https://git.crash.work/api/packages/ShrinkSDK/nuget/index.json --skip-duplicate
|
||||||
|
if compgen -G 'packages/*.snupkg' > /dev/null; then
|
||||||
|
dotnet nuget push 'packages/*.snupkg' --api-key "$NUGET_AUTH_TOKEN" --source https://git.crash.work/api/packages/ShrinkSDK/nuget/index.json --skip-duplicate
|
||||||
|
fi
|
||||||
@@ -15,22 +15,39 @@ jobs:
|
|||||||
env:
|
env:
|
||||||
NODE_AUTH_TOKEN: ${{ secrets.SHRINKSDK_PACKAGE_TOKEN }}
|
NODE_AUTH_TOKEN: ${{ secrets.SHRINKSDK_PACKAGE_TOKEN }}
|
||||||
steps:
|
steps:
|
||||||
- name: Fetch tagged revision
|
- name: Fetch exact tagged release archive
|
||||||
|
env:
|
||||||
|
GITEA_REF: ${{ gitea.ref }}
|
||||||
shell: bash
|
shell: bash
|
||||||
run: |
|
run: |
|
||||||
set -eu
|
set -eu
|
||||||
ref="${{ gitea.sha }}"
|
tag="${GITEA_REF#refs/tags/}"
|
||||||
test -n "$ref"
|
case "$tag" in
|
||||||
git init .
|
v[0-9]*) ;;
|
||||||
git remote add origin "https://git.crash.work/ShrinkSDK/ShrinkNetwork.git"
|
*) echo "Expected a version tag ref, got: $GITEA_REF" >&2; exit 1 ;;
|
||||||
git fetch --depth=1 origin "$ref"
|
esac
|
||||||
git checkout --detach FETCH_HEAD
|
export SHRINKSDK_ARCHIVE_URL="https://git.crash.work/ShrinkSDK/ShrinkNetwork/archive/${tag}.tar.gz"
|
||||||
|
node --input-type=module <<'NODE'
|
||||||
|
import { writeFile } from 'node:fs/promises';
|
||||||
|
|
||||||
|
const response = await fetch(process.env.SHRINKSDK_ARCHIVE_URL);
|
||||||
|
if (!response.ok) {
|
||||||
|
throw new Error(`Release archive download failed: ${response.status} ${response.statusText}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
await writeFile('release.tar.gz', new Uint8Array(await response.arrayBuffer()));
|
||||||
|
NODE
|
||||||
|
mkdir release
|
||||||
|
tar -xzf release.tar.gz --strip-components=1 -C release
|
||||||
|
rm -f release.tar.gz
|
||||||
|
printf '%s' "$tag" > release/.shrink-sdk-release-tag
|
||||||
|
|
||||||
- name: Validate immutable release version
|
- name: Validate immutable release version
|
||||||
shell: bash
|
shell: bash
|
||||||
run: |
|
run: |
|
||||||
set -eu
|
set -eu
|
||||||
tag="$(git describe --exact-match --tags HEAD)"
|
cd release
|
||||||
|
tag="$(cat .shrink-sdk-release-tag)"
|
||||||
version="$(node -p "require('./package.json').version")"
|
version="$(node -p "require('./package.json').version")"
|
||||||
test "$tag" = "v$version"
|
test "$tag" = "v$version"
|
||||||
npm pack --dry-run
|
npm pack --dry-run
|
||||||
@@ -40,6 +57,7 @@ jobs:
|
|||||||
run: |
|
run: |
|
||||||
set -eu
|
set -eu
|
||||||
: "${NODE_AUTH_TOKEN:?SHRINKSDK_PACKAGE_TOKEN is required}"
|
: "${NODE_AUTH_TOKEN:?SHRINKSDK_PACKAGE_TOKEN is required}"
|
||||||
|
cd release
|
||||||
npmrc="$HOME/.npmrc"
|
npmrc="$HOME/.npmrc"
|
||||||
cleanup() { rm -f "$npmrc"; }
|
cleanup() { rm -f "$npmrc"; }
|
||||||
trap cleanup EXIT
|
trap cleanup EXIT
|
||||||
|
|||||||
@@ -26,12 +26,17 @@ jobs:
|
|||||||
shell: bash
|
shell: bash
|
||||||
run: |
|
run: |
|
||||||
set -eu
|
set -eu
|
||||||
|
machine_id_file="/root/.local/share/unity3d/Unity/.machine-id"
|
||||||
|
if test -s "$machine_id_file"; then
|
||||||
|
cat "$machine_id_file" > /etc/machine-id
|
||||||
|
echo "Unity machine identity restored"
|
||||||
|
fi
|
||||||
|
git config --global url."https://ghfast.top/https://github.com/".insteadOf "https://github.com/"
|
||||||
unity_bin="$(command -v unity-editor || command -v unity || command -v Unity || true)"
|
unity_bin="$(command -v unity-editor || command -v unity || command -v Unity || true)"
|
||||||
test -n "$unity_bin"
|
test -n "$unity_bin"
|
||||||
"$unity_bin" \
|
"$unity_bin" \
|
||||||
-batchmode \
|
-batchmode \
|
||||||
-nographics \
|
-nographics \
|
||||||
-quit \
|
|
||||||
-projectPath "$PWD/Development~/UnityProject" \
|
-projectPath "$PWD/Development~/UnityProject" \
|
||||||
-runTests \
|
-runTests \
|
||||||
-testPlatform EditMode \
|
-testPlatform EditMode \
|
||||||
|
|||||||
+12
@@ -8,3 +8,15 @@
|
|||||||
/Tools~/**/[Oo]bj/
|
/Tools~/**/[Oo]bj/
|
||||||
*.user
|
*.user
|
||||||
*.DotSettings.user
|
*.DotSettings.user
|
||||||
|
/DotNet~/**/[Bb]in/
|
||||||
|
/DotNet~/**/[Oo]bj/
|
||||||
|
/Godot~/**/[Bb]in/
|
||||||
|
/Godot~/**/[Oo]bj/
|
||||||
|
/artifacts/
|
||||||
|
/packages/
|
||||||
|
!DotNet~/**/*.csproj
|
||||||
|
!Godot~/**/*.csproj
|
||||||
|
|
||||||
|
/Adapters~/**/bin/
|
||||||
|
/Adapters~/**/obj/
|
||||||
|
!Adapters~/**/*.csproj
|
||||||
|
|||||||
@@ -6,3 +6,9 @@ Tools~/
|
|||||||
*.sln
|
*.sln
|
||||||
*.user
|
*.user
|
||||||
*.DotSettings.user
|
*.DotSettings.user
|
||||||
|
DotNet~/
|
||||||
|
Godot~/
|
||||||
|
NuGet.Config
|
||||||
|
Directory.Build.props
|
||||||
|
NuGet.Config.meta
|
||||||
|
Directory.Build.props.meta
|
||||||
|
|||||||
@@ -0,0 +1,9 @@
|
|||||||
|
# Generated MessagePack adapter
|
||||||
|
|
||||||
|
独立可选包。Unity 先安装 MessagePack-CSharp 3.1.8(含其生成器及依赖),再通过 UPM 引用本目录。此目录放在 `Adapters~` 下,不让未安装 MessagePack 的项目产生缺失程序集错误。.NET 使用 `ShrinkSDK.Network.MessagePack`。
|
||||||
|
|
||||||
|
消息使用 `[MessagePackObject]`、稳定 `[Key(n)]` 和生成的 resolver。构造 `ShrinkNetwork.MessagePack.ShrinkMessagePackNetworkSerializer` 时显式传入 resolver,再为每个消息调用 `Register<T>()`。未注册或没有 formatter 的类型明确报错;不回退到 Contractless/动态 IL/反射调用。
|
||||||
|
|
||||||
|
在绑定 transport 前注册;停止使用相关消息后才能撤回 codec。共享合同不能依赖 Unity 对象或引擎专属类型。双方必须使用相同合同与消息编码。
|
||||||
|
|
||||||
|
可编译示例与生成式互通测试:Workspace `Tools/AgentSupport/RuntimeTests/NetworkCodecTests.cs`。
|
||||||
@@ -0,0 +1,7 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 2ad6cd0d17394435bfb2d7bc04e78c0a
|
||||||
|
DefaultImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -0,0 +1,8 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: b29018215be842d782c6e5b4fc53cf7e
|
||||||
|
folderAsset: yes
|
||||||
|
DefaultImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -0,0 +1,47 @@
|
|||||||
|
#nullable enable
|
||||||
|
using System;
|
||||||
|
using System.Buffers;
|
||||||
|
using MessagePack;
|
||||||
|
using MessagePack.Formatters;
|
||||||
|
|
||||||
|
namespace ShrinkNetwork.MessagePack
|
||||||
|
{
|
||||||
|
/// <summary>Pass a generated-only resolver; register all message types before starting transports.</summary>
|
||||||
|
public sealed class ShrinkMessagePackNetworkSerializer : ShrinkRegisteredNetworkSerializer
|
||||||
|
{
|
||||||
|
private readonly MessagePackSerializerOptions _options;
|
||||||
|
public ShrinkMessagePackNetworkSerializer(IFormatterResolver generatedResolver)
|
||||||
|
{
|
||||||
|
_options = MessagePackSerializerOptions.Standard
|
||||||
|
.WithResolver(global::MessagePack.Resolvers.CompositeResolver.Create(
|
||||||
|
generatedResolver ?? throw new ArgumentNullException(nameof(generatedResolver)),
|
||||||
|
global::MessagePack.Resolvers.BuiltinResolver.Instance))
|
||||||
|
.WithSecurity(MessagePackSecurity.UntrustedData);
|
||||||
|
}
|
||||||
|
public void Register<T>() => Register(new FormatterCodec<T>(_options));
|
||||||
|
|
||||||
|
private sealed class FormatterCodec<T> : IShrinkMessageCodec<T>
|
||||||
|
{
|
||||||
|
private readonly IMessagePackFormatter<T> _formatter;
|
||||||
|
private readonly MessagePackSerializerOptions _options;
|
||||||
|
public FormatterCodec(MessagePackSerializerOptions options)
|
||||||
|
{
|
||||||
|
_options = options;
|
||||||
|
_formatter = options.Resolver.GetFormatter<T>() ?? throw new InvalidOperationException($"No generated MessagePack formatter for {typeof(T).FullName}.");
|
||||||
|
}
|
||||||
|
public void Write(IBufferWriter<byte> buffer, T value)
|
||||||
|
{
|
||||||
|
var writer = new MessagePackWriter(buffer);
|
||||||
|
_formatter.Serialize(ref writer, value, _options);
|
||||||
|
writer.Flush();
|
||||||
|
}
|
||||||
|
public T Read(ReadOnlyMemory<byte> payload)
|
||||||
|
{
|
||||||
|
var reader = new MessagePackReader(payload);
|
||||||
|
var result = _formatter.Deserialize(ref reader, _options);
|
||||||
|
if (!reader.End) throw new System.IO.InvalidDataException("Trailing MessagePack payload.");
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+1
-1
@@ -1,5 +1,5 @@
|
|||||||
fileFormatVersion: 2
|
fileFormatVersion: 2
|
||||||
guid: 88c6239e9fc049f43b2999d7311eec1d
|
guid: c8394ccb85a14b23b66d1c9ec3a40dc3
|
||||||
MonoImporter:
|
MonoImporter:
|
||||||
externalObjects: {}
|
externalObjects: {}
|
||||||
serializedVersion: 2
|
serializedVersion: 2
|
||||||
@@ -0,0 +1,6 @@
|
|||||||
|
{
|
||||||
|
"name": "ShrinkNetwork.MessagePack",
|
||||||
|
"references": ["ShrinkNetwork.Runtime"],
|
||||||
|
"autoReferenced": true,
|
||||||
|
"noEngineReferences": true
|
||||||
|
}
|
||||||
@@ -0,0 +1,7 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 74898e5de85247f6b0228c06ee3bdf6e
|
||||||
|
AssemblyDefinitionImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -0,0 +1,15 @@
|
|||||||
|
<Project Sdk="Microsoft.NET.Sdk">
|
||||||
|
<PropertyGroup>
|
||||||
|
<TargetFramework>netstandard2.1</TargetFramework>
|
||||||
|
<LangVersion>latest</LangVersion>
|
||||||
|
<Nullable>enable</Nullable>
|
||||||
|
<PackageId>ShrinkSDK.Network.MessagePack</PackageId>
|
||||||
|
<Version>0.1.0</Version>
|
||||||
|
<EnableDefaultCompileItems>false</EnableDefaultCompileItems>
|
||||||
|
</PropertyGroup>
|
||||||
|
<ItemGroup>
|
||||||
|
<Compile Include="Runtime/**/*.cs" />
|
||||||
|
<PackageReference Include="MessagePack" Version="3.1.8" />
|
||||||
|
<PackageReference Include="ShrinkSDK.Network" Version="0.4.2" />
|
||||||
|
</ItemGroup>
|
||||||
|
</Project>
|
||||||
@@ -0,0 +1,10 @@
|
|||||||
|
{
|
||||||
|
"name": "com.cneicy.shrink-network-messagepack",
|
||||||
|
"version": "0.1.0",
|
||||||
|
"displayName": "Shrink Network MessagePack",
|
||||||
|
"description": "Explicit generated MessagePack codecs for ShrinkNetwork protocol v2. Requires MessagePack-CSharp 3.1.8.",
|
||||||
|
"unity": "2022.3",
|
||||||
|
"dependencies": {
|
||||||
|
"com.cneicy.shrink-network": "0.4.2"
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,7 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 7e1dd2e113fe4a74b6afe398b315db86
|
||||||
|
DefaultImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -10,6 +10,10 @@
|
|||||||
],
|
],
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"com.unity.test-framework": "1.1.33",
|
"com.unity.test-framework": "1.1.33",
|
||||||
"com.cneicy.shrink-network": "file:../../.."
|
"com.cneicy.shrink-network": "file:../../..",
|
||||||
}
|
"com.cysharp.unitask": "https://github.com/Cysharp/UniTask.git?path=src/UniTask/Assets/Plugins/UniTask#7c0f199fe0d3fc528024488ccd671e6c7b27745b"
|
||||||
|
},
|
||||||
|
"testables": [
|
||||||
|
"com.cneicy.shrink-network"
|
||||||
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,13 @@
|
|||||||
|
<Project>
|
||||||
|
<PropertyGroup>
|
||||||
|
<LangVersion>latest</LangVersion>
|
||||||
|
<Nullable>enable</Nullable>
|
||||||
|
<Deterministic>true</Deterministic>
|
||||||
|
<ContinuousIntegrationBuild>true</ContinuousIntegrationBuild>
|
||||||
|
<Authors>ShrinkSDK</Authors>
|
||||||
|
<Company>ShrinkSDK</Company>
|
||||||
|
<RepositoryUrl>https://git.crash.work/ShrinkSDK</RepositoryUrl>
|
||||||
|
<IncludeSymbols>true</IncludeSymbols>
|
||||||
|
<SymbolPackageFormat>snupkg</SymbolPackageFormat>
|
||||||
|
</PropertyGroup>
|
||||||
|
</Project>
|
||||||
@@ -0,0 +1,7 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 4b4f62aaca0a7d3419ac467e74201fce
|
||||||
|
DefaultImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -0,0 +1,21 @@
|
|||||||
|
<Project Sdk="Microsoft.NET.Sdk">
|
||||||
|
<PropertyGroup>
|
||||||
|
<TargetFramework>netstandard2.1</TargetFramework>
|
||||||
|
<EnableDefaultCompileItems>false</EnableDefaultCompileItems>
|
||||||
|
<AssemblyName>ShrinkNetwork.Runtime</AssemblyName>
|
||||||
|
<RootNamespace>ShrinkNetwork</RootNamespace>
|
||||||
|
<PackageId>ShrinkSDK.Network</PackageId>
|
||||||
|
<Version>0.4.2</Version>
|
||||||
|
<Description>ShrinkSDK messaging, RPC, permission and transport runtime.</Description>
|
||||||
|
<AllowUnsafeBlocks>true</AllowUnsafeBlocks>
|
||||||
|
<ShrinkCodeGenEnabled>false</ShrinkCodeGenEnabled>
|
||||||
|
</PropertyGroup>
|
||||||
|
<ItemGroup>
|
||||||
|
<Compile Include="..\Runtime\**\*.cs" />
|
||||||
|
<PackageReference Include="Kcp-CSharp" Version="1.0.8" />
|
||||||
|
<PackageReference Include="Newtonsoft.Json" Version="13.0.3" />
|
||||||
|
<PackageReference Include="UniTask" Version="2.5.10" />
|
||||||
|
<PackageReference Include="ShrinkSDK.Runtime.Abstractions" Version="0.1.0" />
|
||||||
|
<PackageReference Include="ShrinkSDK.CodeGen" Version="0.2.0" PrivateAssets="compile;runtime;contentfiles;native" />
|
||||||
|
</ItemGroup>
|
||||||
|
</Project>
|
||||||
@@ -11,33 +11,10 @@
|
|||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
<PackageReference Include="Kcp-CSharp" Version="1.0.8" />
|
<PackageReference Include="Kcp-CSharp" Version="1.0.8" />
|
||||||
<PackageReference Include="UniTask" Version="2.5.10" />
|
<PackageReference Include="UniTask" Version="2.5.10" />
|
||||||
<PackageReference Include="MessagePack" Version="3.1.4" />
|
|
||||||
<PackageReference Include="Newtonsoft.Json" Version="13.0.3" />
|
<PackageReference Include="Newtonsoft.Json" Version="13.0.3" />
|
||||||
</ItemGroup>
|
</ItemGroup>
|
||||||
|
|
||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Serialization\IShrinkNetworkSerializer.cs" Link="Runtime\Serialization\IShrinkNetworkSerializer.cs" />
|
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\**\*.cs" Link="Runtime\%(RecursiveDir)%(Filename)%(Extension)" />
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Transport\Abstractions\IShrinkNetworkAsyncTransport.cs" Link="Runtime\Transport\Abstractions\IShrinkNetworkAsyncTransport.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Transport\Abstractions\IShrinkNetworkTransport.cs" Link="Runtime\Transport\Abstractions\IShrinkNetworkTransport.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Transport\Tcp\ShrinkTcpTlsOptions.cs" Link="Runtime\Transport\Tcp\ShrinkTcpTlsOptions.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Serialization\ShrinkJsonNetworkSerializer.cs" Link="Runtime\Serialization\ShrinkJsonNetworkSerializer.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Transport\Kcp\ShrinkKcpPeer.cs" Link="Runtime\Transport\Kcp\ShrinkKcpPeer.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Transport\Kcp\ShrinkKcpTransportOptions.cs" Link="Runtime\Transport\Kcp\ShrinkKcpTransportOptions.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Transport\Kcp\ShrinkKcpTransportProtocol.cs" Link="Runtime\Transport\Kcp\ShrinkKcpTransportProtocol.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Serialization\ShrinkMessagePackNetworkSerializer.cs" Link="Runtime\Serialization\ShrinkMessagePackNetworkSerializer.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Metadata\ShrinkNetworkAttributes.cs" Link="Runtime\Metadata\ShrinkNetworkAttributes.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Core\ShrinkNetworkContext.cs" Link="Runtime\Core\ShrinkNetworkContext.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Core\ShrinkNetworkLogger.cs" Link="Runtime\Core\ShrinkNetworkLogger.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Metadata\ShrinkNetworkMessageContracts.cs" Link="Runtime\Metadata\ShrinkNetworkMessageContracts.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Routing\ShrinkNetworkMessageRegistry.cs" Link="Runtime\Routing\ShrinkNetworkMessageRegistry.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Metadata\ShrinkNetworkPacket.cs" Link="Runtime\Metadata\ShrinkNetworkPacket.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Metadata\ShrinkNetworkPermissions.cs" Link="Runtime\Metadata\ShrinkNetworkPermissions.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Routing\ShrinkNetworkRegHelper.cs" Link="Runtime\Routing\ShrinkNetworkRegHelper.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Routing\ShrinkNetworkGeneratedRegistry.cs" Link="Runtime\Routing\ShrinkNetworkGeneratedRegistry.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Routing\ShrinkNetworkRouter.cs" Link="Runtime\Routing\ShrinkNetworkRouter.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Metadata\ShrinkNetworkRpc.cs" Link="Runtime\Metadata\ShrinkNetworkRpc.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Core\ShrinkNetworkService.cs" Link="Runtime\Core\ShrinkNetworkService.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Core\ShrinkNetworkSession.cs" Link="Runtime\Core\ShrinkNetworkSession.cs" />
|
|
||||||
<Compile Include="..\..\Assets\Modules\ShrinkNetwork\Runtime\Metadata\ShrinkNetworkTransportEvent.cs" Link="Runtime\Metadata\ShrinkNetworkTransportEvent.cs" />
|
|
||||||
</ItemGroup>
|
</ItemGroup>
|
||||||
</Project>
|
</Project>
|
||||||
|
|||||||
@@ -11,9 +11,8 @@ using UnityEngine;
|
|||||||
|
|
||||||
public static class ShrinkDedicatedServerScaffoldGenerator
|
public static class ShrinkDedicatedServerScaffoldGenerator
|
||||||
{
|
{
|
||||||
private const string LegacyMenuPath = "ShrinkSDK/Network/生成独立服务器脚手架";
|
private const string FullProjectMenuPath = "ShrinkSDK/网络/生成完整独立服务器工程";
|
||||||
private const string FullProjectMenuPath = "ShrinkSDK/Network/生成完整独立服务器工程";
|
private const string RefreshGeneratedMenuPath = "ShrinkSDK/网络/刷新独立服务器生成合同";
|
||||||
private const string RefreshGeneratedMenuPath = "ShrinkSDK/Network/刷新独立服务器 Generated 合同";
|
|
||||||
private const string TemplateRelativePath = "Assets/Modules/ShrinkNetwork/Editor/Scaffolding/ServerProjectTemplate";
|
private const string TemplateRelativePath = "Assets/Modules/ShrinkNetwork/Editor/Scaffolding/ServerProjectTemplate";
|
||||||
private const string DefaultGeneratedServerRelativePath = "GeneratedServers/ShrinkNetwork.ServerHost";
|
private const string DefaultGeneratedServerRelativePath = "GeneratedServers/ShrinkNetwork.ServerHost";
|
||||||
private const string TemplateStampFileName = ".shrink-server-template.json";
|
private const string TemplateStampFileName = ".shrink-server-template.json";
|
||||||
@@ -115,7 +114,6 @@ public static class ShrinkDedicatedServerScaffoldGenerator
|
|||||||
"/heartbeat"
|
"/heartbeat"
|
||||||
};
|
};
|
||||||
|
|
||||||
[MenuItem(LegacyMenuPath)]
|
|
||||||
[MenuItem(FullProjectMenuPath)]
|
[MenuItem(FullProjectMenuPath)]
|
||||||
public static void GenerateFullProject()
|
public static void GenerateFullProject()
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -0,0 +1,8 @@
|
|||||||
|
<?xml version="1.0" encoding="utf-8"?>
|
||||||
|
<configuration>
|
||||||
|
<packageSources>
|
||||||
|
<clear />
|
||||||
|
<add key="ShrinkSDK" value="https://git.crash.work/api/packages/ShrinkSDK/nuget/index.json" />
|
||||||
|
<add key="nuget.org" value="https://api.nuget.org/v3/index.json" protocolVersion="3" />
|
||||||
|
</packageSources>
|
||||||
|
</configuration>
|
||||||
@@ -0,0 +1,32 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 3fd1dd39ba4a5c84fae6852e5d3fe643
|
||||||
|
PluginImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
serializedVersion: 2
|
||||||
|
iconMap: {}
|
||||||
|
executionOrder: {}
|
||||||
|
defineConstraints: []
|
||||||
|
isPreloaded: 0
|
||||||
|
isOverridable: 0
|
||||||
|
isExplicitlyReferenced: 0
|
||||||
|
validateReferences: 1
|
||||||
|
platformData:
|
||||||
|
- first:
|
||||||
|
Any:
|
||||||
|
second:
|
||||||
|
enabled: 0
|
||||||
|
settings: {}
|
||||||
|
- first:
|
||||||
|
Editor: Editor
|
||||||
|
second:
|
||||||
|
enabled: 0
|
||||||
|
settings:
|
||||||
|
DefaultValueInitialized: true
|
||||||
|
- first:
|
||||||
|
Windows Store Apps: WindowsStoreApps
|
||||||
|
second:
|
||||||
|
enabled: 1
|
||||||
|
settings: {}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -1,6 +1,8 @@
|
|||||||
# ShrinkNetwork
|
# ShrinkNetwork
|
||||||
|
|
||||||
一个为 Unity C# 项目设计的轻量网络框架。重点不是“再包一层 Socket”,而是把消息声明、权限校验、RPC、序列化、传输层,以及独立服务器生成流程收拢到一套统一范式里。
|
一个可用于 Unity、Godot 和普通 .NET 宿主的轻量网络框架。重点不是“再包一层 Socket”,而是把消息声明、权限校验、RPC、序列化、传输层,以及独立服务器生成流程收拢到一套统一范式里。
|
||||||
|
|
||||||
|
Godot 和普通 .NET 项目安装 `ShrinkSDK.Network`;可打包源码位于 `DotNet~`,消息合同和 handler 注册由 `ShrinkSDK.CodeGen` 在构建时织入。
|
||||||
|
|
||||||
## ✨ 特性概览
|
## ✨ 特性概览
|
||||||
|
|
||||||
@@ -90,7 +92,7 @@ public static class PingHandlers
|
|||||||
|
|
||||||
```csharp
|
```csharp
|
||||||
var service = new ShrinkNetworkService(
|
var service = new ShrinkNetworkService(
|
||||||
new ShrinkMessagePackNetworkSerializer(),
|
new ShrinkJsonNetworkSerializer(),
|
||||||
new ShrinkNetworkMessageRegistry(),
|
new ShrinkNetworkMessageRegistry(),
|
||||||
new ShrinkNetworkRouter());
|
new ShrinkNetworkRouter());
|
||||||
|
|
||||||
@@ -220,8 +222,8 @@ public sealed class LanPlayerStateDelta : IShrinkNetworkMessage
|
|||||||
|
|
||||||
Unity 菜单:
|
Unity 菜单:
|
||||||
|
|
||||||
- `ShrinkSDK/Network/生成完整独立服务器工程`
|
- `ShrinkSDK/网络/生成完整独立服务器工程`
|
||||||
- `ShrinkSDK/Network/刷新独立服务器 Generated 合同`
|
- `ShrinkSDK/网络/刷新独立服务器生成合同`
|
||||||
|
|
||||||
生成结果默认输出到:
|
生成结果默认输出到:
|
||||||
|
|
||||||
@@ -283,9 +285,10 @@ new ShrinkRpcCallOptions
|
|||||||
|
|
||||||
```csharp
|
```csharp
|
||||||
new ShrinkJsonNetworkSerializer()
|
new ShrinkJsonNetworkSerializer()
|
||||||
new ShrinkMessagePackNetworkSerializer()
|
|
||||||
```
|
```
|
||||||
|
|
||||||
|
生成式 MessagePack 通过可选包 `Adapters~/MessagePack` 接入。使用 `ShrinkNetwork.MessagePack.ShrinkMessagePackNetworkSerializer(GeneratedResolver.Instance)` 并显式 `Register<T>()`;旧无参反射适配器已移除。消息体只编码一次,外层统一使用协议 v2 的 SHK2 二进制封装。
|
||||||
|
|
||||||
## ✅ 最佳实践
|
## ✅ 最佳实践
|
||||||
|
|
||||||
- 所有消息统一用 `[ShrinkNetworkMessage]` 声明,不要只靠手写 `RegisterMessage<T>()`
|
- 所有消息统一用 `[ShrinkNetworkMessage]` 声明,不要只靠手写 `RegisterMessage<T>()`
|
||||||
@@ -365,3 +368,13 @@ var serverTransport = new ShrinkTcpServerTransport(
|
|||||||
## 📄 License
|
## 📄 License
|
||||||
|
|
||||||
[MIT](LICENSE)
|
[MIT](LICENSE)
|
||||||
|
|
||||||
|
## 可复用能力
|
||||||
|
|
||||||
|
<!-- shrink:capabilities -->
|
||||||
|
Capability: 消息合同、RPC、权限、TCP/KCP/Loopback 传输与有预算的串行排队
|
||||||
|
Aliases: 联机 网络 远程 rpc 广播 拥塞 背压 同步
|
||||||
|
Limits: 协议 v2 拒绝旧封装;TCP/KCP 均为可靠通道;可替换状态仅合并排队项,不用于快照分片
|
||||||
|
Extension: RegisterMessage/RegisterHandler;IShrinkNetworkBufferSerializer;ShrinkNetworkWorkQueue 状态键
|
||||||
|
Evidence: [ShrinkNetworkService](Runtime/Core/ShrinkNetworkService.cs); [ShrinkPacketCodec](Runtime/Serialization/ShrinkPacketCodec.cs); [ShrinkNetworkWorkQueue](Runtime/Core/ShrinkNetworkWorkQueue.cs)
|
||||||
|
<!-- /shrink:capabilities -->
|
||||||
|
|||||||
@@ -1,17 +1,21 @@
|
|||||||
|
#if UNITY_5_3_OR_NEWER
|
||||||
using UnityEngine;
|
using UnityEngine;
|
||||||
|
#endif
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
public static class ShrinkNetworkRuntime
|
public static class ShrinkNetworkRuntime
|
||||||
{
|
{
|
||||||
public static ShrinkNetworkService Default { get; private set; }
|
public static ShrinkNetworkService Default { get; private set; } = null!;
|
||||||
|
|
||||||
static ShrinkNetworkRuntime()
|
static ShrinkNetworkRuntime()
|
||||||
{
|
{
|
||||||
RebuildDefault();
|
RebuildDefault();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#if UNITY_5_3_OR_NEWER
|
||||||
[RuntimeInitializeOnLoadMethod(RuntimeInitializeLoadType.SubsystemRegistration)]
|
[RuntimeInitializeOnLoadMethod(RuntimeInitializeLoadType.SubsystemRegistration)]
|
||||||
|
#endif
|
||||||
private static void ResetOnPlayModeEnter()
|
private static void ResetOnPlayModeEnter()
|
||||||
{
|
{
|
||||||
RebuildDefault();
|
RebuildDefault();
|
||||||
|
|||||||
@@ -54,7 +54,7 @@ namespace ShrinkNetwork
|
|||||||
Serializer = serializer ?? throw new ArgumentNullException(nameof(serializer));
|
Serializer = serializer ?? throw new ArgumentNullException(nameof(serializer));
|
||||||
MessageRegistry = messageRegistry ?? throw new ArgumentNullException(nameof(messageRegistry));
|
MessageRegistry = messageRegistry ?? throw new ArgumentNullException(nameof(messageRegistry));
|
||||||
Router = router ?? throw new ArgumentNullException(nameof(router));
|
Router = router ?? throw new ArgumentNullException(nameof(router));
|
||||||
_dispatchScheduler = dispatchScheduler ?? ShrinkNetworkDispatchSchedulers.Inline;
|
_dispatchScheduler = dispatchScheduler ?? new ShrinkNetworkInlineDispatchScheduler();
|
||||||
}
|
}
|
||||||
|
|
||||||
public IShrinkNetworkSerializer Serializer { get; }
|
public IShrinkNetworkSerializer Serializer { get; }
|
||||||
@@ -74,6 +74,8 @@ namespace ShrinkNetwork
|
|||||||
set => _dispatchScheduler = value ?? throw new ArgumentNullException(nameof(value));
|
set => _dispatchScheduler = value ?? throw new ArgumentNullException(nameof(value));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>Reliable receive work could not be processed. Session-control transports also disconnect the peer.</summary>
|
||||||
|
public event Action<long>? OnDispatchRejected;
|
||||||
public event Action<ShrinkNetworkSession>? OnSessionConnected;
|
public event Action<ShrinkNetworkSession>? OnSessionConnected;
|
||||||
public event Action<ShrinkNetworkSession>? OnSessionDisconnected;
|
public event Action<ShrinkNetworkSession>? OnSessionDisconnected;
|
||||||
|
|
||||||
@@ -273,11 +275,11 @@ namespace ShrinkNetwork
|
|||||||
if (messageType == null)
|
if (messageType == null)
|
||||||
throw new ArgumentNullException(nameof(messageType));
|
throw new ArgumentNullException(nameof(messageType));
|
||||||
|
|
||||||
return SendPacketAsync(session, messageType, kind, requestToken, route, Serializer.Serialize(message));
|
return SendPacketAsync(session, messageType, kind, requestToken, route, Array.Empty<byte>(), message);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async UniTask SendPacketAsync(ShrinkNetworkSession session, Type messageType,
|
private async UniTask SendPacketAsync(ShrinkNetworkSession session, Type messageType,
|
||||||
ShrinkNetworkPacketKind kind, ShrinkRequestToken requestToken, string? route, byte[] payload)
|
ShrinkNetworkPacketKind kind, ShrinkRequestToken requestToken, string? route, byte[] payload, object? message = null)
|
||||||
{
|
{
|
||||||
if (session == null)
|
if (session == null)
|
||||||
throw new ArgumentNullException(nameof(session));
|
throw new ArgumentNullException(nameof(session));
|
||||||
@@ -299,16 +301,22 @@ namespace ShrinkNetwork
|
|||||||
Payload = payload
|
Payload = payload
|
||||||
};
|
};
|
||||||
|
|
||||||
var packetData = Serializer.Serialize(packet);
|
using var encoded = ShrinkPacketCodec.Encode(packet, Serializer, message);
|
||||||
|
var packetData = encoded.WrittenMemory;
|
||||||
Interlocked.Increment(ref _packetsSent);
|
Interlocked.Increment(ref _packetsSent);
|
||||||
Interlocked.Add(ref _bytesSent, packetData.Length);
|
Interlocked.Add(ref _bytesSent, packetData.Length);
|
||||||
|
if (transport is IShrinkNetworkMemoryTransport memoryTransport)
|
||||||
|
{
|
||||||
|
await memoryTransport.SendAsync(session.SessionId, packetData);
|
||||||
|
return;
|
||||||
|
}
|
||||||
if (transport is IShrinkNetworkAsyncTransport asyncTransport)
|
if (transport is IShrinkNetworkAsyncTransport asyncTransport)
|
||||||
{
|
{
|
||||||
await asyncTransport.SendAsync(session.SessionId, packetData);
|
await asyncTransport.SendAsync(session.SessionId, packetData.ToArray());
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
transport.Send(session.SessionId, packetData);
|
transport.Send(session.SessionId, packetData.ToArray());
|
||||||
}
|
}
|
||||||
|
|
||||||
private void OnTransportEvent(ShrinkNetworkTransportEvent evt)
|
private void OnTransportEvent(ShrinkNetworkTransportEvent evt)
|
||||||
@@ -330,6 +338,12 @@ namespace ShrinkNetwork
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Responses complete pending RPCs independently of the serial handler queue, including nested calls.
|
||||||
|
if (evt.Type == ShrinkNetworkTransportEventType.Packet && ShrinkPacketCodec.IsResponse(evt.PacketData))
|
||||||
|
{
|
||||||
|
HandlePacketAsync(evt).Forget();
|
||||||
|
return;
|
||||||
|
}
|
||||||
ScheduleTransportEventAsync(evt).Forget();
|
ScheduleTransportEventAsync(evt).Forget();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -337,9 +351,23 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
var scheduled = await DispatchScheduler.ScheduleAsync(() => HandleTransportEventAsync(evt));
|
var scheduler = DispatchScheduler;
|
||||||
|
_sessions.TryGetValue(evt.SessionId, out var receivedSession);
|
||||||
|
UniTask DispatchCurrent()
|
||||||
|
{
|
||||||
|
return receivedSession != null && _sessions.TryGetValue(evt.SessionId, out var current) && ReferenceEquals(current, receivedSession)
|
||||||
|
? HandleTransportEventAsync(evt) : UniTask.CompletedTask;
|
||||||
|
}
|
||||||
|
var scheduled = scheduler is IShrinkNetworkPacketDispatchScheduler queue
|
||||||
|
? await queue.ScheduleAsync(evt.SessionId, evt.PacketData?.Length ?? 0, DispatchCurrent)
|
||||||
|
: await scheduler.ScheduleAsync(DispatchCurrent);
|
||||||
if (!scheduled)
|
if (!scheduled)
|
||||||
|
{
|
||||||
Interlocked.Increment(ref _dispatchQueueRejectedCount);
|
Interlocked.Increment(ref _dispatchQueueRejectedCount);
|
||||||
|
OnDispatchRejected?.Invoke(evt.SessionId);
|
||||||
|
if (_transport is IShrinkNetworkSessionControlTransport control)
|
||||||
|
control.DisconnectSession(evt.SessionId, "SHRINK-NET-CONGESTION: reliable receive queue rejected work.");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
catch (Exception ex)
|
catch (Exception ex)
|
||||||
{
|
{
|
||||||
@@ -395,12 +423,12 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
Interlocked.Increment(ref _packetsReceived);
|
Interlocked.Increment(ref _packetsReceived);
|
||||||
Interlocked.Add(ref _bytesReceived, evt.PacketData.Length);
|
Interlocked.Add(ref _bytesReceived, evt.PacketData.Length);
|
||||||
var packet = Serializer.Deserialize<ShrinkNetworkPacket>(evt.PacketData);
|
var packet = ShrinkPacketCodec.Decode(evt.PacketData);
|
||||||
|
|
||||||
if (!_sessions.TryGetValue(evt.SessionId, out var session))
|
if (!_sessions.TryGetValue(evt.SessionId, out var session))
|
||||||
{
|
{
|
||||||
session = _sessions.GetOrAdd(evt.SessionId,
|
// Queued packets from a disconnected session cannot resurrect its state.
|
||||||
id => new ShrinkNetworkSession(id, evt.RemoteAddress, this));
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!ValidatePacketCompatibility(packet, evt.SessionId))
|
if (!ValidatePacketCompatibility(packet, evt.SessionId))
|
||||||
@@ -410,7 +438,7 @@ namespace ShrinkNetwork
|
|||||||
|
|
||||||
if (packet.Kind == ShrinkNetworkPacketKind.Response)
|
if (packet.Kind == ShrinkNetworkPacketKind.Response)
|
||||||
{
|
{
|
||||||
HandleResponse(packet);
|
HandleResponse(session, packet);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -425,7 +453,7 @@ namespace ShrinkNetwork
|
|||||||
}
|
}
|
||||||
|
|
||||||
var resolvedMeta = meta!;
|
var resolvedMeta = meta!;
|
||||||
var message = Serializer.Deserialize(packet.Payload, resolvedMeta.MessageType);
|
var message = DeserializePayload(packet.Payload, resolvedMeta.MessageType);
|
||||||
if (message == null)
|
if (message == null)
|
||||||
{
|
{
|
||||||
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] Failed to deserialize message for opcode {packet.Opcode}.");
|
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] Failed to deserialize message for opcode {packet.Opcode}.");
|
||||||
@@ -442,6 +470,12 @@ namespace ShrinkNetwork
|
|||||||
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] No handler found for {resolvedMeta.MessageType.FullName}");
|
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] No handler found for {resolvedMeta.MessageType.FullName}");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
catch (ShrinkProtocolException ex)
|
||||||
|
{
|
||||||
|
Interlocked.Increment(ref _protocolViolations);
|
||||||
|
if (DisconnectOnProtocolViolation && _transport is IShrinkNetworkSessionControlTransport control)
|
||||||
|
control.DisconnectSession(evt.SessionId, ex.Message);
|
||||||
|
}
|
||||||
catch (Exception ex)
|
catch (Exception ex)
|
||||||
{
|
{
|
||||||
Interlocked.Increment(ref _serializationErrorCount);
|
Interlocked.Increment(ref _serializationErrorCount);
|
||||||
@@ -450,9 +484,10 @@ namespace ShrinkNetwork
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private void HandleResponse(ShrinkNetworkPacket packet)
|
private void HandleResponse(ShrinkNetworkSession session, ShrinkNetworkPacket packet)
|
||||||
{
|
{
|
||||||
if (!_pendingRequests.TryRemove(packet.RequestToken, out var pending))
|
if (!_pendingRequests.TryGetValue(packet.RequestToken, out var expected) || expected.SessionId != session.SessionId ||
|
||||||
|
!_pendingRequests.TryRemove(packet.RequestToken, out var pending))
|
||||||
{
|
{
|
||||||
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] Pending request not found. RequestToken={packet.RequestToken}");
|
ShrinkNetworkLogger.Warn($"[ShrinkNetwork] Pending request not found. RequestToken={packet.RequestToken}");
|
||||||
return;
|
return;
|
||||||
@@ -460,7 +495,7 @@ namespace ShrinkNetwork
|
|||||||
|
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
var response = Serializer.Deserialize(packet.Payload, pending.ResponseType);
|
var response = DeserializePayload(packet.Payload, pending.ResponseType);
|
||||||
pending.CompletionSource.TrySetResult(response);
|
pending.CompletionSource.TrySetResult(response);
|
||||||
}
|
}
|
||||||
catch (Exception ex)
|
catch (Exception ex)
|
||||||
@@ -469,6 +504,10 @@ namespace ShrinkNetwork
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private object DeserializePayload(ReadOnlyMemory<byte> payload, Type type) =>
|
||||||
|
Serializer is IShrinkNetworkBufferSerializer buffered
|
||||||
|
? buffered.Deserialize(payload, type) : Serializer.Deserialize(payload.ToArray(), type);
|
||||||
|
|
||||||
private async UniTask<TResponse> WaitForPendingResponse<TResponse>(ShrinkRequestToken requestToken, PendingRequest pending,
|
private async UniTask<TResponse> WaitForPendingResponse<TResponse>(ShrinkRequestToken requestToken, PendingRequest pending,
|
||||||
ShrinkRpcCallOptions? options)
|
ShrinkRpcCallOptions? options)
|
||||||
where TResponse : class, IShrinkNetworkResponse
|
where TResponse : class, IShrinkNetworkResponse
|
||||||
@@ -627,20 +666,42 @@ namespace ShrinkNetwork
|
|||||||
|
|
||||||
public static class ShrinkNetworkDispatchSchedulers
|
public static class ShrinkNetworkDispatchSchedulers
|
||||||
{
|
{
|
||||||
public static IShrinkNetworkDispatchScheduler Inline { get; } =
|
public static IShrinkNetworkDispatchScheduler Inline => new ShrinkNetworkInlineDispatchScheduler();
|
||||||
new ShrinkNetworkInlineDispatchScheduler();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public sealed class ShrinkNetworkInlineDispatchScheduler : IShrinkNetworkDispatchScheduler
|
public interface IShrinkNetworkPacketDispatchScheduler : IShrinkNetworkDispatchScheduler
|
||||||
{
|
{
|
||||||
public async UniTask<bool> ScheduleAsync(Func<UniTask> callback)
|
UniTask<bool> ScheduleAsync(long sessionId, int bytes, Func<UniTask> callback);
|
||||||
{
|
}
|
||||||
if (callback == null)
|
|
||||||
throw new ArgumentNullException(nameof(callback));
|
|
||||||
|
|
||||||
await callback();
|
/// <summary>Automatically drains a bounded serial queue. Callbacks start on the initiating transport/continuation thread.</summary>
|
||||||
return true;
|
public sealed class ShrinkNetworkInlineDispatchScheduler : IShrinkNetworkPacketDispatchScheduler, IDisposable
|
||||||
|
{
|
||||||
|
private readonly object _gate = new();
|
||||||
|
private readonly ShrinkNetworkWorkQueue _queue = new(4096, 64 * 1024 * 1024);
|
||||||
|
private bool _draining;
|
||||||
|
public ShrinkNetworkQueueDiagnostics CaptureDiagnostics() => _queue.CaptureDiagnostics();
|
||||||
|
public UniTask<bool> ScheduleAsync(Func<UniTask> callback) => ScheduleAsync(0, 0, callback);
|
||||||
|
public async UniTask<bool> ScheduleAsync(long sessionId, int bytes, Func<UniTask> callback)
|
||||||
|
{
|
||||||
|
var completion = _queue.EnqueueAsync(sessionId, "receive", bytes, callback);
|
||||||
|
var start = false;
|
||||||
|
lock (_gate) { if (!_draining) { _draining = true; start = true; } }
|
||||||
|
if (start) DrainAsync().Forget();
|
||||||
|
return await completion == ShrinkNetworkQueueResult.Completed;
|
||||||
}
|
}
|
||||||
|
private async UniTask DrainAsync()
|
||||||
|
{
|
||||||
|
while (true)
|
||||||
|
{
|
||||||
|
await _queue.PumpAsync(int.MaxValue);
|
||||||
|
lock (_gate)
|
||||||
|
{
|
||||||
|
if (_queue.PendingCount == 0) { _draining = false; return; }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
public void Dispose() => _queue.Dispose();
|
||||||
}
|
}
|
||||||
|
|
||||||
public enum ShrinkNetworkDispatchOverflowPolicy
|
public enum ShrinkNetworkDispatchOverflowPolicy
|
||||||
@@ -650,163 +711,30 @@ namespace ShrinkNetwork
|
|||||||
DropOldest = 2
|
DropOldest = 2
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>Caller-pumped serial receive queue; rejects overflow without silently dropping reliable packets.</summary>
|
||||||
/// A caller-pumped, bounded dispatch queue. Unity can pump it from Update
|
public sealed class ShrinkNetworkDispatchQueue : IShrinkNetworkPacketDispatchScheduler, IDisposable
|
||||||
/// while a dedicated server can keep the default inline scheduler.
|
|
||||||
/// </summary>
|
|
||||||
public sealed class ShrinkNetworkDispatchQueue : IShrinkNetworkDispatchScheduler, IDisposable
|
|
||||||
{
|
{
|
||||||
private sealed class WorkItem
|
private readonly ShrinkNetworkWorkQueue _queue;
|
||||||
{
|
|
||||||
public Func<UniTask> Callback = null!;
|
|
||||||
public UniTaskCompletionSource<bool> Completion = null!;
|
|
||||||
}
|
|
||||||
|
|
||||||
private readonly ConcurrentQueue<WorkItem> _queue = new();
|
|
||||||
private readonly object _lifecycleLock = new();
|
|
||||||
private readonly int _capacity;
|
|
||||||
private readonly ShrinkNetworkDispatchOverflowPolicy _overflowPolicy;
|
|
||||||
private int _queuedCount;
|
|
||||||
private int _pumping;
|
|
||||||
private int _disposed;
|
|
||||||
private long _rejectedCount;
|
|
||||||
private long _droppedCount;
|
|
||||||
|
|
||||||
public ShrinkNetworkDispatchQueue(int capacity,
|
public ShrinkNetworkDispatchQueue(int capacity,
|
||||||
ShrinkNetworkDispatchOverflowPolicy overflowPolicy = ShrinkNetworkDispatchOverflowPolicy.Reject)
|
ShrinkNetworkDispatchOverflowPolicy overflowPolicy = ShrinkNetworkDispatchOverflowPolicy.Reject,
|
||||||
|
long byteCapacity = 64 * 1024 * 1024, int perSessionCapacity = int.MaxValue)
|
||||||
{
|
{
|
||||||
if (capacity <= 0)
|
if (overflowPolicy != ShrinkNetworkDispatchOverflowPolicy.Reject)
|
||||||
throw new ArgumentOutOfRangeException(nameof(capacity));
|
throw new ArgumentException("Protocol v2 reliable dispatch requires Reject. Use an explicit state key on ShrinkNetworkWorkQueue for replaceable state.", nameof(overflowPolicy));
|
||||||
|
Capacity = capacity;
|
||||||
_capacity = capacity;
|
_queue = new ShrinkNetworkWorkQueue(capacity, byteCapacity, perSessionCapacity);
|
||||||
_overflowPolicy = overflowPolicy;
|
|
||||||
}
|
|
||||||
|
|
||||||
public int Capacity => _capacity;
|
|
||||||
public int PendingCount => Volatile.Read(ref _queuedCount);
|
|
||||||
public long RejectedCount => Volatile.Read(ref _rejectedCount);
|
|
||||||
public long DroppedCount => Volatile.Read(ref _droppedCount);
|
|
||||||
|
|
||||||
public UniTask<bool> ScheduleAsync(Func<UniTask> callback)
|
|
||||||
{
|
|
||||||
if (callback == null)
|
|
||||||
throw new ArgumentNullException(nameof(callback));
|
|
||||||
if (Volatile.Read(ref _disposed) != 0)
|
|
||||||
return UniTask.FromException<bool>(new ObjectDisposedException(nameof(ShrinkNetworkDispatchQueue)));
|
|
||||||
|
|
||||||
var item = new WorkItem
|
|
||||||
{
|
|
||||||
Callback = callback,
|
|
||||||
Completion = new UniTaskCompletionSource<bool>()
|
|
||||||
};
|
|
||||||
|
|
||||||
while (true)
|
|
||||||
{
|
|
||||||
if (Volatile.Read(ref _queuedCount) >= _capacity)
|
|
||||||
{
|
|
||||||
switch (_overflowPolicy)
|
|
||||||
{
|
|
||||||
case ShrinkNetworkDispatchOverflowPolicy.Reject:
|
|
||||||
Interlocked.Increment(ref _rejectedCount);
|
|
||||||
return UniTask.FromResult(false);
|
|
||||||
case ShrinkNetworkDispatchOverflowPolicy.DropNewest:
|
|
||||||
Interlocked.Increment(ref _droppedCount);
|
|
||||||
return UniTask.FromResult(false);
|
|
||||||
case ShrinkNetworkDispatchOverflowPolicy.DropOldest:
|
|
||||||
if (_queue.TryDequeue(out var dropped))
|
|
||||||
{
|
|
||||||
Interlocked.Decrement(ref _queuedCount);
|
|
||||||
Interlocked.Increment(ref _droppedCount);
|
|
||||||
dropped.Completion.TrySetResult(false);
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
Thread.Yield();
|
|
||||||
continue;
|
|
||||||
default:
|
|
||||||
throw new ArgumentOutOfRangeException();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
var currentCount = Volatile.Read(ref _queuedCount);
|
|
||||||
if (currentCount >= _capacity ||
|
|
||||||
Interlocked.CompareExchange(ref _queuedCount, currentCount + 1, currentCount) != currentCount)
|
|
||||||
{
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
lock (_lifecycleLock)
|
|
||||||
{
|
|
||||||
if (Volatile.Read(ref _disposed) != 0)
|
|
||||||
{
|
|
||||||
Interlocked.Decrement(ref _queuedCount);
|
|
||||||
item.Completion.TrySetResult(false);
|
|
||||||
return item.Completion.Task;
|
|
||||||
}
|
|
||||||
|
|
||||||
_queue.Enqueue(item);
|
|
||||||
return item.Completion.Task;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
public UniTask<int> PumpAsync(int maxItems)
|
|
||||||
{
|
|
||||||
if (maxItems <= 0)
|
|
||||||
throw new ArgumentOutOfRangeException(nameof(maxItems));
|
|
||||||
if (Interlocked.Exchange(ref _pumping, 1) == 1)
|
|
||||||
return UniTask.FromResult(0);
|
|
||||||
|
|
||||||
return PumpCoreAsync(maxItems);
|
|
||||||
}
|
|
||||||
|
|
||||||
public void Dispose()
|
|
||||||
{
|
|
||||||
lock (_lifecycleLock)
|
|
||||||
{
|
|
||||||
if (Interlocked.Exchange(ref _disposed, 1) != 0)
|
|
||||||
return;
|
|
||||||
|
|
||||||
while (_queue.TryDequeue(out var item))
|
|
||||||
{
|
|
||||||
Interlocked.Decrement(ref _queuedCount);
|
|
||||||
item.Completion.TrySetResult(false);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private async UniTask<int> PumpCoreAsync(int maxItems)
|
|
||||||
{
|
|
||||||
var processed = 0;
|
|
||||||
try
|
|
||||||
{
|
|
||||||
while (processed < maxItems && _queue.TryDequeue(out var item))
|
|
||||||
{
|
|
||||||
Interlocked.Decrement(ref _queuedCount);
|
|
||||||
await ExecuteItemAsync(item);
|
|
||||||
processed++;
|
|
||||||
}
|
|
||||||
|
|
||||||
return processed;
|
|
||||||
}
|
|
||||||
finally
|
|
||||||
{
|
|
||||||
Volatile.Write(ref _pumping, 0);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private static async UniTask ExecuteItemAsync(WorkItem item)
|
|
||||||
{
|
|
||||||
try
|
|
||||||
{
|
|
||||||
await item.Callback();
|
|
||||||
item.Completion.TrySetResult(true);
|
|
||||||
}
|
|
||||||
catch (Exception ex)
|
|
||||||
{
|
|
||||||
item.Completion.TrySetException(ex);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
public int Capacity { get; }
|
||||||
|
public int PendingCount => _queue.PendingCount;
|
||||||
|
public long RejectedCount => CaptureDiagnostics().Rejected;
|
||||||
|
public long DroppedCount => 0;
|
||||||
|
public ShrinkNetworkQueueDiagnostics CaptureDiagnostics() => _queue.CaptureDiagnostics();
|
||||||
|
public UniTask<bool> ScheduleAsync(Func<UniTask> callback) => ScheduleAsync(0, 0, callback);
|
||||||
|
public async UniTask<bool> ScheduleAsync(long sessionId, int bytes, Func<UniTask> callback) =>
|
||||||
|
await _queue.EnqueueAsync(sessionId, "receive", bytes, callback) == ShrinkNetworkQueueResult.Completed;
|
||||||
|
public UniTask<int> PumpAsync(int maxItems) => _queue.PumpAsync(maxItems);
|
||||||
|
public UniTask<int> PumpAsync(int maxItems, long maxBytes, TimeSpan timeBudget) => _queue.PumpAsync(maxItems, maxBytes, timeBudget);
|
||||||
|
public void Dispose() => _queue.Dispose();
|
||||||
}
|
}
|
||||||
|
|
||||||
public sealed class ShrinkNetworkServiceDiagnosticsSnapshot
|
public sealed class ShrinkNetworkServiceDiagnosticsSnapshot
|
||||||
|
|||||||
@@ -0,0 +1,159 @@
|
|||||||
|
#nullable enable
|
||||||
|
using System;
|
||||||
|
using System.Collections.Generic;
|
||||||
|
using System.Diagnostics;
|
||||||
|
using System.Threading;
|
||||||
|
using Cysharp.Threading.Tasks;
|
||||||
|
|
||||||
|
namespace ShrinkNetwork
|
||||||
|
{
|
||||||
|
public enum ShrinkNetworkQueueResult { Completed, Rejected, Replaced, Canceled }
|
||||||
|
|
||||||
|
public sealed class ShrinkNetworkQueueDiagnostics
|
||||||
|
{
|
||||||
|
public int PendingCount { get; internal set; }
|
||||||
|
public long PendingBytes { get; internal set; }
|
||||||
|
public long Rejected { get; internal set; }
|
||||||
|
public long Replaced { get; internal set; }
|
||||||
|
public long Completed { get; internal set; }
|
||||||
|
public double OldestWaitMilliseconds { get; internal set; }
|
||||||
|
public double LastWaitMilliseconds { get; internal set; }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>Caller-pumped serial execution, fair across session/channel partitions. Only queued state is replaceable.</summary>
|
||||||
|
public sealed class ShrinkNetworkWorkQueue : IDisposable
|
||||||
|
{
|
||||||
|
private sealed class Work
|
||||||
|
{
|
||||||
|
public Func<UniTask> Callback = null!;
|
||||||
|
public UniTaskCompletionSource<ShrinkNetworkQueueResult> Completion = new();
|
||||||
|
public string? StateKey;
|
||||||
|
public int Bytes;
|
||||||
|
public long Enqueued = Stopwatch.GetTimestamp();
|
||||||
|
public CancellationToken Cancellation;
|
||||||
|
}
|
||||||
|
private sealed class Partition
|
||||||
|
{
|
||||||
|
public readonly LinkedList<Work> Items = new();
|
||||||
|
public readonly Dictionary<string, LinkedListNode<Work>> States = new(StringComparer.Ordinal);
|
||||||
|
}
|
||||||
|
private readonly object _gate = new();
|
||||||
|
private readonly Dictionary<(long, string), Partition> _partitions = new();
|
||||||
|
private readonly Queue<(long, string)> _ready = new();
|
||||||
|
private readonly int _capacity;
|
||||||
|
private readonly long _byteCapacity;
|
||||||
|
private readonly int _partitionCapacity;
|
||||||
|
private int _count, _pumping;
|
||||||
|
private long _bytes, _rejected, _replaced, _completed;
|
||||||
|
private double _lastWait;
|
||||||
|
private bool _disposed;
|
||||||
|
|
||||||
|
public int PendingCount { get { lock (_gate) return _count; } }
|
||||||
|
|
||||||
|
public ShrinkNetworkWorkQueue(int capacity, long byteCapacity, int perPartitionCapacity = int.MaxValue)
|
||||||
|
{
|
||||||
|
if (capacity <= 0 || byteCapacity <= 0 || perPartitionCapacity <= 0) throw new ArgumentOutOfRangeException(nameof(capacity));
|
||||||
|
_capacity = capacity; _byteCapacity = byteCapacity; _partitionCapacity = perPartitionCapacity;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>A nonempty stateKey explicitly permits replacing an unsent state in this session/channel. Never use for RPC or snapshot fragments.</summary>
|
||||||
|
public UniTask<ShrinkNetworkQueueResult> EnqueueAsync(long sessionId, string channel, int byteCount,
|
||||||
|
Func<UniTask> callback, string? stateKey = null, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
if (callback == null) throw new ArgumentNullException(nameof(callback));
|
||||||
|
if (channel == null) throw new ArgumentNullException(nameof(channel));
|
||||||
|
if (byteCount < 0) throw new ArgumentOutOfRangeException(nameof(byteCount));
|
||||||
|
Work? replaced = null;
|
||||||
|
Work item;
|
||||||
|
lock (_gate)
|
||||||
|
{
|
||||||
|
if (_disposed || cancellationToken.IsCancellationRequested) return UniTask.FromResult(ShrinkNetworkQueueResult.Canceled);
|
||||||
|
var key = (sessionId, channel);
|
||||||
|
_partitions.TryGetValue(key, out var partition);
|
||||||
|
LinkedListNode<Work>? old = null;
|
||||||
|
if (!string.IsNullOrEmpty(stateKey)) partition?.States.TryGetValue(stateKey!, out old);
|
||||||
|
var nextBytes = _bytes - (old?.Value.Bytes ?? 0) + byteCount;
|
||||||
|
if (nextBytes > _byteCapacity || (old == null && (_count >= _capacity || (partition?.Items.Count ?? 0) >= _partitionCapacity)))
|
||||||
|
{ _rejected++; return UniTask.FromResult(ShrinkNetworkQueueResult.Rejected); }
|
||||||
|
item = new Work { Callback = callback, Bytes = byteCount, StateKey = string.IsNullOrEmpty(stateKey) ? null : stateKey, Cancellation = cancellationToken };
|
||||||
|
if (partition == null) { partition = new Partition(); _partitions.Add(key, partition); _ready.Enqueue(key); }
|
||||||
|
if (old != null)
|
||||||
|
{
|
||||||
|
replaced = old.Value;
|
||||||
|
// Move a replacement to the tail: later state must not jump ahead of intervening reliable operations.
|
||||||
|
partition.Items.Remove(old); _replaced++;
|
||||||
|
}
|
||||||
|
else _count++;
|
||||||
|
var node = partition.Items.AddLast(item);
|
||||||
|
if (item.StateKey != null) partition.States[item.StateKey] = node;
|
||||||
|
_bytes = nextBytes;
|
||||||
|
}
|
||||||
|
replaced?.Completion.TrySetResult(ShrinkNetworkQueueResult.Replaced);
|
||||||
|
return item.Completion.Task;
|
||||||
|
}
|
||||||
|
|
||||||
|
public async UniTask<int> PumpAsync(int maxItems, long maxBytes = long.MaxValue, TimeSpan? timeBudget = null)
|
||||||
|
{
|
||||||
|
if (maxItems <= 0 || maxBytes <= 0 || (timeBudget.HasValue && timeBudget.Value <= TimeSpan.Zero)) throw new ArgumentOutOfRangeException(nameof(maxItems));
|
||||||
|
if (Interlocked.Exchange(ref _pumping, 1) != 0) return 0;
|
||||||
|
var started = Stopwatch.GetTimestamp();
|
||||||
|
var processed = 0;
|
||||||
|
long bytes = 0;
|
||||||
|
try
|
||||||
|
{
|
||||||
|
while (processed < maxItems && (!timeBudget.HasValue || Elapsed(started) < timeBudget.Value.TotalMilliseconds))
|
||||||
|
{
|
||||||
|
Work item;
|
||||||
|
lock (_gate)
|
||||||
|
{
|
||||||
|
if (_disposed || _ready.Count == 0) break;
|
||||||
|
var key = _ready.Peek();
|
||||||
|
var partition = _partitions[key];
|
||||||
|
item = partition.Items.First!.Value;
|
||||||
|
// Allow one oversized item so a byte budget cannot permanently starve a valid packet.
|
||||||
|
if (processed > 0 && item.Bytes > maxBytes - bytes) break;
|
||||||
|
_ready.Dequeue(); partition.Items.RemoveFirst();
|
||||||
|
if (item.StateKey != null) partition.States.Remove(item.StateKey);
|
||||||
|
if (partition.Items.Count == 0) _partitions.Remove(key); else _ready.Enqueue(key);
|
||||||
|
_count--; _bytes -= item.Bytes; _lastWait = Elapsed(item.Enqueued);
|
||||||
|
}
|
||||||
|
try
|
||||||
|
{
|
||||||
|
if (item.Cancellation.IsCancellationRequested) item.Completion.TrySetResult(ShrinkNetworkQueueResult.Canceled);
|
||||||
|
else { await item.Callback(); item.Completion.TrySetResult(ShrinkNetworkQueueResult.Completed); lock (_gate) _completed++; }
|
||||||
|
}
|
||||||
|
catch (OperationCanceledException) { item.Completion.TrySetResult(ShrinkNetworkQueueResult.Canceled); }
|
||||||
|
catch (Exception ex) { item.Completion.TrySetException(ex); }
|
||||||
|
processed++; bytes += item.Bytes;
|
||||||
|
}
|
||||||
|
return processed;
|
||||||
|
}
|
||||||
|
finally { Volatile.Write(ref _pumping, 0); }
|
||||||
|
}
|
||||||
|
|
||||||
|
public ShrinkNetworkQueueDiagnostics CaptureDiagnostics()
|
||||||
|
{
|
||||||
|
lock (_gate)
|
||||||
|
{
|
||||||
|
long oldest = Stopwatch.GetTimestamp();
|
||||||
|
foreach (var partition in _partitions.Values)
|
||||||
|
if (partition.Items.First != null) oldest = Math.Min(oldest, partition.Items.First.Value.Enqueued);
|
||||||
|
return new ShrinkNetworkQueueDiagnostics { PendingCount = _count, PendingBytes = _bytes, Rejected = _rejected, Replaced = _replaced,
|
||||||
|
Completed = _completed, LastWaitMilliseconds = _lastWait, OldestWaitMilliseconds = _count == 0 ? 0 : Elapsed(oldest) };
|
||||||
|
}
|
||||||
|
}
|
||||||
|
private static double Elapsed(long start) => (Stopwatch.GetTimestamp() - start) * 1000d / Stopwatch.Frequency;
|
||||||
|
public void Dispose()
|
||||||
|
{
|
||||||
|
List<Work> canceled = new();
|
||||||
|
lock (_gate)
|
||||||
|
{
|
||||||
|
if (_disposed) return;
|
||||||
|
_disposed = true;
|
||||||
|
foreach (var partition in _partitions.Values) canceled.AddRange(partition.Items);
|
||||||
|
_partitions.Clear(); _ready.Clear(); _count = 0; _bytes = 0;
|
||||||
|
}
|
||||||
|
foreach (var item in canceled) item.Completion.TrySetResult(ShrinkNetworkQueueResult.Canceled);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 3b35a44b897396549b5a1dec58ae05c0
|
||||||
|
MonoImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
serializedVersion: 2
|
||||||
|
defaultReferences: []
|
||||||
|
executionOrder: 0
|
||||||
|
icon: {instanceID: 0}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -6,7 +6,7 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
public static class ShrinkNetworkProtocol
|
public static class ShrinkNetworkProtocol
|
||||||
{
|
{
|
||||||
public const int CurrentProtocolVersion = 1;
|
public const int CurrentProtocolVersion = 2;
|
||||||
public const int CurrentSchemaVersion = 1;
|
public const int CurrentSchemaVersion = 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -28,6 +28,6 @@ namespace ShrinkNetwork
|
|||||||
public long SessionTokenExpiresAtUnixTimeSeconds;
|
public long SessionTokenExpiresAtUnixTimeSeconds;
|
||||||
public string? Route;
|
public string? Route;
|
||||||
public ShrinkNetworkPacketKind Kind;
|
public ShrinkNetworkPacketKind Kind;
|
||||||
public byte[] Payload = Array.Empty<byte>();
|
public ReadOnlyMemory<byte> Payload = ReadOnlyMemory<byte>.Empty;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,6 +9,8 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
private readonly Dictionary<int, ShrinkNetworkMessageMeta> _opcodeToMeta = new();
|
private readonly Dictionary<int, ShrinkNetworkMessageMeta> _opcodeToMeta = new();
|
||||||
private readonly Dictionary<Type, ShrinkNetworkMessageMeta> _typeToMeta = new();
|
private readonly Dictionary<Type, ShrinkNetworkMessageMeta> _typeToMeta = new();
|
||||||
|
private readonly HashSet<string> _routes = new(StringComparer.Ordinal);
|
||||||
|
public IReadOnlyCollection<ShrinkNetworkMessageMeta> Registrations => _typeToMeta.Values;
|
||||||
|
|
||||||
public void Register<TMessage>(int opcode, string? route = null) where TMessage : IShrinkNetworkMessage
|
public void Register<TMessage>(int opcode, string? route = null) where TMessage : IShrinkNetworkMessage
|
||||||
=> Register(typeof(TMessage), opcode, route);
|
=> Register(typeof(TMessage), opcode, route);
|
||||||
@@ -24,6 +26,8 @@ namespace ShrinkNetwork
|
|||||||
if (_typeToMeta.ContainsKey(messageType))
|
if (_typeToMeta.ContainsKey(messageType))
|
||||||
throw new InvalidOperationException($"Message type {messageType.FullName} is already registered.");
|
throw new InvalidOperationException($"Message type {messageType.FullName} is already registered.");
|
||||||
|
|
||||||
|
route = string.IsNullOrWhiteSpace(route) ? null : route.Trim();
|
||||||
|
if (route != null && !_routes.Add(route)) throw new InvalidOperationException($"SHRINK002 Route '{route}' is already registered.");
|
||||||
var meta = new ShrinkNetworkMessageMeta(opcode, messageType, route);
|
var meta = new ShrinkNetworkMessageMeta(opcode, messageType, route);
|
||||||
_opcodeToMeta.Add(opcode, meta);
|
_opcodeToMeta.Add(opcode, meta);
|
||||||
_typeToMeta.Add(messageType, meta);
|
_typeToMeta.Add(messageType, meta);
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
|
|
||||||
using System;
|
using System;
|
||||||
using System.Collections.Generic;
|
using System.Collections.Generic;
|
||||||
|
using System.Reflection;
|
||||||
using Cysharp.Threading.Tasks;
|
using Cysharp.Threading.Tasks;
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
@@ -21,6 +22,8 @@ namespace ShrinkNetwork
|
|||||||
public Func<ShrinkNetworkContext, object, UniTask<object?>> Handler = null!;
|
public Func<ShrinkNetworkContext, object, UniTask<object?>> Handler = null!;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private readonly Dictionary<Type, MethodInfo> _bindingMethods = new();
|
||||||
|
public IReadOnlyDictionary<Type, MethodInfo> CaptureBindings() => new Dictionary<Type, MethodInfo>(_bindingMethods);
|
||||||
private readonly Dictionary<Type, MessageHandlerRegistration> _messageHandlers = new();
|
private readonly Dictionary<Type, MessageHandlerRegistration> _messageHandlers = new();
|
||||||
private readonly Dictionary<Type, RequestHandlerRegistration> _requestHandlers = new();
|
private readonly Dictionary<Type, RequestHandlerRegistration> _requestHandlers = new();
|
||||||
|
|
||||||
@@ -28,7 +31,9 @@ namespace ShrinkNetwork
|
|||||||
ShrinkNetworkPermissionRequirement requirement = default)
|
ShrinkNetworkPermissionRequirement requirement = default)
|
||||||
where TMessage : IShrinkNetworkMessage
|
where TMessage : IShrinkNetworkMessage
|
||||||
{
|
{
|
||||||
|
if (handler == null) throw new ArgumentNullException(nameof(handler));
|
||||||
RegisterHandler(typeof(TMessage), (context, message) => handler(context, (TMessage)message), requirement);
|
RegisterHandler(typeof(TMessage), (context, message) => handler(context, (TMessage)message), requirement);
|
||||||
|
_bindingMethods[typeof(TMessage)] = handler.Method;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void RegisterHandler(Type messageType, Func<ShrinkNetworkContext, object, UniTask> handler,
|
public void RegisterHandler(Type messageType, Func<ShrinkNetworkContext, object, UniTask> handler,
|
||||||
@@ -41,6 +46,7 @@ namespace ShrinkNetwork
|
|||||||
if (_messageHandlers.ContainsKey(messageType) || _requestHandlers.ContainsKey(messageType))
|
if (_messageHandlers.ContainsKey(messageType) || _requestHandlers.ContainsKey(messageType))
|
||||||
throw new InvalidOperationException($"Handler already exists for {messageType.FullName}.");
|
throw new InvalidOperationException($"Handler already exists for {messageType.FullName}.");
|
||||||
|
|
||||||
|
_bindingMethods[messageType] = handler.Method;
|
||||||
_messageHandlers.Add(messageType, new MessageHandlerRegistration
|
_messageHandlers.Add(messageType, new MessageHandlerRegistration
|
||||||
{
|
{
|
||||||
Requirement = requirement,
|
Requirement = requirement,
|
||||||
@@ -53,8 +59,10 @@ namespace ShrinkNetwork
|
|||||||
where TRequest : IShrinkNetworkRequest
|
where TRequest : IShrinkNetworkRequest
|
||||||
where TResponse : class, IShrinkNetworkResponse
|
where TResponse : class, IShrinkNetworkResponse
|
||||||
{
|
{
|
||||||
|
if (handler == null) throw new ArgumentNullException(nameof(handler));
|
||||||
RegisterRequestHandler(typeof(TRequest), typeof(TResponse),
|
RegisterRequestHandler(typeof(TRequest), typeof(TResponse),
|
||||||
async (context, message) => await handler(context, (TRequest)message), requirement);
|
async (context, message) => await handler(context, (TRequest)message), requirement);
|
||||||
|
_bindingMethods[typeof(TRequest)] = handler.Method;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void RegisterRequestHandler(Type requestType, Type responseType,
|
public void RegisterRequestHandler(Type requestType, Type responseType,
|
||||||
@@ -70,6 +78,7 @@ namespace ShrinkNetwork
|
|||||||
if (_messageHandlers.ContainsKey(requestType) || _requestHandlers.ContainsKey(requestType))
|
if (_messageHandlers.ContainsKey(requestType) || _requestHandlers.ContainsKey(requestType))
|
||||||
throw new InvalidOperationException($"Handler already exists for {requestType.FullName}.");
|
throw new InvalidOperationException($"Handler already exists for {requestType.FullName}.");
|
||||||
|
|
||||||
|
_bindingMethods[requestType] = handler.Method;
|
||||||
_requestHandlers.Add(requestType, new RequestHandlerRegistration
|
_requestHandlers.Add(requestType, new RequestHandlerRegistration
|
||||||
{
|
{
|
||||||
ResponseType = responseType,
|
ResponseType = responseType,
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
using System;
|
using System;
|
||||||
|
using System.Buffers;
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
@@ -8,4 +9,10 @@ namespace ShrinkNetwork
|
|||||||
object Deserialize(byte[] payload, Type type);
|
object Deserialize(byte[] payload, Type type);
|
||||||
T Deserialize<T>(byte[] payload);
|
T Deserialize<T>(byte[] payload);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public interface IShrinkNetworkBufferSerializer : IShrinkNetworkSerializer
|
||||||
|
{
|
||||||
|
void Serialize(IBufferWriter<byte> writer, object value);
|
||||||
|
object Deserialize(ReadOnlyMemory<byte> payload, Type type);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,40 @@
|
|||||||
|
#nullable enable
|
||||||
|
using System;
|
||||||
|
using System.Buffers;
|
||||||
|
|
||||||
|
namespace ShrinkNetwork
|
||||||
|
{
|
||||||
|
/// <summary>Owns rented memory until Dispose. A sender must await completion before disposing.</summary>
|
||||||
|
public sealed class ShrinkBufferWriter : IBufferWriter<byte>, IDisposable
|
||||||
|
{
|
||||||
|
private byte[]? _buffer;
|
||||||
|
public ShrinkBufferWriter(int initialCapacity = 256) => _buffer = ArrayPool<byte>.Shared.Rent(Math.Max(1, initialCapacity));
|
||||||
|
public int WrittenCount { get; private set; }
|
||||||
|
public ReadOnlyMemory<byte> WrittenMemory => Buffer.AsMemory(0, WrittenCount);
|
||||||
|
internal Span<byte> WrittenSpan => Buffer.AsSpan(0, WrittenCount);
|
||||||
|
private byte[] Buffer => _buffer ?? throw new ObjectDisposedException(nameof(ShrinkBufferWriter));
|
||||||
|
public void Advance(int count)
|
||||||
|
{
|
||||||
|
if (count < 0 || count > Buffer.Length - WrittenCount) throw new ArgumentOutOfRangeException(nameof(count));
|
||||||
|
WrittenCount += count;
|
||||||
|
}
|
||||||
|
public Memory<byte> GetMemory(int sizeHint = 0) { Ensure(sizeHint); return Buffer.AsMemory(WrittenCount); }
|
||||||
|
public Span<byte> GetSpan(int sizeHint = 0) { Ensure(sizeHint); return Buffer.AsSpan(WrittenCount); }
|
||||||
|
private void Ensure(int sizeHint)
|
||||||
|
{
|
||||||
|
if (sizeHint < 0) throw new ArgumentOutOfRangeException(nameof(sizeHint));
|
||||||
|
sizeHint = Math.Max(1, sizeHint);
|
||||||
|
if (sizeHint <= Buffer.Length - WrittenCount) return;
|
||||||
|
var next = ArrayPool<byte>.Shared.Rent(checked(Math.Max(Buffer.Length * 2, WrittenCount + sizeHint)));
|
||||||
|
Buffer.AsSpan(0, WrittenCount).CopyTo(next);
|
||||||
|
ArrayPool<byte>.Shared.Return(Buffer, clearArray: true);
|
||||||
|
_buffer = next;
|
||||||
|
}
|
||||||
|
public void Dispose()
|
||||||
|
{
|
||||||
|
var buffer = _buffer;
|
||||||
|
_buffer = null;
|
||||||
|
if (buffer != null) ArrayPool<byte>.Shared.Return(buffer, clearArray: true);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 9ae3dee10fcf29f4da1b7ace3063bf92
|
||||||
|
MonoImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
serializedVersion: 2
|
||||||
|
defaultReferences: []
|
||||||
|
executionOrder: 0
|
||||||
|
icon: {instanceID: 0}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -2,11 +2,13 @@
|
|||||||
|
|
||||||
using System;
|
using System;
|
||||||
using System.Text;
|
using System.Text;
|
||||||
|
using System.Buffers;
|
||||||
|
using System.IO;
|
||||||
using Newtonsoft.Json;
|
using Newtonsoft.Json;
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
public sealed class ShrinkJsonNetworkSerializer : IShrinkNetworkSerializer
|
public sealed class ShrinkJsonNetworkSerializer : IShrinkNetworkBufferSerializer
|
||||||
{
|
{
|
||||||
private static readonly JsonSerializerSettings Settings = new()
|
private static readonly JsonSerializerSettings Settings = new()
|
||||||
{
|
{
|
||||||
@@ -25,5 +27,33 @@ namespace ShrinkNetwork
|
|||||||
public T Deserialize<T>(byte[] payload)
|
public T Deserialize<T>(byte[] payload)
|
||||||
=> JsonConvert.DeserializeObject<T>(Encoding.UTF8.GetString(payload), Settings)
|
=> JsonConvert.DeserializeObject<T>(Encoding.UTF8.GetString(payload), Settings)
|
||||||
?? throw new JsonSerializationException($"Failed to deserialize payload into {typeof(T).FullName}.");
|
?? throw new JsonSerializationException($"Failed to deserialize payload into {typeof(T).FullName}.");
|
||||||
|
|
||||||
|
public void Serialize(IBufferWriter<byte> writer, object value)
|
||||||
|
{
|
||||||
|
using var stream = new BufferStream(writer);
|
||||||
|
using var text = new StreamWriter(stream, new UTF8Encoding(false), 1024, leaveOpen: true);
|
||||||
|
using var json = new JsonTextWriter(text);
|
||||||
|
JsonSerializer.Create(Settings).Serialize(json, value);
|
||||||
|
}
|
||||||
|
|
||||||
|
public object Deserialize(ReadOnlyMemory<byte> payload, Type type) =>
|
||||||
|
JsonConvert.DeserializeObject(Encoding.UTF8.GetString(payload.Span), type, Settings)
|
||||||
|
?? throw new JsonSerializationException($"Failed to deserialize payload into {type.FullName}.");
|
||||||
|
|
||||||
|
private sealed class BufferStream : Stream
|
||||||
|
{
|
||||||
|
private readonly IBufferWriter<byte> _writer;
|
||||||
|
public BufferStream(IBufferWriter<byte> writer) => _writer = writer;
|
||||||
|
public override bool CanRead => false;
|
||||||
|
public override bool CanSeek => false;
|
||||||
|
public override bool CanWrite => true;
|
||||||
|
public override long Length => throw new NotSupportedException();
|
||||||
|
public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); }
|
||||||
|
public override void Flush() { }
|
||||||
|
public override void Write(byte[] buffer, int offset, int count) { buffer.AsSpan(offset, count).CopyTo(_writer.GetSpan(count)); _writer.Advance(count); }
|
||||||
|
public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException();
|
||||||
|
public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();
|
||||||
|
public override void SetLength(long value) => throw new NotSupportedException();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,150 +0,0 @@
|
|||||||
#nullable enable
|
|
||||||
|
|
||||||
using System;
|
|
||||||
using System.Linq;
|
|
||||||
using System.Reflection;
|
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
|
||||||
{
|
|
||||||
public sealed class ShrinkMessagePackNetworkSerializer : IShrinkNetworkSerializer
|
|
||||||
{
|
|
||||||
private readonly MethodInfo _serializeMethod;
|
|
||||||
private readonly MethodInfo _deserializeMethod;
|
|
||||||
private readonly object? _serializerOptions;
|
|
||||||
|
|
||||||
public ShrinkMessagePackNetworkSerializer()
|
|
||||||
{
|
|
||||||
var serializerType = Type.GetType("MessagePack.MessagePackSerializer, MessagePack");
|
|
||||||
if (serializerType == null)
|
|
||||||
{
|
|
||||||
throw new InvalidOperationException(
|
|
||||||
"MessagePack assembly was not found. Please install MessagePack-CSharp before using ShrinkMessagePackNetworkSerializer.");
|
|
||||||
}
|
|
||||||
|
|
||||||
_serializerOptions = ResolveSerializerOptions(serializerType.Assembly);
|
|
||||||
|
|
||||||
var serializeMethod = serializerType
|
|
||||||
.GetMethods(BindingFlags.Public | BindingFlags.Static)
|
|
||||||
.FirstOrDefault(m =>
|
|
||||||
{
|
|
||||||
if (m.Name != "Serialize")
|
|
||||||
return false;
|
|
||||||
var parameters = m.GetParameters();
|
|
||||||
return parameters.Length >= 2 &&
|
|
||||||
parameters[0].ParameterType == typeof(Type) &&
|
|
||||||
parameters[1].ParameterType == typeof(object);
|
|
||||||
});
|
|
||||||
|
|
||||||
var deserializeMethod = serializerType
|
|
||||||
.GetMethods(BindingFlags.Public | BindingFlags.Static)
|
|
||||||
.FirstOrDefault(m =>
|
|
||||||
{
|
|
||||||
if (m.Name != "Deserialize")
|
|
||||||
return false;
|
|
||||||
var parameters = m.GetParameters();
|
|
||||||
return parameters.Length >= 2 &&
|
|
||||||
parameters[0].ParameterType == typeof(Type) &&
|
|
||||||
(parameters[1].ParameterType == typeof(byte[]) ||
|
|
||||||
parameters[1].ParameterType == typeof(ReadOnlyMemory<byte>));
|
|
||||||
});
|
|
||||||
|
|
||||||
if (serializeMethod == null || deserializeMethod == null)
|
|
||||||
throw new MissingMethodException("MessagePack serialize/deserialize API not found.");
|
|
||||||
|
|
||||||
_serializeMethod = serializeMethod;
|
|
||||||
_deserializeMethod = deserializeMethod;
|
|
||||||
}
|
|
||||||
|
|
||||||
public byte[] Serialize(object value)
|
|
||||||
{
|
|
||||||
if (value == null)
|
|
||||||
return Array.Empty<byte>();
|
|
||||||
|
|
||||||
var parameters = BuildParameters(_serializeMethod, value.GetType(), value, _serializerOptions);
|
|
||||||
return (byte[])_serializeMethod.Invoke(null, parameters)!;
|
|
||||||
}
|
|
||||||
|
|
||||||
public object Deserialize(byte[] payload, Type type)
|
|
||||||
{
|
|
||||||
payload ??= Array.Empty<byte>();
|
|
||||||
var parameters = BuildParameters(_deserializeMethod, type, payload, _serializerOptions);
|
|
||||||
return _deserializeMethod.Invoke(null, parameters)
|
|
||||||
?? throw new InvalidOperationException($"MessagePack returned null for type {type.FullName}.");
|
|
||||||
}
|
|
||||||
|
|
||||||
public T Deserialize<T>(byte[] payload)
|
|
||||||
{
|
|
||||||
return (T)Deserialize(payload, typeof(T));
|
|
||||||
}
|
|
||||||
|
|
||||||
private static object?[] BuildParameters(MethodInfo method, Type type, object valueOrBytes, object? serializerOptions)
|
|
||||||
{
|
|
||||||
var parameters = method.GetParameters();
|
|
||||||
var args = new object?[parameters.Length];
|
|
||||||
|
|
||||||
if (parameters.Length > 0)
|
|
||||||
args[0] = type;
|
|
||||||
if (parameters.Length > 1)
|
|
||||||
args[1] = ConvertPrimaryArgument(parameters[1].ParameterType, valueOrBytes);
|
|
||||||
|
|
||||||
for (var i = 2; i < parameters.Length; i++)
|
|
||||||
{
|
|
||||||
args[i] = ResolveAdditionalArgument(parameters[i], serializerOptions);
|
|
||||||
}
|
|
||||||
|
|
||||||
return args;
|
|
||||||
}
|
|
||||||
|
|
||||||
private static object? ResolveAdditionalArgument(ParameterInfo parameter, object? serializerOptions)
|
|
||||||
{
|
|
||||||
if (serializerOptions != null && parameter.ParameterType.IsInstanceOfType(serializerOptions))
|
|
||||||
return serializerOptions;
|
|
||||||
|
|
||||||
return parameter.HasDefaultValue
|
|
||||||
? parameter.DefaultValue
|
|
||||||
: GetDefault(parameter.ParameterType);
|
|
||||||
}
|
|
||||||
|
|
||||||
private static object ConvertPrimaryArgument(Type parameterType, object value)
|
|
||||||
{
|
|
||||||
if (parameterType == typeof(ReadOnlyMemory<byte>) && value is byte[] bytes)
|
|
||||||
return new ReadOnlyMemory<byte>(bytes);
|
|
||||||
|
|
||||||
return value;
|
|
||||||
}
|
|
||||||
|
|
||||||
private static object? ResolveSerializerOptions(Assembly serializerAssembly)
|
|
||||||
{
|
|
||||||
var contractlessResolverType = serializerAssembly.GetType("MessagePack.Resolvers.ContractlessStandardResolver");
|
|
||||||
if (contractlessResolverType != null)
|
|
||||||
{
|
|
||||||
var optionsField = contractlessResolverType.GetField("Options",
|
|
||||||
BindingFlags.Public | BindingFlags.NonPublic | BindingFlags.Static);
|
|
||||||
var options = optionsField?.GetValue(null);
|
|
||||||
if (options != null)
|
|
||||||
return options;
|
|
||||||
|
|
||||||
var instanceField = contractlessResolverType.GetField("Instance",
|
|
||||||
BindingFlags.Public | BindingFlags.NonPublic | BindingFlags.Static);
|
|
||||||
var instance = instanceField?.GetValue(null);
|
|
||||||
if (instance != null)
|
|
||||||
{
|
|
||||||
var optionsType = serializerAssembly.GetType("MessagePack.MessagePackSerializerOptions");
|
|
||||||
var standardProperty = optionsType?.GetProperty("Standard", BindingFlags.Public | BindingFlags.Static);
|
|
||||||
var standardOptions = standardProperty?.GetValue(null);
|
|
||||||
var withResolverMethod = optionsType?.GetMethod("WithResolver", BindingFlags.Public | BindingFlags.Instance);
|
|
||||||
var resolvedOptions = withResolverMethod?.Invoke(standardOptions, new[] { instance });
|
|
||||||
if (resolvedOptions != null)
|
|
||||||
return resolvedOptions;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
private static object? GetDefault(Type type)
|
|
||||||
{
|
|
||||||
return type.IsValueType ? Activator.CreateInstance(type) : null;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -0,0 +1,105 @@
|
|||||||
|
#nullable enable
|
||||||
|
using System;
|
||||||
|
using System.Buffers.Binary;
|
||||||
|
using System.IO;
|
||||||
|
using System.Text;
|
||||||
|
|
||||||
|
namespace ShrinkNetwork
|
||||||
|
{
|
||||||
|
public sealed class ShrinkProtocolException : IOException
|
||||||
|
{
|
||||||
|
public ShrinkProtocolException(string message) : base(message) { }
|
||||||
|
public ShrinkProtocolException(string message, Exception inner) : base(message, inner) { }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>V2 little-endian framing, independent of the payload serializer. Decode borrows input memory.</summary>
|
||||||
|
public static class ShrinkPacketCodec
|
||||||
|
{
|
||||||
|
private const uint Magic = 0x324B4853; // SHK2
|
||||||
|
public const int HeaderSize = 33;
|
||||||
|
public const int MaximumPacketBytes = 16 * 1024 * 1024;
|
||||||
|
private static readonly UTF8Encoding Utf8 = new(false, true);
|
||||||
|
|
||||||
|
public static ShrinkBufferWriter Encode(ShrinkNetworkPacket packet, IShrinkNetworkSerializer? serializer = null, object? message = null)
|
||||||
|
{
|
||||||
|
var route = packet.Route ?? string.Empty;
|
||||||
|
var token = packet.SessionToken ?? string.Empty;
|
||||||
|
var routeBytes = Utf8.GetByteCount(route);
|
||||||
|
var tokenBytes = Utf8.GetByteCount(token);
|
||||||
|
if (routeBytes > ushort.MaxValue || tokenBytes > ushort.MaxValue) throw new ShrinkProtocolException("Route or session token exceeds framing limit.");
|
||||||
|
if (packet.ProtocolVersion != ShrinkNetworkProtocol.CurrentProtocolVersion || packet.SchemaVersion is < 0 or > ushort.MaxValue || packet.Kind < 0 || packet.Kind > ShrinkNetworkPacketKind.Response)
|
||||||
|
throw new ShrinkProtocolException("Invalid packet header.");
|
||||||
|
var writer = new ShrinkBufferWriter(HeaderSize + routeBytes + tokenBytes);
|
||||||
|
try
|
||||||
|
{
|
||||||
|
var header = writer.GetSpan(HeaderSize + routeBytes + tokenBytes);
|
||||||
|
BinaryPrimitives.WriteUInt32LittleEndian(header, Magic);
|
||||||
|
BinaryPrimitives.WriteUInt16LittleEndian(header.Slice(4), (ushort)packet.ProtocolVersion);
|
||||||
|
BinaryPrimitives.WriteUInt16LittleEndian(header.Slice(6), (ushort)packet.SchemaVersion);
|
||||||
|
BinaryPrimitives.WriteInt32LittleEndian(header.Slice(8), packet.Opcode);
|
||||||
|
BinaryPrimitives.WriteInt32LittleEndian(header.Slice(12), packet.RequestToken.Value);
|
||||||
|
header[16] = (byte)packet.Kind;
|
||||||
|
BinaryPrimitives.WriteInt64LittleEndian(header.Slice(17), packet.SessionTokenExpiresAtUnixTimeSeconds);
|
||||||
|
BinaryPrimitives.WriteUInt16LittleEndian(header.Slice(25), (ushort)routeBytes);
|
||||||
|
BinaryPrimitives.WriteUInt16LittleEndian(header.Slice(27), (ushort)tokenBytes);
|
||||||
|
Utf8.GetBytes(route.AsSpan(), header.Slice(HeaderSize, routeBytes));
|
||||||
|
Utf8.GetBytes(token.AsSpan(), header.Slice(HeaderSize + routeBytes, tokenBytes));
|
||||||
|
writer.Advance(HeaderSize + routeBytes + tokenBytes);
|
||||||
|
var payloadStart = writer.WrittenCount;
|
||||||
|
if (message != null)
|
||||||
|
{
|
||||||
|
if (serializer is IShrinkNetworkBufferSerializer buffered) buffered.Serialize(writer, message);
|
||||||
|
else
|
||||||
|
{
|
||||||
|
var bytes = (serializer ?? throw new ArgumentNullException(nameof(serializer))).Serialize(message);
|
||||||
|
bytes.CopyTo(writer.GetSpan(bytes.Length)); writer.Advance(bytes.Length);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
packet.Payload.Span.CopyTo(writer.GetSpan(packet.Payload.Length));
|
||||||
|
writer.Advance(packet.Payload.Length);
|
||||||
|
}
|
||||||
|
if (writer.WrittenCount > MaximumPacketBytes) throw new ShrinkProtocolException("Packet exceeds framing limit.");
|
||||||
|
BinaryPrimitives.WriteInt32LittleEndian(writer.WrittenSpan.Slice(29), writer.WrittenCount - payloadStart);
|
||||||
|
return writer;
|
||||||
|
}
|
||||||
|
catch { writer.Dispose(); throw; }
|
||||||
|
}
|
||||||
|
|
||||||
|
internal static bool IsResponse(ReadOnlyMemory<byte> memory) => memory.Length >= HeaderSize &&
|
||||||
|
BinaryPrimitives.ReadUInt32LittleEndian(memory.Span) == Magic && memory.Span[16] == (byte)ShrinkNetworkPacketKind.Response;
|
||||||
|
|
||||||
|
public static ShrinkNetworkPacket Decode(ReadOnlyMemory<byte> memory)
|
||||||
|
{
|
||||||
|
var span = memory.Span;
|
||||||
|
if (span.Length < HeaderSize || span.Length > MaximumPacketBytes || BinaryPrimitives.ReadUInt32LittleEndian(span) != Magic)
|
||||||
|
throw new ShrinkProtocolException("Expected ShrinkNetwork protocol v2 binary envelope. Legacy JSON/MessagePack envelopes are unsupported.");
|
||||||
|
var protocol = BinaryPrimitives.ReadUInt16LittleEndian(span.Slice(4));
|
||||||
|
if (protocol != ShrinkNetworkProtocol.CurrentProtocolVersion) throw new ShrinkProtocolException("Unsupported protocol version " + protocol);
|
||||||
|
var routeLength = BinaryPrimitives.ReadUInt16LittleEndian(span.Slice(25));
|
||||||
|
var tokenLength = BinaryPrimitives.ReadUInt16LittleEndian(span.Slice(27));
|
||||||
|
var payloadLength = BinaryPrimitives.ReadInt32LittleEndian(span.Slice(29));
|
||||||
|
var start = HeaderSize + routeLength + tokenLength;
|
||||||
|
if (payloadLength < 0 || start > span.Length || payloadLength != span.Length - start || span[16] > (byte)ShrinkNetworkPacketKind.Response)
|
||||||
|
throw new ShrinkProtocolException("Invalid packet lengths or kind.");
|
||||||
|
string route, token;
|
||||||
|
try
|
||||||
|
{
|
||||||
|
route = Utf8.GetString(span.Slice(HeaderSize, routeLength));
|
||||||
|
token = Utf8.GetString(span.Slice(HeaderSize + routeLength, tokenLength));
|
||||||
|
}
|
||||||
|
catch (DecoderFallbackException exception)
|
||||||
|
{
|
||||||
|
throw new ShrinkProtocolException("Invalid UTF-8 in packet header.", exception);
|
||||||
|
}
|
||||||
|
return new ShrinkNetworkPacket {
|
||||||
|
ProtocolVersion = protocol, SchemaVersion = BinaryPrimitives.ReadUInt16LittleEndian(span.Slice(6)),
|
||||||
|
Opcode = BinaryPrimitives.ReadInt32LittleEndian(span.Slice(8)), RequestToken = new ShrinkRequestToken(BinaryPrimitives.ReadInt32LittleEndian(span.Slice(12))),
|
||||||
|
Kind = (ShrinkNetworkPacketKind)span[16], SessionTokenExpiresAtUnixTimeSeconds = BinaryPrimitives.ReadInt64LittleEndian(span.Slice(17)),
|
||||||
|
Route = route, SessionToken = token,
|
||||||
|
Payload = memory.Slice(start, payloadLength)
|
||||||
|
};
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: 8a01225f12864f142b0d0c6395ef381b
|
||||||
|
MonoImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
serializedVersion: 2
|
||||||
|
defaultReferences: []
|
||||||
|
executionOrder: 0
|
||||||
|
icon: {instanceID: 0}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -0,0 +1,49 @@
|
|||||||
|
#nullable enable
|
||||||
|
using System;
|
||||||
|
using System.Buffers;
|
||||||
|
using System.Collections.Generic;
|
||||||
|
|
||||||
|
namespace ShrinkNetwork
|
||||||
|
{
|
||||||
|
public interface IShrinkMessageCodec<T>
|
||||||
|
{
|
||||||
|
void Write(IBufferWriter<byte> writer, T value);
|
||||||
|
T Read(ReadOnlyMemory<byte> payload);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>Explicit codecs work with AOT and external modules without runtime generic reflection.</summary>
|
||||||
|
public class ShrinkRegisteredNetworkSerializer : IShrinkNetworkBufferSerializer
|
||||||
|
{
|
||||||
|
private interface ICodec
|
||||||
|
{
|
||||||
|
void Write(IBufferWriter<byte> writer, object value);
|
||||||
|
object Read(ReadOnlyMemory<byte> payload);
|
||||||
|
}
|
||||||
|
private sealed class Codec<T> : ICodec
|
||||||
|
{
|
||||||
|
private readonly IShrinkMessageCodec<T> _codec;
|
||||||
|
public Codec(IShrinkMessageCodec<T> codec) => _codec = codec;
|
||||||
|
public void Write(IBufferWriter<byte> writer, object value) => _codec.Write(writer, (T)value);
|
||||||
|
public object Read(ReadOnlyMemory<byte> payload) => _codec.Read(payload)!;
|
||||||
|
}
|
||||||
|
private readonly Dictionary<Type, ICodec> _codecs = new();
|
||||||
|
public void Register<T>(IShrinkMessageCodec<T> codec)
|
||||||
|
{
|
||||||
|
if (codec == null) throw new ArgumentNullException(nameof(codec));
|
||||||
|
_codecs.Add(typeof(T), new Codec<T>(codec));
|
||||||
|
}
|
||||||
|
public bool Unregister<T>() => _codecs.Remove(typeof(T));
|
||||||
|
private ICodec Resolve(Type type) => _codecs.TryGetValue(type, out var codec) ? codec :
|
||||||
|
throw new InvalidOperationException($"SHRINK-NET-CODEC: No codec registered for {type.FullName}. Register a generated formatter before binding the transport.");
|
||||||
|
public void Serialize(IBufferWriter<byte> writer, object value) => Resolve(value.GetType()).Write(writer, value);
|
||||||
|
public object Deserialize(ReadOnlyMemory<byte> payload, Type type) => Resolve(type).Read(payload);
|
||||||
|
public byte[] Serialize(object value)
|
||||||
|
{
|
||||||
|
using var writer = new ShrinkBufferWriter();
|
||||||
|
Serialize(writer, value);
|
||||||
|
return writer.WrittenMemory.ToArray();
|
||||||
|
}
|
||||||
|
public object Deserialize(byte[] payload, Type type) => Deserialize((ReadOnlyMemory<byte>)payload, type);
|
||||||
|
public T Deserialize<T>(byte[] payload) => (T)Deserialize(payload, typeof(T));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
fileFormatVersion: 2
|
||||||
|
guid: b0225bad560778f43827ddefb85fd61c
|
||||||
|
MonoImporter:
|
||||||
|
externalObjects: {}
|
||||||
|
serializedVersion: 2
|
||||||
|
defaultReferences: []
|
||||||
|
executionOrder: 0
|
||||||
|
icon: {instanceID: 0}
|
||||||
|
userData:
|
||||||
|
assetBundleName:
|
||||||
|
assetBundleVariant:
|
||||||
@@ -1,4 +1,6 @@
|
|||||||
using Cysharp.Threading.Tasks;
|
using Cysharp.Threading.Tasks;
|
||||||
|
using System;
|
||||||
|
using System.Threading;
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
@@ -6,4 +8,10 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
UniTask SendAsync(long sessionId, byte[] packetData);
|
UniTask SendAsync(long sessionId, byte[] packetData);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>Borrowed memory remains owned by caller until SendAsync completes, faults or cancels.</summary>
|
||||||
|
public interface IShrinkNetworkMemoryTransport : IShrinkNetworkAsyncTransport
|
||||||
|
{
|
||||||
|
UniTask SendAsync(long sessionId, ReadOnlyMemory<byte> packetData, CancellationToken cancellationToken = default);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,3 +1,5 @@
|
|||||||
|
#nullable enable
|
||||||
|
|
||||||
using System;
|
using System;
|
||||||
using System.Net;
|
using System.Net;
|
||||||
using System.Net.Sockets;
|
using System.Net.Sockets;
|
||||||
@@ -16,14 +18,14 @@ namespace ShrinkNetwork
|
|||||||
private readonly ShrinkKcpTransportOptions _options;
|
private readonly ShrinkKcpTransportOptions _options;
|
||||||
private readonly object _syncRoot = new();
|
private readonly object _syncRoot = new();
|
||||||
|
|
||||||
private CancellationTokenSource _cts;
|
private CancellationTokenSource? _cts;
|
||||||
private UdpClient _udpClient;
|
private UdpClient? _udpClient;
|
||||||
private ShrinkKcpPeer _peer;
|
private ShrinkKcpPeer? _peer;
|
||||||
private long _handshakeNonce;
|
private long _handshakeNonce;
|
||||||
private int _started;
|
private int _started;
|
||||||
private int _connected;
|
private int _connected;
|
||||||
|
|
||||||
public ShrinkKcpClientTransport(string host, int port, ShrinkKcpTransportOptions options = null, long sessionId = 1)
|
public ShrinkKcpClientTransport(string host, int port, ShrinkKcpTransportOptions? options = null, long sessionId = 1)
|
||||||
{
|
{
|
||||||
if (string.IsNullOrWhiteSpace(host))
|
if (string.IsNullOrWhiteSpace(host))
|
||||||
throw new ArgumentException("Host cannot be empty.", nameof(host));
|
throw new ArgumentException("Host cannot be empty.", nameof(host));
|
||||||
@@ -39,7 +41,7 @@ namespace ShrinkNetwork
|
|||||||
|
|
||||||
public bool IsStarted => Volatile.Read(ref _started) == 1;
|
public bool IsStarted => Volatile.Read(ref _started) == 1;
|
||||||
|
|
||||||
public event Action<ShrinkNetworkTransportEvent> OnEvent;
|
public event Action<ShrinkNetworkTransportEvent>? OnEvent;
|
||||||
|
|
||||||
public void Start()
|
public void Start()
|
||||||
{
|
{
|
||||||
@@ -52,7 +54,7 @@ namespace ShrinkNetwork
|
|||||||
_udpClient.Connect(_host, _port);
|
_udpClient.Connect(_host, _port);
|
||||||
|
|
||||||
_handshakeNonce = CreateHandshakeNonce();
|
_handshakeNonce = CreateHandshakeNonce();
|
||||||
ReceiveLoopAsync(_cts.Token).Forget();
|
ReceiveLoopAsync(_udpClient, _cts.Token).Forget();
|
||||||
HandshakeLoopAsync(_cts.Token).Forget();
|
HandshakeLoopAsync(_cts.Token).Forget();
|
||||||
UpdateLoopAsync(_cts.Token).Forget();
|
UpdateLoopAsync(_cts.Token).Forget();
|
||||||
}
|
}
|
||||||
@@ -145,7 +147,7 @@ namespace ShrinkNetwork
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private async UniTaskVoid ReceiveLoopAsync(CancellationToken cancellationToken)
|
private async UniTaskVoid ReceiveLoopAsync(UdpClient udpClient, CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
@@ -154,7 +156,7 @@ namespace ShrinkNetwork
|
|||||||
UdpReceiveResult result;
|
UdpReceiveResult result;
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
result = await _udpClient.ReceiveAsync();
|
result = await udpClient.ReceiveAsync();
|
||||||
}
|
}
|
||||||
catch (ObjectDisposedException)
|
catch (ObjectDisposedException)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -37,7 +37,6 @@ namespace ShrinkNetwork
|
|||||||
private CancellationTokenSource? _cts;
|
private CancellationTokenSource? _cts;
|
||||||
private UdpClient? _udpClient;
|
private UdpClient? _udpClient;
|
||||||
private long _sessionIdGenerator;
|
private long _sessionIdGenerator;
|
||||||
private int _conversationIdGenerator;
|
|
||||||
|
|
||||||
public ShrinkKcpServerTransport(IPAddress listeningAddress, int port, ShrinkKcpTransportOptions? options = null)
|
public ShrinkKcpServerTransport(IPAddress listeningAddress, int port, ShrinkKcpTransportOptions? options = null)
|
||||||
{
|
{
|
||||||
@@ -164,9 +163,9 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
while (!cancellationToken.IsCancellationRequested)
|
while (!cancellationToken.IsCancellationRequested)
|
||||||
{
|
{
|
||||||
var sessions = _sessions.Values.ToArray();
|
foreach (var pair in _sessions)
|
||||||
foreach (var session in sessions)
|
|
||||||
{
|
{
|
||||||
|
var session = pair.Value;
|
||||||
var shouldDisconnect = false;
|
var shouldDisconnect = false;
|
||||||
lock (session.SyncRoot)
|
lock (session.SyncRoot)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -1,11 +1,12 @@
|
|||||||
#nullable enable
|
#nullable enable
|
||||||
using System;
|
using System;
|
||||||
using System.Collections.Generic;
|
using System.Collections.Generic;
|
||||||
|
using System.Threading;
|
||||||
using Cysharp.Threading.Tasks;
|
using Cysharp.Threading.Tasks;
|
||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
public sealed class ShrinkLoopbackTransport : IShrinkNetworkAsyncTransport, IShrinkNetworkSessionControlTransport
|
public sealed class ShrinkLoopbackTransport : IShrinkNetworkMemoryTransport, IShrinkNetworkSessionControlTransport
|
||||||
{
|
{
|
||||||
private readonly HashSet<long> _openedSessions = new();
|
private readonly HashSet<long> _openedSessions = new();
|
||||||
private ShrinkLoopbackTransport? _peer;
|
private ShrinkLoopbackTransport? _peer;
|
||||||
@@ -73,5 +74,12 @@ namespace ShrinkNetwork
|
|||||||
_peer.OnEvent?.Invoke(ShrinkNetworkTransportEvent.Packet(sessionId, packetData));
|
_peer.OnEvent?.Invoke(ShrinkNetworkTransportEvent.Packet(sessionId, packetData));
|
||||||
return UniTask.CompletedTask;
|
return UniTask.CompletedTask;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public UniTask SendAsync(long sessionId, ReadOnlyMemory<byte> packetData, CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
cancellationToken.ThrowIfCancellationRequested();
|
||||||
|
// Receivers may enqueue the event beyond this call; transfer a dedicated copy.
|
||||||
|
return SendAsync(sessionId, packetData.ToArray());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
#nullable enable
|
#nullable enable
|
||||||
using System;
|
using System;
|
||||||
using System.Buffers.Binary;
|
using System.Buffers.Binary;
|
||||||
|
using System.Buffers;
|
||||||
using System.IO;
|
using System.IO;
|
||||||
using System.Net.Sockets;
|
using System.Net.Sockets;
|
||||||
using System.Net.Security;
|
using System.Net.Security;
|
||||||
@@ -10,7 +11,7 @@ using Cysharp.Threading.Tasks;
|
|||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
public sealed class ShrinkTcpClientTransport : IShrinkNetworkAsyncTransport
|
public sealed class ShrinkTcpClientTransport : IShrinkNetworkMemoryTransport
|
||||||
{
|
{
|
||||||
private readonly string _host;
|
private readonly string _host;
|
||||||
private readonly int _port;
|
private readonly int _port;
|
||||||
@@ -86,7 +87,9 @@ namespace ShrinkNetwork
|
|||||||
SendAsync(sessionId, packetData).Forget();
|
SendAsync(sessionId, packetData).Forget();
|
||||||
}
|
}
|
||||||
|
|
||||||
public UniTask SendAsync(long sessionId, byte[] packetData)
|
public UniTask SendAsync(long sessionId, byte[] packetData) => SendAsync(sessionId, (ReadOnlyMemory<byte>)(packetData ?? Array.Empty<byte>()));
|
||||||
|
|
||||||
|
public UniTask SendAsync(long sessionId, ReadOnlyMemory<byte> packetData, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
if (!IsStarted)
|
if (!IsStarted)
|
||||||
throw new InvalidOperationException("Transport is not started.");
|
throw new InvalidOperationException("Transport is not started.");
|
||||||
@@ -94,11 +97,11 @@ namespace ShrinkNetwork
|
|||||||
throw new InvalidOperationException($"Unsupported session id {sessionId}. This transport only supports {_sessionId}.");
|
throw new InvalidOperationException($"Unsupported session id {sessionId}. This transport only supports {_sessionId}.");
|
||||||
if (_client == null || !_client.Connected || _stream == null)
|
if (_client == null || !_client.Connected || _stream == null)
|
||||||
throw new InvalidOperationException("TCP client is not connected.");
|
throw new InvalidOperationException("TCP client is not connected.");
|
||||||
if ((packetData?.Length ?? 0) > _maxPacketSize)
|
if (packetData.Length > _maxPacketSize)
|
||||||
throw new InvalidOperationException($"TCP packet is too large. Size={(packetData?.Length ?? 0)}, Limit={_maxPacketSize}.");
|
throw new InvalidOperationException($"TCP packet is too large. Size={packetData.Length}, Limit={_maxPacketSize}.");
|
||||||
|
|
||||||
return SendInternalAsync(_stream, _sendLock, packetData ?? Array.Empty<byte>(),
|
return SendInternalAsync(_stream, _sendLock, packetData,
|
||||||
_cts?.Token ?? CancellationToken.None);
|
cancellationToken.CanBeCanceled ? cancellationToken : _cts?.Token ?? CancellationToken.None);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async UniTaskVoid ConnectAsync(CancellationToken cancellationToken)
|
private async UniTaskVoid ConnectAsync(CancellationToken cancellationToken)
|
||||||
@@ -209,18 +212,20 @@ namespace ShrinkNetwork
|
|||||||
ex.SocketErrorCode == SocketError.Shutdown;
|
ex.SocketErrorCode == SocketError.Shutdown;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static async UniTask SendInternalAsync(Stream stream, SemaphoreSlim sendLock, byte[] packetData,
|
private static async UniTask SendInternalAsync(Stream stream, SemaphoreSlim sendLock, ReadOnlyMemory<byte> packetData,
|
||||||
CancellationToken cancellationToken)
|
CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
await sendLock.WaitAsync(cancellationToken);
|
await sendLock.WaitAsync(cancellationToken);
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
var header = new byte[4];
|
var frame = ArrayPool<byte>.Shared.Rent(checked(4 + packetData.Length));
|
||||||
BinaryPrimitives.WriteInt32LittleEndian(header, packetData.Length);
|
try
|
||||||
await stream.WriteAsync(header, cancellationToken);
|
{
|
||||||
if (packetData.Length > 0)
|
BinaryPrimitives.WriteInt32LittleEndian(frame, packetData.Length);
|
||||||
await stream.WriteAsync(packetData, cancellationToken);
|
packetData.CopyTo(frame.AsMemory(4));
|
||||||
await stream.FlushAsync(cancellationToken);
|
await stream.WriteAsync(frame.AsMemory(0, 4 + packetData.Length), cancellationToken);
|
||||||
|
}
|
||||||
|
finally { ArrayPool<byte>.Shared.Return(frame, clearArray: true); }
|
||||||
}
|
}
|
||||||
finally
|
finally
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
#nullable enable
|
#nullable enable
|
||||||
using System;
|
using System;
|
||||||
using System.Buffers.Binary;
|
using System.Buffers.Binary;
|
||||||
|
using System.Buffers;
|
||||||
using System.Collections.Concurrent;
|
using System.Collections.Concurrent;
|
||||||
using System.IO;
|
using System.IO;
|
||||||
using System.Net;
|
using System.Net;
|
||||||
@@ -13,7 +14,7 @@ using Cysharp.Threading.Tasks;
|
|||||||
|
|
||||||
namespace ShrinkNetwork
|
namespace ShrinkNetwork
|
||||||
{
|
{
|
||||||
public sealed class ShrinkTcpServerTransport : IShrinkNetworkAsyncTransport, IShrinkNetworkSessionControlTransport
|
public sealed class ShrinkTcpServerTransport : IShrinkNetworkMemoryTransport, IShrinkNetworkSessionControlTransport
|
||||||
{
|
{
|
||||||
private readonly ConcurrentDictionary<long, TcpClient> _clients = new();
|
private readonly ConcurrentDictionary<long, TcpClient> _clients = new();
|
||||||
private readonly ConcurrentDictionary<long, SemaphoreSlim> _sendLocks = new();
|
private readonly ConcurrentDictionary<long, SemaphoreSlim> _sendLocks = new();
|
||||||
@@ -92,11 +93,8 @@ namespace ShrinkNetwork
|
|||||||
|
|
||||||
_streams.Clear();
|
_streams.Clear();
|
||||||
|
|
||||||
foreach (var pair in _sendLocks)
|
// In-flight writers still release these managed semaphores after the stream closes.
|
||||||
{
|
// No WaitHandle is allocated; let their final users release and then collect them.
|
||||||
pair.Value.Dispose();
|
|
||||||
}
|
|
||||||
|
|
||||||
_sendLocks.Clear();
|
_sendLocks.Clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -125,7 +123,9 @@ namespace ShrinkNetwork
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
public async UniTask SendAsync(long sessionId, byte[] packetData)
|
public UniTask SendAsync(long sessionId, byte[] packetData) => SendAsync(sessionId, (ReadOnlyMemory<byte>)(packetData ?? Array.Empty<byte>()));
|
||||||
|
|
||||||
|
public async UniTask SendAsync(long sessionId, ReadOnlyMemory<byte> packetData, CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
if (!_clients.TryGetValue(sessionId, out var client))
|
if (!_clients.TryGetValue(sessionId, out var client))
|
||||||
throw new InvalidOperationException($"Session {sessionId} is not connected.");
|
throw new InvalidOperationException($"Session {sessionId} is not connected.");
|
||||||
@@ -134,10 +134,10 @@ namespace ShrinkNetwork
|
|||||||
|
|
||||||
if (!_sendLocks.TryGetValue(sessionId, out var sendLock))
|
if (!_sendLocks.TryGetValue(sessionId, out var sendLock))
|
||||||
throw new InvalidOperationException($"Session {sessionId} send lock is not initialized.");
|
throw new InvalidOperationException($"Session {sessionId} send lock is not initialized.");
|
||||||
if ((packetData?.Length ?? 0) > _maxPacketSize)
|
if (packetData.Length > _maxPacketSize)
|
||||||
throw new InvalidOperationException($"TCP packet is too large. Size={(packetData?.Length ?? 0)}, Limit={_maxPacketSize}.");
|
throw new InvalidOperationException($"TCP packet is too large. Size={packetData.Length}, Limit={_maxPacketSize}.");
|
||||||
|
|
||||||
await SendInternalAsync(stream, sendLock, packetData ?? Array.Empty<byte>(), CancellationToken.None);
|
await SendInternalAsync(stream, sendLock, packetData, cancellationToken);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async Task AcceptLoopAsync(CancellationToken cancellationToken)
|
private async Task AcceptLoopAsync(CancellationToken cancellationToken)
|
||||||
@@ -244,8 +244,7 @@ namespace ShrinkNetwork
|
|||||||
{
|
{
|
||||||
if (_streams.TryRemove(sessionId, out var ownedStream))
|
if (_streams.TryRemove(sessionId, out var ownedStream))
|
||||||
ownedStream.Dispose();
|
ownedStream.Dispose();
|
||||||
if (_sendLocks.TryRemove(sessionId, out var sendLock))
|
_sendLocks.TryRemove(sessionId, out _);
|
||||||
sendLock.Dispose();
|
|
||||||
|
|
||||||
var remoteAddress = removed.Client.RemoteEndPoint?.ToString() ?? "unknown";
|
var remoteAddress = removed.Client.RemoteEndPoint?.ToString() ?? "unknown";
|
||||||
try
|
try
|
||||||
@@ -261,18 +260,20 @@ namespace ShrinkNetwork
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private static async Task SendInternalAsync(Stream stream, SemaphoreSlim sendLock, byte[] packetData,
|
private static async Task SendInternalAsync(Stream stream, SemaphoreSlim sendLock, ReadOnlyMemory<byte> packetData,
|
||||||
CancellationToken cancellationToken)
|
CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
await sendLock.WaitAsync(cancellationToken);
|
await sendLock.WaitAsync(cancellationToken);
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
var header = new byte[4];
|
var frame = ArrayPool<byte>.Shared.Rent(checked(4 + packetData.Length));
|
||||||
BinaryPrimitives.WriteInt32LittleEndian(header, packetData.Length);
|
try
|
||||||
await stream.WriteAsync(header, cancellationToken);
|
{
|
||||||
if (packetData.Length > 0)
|
BinaryPrimitives.WriteInt32LittleEndian(frame, packetData.Length);
|
||||||
await stream.WriteAsync(packetData, cancellationToken);
|
packetData.CopyTo(frame.AsMemory(4));
|
||||||
await stream.FlushAsync(cancellationToken);
|
await stream.WriteAsync(frame.AsMemory(0, 4 + packetData.Length), cancellationToken);
|
||||||
|
}
|
||||||
|
finally { ArrayPool<byte>.Shared.Return(frame, clearArray: true); }
|
||||||
}
|
}
|
||||||
finally
|
finally
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -59,7 +59,7 @@ public sealed class ShrinkNetworkPingClientExample : MonoBehaviour
|
|||||||
Disconnect();
|
Disconnect();
|
||||||
|
|
||||||
_service = new ShrinkNetworkService(
|
_service = new ShrinkNetworkService(
|
||||||
new ShrinkMessagePackNetworkSerializer(),
|
new ShrinkJsonNetworkSerializer(),
|
||||||
new ShrinkNetworkMessageRegistry(),
|
new ShrinkNetworkMessageRegistry(),
|
||||||
new ShrinkNetworkRouter());
|
new ShrinkNetworkRouter());
|
||||||
_service.AutoRegisterAttributedMessages();
|
_service.AutoRegisterAttributedMessages();
|
||||||
|
|||||||
+2
-2
@@ -1,11 +1,11 @@
|
|||||||
{
|
{
|
||||||
"name": "com.cneicy.shrink-network",
|
"name": "com.cneicy.shrink-network",
|
||||||
"version": "0.2.0",
|
"version": "0.4.2",
|
||||||
"displayName": "ShrinkNetwork",
|
"displayName": "ShrinkNetwork",
|
||||||
"description": "面向 Unity 的轻量网络框架,提供会话、消息注册、RPC、权限控制,以及 TCP/KCP 传输抽象。",
|
"description": "面向 Unity 的轻量网络框架,提供会话、消息注册、RPC、权限控制,以及 TCP/KCP 传输抽象。",
|
||||||
"unity": "2022.3",
|
"unity": "2022.3",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"com.cneicy.shrink-shared-codegen": "0.1.0",
|
"com.cneicy.shrink-shared-codegen": "0.2.1",
|
||||||
"com.cysharp.unitask": "2.5.10",
|
"com.cysharp.unitask": "2.5.10",
|
||||||
"com.unity.nuget.newtonsoft-json": "3.2.2"
|
"com.unity.nuget.newtonsoft-json": "3.2.2"
|
||||||
},
|
},
|
||||||
|
|||||||
Reference in New Issue
Block a user