Attempt combobulation
All checks were successful
continuous-integration/drone/push Build is passing

This commit is contained in:
2026-03-01 01:19:27 -05:00
parent 7bec80a2fe
commit e83c9db866
10 changed files with 3385 additions and 85 deletions

View File

@@ -32,6 +32,15 @@ steps:
- master
- staging
- name: build-combobulator
image: clux/muslrust:1.91.0-stable
commands:
- make build-combobulator
when:
branch:
- master
- staging
- name: build-frontend
image: oven/bun:1.3.3
commands:
@@ -112,6 +121,29 @@ steps:
event:
- push
- name: image-combobulator
image: plugins/docker
settings:
registry: registry.itzana.me
repo: registry.itzana.me/strafesnet/maptest-combobulator
tags:
- ${DRONE_BRANCH}-${DRONE_BUILD_NUMBER}
- ${DRONE_BRANCH}
username:
from_secret: REGISTRY_USER
password:
from_secret: REGISTRY_PASS
dockerfile: combobulator/Containerfile
context: .
depends_on:
- build-combobulator
when:
branch:
- master
- staging
event:
- push
- name: deploy
image: argoproj/argocd:latest
commands:
@@ -119,6 +151,7 @@ steps:
- argocd app --grpc-web set ${DRONE_BRANCH}-maps-service --kustomize-image registry.itzana.me/strafesnet/maptest-api:${DRONE_BRANCH}-${DRONE_BUILD_NUMBER}
- argocd app --grpc-web set ${DRONE_BRANCH}-maps-service --kustomize-image registry.itzana.me/strafesnet/maptest-frontend:${DRONE_BRANCH}-${DRONE_BUILD_NUMBER}
- argocd app --grpc-web set ${DRONE_BRANCH}-maps-service --kustomize-image registry.itzana.me/strafesnet/maptest-validator:${DRONE_BRANCH}-${DRONE_BUILD_NUMBER}
- argocd app --grpc-web set ${DRONE_BRANCH}-maps-service --kustomize-image registry.itzana.me/strafesnet/maptest-combobulator:${DRONE_BRANCH}-${DRONE_BUILD_NUMBER}
environment:
USERNAME:
from_secret: ARGO_USER
@@ -128,6 +161,7 @@ steps:
- image-backend
- image-frontend
- image-validator
- image-combobulator
when:
branch:
- master
@@ -143,12 +177,13 @@ steps:
depends_on:
- build-backend
- build-validator
- build-combobulator
- build-frontend
when:
event:
- pull_request
---
kind: signature
hmac: 6de9d4b91f14b30561856daf275d1fd523e1ce7a5a3651b660f0d8907b4692fb
hmac: 2d2a3b50b5864bd79efacf31f71b5a409a1782f6dbfb4669a418f577cc5517bd
...

2961
Cargo.lock generated

File diff suppressed because it is too large Load Diff

View File

@@ -1,5 +1,6 @@
[workspace]
members = [
"combobulator",
"validation",
"submissions-api-rs",
]

View File

@@ -9,12 +9,15 @@ build-backend:
build-validator:
cargo build --release --target x86_64-unknown-linux-musl --bin maps-validation
build-combobulator:
cargo build --release --target x86_64-unknown-linux-musl --bin maps-combobulator
build-frontend:
rm -rf web/build
cd web && bun install --frozen-lockfile
cd web && bun run build
build: build-backend build-validator build-frontend
build: build-backend build-validator build-combobulator build-frontend
# image
image-backend:
@@ -23,6 +26,9 @@ image-backend:
image-validator:
docker build . -f validation/Containerfile -t maptest-validator
image-combobulator:
docker build . -f combobulator/Containerfile -t maptest-combobulator
image-frontend:
docker build web -f web/Containerfile -t maptest-frontend
@@ -33,9 +39,12 @@ docker-backend:
docker-validator:
make build-validator
make image-validator
docker-combobulator:
make build-combobulator
make image-combobulator
docker-frontend:
make image-frontend
docker: docker-backend docker-validator docker-frontend
docker: docker-backend docker-validator docker-combobulator docker-frontend
.PHONY: clean build-backend build-validator build-frontend build image-backend image-validator image-frontend docker-backend docker-validator docker-frontend docker
.PHONY: clean build-backend build-validator build-combobulator build-frontend build image-backend image-validator image-combobulator image-frontend docker-backend docker-validator docker-combobulator docker-frontend docker

15
combobulator/Cargo.toml Normal file
View File

@@ -0,0 +1,15 @@
[package]
name = "maps-combobulator"
version = "0.1.0"
edition = "2024"
[dependencies]
async-nats = "0.45.0"
aws-config = { version = "1", features = ["behavior-version-latest"] }
aws-sdk-s3 = "1"
map-tool = { version = "2.0.0", registry = "strafesnet" }
rbx_asset = { version = "0.5.0", features = ["gzip", "rustls-tls"], default-features = false, registry = "strafesnet" }
serde = { version = "1.0.215", features = ["derive"] }
serde_json = "1.0.133"
tokio = { version = "1.41.1", features = ["macros", "rt-multi-thread", "signal"] }
tokio-stream = "0.1"

View File

@@ -0,0 +1,3 @@
FROM alpine:3.21 AS runtime
COPY /target/x86_64-unknown-linux-musl/release/maps-combobulator /
ENTRYPOINT ["/maps-combobulator"]

169
combobulator/src/main.rs Normal file
View File

@@ -0,0 +1,169 @@
use tokio_stream::StreamExt;
mod nats_types;
mod process;
mod s3;
const SUBJECT_MAPFIX_RELEASE:&str="maptest.mapfixes.release";
const SUBJECT_SUBMISSION_BATCHRELEASE:&str="maptest.submissions.batchrelease";
const SUBJECT_SUBMISSION_RELEASE:&str="maptest.combobulator.submissions.release";
#[derive(Debug)]
pub enum StartupError{
NatsConnect(async_nats::ConnectError),
NatsGetStream(async_nats::jetstream::context::GetStreamError),
NatsConsumer(async_nats::jetstream::stream::ConsumerError),
NatsConsumerUpdate(async_nats::jetstream::stream::ConsumerUpdateError),
NatsStream(async_nats::jetstream::consumer::StreamError),
}
impl std::fmt::Display for StartupError{
fn fmt(&self,f:&mut std::fmt::Formatter<'_>)->std::fmt::Result{
write!(f,"{self:?}")
}
}
impl std::error::Error for StartupError{}
#[expect(dead_code)]
#[derive(Debug)]
enum HandleMessageError{
Json(serde_json::Error),
UnknownSubject(String),
Process(process::Error),
Ack(async_nats::Error),
Publish(async_nats::jetstream::context::PublishError),
}
impl std::fmt::Display for HandleMessageError{
fn fmt(&self,f:&mut std::fmt::Formatter<'_>)->std::fmt::Result{
write!(f,"{self:?}")
}
}
impl std::error::Error for HandleMessageError{}
fn from_slice<'a,T:serde::de::Deserialize<'a>>(slice:&'a [u8])->Result<T,HandleMessageError>{
serde_json::from_slice(slice).map_err(HandleMessageError::Json)
}
async fn handle_message(
processor:&process::Processor,
jetstream:&async_nats::jetstream::Context,
message:async_nats::jetstream::Message,
)->Result<(),HandleMessageError>{
match message.subject.as_str(){
SUBJECT_MAPFIX_RELEASE=>{
let request:nats_types::ReleaseMapfixRequest=from_slice(&message.payload)?;
processor.handle_mapfix_release(request).await.map_err(HandleMessageError::Process)?;
message.ack().await.map_err(HandleMessageError::Ack)?;
},
SUBJECT_SUBMISSION_BATCHRELEASE=>{
// split batch into individual messages and republish
let batch:nats_types::ReleaseSubmissionsBatchRequest=from_slice(&message.payload)?;
println!("[combobulator] Splitting batch release (operation {}, {} submissions)",
batch.OperationID,batch.Submissions.len());
for submission in batch.Submissions{
let payload=serde_json::to_vec(&submission).map_err(HandleMessageError::Json)?;
jetstream.publish(SUBJECT_SUBMISSION_RELEASE,payload.into())
.await.map_err(HandleMessageError::Publish)?;
println!("[combobulator] Published individual release for submission {}",submission.SubmissionID);
}
// ack the batch now that all individual messages are queued
message.ack().await.map_err(HandleMessageError::Ack)?;
},
SUBJECT_SUBMISSION_RELEASE=>{
let request:nats_types::ReleaseSubmissionRequest=from_slice(&message.payload)?;
processor.handle_submission_release(request).await.map_err(HandleMessageError::Process)?;
message.ack().await.map_err(HandleMessageError::Ack)?;
},
other=>return Err(HandleMessageError::UnknownSubject(other.to_owned())),
}
println!("[combobulator] Message processed and acked");
Ok(())
}S
#[tokio::main]
async fn main()->Result<(),StartupError>{
// roblox cloud api for downloading models
let api_key=std::env::var("RBX_API_KEY").expect("RBX_API_KEY env required");
let cloud_context=rbx_asset::cloud::Context::new(rbx_asset::cloud::ApiKey::new(api_key));
// roblox cookie api for downloading assets (textures, meshes, unions)
let cookie=std::env::var("RBXCOOKIE").expect("RBXCOOKIE env required");
let cookie_context=rbx_asset::cookie::Context::new(rbx_asset::cookie::Cookie::new(cookie));
// s3
let s3_bucket=std::env::var("S3_BUCKET").expect("S3_BUCKET env required");
let s3_config=aws_config::load_defaults(aws_config::BehaviorVersion::latest()).await;
let s3_client=aws_sdk_s3::Client::new(&s3_config);
let s3_cache=s3::S3Cache::new(s3_client,s3_bucket);
let processor=process::Processor{
cloud_context,
cookie_context,
s3:s3_cache,
};
// nats
let nats_host=std::env::var("NATS_HOST").expect("NATS_HOST env required");
const STREAM_NAME:&str="maptest";
const DURABLE_NAME:&str="combobulator";
let filter_subjects=vec![
SUBJECT_MAPFIX_RELEASE.to_owned(),
SUBJECT_SUBMISSION_BATCHRELEASE.to_owned(),
SUBJECT_SUBMISSION_RELEASE.to_owned(),
];
let nats_config=async_nats::jetstream::consumer::pull::Config{
name:Some(DURABLE_NAME.to_owned()),
durable_name:Some(DURABLE_NAME.to_owned()),
filter_subjects:filter_subjects.clone(),
ack_wait:std::time::Duration::from_secs(300), // 5 minutes for processing
max_deliver:5, // retry up to 5 times
..Default::default()
};
let nasty=async_nats::connect(nats_host).await.map_err(StartupError::NatsConnect)?;
let jetstream=async_nats::jetstream::new(nasty);
let stream=jetstream.get_stream(STREAM_NAME).await.map_err(StartupError::NatsGetStream)?;
let consumer=stream.get_or_create_consumer(DURABLE_NAME,nats_config.clone()).await.map_err(StartupError::NatsConsumer)?;
// update consumer config if filter subjects changed
if consumer.cached_info().config.filter_subjects!=filter_subjects{
stream.update_consumer(nats_config).await.map_err(StartupError::NatsConsumerUpdate)?;
}
let mut messages=consumer.messages().await.map_err(StartupError::NatsStream)?;
// SIGTERM graceful shutdown
let mut sig_term=tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("Failed to create SIGTERM signal listener");
println!("[combobulator] Started, waiting for messages...");
// sequential processing loop - one message at a time
let main_loop=async{
while let Some(message_result)=messages.next().await{
match message_result{
Ok(message)=>{
match handle_message(&processor,&jetstream,message).await{
Ok(())=>println!("[combobulator] Success"),
Err(e)=>println!("[combobulator] Error: {e}"),
}
},
Err(e)=>println!("[combobulator] Message stream error: {e}"),
}
}
};
tokio::select!{
_=sig_term.recv()=>{
println!("[combobulator] Received SIGTERM, shutting down");
},
_=main_loop=>{
println!("[combobulator] Message stream ended");
},
};
Ok(())
}

View File

@@ -0,0 +1,29 @@
#[expect(nonstandard_style,dead_code)]
#[derive(serde::Deserialize)]
pub struct ReleaseMapfixRequest{
pub MapfixID:u64,
pub ModelID:u64,
pub ModelVersion:u64,
pub TargetAssetID:u64,
}
#[expect(nonstandard_style)]
#[derive(serde::Deserialize,serde::Serialize)]
pub struct ReleaseSubmissionRequest{
pub SubmissionID:u64,
pub ReleaseDate:i64,
pub ModelID:u64,
pub ModelVersion:u64,
pub UploadedAssetID:u64,
pub DisplayName:String,
pub Creator:String,
pub GameID:u32,
pub Submitter:u64,
}
#[expect(nonstandard_style)]
#[derive(serde::Deserialize)]
pub struct ReleaseSubmissionsBatchRequest{
pub Submissions:Vec<ReleaseSubmissionRequest>,
pub OperationID:u32,
}

144
combobulator/src/process.rs Normal file
View File

@@ -0,0 +1,144 @@
use crate::nats_types::ReleaseMapfixRequest;
use crate::s3::S3Cache;
#[expect(dead_code)]
#[derive(Debug)]
pub enum Error{
Download(rbx_asset::cloud::GetError),
NonFreeModel,
GetAssets(map_tool::roblox::UniqueAssetError),
DownloadAsset(map_tool::roblox::DownloadAssetError),
ConvertTexture(map_tool::roblox::ConvertTextureError),
ConvertSnf(map_tool::roblox::ConvertError),
S3Get(crate::s3::GetError),
S3Put(crate::s3::PutError),
}
impl std::fmt::Display for Error{
fn fmt(&self,f:&mut std::fmt::Formatter<'_>)->std::fmt::Result{
write!(f,"{self:?}")
}
}
impl std::error::Error for Error{}
pub struct Processor{
pub cloud_context:rbx_asset::cloud::Context,
pub cookie_context:rbx_asset::cookie::Context,
pub s3:S3Cache,
}
impl Processor{
/// Download a model version from Roblox cloud API.
async fn download_model(&self,model_id:u64,model_version:u64)->Result<Vec<u8>,Error>{
let location=self.cloud_context.get_asset_version_location(
rbx_asset::cloud::GetAssetVersionRequest{
asset_id:model_id,
version:model_version,
}
).await.map_err(Error::Download)?;
let location=location.location.ok_or(Error::NonFreeModel)?;
let maybe_gzip=self.cloud_context.get_asset(&location).await.map_err(Error::Download)?;
Ok(maybe_gzip.into_inner().to_vec())
}
/// Process a single model: extract assets, cache to S3, build SNF.
async fn process_model(&self,model_id:u64,model_version:u64)->Result<(),Error>{
println!("[combobulator] Downloading model {model_id} v{model_version}");
let rbxl_bytes=self.download_model(model_id,model_version).await?;
// extract unique assets from the file
let assets=map_tool::roblox::get_unique_assets_from_file(&rbxl_bytes)
.map_err(Error::GetAssets)?;
// process textures: download, cache, convert to DDS
for id in &assets.textures{
let asset_id=id.0;
let dds_key=S3Cache::texture_dds_key(asset_id);
// skip if DDS already cached
if self.s3.get(&dds_key).await.map_err(Error::S3Get)?.is_some(){
println!("[combobulator] Texture {asset_id} already cached, skipping");
continue;
}
// check raw cache, download if missing
let raw_key=S3Cache::texture_raw_key(asset_id);
let raw_data=match self.s3.get(&raw_key).await.map_err(Error::S3Get)?{
Some(cached)=>cached,
None=>{
println!("[combobulator] Downloading texture {asset_id}");
let data=map_tool::roblox::download_asset(&self.cookie_context,asset_id)
.await.map_err(Error::DownloadAsset)?;
self.s3.put(&raw_key,data.clone()).await.map_err(Error::S3Put)?;
data
},
};
// convert to DDS and upload
let dds=map_tool::roblox::convert_texture_to_dds(&raw_data)
.map_err(Error::ConvertTexture)?;
self.s3.put(&dds_key,dds).await.map_err(Error::S3Put)?;
println!("[combobulator] Texture {asset_id} processed");
}
// process meshes
for id in &assets.meshes{
let asset_id=id.0;
let mesh_key=S3Cache::mesh_key(asset_id);
if self.s3.get(&mesh_key).await.map_err(Error::S3Get)?.is_some(){
println!("[combobulator] Mesh {asset_id} already cached, skipping");
continue;
}
println!("[combobulator] Downloading mesh {asset_id}");
let data=map_tool::roblox::download_asset(&self.cookie_context,asset_id)
.await.map_err(Error::DownloadAsset)?;
self.s3.put(&mesh_key,data).await.map_err(Error::S3Put)?;
println!("[combobulator] Mesh {asset_id} processed");
}
// process unions
for id in &assets.unions{
let asset_id=id.0;
let union_key=S3Cache::union_key(asset_id);
if self.s3.get(&union_key).await.map_err(Error::S3Get)?.is_some(){
println!("[combobulator] Union {asset_id} already cached, skipping");
continue;
}
println!("[combobulator] Downloading union {asset_id}");
let data=map_tool::roblox::download_asset(&self.cookie_context,asset_id)
.await.map_err(Error::DownloadAsset)?;
self.s3.put(&union_key,data).await.map_err(Error::S3Put)?;
println!("[combobulator] Union {asset_id} processed");
}
// convert to SNF and upload
println!("[combobulator] Converting to SNF");
let output=map_tool::roblox::convert_to_snf(&rbxl_bytes)
.map_err(Error::ConvertSnf)?;
let snf_key=S3Cache::snf_key(model_id,model_version);
self.s3.put(&snf_key,output.snf).await.map_err(Error::S3Put)?;
println!("[combobulator] SNF uploaded to {snf_key}");
Ok(())
}
/// Handle a mapfix release message.
pub async fn handle_mapfix_release(&self,request:ReleaseMapfixRequest)->Result<(),Error>{
println!("[combobulator] Processing mapfix {} (model {} v{})",
request.MapfixID,request.ModelID,request.ModelVersion);
self.process_model(request.ModelID,request.ModelVersion).await
}
/// Handle an individual submission release message.
pub async fn handle_submission_release(&self,request:crate::nats_types::ReleaseSubmissionRequest)->Result<(),Error>{
println!("[combobulator] Processing submission {} (model {} v{})",
request.SubmissionID,request.ModelID,request.ModelVersion);
self.process_model(request.ModelID,request.ModelVersion).await
}
}

96
combobulator/src/s3.rs Normal file
View File

@@ -0,0 +1,96 @@
use aws_sdk_s3::Client;
use aws_sdk_s3::primitives::ByteStream;
#[expect(dead_code)]
#[derive(Debug)]
pub enum GetError{
Get(aws_sdk_s3::error::SdkError<aws_sdk_s3::operation::get_object::GetObjectError>),
Collect(aws_sdk_s3::primitives::ByteStreamError),
}
impl std::fmt::Display for GetError{
fn fmt(&self,f:&mut std::fmt::Formatter<'_>)->std::fmt::Result{
write!(f,"{self:?}")
}
}
impl std::error::Error for GetError{}
#[expect(dead_code)]
#[derive(Debug)]
pub enum PutError{
Put(aws_sdk_s3::error::SdkError<aws_sdk_s3::operation::put_object::PutObjectError>),
}
impl std::fmt::Display for PutError{
fn fmt(&self,f:&mut std::fmt::Formatter<'_>)->std::fmt::Result{
write!(f,"{self:?}")
}
}
impl std::error::Error for PutError{}
pub struct S3Cache{
client:Client,
bucket:String,
}
impl S3Cache{
pub fn new(client:Client,bucket:String)->Self{
Self{client,bucket}
}
/// Try to get a cached object. Returns None if the key doesn't exist.
pub async fn get(&self,key:&str)->Result<Option<Vec<u8>>,GetError>{
match self.client.get_object()
.bucket(&self.bucket)
.key(key)
.send()
.await
{
Ok(output)=>{
let bytes=output.body.collect().await.map_err(GetError::Collect)?;
Ok(Some(bytes.to_vec()))
},
Err(e)=>{
// check if it's a NoSuchKey error
if let aws_sdk_s3::error::SdkError::ServiceError(ref service_err)=e{
if service_err.err().is_no_such_key(){
return Ok(None);
}
}
Err(GetError::Get(e))
},
}
}
/// Put an object into S3.
pub async fn put(&self,key:&str,data:Vec<u8>)->Result<(),PutError>{
self.client.put_object()
.bucket(&self.bucket)
.key(key)
.body(ByteStream::from(data))
.send()
.await
.map_err(PutError::Put)?;
Ok(())
}
// S3 key helpers
pub fn texture_raw_key(asset_id:u64)->String{
format!("assets/textures/{asset_id}.raw")
}
pub fn texture_dds_key(asset_id:u64)->String{
format!("assets/textures/{asset_id}.dds")
}
pub fn mesh_key(asset_id:u64)->String{
format!("assets/meshes/{asset_id}")
}
pub fn union_key(asset_id:u64)->String{
format!("assets/unions/{asset_id}")
}
pub fn snf_key(model_id:u64,model_version:u64)->String{
format!("maps/{model_id}/v{model_version}/map.snfm")
}
}