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
9 changes: 7 additions & 2 deletions .github/workflows/ci.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,11 @@ jobs:
profile: minimal
override: true

- name: Set up Node.js (for websocket interop tests)
uses: actions/setup-node@v4
with:
node-version: "20"

- name: Create dummy local.properties
run: |
echo 'sdk.dir=/fake' >> local.properties
Expand Down Expand Up @@ -173,14 +178,14 @@ jobs:
env:
JAVA_HOME: ${{ env.JAVA_HOME_FOR_BUILD }}
with:
arguments: lib:assemble -Penv=dev --info
arguments: lib:assemble websocket:assemble -Penv=dev --info

- name: Test with Java 21 runtime
uses: gradle/gradle-build-action@749f47bda3e44aa060e82d7b3ef7e40d953bd629
env:
JAVA_HOME: ${{ env.JAVA_HOME_FOR_BUILD }}
with:
arguments: lib:test -Penv=dev --info
arguments: lib:test websocket:test -Penv=dev --info

- name: Test with Java 8 runtime (backward compatibility)
uses: gradle/gradle-build-action@749f47bda3e44aa060e82d7b3ef7e40d953bd629
Expand Down
1 change: 1 addition & 0 deletions settings.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -20,5 +20,6 @@ pluginManagement {

rootProject.name = "automerge-java"
include("lib")
include("websocket")
include("android")
include("android-test-app")
97 changes: 97 additions & 0 deletions websocket/build.gradle.kts
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
plugins {
`java-library`
id("org.danilopianini.publish-on-central")
id("com.diffplug.spotless")
}

java {
withJavadocJar()
withSourcesJar()
toolchain {
languageVersion.set(JavaLanguageVersion.of(21))
}
}

repositories {
mavenCentral()
}

dependencies {
api(project(":lib"))
implementation("org.java-websocket:Java-WebSocket:1.5.7")
implementation("org.slf4j:slf4j-api:2.0.9")

testImplementation(platform("org.junit:junit-bom:5.11.4"))
testImplementation("org.junit.jupiter:junit-jupiter-api")
testRuntimeOnly("org.junit.jupiter:junit-jupiter-engine")
testRuntimeOnly("org.junit.platform:junit-platform-launcher")
testImplementation("org.slf4j:slf4j-simple:2.0.9")
}

spotless {
java {
importOrder()
targetExclude("src/templates/*")
removeUnusedImports()
cleanthat()
eclipse().configFile("${project.rootDir}/spotless.eclipseformat.xml")
formatAnnotations()
}
}

publishOnCentral {
projectDescription.set("WebSocket network transport for Automerge repos")
projectLongName.set("Automerge WebSocket")
}

publishing {
publications {
withType<MavenPublication> {
artifactId = "automerge-websocket"
}
}
}

val env = providers.gradleProperty("env").getOrElse("release")
val isDev = env == "dev"

if (isDev) {
tasks.register<Exec>("compileRustForTest") {
workingDir = File("../rust")
commandLine = listOf("cargo", "build")
}

val version = (project.extra.get("libVersionSuffix") as String)

tasks.register("createVersionedLibForTest") {
dependsOn("compileRustForTest")
val debugDir = file("../rust/target/debug")
doLast {
listOf("libautomerge_jni" to "so", "libautomerge_jni" to "dylib", "automerge_jni" to "dll").forEach { (base, ext) ->
val src = debugDir.resolve("$base.$ext")
if (src.exists()) {
src.copyTo(debugDir.resolve("${base}_$version.$ext"), overwrite = true)
}
}
}
}

tasks.withType<Test> {
dependsOn("createVersionedLibForTest")
systemProperty("java.library.path", file("../rust/target/debug").absolutePath)
}
}

tasks.compileJava {
options.release = 8
}

tasks.test {
useJUnitPlatform()
// Run tests with the websocket module dir as working dir so that
// JsServerWrapper can find interop-test-server/ via a relative path.
workingDir = projectDir
testLogging {
showStandardStreams = true
}
}
2 changes: 2 additions & 0 deletions websocket/interop-test-server/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
node_modules
*.js
222 changes: 222 additions & 0 deletions websocket/interop-test-server/client.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,222 @@
import {
Repo,
isValidAutomergeUrl,
parseAutomergeUrl,
} from "@automerge/automerge-repo";
import { BrowserWebSocketClientAdapter } from "@automerge/automerge-repo-network-websocket";
import { command, run, string, positional, number, subcommands } from "cmd-ts";
import { next as A } from "@automerge/automerge";

const create = command({
name: "create",
args: {
port: positional({
type: number,
displayName: "port",
description: "The port to connect to",
}),
},
handler: ({ port }) => {
const repo = new Repo({
network: [new BrowserWebSocketClientAdapter(`ws://localhost:${port}`)],
});
const doc = repo.create<{ foo: string }>();
doc.change((d) => (d.foo = "bar"));
console.log(doc.url);
console.log(A.getHeads(doc.doc()).join(","));
},
});

const fetch = command({
name: "fetch",
args: {
port: positional({
type: number,
displayName: "port",
description: "The port to connect to",
}),
docUrl: positional({
type: string,
displayName: "docUrl",
description: "The document url to fetch",
}),
},
handler: ({ port, docUrl }) => {
const repo = new Repo({
network: [new BrowserWebSocketClientAdapter(`ws://localhost:${port}`)],
});
if (isValidAutomergeUrl(docUrl)) {
} else {
throw new Error("Invalid docUrl");
}
const doc = repo.find(docUrl);
repo.find(docUrl).then((d) => console.log(A.getHeads(d.doc()).join(",")));
},
});

const sendEphemeral = command({
name: "send-ephemeral",
args: {
port: positional({
type: number,
displayName: "port",
description: "The port to connect to",
}),
docUrl: positional({
type: string,
displayName: "docUrl",
description: "The document url to fetch",
}),
message: positional({
type: string,
displayName: "message",
description: "The message to send",
}),
},
handler: ({ port, docUrl, message }) => {
const repo = new Repo({
network: [new BrowserWebSocketClientAdapter(`ws://localhost:${port}`)],
});
if (!isValidAutomergeUrl(docUrl)) {
throw new Error("Invalid docUrl");
}
repo
.find(docUrl)
.then((doc) => {
doc.broadcast({ message });
process.exit(0);
})
.catch((e) => {
console.error(e);
process.exit(1);
});
},
});

const receiveEphemeral = command({
name: "receive-ephemeral",
args: {
port: positional({
type: number,
displayName: "port",
description: "The port to connect to",
}),
docUrl: positional({
type: string,
displayName: "docUrl",
description: "The document url to fetch",
}),
},
handler: async ({ port, docUrl }) => {
const repo = new Repo({
network: [new BrowserWebSocketClientAdapter(`ws://localhost:${port}`)],
});
if (!isValidAutomergeUrl(docUrl)) {
throw new Error("Invalid docUrl");
}
repo.find(docUrl).then((doc) => {
doc.on("ephemeral-message", ({ message }) => {
if (typeof message === "object" && "message" in message) {
console.log(message.message);
}
});
// Signal that we're ready to receive ephemeral messages
console.log("ready");
});
},
});

// This command connects to two servers: first a JS server (which has a storage ID), syncs a
// document with it, then connects to a second server (the Java server). When the second server
// becomes a "generous peer", addGenerousPeer fires and sends a `remote-heads-changed` message
// containing the stored remote heads info with a `Date.now()` timestamp (encoded as f64 by cbor-x).
const createAndRelayHeads = command({
name: "create-and-relay-heads",
args: {
jsServerPort: positional({
type: number,
displayName: "jsServerPort",
description: "The port of the JS server to sync with first",
}),
javaServerPort: positional({
type: number,
displayName: "javaServerPort",
description: "The port of the Java server to connect to second",
}),
},
handler: async ({ jsServerPort, javaServerPort }) => {
const jsAdapter = new BrowserWebSocketClientAdapter(
`ws://localhost:${jsServerPort}`,
);
const repo = new Repo({
network: [jsAdapter],
enableRemoteHeadsGossiping: true,
});

// Subscribe to a dummy storage ID so that when the Java server connects,
// addGenerousPeer sends a remote-subscription-change message too.
repo.subscribeToRemotes(["dummy-storage-id" as any]);

const doc = repo.create<{ foo: string }>();
doc.change((d) => (d.foo = "bar"));

// Wait for sync with JS server to complete. This stores remote heads info
// (with a Date.now() timestamp) in RemoteHeadsSubscriptions.#syncInfoByDocId.
await new Promise<void>((resolve) => setTimeout(resolve, 500));

// Now connect to the Java server. When the peer event fires, sharePolicy
// returns true (default), so addGenerousPeer is called. This sends a
// `remote-heads-changed` message with the stored f64 timestamp.
const javaAdapter = new BrowserWebSocketClientAdapter(
`ws://localhost:${javaServerPort}`,
);
repo.networkSubsystem.addNetworkAdapter(javaAdapter);

// Wait for the Java server connection to establish and messages to be sent.
await new Promise<void>((resolve) => setTimeout(resolve, 500));

console.log(doc.url);
console.log(A.getHeads(doc.doc()).join(","));
},
});

const subscribeAndCreate = command({
name: "subscribe-and-create",
args: {
port: positional({
type: number,
displayName: "port",
description: "The port to connect to",
}),
storageId: positional({
type: string,
displayName: "storageId",
description: "The storage ID to subscribe to remote heads for",
}),
},
handler: ({ port, storageId }) => {
const repo = new Repo({
network: [new BrowserWebSocketClientAdapter(`ws://localhost:${port}`)],
enableRemoteHeadsGossiping: true,
});
repo.subscribeToRemotes([storageId as any]);
const doc = repo.create<{ foo: string }>();
doc.change((d) => (d.foo = "bar"));
console.log(doc.url);
console.log(A.getHeads(doc.doc()).join(","));
},
});

const app = subcommands({
name: "client",
cmds: {
create,
fetch,
"send-ephemeral": sendEphemeral,
"receive-ephemeral": receiveEphemeral,
"subscribe-and-create": subscribeAndCreate,
"create-and-relay-heads": createAndRelayHeads,
},
});

run(app, process.argv.slice(2));
Loading
Loading