* Add default workspace URL

* R1 WIP

* Improve help docs slightly

* Rework tracking state

* WIP Rework
Remaining bug: Not returning state-only files (no local file) from *getFiles()

* Create newly found files

* Finish ZIP & new tracking code

* Fix two minor bugs

* do not consider conflict if same content

* add more logs to cli writing

* progress

* progress

* iteration

* Add most basic App support

* fix folder frontend bug

* fix folder frontend bug

* init done by default

* sqlx merge

---------

Co-authored-by: Ruben Fiszel <ruben@rubenfiszel.com>
This commit is contained in:
Kai Jellinghaus
2023-02-21 10:13:21 +01:00
committed by GitHub
parent 7cee6e498d
commit 908e7e5ba2
15 changed files with 586 additions and 227 deletions
+68
View File
@@ -35,6 +35,74 @@
},
"query": "SELECT app.id FROM app\n WHERE app.path = $1 AND app.workspace_id = $2"
},
"023fffd0a042a28b5be991169a506aff92f64f84e49b4c041cd369b045c31e73": {
"describe": {
"columns": [
{
"name": "id",
"ordinal": 0,
"type_info": "Int8"
},
{
"name": "path",
"ordinal": 1,
"type_info": "Varchar"
},
{
"name": "summary",
"ordinal": 2,
"type_info": "Varchar"
},
{
"name": "versions",
"ordinal": 3,
"type_info": "Int8Array"
},
{
"name": "policy",
"ordinal": 4,
"type_info": "Jsonb"
},
{
"name": "extra_perms",
"ordinal": 5,
"type_info": "Jsonb"
},
{
"name": "value",
"ordinal": 6,
"type_info": "Jsonb"
},
{
"name": "created_at",
"ordinal": 7,
"type_info": "Timestamptz"
},
{
"name": "created_by",
"ordinal": 8,
"type_info": "Varchar"
}
],
"nullable": [
false,
false,
false,
false,
false,
false,
false,
false,
false
],
"parameters": {
"Left": [
"Text"
]
}
},
"query": "SELECT app.id, app.path, app.summary, app.versions, app.policy,\n app.extra_perms, app_version.value, \n app_version.created_at, app_version.created_by from app, app_version \n WHERE app.workspace_id = $1 AND app_version.id = app.versions[array_upper(app.versions, 1)]"
},
"0355b53b1d45955ca56b2829372ce9c656d7f0ad7b8d0709161047f0d8cdc4f4": {
"describe": {
"columns": [
+8
View File
@@ -1092,6 +1092,10 @@ paths:
- variable
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- name: already_encrypted
in: query
schema:
type: boolean
requestBody:
description: new variable
required: true
@@ -1133,6 +1137,10 @@ paths:
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/Path"
- name: already_encrypted
in: query
schema:
type: boolean
requestBody:
description: updated variable
required: true
+9 -2
View File
@@ -209,12 +209,13 @@ async fn create_variable(
Extension(user_db): Extension<UserDB>,
Extension(webhook): Extension<WebhookShared>,
Path(w_id): Path<String>,
Query(AlreadyEncrypted { already_encrypted }): Query<AlreadyEncrypted>,
Json(variable): Json<CreateVariable>,
) -> Result<(StatusCode, String)> {
let mut tx = user_db.begin(&authed).await?;
check_path_conflict(&mut tx, &w_id, &variable.path).await?;
let value = if variable.is_secret {
let value = if variable.is_secret && !already_encrypted.unwrap_or(false) {
let mc = build_crypt(&mut tx, &w_id).await?;
encrypt(&mc, &variable.value)
} else {
@@ -312,12 +313,18 @@ struct EditVariable {
description: Option<String>,
}
#[derive(Deserialize)]
struct AlreadyEncrypted {
already_encrypted: Option<bool>,
}
async fn update_variable(
authed: Authed,
Extension(user_db): Extension<UserDB>,
Extension(webhook): Extension<WebhookShared>,
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
Query(AlreadyEncrypted { already_encrypted }): Query<AlreadyEncrypted>,
Json(ns): Json<EditVariable>,
) -> Result<String> {
use sql_builder::prelude::*;
@@ -343,7 +350,7 @@ async fn update_variable(
.await?
.unwrap_or(false);
let value = if is_secret {
let value = if is_secret && !already_encrypted.unwrap_or(false) {
let mc = build_crypt(&mut tx, &w_id).await?;
encrypt(&mc, &nvalue)
} else {
+50 -6
View File
@@ -9,6 +9,7 @@
use std::str::FromStr;
use crate::{
apps::AppWithLastVersion,
db::{UserDB, DB},
folders::Folder,
resources::{Resource, ResourceType},
@@ -24,6 +25,7 @@ use axum::{
routing::{delete, get, post},
Json, Router,
};
use serde_json::to_string_pretty;
use stripe::CustomerId;
use windmill_audit::{audit_log, ActionKind};
use windmill_common::{
@@ -1098,6 +1100,28 @@ struct ArchiveQueryParams {
archive_type: Option<String>,
}
#[inline]
pub fn to_string_without_metadata<T>(value: &T) -> Result<String>
where
T: ?Sized + Serialize,
{
let value = serde_json::to_value(value).map_err(to_anyhow)?;
value
.as_object()
.map(|obj| {
let mut obj = obj.clone();
for key in ["workspace_id", "path", "name"] {
if obj.contains_key(key) {
obj.remove(key);
}
}
serde_json::to_string_pretty(&obj).ok()
})
.flatten()
.ok_or_else(|| Error::BadRequest("Impossible to serialize value".to_string()))
}
async fn tarball_workspace(
authed: Authed,
Extension(db): Extension<DB>,
@@ -1129,7 +1153,7 @@ async fn tarball_workspace(
for folder in folders {
archive
.write_to_archive(
&serde_json::to_string_pretty(&folder).unwrap(),
&to_string_without_metadata(&folder).unwrap(),
&format!("f/{}/folder.meta.json", folder.name),
)
.await?;
@@ -1187,7 +1211,7 @@ async fn tarball_workspace(
.await?;
for resource in resources {
let resource_str = serde_json::to_string_pretty(&resource).unwrap();
let resource_str = &to_string_without_metadata(&resource).unwrap();
archive
.write_to_archive(&resource_str, &format!("{}.resource.json", resource.path))
.await?;
@@ -1204,7 +1228,7 @@ async fn tarball_workspace(
.await?;
for resource_type in resource_types {
let resource_str = serde_json::to_string_pretty(&resource_type).unwrap();
let resource_str = &to_string_without_metadata(&resource_type).unwrap();
archive
.write_to_archive(
&resource_str,
@@ -1223,7 +1247,7 @@ async fn tarball_workspace(
.await?;
for flow in flows {
let flow_str = serde_json::to_string_pretty(&flow).unwrap();
let flow_str = &to_string_without_metadata(&flow).unwrap();
archive
.write_to_archive(&flow_str, &format!("{}.flow.json", flow.path))
.await?;
@@ -1232,19 +1256,39 @@ async fn tarball_workspace(
{
let variables = sqlx::query_as::<_, ExportableListableVariable>(
"SELECT *, false as is_expired FROM variable WHERE workspace_id = $1 AND is_secret = false",
"SELECT *, false as is_expired FROM variable WHERE workspace_id = $1",
)
.bind(&w_id)
.fetch_all(&db)
.await?;
for var in variables {
let flow_str = serde_json::to_string_pretty(&var).unwrap();
let flow_str = &to_string_without_metadata(&var).unwrap();
archive
.write_to_archive(&flow_str, &format!("{}.variable.json", var.path))
.await?;
}
}
{
let apps = sqlx::query_as!(
AppWithLastVersion,
"SELECT app.id, app.path, app.summary, app.versions, app.policy,
app.extra_perms, app_version.value,
app_version.created_at, app_version.created_by from app, app_version
WHERE app.workspace_id = $1 AND app_version.id = app.versions[array_upper(app.versions, 1)]",
&w_id
)
.fetch_all(&db)
.await?;
for app in apps {
let flow_str = &to_string_pretty(&app).unwrap();
archive
.write_to_archive(&flow_str, &format!("{}.app.json", app.path))
.await?;
}
}
archive.finish().await?;
let file = tokio::fs::File::open(file_path).await?;
+105
View File
@@ -0,0 +1,105 @@
import { Any, model, property } from "./decoverto.ts";
import {
AppService,
AppWithLastVersion,
colors,
microdiff,
Policy,
} from "./deps.ts";
import { Difference, PushDiffs, Resource, setValueByPath } from "./types.ts";
@model()
export class AppFile implements Resource, PushDiffs {
@property(Any)
value: any;
@property(() => String)
summary: string;
@property(Any)
policy: Policy;
constructor(value: string, summary: string, policy: Policy) {
this.value = value;
this.summary = summary;
this.policy = policy;
}
async pushDiffs(
workspace: string,
remotePath: string,
diffs: Difference[],
): Promise<void> {
if (await AppService.existsApp({ workspace, path: remotePath })) {
console.log(
colors.bold.yellow(
`Applying ${diffs.length} diffs to existing app...`,
),
);
const changeset: {
path?: string | undefined;
summary?: string | undefined;
value?: any;
policy?: Policy | undefined;
} = {};
for (const diff of diffs) {
if (
diff.type !== "REMOVE" &&
(
diff.path[0] !== "value" && diff.path[0] !== "policy" && (
diff.path.length !== 1 ||
!["path", "summary"].includes(
diff.path[0] as string,
)
)
)
) {
throw new Error("Invalid app diff with path " + diff.path);
}
if (diff.type === "CREATE" || diff.type === "CHANGE") {
setValueByPath(changeset, diff.path, diff.value);
} else if (diff.type === "REMOVE") {
setValueByPath(changeset, diff.path, null);
}
}
const hasChanges = Object.values(changeset).some((v) =>
v !== null && typeof v !== "undefined"
);
if (!hasChanges) {
console.log(colors.yellow("! Skipping empty changeset"));
return;
}
await AppService.updateApp({
workspace,
path: remotePath,
requestBody: changeset,
});
} else {
console.log(colors.yellow("Creating new app..."));
await AppService.createApp({
workspace,
requestBody: {
path: remotePath,
policy: this.policy,
summary: this.summary,
value: this.value,
},
});
}
}
async push(workspace: string, remotePath: string): Promise<void> {
let existing: AppWithLastVersion | undefined;
try {
existing = await AppService.getAppByPath({
workspace: workspace,
path: remotePath,
});
} catch {
existing = undefined;
}
await this.pushDiffs(
workspace,
remotePath,
microdiff(existing ?? {}, this, { cyclesFix: false }),
);
}
}
+2 -1
View File
@@ -16,7 +16,6 @@ export {
} from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/mod.ts";
// std
export { Untar } from "https://deno.land/std@0.176.0/archive/untar.ts";
export * as path from "https://deno.land/std@0.176.0/path/mod.ts";
export { ensureDir } from "https://deno.land/std@0.176.0/fs/ensure_dir.ts";
export {
@@ -38,3 +37,5 @@ export {
default as microdiff,
} from "https://deno.land/x/microdiff@v1.3.1/index.ts";
export { default as objectHash } from "https://deno.land/x/object_hash@2.0.3.1/mod.ts";
export { default as gitignore_parser } from "npm:gitignore-parser";
export { default as JSZip } from "npm:jszip@3.7.1";
+5 -1
View File
@@ -76,6 +76,7 @@ export class FlowFile implements Resource, PushDiffs {
description: this.description, // This is OpenAPIed as optional, but isn't
schema: this.schema, // Same
};
const base_changeset = { ...changeset };
for (const diff of diffs) {
if (
diff.type !== "REMOVE" &&
@@ -107,7 +108,10 @@ export class FlowFile implements Resource, PushDiffs {
await FlowService.updateFlow({
workspace: workspace,
path: remotePath,
requestBody: changeset,
requestBody: {
...changeset,
...base_changeset,
},
});
} else {
console.log(colors.bold.yellow("Creating new flow..."));
+3 -2
View File
@@ -105,12 +105,13 @@ export class FolderFile implements Resource, PushDiffs {
requestBody: changeset,
});
} else {
console.log(colors.bold.yellow("Creating new folder..."));
console.log(colors.bold.yellow("Creating new folder: " + remotePath));
console.log(this.owners, this.extra_perms)
await FolderService.createFolder({
workspace: workspace,
requestBody: {
name: remotePath,
extra_perms: this.extra_perms,
extra_perms: Object.fromEntries(this.extra_perms?.entries() ?? []),
owners: this.owners,
},
});
+1
View File
@@ -18,6 +18,7 @@ const VERSION = "v1.64.0";
let command: any = new Command()
.name("wmill")
.description("A simple CLI tool for windmill.")
.action(() => command.showHelp())
.globalOption(
"--workspace <workspace:string>",
"Specify the target workspace. This overrides the default workspace.",
+10 -18
View File
@@ -1,43 +1,35 @@
// deno-lint-ignore-file no-explicit-any
import { GlobalOptions } from "./types.ts";
import { colors, Command, readerFromStreamReader, Untar } from "./deps.ts";
import { colors, Command, JSZip } from "./deps.ts";
import { Workspace } from "./workspace.ts";
export async function downloadTar(
export async function downloadZip(
workspace: Workspace,
): Promise<Untar | undefined> {
): Promise<JSZip | undefined> {
const requestHeaders: HeadersInit = new Headers();
requestHeaders.set("Authorization", "Bearer " + workspace.token);
requestHeaders.set("Content-Type", "application/octet-stream");
const tarResponse = await fetch(
const zipResponse = await fetch(
workspace.remote + "api/w/" + workspace.workspaceId +
"/workspaces/tarball",
"/workspaces/tarball?archive_type=zip",
{
headers: requestHeaders,
method: "GET",
},
);
if (!tarResponse.ok) {
if (!zipResponse.ok) {
console.log(
colors.red(
"Failed to request tarball from API " + tarResponse.statusText,
"Failed to request tarball from API " + zipResponse.statusText,
),
);
console.log(await tarResponse.text());
console.log(await zipResponse.text());
return undefined;
}
const streamReader = tarResponse.body?.getReader();
if (!streamReader) {
console.log(colors.red("Failed to read tar request body"));
return undefined;
}
console.log(colors.yellow("Streaming tarball to disk..."));
const denoReader = readerFromStreamReader(streamReader);
const untar = new Untar(denoReader);
return untar;
const blob = await zipResponse.blob();
return await JSZip.loadAsync(blob);
}
async function stub(
+1 -1
View File
@@ -60,7 +60,7 @@ export class ResourceTypeFile implements ResourceI, PushDiffs {
) {
console.log(
"Resource type " + remotePath +
" is already taken for the current workspace, but cannot be updated. Is this a conflict with starter?",
" is already taken for the current workspace, but cannot be updated. Is this a conflict with starter?",
);
return;
}
+300 -190
View File
@@ -6,9 +6,9 @@ import {
colors,
Command,
Confirm,
copy,
ensureDir,
iterateReader,
gitignore_parser,
JSZip,
microdiff,
nanoid,
objectHash,
@@ -22,7 +22,7 @@ import {
inferTypeFromPath,
setValueByPath,
} from "./types.ts";
import { downloadTar } from "./pull.ts";
import { downloadZip } from "./pull.ts";
import { FolderFile } from "./folder.ts";
import { ResourceTypeFile } from "./resource-type.ts";
import {
@@ -41,15 +41,13 @@ const TrackedId = String;
const CONTENT_ENCODER: cbor.Encoder = new cbor.Encoder({ pack: true });
export class Tracked {
#id: TrackedId;
#parent: State;
path: string;
#content_cache?: string;
constructor(id: TrackedId, parent: State, path: string) {
this.#id = id;
constructor(parent: State, path: string) {
this.#parent = parent;
this.path = path;
}
@@ -59,29 +57,146 @@ export class Tracked {
return this.#content_cache;
}
if (!this.#parent.stateRoot) {
throw new Error("Parent uninitialized");
}
const f = this.#parent.contentFile(this.#id);
const f = this.#parent.stateContentFor(this.path);
if (!f) {
return undefined;
}
const data = await Deno.readFile(
path.join(this.#parent.stateRoot, ".wmill", f),
State.getInternalWmillFolder(f),
);
const content = CONTENT_ENCODER.decode(data);
this.#content_cache = content;
return content;
}
getHash(): string {
return this.#parent.hashes.get(this.#id)!;
getHash(): string | undefined {
return this.#parent.hashes.get(objectHash(this.path));
}
}
getId(): TrackedId {
return this.#id;
type DynFSElement = {
isDirectory: boolean;
path: string;
getContentBytes(): Promise<Uint8Array>;
getContentText(): Promise<string>;
getChildren(): AsyncIterable<DynFSElement>;
};
async function FSFSElement(p: string): Promise<DynFSElement> {
function _internal_element(
p: string,
isDir: boolean,
): DynFSElement {
return {
isDirectory: isDir,
path: p,
async *getChildren(): AsyncIterable<DynFSElement> {
for await (const e of Deno.readDir(p)) {
yield _internal_element(path.join(p, e.name), e.isDirectory);
}
},
async getContentBytes(): Promise<Uint8Array> {
return await Deno.readFile(p);
},
async getContentText(): Promise<string> {
return await Deno.readTextFile(p);
},
};
}
return _internal_element(p, (await Deno.stat(p)).isDirectory);
}
function ZipFSElement(zip: JSZip): DynFSElement {
function _internal_file(p: string, f: JSZip.JSZipObject): DynFSElement {
return {
isDirectory: false,
path: p,
// deno-lint-ignore require-yield
async *getChildren(): AsyncIterable<DynFSElement> {
throw new Error("Cannot get children of file");
},
async getContentBytes(): Promise<Uint8Array> {
return await f.async("uint8array");
},
async getContentText(): Promise<string> {
return await f.async("text");
},
};
}
function _internal_folder(p: string, zip: JSZip): DynFSElement {
return {
isDirectory: true,
path: p,
async *getChildren(): AsyncIterable<DynFSElement> {
for (const filename in zip.files) {
const file = zip.files[filename];
const totalPath = path.join(p, filename);
if (file.dir) {
const e = zip.folder(file.name)!;
yield _internal_folder(totalPath, e);
} else {
yield _internal_file(totalPath, file);
}
}
},
async getContentBytes(): Promise<Uint8Array> {
throw new Error("Cannot get content of folder");
},
async getContentText(): Promise<string> {
throw new Error("Cannot get content of folder");
},
};
}
return _internal_folder("./", zip);
}
async function* readDirRecursiveWithIgnore(
ignore: (path: string) => boolean,
root: DynFSElement,
): AsyncGenerator<
{
path: string;
ignored: boolean;
isDirectory: boolean;
getContentBytes(): Promise<Uint8Array>;
getContentText(): Promise<string>;
}
> {
const stack: {
path: string;
isDirectory: boolean;
ignored: boolean;
c(): AsyncIterable<DynFSElement>;
getContentBytes(): Promise<Uint8Array>;
getContentText(): Promise<string>;
}[] = [{
path: root.path,
ignored: ignore(root.path),
isDirectory: root.isDirectory,
c: root.getChildren,
getContentBytes(): Promise<Uint8Array> {
throw undefined;
},
getContentText(): Promise<string> {
throw undefined;
},
}];
while (stack.length > 0) {
const e = stack.pop()!;
yield e;
if (!e.isDirectory) continue;
for await (const e2 of e.c()) {
stack.push({
path: e2.path,
ignored: e.ignored || ignore(e2.path),
isDirectory: e2.isDirectory,
getContentBytes: e2.getContentBytes,
getContentText: e2.getContentText,
c: e2.getChildren,
});
}
}
}
@@ -97,85 +212,97 @@ export class State {
@property(map(() => TrackedId, () => String, { shape: MapShape.Object }))
contentFiles: Map<TrackedId, string>;
@property(map(() => String, () => TrackedId, { shape: MapShape.Object }))
tracked: Map<string, TrackedId>;
@property(() => String)
workspaceId: string;
@property(() => String)
remoteUrl: string;
stateRoot?: string;
constructor(
hashes: Map<TrackedId, string>,
tracked: Map<string, TrackedId>,
contentFiles: Map<TrackedId, string>,
workspaceId: string,
remoteUrl: string,
) {
this.hashes = hashes;
this.tracked = tracked;
this.contentFiles = contentFiles;
this.workspaceId = workspaceId;
this.remoteUrl = remoteUrl;
}
add(path: string) {
if (this.tracked.get(path)) {
throw new Error("Cannot newly track already tracked paths");
} else {
this.tracked.set(path, nanoid());
}
}
public forget(path: string) {
const id = this.tracked.get(path);
if (id) {
this.tracked.delete(path);
this.hashes.delete(id);
}
}
public contentFile(id: TrackedId): string | undefined {
return this.contentFiles.get(id);
public stateContentFor(path: string) {
const hashOfPath = objectHash(path);
const file = this.contentFiles.get(hashOfPath);
return file;
}
public get(path: string): Tracked {
const id = this.tracked.get(path);
if (!id) {
throw new Error("Could not resolve path " + path);
// TODO: Normalize path
return new Tracked(this, path);
}
public static getInternalWmillFolder(...subpath: string[]): string {
return path.join(Deno.cwd(), ".wmill", ...subpath);
}
public async *getFiles(): AsyncGenerator<{
localFile: string | undefined;
stateFile: Tracked;
isIgnored: boolean;
path: string;
}> {
// not sure why the auto-typing doesn't work here, see <https://www.npmjs.com/package/gitignore-parser> (a @types package is available...)
const ignore: {
accepts(file: string): boolean;
denies(file: string): boolean;
} = gitignore_parser.compile(
await Deno.readTextFile(".wmillignore"),
);
const base = Deno.cwd();
for await (
const { ignored, path, getContentText, isDirectory }
of readDirRecursiveWithIgnore(
ignore.denies,
await FSFSElement(base),
)
) {
if (isDirectory || path.includes(".wmill")) {
continue;
}
const path2 = path.substring(base.length + 1);
const stateFile = this.get(path2);
let localFile: string | undefined;
try {
localFile = await getContentText();
} catch {
localFile = undefined;
}
yield { path: path2, isIgnored: ignored, stateFile, localFile };
}
return new Tracked(id, this, path);
}
public async save(): Promise<void> {
if (!this.stateRoot) {
throw new Error("Uninitialized state root");
}
const encoder = new cbor.Encoder({});
const plain = decoverto.type(State).instanceToPlain(this);
const result: Uint8Array = encoder.encode(plain, {});
await Deno.writeFile(path.join(this.stateRoot, ".wmill", "main"), result, {
await Deno.writeFile(State.getInternalWmillFolder("main"), result, {
create: true,
});
}
public static async loadState(dir: string): Promise<State> {
const source = await Deno.readFile(path.join(dir, ".wmill", "main"));
public static async loadState(): Promise<State> {
const source = await Deno.readFile(this.getInternalWmillFolder("main"));
const encoder = new cbor.Encoder({});
const raw = encoder.decode(source);
const state = decoverto.type(State).plainToInstance(
raw,
);
state.stateRoot = dir;
return Object.freeze(state);
}
}
async function getState(opts: GlobalOptions) {
const existingState = await State.loadState(Deno.cwd());
const existingState = await State.loadState();
const workspaceStream = await getWorkspaceStream();
const reader = workspaceStream.getReader();
@@ -203,50 +330,63 @@ async function updateStateFromRemote(
state: State,
callback: (filename: string) => PromiseLike<boolean> | boolean,
) {
const untar = await downloadTar(workspace);
if (!untar) throw new Error("Failed to pull Tar");
const zipDir = await downloadZip(workspace);
if (!zipDir) throw new Error("Failed to pull Zip");
const decoder = new TextDecoder();
for await (const entry of untar) {
const id = state.tracked.get(entry.fileName);
if (id) {
let val = "";
for await (const e of iterateReader(entry)) {
const tmp = decoder.decode(e);
val += tmp;
}
if (entry.fileName.endsWith(".json")) {
for await (
const entry of readDirRecursiveWithIgnore(
(_) => false,
ZipFSElement(zipDir),
)
) {
if (entry.isDirectory || entry.ignored) continue;
const e = state.get(entry.path);
if (e) {
const val = await entry.getContentText();
if (entry.path.endsWith(".json")) {
const parsed = JSON.parse(val);
const typed = inferTypeFromPath(entry.fileName, parsed);
const typed = inferTypeFromPath(entry.path, parsed);
const oldHash = state.hashes.get(id);
const oldHash = e.getHash();
const newHash = objectHash(typed);
if (!oldHash || oldHash !== newHash) {
if (!await callback(entry.fileName)) {
if (!await callback(entry.path)) {
return; // notice that we are not saving
}
state.hashes.set(id, newHash);
try {
await Deno.stat(entry.path);
} catch {
await ensureDir(path.dirname(e.path));
await Deno.writeTextFile(e.path, "{}");
}
state.hashes.set(objectHash(e.path), newHash);
const encoded = CONTENT_ENCODER.encode(typed);
const fileName = nanoid();
await Deno.writeFile(
path.join(state.stateRoot!, ".wmill", fileName),
State.getInternalWmillFolder(fileName),
encoded,
{ create: true },
);
state.contentFiles.set(id, fileName);
state.contentFiles.set(objectHash(e.path), fileName);
}
} else {
const fileName = nanoid();
await Deno.writeTextFile(
path.join(state.stateRoot!, ".wmill", fileName),
State.getInternalWmillFolder(fileName),
val,
{ create: true },
);
state.contentFiles.set(id, fileName);
state.contentFiles.set(objectHash(e.path), fileName);
try {
await Deno.stat(entry.path);
} catch {
await ensureDir(path.dirname(e.path));
await Deno.writeTextFile(e.path, "tmp_file");
}
}
}
}
@@ -262,11 +402,23 @@ async function pull(
opts2.override = opts2.rawOverride;
opts2.raw = undefined;
opts2.rawOverride = undefined;
await pullRaw(opts2, Deno.cwd());
await pullRaw(opts2);
return;
}
try {
await Deno.stat(State.getInternalWmillFolder("main"))
} catch {
console.log(
colors.yellow(
"No state file found, creating empty state folder at .wmill",
),
);
init()
}
const state = await getState(opts);
console.log(state);
const workspace = await resolveWorkspace(opts);
await requireLogin(opts);
@@ -282,10 +434,11 @@ async function pull(
for await (const diff of diffs) {
await applyDiff(
diff.diff,
path.join(state.stateRoot!, diff.localPath),
path.join(Deno.cwd(), diff.localPath),
);
}
console.log(colors.green.underline("Done! All changes applied."));
console.log(state);
async function applyDiff(diffs: Difference[], file: string) {
ensureDir(path.dirname(file));
@@ -315,18 +468,20 @@ async function pull(
}
async function copyNonJsonFiles(state: State): Promise<void> {
for (const t of state.tracked.keys()) {
if (t.endsWith(".json")) {
for await (
const { path, localFile, stateFile, isIgnored } of state.getFiles()
) {
if (path.endsWith(".json") || isIgnored || !localFile) {
continue;
}
const entry = state.get(t);
const target = await Deno.open(entry.path, { create: true, write: true });
const target = await Deno.open(stateFile.path, {
create: true,
write: true,
});
const source = await Deno.open(
path.join(
state.stateRoot!,
".wmill",
state.contentFile(entry.getId())!,
State.getInternalWmillFolder(
state.stateContentFor(stateFile.path)!,
),
{
read: true,
@@ -339,24 +494,18 @@ async function pull(
}
async function* diffState(state: State): AsyncGenerator<StateDiff, void, void> {
for (const t of state.tracked.keys()) {
if (!t.endsWith(".json")) {
for await (
const { path, localFile, stateFile, isIgnored } of state.getFiles()
) {
if (!path.endsWith(".json") || isIgnored) {
continue;
}
const entry = state.get(t);
let fileText;
try {
fileText = await Deno.readTextFile(
path.join(state.stateRoot!, entry.path),
);
} catch {
fileText = "{}";
}
const fileText = localFile ?? "{}";
const old = JSON.parse(fileText);
const fileHash = objectHash(old);
if (fileHash !== entry.getHash()) {
const stateContent = await entry.getContent() as any;
if (fileHash !== stateFile.getHash()) {
const stateContent = await stateFile.getContent() as any;
const diff = microdiff(
old,
@@ -364,7 +513,7 @@ async function* diffState(state: State): AsyncGenerator<StateDiff, void, void> {
{ cyclesFix: false },
);
yield new StateDiff(entry.getId(), entry.path, diff);
yield new StateDiff(stateFile.path, diff);
}
}
}
@@ -414,7 +563,7 @@ async function push(opts: GlobalOptions & { raw: boolean }) {
let fileJSON;
try {
fileJSON = JSON.parse(
await Deno.readTextFile(path.join(state.stateRoot!, e.path)),
await Deno.readTextFile(path.join(Deno.cwd(), e.path)),
);
} catch {
fileJSON = {};
@@ -444,32 +593,27 @@ async function push(opts: GlobalOptions & { raw: boolean }) {
return;
}
for (const p of state.tracked.keys()) {
if (!p.endsWith(".json")) continue;
const entry = state.get(p);
let fileJSON;
try {
fileJSON = JSON.parse(
await Deno.readTextFile(path.join(state.stateRoot!, entry.path)),
);
} catch {
fileJSON = {};
}
const file = inferTypeFromPath(entry.path, fileJSON);
const eContent =
(inferTypeFromPath(entry.path, await entry.getContent())) ?? {};
for await (
const { path, isIgnored, stateFile, localFile } of state.getFiles()
) {
if (!path.endsWith(".json") || isIgnored) continue;
const fileJSON = JSON.parse(localFile ?? "{}");
const file = inferTypeFromPath(path, fileJSON);
const eContent = (inferTypeFromPath(path, await stateFile.getContent())) ??
{};
const fileHash = objectHash(file);
const eHash = objectHash(eContent);
if (fileHash !== eHash) {
const remotePath = entry.path.split(".")[0];
const type = getTypeStrFromPath(entry.path);
const remotePath = stateFile.path.split(".")[0];
const type = getTypeStrFromPath(stateFile.path);
if (type === "script") {
// Diffing makes no sense for scripts - instead fetch parent hash & check hash again.
// If hash is still missmatched - create new script as child.
const typed = decoverto.type(ScriptFile).plainToInstance(file);
const contentPath = await findContentFile(entry.path);
const contentPath = await findContentFile(stateFile.path);
const language = inferContentTypeFromFilePath(contentPath);
const content = await Deno.readTextFile(contentPath);
try {
@@ -562,65 +706,21 @@ async function push(opts: GlobalOptions & { raw: boolean }) {
}
class StateDiff {
trackedId: TrackedId;
localPath: string;
diff: Difference[];
constructor(
trackedId: TrackedId,
localPath: string,
diff: Difference[],
) {
this.trackedId = trackedId;
this.localPath = localPath;
this.diff = diff;
}
}
async function add(opts: GlobalOptions, path: string) {
const state = await getState(opts);
// TODO: Automatically check whether this path exists either locally or on the remote
state.add(path);
if (path.endsWith(".script.json")) {
try {
const f = await findContentFile(path);
state.add(f);
} catch {
const workspace = await resolveWorkspace(opts);
await requireLogin(opts);
try {
const old = await ScriptService.getScriptByPath({
workspace: workspace.workspaceId,
path: path.split(".")[0],
});
if (old.language === "python3") {
state.add(path.replace(".script.json", ".py"));
} else if (old.language === "bash") {
state.add(path.replace(".script.json", ".sh"));
} else if (old.language === "deno") {
state.add(path.replace(".script.json", ".ts"));
} else if (old.language === "go") {
state.add(path.replace(".script.json", ".go"));
} else {
throw new Error("Remote returned invalid language?! " + old.language);
}
} catch {
throw new Error(
"Could not infer script language from local or remote. Exiting.",
);
}
}
}
await state.save();
}
async function init(opts: GlobalOptions) {
const root = Deno.cwd();
try {
await Deno.mkdir(path.join(root, ".wmill"));
await Deno.mkdir(State.getInternalWmillFolder());
} catch {
console.log(
colors.red(
@@ -630,33 +730,41 @@ async function init(opts: GlobalOptions) {
return;
}
await Deno.writeTextFile(
".wmillignore",
"# Write any ignores here as you see fit. Same syntax as .gitignore.",
);
const workspace = await resolveWorkspace(opts);
const newState = new State(
new Map(),
new Map(),
new Map(),
workspace.workspaceId,
workspace.remote,
);
newState.stateRoot = root;
await newState.save();
}
async function pullRaw(
opts: GlobalOptions & { override: boolean },
dir: string,
) {
const workspace = await resolveWorkspace(opts);
const untar = await downloadTar(workspace);
if (!untar) return;
const zipDir = await downloadZip(workspace);
if (!zipDir) return;
for await (const entry of untar) {
console.log(entry.fileName);
const filePath = path.resolve(dir, entry.fileName);
if (entry.type === "directory") {
// TODO use ZipFSElement & readDirRecursiveWithIgnore here
// TODO also remember to read content via entry methods now instead of direct I/O
for await (
const entry of readDirRecursiveWithIgnore(
(_) => false,
ZipFSElement(zipDir),
)
) {
const filePath = entry.path;
if (entry.isDirectory) {
await ensureDir(filePath);
continue;
}
@@ -670,21 +778,22 @@ async function pullRaw(
exists = false;
}
if (exists) {
if (await Deno.readTextFileSync(filePath) === await entry.getContentText()) {
continue;
}
if (
!(await Confirm.prompt(
"Conflict at " +
filePath +
" do you want to override the local version?",
filePath +
" do you want to override the local version?",
))
) {
continue;
}
}
}
const file = await Deno.open(filePath, { write: true, create: true });
const len = await copy(entry, file);
await file.truncate(len);
file.close();
console.log("Writing " + filePath)
await Deno.writeFile(filePath, await entry.getContentBytes());
}
console.log(colors.green("Done. Wrote all files to disk."));
}
@@ -738,7 +847,7 @@ async function pushRawFindCandidateFiles(
namespaceName,
path: path + "folder.meta.json",
});
} catch {}
} catch { }
}
while (stack.length > 0) {
@@ -751,8 +860,8 @@ async function pushRawFindCandidateFiles(
namespaceKind: e.name == "g"
? "group"
: e.name == "u"
? "user"
: "folder",
? "user"
: "folder",
namespaceName: namespaceName,
});
} else {
@@ -804,7 +913,7 @@ async function pushRaw(opts: GlobalOptions, dir?: string) {
console.log(
colors.blue(
"Found " + (normal.length + resourceTypes.length + folders.length) +
" candidates",
" candidates",
),
);
for (const resourceType of resourceTypes) {
@@ -883,8 +992,8 @@ async function pushRaw(opts: GlobalOptions, dir?: string) {
console.log(
colors.yellow(
"Found resource type file at " +
candidate.path +
" this appears to be inside a path folder. Resource types are not addressed by path. Place them at the root or inside only an organizational folder. Ignoring this file!",
candidate.path +
" this appears to be inside a path folder. Resource types are not addressed by path. Place them at the root or inside only an organizational folder. Ignoring this file!",
),
);
continue;
@@ -938,7 +1047,7 @@ async function pushRaw(opts: GlobalOptions, dir?: string) {
remotePath,
);
} else {
typed.push(workspace.workspaceId, remotePath);
await typed.push(workspace.workspaceId, remotePath);
}
}
console.log(colors.underline.bold.green("Successfully Pushed all files."));
@@ -947,22 +1056,23 @@ async function pushRaw(opts: GlobalOptions, dir?: string) {
const command = new Command()
.command("init")
.description(
"Initialize this folder as sync root for the currently selected workspace & remote.",
"Initialize this folder as sync root for the currently selected workspace & remote." +
"\nBegin by initializing state tracking using `init` & add files you want to track using `add`. `push` & `pull` will then use local state to accurately track changes required on the remote.",
)
.action(init as any)
.command("add")
.description("Add a local file for tracking")
.arguments("<path:string>")
.action(add as any)
.command("pull")
.description("Pull any remote changes and apply them locally")
.description(
"Pull any remote changes and apply them locally. Use --raw for usage without local state tracking.",
)
.option("--raw", "Pull without using state.")
.option("--raw-override", "Always override local files with remote.", {
depends: ["raw"],
})
.action(pull as any)
.command("push")
.description("Push any local changes and apply them remotely")
.description(
"Push any local changes and apply them remotely. Use --raw for usage without local state tracking.",
)
.option("--raw", "Push without using state.")
.action(push as any);
+14 -3
View File
@@ -6,6 +6,7 @@ import { ScriptFile } from "./script.ts";
import { VariableFile } from "./variable.ts";
import { path } from "./deps.ts";
import { FolderFile } from "./folder.ts";
import { AppFile } from "./apps.ts";
// TODO: Remove this & replace with a "pull" that lets the object either pull the remote version or return undefined.
// Then combine those with diffing, which then gives the new push impl
@@ -90,7 +91,8 @@ export function inferTypeFromPath(
| FlowFile
| ResourceFile
| ResourceTypeFile
| FolderFile {
| FolderFile
| AppFile {
const typeEnding = getTypeStrFromPath(p);
if (typeEnding === "folder") {
@@ -105,6 +107,8 @@ export function inferTypeFromPath(
return decoverto.type(ResourceFile).plainToInstance(obj);
} else if (typeEnding === "resource-type") {
return decoverto.type(ResourceTypeFile).plainToInstance(obj);
} else if (typeEnding === "app") {
return decoverto.type(AppFile).plainToInstance(obj);
} else {
throw new Error("infer type unreachable");
}
@@ -112,7 +116,14 @@ export function inferTypeFromPath(
export function getTypeStrFromPath(
p: string,
): "script" | "variable" | "flow" | "resource" | "resource-type" | "folder" {
):
| "script"
| "variable"
| "flow"
| "resource"
| "resource-type"
| "folder"
| "app" {
const parsed = path.parse(p);
if (parsed.ext !== ".json") {
throw new Error(
@@ -128,7 +139,7 @@ export function getTypeStrFromPath(
if (
typeEnding === "script" || typeEnding === "variable" ||
typeEnding === "flow" || typeEnding === "resource" ||
typeEnding === "resource-type"
typeEnding === "resource-type" || typeEnding === "app"
) {
return typeEnding;
} else {
+8 -2
View File
@@ -191,7 +191,7 @@ export async function add(
}
if (!workspaceId) {
workspaceId = await Input.prompt("Enter the ID of this workspace");
workspaceId = await Input.prompt({ message: "Enter the ID of this workspace", default: workspaceName, suggestions: [workspaceName] });
}
if (!remote) {
@@ -203,7 +203,13 @@ export async function add(
remote = url.toString();
} catch {
// not a url
remote = new URL(await Input.prompt("Enter the Remote URL")).toString();
remote = new URL(
await Input.prompt({
message: "Enter the Remote URL",
suggestions: ["https://app.windmill.dev/"],
default: "https://app.windmill.dev/"
}),
).toString();
}
}
remote = new URL(remote).toString(); // add trailing slash in all cases!
@@ -168,7 +168,8 @@
path: name,
kind: 'folder',
requestBody: {
owner: owner_name
owner: owner_name,
write: true
}
})
} else if (role == 'writer') {