Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 15 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "a3s-boot"
version = "0.1.2"
version = "0.1.3"
edition = "2021"
authors = ["A3S Lab"]
license = "MIT"
Expand Down Expand Up @@ -28,6 +28,14 @@ events = ["dep:a3s-event"]
file-upload = ["dep:bytes", "dep:multer"]
health = []
http-client = ["dep:reqwest"]
ilink = [
"dep:async-trait",
"dep:base64",
"dep:rand",
"dep:reqwest",
"dep:url",
"dep:zeroize",
]
logging = []
macros = ["dep:a3s-boot-macros"]
openapi-schemas = ["dep:schemars"]
Expand All @@ -51,8 +59,10 @@ a3s-acl = { version = "0.2.1", optional = true }
a3s-boot-macros = { version = "0.1.2", path = "macros", optional = true }
a3s-event = { version = "0.3.0", default-features = false, optional = true }
a3s-lane = { version = "0.5.1", default-features = false, optional = true }
async-trait = { version = "0.1", optional = true }
async-nats = { version = "0.49.1", default-features = false, optional = true }
axum = { version = "0.8", features = ["ws"], optional = true }
base64 = { version = "0.22", optional = true }
bytes = { version = "1", optional = true }
chrono = { version = "0.4", optional = true }
cron = { version = "0.17.0", optional = true }
Expand All @@ -65,8 +75,9 @@ lapin = { version = "4.10.0", default-features = false, features = ["tokio"], op
multer = { version = "3", optional = true }
percent-encoding = "2"
prost = { version = "0.14.4", optional = true }
rand = { version = "0.8", optional = true }
redis = { version = "1.3.0", default-features = false, features = ["tokio-comp"], optional = true }
reqwest = { version = "0.12", default-features = false, features = ["rustls-tls"], optional = true }
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls", "stream"], optional = true }
rskafka = { version = "0.6.0", default-features = false, optional = true }
rumqttc = { version = "0.25.1", default-features = false, optional = true }
schemars = { version = "1", optional = true }
Expand All @@ -78,6 +89,8 @@ thiserror = "2"
tokio = { version = "1", features = ["fs", "io-util", "net", "rt", "sync", "time"], optional = true }
tonic = { version = "0.14.6", default-features = false, features = ["codegen", "transport"], optional = true }
tonic-prost = { version = "0.14.6", optional = true }
url = { version = "2", optional = true }
zeroize = { version = "1", features = ["derive"], optional = true }

[dev-dependencies]
tokio = { version = "1", features = ["macros", "rt", "sync", "time"] }
Expand Down
28 changes: 25 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,7 @@ are opt-in.
| Scheduling | `schedule` | Cron, interval, and timeout jobs |
| Observability | `logging`, `health` | Structured logging and health indicators |
| HTTP utilities | `http-client`, `compression` | Outbound HTTP and gzip responses |
| Channels | `ilink` | Tencent Weixin iLink QR login, polling, messaging, and lifecycle client |
| Content | `file-upload`, `static` | Multipart uploads and static files |
| Context | `request-context` | Task-local access to the current request |
| OpenAPI | `openapi-schemas` | `schemars`-based component schemas |
Expand All @@ -141,22 +142,22 @@ is not a claim of support for every production backend.

```toml
[dependencies]
a3s-boot = "0.1.2"
a3s-boot = "0.1.3"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
```

For a core-only build without Axum, macros, or shutdown signal handling:

```toml
[dependencies]
a3s-boot = { version = "0.1.2", default-features = false }
a3s-boot = { version = "0.1.3", default-features = false }
```

Enable only the optional modules an application uses:

```toml
[dependencies]
a3s-boot = { version = "0.1.2", features = ["auth", "security", "openapi-schemas"] }
a3s-boot = { version = "0.1.3", features = ["auth", "security", "openapi-schemas"] }
serde = { version = "1", features = ["derive"] }
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
```
Expand Down Expand Up @@ -322,6 +323,27 @@ Transport implementations share typed payload handling, scoped providers,
validation, guards, interceptors, pipes, exception filters, and client APIs.
Protocol delivery and durability semantics still depend on the selected backend.

### Weixin iLink

The optional `ilink` feature provides the native Rust protocol boundary used by
the Tencent Weixin channel. `IlinkModule` exports a typed `IlinkClient`
provider; the client owns QR login requests, authenticated headers, strict
server URL validation, update polling, text replies, typing calls, and channel
start/stop notifications.

```rust
use a3s_boot::ilink::IlinkModule;

let module = IlinkModule::weixin("A3S/0.10.1");
```

The wire defaults are compatible with Tencent `openclaw-weixin` v2.4.6:
`iLink-App-Id: bot`, `bot_type=3`, and packed client version `2.4.6`. The
product-specific `bot_agent` remains `A3S/<version>` so upstream diagnostics do
not misidentify the caller. Boot deliberately does not own browser APIs,
credential persistence, owner authorization, or agent/session commands; those
policies stay in the host application.

## Architecture

The application core is independent of its HTTP server and message broker:
Expand Down
5 changes: 5 additions & 0 deletions ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,11 @@ Implemented today:
topics plus optional gRPC unary request/reply and event calls. Transport error
envelopes round-trip through the same `BootError` HTTP exception mapping used
by HTTP routes.
- Optional native Tencent Weixin iLink support with an injectable client,
QR-login and redirect handling, authenticated update polling, text replies,
typing calls, lifecycle notifications, bounded responses, and strict URL
validation. Browser APIs, credential storage, and product authorization
remain host-application responsibilities.
- ACL-backed typed configuration modules with `ConfigModule`, named/global
provider exports, environment/default function support, and validation hooks.
- Provider-backed outbound HTTP clients with `HttpModule`, `HttpService`,
Expand Down
93 changes: 93 additions & 0 deletions src/ilink/auth.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
//! Secret handling and client-version encoding for the iLink wire protocol.

use std::fmt;

use base64::Engine as _;
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use thiserror::Error;
use zeroize::{Zeroize, ZeroizeOnDrop};

const MAX_SECRET_BYTES: usize = 64 * 1024;

#[derive(Clone, PartialEq, Eq, Zeroize, ZeroizeOnDrop)]
pub struct SecretValue(String);

impl SecretValue {
pub fn new(value: impl Into<String>) -> Result<Self, SecretValueError> {
let value = value.into();
if value.is_empty() {
return Err(SecretValueError::Empty);
}
if value.len() > MAX_SECRET_BYTES {
return Err(SecretValueError::TooLarge);
}
Ok(Self(value))
}

pub fn expose(&self) -> &str {
&self.0
}
}

impl fmt::Debug for SecretValue {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("SecretValue([REDACTED])")
}
}

impl Serialize for SecretValue {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(self.expose())
}
}

impl<'de> Deserialize<'de> for SecretValue {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let value = String::deserialize(deserializer)?;
Self::new(value).map_err(serde::de::Error::custom)
}
}

#[derive(Clone, Copy, Debug, Error, PartialEq, Eq)]
pub enum SecretValueError {
#[error("secret value is empty")]
Empty,
#[error("secret value exceeds the protocol size limit")]
TooLarge,
}

#[derive(Clone, Debug, Error, PartialEq, Eq)]
pub enum ClientVersionError {
#[error("client version must contain exactly three numeric components")]
InvalidShape,
#[error("client version component exceeds 255")]
ComponentOutOfRange,
}

pub(super) fn pack_client_version(version: &str) -> Result<u32, ClientVersionError> {
let components = version.split('.').collect::<Vec<_>>();
if components.len() != 3 || components.iter().any(|component| component.is_empty()) {
return Err(ClientVersionError::InvalidShape);
}
let mut parsed = [0u32; 3];
for (index, component) in components.into_iter().enumerate() {
let value = component
.parse::<u32>()
.map_err(|_| ClientVersionError::InvalidShape)?;
if value > u8::MAX as u32 {
return Err(ClientVersionError::ComponentOutOfRange);
}
parsed[index] = value;
}
Ok((parsed[0] << 16) | (parsed[1] << 8) | parsed[2])
}

pub(super) fn random_wechat_uin() -> String {
base64::engine::general_purpose::STANDARD.encode(rand::random::<u32>().to_string())
}
172 changes: 172 additions & 0 deletions src/ilink/client.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
//! iLink client identity, authentication, and protocol errors.

use std::fmt;
use std::time::Duration;

use reqwest::header::{HeaderMap, HeaderName, HeaderValue, AUTHORIZATION, CONTENT_TYPE};
use thiserror::Error;

use super::auth::{pack_client_version, random_wechat_uin, ClientVersionError, SecretValue};
use super::url_policy::{IlinkUrlError, ValidatedBaseUrl};

pub(super) const STALE_TOKEN_ERROR_CODE: i64 = -14;
pub(super) const DEFAULT_API_TIMEOUT: Duration = Duration::from_secs(15);
pub(super) const DEFAULT_CONFIG_TIMEOUT: Duration = Duration::from_secs(10);
pub(super) const DEFAULT_QR_POLL_TIMEOUT: Duration = Duration::from_secs(35);
pub(super) const MAX_LONG_POLL_TIMEOUT: Duration = Duration::from_secs(60);
pub(super) const MAX_RESPONSE_BYTES: usize = 1024 * 1024;

const AUTHORIZATION_TYPE: HeaderName = HeaderName::from_static("authorizationtype");
const WECHAT_UIN: HeaderName = HeaderName::from_static("x-wechat-uin");
const ILINK_APP_ID: HeaderName = HeaderName::from_static("ilink-app-id");
const ILINK_APP_CLIENT_VERSION: HeaderName = HeaderName::from_static("ilink-app-clientversion");

#[derive(Clone, Debug)]
pub(super) struct IlinkClientIdentity {
app_id: String,
packed_client_version: u32,
pub(super) bot_type: String,
pub(super) channel_version: String,
pub(super) bot_agent: String,
}

impl IlinkClientIdentity {
pub(super) fn new(
app_id: impl Into<String>,
bot_type: impl Into<String>,
client_version: &str,
bot_agent: impl Into<String>,
) -> Result<Self, IlinkError> {
let app_id = bounded_ascii(app_id.into(), "app id", 128)?;
let bot_type = bounded_ascii(bot_type.into(), "bot type", 16)?;
let channel_version = bounded_ascii(client_version.to_string(), "client version", 32)?;
let bot_agent = bounded_ascii(bot_agent.into(), "bot agent", 256)?;
let packed_client_version = pack_client_version(client_version)?;
Ok(Self {
app_id,
packed_client_version,
bot_type,
channel_version,
bot_agent,
})
}

pub(super) fn base_info(&self) -> super::types::BaseInfo {
super::types::BaseInfo {
channel_version: Some(self.channel_version.clone()),
bot_agent: Some(self.bot_agent.clone()),
}
}

pub(super) fn application_headers(&self) -> Result<HeaderMap, IlinkError> {
let mut headers = HeaderMap::new();
headers.insert(ILINK_APP_ID, safe_header_value(&self.app_id, "app id")?);
headers.insert(
ILINK_APP_CLIENT_VERSION,
safe_header_value(&self.packed_client_version.to_string(), "client version")?,
);
Ok(headers)
}

pub(super) fn post_headers(
&self,
token: Option<&SecretValue>,
) -> Result<HeaderMap, IlinkError> {
let mut headers = self.application_headers()?;
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
headers.insert(
AUTHORIZATION_TYPE,
HeaderValue::from_static("ilink_bot_token"),
);
headers.insert(
WECHAT_UIN,
safe_header_value(&random_wechat_uin(), "Weixin UIN")?,
);
if let Some(token) = token {
headers.insert(
AUTHORIZATION,
safe_header_value(&format!("Bearer {}", token.expose()), "authorization")?,
);
}
Ok(headers)
}
}

#[derive(Clone)]
pub struct IlinkAuth {
pub(super) base_url: ValidatedBaseUrl,
pub(super) bot_token: SecretValue,
}

impl IlinkAuth {
pub fn new(base_url: ValidatedBaseUrl, bot_token: SecretValue) -> Self {
Self {
base_url,
bot_token,
}
}
}

impl fmt::Debug for IlinkAuth {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("IlinkAuth")
.field("base_url", &self.base_url)
.field("bot_token", &self.bot_token)
.finish()
}
}

fn bounded_ascii(value: String, field: &'static str, max: usize) -> Result<String, IlinkError> {
if value.is_empty()
|| value.len() > max
|| !value.is_ascii()
|| value.chars().any(char::is_control)
{
return Err(IlinkError::InvalidConfiguration(field));
}
Ok(value)
}

fn safe_header_value(value: &str, field: &'static str) -> Result<HeaderValue, IlinkError> {
HeaderValue::from_str(value).map_err(|_| IlinkError::InvalidConfiguration(field))
}

pub(super) fn ensure_api_success(
operation: &'static str,
ret: Option<i64>,
errcode: Option<i64>,
) -> Result<(), IlinkError> {
let code = errcode
.filter(|code| *code != 0)
.or_else(|| ret.filter(|code| *code != 0));
match code {
None => Ok(()),
Some(STALE_TOKEN_ERROR_CODE) => Err(IlinkError::StaleCredential),
Some(code) => Err(IlinkError::Protocol { operation, code }),
}
}

#[derive(Clone, Debug, Error, PartialEq, Eq)]
pub enum IlinkError {
#[error("iLink configuration field is invalid: {0}")]
InvalidConfiguration(&'static str),
#[error("iLink URL policy rejected the request")]
Url(#[from] IlinkUrlError),
#[error("iLink client version is invalid")]
ClientVersion(#[from] ClientVersionError),
#[error("iLink request timed out")]
Timeout,
#[error("iLink transport failed")]
Transport,
#[error("iLink returned HTTP status {0}")]
HttpStatus(u16),
#[error("iLink response exceeds the size limit")]
ResponseTooLarge,
#[error("iLink returned an invalid response for {0}")]
InvalidResponse(&'static str),
#[error("iLink credential is stale")]
StaleCredential,
#[error("iLink operation {operation} failed with code {code}")]
Protocol { operation: &'static str, code: i64 },
}
Loading
Loading